chore(sync): classify transient watchdog recovery

Production at 57bc76ea forced exactly three relay.ngit.dev transient subscriptions through the 120-second watchdog every two minutes: 30 CLOSEs in the first 20 minutes. Subscription accounting remained within bounds, but the held-permit state and warning retained only random subscription IDs, so the stalled lifecycle could not be distinguished without guessing.

Require every auto-close REQ caller to provide one of six bounded request classes, retain that class with its per-session permit, and expose it in watchdog logs and a class/outcome Prometheus counter. Pagination distinguishes ordinary continuation from NIP-11 hint verification using its existing state; exact-ID purgatory polling remains correctly outside these labels because it uses the separately timed fetch_events path.

This is observability only: timeout, retry, pagination, request admission, filter contents, and slot-release behaviour are unchanged. Relay URLs and subscription IDs are deliberately excluded from metric labels, and no configuration surface is added.

Validated by the 646-test library suite, including fixed-label uniqueness and hint-verification classification tests, plus rustfmt and git diff --check.
This commit is contained in:
DanConwayDev
2026-08-07 13:21:45 +00:00
parent 64f33dcda7
commit 58ef8682f3
5 changed files with 180 additions and 35 deletions
+2
View File
@@ -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 - Expanded `grasp-audit` with stable machine-readable results, full JSON
reports, explicit audit identities, and hardened discovered-server probes. 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 ### Changed
+5 -1
View File
@@ -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 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 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 (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, transient REQ,
and purgatory exact-ID polling share only the remaining capacity. and purgatory exact-ID polling share only the remaining capacity.
- Permit acquisition checks relay health first: while a rate-limit or - Permit acquisition checks relay health first: while a rate-limit or
+20
View File
@@ -202,6 +202,26 @@ lazy_static! {
.expect("register purgatory git fetch oids metric"); .expect("register purgatory git fetch oids metric");
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. /// Record one completed outbound purgatory git fetch pass.
+49 -7
View File
@@ -34,7 +34,9 @@ pub use rejected_index::{EventType, RejectionReason};
// Current code still uses the simple HashSet type alias below // Current code still uses the simple HashSet type alias below
// Re-export relay connection types // 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 // Re-export self-subscriber types
pub use self_subscriber::SelfSubscriber; pub use self_subscriber::SelfSubscriber;
@@ -466,6 +468,18 @@ impl PaginationState {
.collect() .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<Self> { fn next_page(mut self, session: &mut RelayPaginationSession) -> Option<Self> {
let mut next = Vec::new(); let mut next = Vec::new();
for mut state in self.filters.drain(..) { for mut state in self.filters.drain(..) {
@@ -1491,7 +1505,10 @@ impl SyncManager {
let mut next_page_started = false; let mut next_page_started = false;
if let Some(conn) = self.connections.get(&relay_url_for_pagination) { 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) => { Ok(new_sub_id) => {
let mut pending = self.pending_sync_index.write().await; let mut pending = self.pending_sync_index.write().await;
if let Some(batches) = pending.get_mut(&relay_url_for_pagination) { if let Some(batches) = pending.get_mut(&relay_url_for_pagination) {
@@ -1637,7 +1654,13 @@ impl SyncManager {
let mut new_sub_ids = HashSet::new(); let mut new_sub_ids = HashSet::new();
if let Some(conn) = self.connections.get(&relay_url_for_fallback) { if let Some(conn) = self.connections.get(&relay_url_for_fallback) {
for filter_group in group_filters_for_req(&fallback_filters) { 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) => { Ok(sub_id) => {
new_sub_ids.insert(sub_id); new_sub_ids.insert(sub_id);
} }
@@ -1808,7 +1831,10 @@ impl SyncManager {
let mut new_sub_ids = HashSet::new(); let mut new_sub_ids = HashSet::new();
if let Some(conn) = self.connections.get(&relay_url_for_retry) { if let Some(conn) = self.connections.get(&relay_url_for_retry) {
for filter in retry_filters { for filter in retry_filters {
match conn.subscribe_filter(filter, true).await { match conn
.subscribe_filter(filter, TransientRequestClass::NegentropyRetry)
.await
{
Ok(sub_id) => { Ok(sub_id) => {
new_sub_ids.insert(sub_id); new_sub_ids.insert(sub_id);
} }
@@ -1958,7 +1984,7 @@ impl SyncManager {
} }
match connection match connection
.subscribe_filters(next_filters.clone(), true) .subscribe_filters(next_filters.clone(), next_page.request_class())
.await .await
{ {
Ok(new_sub_id) => { Ok(new_sub_id) => {
@@ -5573,7 +5599,13 @@ impl SyncManager {
let mut subscription_ids = HashSet::new(); let mut subscription_ids = HashSet::new();
for (idx, filter) in ids_filters.iter().enumerate() { for (idx, filter) in ids_filters.iter().enumerate() {
if let Some(conn) = self.connections.get(relay_url) { 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) => { Ok(sub_id) => {
subscription_ids.insert(sub_id); subscription_ids.insert(sub_id);
} }
@@ -5646,7 +5678,13 @@ impl SyncManager {
if let Some(conn) = self.connections.get(relay_url) { if let Some(conn) = self.connections.get(relay_url) {
let grouped_filters = filter_group; 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) => { Ok(sub_id) => {
subscription_ids.insert(sub_id.clone()); subscription_ids.insert(sub_id.clone());
pagination_state.insert(sub_id, PaginationState::new(grouped_filters)); pagination_state.insert(sub_id, PaginationState::new(grouped_filters));
@@ -6195,6 +6233,10 @@ mod tests {
.next_page(&mut session) .next_page(&mut session)
.expect("a suspiciously short hinted page needs verification"); .expect("a suspiciously short hinted page needs verification");
assert_eq!(session.hint, PaginationHint::Verifying(1000)); assert_eq!(session.hint, PaginationHint::Verifying(1000));
assert_eq!(
verification.request_class(),
TransientRequestClass::PaginationVerification
);
let unseen = EventBuilder::new(Kind::TextNote, "older unseen event") let unseen = EventBuilder::new(Kind::TextNote, "older unseen event")
.custom_created_at(Timestamp::from_secs(99)) .custom_created_at(Timestamp::from_secs(99))
+104 -27
View File
@@ -119,8 +119,43 @@ struct HeldTransientPermits {
_class_cap: tokio::sync::OwnedSemaphorePermit, _class_cap: tokio::sync::OwnedSemaphorePermit,
_ledger_slot: tokio::sync::OwnedSemaphorePermit, _ledger_slot: tokio::sync::OwnedSemaphorePermit,
generation: u64, 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)] #[derive(Clone)]
struct SubscriptionLedgerSession { struct SubscriptionLedgerSession {
generation: u64, generation: u64,
@@ -655,7 +690,10 @@ impl RelayConnection {
/// Acquire both transient constraints without waiting while holding only /// Acquire both transient constraints without waiting while holding only
/// one of them. This prevents class-cap waiters from pinning ledger slots /// one of them. This prevents class-cap waiters from pinning ledger slots
/// and ledger waiters from pinning the class cap. /// and ledger waiters from pinning the class cap.
async fn acquire_transient_permits(&self) -> Result<HeldTransientPermits, String> { async fn acquire_transient_permits(
&self,
request_class: TransientRequestClass,
) -> Result<HeldTransientPermits, String> {
loop { loop {
let ledger_slot = self.acquire_subscription_slots(1).await?; let ledger_slot = self.acquire_subscription_slots(1).await?;
if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() { if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() {
@@ -663,6 +701,7 @@ impl RelayConnection {
_class_cap: class_cap, _class_cap: class_cap,
generation: ledger_slot.generation, generation: ledger_slot.generation,
_ledger_slot: ledger_slot.permit, _ledger_slot: ledger_slot.permit,
request_class,
}); });
} }
drop(ledger_slot); drop(ledger_slot);
@@ -678,6 +717,7 @@ impl RelayConnection {
_class_cap: class_cap, _class_cap: class_cap,
generation: ledger_slot.generation, generation: ledger_slot.generation,
_ledger_slot: ledger_slot.permit, _ledger_slot: ledger_slot.permit,
request_class,
}); });
} }
drop(class_cap); drop(class_cap);
@@ -741,7 +781,7 @@ impl RelayConnection {
filter_groups: Vec<Vec<Filter>>, filter_groups: Vec<Vec<Filter>>,
) -> Result<Vec<SubscriptionId>, String> { ) -> Result<Vec<SubscriptionId>, String> {
self.replace_live_filter_groups_with(filter_groups, |filters, permit| { 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 .await
} }
@@ -1009,7 +1049,7 @@ impl RelayConnection {
/// ///
/// # Arguments /// # Arguments
/// * `filter` - The filter to subscribe to /// * `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 /// # Returns
/// * `Ok(SubscriptionId)` - The subscription ID on success /// * `Ok(SubscriptionId)` - The subscription ID on success
@@ -1017,9 +1057,9 @@ impl RelayConnection {
pub async fn subscribe_filter( pub async fn subscribe_filter(
&self, &self,
filter: Filter, filter: Filter,
auto_close: bool, request_class: TransientRequestClass,
) -> Result<SubscriptionId, String> { ) -> Result<SubscriptionId, String> {
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. /// Subscribe to several OR filters under one NIP-01 subscription ID.
@@ -1030,9 +1070,9 @@ impl RelayConnection {
pub async fn subscribe_filters( pub async fn subscribe_filters(
&self, &self,
filters: Vec<Filter>, filters: Vec<Filter>,
auto_close: bool, request_class: TransientRequestClass,
) -> Result<SubscriptionId, String> { ) -> Result<SubscriptionId, String> {
self.subscribe_filters_with_live_permit(filters, auto_close, None) self.subscribe_filters_with_live_permit(filters, Some(request_class), None)
.await .await
} }
@@ -1044,7 +1084,7 @@ impl RelayConnection {
filter_groups: Vec<Vec<Filter>>, filter_groups: Vec<Vec<Filter>>,
) -> Result<Vec<SubscriptionId>, String> { ) -> Result<Vec<SubscriptionId>, String> {
self.subscribe_live_filter_groups_with(filter_groups, |filters, permit| { 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 .await
} }
@@ -1101,7 +1141,7 @@ impl RelayConnection {
async fn subscribe_filters_with_live_permit( async fn subscribe_filters_with_live_permit(
&self, &self,
filters: Vec<Filter>, filters: Vec<Filter>,
auto_close: bool, transient_class: Option<TransientRequestClass>,
reserved_live_permit: Option<SessionPermit>, reserved_live_permit: Option<SessionPermit>,
) -> Result<SubscriptionId, String> { ) -> Result<SubscriptionId, String> {
if filters.is_empty() { if filters.is_empty() {
@@ -1112,7 +1152,7 @@ impl RelayConnection {
relay = %self.url, relay = %self.url,
filter_count = filters.len(), filter_count = filters.len(),
filters = ?filters, filters = ?filters,
auto_close = auto_close, auto_close = transient_class.is_some(),
"subscribe_filters called" "subscribe_filters called"
); );
@@ -1121,7 +1161,7 @@ impl RelayConnection {
// exceeding relay subscription budgets. The permit is registered // exceeding relay subscription budgets. The permit is registered
// against the subscription id on success and released when the // against the subscription id on success and released when the
// subscription's EOSE or CLOSED arrives (see run_event_loop). // 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() { if self.historic_capacity_consumed_by_live() {
tracing::warn!( tracing::warn!(
relay = %self.url, relay = %self.url,
@@ -1132,7 +1172,10 @@ impl RelayConnection {
self.url self.url
)); ));
} }
(Some(self.acquire_transient_permits().await?), None) (
Some(self.acquire_transient_permits(request_class).await?),
None,
)
} else { } else {
let acquired = match reserved_live_permit { let acquired = match reserved_live_permit {
Some(permit) => Ok(permit), Some(permit) => Ok(permit),
@@ -1162,7 +1205,7 @@ impl RelayConnection {
// Transient permits are acquired through the same session helper; a // Transient permits are acquired through the same session helper; a
// reset closes queued acquisitions and the helper checks generation. // reset closes queued acquisitions and the helper checks generation.
let retained_filters = filters.clone(); let retained_filters = filters.clone();
let output = if auto_close { let output = if transient_class.is_some() {
self.client self.client
.subscribe(filters) .subscribe(filters)
.close_on( .close_on(
@@ -1184,7 +1227,7 @@ impl RelayConnection {
return Err(format!("Failed to subscribe on {}: {}", self.url, failures)); return Err(format!("Failed to subscribe on {}: {}", self.url, failures));
} }
if auto_close { if transient_class.is_some() {
if let Some(permit) = transient_permit { if let Some(permit) = transient_permit {
self.hold_transient_req_permit(output.value.clone(), permit); self.hold_transient_req_permit(output.value.clone(), permit);
} }
@@ -1323,33 +1366,38 @@ impl RelayConnection {
let relay_url = self.url.clone(); let relay_url = self.url.clone();
tokio::spawn(async move { tokio::spawn(async move {
tokio::time::sleep(TRANSIENT_REQ_PERMIT_TIMEOUT).await; tokio::time::sleep(TRANSIENT_REQ_PERMIT_TIMEOUT).await;
let generation = held let held_request = held
.lock() .lock()
.expect("transient permit map poisoned") .expect("transient permit map poisoned")
.get(&sub_id) .get(&sub_id)
.map(|permits| permits.generation); .map(|permits| (permits.generation, permits.request_class));
if let (Some(generation), Ok(Some(relay))) = (generation, client.relay(&url).await) { if let (Some((generation, request_class)), Ok(Some(relay))) =
(held_request, client.relay(&url).await)
{
if relay if relay
.send_msg(ClientMessage::close(sub_id.clone())) .send_msg(ClientMessage::close(sub_id.clone()))
.await .await
.is_ok() .is_ok()
{ {
Self::release_transient_req_permit_for_generation( Self::release_transient_req_permit_for_generation(&held, &sub_id, generation);
&held,
&sub_id,
generation,
);
tracing::warn!( tracing::warn!(
relay = %relay_url, relay = %relay_url,
sub_id = %sub_id, sub_id = %sub_id,
request_class = request_class.as_str(),
"Transient REQ watchdog sent CLOSE after missing EOSE/CLOSED" "Transient REQ watchdog sent CLOSE after missing EOSE/CLOSED"
); );
crate::metrics::record_transient_req_watchdog(request_class.as_str(), "closed");
} else { } else {
tracing::warn!( tracing::warn!(
relay = %relay_url, relay = %relay_url,
sub_id = %sub_id, sub_id = %sub_id,
request_class = request_class.as_str(),
"Transient REQ watchdog could not send CLOSE; retaining subscription slot" "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!(ids.len(), 8);
assert_eq!(connection.subscription_budget().available_permits(), 0); assert_eq!(connection.subscription_budget().available_permits(), 0);
let error = connection let error = connection
.subscribe_filter(Filter::new().kind(Kind::TextNote), true) .subscribe_filter(
Filter::new().kind(Kind::TextNote),
TransientRequestClass::HistoricPage,
)
.await .await
.expect_err("history must defer when complete live coverage fills usable capacity"); .expect_err("history must defer when complete live coverage fills usable capacity");
assert!( assert!(
@@ -2368,6 +2419,7 @@ mod tests {
let generation = connection.current_subscription_generation(); let generation = connection.current_subscription_generation();
let held = HeldTransientPermits { let held = HeldTransientPermits {
generation, generation,
request_class: TransientRequestClass::HistoricPage,
_class_cap: connection _class_cap: connection
.transient_req_permits .transient_req_permits
.clone() .clone()
@@ -2449,8 +2501,11 @@ mod tests {
.await .await
.expect("occupy transient class cap"); .expect("occupy transient class cap");
let class_waiter_connection = connection.clone(); let class_waiter_connection = connection.clone();
let mut class_waiter = let mut class_waiter = tokio::spawn(async move {
tokio::spawn(async move { class_waiter_connection.acquire_transient_permits().await }); class_waiter_connection
.acquire_transient_permits(TransientRequestClass::HistoricPage)
.await
});
assert!( assert!(
tokio::time::timeout(Duration::from_millis(50), &mut class_waiter) tokio::time::timeout(Duration::from_millis(50), &mut class_waiter)
.await .await
@@ -2471,8 +2526,11 @@ mod tests {
.await .await
.expect("occupy subscription ledger"); .expect("occupy subscription ledger");
let ledger_waiter_connection = connection.clone(); let ledger_waiter_connection = connection.clone();
let mut ledger_waiter = let mut ledger_waiter = tokio::spawn(async move {
tokio::spawn(async move { ledger_waiter_connection.acquire_transient_permits().await }); ledger_waiter_connection
.acquire_transient_permits(TransientRequestClass::HistoricPage)
.await
});
assert!( assert!(
tokio::time::timeout(Duration::from_millis(50), &mut ledger_waiter) tokio::time::timeout(Duration::from_millis(50), &mut ledger_waiter)
.await .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::<std::collections::HashSet<_>>();
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] #[test]
fn test_normalize_url_real_world_example() { fn test_normalize_url_real_world_example() {
// Test the exact case from the bug report // Test the exact case from the bug report