From 4444e01132200ed9918717410cc5770ad3dff661 Mon Sep 17 00:00:00 2001 From: Dawn Date: Sat, 15 Aug 2026 22:40:06 -0400 Subject: [PATCH] refactor(desktop): split native_relay_client tests into a sibling file The lifecycle tests added for the retry-eviction fix pushed `native_relay_client.rs` to 1292 lines, over the desktop file-size ratchet's 1000-line limit for new files. Moves `closed_recovery_tests` into `native_relay_client_tests.rs` behind `#[cfg(test)] #[path = ...]`, the same convention `archive/sync.rs` uses for `sync_tests.rs`. Source drops to 685 lines and the tests to 611. Production code is byte-identical and the test module moved verbatim apart from dedenting; a `#[path]` module is still a child module, so `use super::*` continues to reach the private items the tests drive and no visibility changed. The mutant matrix was re-run at these bytes rather than inherited: all six mutants stay killed with the same attribution. Co-authored-by: Tyler Longwell Signed-off-by: Tyler Longwell --- desktop/src-tauri/src/native_relay_client.rs | 609 +---------------- .../src/native_relay_client_tests.rs | 611 ++++++++++++++++++ 2 files changed, 613 insertions(+), 607 deletions(-) create mode 100644 desktop/src-tauri/src/native_relay_client_tests.rs diff --git a/desktop/src-tauri/src/native_relay_client.rs b/desktop/src-tauri/src/native_relay_client.rs index 3c602d88d..4f36801c5 100644 --- a/desktop/src-tauri/src/native_relay_client.rs +++ b/desktop/src-tauri/src/native_relay_client.rs @@ -585,613 +585,8 @@ fn is_read_timeout(error: &buzz_ws_client_pkg::WsClientError) -> bool { } #[cfg(test)] -mod closed_recovery_tests { - use super::*; - use futures_util::{SinkExt, StreamExt}; - use nostr::EventBuilder; - use tokio_tungstenite::tungstenite::protocol::Message; - - /// The subscription id every test below drives. - const PROBE_ID: &str = "archive:probe"; - - /// Minimal relay that completes the NIP-42 handshake, records every REQ, - /// and sends a CLOSED only when the test asks it to. - /// - /// A real socket rather than a fake `NostrWsConnection`, because the bug - /// this covers lives in the lifecycle between frames — the loop's only - /// reconcile triggers — and a fake that hands the loop a `Closed` value - /// cannot show that a REQ went back out over the wire afterwards. Same - /// `accept_async` stub shape as `native_websocket.rs`'s live-TCP tests. - /// - /// CLOSED is test-driven rather than a scripted reply to the first REQ so - /// the test can wait for the session to go quiet first. `set_subscriptions` - /// queues a wake that may still be pending when an immediate CLOSED lands, - /// and that wake reopens the subscription on its own — which made the first - /// version of this test pass against the unfixed code. - /// - /// `frames` reports REQ and CLOSE in wire order, not REQ alone: the - /// lifecycle tests below assert that a CLOSE was sent before the REQ that - /// follows it, which a REQ-only channel cannot express. - async fn stub_relay() -> (String, mpsc::Receiver, mpsc::Sender) { - let listener = tokio::net::TcpListener::bind("127.0.0.1:0") - .await - .expect("bind stub relay"); - let address = listener.local_addr().expect("stub relay address"); - let (req_tx, req_rx) = mpsc::channel(16); - let (closed_tx, mut closed_rx) = mpsc::channel::(4); - - tokio::spawn(async move { - let (stream, _) = listener.accept().await.expect("accept"); - let mut socket = tokio_tungstenite::accept_async(stream) - .await - .expect("websocket handshake"); - - socket - .send(Message::Text(r#"["AUTH","stub-challenge"]"#.into())) - .await - .expect("send challenge"); - - loop { - tokio::select! { - incoming = socket.next() => { - let Some(Ok(Message::Text(text))) = incoming else { return }; - let Ok(frame) = serde_json::from_str::(&text) else { - continue; - }; - match frame[0].as_str() { - Some("AUTH") => { - let id = frame[1]["id"].as_str().unwrap_or_default(); - socket - .send(Message::Text( - serde_json::json!(["OK", id, true, ""]).to_string().into(), - )) - .await - .expect("send auth ok"); - } - Some("REQ") => { - let id = frame[1].as_str().unwrap_or_default().to_string(); - if req_tx.send(Frame::Req(id)).await.is_err() { - return; - } - } - Some("CLOSE") => { - let id = frame[1].as_str().unwrap_or_default().to_string(); - if req_tx.send(Frame::Close(id)).await.is_err() { - return; - } - } - _ => {} - } - } - Some(command) = closed_rx.recv() => { - let frame = match command { - StubCommand::Closed(id, message) => { - serde_json::json!(["CLOSED", id, message]) - } - StubCommand::Eose(id) => serde_json::json!(["EOSE", id]), - StubCommand::Event(id, event) => { - serde_json::json!(["EVENT", id, event]) - } - }; - socket - .send(Message::Text(frame.to_string().into())) - .await - .expect("send stub frame"); - } - } - } - }); - - (format!("ws://{address}"), req_rx, closed_tx) - } - - /// A client→relay frame the stub observed, in wire order. - #[derive(Debug, PartialEq, Eq)] - enum Frame { - Req(String), - Close(String), - } - - /// A relay→client frame the test asks the stub to emit. - enum StubCommand { - Closed(String, String), - Eose(String), - Event(String, serde_json::Value), - } - - fn probe_subscription() -> Subscription { - Subscription { - id: PROBE_ID.to_string(), - filter: serde_json::json!({ "kinds": [1], "limit": 0 }), - } - } - - async fn next_frame(frames: &mut mpsc::Receiver, label: &str) -> Frame { - tokio::time::timeout(Duration::from_secs(10), frames.recv()) - .await - .unwrap_or_else(|_| panic!("timed out waiting for {label}")) - .unwrap_or_else(|| panic!("stub relay closed before {label}")) - } - - /// Waits for the next REQ, tolerating the CLOSE frames a reconcile sends - /// first. Asserting on `Frame::Req` directly would couple every test to - /// whether a particular reconcile also had cleanup to do. - async fn next_req(frames: &mut mpsc::Receiver, label: &str) -> String { - loop { - if let Frame::Req(id) = next_frame(frames, label).await { - return id; - } - } - } - - /// Waits out the wake `set_subscriptions` queued, so a CLOSED sent after - /// this cannot be reopened by anything but the CLOSED path itself. - /// - /// A pending wake is harmless while the subscription is still open — that - /// reconcile is a no-op — so draining it before the CLOSED is what makes - /// the assertion below attributable. - async fn settle() { - tokio::time::sleep(Duration::from_millis(500)).await; - } - - /// The blocker: a CLOSED with the desired set never changing again must - /// still reopen the subscription. - /// - /// Before the fix the loop removed the id from `open` and waited on a wake - /// that only `set_subscriptions` can produce, so a stable desired set left - /// the subscription dead for the life of the socket — silent permanent - /// loss for ephemeral kind 24200. - #[tokio::test] - async fn a_closed_subscription_reopens_without_a_desired_set_change() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - // Retryable class, sent once: the reopen is answered normally, so a - // failure here means "never retried" rather than "retried into another - // rejection". - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "error: temporary".into(), - )) - .await - .expect("stub relay accepts the closed command"); - - // No `set_subscriptions` between the two REQs: the reopen must come - // from the CLOSED itself, which is exactly the edge that was missing. - assert_eq!(next_req(&mut frames, "the reopened REQ").await, PROBE_ID); - - session.shutdown(); - } - - /// A relay that rejects on policy must not be re-asked in a tight loop. - #[tokio::test] - async fn a_terminal_closed_is_not_retried_on_the_same_socket() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "restricted: not authorized".into(), - )) - .await - .expect("stub relay accepts the closed command"); - - // Long enough that a retryable class (1s base) would have reopened - // several times, so this asserts suppression rather than just slowness. - let retried = tokio::time::timeout(Duration::from_secs(5), frames.recv()).await; - assert!( - retried.is_err(), - "a terminal CLOSED must not be retried on this socket, got {retried:?}" - ); - - session.shutdown(); - } - - /// M18: a subscription deleted and recreated must get a fresh REQ, even - /// though its terminal latch says never to retry. - /// - /// The latch is scoped to the subscription that earned it. Recreating the - /// id is a new subscription that happens to share a name — `archive::sync` - /// derives the id from scope and kinds, so a delete/recreate of the same - /// saved subscription produces a byte-identical id and would otherwise - /// inherit a permanent suppression for the life of the socket. - #[tokio::test] - async fn a_recreated_subscription_does_not_inherit_a_terminal_latch() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "restricted: not authorized".into(), - )) - .await - .expect("stub relay accepts the closed command"); - settle().await; - - // Delete, then recreate — each observed as its own reconcile. - session.set_subscriptions(vec![]).await; - settle().await; - session.set_subscriptions(vec![probe_subscription()]).await; - - assert_eq!( - next_req(&mut frames, "the REQ for the recreated subscription").await, - PROBE_ID, - ); - - session.shutdown(); - } - - /// M19: the same schedule, with both writes landing before the loop - /// consumes its single wake. - /// - /// This is the mutant that discriminates the mechanism. The wake channel - /// has capacity 1 and `set_subscriptions` only ever queues "reconcile - /// pending", so the delete and the recreate collapse into ONE observed - /// reconcile whose desired set already contains the id again. A prune that - /// reads only the current desired set never sees the id absent and leaves - /// the latch in place — passing the test above while failing this one. - /// The departure is therefore recorded at write time, where it is visible. - /// - /// No `settle()` between the two writes: that gap is the whole point, and - /// adding one would silently convert this into a duplicate of M18. - #[tokio::test] - async fn a_recreated_subscription_is_not_suppressed_when_the_writes_coalesce() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "restricted: not authorized".into(), - )) - .await - .expect("stub relay accepts the closed command"); - settle().await; - - session.set_subscriptions(vec![]).await; - session.set_subscriptions(vec![probe_subscription()]).await; - - assert_eq!( - next_req(&mut frames, "the REQ for the recreated subscription").await, - PROBE_ID, - ); - - session.shutdown(); - } - - /// M20: pruning must be scoped to departures, not run every pass. - /// - /// A reconcile triggered while the id is still desired must leave its - /// pending backoff alone. Clearing wholesale would collapse the CLOSED - /// backoff — every unrelated subscription change would re-ask a relay that - /// just rejected us, at the speed of the event loop. - #[tokio::test] - async fn a_reconcile_preserves_the_backoff_of_a_still_desired_subscription() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - // Rate-limited: a long, unambiguously pending backoff, so a reopen - // inside the window is the prune and not the timer. - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "rate-limited: slow down; retry in 30s".into(), - )) - .await - .expect("stub relay accepts the closed command"); - settle().await; - - // A change that adds an unrelated subscription. The probe never leaves - // the desired set, so its backoff must survive this reconcile. - session - .set_subscriptions(vec![ - probe_subscription(), - Subscription { - id: "archive:other".to_string(), - filter: serde_json::json!({ "kinds": [7], "limit": 0 }), - }, - ]) - .await; - - assert_eq!( - next_req(&mut frames, "the REQ for the newly added subscription").await, - "archive:other", - ); - let reopened = tokio::time::timeout(Duration::from_secs(3), frames.recv()).await; - assert!( - reopened.is_err(), - "a still-desired subscription must keep its pending backoff across a \ - reconcile, got {reopened:?}" - ); - - crate::relay_admission::reset_rate_limit_gate(); - session.shutdown(); - } - - /// M21: a CLOSED that arrives after we stopped running the subscription is - /// stale and must mint nothing. - /// - /// Our CLOSE races the relay's in-flight frames — the EVENT arm already - /// guards this. Without the same guard on CLOSED, the frame recreates the - /// retry entry the drain just removed, and nothing can evict it: the id is - /// gone from the desired set, so no future departure records it again. - #[tokio::test] - async fn a_closed_arriving_after_removal_does_not_mint_retry_state() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, _events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - // Delete first, and wait for our CLOSE to reach the wire: that ordering - // is what makes the CLOSED below arrive after the drain rather than - // before it, which is the schedule M18 and M19 do not cover. - session.set_subscriptions(vec![]).await; - assert_eq!( - next_frame(&mut frames, "the CLOSE for the deleted subscription").await, - Frame::Close(PROBE_ID.to_string()), - ); - - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "restricted: not authorized".into(), - )) - .await - .expect("stub relay accepts the closed command"); - settle().await; - - session.set_subscriptions(vec![probe_subscription()]).await; - - assert_eq!( - next_req(&mut frames, "the REQ for the recreated subscription").await, - PROBE_ID, - ); - - session.shutdown(); - } - - /// M22: a stale *terminal* CLOSED landing after the id was recreated must - /// not blackhole the live subscription. - /// - /// This one survives every defense above. The CLOSED is legitimately - /// attributed — the id is open again, so the M21 guard passes it — and - /// terminal means no `due_at`, so the timer arm is disabled and no wake is - /// pending. `open` loses the id while the relay keeps delivering, and the - /// EVENT arm drops every frame in silence. - /// - /// EOSE is the recovery edge because it is the only ordered fence - /// available: frames on one socket are totally ordered, so the previous - /// generation's CLOSED necessarily precedes the new generation's EOSE. - #[tokio::test] - async fn a_stale_terminal_closed_does_not_blackhole_a_recreated_subscription() { - let (relay_url, mut frames, closed) = stub_relay().await; - let (session, mut events) = start(relay_url, Keys::generate(), None); - - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); - settle().await; - - // Delete and recreate, so the id is open again under a new generation. - session.set_subscriptions(vec![]).await; - assert_eq!( - next_frame(&mut frames, "the CLOSE for the deleted subscription").await, - Frame::Close(PROBE_ID.to_string()), - ); - session.set_subscriptions(vec![probe_subscription()]).await; - assert_eq!( - next_req(&mut frames, "the REQ for the recreated subscription").await, - PROBE_ID, - ); - settle().await; - - // The old generation's terminal CLOSED, delayed past the new REQ. - closed - .send(StubCommand::Closed( - PROBE_ID.into(), - "restricted: not authorized".into(), - )) - .await - .expect("stub relay accepts the closed command"); - // The new generation's EOSE, which the wire orders after it. - closed - .send(StubCommand::Eose(PROBE_ID.into())) - .await - .expect("stub relay accepts the eose command"); - - // The EOSE found the id closed, so it must drive a reconcile that - // reopens it. Nothing else can: terminal schedules no timer, and the - // desired set is stable. - assert_eq!( - next_req(&mut frames, "the REQ healing the open-map mismatch").await, - PROBE_ID, - ); - - // And the heal converges rather than storming: the replacement EOSE - // finds the id open, so it wakes nothing. - closed - .send(StubCommand::Eose(PROBE_ID.into())) - .await - .expect("stub relay accepts the second eose command"); - let extra = tokio::time::timeout(Duration::from_secs(3), frames.recv()).await; - assert!( - extra.is_err(), - "an EOSE for an already-open subscription must not re-reconcile, got {extra:?}" - ); - - // The point of the heal: events flow again. - let event = EventBuilder::text_note("post-heal") - .sign_with_keys(&Keys::generate()) - .expect("sign event"); - let event_id = event.id.to_hex(); - closed - .send(StubCommand::Event( - PROBE_ID.into(), - serde_json::to_value(&event).expect("serialize event"), - )) - .await - .expect("stub relay accepts the event command"); - - let delivered = tokio::time::timeout(Duration::from_secs(10), events.recv()) - .await - .expect("timed out waiting for an event after the heal") - .expect("session channel closed"); - assert_eq!( - delivered.event.id.to_hex(), - event_id, - "events must flow again once the open map is healed" - ); - - session.shutdown(); - } - - /// M23: reusing an id for a changed filter must be *detected*. - /// - /// This test pins detection and nothing else. Post-violation behavior — - /// whether the subscription reopens, what happens to its retry state, what - /// the relay is sent — is unspecified by design, because the wire carries - /// only the id and an in-flight CLOSED from the old filter is - /// indistinguishable from one caused by the new one. Asserting any of that - /// would turn an unsupported input into a supported one. - /// - /// It exists because the `(id, filter)` departure diff is otherwise - /// unpinned: on every supported path it is byte-equivalent to an id-only - /// diff, so a refactor could revert it, pass every other test here, and - /// silently remove the one signal that tells C and D they broke the - /// contract. - #[test] - fn a_filter_change_under_a_reused_id_is_reported_as_a_contract_violation() { - let mut state = SessionState::default(); - - assert!( - state.replace_desired(vec![probe_subscription()]).is_empty(), - "a first desired set violates nothing" - ); - assert!( - state.replace_desired(vec![probe_subscription()]).is_empty(), - "an unchanged subscription is not a filter change" - ); - - let violations = state.replace_desired(vec![Subscription { - id: PROBE_ID.to_string(), - filter: serde_json::json!({ "kinds": [7], "limit": 0 }), - }]); - - assert_eq!( - violations, - vec![PROBE_ID.to_string()], - "a filter changed under a reused id must be reported" - ); - } - - #[test] - fn closed_messages_classify_like_the_renderer_policy() { - assert_eq!( - classify_closed("rate-limited: quota exceeded; retry in 4s"), - ClosedClass::RateLimited - ); - assert_eq!( - classify_closed("restricted: not authorized"), - ClosedClass::Terminal - ); - assert_eq!( - classify_closed("error: too many subscriptions"), - ClosedClass::Terminal - ); - // Transient AUTH race, not a permanent rejection — the one prefix that - // looks terminal and deliberately is not. - assert_eq!( - classify_closed("auth-required: we can't serve unauthenticated"), - ClosedClass::Retryable - ); - assert_eq!(classify_closed(""), ClosedClass::Retryable); - // Case and padding come from the relay, not from us. - assert_eq!( - classify_closed(" RESTRICTED: nope "), - ClosedClass::Terminal - ); - } - - #[test] - fn retry_delay_grows_and_stops_at_the_ceiling() { - let mut retry = ClosedRetry::default(); - assert_eq!(retry.backoff(), CLOSED_RETRY_BASE_DELAY); - - retry.schedule("error: temporary"); - assert_eq!(retry.backoff(), CLOSED_RETRY_BASE_DELAY * 2); - - for _ in 0..40 { - retry.schedule("error: temporary"); - } - assert_eq!( - retry.backoff(), - CLOSED_RETRY_MAX_DELAY, - "backoff must saturate at the ceiling rather than wrapping" - ); - } - - #[test] - fn a_rate_limited_closed_waits_at_least_the_relay_hint() { - let mut retry = ClosedRetry::default(); - retry.schedule("rate-limited: quota exceeded; retry in 12s"); - - let due = retry.due_at.expect("rate-limited must schedule a reopen"); - // The hint dominates the 1s first backoff, so this asserts the hint was - // honored rather than that anything at all was scheduled. - assert!( - due >= Instant::now() + Duration::from_secs(11), - "a 12s hint must not be undercut by the base backoff" - ); - crate::relay_admission::reset_rate_limit_gate(); - } - - #[test] - fn a_hintless_rate_limited_closed_uses_the_shared_default() { - let mut retry = ClosedRetry::default(); - retry.schedule("rate-limited: quota exceeded"); - - let due = retry.due_at.expect("rate-limited must schedule a reopen"); - assert!( - due >= Instant::now() + CLOSED_RATE_LIMIT_DEFAULT - Duration::from_secs(1), - "a hintless rate-limit must fall back to the shared default window" - ); - crate::relay_admission::reset_rate_limit_gate(); - } - - #[test] - fn retry_hints_parse_the_relays_canonical_format() { - assert_eq!( - parse_retry_in_seconds("rate-limited: quota exceeded; retry in 4s"), - Some(4) - ); - assert_eq!(parse_retry_in_seconds("rate-limited: quota exceeded"), None); - assert_eq!(parse_retry_in_seconds("retry in s"), None); - } -} +#[path = "native_relay_client_tests.rs"] +mod closed_recovery_tests; #[cfg(test)] mod relay_backed_tests { diff --git a/desktop/src-tauri/src/native_relay_client_tests.rs b/desktop/src-tauri/src/native_relay_client_tests.rs new file mode 100644 index 000000000..b1b8a632a --- /dev/null +++ b/desktop/src-tauri/src/native_relay_client_tests.rs @@ -0,0 +1,611 @@ +//! Lifecycle tests for [`super`]'s CLOSED recovery and subscription bookkeeping. +//! +//! Split out of `native_relay_client.rs` to keep that file under the desktop +//! file-size ratchet. Same `#[path]` sibling-module convention as +//! `archive/sync.rs` and its `sync_tests.rs`. + +use super::*; +use futures_util::{SinkExt, StreamExt}; +use nostr::EventBuilder; +use tokio_tungstenite::tungstenite::protocol::Message; + +/// The subscription id every test below drives. +const PROBE_ID: &str = "archive:probe"; + +/// Minimal relay that completes the NIP-42 handshake, records every REQ, +/// and sends a CLOSED only when the test asks it to. +/// +/// A real socket rather than a fake `NostrWsConnection`, because the bug +/// this covers lives in the lifecycle between frames — the loop's only +/// reconcile triggers — and a fake that hands the loop a `Closed` value +/// cannot show that a REQ went back out over the wire afterwards. Same +/// `accept_async` stub shape as `native_websocket.rs`'s live-TCP tests. +/// +/// CLOSED is test-driven rather than a scripted reply to the first REQ so +/// the test can wait for the session to go quiet first. `set_subscriptions` +/// queues a wake that may still be pending when an immediate CLOSED lands, +/// and that wake reopens the subscription on its own — which made the first +/// version of this test pass against the unfixed code. +/// +/// `frames` reports REQ and CLOSE in wire order, not REQ alone: the +/// lifecycle tests below assert that a CLOSE was sent before the REQ that +/// follows it, which a REQ-only channel cannot express. +async fn stub_relay() -> (String, mpsc::Receiver, mpsc::Sender) { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind stub relay"); + let address = listener.local_addr().expect("stub relay address"); + let (req_tx, req_rx) = mpsc::channel(16); + let (closed_tx, mut closed_rx) = mpsc::channel::(4); + + tokio::spawn(async move { + let (stream, _) = listener.accept().await.expect("accept"); + let mut socket = tokio_tungstenite::accept_async(stream) + .await + .expect("websocket handshake"); + + socket + .send(Message::Text(r#"["AUTH","stub-challenge"]"#.into())) + .await + .expect("send challenge"); + + loop { + tokio::select! { + incoming = socket.next() => { + let Some(Ok(Message::Text(text))) = incoming else { return }; + let Ok(frame) = serde_json::from_str::(&text) else { + continue; + }; + match frame[0].as_str() { + Some("AUTH") => { + let id = frame[1]["id"].as_str().unwrap_or_default(); + socket + .send(Message::Text( + serde_json::json!(["OK", id, true, ""]).to_string().into(), + )) + .await + .expect("send auth ok"); + } + Some("REQ") => { + let id = frame[1].as_str().unwrap_or_default().to_string(); + if req_tx.send(Frame::Req(id)).await.is_err() { + return; + } + } + Some("CLOSE") => { + let id = frame[1].as_str().unwrap_or_default().to_string(); + if req_tx.send(Frame::Close(id)).await.is_err() { + return; + } + } + _ => {} + } + } + Some(command) = closed_rx.recv() => { + let frame = match command { + StubCommand::Closed(id, message) => { + serde_json::json!(["CLOSED", id, message]) + } + StubCommand::Eose(id) => serde_json::json!(["EOSE", id]), + StubCommand::Event(id, event) => { + serde_json::json!(["EVENT", id, event]) + } + }; + socket + .send(Message::Text(frame.to_string().into())) + .await + .expect("send stub frame"); + } + } + } + }); + + (format!("ws://{address}"), req_rx, closed_tx) +} + +/// A client→relay frame the stub observed, in wire order. +#[derive(Debug, PartialEq, Eq)] +enum Frame { + Req(String), + Close(String), +} + +/// A relay→client frame the test asks the stub to emit. +enum StubCommand { + Closed(String, String), + Eose(String), + Event(String, serde_json::Value), +} + +fn probe_subscription() -> Subscription { + Subscription { + id: PROBE_ID.to_string(), + filter: serde_json::json!({ "kinds": [1], "limit": 0 }), + } +} + +async fn next_frame(frames: &mut mpsc::Receiver, label: &str) -> Frame { + tokio::time::timeout(Duration::from_secs(10), frames.recv()) + .await + .unwrap_or_else(|_| panic!("timed out waiting for {label}")) + .unwrap_or_else(|| panic!("stub relay closed before {label}")) +} + +/// Waits for the next REQ, tolerating the CLOSE frames a reconcile sends +/// first. Asserting on `Frame::Req` directly would couple every test to +/// whether a particular reconcile also had cleanup to do. +async fn next_req(frames: &mut mpsc::Receiver, label: &str) -> String { + loop { + if let Frame::Req(id) = next_frame(frames, label).await { + return id; + } + } +} + +/// Waits out the wake `set_subscriptions` queued, so a CLOSED sent after +/// this cannot be reopened by anything but the CLOSED path itself. +/// +/// A pending wake is harmless while the subscription is still open — that +/// reconcile is a no-op — so draining it before the CLOSED is what makes +/// the assertion below attributable. +async fn settle() { + tokio::time::sleep(Duration::from_millis(500)).await; +} + +/// The blocker: a CLOSED with the desired set never changing again must +/// still reopen the subscription. +/// +/// Before the fix the loop removed the id from `open` and waited on a wake +/// that only `set_subscriptions` can produce, so a stable desired set left +/// the subscription dead for the life of the socket — silent permanent +/// loss for ephemeral kind 24200. +#[tokio::test] +async fn a_closed_subscription_reopens_without_a_desired_set_change() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, _events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + // Retryable class, sent once: the reopen is answered normally, so a + // failure here means "never retried" rather than "retried into another + // rejection". + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "error: temporary".into(), + )) + .await + .expect("stub relay accepts the closed command"); + + // No `set_subscriptions` between the two REQs: the reopen must come + // from the CLOSED itself, which is exactly the edge that was missing. + assert_eq!(next_req(&mut frames, "the reopened REQ").await, PROBE_ID); + + session.shutdown(); +} + +/// A relay that rejects on policy must not be re-asked in a tight loop. +#[tokio::test] +async fn a_terminal_closed_is_not_retried_on_the_same_socket() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, _events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "restricted: not authorized".into(), + )) + .await + .expect("stub relay accepts the closed command"); + + // Long enough that a retryable class (1s base) would have reopened + // several times, so this asserts suppression rather than just slowness. + let retried = tokio::time::timeout(Duration::from_secs(5), frames.recv()).await; + assert!( + retried.is_err(), + "a terminal CLOSED must not be retried on this socket, got {retried:?}" + ); + + session.shutdown(); +} + +/// M18: a subscription deleted and recreated must get a fresh REQ, even +/// though its terminal latch says never to retry. +/// +/// The latch is scoped to the subscription that earned it. Recreating the +/// id is a new subscription that happens to share a name — `archive::sync` +/// derives the id from scope and kinds, so a delete/recreate of the same +/// saved subscription produces a byte-identical id and would otherwise +/// inherit a permanent suppression for the life of the socket. +#[tokio::test] +async fn a_recreated_subscription_does_not_inherit_a_terminal_latch() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, _events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "restricted: not authorized".into(), + )) + .await + .expect("stub relay accepts the closed command"); + settle().await; + + // Delete, then recreate — each observed as its own reconcile. + session.set_subscriptions(vec![]).await; + settle().await; + session.set_subscriptions(vec![probe_subscription()]).await; + + assert_eq!( + next_req(&mut frames, "the REQ for the recreated subscription").await, + PROBE_ID, + ); + + session.shutdown(); +} + +/// M19: the same schedule, with both writes landing before the loop +/// consumes its single wake. +/// +/// This is the mutant that discriminates the mechanism. The wake channel +/// has capacity 1 and `set_subscriptions` only ever queues "reconcile +/// pending", so the delete and the recreate collapse into ONE observed +/// reconcile whose desired set already contains the id again. A prune that +/// reads only the current desired set never sees the id absent and leaves +/// the latch in place — passing the test above while failing this one. +/// The departure is therefore recorded at write time, where it is visible. +/// +/// No `settle()` between the two writes: that gap is the whole point, and +/// adding one would silently convert this into a duplicate of M18. +#[tokio::test] +async fn a_recreated_subscription_is_not_suppressed_when_the_writes_coalesce() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, _events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "restricted: not authorized".into(), + )) + .await + .expect("stub relay accepts the closed command"); + settle().await; + + session.set_subscriptions(vec![]).await; + session.set_subscriptions(vec![probe_subscription()]).await; + + assert_eq!( + next_req(&mut frames, "the REQ for the recreated subscription").await, + PROBE_ID, + ); + + session.shutdown(); +} + +/// M20: pruning must be scoped to departures, not run every pass. +/// +/// A reconcile triggered while the id is still desired must leave its +/// pending backoff alone. Clearing wholesale would collapse the CLOSED +/// backoff — every unrelated subscription change would re-ask a relay that +/// just rejected us, at the speed of the event loop. +#[tokio::test] +async fn a_reconcile_preserves_the_backoff_of_a_still_desired_subscription() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, _events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + // Rate-limited: a long, unambiguously pending backoff, so a reopen + // inside the window is the prune and not the timer. + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "rate-limited: slow down; retry in 30s".into(), + )) + .await + .expect("stub relay accepts the closed command"); + settle().await; + + // A change that adds an unrelated subscription. The probe never leaves + // the desired set, so its backoff must survive this reconcile. + session + .set_subscriptions(vec![ + probe_subscription(), + Subscription { + id: "archive:other".to_string(), + filter: serde_json::json!({ "kinds": [7], "limit": 0 }), + }, + ]) + .await; + + assert_eq!( + next_req(&mut frames, "the REQ for the newly added subscription").await, + "archive:other", + ); + let reopened = tokio::time::timeout(Duration::from_secs(3), frames.recv()).await; + assert!( + reopened.is_err(), + "a still-desired subscription must keep its pending backoff across a \ + reconcile, got {reopened:?}" + ); + + crate::relay_admission::reset_rate_limit_gate(); + session.shutdown(); +} + +/// M21: a CLOSED that arrives after we stopped running the subscription is +/// stale and must mint nothing. +/// +/// Our CLOSE races the relay's in-flight frames — the EVENT arm already +/// guards this. Without the same guard on CLOSED, the frame recreates the +/// retry entry the drain just removed, and nothing can evict it: the id is +/// gone from the desired set, so no future departure records it again. +#[tokio::test] +async fn a_closed_arriving_after_removal_does_not_mint_retry_state() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, _events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + // Delete first, and wait for our CLOSE to reach the wire: that ordering + // is what makes the CLOSED below arrive after the drain rather than + // before it, which is the schedule M18 and M19 do not cover. + session.set_subscriptions(vec![]).await; + assert_eq!( + next_frame(&mut frames, "the CLOSE for the deleted subscription").await, + Frame::Close(PROBE_ID.to_string()), + ); + + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "restricted: not authorized".into(), + )) + .await + .expect("stub relay accepts the closed command"); + settle().await; + + session.set_subscriptions(vec![probe_subscription()]).await; + + assert_eq!( + next_req(&mut frames, "the REQ for the recreated subscription").await, + PROBE_ID, + ); + + session.shutdown(); +} + +/// M22: a stale *terminal* CLOSED landing after the id was recreated must +/// not blackhole the live subscription. +/// +/// This one survives every defense above. The CLOSED is legitimately +/// attributed — the id is open again, so the M21 guard passes it — and +/// terminal means no `due_at`, so the timer arm is disabled and no wake is +/// pending. `open` loses the id while the relay keeps delivering, and the +/// EVENT arm drops every frame in silence. +/// +/// EOSE is the recovery edge because it is the only ordered fence +/// available: frames on one socket are totally ordered, so the previous +/// generation's CLOSED necessarily precedes the new generation's EOSE. +#[tokio::test] +async fn a_stale_terminal_closed_does_not_blackhole_a_recreated_subscription() { + let (relay_url, mut frames, closed) = stub_relay().await; + let (session, mut events) = start(relay_url, Keys::generate(), None); + + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!(next_req(&mut frames, "the initial REQ").await, PROBE_ID); + settle().await; + + // Delete and recreate, so the id is open again under a new generation. + session.set_subscriptions(vec![]).await; + assert_eq!( + next_frame(&mut frames, "the CLOSE for the deleted subscription").await, + Frame::Close(PROBE_ID.to_string()), + ); + session.set_subscriptions(vec![probe_subscription()]).await; + assert_eq!( + next_req(&mut frames, "the REQ for the recreated subscription").await, + PROBE_ID, + ); + settle().await; + + // The old generation's terminal CLOSED, delayed past the new REQ. + closed + .send(StubCommand::Closed( + PROBE_ID.into(), + "restricted: not authorized".into(), + )) + .await + .expect("stub relay accepts the closed command"); + // The new generation's EOSE, which the wire orders after it. + closed + .send(StubCommand::Eose(PROBE_ID.into())) + .await + .expect("stub relay accepts the eose command"); + + // The EOSE found the id closed, so it must drive a reconcile that + // reopens it. Nothing else can: terminal schedules no timer, and the + // desired set is stable. + assert_eq!( + next_req(&mut frames, "the REQ healing the open-map mismatch").await, + PROBE_ID, + ); + + // And the heal converges rather than storming: the replacement EOSE + // finds the id open, so it wakes nothing. + closed + .send(StubCommand::Eose(PROBE_ID.into())) + .await + .expect("stub relay accepts the second eose command"); + let extra = tokio::time::timeout(Duration::from_secs(3), frames.recv()).await; + assert!( + extra.is_err(), + "an EOSE for an already-open subscription must not re-reconcile, got {extra:?}" + ); + + // The point of the heal: events flow again. + let event = EventBuilder::text_note("post-heal") + .sign_with_keys(&Keys::generate()) + .expect("sign event"); + let event_id = event.id.to_hex(); + closed + .send(StubCommand::Event( + PROBE_ID.into(), + serde_json::to_value(&event).expect("serialize event"), + )) + .await + .expect("stub relay accepts the event command"); + + let delivered = tokio::time::timeout(Duration::from_secs(10), events.recv()) + .await + .expect("timed out waiting for an event after the heal") + .expect("session channel closed"); + assert_eq!( + delivered.event.id.to_hex(), + event_id, + "events must flow again once the open map is healed" + ); + + session.shutdown(); +} + +/// M23: reusing an id for a changed filter must be *detected*. +/// +/// This test pins detection and nothing else. Post-violation behavior — +/// whether the subscription reopens, what happens to its retry state, what +/// the relay is sent — is unspecified by design, because the wire carries +/// only the id and an in-flight CLOSED from the old filter is +/// indistinguishable from one caused by the new one. Asserting any of that +/// would turn an unsupported input into a supported one. +/// +/// It exists because the `(id, filter)` departure diff is otherwise +/// unpinned: on every supported path it is byte-equivalent to an id-only +/// diff, so a refactor could revert it, pass every other test here, and +/// silently remove the one signal that tells C and D they broke the +/// contract. +#[test] +fn a_filter_change_under_a_reused_id_is_reported_as_a_contract_violation() { + let mut state = SessionState::default(); + + assert!( + state.replace_desired(vec![probe_subscription()]).is_empty(), + "a first desired set violates nothing" + ); + assert!( + state.replace_desired(vec![probe_subscription()]).is_empty(), + "an unchanged subscription is not a filter change" + ); + + let violations = state.replace_desired(vec![Subscription { + id: PROBE_ID.to_string(), + filter: serde_json::json!({ "kinds": [7], "limit": 0 }), + }]); + + assert_eq!( + violations, + vec![PROBE_ID.to_string()], + "a filter changed under a reused id must be reported" + ); +} + +#[test] +fn closed_messages_classify_like_the_renderer_policy() { + assert_eq!( + classify_closed("rate-limited: quota exceeded; retry in 4s"), + ClosedClass::RateLimited + ); + assert_eq!( + classify_closed("restricted: not authorized"), + ClosedClass::Terminal + ); + assert_eq!( + classify_closed("error: too many subscriptions"), + ClosedClass::Terminal + ); + // Transient AUTH race, not a permanent rejection — the one prefix that + // looks terminal and deliberately is not. + assert_eq!( + classify_closed("auth-required: we can't serve unauthenticated"), + ClosedClass::Retryable + ); + assert_eq!(classify_closed(""), ClosedClass::Retryable); + // Case and padding come from the relay, not from us. + assert_eq!( + classify_closed(" RESTRICTED: nope "), + ClosedClass::Terminal + ); +} + +#[test] +fn retry_delay_grows_and_stops_at_the_ceiling() { + let mut retry = ClosedRetry::default(); + assert_eq!(retry.backoff(), CLOSED_RETRY_BASE_DELAY); + + retry.schedule("error: temporary"); + assert_eq!(retry.backoff(), CLOSED_RETRY_BASE_DELAY * 2); + + for _ in 0..40 { + retry.schedule("error: temporary"); + } + assert_eq!( + retry.backoff(), + CLOSED_RETRY_MAX_DELAY, + "backoff must saturate at the ceiling rather than wrapping" + ); +} + +#[test] +fn a_rate_limited_closed_waits_at_least_the_relay_hint() { + let mut retry = ClosedRetry::default(); + retry.schedule("rate-limited: quota exceeded; retry in 12s"); + + let due = retry.due_at.expect("rate-limited must schedule a reopen"); + // The hint dominates the 1s first backoff, so this asserts the hint was + // honored rather than that anything at all was scheduled. + assert!( + due >= Instant::now() + Duration::from_secs(11), + "a 12s hint must not be undercut by the base backoff" + ); + crate::relay_admission::reset_rate_limit_gate(); +} + +#[test] +fn a_hintless_rate_limited_closed_uses_the_shared_default() { + let mut retry = ClosedRetry::default(); + retry.schedule("rate-limited: quota exceeded"); + + let due = retry.due_at.expect("rate-limited must schedule a reopen"); + assert!( + due >= Instant::now() + CLOSED_RATE_LIMIT_DEFAULT - Duration::from_secs(1), + "a hintless rate-limit must fall back to the shared default window" + ); + crate::relay_admission::reset_rate_limit_gate(); +} + +#[test] +fn retry_hints_parse_the_relays_canonical_format() { + assert_eq!( + parse_retry_in_seconds("rate-limited: quota exceeded; retry in 4s"), + Some(4) + ); + assert_eq!(parse_retry_in_seconds("rate-limited: quota exceeded"), None); + assert_eq!(parse_retry_in_seconds("retry in s"), None); +}