From 4765f7ce206d16f348bfbcd1fa5c6f4e9d842dfd Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 2 Oct 2026 07:16:19 +0000 Subject: [PATCH 1/2] fix(sync): retire SDK subscriptions before reconnecting A disconnected session loses its local permits, but rust-nostr retains and replays its live REQs. Rebuilding coverage then duplicates those subscriptions outside the new ledger, where consolidation cannot close them. Production logs showed repeated retained-filter capacity refusals despite bounded rebuilds. Replace the ended SDK relay before dialing so old subscriptions and queued messages cannot enter the new session. Preserve connection-level learned hints; the manager remains responsible for rebuilding live and historic coverage. This does not change peer limits, backoff, or the coverage selection policy. Validation: a real WebSocket regression fails on the previous implementation and passes across peer-initiated and explicit disconnects. All 428 sync unit tests pass. The test checks wire IDs and SDK inventory under a one-slot budget. Assisted-by: GPT-6 --- CHANGELOG.md | 5 + docs/explanation/sync-scaling-constraints.md | 6 ++ src/sync/relay_connection.rs | 96 ++++++++++++++++++++ 3 files changed, 107 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index db6cdec..00a86ca 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- Retire old relay subscriptions before reconnecting so SDK replay cannot + accumulate untracked subscriptions or exhaust peers' retained-filter limits. + ## [3.0.6] - 2026-09-30 Keep unsigned Git uploads temporary, reclaim abandoned PR data, and preserve diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index ecdb3ae..e7b4484 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -400,6 +400,12 @@ sent, so auxiliary coverage cannot silently consume capacity needed by newly required core filters. An auxiliary CLOSED retires its remaining group and falls back to EOSE-closing history without rebuilding healthy core coverage. +Before reconnecting, the connection replaces the ended SDK relay session. The +manager rebuilds live coverage under a fresh ledger; keeping the SDK's old +subscriptions would let its automatic replay bypass that ledger and accumulate +untracked remote subscriptions. Learned connection-level limits are preserved, +while SDK subscriptions and queued messages belong only to their old session. + Some relays additionally cap the cumulative serialized REQ state retained by one connection. NIP-11 has no field for this limit, so it cannot be negotiated before the first refusal. A CLOSED reason of the rust-nostr form `active diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index f7da78e..b8a23f4 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -973,6 +973,24 @@ impl RelayConnection { }; let relay = tokio::time::timeout_at(connection_deadline, async { + // The manager rebuilds coverage under a new ledger after every + // connection. The SDK otherwise replays its previous live REQs, + // which no longer own permits and cannot be retired by consolidation. + // Replace the ended SDK session before dialing, including its queued + // messages and auto-closing requests; retain connection-level hints. + if self + .client + .relay(&self.url) + .await + .map_err(|e| e.to_string())? + .is_some() + { + self.client + .remove_relay(&self.url) + .force() + .await + .map_err(|e| format!("Failed to retire relay session {}: {e}", self.url))?; + } self.client .add_relay(&self.url) .reconnect(false) @@ -3664,6 +3682,84 @@ mod tests { drain.abort(); } + #[tokio::test] + async fn reconnect_does_not_replay_subscriptions_outside_the_new_ledger() { + use futures_util::SinkExt; + use tokio_tungstenite::tungstenite::Message; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + let (requests_tx, mut requests_rx) = tokio::sync::mpsc::unbounded_channel(); + let server = tokio::spawn(async move { + for session in 0..3 { + let (socket, _) = listener.accept().await.unwrap(); + let mut ws = tokio_tungstenite::accept_async(socket).await.unwrap(); + while let Some(Ok(frame)) = ws.next().await { + if frame.is_close() { + break; + } + if !frame.is_text() { + continue; + } + let request: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + if request[0] == "REQ" { + requests_tx + .send((session, request[1].as_str().unwrap().to_owned())) + .unwrap(); + ws.send(Message::Text( + serde_json::json!(["EOSE", request[1]]).to_string().into(), + )) + .await + .unwrap(); + if session == 0 { + ws.close(None).await.unwrap(); + break; + } + } + } + } + }); + let connection = permissive_connection(&url, Keys::generate()); + for session in 0..3 { + connection.connect(3).await.unwrap(); + connection.reset_subscription_budget(Some(3)); + let ids = connection + .subscribe_live_filter_groups(vec![vec![Filter::new() + .kind(Kind::TextNote) + .limit(0)]]) + .await + .unwrap(); + let (wire_session, wire_id) = + tokio::time::timeout(Duration::from_secs(3), requests_rx.recv()) + .await + .unwrap() + .unwrap(); + assert_eq!(wire_session, session); + assert_eq!( + wire_id, + ids[0].to_string(), + "only the new ledger's REQ may be sent" + ); + assert_eq!(connection.subscription_count().await, 1); + if session == 0 { + tokio::time::timeout(Duration::from_secs(3), async { + while connection.is_connected().await { + tokio::task::yield_now().await; + } + }) + .await + .expect("peer disconnect must become observable"); + } else { + connection.disconnect().await; + } + } + tokio::time::timeout(Duration::from_secs(3), server) + .await + .unwrap() + .unwrap(); + } + #[tokio::test] async fn peer_closed_subscription_is_removed_from_sdk_registry() { let relay = TestRelay::start(LocalRelayBuilder::default()).await; From 71a248ba2362d4a1e6e7596a04193e35bb476bd9 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 2 Oct 2026 07:18:40 +0000 Subject: [PATCH 2/2] fix(sync): unregister rolled-back live subscriptions Partial live-set rollback sent raw CLOSE messages and released local permits, but left the SDK's desired subscriptions registered. Authentication can replay those groups without ledger ownership even when the socket never reconnects. Use SDK unsubscribe for each rolled-back live group before returning its slot. This preserves the existing best-effort close semantics and leaves transient request handling and coverage selection unchanged. Validation: the real-relay regression fails with one retained SDK subscription before this fix and passes with none afterward. All 429 sync unit tests pass. Assisted-by: GPT-6 --- CHANGELOG.md | 2 ++ docs/explanation/sync-scaling-constraints.md | 2 ++ src/sync/relay_connection.rs | 27 +++++++++++++++++--- 3 files changed, 27 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 00a86ca..e7d003b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Retire old relay subscriptions before reconnecting so SDK replay cannot accumulate untracked subscriptions or exhaust peers' retained-filter limits. +- Remove rolled-back live subscriptions from the SDK registry so authentication + cannot replay groups whose local permits have already been released. ## [3.0.6] - 2026-09-30 diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index e7b4484..c6a1b22 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -405,6 +405,8 @@ manager rebuilds live coverage under a fresh ledger; keeping the SDK's old subscriptions would let its automatic replay bypass that ledger and accumulate untracked remote subscriptions. Learned connection-level limits are preserved, while SDK subscriptions and queued messages belong only to their old session. +Rolling back a partial live set also unregisters its SDK subscriptions, so an +authentication retry cannot replay coverage whose local permits were released. Some relays additionally cap the cumulative serialized REQ state retained by one connection. NIP-11 has no field for this limit, so it cannot be negotiated diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index b8a23f4..419f991 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -2702,11 +2702,10 @@ impl RelayConnection { /// the corresponding ledger slot is returned, keeping local and relay /// subscription accounting in lockstep. async fn close_and_release_live_req_permits(&self, sub_ids: &[SubscriptionId]) { - let relay = self.client.relay(&self.url).await.ok().flatten(); for sub_id in sub_ids { - if let Some(relay) = &relay { - let _ = relay.send_msg(ClientMessage::close(sub_id.clone())).await; - } + // A raw CLOSE leaves the SDK's desired subscription registered, + // allowing authentication to replay a rolled-back live group. + let _ = self.client.unsubscribe(sub_id).await; let _ = self.release_live_req_permit(sub_id); } } @@ -3760,6 +3759,26 @@ mod tests { .unwrap(); } + #[tokio::test] + async fn rolled_back_live_groups_are_removed_from_the_sdk_registry() { + let relay = TestRelay::start(LocalRelayBuilder::default()).await; + let connection = permissive_connection(&relay.url().await.to_string(), Keys::generate()); + connection.connect(3).await.unwrap(); + let ids = connection + .subscribe_live_filter_groups(vec![vec![Filter::new().kind(Kind::TextNote).limit(0)]]) + .await + .unwrap(); + connection.close_and_release_live_req_permits(&ids).await; + assert_eq!( + connection.subscription_count().await, + 0, + "rollback must remove SDK replay eligibility, not just send CLOSE" + ); + assert!(!connection.holds_subscription_permit(&ids[0])); + connection.disconnect().await; + relay.shutdown(); + } + #[tokio::test] async fn peer_closed_subscription_is_removed_from_sdk_registry() { let relay = TestRelay::start(LocalRelayBuilder::default()).await;