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).
This commit is contained in:
DanConwayDev
2026-08-05 07:50:45 +00:00
parent 680eaf2f4a
commit e942199a55
9 changed files with 851 additions and 1 deletions
@@ -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 |
+145 -1
View File
@@ -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<SubscriptionId, tokio::sync::OwnedSemaphorePermit>,
>,
>;
/// 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<tokio::sync::Semaphore>,
/// Bounds concurrent transient (auto-close) REQ subscriptions on this
/// connection (shared across clones; see
/// [`MAX_CONCURRENT_TRANSIENT_REQS`])
transient_req_permits: std::sync::Arc<tokio::sync::Semaphore>,
/// 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<SubscriptionId> 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
// =========================================================================
+1
View File
@@ -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;
+14
View File
@@ -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<PortReservation> {
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 {
+68
View File
@@ -62,6 +62,9 @@ pub struct TestRelay {
_relay_data_dir: Option<tempfile::TempDir>,
/// 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<String>,
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
+324
View File
@@ -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<AtomicUsize>,
rejected: Arc<AtomicUsize>,
opened: Arc<AtomicUsize>,
shutdown_tx: Option<oneshot::Sender<()>>,
handle: Option<tokio::task::JoinHandle<()>>,
}
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<AtomicUsize>,
rejected: Arc<AtomicUsize>,
opened: Arc<AtomicUsize>,
) -> 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<String> = 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", <subid>, <filters...>]` - transient, counted.
Req(String),
/// `["REQ", ...]` where every filter carries `limit: 0` - a persistent
/// live subscription, exempt from the transient bound.
LiveReq,
/// `["CLOSE", <subid>]` - ends the subscription.
Close(String),
/// `["NEG-OPEN", <subid>, ...]` - rejected as unsupported.
NegOpen(String),
}
/// Classify a client text frame if the proxy tracks it, else `None`.
fn parse_client_frame(text: &str) -> Option<ClientFrame> {
if !(text.starts_with("[\"REQ")
|| text.starts_with("[\"CLOSE")
|| text.starts_with("[\"NEG-OPEN"))
{
return None;
}
let value = serde_json::from_str::<serde_json::Value>(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", <subid>]` frame.
fn parse_backend_eose(text: &str) -> Option<String> {
if !text.starts_with("[\"EOSE") {
return None;
}
let value = serde_json::from_str::<serde_json::Value>(text).ok()?;
let array = value.as_array()?;
if array.first()?.as_str()? != "EOSE" {
return None;
}
Some(array.get(1)?.as_str()?.to_string())
}
+1
View File
@@ -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;
}
+1
View File
@@ -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;
+289
View File
@@ -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::<Vec<_>>(),
&relay_urls.iter().map(|s| s.as_str()).collect::<Vec<_>>(),
)
.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;
}