chore(sync): diagnose missing transient terminals

Production at 58ef8682 showed auto-close requests reaching the 120-second watchdog across multiple independent relays and request classes. The existing warning cannot distinguish a peer that omitted EOSE/CLOSED from an SDK or application notification delay, so a behavioural recovery change would be speculative and could prematurely complete historic pages.

Retain bounded per-request diagnostics with the existing held permit: delivered-event count, open time, and last-event time. At watchdog recovery, also query whether rust-nostr still considers the subscription active and include those values in the existing warning. Normal EOSE/CLOSED, admission, retry, timeout, and slot-release behaviour are unchanged.

The fields are bounded by the existing transient subscription cap and contain no event IDs or new metric labels. Synthetic completion and watchdog-policy changes are deliberately excluded until production identifies which lifecycle layer lost the terminal signal.

Validated with the 39 relay_connection unit tests, a focused diagnostic-count assertion, rustfmt on the changed source, and git diff --check. The repository-wide cargo fmt --check remains affected by pre-existing formatter drift in integration fixtures under Rust 1.96.
This commit is contained in:
DanConwayDev
2026-08-07 14:48:01 +00:00
parent e5e088e306
commit b55d2de404
+59
View File
@@ -120,6 +120,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 +705,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 +724,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);
@@ -915,6 +924,7 @@ impl RelayConnection {
event,
subscription_id,
} => {
self.record_transient_req_event(&subscription_id);
tracing::trace!(
relay = %url,
event_id = %event.id,
@@ -1374,6 +1384,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 +1409,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 +1421,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 +1436,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
@@ -2420,6 +2465,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 +2481,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(