From a0d8b70bd49adc77fe7a0834b01001d86462046a Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 13:44:54 +0000 Subject: [PATCH] fix(sync): keep lifecycle bursts off the data lane The instrumented archive reproduced relay.ngit.dev's 50,591-ID hydration as 169 subscriptions. Ordinary duplicate processing averaged 0.02-0.03 ms with an empty queue, but the manager's bounded 100-item EOSE inbox filled while the actor held its lock to open the paced batch. The event processor then waited while forwarding EOSE, its 1,000-item EVENT queue filled, and later EOSE messages reached batch accounting 77-82 seconds late. Use non-blocking unbounded actor inboxes for EOSE and CLOSED notifications. Their producers remain bounded by the per-session subscription ledger, so this removes an accidental second capacity limit rather than allowing unbounded wire work. The ordered EVENT queue remains fixed at 1,000 as the peer-facing memory boundary. A regression test queues the observed 169-terminal burst without an actor receiver. This does not raise subscription concurrency, alter query pacing, reorder EVENT processing, or change terminal permit release. The separate rust-nostr terminal listener still closes and releases wire ownership immediately; these inboxes carry later serialized batch and live-coverage accounting. Validation: nix develop -c cargo test --lib (688 passed before the focused channel regression); nix develop -c cargo test --lib lifecycle_inbox_accepts_production_sized_terminal_burst (passed); git diff --check passed. Production validation will repeat the populated relay.ngit.dev reconciliation on this exact tip. --- docs/explanation/sync-scaling-constraints.md | 9 ++-- src/sync/mod.rs | 44 ++++++++++++++++---- 2 files changed, 42 insertions(+), 11 deletions(-) diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index bdf455e..23f4c4c 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -407,9 +407,12 @@ The per-relay event processor retains its 1,000-message bounded data queue. Permit release does not depend on that queue draining: a separate listener on rust-nostr's broadcast relay notifications consumes only EOSE/CLOSED terminals and closes/releases transient ownership. The processor-facing listener still -delivers ordered EVENT and lifecycle work to the sync actor, but sustained page -traffic cannot hide terminal accounting behind EVENT backpressure. The data -queue remains finite as a memory-safety boundary for non-conforming peers. +delivers ordered EVENT and lifecycle work to the sync actor. EOSE and CLOSED +use non-blocking actor inboxes: their production is bounded by the session +subscription ledger, and a large paced historic batch can keep the actor busy +longer than a fixed lifecycle inbox could safely absorb. This prevents actor +backpressure from stopping the ordered EVENT processor, while the EVENT queue +remains finite as the memory-safety boundary for non-conforming peers. The levers, in the order we reach for them: diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 894d015..8bf0499 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1157,6 +1157,13 @@ struct SubscriptionClosedNotification { live_filter_count: Option, } +fn lifecycle_notification_channel() -> ( + tokio::sync::mpsc::UnboundedSender, + tokio::sync::mpsc::UnboundedReceiver, +) { + tokio::sync::mpsc::unbounded_channel() +} + fn is_rate_limit_message(message: &str) -> bool { let message = message.to_lowercase(); (message.contains("rate") && message.contains("limit")) @@ -1787,9 +1794,10 @@ pub struct SyncManager { /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, /// Channel for EOSE notifications (set during run) - eose_tx: Option>, + eose_tx: Option>, /// Serializes CLOSED recovery and pending-batch cleanup through the actor. - subscription_closed_tx: Option>, + subscription_closed_tx: + Option>, /// Returns connection outcomes to the sync actor for serialized state changes. connect_attempt_result_tx: Option>, /// Wakes the actor when a batch completion may unblock consolidation. @@ -3251,9 +3259,13 @@ impl SyncManager { let (disconnect_tx, mut disconnect_rx) = mpsc::channel::(100); // 3. Create EOSE channel for spawned tasks -> manager communication - let (eose_tx, mut eose_rx) = mpsc::channel::(100); + // Lifecycle notifications must never backpressure the ordered EVENT + // processor. Their production is already bounded by the subscription + // ledger, while the actor may legitimately stay busy opening a large, + // paced historic batch for longer than a fixed inbox can absorb. + let (eose_tx, mut eose_rx) = lifecycle_notification_channel::(); let (subscription_closed_tx, mut subscription_closed_rx) = - mpsc::channel::(100); + lifecycle_notification_channel::(); // 4. Connection workers never mutate manager state. Their unbounded // result channel cannot make a completed worker wait behind the actor. @@ -3778,8 +3790,7 @@ impl SyncManager { .send(EoseNotification { relay_url: relay_url_clone.clone(), sub_id, - }) - .await; + }); } RelayEvent::Notice(notice) => { if is_rate_limit_message(¬ice) { @@ -3854,8 +3865,7 @@ impl SyncManager { reason, generation: live_generation, live_filter_count, - }) - .await; + }); } RelayEvent::Shutdown => { tracing::info!(relay = %relay_url_clone, "Relay shutdown detected"); @@ -7165,6 +7175,24 @@ mod tests { ); } + #[test] + fn lifecycle_inbox_accepts_production_sized_terminal_burst() { + let (tx, mut rx) = lifecycle_notification_channel::(); + for index in 0..169 { + tx.send(EoseNotification { + relay_url: "wss://relay.example".to_string(), + sub_id: SubscriptionId::new(format!("hydration-{index}")), + }) + .expect("the actor-side receiver remains alive"); + } + + assert_eq!(rx.len(), 169); + for _ in 0..169 { + rx.try_recv().expect("every terminal remains queued"); + } + assert!(rx.try_recv().is_err()); + } + #[test] fn descendant_frontier_derives_replaceable_and_addressable_coordinates() { let keys = Keys::generate();