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
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
+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
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
+20
View File
@@ -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.
+49 -7
View File
@@ -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<Self> {
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))
+104 -27
View File
@@ -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<HeldTransientPermits, String> {
async fn acquire_transient_permits(
&self,
request_class: TransientRequestClass,
) -> Result<HeldTransientPermits, String> {
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<Vec<Filter>>,
) -> Result<Vec<SubscriptionId>, 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<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.
@@ -1030,9 +1070,9 @@ impl RelayConnection {
pub async fn subscribe_filters(
&self,
filters: Vec<Filter>,
auto_close: bool,
request_class: TransientRequestClass,
) -> 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
}
@@ -1044,7 +1084,7 @@ impl RelayConnection {
filter_groups: Vec<Vec<Filter>>,
) -> Result<Vec<SubscriptionId>, 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<Filter>,
auto_close: bool,
transient_class: Option<TransientRequestClass>,
reserved_live_permit: Option<SessionPermit>,
) -> Result<SubscriptionId, String> {
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::<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]
fn test_normalize_url_real_world_example() {
// Test the exact case from the bug report