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