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.
This commit is contained in:
DanConwayDev
2026-08-04 21:13:17 +00:00
parent 5e840dc1ef
commit a0520c3d6b
6 changed files with 553 additions and 0 deletions
+44
View File
@@ -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<std::sync::atomic::AtomicU32>,
/// Deadline before which negentropy is not attempted (transient-failure cooldown)
nip77_cooldown_until: std::sync::Arc<std::sync::Mutex<Option<tokio::time::Instant>>>,
/// 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>,
}
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<nostr_sdk::client::SyncSummary, String> {
// 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();
+1
View File
@@ -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;
+280
View File
@@ -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<AtomicUsize>,
rejected: Arc<AtomicUsize>,
opened: Arc<AtomicUsize>,
shutdown_tx: Option<oneshot::Sender<()>>,
handle: Option<tokio::task::JoinHandle<()>>,
}
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<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();
// NEG 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_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", <subid>, ...]` from the client.
Open(String),
/// `["NEG-CLOSE", <subid>]` from the client, or `["NEG-ERR", <subid>, ..]`
/// 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<NegFrame> {
if !text.starts_with("[\"NEG-") {
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)
.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,
}
}
+1
View File
@@ -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;
}
+1
View File
@@ -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;
+226
View File
@@ -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::<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, "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;
}