From 18213db4cb000f096510141a614f7de765bf22de Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Wed, 12 Aug 2026 22:37:28 +0000 Subject: [PATCH] feat(metrics): expose retained sync state totals The peer-state audit found that important long-lived collections were visible only through incidental logs. That made slow cardinality growth difficult to distinguish from expected event throughput during production soaks. Export one fixed-label aggregate gauge and refresh eight operational classes from the existing two-second health pass: pending batches and subscriptions, auth retries, connection attempts, purgatory dependency attempts, temporary dependency relays, deferred consolidations, and descendant rotations. No relay, event, or repository identifier enters a label. The metric is observational only and introduces no admission limit. Its class vocabulary is deliberately fixed in the sync manager so monitoring cannot recreate the peer-controlled cardinality problem removed by the preceding metrics fix. Validated with a registry test for aggregate values and series count, cargo fmt, and cargo clippy --locked --lib -- -D warnings. --- src/sync/metrics.rs | 44 ++++++++++++++++++++++++++++++++++++++++++++ src/sync/mod.rs | 24 ++++++++++++++++++++++++ 2 files changed, 68 insertions(+) diff --git a/src/sync/metrics.rs b/src/sync/metrics.rs index ecda789..fe063b3 100644 --- a/src/sync/metrics.rs +++ b/src/sync/metrics.rs @@ -42,6 +42,8 @@ pub struct SyncMetrics { relays_connected_total: IntGauge, /// Relays marked as dead relays_dead_total: IntGauge, + /// Aggregate cardinality of long-lived sync-manager state by fixed class. + retained_state_current: IntGaugeVec, // === Rejected Events Index Metrics (unified with event_type label) === /// Current number of entries in hot cache (by event_type: announcement, state) @@ -144,6 +146,15 @@ impl SyncMetrics { ))?; registry.register(Box::new(relays_dead_total.clone()))?; + let retained_state_current = IntGaugeVec::new( + Opts::new( + "ngit_sync_retained_state_current", + "Current retained sync-manager entries by fixed state class", + ), + &["class"], + )?; + registry.register(Box::new(retained_state_current.clone()))?; + // Rejected events metrics (unified with event_type label) let rejected_hot_cache_current = IntGaugeVec::new( Opts::new( @@ -228,6 +239,7 @@ impl SyncMetrics { relays_tracked_total, relays_connected_total, relays_dead_total, + retained_state_current, rejected_hot_cache_current, rejected_hot_cache_hits_total, rejected_hot_cache_misses_total, @@ -425,6 +437,14 @@ impl SyncMetrics { self.relays_dead_total.get() } + /// Record one aggregate long-lived state cardinality. Callers must use the + /// fixed class vocabulary documented by the sync manager. + pub fn set_retained_state(&self, class: &str, count: usize) { + self.retained_state_current + .with_label_values(&[class]) + .set(count as i64); + } + // === Rejected Events Recording Methods (unified with event_type parameter) === /// Update hot cache current size gauge for a specific event type. @@ -663,6 +683,30 @@ mod tests { assert!(metrics2.is_err()); } + #[test] + fn retained_state_metric_uses_fixed_aggregate_classes() { + let registry = create_test_registry(); + let metrics = SyncMetrics::register(®istry).unwrap(); + + metrics.set_retained_state("pending_batches", 12); + metrics.set_retained_state("queued_connection_attempts", 3); + + let family = registry + .gather() + .into_iter() + .find(|family| family.name() == "ngit_sync_retained_state_current") + .expect("retained-state metric must be exported"); + assert_eq!(family.get_metric().len(), 2); + assert_eq!( + family + .get_metric() + .iter() + .map(|metric| metric.get_gauge().value() as i64) + .sum::(), + 15 + ); + } + #[test] fn test_rejected_events_metrics() { let registry = create_test_registry(); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index c5cf41f..a796af4 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1803,6 +1803,30 @@ async fn run_health_and_metrics_checker( let entries = naughty_list.get_all(); metrics.update_naughty_list(entries); } + + let (pending_batches, pending_subscriptions) = { + let pending = manager.pending_sync_index.read().await; + ( + pending.values().map(Vec::len).sum(), + pending + .values() + .flatten() + .map(|batch| batch.outstanding_subs.len()) + .sum(), + ) + }; + for (class, count) in [ + ("pending_batches", pending_batches), + ("pending_subscriptions", pending_subscriptions), + ("auth_retries", manager.auth_required_attempts.len()), + ("queued_connection_attempts", manager.in_flight_connect_attempts.len()), + ("purgatory_dependencies", manager.purgatory_dependency_attempts.len()), + ("dependency_relays", manager.dependency_relay_deadlines.len()), + ("deferred_consolidations", manager.deferred_consolidations.relays.len()), + ("descendant_rotations", manager.descendant_sync_rotations.len()), + ] { + metrics.set_retained_state(class, count); + } } } _ = shutdown_rx.recv() => {