Merge #97556042: chore(sync): diagnose missing transient terminals

nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsfw4tqgfrz874nyp8ae9h08dsrj0zc44w7p5qcgc5xgx7m5epu4rctceh6k

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

CoverNote:

Production at `58ef8682` repeatedly held transient subscriptions until the 120-second watchdog across independent relays and request classes. Diagnostic tip `b55d2de` was deployed to gitnostr.com at 2026-08-07 14:52:34 UTC against a directly comparable startup workload.

The diagnostics established that expired subscriptions were already absent from rust-nostr's active subscription map, so the relay terminal had arrived but ngit-grasp had not released its local ledger ownership. An initial 4,096-message queue removed watchdogs from a 36,054-ID relay.ngit.dev burst, but production falsified queue sizing as a general solution: the five-per-cycle signature moved to ngit.danconwaydev.com and then git.shakespeare.diy, including four empty responses on each peer. Empty responses exposed a second ordering hole, while the shifting high-volume peer showed that any finite EVENT queue can remain congested under sustained pages.

The final proposal fixes both lifecycle dependencies by construction:

- register a second rust-nostr broadcast receiver synchronously as a terminal-control lane; it consumes EOSE/CLOSED and connection terminal status independently of the processor-facing EVENT queue;
- pre-generate every transient subscription ID and register its generation-scoped permit before sending the REQ, pass that ID to rust-nostr, and roll registration back on subscribe failure.

The terminal listener is the sole transient-release path during a connected session. The processor listener still delivers ordered EVENTs, pagination terminals, and live CLOSED restoration. EOSE still enqueues CLOSE before returning the slot; CLOSED and connection teardown remain definitive boundaries. The EVENT queue stays at its original bounded 1,000 messages—correctness no longer depends on buying lifecycle headroom with memory.

Subscription concurrency, filter contents, pagination, watchdog timing, retry behaviour, live lifecycle, and the configuration surface are unchanged. The bounded diagnostic fields remain useful for attributing any residual watchdog.

Regression coverage deliberately blocks a capacity-one EVENT lane behind 1,200 LocalRelay events and requires all transient permits to release within three seconds. A second scenario opens 25 immediate empty queries and requires every pre-registered permit to release within the same bounded deadline. The architecture document records the implemented control/data separation and registration ordering.

Pre-deployment validation of final tip `04b93c9`:

- all 648 library tests pass;
- both focused lifecycle regressions pass;
- standalone `startup_historic_sync_stays_within_relay_req_concurrency_limit` passes with zero rejected REQs;
- changed-source rustfmt and `git diff --check` pass.

Production verification of exact tip `04b93c9a1df9d28fc6af96b131d0ebc5646552ca`:

- remotely built and activated on gitnostr.com at 2026-08-07 15:37:54 UTC;
- immediately reproduced a larger comparable workload: relay.ngit.dev had 994 missing IDs/4 chunks and 38,518 IDs/129 chunks;
- through 15:45:27 (7m33s, more than three complete 120-second recurrence windows), there were zero watchdogs with `sdk_subscription_active=false`, including zero at relay.ngit.dev, ngit.danconwaydev.com, and git.shakespeare.diy;
- four watchdogs did fire correctly for genuinely active relay-side subscriptions: three historic pages at relay.cyberguy.fyi and one at nostr.sebastix.dev, all `sdk_subscription_active=true`;
- the target relay continued useful progress: its generic announcement batch completed and later REQ/EOSE batches had confirmed 86 state-only repositories by 15:40:45 despite the pre-existing startup rate-limit CLOSED burst;
- service remained active with zero restarts, zero panics, zero budget-exhaustion/concurrency-limit messages, and 1.24 GiB memory at the final sample.

