From e942199a5532e2074497c9b44b455c78dcdef51c Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 5 Aug 2026 07:50:45 +0000 Subject: [PATCH] fix(sync): bound concurrent transient REQ subscriptions per relay connection Historic REQ+EOSE sync subscribed every byte-budgeted filter group of a batch in one loop and left them all awaiting EOSE concurrently, as did the REQ+EOSE fallback, negentropy ID fetches and missing-event retries. On strfry-family relays these REQs share maxSubsPerConnection with negentropy views and live subscriptions; nos.lol (budget 20) answered each gitnostr.com startup with a burst of 'ERROR: too many concurrent REQs' NOTICEs (6 in one second at the 2026-08-05 06:41 startup; 7,739 over the prior three days). A per-connection semaphore (5 permits, shared across clones) now gates every auto-close subscription inside subscribe_filters: the permit is registered against the subscription id on success and released when the connection's own event loop sees that subscription's EOSE or CLOSED frame - before forwarding the notification, so release never depends on downstream channel consumers. This keeps a plain blocking acquire deadlock-free even though the SyncManager actor both creates subscriptions and processes EOSE: bursts pipeline at five in flight, matching the precedent of the actor already stalling inline for negentropy batches. Permits are additionally freed on unsubscribe, disconnect, event-loop termination, and by a 30-second watchdog so a relay that never answers cannot starve later subscriptions. The negentropy semaphore stays separate (different lifetimes); live subscriptions are not gated. Budget: 4 NEG + 5 REQ + 2 margin leaves at least nine slots of the tightest observed budget (20) for live subscriptions. Scheduling is deliberately per-REQ rather than per-core-filter: a paginating filter chain holds no slot between pages, so queued groups interleave breadth-first. Relay-visible concurrency is identical either way and pagination chains have no durable identity across disconnects. Reproduction: a new proxy fixture mimics the strfry limit (rejecting REQs beyond 5 with the production NOTICE, exempting limit:0 live subscriptions, delaying EOSE so REQs provably overlap), and a scenario test syncs 2500 root events into persistent storage, restarts the relay - startup recomputes filters from the full index, the shape that bursts in production - and asserts zero rejections with overlap retained. Unfixed: 3 REQs rejected (opened 91, peak pinned at the limit). Fixed: zero rejections, peak <= 5, sync completes. The TestRelay fixture gains same-port restart support with persistent LMDB storage for this. Deliberately excluded: the unified budget ledger (NIP-11-aware B/M, per-query result caps), raising the 300-ID exact-ID chunks, and any configuration surface for the bound. Validated with the full test suite (nix develop -c cargo test). --- docs/explanation/sync-scaling-constraints.md | 8 + src/sync/relay_connection.rs | 146 ++++++++- tests/common/mod.rs | 1 + tests/common/port.rs | 14 + tests/common/relay.rs | 68 ++++ tests/common/req_limiting_proxy.rs | 324 +++++++++++++++++++ tests/sync.rs | 1 + tests/sync/mod.rs | 1 + tests/sync/req_concurrency.rs | 289 +++++++++++++++++ 9 files changed, 851 insertions(+), 1 deletion(-) create mode 100644 tests/common/req_limiting_proxy.rs create mode 100644 tests/sync/req_concurrency.rs diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 7bb2d83..c3d7aee 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -171,6 +171,13 @@ So concurrency is not a free scaling axis; it is the residual of the ledger: - Rounds queue behind a per-connection semaphore; each completion releases the next. No timed batches or sleeps — throughput degrades smoothly instead of bursting into rejections. +- Transient REQ+EOSE subscriptions — historic sync groups, fallback + filters, exact-ID fetches, retries, and pagination pages — queue behind + their own per-connection semaphore (5 permits): a permit is acquired when + the auto-close REQ is sent and released when its EOSE or CLOSED arrives + (with a 30 s watchdog against relays that never answer). Live + subscriptions are not gated. 4 NEG + 5 REQ + 2 margin leaves at least + nine slots of the B = 20 floor for live subscriptions. - Permit acquisition checks relay health first: while a rate-limit or transient-failure cooldown is active, queued rounds take the REQ+EOSE fallback path (which is itself budget-accounted) instead of firing into a @@ -282,6 +289,7 @@ with the heuristics demoted to backstop. | Lever | Status | | --- | --- | | 3 — bounded NEG concurrency | Stabilisation cycle 3 (in flight) | +| 3 — bounded transient REQ+EOSE concurrency | Landed with cycle 3 (same PR) | | 1 + 2 — byte-budgeted chunking and REQ packing | Landed with cycle 3 (same PR) | | Ledger unification (live + historic + fallback against one budget, NIP-11-aware B/M) | Design accepted here; implement after cycles 3–4 | | 4 — multi-connection sharding | Deferred until a relay's target set approaches the single-connection ceiling | diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index cdedd76..708e222 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -56,6 +56,39 @@ const NEGENTROPY_TRANSIENT_BACKOFF: [Duration; 4] = [ /// docs/explanation/sync-scaling-constraints.md. const MAX_CONCURRENT_NEG_DIFFS: usize = 4; +/// Maximum concurrent transient (auto-close) REQ subscriptions per relay +/// connection. +/// +/// Historic REQ+EOSE sync, negentropy ID fetches, missing-event retries, +/// fallback subscriptions and pagination pages each open an auto-close REQ +/// that stays open until EOSE. Relays bound concurrent subscriptions per +/// connection (`maxSubsPerConnection` on strfry-family relays; nos.lol +/// advertises 20, shared between live subscriptions, negentropy views and +/// REQs), so an unbounded startup batch bursts past the tightest common +/// budget — observed in production as "ERROR: too many concurrent REQs" +/// NOTICEs from nos.lol. Five permits alongside the four negentropy +/// permits and a two-slot margin leaves at least nine slots of that +/// budget for live subscriptions; see +/// docs/explanation/sync-scaling-constraints.md. +const MAX_CONCURRENT_TRANSIENT_REQS: usize = 5; + +/// Upper bound on how long a transient-REQ permit may be held. +/// +/// Permits are normally released when the subscription's EOSE (or CLOSED) +/// arrives. A relay that never answers would otherwise pin its permits +/// forever and starve every later transient subscription on the +/// connection, so a watchdog reclaims the permit after this deadline +/// (twice the negentropy round timeout, generous for large REQ pages). +const TRANSIENT_REQ_PERMIT_TIMEOUT: Duration = Duration::from_secs(30); + +/// Permits held by in-flight transient REQ subscriptions, keyed by +/// subscription id and shared across connection clones. +type TransientReqPermitMap = std::sync::Arc< + std::sync::Mutex< + std::collections::HashMap, + >, +>; + /// How a failed negentropy diff should affect future NIP-77 attempts. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum NegentropyFailure { @@ -175,6 +208,13 @@ pub struct RelayConnection { /// Bounds concurrent negentropy diff rounds on this connection /// (shared across clones; see [`MAX_CONCURRENT_NEG_DIFFS`]) neg_diff_permits: std::sync::Arc, + /// Bounds concurrent transient (auto-close) REQ subscriptions on this + /// connection (shared across clones; see + /// [`MAX_CONCURRENT_TRANSIENT_REQS`]) + transient_req_permits: std::sync::Arc, + /// Permits held by open transient subscriptions, keyed by subscription + /// id; released on EOSE/CLOSED, connection teardown, or the watchdog + transient_req_permits_held: TransientReqPermitMap, } impl RelayConnection { @@ -231,6 +271,12 @@ impl RelayConnection { neg_diff_permits: std::sync::Arc::new(tokio::sync::Semaphore::new( MAX_CONCURRENT_NEG_DIFFS, )), + transient_req_permits: std::sync::Arc::new(tokio::sync::Semaphore::new( + MAX_CONCURRENT_TRANSIENT_REQS, + )), + transient_req_permits_held: std::sync::Arc::new(std::sync::Mutex::new( + std::collections::HashMap::new(), + )), } } @@ -266,6 +312,12 @@ impl RelayConnection { neg_diff_permits: std::sync::Arc::new(tokio::sync::Semaphore::new( MAX_CONCURRENT_NEG_DIFFS, )), + transient_req_permits: std::sync::Arc::new(tokio::sync::Semaphore::new( + MAX_CONCURRENT_TRANSIENT_REQS, + )), + transient_req_permits_held: std::sync::Arc::new(std::sync::Mutex::new( + std::collections::HashMap::new(), + )), } } @@ -436,6 +488,12 @@ 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.release_transient_req_permit(&owned_sub_id); if event_sender .send(RelayEvent::EndOfStoredEvents(owned_sub_id)) .await @@ -469,7 +527,13 @@ impl RelayConnection { let _ = event_sender.send(RelayEvent::Notice(msg.to_string())).await; // Don't break - continue processing events } - RelayMessage::Closed { message: msg, .. } => { + RelayMessage::Closed { + subscription_id, + message: msg, + } => { + // A CLOSED subscription can no longer produce EOSE; + // free its transient permit if one is held. + self.release_transient_req_permit(subscription_id.as_ref()); tracing::info!(relay = %url, message = %msg, "Relay closed subscription"); let _ = event_sender.send(RelayEvent::Closed(msg.to_string())).await; // Don't break - CLOSED is subscription-specific, not connection-specific @@ -514,6 +578,10 @@ impl RelayConnection { } } + // The connection is going away; every outstanding transient + // subscription dies with it, so free their permits. + self.clear_transient_req_permits(); + tracing::debug!(relay = %url, "Notification stream ended, event loop terminated"); } @@ -560,6 +628,25 @@ impl RelayConnection { "subscribe_filters called" ); + // Transient (auto-close) subscriptions share a bounded number of + // per-connection slots so historic bursts queue instead of + // exceeding relay subscription budgets. The permit is registered + // against the subscription id on success and released when the + // subscription's EOSE or CLOSED arrives (see run_event_loop). + let permit = if auto_close { + match self.transient_req_permits.clone().acquire_owned().await { + Ok(permit) => Some(permit), + Err(_) => { + return Err(format!( + "Failed to subscribe on {}: transient permit semaphore closed", + self.url + )); + } + } + } else { + None + }; + let output = if auto_close { self.client .subscribe(filters) @@ -582,6 +669,10 @@ impl RelayConnection { return Err(format!("Failed to subscribe on {}: {}", self.url, failures)); } + if let Some(permit) = permit { + self.hold_transient_req_permit(output.value.clone(), permit); + } + tracing::debug!( relay = %self.url, subscription_id = %output.value, @@ -648,6 +739,7 @@ impl RelayConnection { /// Disconnect from the relay pub async fn disconnect(&self) { + self.clear_transient_req_permits(); self.client.disconnect().await; tracing::debug!(relay = %self.url, "Disconnected from relay"); } @@ -658,12 +750,64 @@ impl RelayConnection { /// with consolidated filters. This sends CLOSE messages for all active /// subscriptions on the relay. pub async fn unsubscribe_all(&self) { + self.clear_transient_req_permits(); if let Err(e) = self.client.unsubscribe_all().await { tracing::debug!(relay = %self.url, error = %e, "Failed to unsubscribe from all subscriptions"); } tracing::debug!(relay = %self.url, "Unsubscribed from all subscriptions"); } + /// Register a held transient-REQ permit for `sub_id` and start its + /// watchdog. + fn hold_transient_req_permit( + &self, + sub_id: SubscriptionId, + permit: tokio::sync::OwnedSemaphorePermit, + ) { + self.transient_req_permits_held + .lock() + .expect("transient permit map poisoned") + .insert(sub_id.clone(), permit); + + // Watchdog: a relay that never sends EOSE (or CLOSED) must not pin + // its permits forever. Dropping the permit only frees our own + // slot; the subscription itself remains subject to the normal + // batch and disconnect cleanup paths. + let held = std::sync::Arc::clone(&self.transient_req_permits_held); + let relay_url = self.url.clone(); + tokio::spawn(async move { + tokio::time::sleep(TRANSIENT_REQ_PERMIT_TIMEOUT).await; + let reclaimed = held + .lock() + .expect("transient permit map poisoned") + .remove(&sub_id) + .is_some(); + if reclaimed { + tracing::debug!( + relay = %relay_url, + sub_id = %sub_id, + "Transient REQ permit reclaimed by watchdog (no EOSE within deadline)" + ); + } + }); + } + + /// Release the transient-REQ permit held for `sub_id`, if any. + fn release_transient_req_permit(&self, sub_id: &SubscriptionId) { + self.transient_req_permits_held + .lock() + .expect("transient permit map poisoned") + .remove(sub_id); + } + + /// Release every held transient-REQ permit (connection teardown). + fn clear_transient_req_permits(&self) { + self.transient_req_permits_held + .lock() + .expect("transient permit map poisoned") + .clear(); + } + // ========================================================================= // NIP-77 Negentropy Support // ========================================================================= diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 02f2a73..1259145 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -7,6 +7,7 @@ pub mod git_server; pub mod mock_relay; pub mod neg_limiting_proxy; pub mod nip09_helpers; +pub mod req_limiting_proxy; pub mod port; pub mod purgatory_helpers; pub mod relay; diff --git a/tests/common/port.rs b/tests/common/port.rs index efc16b7..042f6ec 100644 --- a/tests/common/port.rs +++ b/tests/common/port.rs @@ -69,6 +69,20 @@ pub struct PortReservation { _listener: TcpListener, } +/// Reserve a specific loopback port. +/// +/// Used by restart flows that must reuse an address already embedded in +/// published events (e.g. a repository announcement naming the relay's +/// domain). Fails while the port is still bound — callers should wait +/// for the previous holder to exit and retry within a bounded deadline. +pub fn reserve_specific(port: u16) -> std::io::Result { + let listener = TcpListener::bind(("127.0.0.1", port))?; + Ok(PortReservation { + port, + _listener: listener, + }) +} + impl PortReservation { /// The kernel-assigned loopback port number held by this reservation. pub fn port(&self) -> u16 { diff --git a/tests/common/relay.rs b/tests/common/relay.rs index 98e807d..dc0c7eb 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -62,6 +62,9 @@ pub struct TestRelay { _relay_data_dir: Option, /// Path to relay data directory (for test assertions) relay_data_path: PathBuf, + /// Options used at start, retained so [`Self::restart`] can respawn + /// an identical relay on the same port. + options: RelayOptions, } /// Options that the various `start*` constructors fan out into a single @@ -384,6 +387,30 @@ impl TestRelay { .await } + /// Start a syncing relay on a caller-reserved port with persistent + /// LMDB storage in explicit directories, so [`Self::restart`] can + /// resume from the same data on the same port. + pub async fn start_on_reservation_persistent_sync( + reservation: PortReservation, + bootstrap_relay_url: Option, + disable_negentropy: bool, + git_data_path: PathBuf, + relay_data_path: PathBuf, + ) -> Self { + Self::start_internal( + reservation, + RelayOptions { + bootstrap_relay_url, + disable_negentropy, + lmdb_backend: true, + git_data_path: Some(git_data_path), + relay_data_path: Some(relay_data_path), + ..RelayOptions::default() + }, + ) + .await + } + /// Start a relay with every configurable option, on a pre-reserved port. /// /// Prefer the narrower constructors above — this exists so the option @@ -621,6 +648,7 @@ impl TestRelay { git_data_path, _relay_data_dir: relay_data_dir, relay_data_path, + options: options.clone(), }; match relay.wait_for_ready_or_early_exit().await { @@ -650,6 +678,46 @@ impl TestRelay { PathBuf::from(format!("/tmp/relay-{}.log", self.port)) } + /// Stop the relay process and start a fresh one on the same port with + /// the same options and data directories. + /// + /// Only meaningful for relays started with persistent (caller-owned) + /// data directories: an in-memory or tempdir-backed relay would + /// restart empty. The same port is re-reserved so addresses embedded + /// in already-published events stay valid. + pub async fn restart(mut self) -> Self { + let port = self.port; + let options = self.options.clone(); + + let _ = self.process.kill(); + let _ = self.process.wait(); + drop(self); + + // The kernel frees the port once the child is fully gone; retry + // the specific-port reservation within a bounded deadline. + let deadline = Instant::now() + Duration::from_secs(10); + let reservation = loop { + match port::reserve_specific(port) { + Ok(reservation) => break reservation, + Err(error) => { + assert!( + Instant::now() < deadline, + "port {port} not released by stopped relay within deadline: {error}" + ); + sleep(Duration::from_millis(50)).await; + } + } + }; + + match Self::try_start_once(reservation, &options).await { + StartOutcome::Ready(relay) => relay, + StartOutcome::EarlyExit { status } => panic!( + "restarted ngit-grasp subprocess exited early (status: {status:?}); \ + check /tmp/relay-{port}.log" + ), + } + } + /// Get the git data directory path /// /// This is useful for test assertions that need to verify diff --git a/tests/common/req_limiting_proxy.rs b/tests/common/req_limiting_proxy.rs new file mode 100644 index 0000000..df30fc1 --- /dev/null +++ b/tests/common/req_limiting_proxy.rs @@ -0,0 +1,324 @@ +//! Concurrency-Limiting REQ Proxy for Sync Tests +//! +//! A transparent WebSocket proxy that sits between a syncing relay and its +//! bootstrap relay and enforces a strfry-style bound on concurrent transient +//! REQ subscriptions per connection. strfry counts REQ subscriptions against +//! `maxSubsPerConnection` and answers excess REQ frames with +//! `NOTICE ERROR: too many concurrent REQs` — the production behaviour +//! observed from nos.lol (which advertises `max_subscriptions: 20`) during +//! gitnostr.com startup bursts. +//! +//! The proxy models *awaiting-EOSE* occupancy — the quantity the client-side +//! transient-REQ bound controls: a REQ joins the active set when it is +//! opened and leaves it when the client sends CLOSE or the backend answers +//! EOSE, whichever comes first. Live subscriptions (every filter carrying +//! `limit: 0`) are exempt; they are persistent by design and budgeted +//! separately. +//! +//! The proxy additionally: +//! - records the peak number of concurrently open transient REQs, so tests +//! can assert the syncing relay's concurrency bound end to end; +//! - delays backend EOSE frames for counted REQs by a fixed interval, so +//! REQs opened together provably overlap instead of racing the loopback +//! round-trip; +//! - rejects NEG-OPEN frames as unsupported, so a syncing relay with +//! negentropy enabled still exercises the REQ+EOSE path. +//! +//! # Usage +//! +//! ```ignore +//! let source = TestRelay::start().await; +//! let proxy = ReqLimitingProxy::start(source.url(), 5).await; +//! let syncing = TestRelay::start_with_sync_no_negentropy(Some(proxy.url().into())).await; +//! // ... assert proxy.rejected_count() == 0 && proxy.peak_concurrent() <= 5 ... +//! ``` + +use std::collections::HashSet; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +use futures_util::{SinkExt, StreamExt}; +use tokio::net::TcpListener; +use tokio::sync::oneshot; +use tokio_tungstenite::tungstenite::Message; + +/// Fixed delay applied to backend EOSE frames for counted REQs. +/// +/// Long enough that a burst of REQ frames sent together is observed before +/// any subscription can complete; short enough to keep tests fast. +const REQ_RESPONSE_DELAY: Duration = Duration::from_millis(150); + +/// WebSocket proxy bounding concurrent transient REQs like a strfry relay. +pub struct ReqLimitingProxy { + url: String, + peak: Arc, + rejected: Arc, + opened: Arc, + shutdown_tx: Option>, + handle: Option>, +} + +impl ReqLimitingProxy { + /// Start a proxy on a random loopback port, forwarding to `backend_url` + /// and rejecting transient REQ frames that would exceed `limit` + /// concurrent subscriptions on a connection. + pub async fn start(backend_url: &str, limit: usize) -> Self { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("ReqLimitingProxy failed to bind"); + let port = listener + .local_addr() + .expect("ReqLimitingProxy local_addr") + .port(); + + let peak = Arc::new(AtomicUsize::new(0)); + let rejected = Arc::new(AtomicUsize::new(0)); + let opened = Arc::new(AtomicUsize::new(0)); + let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>(); + + let backend_url = backend_url.to_string(); + let accept_peak = peak.clone(); + let accept_rejected = rejected.clone(); + let accept_opened = opened.clone(); + + let handle = tokio::spawn(async move { + loop { + tokio::select! { + accepted = listener.accept() => { + let Ok((stream, _)) = accepted else { break }; + let backend_url = backend_url.clone(); + let peak = accept_peak.clone(); + let rejected = accept_rejected.clone(); + let opened = accept_opened.clone(); + tokio::spawn(async move { + if let Err(error) = proxy_connection( + stream, + &backend_url, + limit, + peak, + rejected, + opened, + ) + .await + { + // Disconnects mid-test are expected; log only. + eprintln!("ReqLimitingProxy connection ended: {error}"); + } + }); + } + _ = &mut shutdown_rx => break, + } + } + }); + + Self { + url: format!("ws://127.0.0.1:{port}"), + peak, + rejected, + opened, + shutdown_tx: Some(shutdown_tx), + handle: Some(handle), + } + } + + /// The ws:// URL the syncing relay should use as its bootstrap relay. + pub fn url(&self) -> &str { + &self.url + } + + /// Highest number of transient REQs observed open at once on any + /// connection. + pub fn peak_concurrent(&self) -> usize { + self.peak.load(Ordering::Relaxed) + } + + /// Number of REQ frames rejected for exceeding the limit. + pub fn rejected_count(&self) -> usize { + self.rejected.load(Ordering::Relaxed) + } + + /// Total transient REQ frames accepted and forwarded to the backend. + pub fn opened_count(&self) -> usize { + self.opened.load(Ordering::Relaxed) + } + + /// Stop the proxy. + pub async fn stop(mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + if let Some(handle) = self.handle.take() { + let _ = handle.await; + } + } +} + +impl Drop for ReqLimitingProxy { + fn drop(&mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + } +} + +/// Forward one client connection to the backend, bounding concurrent +/// transient REQs and recording the observed peak. +/// +/// A single loop owns both writers so rejections can be answered directly +/// to the client without forwarding to the backend. +async fn proxy_connection( + client_stream: tokio::net::TcpStream, + backend_url: &str, + limit: usize, + peak: Arc, + rejected: Arc, + opened: Arc, +) -> Result<(), String> { + let client_ws = tokio_tungstenite::accept_async(client_stream) + .await + .map_err(|e| format!("client handshake failed: {e}"))?; + let (backend_ws, _) = tokio_tungstenite::connect_async(backend_url) + .await + .map_err(|e| format!("backend connect failed: {e}"))?; + + let (mut client_tx, mut client_rx) = client_ws.split(); + let (mut backend_tx, mut backend_rx) = backend_ws.split(); + + // Transient REQ subscription ids currently open on this connection. + let mut active: HashSet = HashSet::new(); + + loop { + tokio::select! { + message = client_rx.next() => { + let Some(message) = message else { break }; + let message = message.map_err(|e| format!("client read: {e}"))?; + if let Message::Text(text) = &message { + match parse_client_frame(text.as_str()) { + Some(ClientFrame::Req(subid)) => { + if active.len() >= limit { + rejected.fetch_add(1, Ordering::Relaxed); + client_tx + .send(Message::text( + r#"["NOTICE","ERROR: too many concurrent REQs"]"# + .to_string(), + )) + .await + .map_err(|e| format!("client write: {e}"))?; + client_tx + .send(Message::text(format!( + r#"["CLOSED","{subid}","blocked: too many concurrent REQs"]"# + ))) + .await + .map_err(|e| format!("client write: {e}"))?; + continue; + } + active.insert(subid); + opened.fetch_add(1, Ordering::Relaxed); + peak.fetch_max(active.len(), Ordering::Relaxed); + } + Some(ClientFrame::LiveReq) => { + // Persistent live subscription (all filters + // limit:0) - budgeted separately, not counted. + } + Some(ClientFrame::Close(subid)) => { + active.remove(&subid); + } + Some(ClientFrame::NegOpen(subid)) => { + // This proxy models a relay without NIP-77 so + // historic sync exercises the REQ+EOSE path. + client_tx + .send(Message::text(format!( + r#"["NEG-ERR","{subid}","blocked: negentropy not supported"]"# + ))) + .await + .map_err(|e| format!("client write: {e}"))?; + continue; + } + None => {} + } + } + backend_tx + .send(message) + .await + .map_err(|e| format!("backend write: {e}"))?; + } + message = backend_rx.next() => { + let Some(message) = message else { break }; + let message = message.map_err(|e| format!("backend read: {e}"))?; + if let Message::Text(text) = &message { + if let Some(subid) = parse_backend_eose(text.as_str()) { + if active.remove(&subid) { + // Hold EOSE frames briefly so simultaneously + // opened REQs provably overlap at the proxy. + tokio::time::sleep(REQ_RESPONSE_DELAY).await; + } + } + } + client_tx + .send(message) + .await + .map_err(|e| format!("client write: {e}"))?; + } + } + } + + Ok(()) +} + +/// Parsed shape of a client-to-relay frame the proxy cares about. +enum ClientFrame { + /// `["REQ", , ]` - transient, counted. + Req(String), + /// `["REQ", ...]` where every filter carries `limit: 0` - a persistent + /// live subscription, exempt from the transient bound. + LiveReq, + /// `["CLOSE", ]` - ends the subscription. + Close(String), + /// `["NEG-OPEN", , ...]` - rejected as unsupported. + NegOpen(String), +} + +/// Classify a client text frame if the proxy tracks it, else `None`. +fn parse_client_frame(text: &str) -> Option { + if !(text.starts_with("[\"REQ") + || text.starts_with("[\"CLOSE") + || text.starts_with("[\"NEG-OPEN")) + { + return None; + } + let value = serde_json::from_str::(text).ok()?; + let array = value.as_array()?; + let kind = array.first()?.as_str()?; + let subid = array.get(1)?.as_str()?.to_string(); + match kind { + "REQ" => { + let filters = &array[2..]; + let live = !filters.is_empty() + && filters + .iter() + .all(|f| f.get("limit").and_then(|l| l.as_u64()) == Some(0)); + if live { + Some(ClientFrame::LiveReq) + } else { + Some(ClientFrame::Req(subid)) + } + } + "CLOSE" => Some(ClientFrame::Close(subid)), + "NEG-OPEN" => Some(ClientFrame::NegOpen(subid)), + _ => None, + } +} + +/// Extract the subscription id from a backend `["EOSE", ]` frame. +fn parse_backend_eose(text: &str) -> Option { + if !text.starts_with("[\"EOSE") { + return None; + } + let value = serde_json::from_str::(text).ok()?; + let array = value.as_array()?; + if array.first()?.as_str()? != "EOSE" { + return None; + } + Some(array.get(1)?.as_str()?.to_string()) +} diff --git a/tests/sync.rs b/tests/sync.rs index 15b708f..f9316a9 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -39,5 +39,6 @@ mod sync { pub mod maintainer_reprocessing; pub mod metrics; pub mod neg_concurrency; + pub mod req_concurrency; pub mod tag_variations; } diff --git a/tests/sync/mod.rs b/tests/sync/mod.rs index 15f0d4a..a351a62 100644 --- a/tests/sync/mod.rs +++ b/tests/sync/mod.rs @@ -136,4 +136,5 @@ pub mod live_sync; pub mod maintainer_reprocessing; pub mod metrics; pub mod neg_concurrency; +pub mod req_concurrency; pub mod tag_variations; \ No newline at end of file diff --git a/tests/sync/req_concurrency.rs b/tests/sync/req_concurrency.rs new file mode 100644 index 0000000..022a969 --- /dev/null +++ b/tests/sync/req_concurrency.rs @@ -0,0 +1,289 @@ +//! Transient REQ Concurrency Bound Tests +//! +//! Regression coverage for a production failure observed on gitnostr.com: +//! historic REQ+EOSE sync subscribed every byte-budgeted filter group of a +//! batch in one fast loop and left them all open concurrently until EOSE. +//! On strfry-family relays REQ subscriptions share `maxSubsPerConnection` +//! with negentropy views and live subscriptions; nos.lol (budget 20) +//! answered each gitnostr.com startup with a burst of +//! "ERROR: too many concurrent REQs" NOTICEs (6 in one second at the +//! 2026-08-05 startup; 7,739 over the prior three days). +//! +//! The scenario drives the real sync path end to end: a genuine ngit-grasp +//! source relay sits behind a proxy that mimics the strfry limit, rejecting +//! REQ frames beyond 5 concurrent transient subscriptions. The syncing +//! relay runs with negentropy disabled so historic sync takes the REQ+EOSE +//! path, and the watched set is sized so that one startup batch needs more +//! REQ groups than the limit (2500 root events → six byte-budgeted chunks → +//! 18 e/E/q filters; two full 32 KB filters fill a 96 KB REQ group, giving +//! ~8 groups). +//! +//! The burst only occurs when startup recomputes filters from an +//! already-populated index — during initial discovery, root events trickle +//! in through short batching windows and batches stay small. The scenario +//! therefore syncs everything into persistent storage first, then restarts +//! the relay, exactly the production startup shape. Unbounded subscription +//! creation draws rejections at that startup; the bounded client must draw +//! none while still overlapping REQs and completing the sync. +//! +//! See docs/explanation/sync-scaling-constraints.md for the budget model. + +use std::time::Duration; + +use nostr_sdk::prelude::*; + +use crate::common::purgatory_helpers::{ + create_state_event, create_test_repo_with_commit, push_to_relay, CommitVariant, +}; +use crate::common::req_limiting_proxy::ReqLimitingProxy; +use crate::common::sync_helpers::{repo_coord, send_to_relay, wait_for_event_on_relay, TestClient}; +use crate::common::{port, TestRelay}; + +/// Concurrency limit enforced by the proxy. +/// +/// Matches `MAX_CONCURRENT_TRANSIENT_REQS` in `sync::relay_connection`: the +/// test relay must stay within its own advertised bound, so a relay +/// enforcing exactly that bound must never reject a REQ. +const PROXY_REQ_LIMIT: usize = 5; + +/// Root events to seed. Byte-budgeted chunking fits ~489 hex event IDs per +/// filter (32 KB of serialized tag values), so 2500 IDs split into six +/// chunks; six chunks × three tag variants (e/E/q) give 18 filters, and two +/// full 32 KB filters fill a 96 KB REQ group — roughly eight groups in a +/// single historic batch, more concurrent REQ subscriptions than the bound, +/// so bounded scheduling is required to avoid rejections. +const ISSUE_COUNT: usize = 2500; + +/// Wait until the proxy has seen at least `min_total` transient REQs +/// (opened + rejected) and the count has been stable for `stable_for`. +/// +/// Returns the final (opened, rejected) pair, or `None` on deadline. +async fn wait_for_req_quiescence( + proxy: &ReqLimitingProxy, + min_total: usize, + stable_for: Duration, + deadline: Duration, +) -> Option<(usize, usize)> { + let end = tokio::time::Instant::now() + deadline; + let mut last_total = 0usize; + let mut stable_since = tokio::time::Instant::now(); + loop { + let opened = proxy.opened_count(); + let rejected = proxy.rejected_count(); + let total = opened + rejected; + if total != last_total { + last_total = total; + stable_since = tokio::time::Instant::now(); + } else if total >= min_total && stable_since.elapsed() >= stable_for { + return Some((opened, rejected)); + } + if tokio::time::Instant::now() >= end { + return None; + } + tokio::time::sleep(Duration::from_millis(200)).await; + } +} + +/// Scenario: +/// 1. Source relay hosts one repository (announcement + state + git data) +/// and 2500 historic issues, each a root event of that repository. +/// 2. A proxy in front of the source enforces the strfry-style limit of 5 +/// concurrent transient REQs and records peak concurrency and +/// rejections. +/// 3. The syncing relay (negentropy disabled, persistent storage) +/// bootstraps through the proxy and syncs all issues. +/// 4. The syncing relay restarts. Startup recomputes sync filters from the +/// persisted index of 2500 root events, producing one historic batch +/// with ~8 REQ groups — more than the limit. +/// 5. The relay must complete that startup sync without a single rejection +/// while still keeping REQ subscriptions overlapping. +#[tokio::test] +async fn startup_historic_sync_stays_within_relay_req_concurrency_limit() { + // 1. Pre-allocate the syncing relay port for announcement tags. + let syncing_reservation = port::reserve_port(); + let syncing_domain = format!("127.0.0.1:{}", syncing_reservation.port()); + + // 2. Source relay with the strfry-style limiting proxy in front. The + // announcement's relays tag lists the proxy (so the syncing relay + // targets the proxied connection), not the source itself, so the + // source runs in archive-all mode to accept it. + let source = TestRelay::start_with_archive_config(true, false).await; + let proxy = ReqLimitingProxy::start(source.url(), PROXY_REQ_LIMIT).await; + + // 3. One hosted repository: announcement + state event + git data. The + // relays tag lists the proxy so the syncing relay targets the proxied + // connection for this repository's sync work. + let git_temp_dir = tempfile::tempdir().expect("create temp dir for git repo"); + let commit_hash = create_test_repo_with_commit(git_temp_dir.path(), CommitVariant::StateTest) + .expect("create test git repo"); + + let keys = Keys::generate(); + let npub = keys.public_key().to_bech32().expect("npub"); + let clone_urls = vec![ + format!("http://{}/{}/req-repo.git", source.domain(), npub), + format!("http://{}/{}/req-repo.git", syncing_domain, npub), + ]; + let relay_urls = vec![proxy.url().to_string(), format!("ws://{}", syncing_domain)]; + + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "Repository state") + .tags(vec![ + Tag::identifier("req-repo"), + Tag::custom("clone", clone_urls.clone()), + Tag::custom("relays", relay_urls.clone()), + ]) + .finalize(&keys) + .expect("sign repo announcement"); + let state_event = create_state_event( + &keys, + "req-repo", + &[("main", &commit_hash)], + &[], + &clone_urls.iter().map(|s| s.as_str()).collect::>(), + &relay_urls.iter().map(|s| s.as_str()).collect::>(), + ) + .expect("create state event"); + + send_to_relay(&source, &announcement) + .await + .expect("send announcement to source"); + send_to_relay(&source, &state_event) + .await + .expect("send state event to source"); + push_to_relay(git_temp_dir.path(), &source.domain(), &npub, "req-repo") + .expect("push git data to source relay"); + + // The push releases the announcement from purgatory; wait until the + // source actually serves it before seeding dependent events. + assert!( + wait_for_event_on_relay( + source.url(), + Filter::new().id(announcement.id), + Duration::from_secs(15), + ) + .await, + "announcement should be released from purgatory on the source" + ); + + // 4. Seed historic issues; each becomes a tracked root event once + // synced. One persistent client keeps seeding fast, and distinct + // descending created_at timestamps let REQ+EOSE pagination cursors + // (until = oldest seen) terminate instead of re-fetching the same + // page forever. + let coordinate = repo_coord(&keys, "req-repo"); + let base_created_at = Timestamp::now().as_secs() - ISSUE_COUNT as u64 - 10; + + // The embedded relay rate-limits notes per minute PER CONNECTION with a + // full initial token bucket, so seed through a fresh connection per + // batch, each batch staying under the initial burst allowance. + const SEED_BATCH: usize = 50; + let mut issue_ids = Vec::with_capacity(ISSUE_COUNT); + for batch_start in (0..ISSUE_COUNT).step_by(SEED_BATCH) { + let seed_client = TestClient::new(source.url(), Keys::generate()) + .await + .expect("connect seed client to source"); + for index in batch_start..(batch_start + SEED_BATCH).min(ISSUE_COUNT) { + let issue = EventBuilder::new(Kind::GitIssue, format!("Historic issue {index}")) + .tags(vec![Tag::custom("a", vec![coordinate.clone()])]) + .custom_created_at(Timestamp::from(base_created_at + index as u64)) + .finalize(&keys) + .expect("sign issue event"); + issue_ids.push(issue.id); + seed_client + .send_event(&issue) + .await + .expect("send issue to source"); + } + seed_client.disconnect().await; + } + + // Seeding sanity: the last issue must be queryable on the source before + // the syncing relay is pointed at it. + assert!( + wait_for_event_on_relay( + source.url(), + Filter::new().id(issue_ids[ISSUE_COUNT - 1]), + Duration::from_secs(15), + ) + .await, + "seeded issues should be accepted by the source relay" + ); + + // 5. Phase 1: the syncing relay (negentropy disabled) bootstraps + // through the proxy and syncs everything into persistent storage. + // Root events trickle in through short batching windows here, so + // historic batches stay small; the burst this test guards against + // comes at the next startup. + let syncing_git_dir = tempfile::tempdir().expect("create syncing git dir"); + let syncing_data_dir = tempfile::tempdir().expect("create syncing relay data dir"); + let syncing = TestRelay::start_on_reservation_persistent_sync( + syncing_reservation, + Some(proxy.url().to_string()), + true, + syncing_git_dir.path().to_path_buf(), + syncing_data_dir.path().to_path_buf(), + ) + .await; + + // 6. The issues must arrive via historic sync (sampled ends of the + // range cover the byte-budgeted chunks). + for issue_id in [issue_ids[0], issue_ids[ISSUE_COUNT / 2], issue_ids[ISSUE_COUNT - 1]] { + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(issue_id), + Duration::from_secs(120), + ) + .await, + "issue should reach the syncing relay through historic sync" + ); + } + + // 7. Let phase-1 REQ activity settle, then restart. Startup recomputes + // sync filters from the full persisted index (2500 root events), + // producing one historic batch with ~8 byte-budgeted REQ groups — + // the production startup burst. + wait_for_req_quiescence(&proxy, 1, Duration::from_secs(3), Duration::from_secs(120)) + .await + .expect("phase-1 REQ activity should settle before restart"); + let before_restart = proxy.opened_count() + proxy.rejected_count(); + + let syncing = syncing.restart().await; + + // 8. Wait for the startup burst to complete and settle: the restarted + // relay re-syncs the repository batch and the ~8 root-event groups. + let (opened, rejected) = wait_for_req_quiescence( + &proxy, + before_restart + 7, + Duration::from_secs(3), + Duration::from_secs(120), + ) + .await + .expect("restart REQ burst should complete and settle"); + + // 9. The regression assertions (cumulative across both phases; the + // trickle-shaped phase 1 must stay within the bound too). + assert_eq!( + rejected, 0, + "relay exceeded the per-connection REQ concurrency limit \ + (opened: {opened}, peak: {})", + proxy.peak_concurrent() + ); + assert!( + proxy.peak_concurrent() <= PROXY_REQ_LIMIT, + "peak concurrent transient REQs {} exceeded the bound {PROXY_REQ_LIMIT}", + proxy.peak_concurrent() + ); + assert!( + proxy.peak_concurrent() >= 2, + "transient REQs should still overlap under the bound, peak: {}", + proxy.peak_concurrent() + ); + assert!( + opened > PROXY_REQ_LIMIT, + "scenario must open more REQs than the bound to exercise queueing, opened: {opened}" + ); + + syncing.stop().await; + proxy.stop().await; + source.stop().await; +}