From 8a141a55d2f90d97c86f0aa8bc12641589aeaacd Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 27 Jul 2026 08:50:09 +0100 Subject: [PATCH] fix(sync): keep relay retries scheduler-owned The bounded connection workers only provide a real global limit if they own every retry. nostr-sdk enables an independent reconnect loop by default, so a relay that disconnected after its first successful invitation sync could reconnect outside the semaphore, health backoff, and actor metrics. Detached workers could also outlive an aborted sync manager while waiting for capacity or a socket. Disable SDK auto-reconnect for managed relay connections and have every queued or active worker observe the manager's shutdown channel. If its result receiver disappears, a worker closes the connection instead of leaving an unowned socket behind. --- CHANGELOG.md | 3 ++ docs/explanation/grasp-02-proactive-sync.md | 3 +- src/sync/mod.rs | 39 +++++++++++++++------ src/sync/relay_connection.rs | 1 + 4 files changed, 35 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9b9a64f..9b99da6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Fix large invitation relay sets repeatedly rebuilding their subscriptions by applying the 70-filter consolidation threshold only to fragmentation above the relay's irreducible desired live-filter baseline. +- Keep invitation relay connection ownership inside the bounded scheduler by + disabling the SDK's independent auto-reconnect loop and cancelling queued or + active connection workers when the sync manager shuts down. - Sideband-aware Git clients now receive periodic progress while GRASP performs post-push purgatory promotion and cross-owner repository alignment, preventing the client I/O timeout from expiring during unusually complex finalization. - Smart HTTP pushes now expose the terminal receive-pack flush only after GRASP has finished promoting the matching repository announcement and state from purgatory. Git progress remains streamed while large packs are resolved and checked, but an immediately following clone or proposal push can now rely on a completed push being queryable on the relay. - Batch the one-time deletion-request lifecycle migration so large production databases do not remain unavailable while LMDB commits every historical request in separate transactions. diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index fd3601e..6cd1754 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -27,7 +27,8 @@ Key Architectural Points: - PendingBatch tracks each new set of filters that may require pagination until they are complete - Websocket handshakes run in at most eight bounded workers outside the sync actor; only the actor applies their results, and subscriptions start only - after the relay reports `Connected` + after the relay reports `Connected`. The sync manager owns retry/backoff + rather than the SDK, and shutdown cancels queued or active workers - Recompute desired filters when connection is established to ensure filters are as consolidated as possible - Consolidation bounds incremental fragmentation to 70 subscriptions above the irreducible desired live-filter baseline, so large desired sets remain diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 05b0187..4f3f124 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -2580,20 +2580,39 @@ impl SyncManager { let health_tracker = Arc::clone(&self.health_tracker); let semaphore = Arc::clone(&self.connect_attempt_semaphore); let timeout = self.health_tracker.base_backoff_secs(); + let Some(mut shutdown_rx) = self.shutdown_tx.as_ref().map(|sender| sender.subscribe()) + else { + tracing::error!(relay = %relay_url, "Connection scheduler has no shutdown signal"); + return; + }; tokio::spawn(async move { - let Some(_permit) = begin_connect_attempt(semaphore, health_tracker, &relay_url).await - else { + let permit = tokio::select! { + permit = begin_connect_attempt(semaphore, health_tracker, &relay_url) => permit, + _ = shutdown_rx.recv() => return, + }; + let Some(_permit) = permit else { return; }; - let outcome = match connection.connect(timeout).await { - Ok(()) => ConnectAttemptOutcome::Connected, - Err(error) => ConnectAttemptOutcome::Failed(error), + let outcome = tokio::select! { + result = connection.connect(timeout) => match result { + Ok(()) => ConnectAttemptOutcome::Connected, + Err(error) => ConnectAttemptOutcome::Failed(error), + }, + _ = shutdown_rx.recv() => { + connection.disconnect().await; + return; + } }; - let _ = result_tx.send(ConnectAttemptResult { - relay_url, - token, - outcome, - }); + if result_tx + .send(ConnectAttemptResult { + relay_url, + token, + outcome, + }) + .is_err() + { + connection.disconnect().await; + } }); } diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 7bb470a..09b6df8 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -208,6 +208,7 @@ impl RelayConnection { let relay = tokio::time::timeout_at(connection_deadline, async { self.client .add_relay(&self.url) + .reconnect(false) .await .map_err(|e| format!("Failed to add relay {}: {}", self.url, e))?;