From b55d2de404d9e61ac61b729eb835cb6224426c2d Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 7 Aug 2026 14:48:01 +0000 Subject: [PATCH] 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. --- src/sync/relay_connection.rs | 59 ++++++++++++++++++++++++++++++++++++ 1 file changed, 59 insertions(+) diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 81f7f88..8606593 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -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, + 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(