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();