From a0520c3d6b2ac0ed0eb5eeb9e22a185331c72a3c Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Tue, 4 Aug 2026 21:13:17 +0000 Subject: [PATCH] fix(sync): bound concurrent negentropy rounds per relay connection Historic sync opened one NIP-77 negentropy diff per filter with no bound: handle_add_filters batches launched every diff simultaneously through an unbounded join_all, so a large watched set (146 filters on the bootstrap relay at startup) burst far past relay per-connection subscription budgets. On strfry-family relays negentropy views share maxSubsPerConnection with ordinary subscriptions; nos.lol (budget 20) answered a gitnostr.com startup with 34 'too many concurrent NEG requests' rejections, 61 per-filter timeouts, and 21 failed fallback subscription creations in two minutes (2026-08-04, PR commit ecb6c8b6 soak). The cycle-2 transient cooldown contains the damage; this removes the cause. Approach: a per-connection tokio semaphore (4 permits, shared across clones) gates negentropy_sync_diff. The permit is held for the whole round including the timeout, so concurrent batches targeting the same relay share one bound; queued rounds re-check supports_negentropy() after acquiring, bailing to the per-batch REQ+EOSE fallback without recording a failure when a cooldown started while they waited. Four permits keeps the tightest commonly observed budget (20, shared with live subscriptions) mostly free; constraint research and the budget model are documented in docs/explanation/sync-scaling-constraints.md. Deliberately excluded: bounding REQ+EOSE fallback subscriptions (needs permit lifetimes spanning EOSE handling in the manager loop; deferred to the budget-ledger work), byte-budgeted filter chunking (next cycle), and any configuration surface for the bound. Validation: new scenario test drives two real relays end to end through a proxy enforcing the strfry limit of 4 with delayed NEG responses so rounds provably overlap; a 150-root-event batch needs six rounds. On the unfixed code the proxy rejected 3 rounds (opened 6, peak 4, production signature reproduced); with the fix, zero rejections, peak <= 4 with overlap retained, and sync completes. Full test suite passes; one unrelated grasp06_pr_hosting test was flaky in the full run and passes standalone. --- src/sync/relay_connection.rs | 44 +++++ tests/common/mod.rs | 1 + tests/common/neg_limiting_proxy.rs | 280 +++++++++++++++++++++++++++++ tests/sync.rs | 1 + tests/sync/mod.rs | 1 + tests/sync/neg_concurrency.rs | 226 +++++++++++++++++++++++ 6 files changed, 553 insertions(+) create mode 100644 tests/common/neg_limiting_proxy.rs create mode 100644 tests/sync/neg_concurrency.rs diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 1b88dd2..cdedd76 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -43,6 +43,19 @@ const NEGENTROPY_TRANSIENT_BACKOFF: [Duration; 4] = [ Duration::from_secs(7200), ]; +/// Maximum concurrent negentropy diff rounds per relay connection. +/// +/// Relays bound concurrent subscriptions per connection, and on +/// strfry-family relays negentropy views count against that same budget +/// (`maxSubsPerConnection`; nos.lol and relay.primal.net advertise 20, +/// shared with live subscriptions). Historic sync opens one diff per +/// filter, so an unbounded batch (146 filters observed in production) +/// bursts far past the tightest common budget and draws +/// "too many concurrent NEG requests" rejections. Four concurrent rounds +/// leaves the shared budget mostly available for live subscriptions; see +/// docs/explanation/sync-scaling-constraints.md. +const MAX_CONCURRENT_NEG_DIFFS: usize = 4; + /// How a failed negentropy diff should affect future NIP-77 attempts. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum NegentropyFailure { @@ -159,6 +172,9 @@ pub struct RelayConnection { nip77_transient_failures: std::sync::Arc, /// Deadline before which negentropy is not attempted (transient-failure cooldown) nip77_cooldown_until: std::sync::Arc>>, + /// Bounds concurrent negentropy diff rounds on this connection + /// (shared across clones; see [`MAX_CONCURRENT_NEG_DIFFS`]) + neg_diff_permits: std::sync::Arc, } impl RelayConnection { @@ -212,6 +228,9 @@ impl RelayConnection { nip77_supported: std::sync::Arc::new(std::sync::atomic::AtomicU8::new(0)), nip77_transient_failures: std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0)), nip77_cooldown_until: std::sync::Arc::new(std::sync::Mutex::new(None)), + neg_diff_permits: std::sync::Arc::new(tokio::sync::Semaphore::new( + MAX_CONCURRENT_NEG_DIFFS, + )), } } @@ -244,6 +263,9 @@ impl RelayConnection { nip77_supported: std::sync::Arc::new(std::sync::atomic::AtomicU8::new(0)), nip77_transient_failures: std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0)), nip77_cooldown_until: std::sync::Arc::new(std::sync::Mutex::new(None)), + neg_diff_permits: std::sync::Arc::new(tokio::sync::Semaphore::new( + MAX_CONCURRENT_NEG_DIFFS, + )), } } @@ -770,6 +792,28 @@ impl RelayConnection { &self, filter: Filter, ) -> Result { + // Bound concurrent rounds: historic sync opens one diff per filter, + // and relays count each open round against a per-connection budget + // shared with live subscriptions (see MAX_CONCURRENT_NEG_DIFFS). + // The permit is held for the whole round, including the timeout. + let _permit = self + .neg_diff_permits + .acquire() + .await + .map_err(|_| format!("Negentropy permits closed for {}", self.url))?; + + // While this round was queued, an earlier round may have started a + // transient-failure cooldown or received an explicit unsupported + // signal. Bail out to the per-batch REQ+EOSE fallback instead of + // opening another round against a relay that just failed. This does + // not record a failure, so it cannot escalate the cooldown. + if !self.supports_negentropy().await { + return Err(format!( + "Negentropy skipped for {}: cooldown active or relay marked unsupported", + self.url + )); + } + // Use dry_run to only identify differences without downloading events let sync_opts = SyncOptions::default().dry_run(); let client = self.client.clone(); diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 802fe94..02f2a73 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -5,6 +5,7 @@ pub mod censoring_proxy; pub mod git_server; pub mod mock_relay; +pub mod neg_limiting_proxy; pub mod nip09_helpers; pub mod port; pub mod purgatory_helpers; diff --git a/tests/common/neg_limiting_proxy.rs b/tests/common/neg_limiting_proxy.rs new file mode 100644 index 0000000..fe55088 --- /dev/null +++ b/tests/common/neg_limiting_proxy.rs @@ -0,0 +1,280 @@ +//! Concurrency-Limiting NEG 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 NIP-77 +//! negentropy rounds per connection. strfry counts negentropy views against +//! `maxSubsPerConnection` and answers excess `NEG-OPEN` frames with +//! `NOTICE ERROR: too many concurrent NEG requests` — the production +//! behaviour observed from nos.lol (which advertises `max_subscriptions: 20`) +//! during gitnostr.com startup bursts. +//! +//! The proxy additionally: +//! - records the peak number of concurrently open NEG rounds, so tests can +//! assert the syncing relay's concurrency bound end to end; +//! - delays backend responses to in-flight NEG rounds by a fixed interval, +//! so rounds opened together provably overlap instead of racing the +//! loopback round-trip. +//! +//! # Usage +//! +//! ```ignore +//! let source = TestRelay::start().await; +//! let proxy = NegLimitingProxy::start(source.url(), 4).await; +//! let syncing = TestRelay::start_with_sync(Some(proxy.url().into())).await; +//! // ... assert proxy.rejected_count() == 0 && proxy.peak_concurrent() <= 4 ... +//! ``` + +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 responses for in-flight NEG rounds. +/// +/// Long enough that a burst of NEG-OPEN frames sent together is observed +/// before any round can complete; short enough to keep tests fast. +const NEG_RESPONSE_DELAY: Duration = Duration::from_millis(100); + +/// WebSocket proxy bounding concurrent NEG rounds like a strfry relay. +pub struct NegLimitingProxy { + url: String, + peak: Arc, + rejected: Arc, + opened: Arc, + shutdown_tx: Option>, + handle: Option>, +} + +impl NegLimitingProxy { + /// Start a proxy on a random loopback port, forwarding to `backend_url` + /// and rejecting NEG-OPEN frames that would exceed `limit` concurrent + /// rounds on a connection. + pub async fn start(backend_url: &str, limit: usize) -> Self { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("NegLimitingProxy failed to bind"); + let port = listener + .local_addr() + .expect("NegLimitingProxy 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!("NegLimitingProxy 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 NEG rounds observed open at once on any connection. + pub fn peak_concurrent(&self) -> usize { + self.peak.load(Ordering::Relaxed) + } + + /// Number of NEG-OPEN frames rejected for exceeding the limit. + pub fn rejected_count(&self) -> usize { + self.rejected.load(Ordering::Relaxed) + } + + /// Total NEG-OPEN 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 NegLimitingProxy { + fn drop(&mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + } +} + +/// Forward one client connection to the backend, bounding concurrent NEG +/// rounds 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(); + + // NEG 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_neg_frame(text.as_str()) { + Some(NegFrame::Open(subid)) => { + if active.len() >= limit { + rejected.fetch_add(1, Ordering::Relaxed); + client_tx + .send(Message::text( + r#"["NOTICE","ERROR: too many concurrent NEG requests"]"# + .to_string(), + )) + .await + .map_err(|e| format!("client write: {e}"))?; + client_tx + .send(Message::text(format!( + r#"["NEG-ERR","{subid}","blocked: too many concurrent NEG requests"]"# + ))) + .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(NegFrame::Close(subid)) => { + active.remove(&subid); + } + Some(NegFrame::Other) | 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 { + match parse_neg_frame(text.as_str()) { + Some(NegFrame::Other) => { + // Hold NEG responses briefly so simultaneously + // opened rounds provably overlap at the proxy. + tokio::time::sleep(NEG_RESPONSE_DELAY).await; + } + Some(NegFrame::Open(subid)) | Some(NegFrame::Close(subid)) => { + // Backend-initiated NEG-ERR ends the round. + active.remove(&subid); + } + None => {} + } + } + client_tx + .send(message) + .await + .map_err(|e| format!("client write: {e}"))?; + } + } + } + + Ok(()) +} + +/// Parsed shape of a NEG-* frame. +enum NegFrame { + /// `["NEG-OPEN", , ...]` from the client. + Open(String), + /// `["NEG-CLOSE", ]` from the client, or `["NEG-ERR", , ..]` + /// from the backend (both end the round). + Close(String), + /// `["NEG-MSG", ...]` — an in-flight reconciliation frame. + Other, +} + +/// Classify a text frame if it is negentropy-related, else `None`. +fn parse_neg_frame(text: &str) -> Option { + if !text.starts_with("[\"NEG-") { + 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) + .and_then(|v| v.as_str()) + .map(|s| s.to_string()) + }; + match kind { + "NEG-OPEN" => Some(NegFrame::Open(subid()?)), + "NEG-CLOSE" | "NEG-ERR" => Some(NegFrame::Close(subid()?)), + "NEG-MSG" => Some(NegFrame::Other), + _ => None, + } +} diff --git a/tests/sync.rs b/tests/sync.rs index a6046e7..15b708f 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -38,5 +38,6 @@ mod sync { pub mod live_sync; pub mod maintainer_reprocessing; pub mod metrics; + pub mod neg_concurrency; pub mod tag_variations; } diff --git a/tests/sync/mod.rs b/tests/sync/mod.rs index 95ea37b..15f0d4a 100644 --- a/tests/sync/mod.rs +++ b/tests/sync/mod.rs @@ -135,4 +135,5 @@ pub mod discovery; pub mod live_sync; pub mod maintainer_reprocessing; pub mod metrics; +pub mod neg_concurrency; pub mod tag_variations; \ No newline at end of file diff --git a/tests/sync/neg_concurrency.rs b/tests/sync/neg_concurrency.rs new file mode 100644 index 0000000..2fb8f2e --- /dev/null +++ b/tests/sync/neg_concurrency.rs @@ -0,0 +1,226 @@ +//! Negentropy Concurrency Bound Tests +//! +//! Regression coverage for a production failure observed on gitnostr.com +//! (2026-08-04): startup historic sync opened one NIP-77 negentropy round per +//! filter with no bound (146 filters in one batch on the bootstrap relay), +//! exceeding relay per-connection subscription budgets. nos.lol — a strfry +//! relay advertising `max_subscriptions: 20`, a budget shared between +//! ordinary subscriptions and negentropy views — answered with 34 +//! "ERROR: too many concurrent NEG requests" NOTICEs, and the burst produced +//! 61 per-filter timeouts and 21 failed fallback-subscription creations in +//! the first two minutes after startup. +//! +//! The scenario drives the real sync path end to end: a genuine ngit-grasp +//! source relay (real NIP-77) sits behind a proxy that mimics the strfry +//! limit, rejecting NEG-OPEN frames beyond 4 concurrent rounds. The watched +//! set is sized so that one historic batch needs more negentropy rounds than +//! the limit (150 root events → two 100-ID chunks × three tag variants = six +//! filters). Unbounded concurrency draws rejections; the bounded client must +//! draw none while still overlapping rounds 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::neg_limiting_proxy::NegLimitingProxy; +use crate::common::purgatory_helpers::{ + create_state_event, create_test_repo_with_commit, push_to_relay, CommitVariant, +}; +use crate::common::sync_helpers::{ + build_layer2_issue_event, repo_coord, send_to_relay, wait_for_event_on_relay, +}; +use crate::common::{port, TestRelay}; + +/// Concurrency limit enforced by the proxy. +/// +/// Matches `MAX_CONCURRENT_NEG_DIFFS` 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 round. +const PROXY_NEG_LIMIT: usize = 4; + +/// Root events to seed. Two 100-ID chunks × three tag variants (e/E/q) give +/// six negentropy filters in a single historic batch — more rounds than the +/// bound, so bounded scheduling is required to avoid rejections. +const ISSUE_COUNT: usize = 150; + +/// Wait until the proxy has seen at least `min_total` NEG rounds +/// (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_neg_quiescence( + proxy: &NegLimitingProxy, + 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 +/// 150 historic issues, each a root event of that repository. +/// 2. A proxy in front of the source enforces the strfry-style limit of 4 +/// concurrent NEG rounds and records peak concurrency and rejections. +/// 3. The syncing relay bootstraps through the proxy. Syncing the issues +/// registers 150 root events, whose next historic batch needs six +/// negentropy rounds — more than the limit. +/// 4. The syncing relay must complete historic sync without a single +/// rejection while still running rounds concurrently. +#[tokio::test] +async fn startup_historic_sync_stays_within_relay_neg_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 = NegLimitingProxy::start(source.url(), PROXY_NEG_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://{}/{}/neg-repo.git", source.domain(), npub), + format!("http://{}/{}/neg-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("neg-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, + "neg-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, "neg-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. + let coordinate = repo_coord(&keys, "neg-repo"); + let mut issue_ids = Vec::with_capacity(ISSUE_COUNT); + for index in 0..ISSUE_COUNT { + let issue = build_layer2_issue_event(&keys, &coordinate, &format!("Issue {index}")) + .expect("build issue event"); + issue_ids.push(issue.id); + send_to_relay(&source, &issue) + .await + .expect("send issue to source"); + } + + // 5. Syncing relay bootstraps through the proxy. + let syncing = TestRelay::start_on_reservation_with_options( + syncing_reservation, + Some(proxy.url().to_string()), + false, + ) + .await; + + // 6. The issues must arrive via the bounded historic sync (sampled ends + // of the range cover both 100-ID 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(90), + ) + .await, + "issue should reach the syncing relay through bounded historic sync" + ); + } + + // 7. Wait for negentropy activity to include the root-event batch and + // settle. Layer-1 plus the repository batch plus the six root-event + // rounds put the floor at eight. + let (opened, rejected) = wait_for_neg_quiescence( + &proxy, + 8, + Duration::from_secs(3), + Duration::from_secs(60), + ) + .await + .expect("negentropy rounds should reach the root-event batch and settle"); + + // 8. The regression assertions. + assert_eq!( + rejected, 0, + "relay exceeded the per-connection NEG concurrency limit \ + (opened: {opened}, peak: {})", + proxy.peak_concurrent() + ); + assert!( + proxy.peak_concurrent() <= PROXY_NEG_LIMIT, + "peak concurrent NEG rounds {} exceeded the bound {PROXY_NEG_LIMIT}", + proxy.peak_concurrent() + ); + assert!( + proxy.peak_concurrent() >= 2, + "negentropy rounds should still overlap under the bound, peak: {}", + proxy.peak_concurrent() + ); + assert!( + opened > PROXY_NEG_LIMIT, + "scenario must run more rounds than the bound to exercise queueing, opened: {opened}" + ); + + syncing.stop().await; + proxy.stop().await; + source.stop().await; +}