diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index b589eff..a091010 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -304,12 +304,24 @@ slot after the SDK accepts the message; valid production startup pages exceeded subscription. If CLOSE cannot be enqueued, the slot remains held until ordinary connection teardown so local accounting cannot run ahead of the relay. Exact-ID purgatory polling uses the same transient class bound and shared ledger as -historic pagination. Unexpected CLOSED for a persistent live subscription is +historic pagination. Transient subscription IDs and their permits are +registered locally before the REQ is sent; this ordering is required because +an empty or cached response can deliver EOSE/CLOSED before the SDK subscribe +call returns. Subscribe failure rolls that pre-registration back. Unexpected +CLOSED for a persistent live subscription is reported to the manager, which recomputes and transactionally reopens complete live coverage. Each reconnect closes the retired ledger and creates a new generation; queued or late borrowers therefore fail before sending on the new SDK session and cannot inflate or bypass its capacity. +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. + The levers, in the order we reach for them: ### Lever 1: Maximise items per filter (byte-budgeted chunking) diff --git a/src/sync/mod.rs b/src/sync/mod.rs index f062ac5..00bbfe0 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -2972,7 +2972,8 @@ impl SyncManager { } }; - let (event_tx, mut event_rx) = mpsc::channel::(1000); + let (event_tx, mut event_rx) = + mpsc::channel::(relay_connection::RELAY_EVENT_BUFFER_CAPACITY); // Spawn event loop task let relay_url_for_loop = relay_url.to_string(); diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 8606593..271948e 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -96,6 +96,12 @@ const MAX_CONCURRENT_NEG_DIFFS: usize = 4; /// docs/explanation/sync-scaling-constraints.md. const MAX_CONCURRENT_TRANSIENT_REQS: usize = 5; +/// Events and lifecycle notifications waiting for the per-relay processor. +/// +/// Transient permit release has a separate terminal-control listener, so this +/// remains a bounded data-processing queue rather than a lifecycle boundary. +pub(crate) const RELAY_EVENT_BUFFER_CAPACITY: usize = 1000; + /// Upper bound on how long a transient-REQ permit may be held. /// /// Permits are normally released when the subscription's EOSE (or CLOSED) @@ -916,6 +922,46 @@ impl RelayConnection { // In nostr-sdk 0.45 this returns a Stream rather than a broadcast receiver. let mut notifications = relay.notifications(); + // EVENT forwarding can block on the bounded processor queue. Subscribe + // independently to terminal messages so EOSE/CLOSED can still return + // transient ledger slots while the data lane drains. Relay + // notifications are broadcast, so this does not consume messages from + // the processor-facing stream below. + let terminal_connection = self.clone(); + // Construct the receiver before spawning so no fast terminal can race + // task scheduling and arrive before the control lane is subscribed. + let mut terminals = relay.notifications(); + let terminal_listener = tokio::spawn(async move { + while let Some(notification) = terminals.next().await { + match notification { + RelayNotification::Message { message } => match *message { + RelayMessage::EndOfStoredEvents(sub_id) => { + terminal_connection + .close_and_release_transient_req_permit(&sub_id.into_owned()) + .await; + } + RelayMessage::Closed { + subscription_id, .. + } => { + terminal_connection + .release_transient_req_permit(&subscription_id.into_owned()); + } + _ => {} + }, + RelayNotification::RelayStatus { status } + if matches!( + status, + RelayStatus::Disconnected | RelayStatus::Terminated + ) => + { + terminal_connection.clear_subscription_permits(); + break; + } + _ => {} + } + } + }); + tracing::debug!(relay = %url, "Starting event loop with relay-level notifications"); while let Some(notification) = notifications.next().await { @@ -945,13 +991,6 @@ impl RelayConnection { tracing::debug!(relay = %url, sub_id = ?sub_id, "Received EOSE"); // Convert Cow to owned SubscriptionId let owned_sub_id = sub_id.into_owned(); - // Release the transient permit BEFORE forwarding: - // the forward can block on channel backpressure and - // permit release must never depend on downstream - // consumers (the SyncManager actor may itself be - // waiting on a permit). - self.close_and_release_transient_req_permit(&owned_sub_id) - .await; if event_sender .send(RelayEvent::EndOfStoredEvents(owned_sub_id)) .await @@ -989,11 +1028,10 @@ impl RelayConnection { subscription_id, message: msg, } => { - // A CLOSED subscription can no longer produce EOSE; - // free whichever ledger ownership it held. + // The terminal-control listener owns transient release; + // this processor-facing path retains live restoration. let subscription_id = subscription_id.into_owned(); - let live_generation = - self.release_subscription_permits_on_closed(&subscription_id); + let live_generation = self.release_live_req_permit(&subscription_id); tracing::info!(relay = %url, message = %msg, "Relay closed subscription"); let _ = event_sender .send(RelayEvent::Closed { @@ -1044,6 +1082,8 @@ impl RelayConnection { } } + terminal_listener.abort(); + // The connection is going away; every outstanding transient // subscription dies with it, so free their permits. self.clear_subscription_permits(); @@ -1215,19 +1255,43 @@ impl RelayConnection { // Transient permits are acquired through the same session helper; a // reset closes queued acquisitions and the helper checks generation. let retained_filters = filters.clone(); + // The relay can answer an empty or cached query before `subscribe` + // returns. Register transient ownership against a caller-chosen ID + // first so an immediate EOSE/CLOSED cannot race past local accounting. + let transient_sub_id = transient_class.map(|_| SubscriptionId::generate()); + if let (Some(sub_id), Some(permit)) = (&transient_sub_id, transient_permit) { + self.hold_transient_req_permit(sub_id.clone(), permit); + } let output = if transient_class.is_some() { self.client .subscribe(filters) + .with_id( + transient_sub_id + .clone() + .expect("transient subscription has a pre-registered id"), + ) .close_on( SubscribeAutoCloseOptions::default().exit_policy(ReqExitPolicy::ExitOnEOSE), ) .await } else { self.client.subscribe(filters).await - } - .map_err(|e| format!("Failed to subscribe on {}: {}", self.url, e))?; + }; + + let output = match output { + Ok(output) => output, + Err(error) => { + if let Some(sub_id) = &transient_sub_id { + self.release_transient_req_permit(sub_id); + } + return Err(format!("Failed to subscribe on {}: {}", self.url, error)); + } + }; if !output.failed.is_empty() { + if let Some(sub_id) = &transient_sub_id { + self.release_transient_req_permit(sub_id); + } let failures = output .failed .values() @@ -1237,11 +1301,7 @@ impl RelayConnection { return Err(format!("Failed to subscribe on {}: {}", self.url, failures)); } - if transient_class.is_some() { - if let Some(permit) = transient_permit { - self.hold_transient_req_permit(output.value.clone(), permit); - } - } else { + if transient_class.is_none() { let live_permit = live_permit.expect("live subscription acquired a ledger slot"); self.live_req_permits_held .lock() @@ -1499,6 +1559,7 @@ impl RelayConnection { .map(|held| held.generation) } + #[cfg(test)] fn release_subscription_permits_on_closed(&self, sub_id: &SubscriptionId) -> Option { self.release_transient_req_permit(sub_id); self.release_live_req_permit(sub_id) @@ -2107,6 +2168,118 @@ mod tests { other.shutdown(); } + #[tokio::test] + async fn terminal_control_lane_releases_permits_behind_event_backpressure() { + const AUTHORS: usize = 3; + const EVENTS_PER_AUTHOR: usize = 400; + + let relay = LocalRelayBuilder::default().build(); + relay.run().await.expect("start burst relay"); + let mut authors = Vec::with_capacity(AUTHORS); + for author_index in 0..AUTHORS { + let keys = Keys::generate(); + authors.push(keys.public_key()); + for event_index in 0..EVENTS_PER_AUTHOR { + let event = EventBuilder::new( + Kind::TextNote, + format!("burst-{author_index}-{event_index}"), + ) + .finalize(&keys) + .expect("build burst event"); + relay.add_event(event).await.expect("seed burst event"); + } + } + + let connection = RelayConnection::new( + relay.url().await.to_string(), + Keys::generate(), + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.expect("connect burst relay"); + // Keep the data lane full after its first event. The independent + // terminal listener must still observe every EOSE and release slots. + let (event_tx, _event_rx) = tokio::sync::mpsc::channel(1); + let event_loop = tokio::spawn(connection.clone().run_event_loop(event_tx)); + + for author in authors { + connection + .subscribe_filter( + Filter::new().author(author), + TransientRequestClass::NegentropyHydration, + ) + .await + .expect("open bounded burst page"); + } + + tokio::time::timeout(Duration::from_secs(3), async { + loop { + if connection + .transient_req_permits_held + .lock() + .expect("transient permit map poisoned") + .is_empty() + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("the terminal lane must release permits behind blocked EVENT delivery"); + + connection.disconnect().await; + relay.shutdown(); + event_loop.abort(); + } + + #[tokio::test] + async fn immediate_empty_eose_cannot_arrive_before_permit_registration() { + let relay = LocalRelayBuilder::default().build(); + relay.run().await.expect("start empty relay"); + let connection = RelayConnection::new( + relay.url().await.to_string(), + Keys::generate(), + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(3).await.expect("connect empty relay"); + let (event_tx, mut event_rx) = tokio::sync::mpsc::channel(RELAY_EVENT_BUFFER_CAPACITY); + let event_loop = tokio::spawn(connection.clone().run_event_loop(event_tx)); + let drain = tokio::spawn(async move { while event_rx.recv().await.is_some() {} }); + + for _ in 0..25 { + connection + .subscribe_filter( + Filter::new().kind(Kind::Custom(65_535)), + TransientRequestClass::NegentropyHydration, + ) + .await + .expect("open empty transient request"); + } + + tokio::time::timeout(Duration::from_secs(3), async { + loop { + if connection + .transient_req_permits_held + .lock() + .expect("transient permit map poisoned") + .is_empty() + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("every immediate empty EOSE must release its registered permit"); + + connection.disconnect().await; + relay.shutdown(); + event_loop.abort(); + drain.abort(); + } + #[tokio::test] async fn fetch_events_reports_an_unregistered_exact_relay() { let connection = permissive_connection("ws://127.0.0.1:1", Keys::generate());