fix(sync): separate transient terminal accounting

Production diagnostics at b55d2de showed watchdog-expired REQs were already absent from rust-nostr's active subscription map. An initial 4,096-message queue removed watchdogs from a 36,054-ID relay.ngit.dev burst but only moved the five-per-cycle signature first to ngit.danconwaydev.com and then git.shakespeare.diy. Empty expired responses disproved event volume as the complete explanation: terminal accounting depended both on a congestible EVENT lane and on permit registration after subscribe returned.

Give transient lifecycle messages an independent control lane by registering a second rust-nostr broadcast receiver before subscriptions can begin. It handles EOSE/CLOSED and connection terminal status without waiting for the bounded processor data queue. Pre-generate transient subscription IDs and register generation-scoped ownership before sending REQ, passing the same ID into rust-nostr and rolling it back on subscribe failure. Keep the original 1,000-message queue; correctness no longer depends on sizing it for traffic.

The terminal listener is the sole transient-release path during a connected session, while the processor listener retains ordered EVENT delivery, pagination signals, and live CLOSED restoration. EOSE still enqueues CLOSE before returning the slot; peer CLOSED and connection teardown are definitive terminal boundaries. Relay notifications are broadcast, so the control listener cannot steal messages from processing.

This deliberately leaves concurrency, filters, pagination, retries, watchdog duration, live lifecycle, and configuration unchanged. It assumes rust-nostr preserves broadcast terminal notifications and that a failed exact-target subscribe sends no REQ; the latter path rolls ownership back.

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