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.
This commit is contained in:
DanConwayDev
2026-08-10 13:44:54 +00:00
parent 88ef365397
commit a0d8b70bd4
2 changed files with 42 additions and 11 deletions
+6 -3
View File
@@ -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 Permit release does not depend on that queue draining: a separate listener on
rust-nostr's broadcast relay notifications consumes only EOSE/CLOSED terminals rust-nostr's broadcast relay notifications consumes only EOSE/CLOSED terminals
and closes/releases transient ownership. The processor-facing listener still and closes/releases transient ownership. The processor-facing listener still
delivers ordered EVENT and lifecycle work to the sync actor, but sustained page delivers ordered EVENT and lifecycle work to the sync actor. EOSE and CLOSED
traffic cannot hide terminal accounting behind EVENT backpressure. The data use non-blocking actor inboxes: their production is bounded by the session
queue remains finite as a memory-safety boundary for non-conforming peers. 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: The levers, in the order we reach for them:
+36 -8
View File
@@ -1157,6 +1157,13 @@ struct SubscriptionClosedNotification {
live_filter_count: Option<usize>, live_filter_count: Option<usize>,
} }
fn lifecycle_notification_channel<T>() -> (
tokio::sync::mpsc::UnboundedSender<T>,
tokio::sync::mpsc::UnboundedReceiver<T>,
) {
tokio::sync::mpsc::unbounded_channel()
}
fn is_rate_limit_message(message: &str) -> bool { fn is_rate_limit_message(message: &str) -> bool {
let message = message.to_lowercase(); let message = message.to_lowercase();
(message.contains("rate") && message.contains("limit")) (message.contains("rate") && message.contains("limit"))
@@ -1787,9 +1794,10 @@ pub struct SyncManager {
/// Channel for disconnect notifications (set during run) /// Channel for disconnect notifications (set during run)
disconnect_tx: Option<tokio::sync::mpsc::Sender<DisconnectNotification>>, disconnect_tx: Option<tokio::sync::mpsc::Sender<DisconnectNotification>>,
/// Channel for EOSE notifications (set during run) /// Channel for EOSE notifications (set during run)
eose_tx: Option<tokio::sync::mpsc::Sender<EoseNotification>>, eose_tx: Option<tokio::sync::mpsc::UnboundedSender<EoseNotification>>,
/// Serializes CLOSED recovery and pending-batch cleanup through the actor. /// Serializes CLOSED recovery and pending-batch cleanup through the actor.
subscription_closed_tx: Option<tokio::sync::mpsc::Sender<SubscriptionClosedNotification>>, subscription_closed_tx:
Option<tokio::sync::mpsc::UnboundedSender<SubscriptionClosedNotification>>,
/// Returns connection outcomes to the sync actor for serialized state changes. /// Returns connection outcomes to the sync actor for serialized state changes.
connect_attempt_result_tx: Option<tokio::sync::mpsc::UnboundedSender<ConnectAttemptResult>>, connect_attempt_result_tx: Option<tokio::sync::mpsc::UnboundedSender<ConnectAttemptResult>>,
/// Wakes the actor when a batch completion may unblock consolidation. /// Wakes the actor when a batch completion may unblock consolidation.
@@ -3251,9 +3259,13 @@ impl SyncManager {
let (disconnect_tx, mut disconnect_rx) = mpsc::channel::<DisconnectNotification>(100); let (disconnect_tx, mut disconnect_rx) = mpsc::channel::<DisconnectNotification>(100);
// 3. Create EOSE channel for spawned tasks -> manager communication // 3. Create EOSE channel for spawned tasks -> manager communication
let (eose_tx, mut eose_rx) = mpsc::channel::<EoseNotification>(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::<EoseNotification>();
let (subscription_closed_tx, mut subscription_closed_rx) = let (subscription_closed_tx, mut subscription_closed_rx) =
mpsc::channel::<SubscriptionClosedNotification>(100); lifecycle_notification_channel::<SubscriptionClosedNotification>();
// 4. Connection workers never mutate manager state. Their unbounded // 4. Connection workers never mutate manager state. Their unbounded
// result channel cannot make a completed worker wait behind the actor. // result channel cannot make a completed worker wait behind the actor.
@@ -3778,8 +3790,7 @@ impl SyncManager {
.send(EoseNotification { .send(EoseNotification {
relay_url: relay_url_clone.clone(), relay_url: relay_url_clone.clone(),
sub_id, sub_id,
}) });
.await;
} }
RelayEvent::Notice(notice) => { RelayEvent::Notice(notice) => {
if is_rate_limit_message(&notice) { if is_rate_limit_message(&notice) {
@@ -3854,8 +3865,7 @@ impl SyncManager {
reason, reason,
generation: live_generation, generation: live_generation,
live_filter_count, live_filter_count,
}) });
.await;
} }
RelayEvent::Shutdown => { RelayEvent::Shutdown => {
tracing::info!(relay = %relay_url_clone, "Relay shutdown detected"); 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::<EoseNotification>();
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] #[test]
fn descendant_frontier_derives_replaceable_and_addressable_coordinates() { fn descendant_frontier_derives_replaceable_and_addressable_coordinates() {
let keys = Keys::generate(); let keys = Keys::generate();