diff --git a/CHANGELOG.md b/CHANGELOG.md index 57392cc..c771a83 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Expanded `grasp-audit` with stable machine-readable results, full JSON reports, explicit audit identities, and hardened discovered-server probes. +- Added bounded request-class labels to transient-sync watchdog logs and + metrics so operators can identify which historic recovery path is stalling. ### Changed diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 82058f0..b589eff 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -372,7 +372,11 @@ So concurrency is not a free scaling axis; it is the residual of the ledger: five-request class cap inside the shared ledger: a slot is acquired when the auto-close REQ is sent and released when its EOSE or CLOSED arrives (with a 120 s watchdog that sends CLOSE for only the unresponsive - subscription before releasing its slot). Live subscriptions are ledgered first, so NEG, + subscription before releasing its slot). Each held permit retains one of a + fixed set of request classes (historic page, pagination page, hint + verification, negentropy hydration, retry, or semantic fallback); watchdog + logs and metrics expose that class without using relay URLs or subscription + IDs as metric labels. Live subscriptions are ledgered first, so NEG, transient REQ, and purgatory exact-ID polling share only the remaining capacity. - Permit acquisition checks relay health first: while a rate-limit or diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 87d03fb..a264a93 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -202,6 +202,26 @@ lazy_static! { .expect("register purgatory git fetch oids metric"); metric }; + static ref TRANSIENT_REQ_WATCHDOG_TOTAL: CounterVec = { + let metric = CounterVec::new( + Opts::new( + "ngit_sync_transient_req_watchdog_total", + "Transient REQ watchdog outcomes by bounded request class", + ), + &["class", "outcome"], + ) + .expect("build transient REQ watchdog metric"); + REGISTRY + .register(Box::new(metric.clone())) + .expect("register transient REQ watchdog metric"); + metric + }; +} + +pub fn record_transient_req_watchdog(class: &str, outcome: &str) { + TRANSIENT_REQ_WATCHDOG_TOTAL + .with_label_values(&[class, outcome]) + .inc(); } /// Record one completed outbound purgatory git fetch pass. diff --git a/src/sync/mod.rs b/src/sync/mod.rs index a68cf8f..f062ac5 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -34,7 +34,9 @@ pub use rejected_index::{EventType, RejectionReason}; // Current code still uses the simple HashSet type alias below // Re-export relay connection types -pub use relay_connection::{NegentropySyncResult, RelayConnection, RelayEvent}; +pub use relay_connection::{ + NegentropySyncResult, RelayConnection, RelayEvent, TransientRequestClass, +}; // Re-export self-subscriber types pub use self_subscriber::SelfSubscriber; @@ -466,6 +468,18 @@ impl PaginationState { .collect() } + fn request_class(&self) -> TransientRequestClass { + if self + .filters + .iter() + .any(|state| state.verification_baseline.is_some()) + { + TransientRequestClass::PaginationVerification + } else { + TransientRequestClass::PaginationPage + } + } + fn next_page(mut self, session: &mut RelayPaginationSession) -> Option { let mut next = Vec::new(); for mut state in self.filters.drain(..) { @@ -1491,7 +1505,10 @@ impl SyncManager { let mut next_page_started = false; if let Some(conn) = self.connections.get(&relay_url_for_pagination) { - match conn.subscribe_filters(next_filters.clone(), true).await { + match conn + .subscribe_filters(next_filters.clone(), next_page.request_class()) + .await + { Ok(new_sub_id) => { let mut pending = self.pending_sync_index.write().await; if let Some(batches) = pending.get_mut(&relay_url_for_pagination) { @@ -1637,7 +1654,13 @@ impl SyncManager { let mut new_sub_ids = HashSet::new(); if let Some(conn) = self.connections.get(&relay_url_for_fallback) { for filter_group in group_filters_for_req(&fallback_filters) { - match conn.subscribe_filters(filter_group, true).await { + match conn + .subscribe_filters( + filter_group, + TransientRequestClass::NegentropyFallback, + ) + .await + { Ok(sub_id) => { new_sub_ids.insert(sub_id); } @@ -1808,7 +1831,10 @@ impl SyncManager { let mut new_sub_ids = HashSet::new(); if let Some(conn) = self.connections.get(&relay_url_for_retry) { for filter in retry_filters { - match conn.subscribe_filter(filter, true).await { + match conn + .subscribe_filter(filter, TransientRequestClass::NegentropyRetry) + .await + { Ok(sub_id) => { new_sub_ids.insert(sub_id); } @@ -1958,7 +1984,7 @@ impl SyncManager { } match connection - .subscribe_filters(next_filters.clone(), true) + .subscribe_filters(next_filters.clone(), next_page.request_class()) .await { Ok(new_sub_id) => { @@ -5573,7 +5599,13 @@ impl SyncManager { let mut subscription_ids = HashSet::new(); for (idx, filter) in ids_filters.iter().enumerate() { if let Some(conn) = self.connections.get(relay_url) { - match conn.subscribe_filter(filter.clone(), true).await { + match conn + .subscribe_filter( + filter.clone(), + TransientRequestClass::NegentropyHydration, + ) + .await + { Ok(sub_id) => { subscription_ids.insert(sub_id); } @@ -5646,7 +5678,13 @@ impl SyncManager { if let Some(conn) = self.connections.get(relay_url) { let grouped_filters = filter_group; - match conn.subscribe_filters(grouped_filters.clone(), true).await { + match conn + .subscribe_filters( + grouped_filters.clone(), + TransientRequestClass::HistoricPage, + ) + .await + { Ok(sub_id) => { subscription_ids.insert(sub_id.clone()); pagination_state.insert(sub_id, PaginationState::new(grouped_filters)); @@ -6195,6 +6233,10 @@ mod tests { .next_page(&mut session) .expect("a suspiciously short hinted page needs verification"); assert_eq!(session.hint, PaginationHint::Verifying(1000)); + assert_eq!( + verification.request_class(), + TransientRequestClass::PaginationVerification + ); let unseen = EventBuilder::new(Kind::TextNote, "older unseen event") .custom_created_at(Timestamp::from_secs(99)) diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 7ce809d..81f7f88 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -119,8 +119,43 @@ struct HeldTransientPermits { _class_cap: tokio::sync::OwnedSemaphorePermit, _ledger_slot: tokio::sync::OwnedSemaphorePermit, generation: u64, + request_class: TransientRequestClass, } +/// Bounded origin of an auto-close REQ, retained for watchdog diagnosis. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TransientRequestClass { + HistoricPage, + PaginationPage, + PaginationVerification, + NegentropyHydration, + NegentropyRetry, + NegentropyFallback, +} + +impl TransientRequestClass { + pub const fn as_str(self) -> &'static str { + match self { + Self::HistoricPage => "historic_page", + Self::PaginationPage => "pagination_page", + Self::PaginationVerification => "pagination_verification", + Self::NegentropyHydration => "negentropy_hydration", + Self::NegentropyRetry => "negentropy_retry", + Self::NegentropyFallback => "negentropy_fallback", + } + } +} + +#[cfg(test)] +const TRANSIENT_REQUEST_CLASSES: [TransientRequestClass; 6] = [ + TransientRequestClass::HistoricPage, + TransientRequestClass::PaginationPage, + TransientRequestClass::PaginationVerification, + TransientRequestClass::NegentropyHydration, + TransientRequestClass::NegentropyRetry, + TransientRequestClass::NegentropyFallback, +]; + #[derive(Clone)] struct SubscriptionLedgerSession { generation: u64, @@ -655,7 +690,10 @@ impl RelayConnection { /// Acquire both transient constraints without waiting while holding only /// one of them. This prevents class-cap waiters from pinning ledger slots /// and ledger waiters from pinning the class cap. - async fn acquire_transient_permits(&self) -> Result { + async fn acquire_transient_permits( + &self, + request_class: TransientRequestClass, + ) -> Result { loop { let ledger_slot = self.acquire_subscription_slots(1).await?; if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() { @@ -663,6 +701,7 @@ impl RelayConnection { _class_cap: class_cap, generation: ledger_slot.generation, _ledger_slot: ledger_slot.permit, + request_class, }); } drop(ledger_slot); @@ -678,6 +717,7 @@ impl RelayConnection { _class_cap: class_cap, generation: ledger_slot.generation, _ledger_slot: ledger_slot.permit, + request_class, }); } drop(class_cap); @@ -741,7 +781,7 @@ impl RelayConnection { filter_groups: Vec>, ) -> Result, String> { self.replace_live_filter_groups_with(filter_groups, |filters, permit| { - self.subscribe_filters_with_live_permit(filters, false, Some(permit)) + self.subscribe_filters_with_live_permit(filters, None, Some(permit)) }) .await } @@ -1009,7 +1049,7 @@ impl RelayConnection { /// /// # Arguments /// * `filter` - The filter to subscribe to - /// * `auto_close` - If true, subscription automatically closes after EOSE (for historic sync). If false, stays open for new events (for live sync). + /// * `request_class` - Bounded historic-sync origin retained for watchdog diagnostics. /// /// # Returns /// * `Ok(SubscriptionId)` - The subscription ID on success @@ -1017,9 +1057,9 @@ impl RelayConnection { pub async fn subscribe_filter( &self, filter: Filter, - auto_close: bool, + request_class: TransientRequestClass, ) -> Result { - self.subscribe_filters(vec![filter], auto_close).await + self.subscribe_filters(vec![filter], request_class).await } /// Subscribe to several OR filters under one NIP-01 subscription ID. @@ -1030,9 +1070,9 @@ impl RelayConnection { pub async fn subscribe_filters( &self, filters: Vec, - auto_close: bool, + request_class: TransientRequestClass, ) -> Result { - self.subscribe_filters_with_live_permit(filters, auto_close, None) + self.subscribe_filters_with_live_permit(filters, Some(request_class), None) .await } @@ -1044,7 +1084,7 @@ impl RelayConnection { filter_groups: Vec>, ) -> Result, String> { self.subscribe_live_filter_groups_with(filter_groups, |filters, permit| { - self.subscribe_filters_with_live_permit(filters, false, Some(permit)) + self.subscribe_filters_with_live_permit(filters, None, Some(permit)) }) .await } @@ -1101,7 +1141,7 @@ impl RelayConnection { async fn subscribe_filters_with_live_permit( &self, filters: Vec, - auto_close: bool, + transient_class: Option, reserved_live_permit: Option, ) -> Result { if filters.is_empty() { @@ -1112,7 +1152,7 @@ impl RelayConnection { relay = %self.url, filter_count = filters.len(), filters = ?filters, - auto_close = auto_close, + auto_close = transient_class.is_some(), "subscribe_filters called" ); @@ -1121,7 +1161,7 @@ impl RelayConnection { // exceeding relay subscription budgets. The permit is registered // against the subscription id on success and released when the // subscription's EOSE or CLOSED arrives (see run_event_loop). - let (transient_permit, live_permit) = if auto_close { + let (transient_permit, live_permit) = if let Some(request_class) = transient_class { if self.historic_capacity_consumed_by_live() { tracing::warn!( relay = %self.url, @@ -1132,7 +1172,10 @@ impl RelayConnection { self.url )); } - (Some(self.acquire_transient_permits().await?), None) + ( + Some(self.acquire_transient_permits(request_class).await?), + None, + ) } else { let acquired = match reserved_live_permit { Some(permit) => Ok(permit), @@ -1162,7 +1205,7 @@ 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(); - let output = if auto_close { + let output = if transient_class.is_some() { self.client .subscribe(filters) .close_on( @@ -1184,7 +1227,7 @@ impl RelayConnection { return Err(format!("Failed to subscribe on {}: {}", self.url, failures)); } - if auto_close { + if transient_class.is_some() { if let Some(permit) = transient_permit { self.hold_transient_req_permit(output.value.clone(), permit); } @@ -1323,33 +1366,38 @@ impl RelayConnection { let relay_url = self.url.clone(); tokio::spawn(async move { tokio::time::sleep(TRANSIENT_REQ_PERMIT_TIMEOUT).await; - let generation = held + let held_request = held .lock() .expect("transient permit map poisoned") .get(&sub_id) - .map(|permits| permits.generation); - if let (Some(generation), Ok(Some(relay))) = (generation, client.relay(&url).await) { + .map(|permits| (permits.generation, permits.request_class)); + if let (Some((generation, request_class)), Ok(Some(relay))) = + (held_request, client.relay(&url).await) + { if relay .send_msg(ClientMessage::close(sub_id.clone())) .await .is_ok() { - Self::release_transient_req_permit_for_generation( - &held, - &sub_id, - generation, - ); + Self::release_transient_req_permit_for_generation(&held, &sub_id, generation); tracing::warn!( relay = %relay_url, sub_id = %sub_id, + request_class = request_class.as_str(), "Transient REQ watchdog sent CLOSE after missing EOSE/CLOSED" ); + crate::metrics::record_transient_req_watchdog(request_class.as_str(), "closed"); } else { tracing::warn!( relay = %relay_url, sub_id = %sub_id, + request_class = request_class.as_str(), "Transient REQ watchdog could not send CLOSE; retaining subscription slot" ); + crate::metrics::record_transient_req_watchdog( + request_class.as_str(), + "close_failed", + ); } } }); @@ -2313,7 +2361,10 @@ mod tests { assert_eq!(ids.len(), 8); assert_eq!(connection.subscription_budget().available_permits(), 0); let error = connection - .subscribe_filter(Filter::new().kind(Kind::TextNote), true) + .subscribe_filter( + Filter::new().kind(Kind::TextNote), + TransientRequestClass::HistoricPage, + ) .await .expect_err("history must defer when complete live coverage fills usable capacity"); assert!( @@ -2368,6 +2419,7 @@ mod tests { let generation = connection.current_subscription_generation(); let held = HeldTransientPermits { generation, + request_class: TransientRequestClass::HistoricPage, _class_cap: connection .transient_req_permits .clone() @@ -2449,8 +2501,11 @@ mod tests { .await .expect("occupy transient class cap"); let class_waiter_connection = connection.clone(); - let mut class_waiter = - tokio::spawn(async move { class_waiter_connection.acquire_transient_permits().await }); + let mut class_waiter = tokio::spawn(async move { + class_waiter_connection + .acquire_transient_permits(TransientRequestClass::HistoricPage) + .await + }); assert!( tokio::time::timeout(Duration::from_millis(50), &mut class_waiter) .await @@ -2471,8 +2526,11 @@ mod tests { .await .expect("occupy subscription ledger"); let ledger_waiter_connection = connection.clone(); - let mut ledger_waiter = - tokio::spawn(async move { ledger_waiter_connection.acquire_transient_permits().await }); + let mut ledger_waiter = tokio::spawn(async move { + ledger_waiter_connection + .acquire_transient_permits(TransientRequestClass::HistoricPage) + .await + }); assert!( tokio::time::timeout(Duration::from_millis(50), &mut ledger_waiter) .await @@ -2600,6 +2658,25 @@ mod tests { ); } + #[test] + fn transient_request_metric_labels_are_bounded_and_unique() { + let labels = TRANSIENT_REQUEST_CLASSES.map(TransientRequestClass::as_str); + let unique = labels.into_iter().collect::>(); + + assert_eq!(unique.len(), TRANSIENT_REQUEST_CLASSES.len()); + assert_eq!( + unique, + std::collections::HashSet::from([ + "historic_page", + "pagination_page", + "pagination_verification", + "negentropy_hydration", + "negentropy_retry", + "negentropy_fallback", + ]) + ); + } + #[test] fn test_normalize_url_real_world_example() { // Test the exact case from the bug report