From 6b373e9ecdbe6711b947d5b237b9c80b2d0961db Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 2 Oct 2026 13:27:00 +0000 Subject: [PATCH] fix(sync): stop self-subscriber when its consumer exits Server shutdown aborts the SyncManager action receiver while its detached self-subscriber may still be sending a batch. The subscriber previously logged one channel-closed error per remaining relay and could ignore shutdown while blocked by channel backpressure. Monitor receiver closure and shutdown around the entire subscriber lifecycle, including startup reconstruction and batch sends. Stop the batch on the first permanent send failure, including races with cancellation. Pending discovery can be reconstructed from persisted events on startup; this does not redesign ownership of other sync workers. Validation: all 13 self-subscriber tests pass, including bounded regressions for consumer closure and shutdown during a backpressured startup batch. cargo clippy --lib --tests -- -D warnings, formatting and git diff --check pass. Shutdown behavior is documented. Assisted-by: GPT-6 --- docs/explanation/grasp-02-proactive-sync.md | 6 + src/sync/self_subscriber.rs | 170 ++++++++++++++------ 2 files changed, 125 insertions(+), 51 deletions(-) diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index fc5fa00..a7c98ad 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -425,6 +425,12 @@ what lets the live feed work under GRASP-08 private mode, where the NIP-42 gate in the HTTP layer refuses a self-dial (see [GRASP-08 design](grasp-08-private-service.md)) +**Shutdown**: The subscriber exits when its shutdown signal arrives or the +SyncManager action receiver closes. This covers startup reconstruction and +blocked batch sends, not just the notification loop. Pending notifications +are not flushed to a stopped consumer; startup reconstructs discovery from +the persisted events on the next run. + **Subscribed kinds**: 30617, 1617, 1618, 1621 (NOT 30618) **Batching**: 5-second window (configurable via `NGIT_SYNC_BATCH_WINDOW_MS`) diff --git a/src/sync/self_subscriber.rs b/src/sync/self_subscriber.rs index 08b41e5..44fce81 100644 --- a/src/sync/self_subscriber.rs +++ b/src/sync/self_subscriber.rs @@ -414,7 +414,29 @@ impl SelfSubscriber { /// /// The optional shutdown receiver allows graceful termination when /// received via the broadcast channel. - pub async fn run(mut self, mut shutdown_rx: Option>) { + pub async fn run(self, mut shutdown_rx: Option>) { + let action_tx = self.action_tx.clone(); + // Cover startup reconstruction and batch sends too: checking only + // between batches cannot interrupt a send blocked on a full channel. + tokio::select! { + biased; + _ = action_tx.closed() => { + tracing::info!("SelfSubscriber action consumer stopped"); + } + _ = async { + match shutdown_rx.as_mut() { + Some(rx) => { let _ = rx.recv().await; } + None => std::future::pending::<()>().await, + } + } => { + tracing::info!("SelfSubscriber received shutdown signal"); + } + _ = self.run_inner() => {} + } + tracing::info!("SelfSubscriber stopped"); + } + + async fn run_inner(mut self) { // Attach to the embedded relay directly rather than dialling our own // public listener. A private instance's NIP-42 gate would refuse that // dial, stalling every runtime discovery until the next restart; the @@ -483,65 +505,38 @@ impl SelfSubscriber { let mut pending = self.load_existing_events().await; // Publish before consuming queued live notifications so newly arriving // roots can resolve their repository through the completed index. - self.process_batch(&mut pending).await; + if let LoopControl::Break = self.process_batch(&mut pending).await { + return; + } // Timer does NOT reset on new events - use interval let mut timer = tokio::time::interval(batch_window); timer.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { - // Build the select based on whether we have a shutdown receiver - if let Some(ref mut rx) = shutdown_rx { - tokio::select! { - notification = notifications.next() => { - match notification { - Some(notification) => { - if let LoopControl::Break = self.process_notification(notification, &mut pending).await { - break; - } - } - None => { - tracing::info!("SelfSubscriber notification stream ended"); + tokio::select! { + notification = notifications.next() => { + match notification { + Some(notification) => { + if let LoopControl::Break = self.process_notification(notification, &mut pending).await { break; } } - } - _ = timer.tick() => { - if !pending.is_empty() { - self.process_batch(&mut pending).await; + None => { + tracing::info!("SelfSubscriber notification stream ended"); + break; } } - _ = rx.recv() => { - tracing::info!("SelfSubscriber received shutdown signal"); - break; - } } - } else { - // No shutdown receiver - original behavior - tokio::select! { - notification = notifications.next() => { - match notification { - Some(notification) => { - if let LoopControl::Break = self.process_notification(notification, &mut pending).await { - break; - } - } - None => { - tracing::info!("SelfSubscriber notification stream ended"); - break; - } - } - } - _ = timer.tick() => { - if !pending.is_empty() { - self.process_batch(&mut pending).await; + _ = timer.tick() => { + if !pending.is_empty() { + if let LoopControl::Break = self.process_batch(&mut pending).await { + break; } } } } } - - tracing::info!("SelfSubscriber stopped"); } /// Handle a root event (1617/1618/1621) @@ -600,11 +595,11 @@ impl SelfSubscriber { /// /// Updates the RepoSyncIndex with discovered repos, then marks only the /// relays changed by this batch for SyncManager recomputation. - async fn process_batch(&self, pending: &mut PendingUpdates) { + async fn process_batch(&self, pending: &mut PendingUpdates) -> LoopControl { let (updates, relay_replacements) = pending.take(); if updates.is_empty() { - return; + return LoopControl::Continue; } tracing::info!( @@ -673,12 +668,11 @@ impl SelfSubscriber { filters: Vec::new(), }; - if let Err(e) = self.action_tx.send(action).await { - tracing::error!( - relay = %relay_url, - error = %e, - "Failed to send AddFilters action" - ); + if self.action_tx.send(action).await.is_err() { + // Closure is permanent. Stop this batch even if it races the + // lifecycle select, rather than retrying every remaining relay. + tracing::debug!("Stopped batch after SyncManager action channel closed"); + return LoopControl::Break; } else { tracing::debug!( relay = %relay_url, @@ -686,6 +680,7 @@ impl SelfSubscriber { ); } } + LoopControl::Continue } } @@ -729,6 +724,79 @@ mod tests { use std::sync::Arc; use tokio::sync::RwLock; + async fn stop_during_startup_batch(close_consumer: bool) { + let database: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded()); + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags([ + Tag::identifier("shutdown-batch"), + Tag::custom( + "relays", + [ + "wss://one.example", + "wss://two.example", + "wss://three.example", + ], + ), + ]) + .finalize(&Keys::generate()) + .unwrap(); + database.save_event(&announcement).await.unwrap(); + let (action_tx, mut action_rx) = mpsc::channel(1); + let (shutdown_tx, shutdown_rx) = broadcast::channel(1); + let subscriber = SelfSubscriber::new( + "ws://127.0.0.1:1".into(), + "127.0.0.1:1".into(), + Arc::new(RwLock::new(HashMap::new())), + Arc::new(RwLock::new(HashMap::new())), + action_tx, + database, + LocalRelay::new(), + ); + // Dropping the set cancels the worker even if an assertion fails. + let mut tasks = tokio::task::JoinSet::new(); + tasks.spawn(subscriber.run(Some(shutdown_rx))); + tokio::time::timeout(Duration::from_secs(5), async { + action_rx + .recv() + .await + .expect("startup emits a relay action"); + // The second action fills the sole slot; the third cannot be sent. + while action_rx.is_empty() { + tokio::task::yield_now().await; + } + }) + .await + .expect("startup batch reaches channel backpressure"); + if close_consumer { + action_rx.close(); + } else { + shutdown_tx.send(()).unwrap(); + } + tokio::time::timeout(Duration::from_secs(2), tasks.join_next()) + .await + .expect("subscriber stops without draining the batch") + .expect("subscriber task exists") + .expect("subscriber did not panic"); + assert_eq!(action_rx.len(), 1, "remaining actions were not flushed"); + action_rx.recv().await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_secs(2), action_rx.recv()) + .await + .expect("worker released its senders") + .is_none() + ); + } + + #[tokio::test] + async fn shutdown_interrupts_a_backpressured_startup_batch() { + stop_during_startup_batch(false).await; + } + + #[tokio::test] + async fn closed_consumer_stops_a_backpressured_startup_batch() { + stop_during_startup_batch(true).await; + } + #[test] fn root_event_repo_ref_requires_a_non_empty_coordinate() { let keys = Keys::generate();