Merge #dc767f66: Retire obsolete SDK subscriptions before replay

nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsdcanlv6qh6gse3syxz2chun73vvw4l2zx65399t3g0yllwy6trkcmur38x

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

PR description:

Disconnected relay sessions kept the SDK's live subscriptions after their local permits were cleared. Reconnecting could replay those old requests alongside newly planned coverage, accumulating untracked subscriptions and exhausting remote capacity. The deployment investigation found repeated 5 MiB retained-subscription refusals followed by supposedly successful bounded rebuilds.

Two independently reviewable fixes retire subscriptions at their ownership boundary:

- Replace the ended SDK relay before reconnecting. The manager rebuilds coverage under a fresh ledger while retaining connection-level learned limits.
- Unregister rolled-back live groups through the SDK instead of sending only CLOSE, so authentication cannot replay them after their permits are released.

Both regressions fail on their prior implementation. The reconnect test uses a real WebSocket peer, exercises peer-initiated and explicit disconnects, and checks wire IDs and SDK inventory under a one-slot budget. The rollback test verifies SDK registry removal using a real local relay. All 429 sync unit tests pass; workspace/all-target Clippy and formatting are checked. Coverage selection, peer limits and backoff are unchanged.
This commit is contained in:
DanConwayDev
2026-10-02 08:36:03 +01:00
3 changed files with 134 additions and 4 deletions
+7
View File
@@ -7,6 +7,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased] ## [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 ## [3.0.6] - 2026-09-30
Keep unsigned Git uploads temporary, reclaim abandoned PR data, and preserve Keep unsigned Git uploads temporary, reclaim abandoned PR data, and preserve
@@ -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 required core filters. An auxiliary CLOSED retires its remaining group and
falls back to EOSE-closing history without rebuilding healthy core coverage. 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 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 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 before the first refusal. A CLOSED reason of the rust-nostr form `active
+119 -4
View File
@@ -973,6 +973,24 @@ impl RelayConnection {
}; };
let relay = tokio::time::timeout_at(connection_deadline, async { 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 self.client
.add_relay(&self.url) .add_relay(&self.url)
.reconnect(false) .reconnect(false)
@@ -2684,11 +2702,10 @@ impl RelayConnection {
/// the corresponding ledger slot is returned, keeping local and relay /// the corresponding ledger slot is returned, keeping local and relay
/// subscription accounting in lockstep. /// subscription accounting in lockstep.
async fn close_and_release_live_req_permits(&self, sub_ids: &[SubscriptionId]) { 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 { for sub_id in sub_ids {
if let Some(relay) = &relay { // A raw CLOSE leaves the SDK's desired subscription registered,
let _ = relay.send_msg(ClientMessage::close(sub_id.clone())).await; // allowing authentication to replay a rolled-back live group.
} let _ = self.client.unsubscribe(sub_id).await;
let _ = self.release_live_req_permit(sub_id); let _ = self.release_live_req_permit(sub_id);
} }
} }
@@ -3664,6 +3681,104 @@ mod tests {
drain.abort(); 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] #[tokio::test]
async fn peer_closed_subscription_is_removed_from_sdk_registry() { async fn peer_closed_subscription_is_removed_from_sdk_registry() {
let relay = TestRelay::start(LocalRelayBuilder::default()).await; let relay = TestRelay::start(LocalRelayBuilder::default()).await;