This exact tip is recommended for merge. Production distinguishes the repaired local-accounting defect from the watchdog's intended recovery of peers that genuinely leave subscriptions open.
This commit is contained in:
DanConwayDev
2026-08-07 16:54:56 +01:00
3 changed files with 265 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();
+250 -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)
@@ -120,6 +126,9 @@ struct HeldTransientPermits {
_ledger_slot: tokio::sync::OwnedSemaphorePermit,
generation: u64,
request_class: TransientRequestClass,
opened_at: std::time::Instant,
last_event_at: Option<std::time::Instant>,
delivered_events: usize,
}
/// Bounded origin of an auto-close REQ, retained for watchdog diagnosis.
@@ -702,6 +711,9 @@ impl RelayConnection {
generation: ledger_slot.generation,
_ledger_slot: ledger_slot.permit,
request_class,
opened_at: std::time::Instant::now(),
last_event_at: None,
delivered_events: 0,
});
}
drop(ledger_slot);
@@ -718,6 +730,9 @@ impl RelayConnection {
generation: ledger_slot.generation,
_ledger_slot: ledger_slot.permit,
request_class,
opened_at: std::time::Instant::now(),
last_event_at: None,
delivered_events: 0,
});
}
drop(class_cap);
@@ -907,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 {
@@ -915,6 +970,7 @@ impl RelayConnection {
event,
subscription_id,
} => {
self.record_transient_req_event(&subscription_id);
tracing::trace!(
relay = %url,
event_id = %event.id,
@@ -935,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
@@ -979,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 {
@@ -1034,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();
@@ -1205,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()
@@ -1227,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()
@@ -1374,6 +1444,21 @@ impl RelayConnection {
if let (Some((generation, request_class)), Ok(Some(relay))) =
(held_request, client.relay(&url).await)
{
let sdk_subscription_active = relay.subscription(&sub_id).await.is_some();
let diagnostic = held
.lock()
.expect("transient permit map poisoned")
.get(&sub_id)
.filter(|permits| permits.generation == generation)
.map(|permits| {
(
permits.delivered_events,
permits.opened_at.elapsed().as_secs(),
permits.last_event_at.map(|at| at.elapsed().as_secs()),
)
});
let (delivered_events, open_seconds, idle_seconds) =
diagnostic.unwrap_or((0, TRANSIENT_REQ_PERMIT_TIMEOUT.as_secs(), None));
if relay
.send_msg(ClientMessage::close(sub_id.clone()))
.await
@@ -1384,6 +1469,10 @@ impl RelayConnection {
relay = %relay_url,
sub_id = %sub_id,
request_class = request_class.as_str(),
sdk_subscription_active,
delivered_events,
open_seconds,
idle_seconds,
"Transient REQ watchdog sent CLOSE after missing EOSE/CLOSED"
);
crate::metrics::record_transient_req_watchdog(request_class.as_str(), "closed");
@@ -1392,6 +1481,10 @@ impl RelayConnection {
relay = %relay_url,
sub_id = %sub_id,
request_class = request_class.as_str(),
sdk_subscription_active,
delivered_events,
open_seconds,
idle_seconds,
"Transient REQ watchdog could not send CLOSE; retaining subscription slot"
);
crate::metrics::record_transient_req_watchdog(
@@ -1403,6 +1496,18 @@ impl RelayConnection {
});
}
fn record_transient_req_event(&self, sub_id: &SubscriptionId) {
if let Some(held) = self
.transient_req_permits_held
.lock()
.expect("transient permit map poisoned")
.get_mut(sub_id)
{
held.delivered_events = held.delivered_events.saturating_add(1);
held.last_event_at = Some(std::time::Instant::now());
}
}
/// 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
@@ -1454,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)
@@ -2062,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());
@@ -2420,6 +2638,9 @@ mod tests {
let held = HeldTransientPermits {
generation,
request_class: TransientRequestClass::HistoricPage,
opened_at: std::time::Instant::now(),
last_event_at: None,
delivered_events: 0,
_class_cap: connection
.transient_req_permits
.clone()
@@ -2433,6 +2654,17 @@ mod tests {
.expect("reserve ledger slot"),
};
connection.hold_transient_req_permit(sub_id.clone(), held);
connection.record_transient_req_event(&sub_id);
{
let permits = connection
.transient_req_permits_held
.lock()
.expect("transient permit map poisoned");
let diagnostic = permits.get(&sub_id).expect("held subscription diagnostic");
assert_eq!(diagnostic.delivered_events, 1);
assert!(diagnostic.last_event_at.is_some());
}
assert!(
!RelayConnection::release_transient_req_permit_for_generation(