diff --git a/CHANGELOG.md b/CHANGELOG.md index db6cdec..e7d003b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,13 @@ 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. +- 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 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..c6a1b22 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -400,6 +400,14 @@ 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. +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 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..419f991 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) @@ -2684,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); } } @@ -3664,6 +3681,104 @@ 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 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;