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
This commit is contained in:
DanConwayDev
2026-10-02 13:27:00 +00:00
parent 0c19c87856
commit 6b373e9ecd
2 changed files with 125 additions and 51 deletions
@@ -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`)
+108 -40
View File
@@ -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<broadcast::Receiver<()>>) {
pub async fn run(self, mut shutdown_rx: Option<broadcast::Receiver<()>>) {
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,15 +505,15 @@ 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 {
@@ -508,41 +530,14 @@ impl SelfSubscriber {
}
_ = timer.tick() => {
if !pending.is_empty() {
self.process_batch(&mut pending).await;
}
}
_ = 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");
if let LoopControl::Break = self.process_batch(&mut pending).await {
break;
}
}
}
_ = timer.tick() => {
if !pending.is_empty() {
self.process_batch(&mut pending).await;
}
}
}
}
}
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();