From 42aa16187a3cb498e5fde48d6ebf2b51215c64b7 Mon Sep 17 00:00:00 2001 From: Aubaid Ahmed Saiyed Date: Sun, 27 Sep 2026 17:55:16 +0530 Subject: [PATCH 1/3] obsv: wire live per-database metrics sampler (fixes #375) --- nodedb-mem/src/governor/metrics.rs | 27 ++++++ nodedb/src/bootstrap/background_loops.rs | 103 +++++++++++++++++++++++ nodedb/src/bridge/dispatch/dispatcher.rs | 16 ++++ nodedb/src/control/metrics/database.rs | 34 ++++++++ nodedb/src/wal/manager/core.rs | 12 +++ nodedb/src/wal/manager/durable_commit.rs | 12 ++- 6 files changed, 201 insertions(+), 3 deletions(-) diff --git a/nodedb-mem/src/governor/metrics.rs b/nodedb-mem/src/governor/metrics.rs index 15a1ba696..b6c6e3757 100644 --- a/nodedb-mem/src/governor/metrics.rs +++ b/nodedb-mem/src/governor/metrics.rs @@ -5,11 +5,25 @@ use std::sync::atomic::Ordering; +use nodedb_types::DatabaseId; + use super::core::MemoryGovernor; use crate::engine::EngineId; use crate::pressure::{PressureLevel, PressureThresholds}; impl MemoryGovernor { + /// Current allocated memory in bytes for a database. + /// + /// Returns 0 if the database has no scoped budget or no active allocations. + pub fn database_usage_bytes(&self, db: DatabaseId) -> usize { + self.database_budgets + .read() + .unwrap_or_else(|p| p.into_inner()) + .get(&db) + .map(|b| b.allocated.load(Ordering::Relaxed)) + .unwrap_or(0) + } + /// Total memory allocated across all engines (engine-layer sum). A /// separate aggregate from [`global_utilization_percent`](Self::global_utilization_percent), /// which reads the global counter admission enforces the ceiling against. @@ -164,4 +178,17 @@ mod tests { assert_eq!(gov.engine_pressure(EngineId::Vector), PressureLevel::Normal); assert_eq!(gov.worst_engine_pressure(), PressureLevel::Critical); } + + #[test] + fn database_usage_bytes_tracks_allocation() { + let gov = MemoryGovernor::new(test_config()).unwrap(); + gov.set_database_budget(db(), 10_000); + assert_eq!(gov.database_usage_bytes(db()), 0); + + let _tok = gov + .try_reserve(db(), tenant(), EngineId::Vector, 512) + .unwrap(); + assert_eq!(gov.database_usage_bytes(db()), 512); + } } + diff --git a/nodedb/src/bootstrap/background_loops.rs b/nodedb/src/bootstrap/background_loops.rs index 17d95db11..83556fb1d 100644 --- a/nodedb/src/bootstrap/background_loops.rs +++ b/nodedb/src/bootstrap/background_loops.rs @@ -142,6 +142,44 @@ pub fn spawn_background_loops( info!("mirror lag monitor running"); } + // Database metrics sampler (10-second interval). + // Samples live per-database metrics across subsystems and updates + // `DatabaseMetricsRegistry` gauges for Prometheus scraping. + { + let shared_sampler = Arc::clone(shared); + crate::control::shutdown::spawn_loop( + &shared.loop_registry, + &shared.shutdown, + "database_metrics_sampler", + crate::control::shutdown::ShutdownPhase::DrainingControlPlane, + move |mut shutdown| async move { + let mut tick = tokio::time::interval(Duration::from_secs(10)); + loop { + tokio::select! { + _ = shutdown.wait_cancelled() => break, + _ = tick.tick() => {} + } + if shutdown.is_cancelled() { + break; + } + let catalog = shared_sampler.credentials.catalog(); + let databases = match catalog.list_databases() { + Ok(d) => d, + Err(e) => { + tracing::warn!(error = %e, "database_metrics_sampler: catalog list error"); + continue; + } + }; + for db in databases { + sample_database_metrics(&shared_sampler, db.id, &db.name); + } + } + }, + ); + info!("database metrics sampler running"); + } + + // Wire stream delivery managers before Event Plane creation can admit // CREATE CHANGE STREAM delivery tasks. shared.webhook_manager.set_state(Arc::clone(shared)); @@ -423,3 +461,68 @@ pub fn spawn_response_poller( }, ); } + +/// Sample live metrics for a single database and update its gauges in `shared.database_metrics`. +pub fn sample_database_metrics( + shared: &SharedState, + db_id: nodedb_types::DatabaseId, + db_name: &str, +) { + // 1. connections: active connections holding an admission permit for this database + // in `shared.admission_registry`. + let connections = shared + .admission_registry + .database_live_connections(db_id) + .map(u64::from) + .unwrap_or(0); + shared + .database_metrics + .set_connections(db_name, connections); + + // 2. memory: resident memory in bytes allocated by this database, read directly + // from the `database_budgets` map in `shared.governor`. + let memory_bytes = shared.governor.database_usage_bytes(db_id) as u64; + shared + .database_metrics + .set_memory_bytes(db_name, memory_bytes); + + // 3. storage: per-database storage usage in bytes from `shared.system_metrics` + // (the same source read by `SHOW DATABASE USAGE`). + let storage_bytes = shared + .system_metrics + .as_ref() + .map(|m| m.database_storage_bytes(db_name)) + .unwrap_or(0); + shared + .database_metrics + .set_storage_bytes(db_name, storage_bytes); + + // 4. bridge_queue_depth: sum of SPSC bridge virtual-queue depths for this database + // across all Data Plane cores in `shared.dispatcher`. + let bridge_queue_depth = match shared.dispatcher.lock() { + Ok(d) => d.virtual_queue_depth(db_id.as_u64()), + Err(p) => p.into_inner().virtual_queue_depth(db_id.as_u64()), + }; + shared + .database_metrics + .set_bridge_queue_depth(db_name, bridge_queue_depth); + + // 5. wal_latency_p99: P99 WAL group-commit fsync latency in microseconds from + // `shared.wal`'s commit latency histogram (falling back to `system_metrics`). + let wal_latency_p99 = { + let from_wal = shared.wal.commit_latency_p99_us(); + if from_wal > 0 { + from_wal + } else { + shared + .system_metrics + .as_ref() + .map(|s| s.wal_fsync_seconds.percentile(99.0)) + .unwrap_or(0) + } + }; + shared + .database_metrics + .set_wal_latency_p99(db_name, wal_latency_p99); +} + diff --git a/nodedb/src/bridge/dispatch/dispatcher.rs b/nodedb/src/bridge/dispatch/dispatcher.rs index 7e1783f3e..f92779699 100644 --- a/nodedb/src/bridge/dispatch/dispatcher.rs +++ b/nodedb/src/bridge/dispatch/dispatcher.rs @@ -336,6 +336,15 @@ impl Dispatcher { .unwrap_or(PressureState::Normal) } + /// Sum of bridge virtual-queue depths for a database across all cores. + pub fn virtual_queue_depth(&self, database_id: u64) -> u64 { + self.cores + .iter() + .map(|c| c.wfq.depth_for(database_id) as u64) + .sum() + } + + /// Poll responses from all Data Plane cores. /// /// A core whose channel has been observed dead contributes a synthesized @@ -835,4 +844,11 @@ mod tests { assert!(r.error_code.is_some()); } } + + #[test] + fn virtual_queue_depth_reporting() { + let (dispatcher, _data_sides) = Dispatcher::new(2, 8); + assert_eq!(dispatcher.virtual_queue_depth(1), 0); + } } + diff --git a/nodedb/src/control/metrics/database.rs b/nodedb/src/control/metrics/database.rs index 75abd4a5d..1f1409f6f 100644 --- a/nodedb/src/control/metrics/database.rs +++ b/nodedb/src/control/metrics/database.rs @@ -135,7 +135,14 @@ impl DatabaseMetricsRegistry { /// Converts to integer microseconds (rounded) to avoid `AtomicF64` complexity. /// Negative or non-finite inputs are clamped to zero so accidental underflow /// in upstream timing arithmetic cannot subtract from the cumulative counter. + /// + /// NOTE: Intentionally unfilled by `database_metrics_sampler`. Unlike the five + /// gauge metrics (connections, memory, storage, bridge queue depth, and WAL + /// latency P99) which represent point-in-time state, this is a cumulative + /// counter tracking task execution time. It remains wired for callers + /// pending a dedicated per-task maintenance completion accounting source. pub fn add_maintenance_cpu_secs(&self, db_name: &str, secs: f64) { + let us = if secs.is_finite() && secs > 0.0 { (secs * 1_000_000.0).round() as u64 } else { @@ -326,4 +333,31 @@ mod tests { }; assert!(m.is_over_quota()); } + + #[test] + fn gauge_setters_update_counters() { + let reg = DatabaseMetricsRegistry::new(); + reg.set_connections("db1", 42); + reg.set_memory_bytes("db1", 1024 * 1024); + reg.set_storage_bytes("db1", 10 * 1024 * 1024); + reg.set_bridge_queue_depth("db1", 7); + reg.set_wal_latency_p99("db1", 1500); + + let c = reg.get_or_create("db1"); + assert_eq!(c.connections.load(Ordering::Relaxed), 42); + assert_eq!(c.memory_bytes.load(Ordering::Relaxed), 1024 * 1024); + assert_eq!(c.storage_bytes.load(Ordering::Relaxed), 10 * 1024 * 1024); + assert_eq!(c.bridge_queue_depth.load(Ordering::Relaxed), 7); + assert_eq!(c.wal_commit_latency_p99_us.load(Ordering::Relaxed), 1500); + } + + #[test] + fn add_maintenance_cpu_secs_accumulates() { + let reg = DatabaseMetricsRegistry::new(); + reg.add_maintenance_cpu_secs("db1", 1.5); + reg.add_maintenance_cpu_secs("db1", 0.5); + let c = reg.get_or_create("db1"); + assert_eq!(c.maintenance_cpu_seconds_total.load(Ordering::Relaxed), 2_000_000); + } } + diff --git a/nodedb/src/wal/manager/core.rs b/nodedb/src/wal/manager/core.rs index 9bd812bea..5a6243154 100644 --- a/nodedb/src/wal/manager/core.rs +++ b/nodedb/src/wal/manager/core.rs @@ -45,6 +45,8 @@ pub struct WalManager { /// Wakes `wait_durable` followers when `durable_lsn` advances (or a leader's /// fsync fails, so they re-attempt and observe the same error). pub(super) durable_notify: tokio::sync::Notify, + /// WAL group-commit fsync latency distribution. + pub(super) commit_latency: crate::control::metrics::AtomicHistogram, } impl WalManager { @@ -63,6 +65,12 @@ impl WalManager { self.durable_lsn.load(std::sync::atomic::Ordering::Acquire) } + /// WAL group-commit fsync latency P99 in microseconds. + pub fn commit_latency_p99_us(&self) -> u64 { + self.commit_latency.percentile(99.0) + } + + /// Return the stable in-memory root for per-user CRDT signing keys. /// The root is persisted only as WAL-key-wrapped ciphertext and is /// rewrapped on rotation, so offline signatures survive chained key @@ -172,9 +180,13 @@ impl WalManager { durable_lsn: AtomicU64::new(0), commit_lock: tokio::sync::Mutex::new(()), durable_notify: tokio::sync::Notify::new(), + commit_latency: crate::control::metrics::AtomicHistogram::with_buckets( + crate::control::metrics::histogram::WAL_FSYNC_BUCKETS_US, + ), }) } + /// Open without `O_DIRECT`. /// /// The in-process test harnesses put their data directories in tempdirs, diff --git a/nodedb/src/wal/manager/durable_commit.rs b/nodedb/src/wal/manager/durable_commit.rs index 452a339b0..a5e8c6a25 100644 --- a/nodedb/src/wal/manager/durable_commit.rs +++ b/nodedb/src/wal/manager/durable_commit.rs @@ -63,18 +63,21 @@ impl WalManager { // advances to exactly what the fsync made durable, never past // it. let wal = std::sync::Arc::clone(&self.wal); - let join = tokio::task::spawn_blocking(move || -> crate::Result { + let join = tokio::task::spawn_blocking(move || -> crate::Result<(u64, u64)> { + let start = std::time::Instant::now(); let mut guard = wal.lock().unwrap_or_else(|p| p.into_inner()); guard.sync().map_err(crate::Error::Wal)?; + let elapsed_us = start.elapsed().as_micros() as u64; // `next_lsn()` is the next LSN to assign; the highest LSN // this sync made durable is one below it. - Ok(guard.next_lsn().saturating_sub(1)) + Ok((guard.next_lsn().saturating_sub(1), elapsed_us)) }) .await; let outcome = match join { - Ok(Ok(durable_through)) => { + Ok(Ok((durable_through, elapsed_us))) => { self.durable_lsn.fetch_max(durable_through, AcqRel); + self.commit_latency.observe(elapsed_us); Ok(()) } // fsync error: do NOT advance `durable_lsn`. @@ -86,6 +89,7 @@ impl WalManager { }), }; + // Wake followers on every leader-exit path so a failed or // panicked fsync never strands them on `notified.await`. On // success `durable_lsn` was advanced first, so they observe @@ -146,10 +150,12 @@ mod tests { ) .expect("append"); wal.wait_durable(lsn).await.expect("first"); + assert!(wal.commit_latency.count() >= 1); // Second call is the fast path — no further fsync required. wal.wait_durable(lsn).await.expect("second"); } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn concurrent_waiters_coalesce() { let dir = tempfile::tempdir().expect("tempdir"); From 23e8899f4a4aee8b590e0ac89fdca175b4777499 Mon Sep 17 00:00:00 2001 From: Aubaid Ahmed Saiyed Date: Sun, 27 Sep 2026 18:38:43 +0530 Subject: [PATCH 2/3] fix(obsv): MaintenanceLease now resolves a database name (via boot hydration and DDL post-apply hooks) and forwards elapsed CPU time to DatabaseMetricsRegistry on drop, completing the sixth metric from #375. --- .../catalog_entry/post_apply/database.rs | 13 +- .../control/catalog_entry/post_apply/quota.rs | 12 +- nodedb/src/control/maintenance/budget.rs | 154 +++++++++++++++++- nodedb/src/control/metrics/database.rs | 7 +- nodedb/src/control/state/init.rs | 11 ++ .../src/control/state/init_prod/post_init.rs | 11 ++ 6 files changed, 196 insertions(+), 12 deletions(-) diff --git a/nodedb/src/control/catalog_entry/post_apply/database.rs b/nodedb/src/control/catalog_entry/post_apply/database.rs index 65964ab61..fb9bff19f 100644 --- a/nodedb/src/control/catalog_entry/post_apply/database.rs +++ b/nodedb/src/control/catalog_entry/post_apply/database.rs @@ -13,12 +13,19 @@ use nodedb_types::DatabaseId; use crate::control::security::catalog::database_types::DatabaseDescriptor; use crate::control::state::SharedState; -/// Post-apply for `PutDatabase` — no in-memory cache to update. -pub fn put(_descriptor: DatabaseDescriptor, _shared: Arc) {} +/// Post-apply for `PutDatabase` — update in-memory maintenance budget name resolution. +pub fn put(descriptor: DatabaseDescriptor, shared: Arc) { + shared + .maintenance_budget + .set_database_name(descriptor.id, &descriptor.name); +} /// Post-apply for `DeleteDatabase` — release the quota caps of the dropped -/// scope, its tenants' caps included. +/// scope, its tenants' caps included, and unregister from maintenance budget. pub fn delete(db_id: u64, shared: Arc) { + shared + .maintenance_budget + .remove_database(DatabaseId::new(db_id)); super::quota::release_database_scope(DatabaseId::new(db_id), &shared); } diff --git a/nodedb/src/control/catalog_entry/post_apply/quota.rs b/nodedb/src/control/catalog_entry/post_apply/quota.rs index 2e81debd4..594ce2bb8 100644 --- a/nodedb/src/control/catalog_entry/post_apply/quota.rs +++ b/nodedb/src/control/catalog_entry/post_apply/quota.rs @@ -18,9 +18,15 @@ use crate::control::state::SharedState; /// Install a database quota into live enforcement. pub fn put_database(db_id: DatabaseId, record: &QuotaRecord, shared: &SharedState) { - shared - .maintenance_budget - .set_cap(db_id, record.maintenance_cpu_pct); + if let Ok(Some(name)) = shared.credentials.catalog().get_database_name_by_id(db_id) { + shared + .maintenance_budget + .set_cap_named(db_id, name, record.maintenance_cpu_pct); + } else { + shared + .maintenance_budget + .set_cap(db_id, record.maintenance_cpu_pct); + } if record.max_memory_bytes > 0 { shared .governor diff --git a/nodedb/src/control/maintenance/budget.rs b/nodedb/src/control/maintenance/budget.rs index 697164948..b42d9a3b0 100644 --- a/nodedb/src/control/maintenance/budget.rs +++ b/nodedb/src/control/maintenance/budget.rs @@ -19,6 +19,8 @@ use std::time::Instant; use nodedb_types::DatabaseId; +use crate::control::metrics::DatabaseMetricsRegistry; + /// Single slot in the 60-bucket sliding window. #[derive(Clone, Default)] struct Bucket { @@ -82,6 +84,10 @@ struct TrackerInner { /// Per-database CPU-seconds cap per minute. /// Derived from `maintenance_cpu_pct / 100.0 * 60.0`. caps: HashMap, + /// Per-database human-readable names for observability and metrics routing. + names: HashMap, + /// Optional metrics registry handle for per-database Prometheus counters. + metrics: Option>, } impl TrackerInner { @@ -89,6 +95,8 @@ impl TrackerInner { Self { windows: HashMap::new(), caps: HashMap::new(), + names: HashMap::new(), + metrics: None, } } @@ -99,6 +107,16 @@ impl TrackerInner { fn window_for_mut(&mut self, db: DatabaseId) -> &mut DbWindow { self.windows.entry(db).or_insert_with(DbWindow::new) } + + fn resolve_name(&self, db: DatabaseId) -> String { + if let Some(name) = self.names.get(&db) { + return name.clone(); + } + if db == DatabaseId::DEFAULT { + return "default".to_string(); + } + format!("db-{}", db.as_u64()) + } } impl std::fmt::Debug for MaintenanceBudgetTracker { @@ -115,6 +133,39 @@ impl MaintenanceBudgetTracker { } } + /// Create a new tracker pre-wired with a database metrics registry. + pub fn with_metrics(metrics: Arc) -> Self { + let tracker = Self::new(); + tracker.set_metrics(metrics); + tracker + } + + /// Wire or replace the database metrics registry handle. + pub fn set_metrics(&self, metrics: Arc) { + let mut inner = self.inner.lock().unwrap_or_else(|p| p.into_inner()); + inner.metrics = Some(metrics); + } + + /// Optional handle to the currently installed metrics registry. + pub fn metrics(&self) -> Option> { + let inner = self.inner.lock().unwrap_or_else(|p| p.into_inner()); + inner.metrics.clone() + } + + /// Associate a human-readable database name with `db`. + pub fn set_database_name(&self, db: DatabaseId, name: impl Into) { + let mut inner = self.inner.lock().unwrap_or_else(|p| p.into_inner()); + inner.names.insert(db, name.into()); + } + + /// Remove a database's budget window, cap, and name mapping (e.g. on database drop). + pub fn remove_database(&self, db: DatabaseId) { + let mut inner = self.inner.lock().unwrap_or_else(|p| p.into_inner()); + inner.names.remove(&db); + inner.caps.remove(&db); + inner.windows.remove(&db); + } + /// Install or replace the maintenance CPU cap for `db`. /// /// `maintenance_cpu_pct` is the `QuotaRecord` field (0–100). @@ -129,8 +180,23 @@ impl MaintenanceBudgetTracker { inner.caps.insert(db, cap); } + /// Install or replace the maintenance CPU cap and register the database name for `db`. + pub fn set_cap_named(&self, db: DatabaseId, name: impl Into, maintenance_cpu_pct: u8) { + let cap = if maintenance_cpu_pct == 0 { + f64::INFINITY + } else { + (maintenance_cpu_pct as f64 / 100.0) * 60.0 + }; + let mut inner = self.inner.lock().unwrap_or_else(|p| p.into_inner()); + inner.caps.insert(db, cap); + inner.names.insert(db, name.into()); + } + /// Attempt to acquire a maintenance lease for `db`. /// + /// Resolves `db_name` from registered database names, falling back to + /// `"default"` for `DatabaseId::DEFAULT` or `"db-{id}"` if unregistered. + /// /// Returns `Some(MaintenanceLease)` when `consumed + estimated_secs ≤ cap` /// for the current 60-second window. Returns `None` when the database is /// over its budget (caller should defer the task to the next window). @@ -142,6 +208,37 @@ impl MaintenanceBudgetTracker { self: &Arc, db: DatabaseId, estimated_secs: f64, + ) -> Option { + self.try_acquire_internal(db, None, estimated_secs, None) + } + + /// Attempt to acquire a maintenance lease for `db` with an explicitly provided database name. + pub fn try_acquire_named( + self: &Arc, + db: DatabaseId, + db_name: &str, + estimated_secs: f64, + ) -> Option { + self.try_acquire_internal(db, Some(db_name.to_string()), estimated_secs, None) + } + + /// Attempt to acquire a maintenance lease for `db` with optional overrides for name and metrics. + pub fn try_acquire_with_metrics( + self: &Arc, + db: DatabaseId, + db_name: Option, + estimated_secs: f64, + metrics: Option>, + ) -> Option { + self.try_acquire_internal(db, db_name, estimated_secs, metrics) + } + + fn try_acquire_internal( + self: &Arc, + db: DatabaseId, + explicit_name: Option, + estimated_secs: f64, + explicit_metrics: Option>, ) -> Option { let now_secs = current_secs(); let mut inner = self.inner.lock().unwrap_or_else(|p| p.into_inner()); @@ -151,10 +248,14 @@ impl MaintenanceBudgetTracker { let consumed = window.window_total(); if consumed + estimated_secs <= cap { + let db_name = explicit_name.unwrap_or_else(|| inner.resolve_name(db)); + let metrics = explicit_metrics.or_else(|| inner.metrics.clone()); Some(MaintenanceLease { tracker: Arc::clone(self), db, + db_name, start: Instant::now(), + metrics, }) } else { None @@ -171,11 +272,26 @@ impl Default for MaintenanceBudgetTracker { /// RAII lease returned by [`MaintenanceBudgetTracker::try_acquire`]. /// /// On drop, the actual elapsed wall-clock seconds are recorded into the -/// sliding window for the database. +/// sliding window for the database and, if a [`DatabaseMetricsRegistry`] +/// is wired, forwarded to `add_maintenance_cpu_secs`. pub struct MaintenanceLease { tracker: Arc, db: DatabaseId, + db_name: String, start: Instant, + metrics: Option>, +} + +impl MaintenanceLease { + /// Database id for which the lease was acquired. + pub fn db(&self) -> DatabaseId { + self.db + } + + /// Resolved database name. + pub fn db_name(&self) -> &str { + &self.db_name + } } impl Drop for MaintenanceLease { @@ -184,6 +300,9 @@ impl Drop for MaintenanceLease { let now_secs = current_secs(); let mut inner = self.tracker.inner.lock().unwrap_or_else(|p| p.into_inner()); inner.window_for_mut(self.db).record(now_secs, elapsed); + if let Some(m) = &self.metrics { + m.add_maintenance_cpu_secs(&self.db_name, elapsed); + } } } @@ -191,6 +310,7 @@ impl std::fmt::Debug for MaintenanceLease { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("MaintenanceLease") .field("db", &self.db) + .field("db_name", &self.db_name) .finish() } } @@ -294,4 +414,36 @@ mod tests { "old consumption should have expired out of the window" ); } + + #[test] + fn lease_drop_increments_prometheus_metrics() { + let reg = Arc::new(DatabaseMetricsRegistry::new()); + let t = Arc::new(MaintenanceBudgetTracker::with_metrics(Arc::clone(®))); + let db = DatabaseId::new(42); + t.set_database_name(db, "analytics"); + t.set_cap(db, 100); + + { + let lease = t.try_acquire(db, 1.0).expect("acquire should succeed"); + assert_eq!(lease.db_name(), "analytics"); + assert_eq!(lease.db(), db); + std::thread::sleep(std::time::Duration::from_millis(15)); + // Lease drops here, forwarding elapsed CPU time to DatabaseMetricsRegistry. + } + + let c = reg.get_or_create("analytics"); + let recorded_us = c + .maintenance_cpu_seconds_total + .load(std::sync::atomic::Ordering::Relaxed); + assert!( + recorded_us > 0, + "cumulative maintenance CPU counter should be > 0 after lease drop" + ); + + let mut out = String::new(); + reg.render_prometheus(&mut out); + assert!(out.contains(r#"database="analytics""#)); + assert!(out.contains("nodedb_database_maintenance_cpu_us_total")); + } } + diff --git a/nodedb/src/control/metrics/database.rs b/nodedb/src/control/metrics/database.rs index 1f1409f6f..d70790cbc 100644 --- a/nodedb/src/control/metrics/database.rs +++ b/nodedb/src/control/metrics/database.rs @@ -136,11 +136,8 @@ impl DatabaseMetricsRegistry { /// Negative or non-finite inputs are clamped to zero so accidental underflow /// in upstream timing arithmetic cannot subtract from the cumulative counter. /// - /// NOTE: Intentionally unfilled by `database_metrics_sampler`. Unlike the five - /// gauge metrics (connections, memory, storage, bridge queue depth, and WAL - /// latency P99) which represent point-in-time state, this is a cumulative - /// counter tracking task execution time. It remains wired for callers - /// pending a dedicated per-task maintenance completion accounting source. + /// Called from [`crate::control::maintenance::MaintenanceLease::drop`] + /// when a maintenance lease finishes and records elapsed wall-clock CPU time. pub fn add_maintenance_cpu_secs(&self, db_name: &str, secs: f64) { let us = if secs.is_finite() && secs > 0.0 { diff --git a/nodedb/src/control/state/init.rs b/nodedb/src/control/state/init.rs index 0184c9413..c108e7261 100644 --- a/nodedb/src/control/state/init.rs +++ b/nodedb/src/control/state/init.rs @@ -521,6 +521,17 @@ impl SharedState { startup: Arc::clone(&startup_gate), }); Self::wire_session_handle_audit(&state); + state + .maintenance_budget + .set_metrics(Arc::clone(&state.database_metrics)); + state + .maintenance_budget + .set_database_name(nodedb_types::DatabaseId::DEFAULT, "default"); + if let Ok(databases) = state.credentials.catalog().list_databases() { + for db in databases { + state.maintenance_budget.set_database_name(db.id, &db.name); + } + } Ok(state) } diff --git a/nodedb/src/control/state/init_prod/post_init.rs b/nodedb/src/control/state/init_prod/post_init.rs index a51bf907a..4eb5e6967 100644 --- a/nodedb/src/control/state/init_prod/post_init.rs +++ b/nodedb/src/control/state/init_prod/post_init.rs @@ -33,6 +33,17 @@ pub(super) fn hydrate_caches(state: &Arc) { "boot: failed to populate idle_timeout_cache from catalog" ); } + state + .maintenance_budget + .set_metrics(Arc::clone(&state.database_metrics)); + state + .maintenance_budget + .set_database_name(nodedb_types::DatabaseId::DEFAULT, "default"); + if let Ok(databases) = catalog.list_databases() { + for db in databases { + state.maintenance_budget.set_database_name(db.id, &db.name); + } + } } /// Spawn the array GC background task. The handle is stored by the caller From 1155e96c6d4a9d2cf1f479efa0190b7598ca0571 Mon Sep 17 00:00:00 2001 From: Aubaid Ahmed Saiyed Date: Mon, 28 Sep 2026 19:42:57 +0530 Subject: [PATCH 3/3] fix(obsv):address PR #375 review for database metrics and sampler --- nodedb-mem/src/governor/metrics.rs | 1 - nodedb-mem/src/governor/reserve.rs | 100 +++++-- nodedb-mem/src/scoped_budget.rs | 4 +- nodedb/src/bootstrap/background_loops.rs | 100 +------ .../src/bootstrap/database_metrics_sampler.rs | 202 +++++++++++++ nodedb/src/bootstrap/mod.rs | 1 + nodedb/src/bridge/dispatch/dispatcher.rs | 44 ++- .../catalog_entry/post_apply/database.rs | 14 +- nodedb/src/control/maintenance/budget.rs | 1 - nodedb/src/control/metrics/database.rs | 58 +++- nodedb/src/control/metrics/histogram.rs | 141 +++++++++ .../src/control/security/catalog/database.rs | 274 ++++++++++++++++++ .../src/control/server/admission/registry.rs | 23 +- .../shared/ddl/neutral/observability.rs | 8 +- nodedb/src/control/state/init.rs | 3 + .../src/control/state/init_prod/post_init.rs | 3 + nodedb/src/wal/manager/core.rs | 25 +- nodedb/src/wal/manager/durable_commit.rs | 23 +- .../tests/inproc/cases/quota_drop_cleanup.rs | 10 +- .../cases/quota_live_enforcement_apply.rs | 10 +- 20 files changed, 866 insertions(+), 179 deletions(-) create mode 100644 nodedb/src/bootstrap/database_metrics_sampler.rs diff --git a/nodedb-mem/src/governor/metrics.rs b/nodedb-mem/src/governor/metrics.rs index b6c6e3757..7593d9eaf 100644 --- a/nodedb-mem/src/governor/metrics.rs +++ b/nodedb-mem/src/governor/metrics.rs @@ -191,4 +191,3 @@ mod tests { assert_eq!(gov.database_usage_bytes(db()), 512); } } - diff --git a/nodedb-mem/src/governor/reserve.rs b/nodedb-mem/src/governor/reserve.rs index 1b3cd5dfb..8335ab49a 100644 --- a/nodedb-mem/src/governor/reserve.rs +++ b/nodedb-mem/src/governor/reserve.rs @@ -13,6 +13,7 @@ use crate::error::{MemError, Result}; use crate::over_release::ReleaseIdentity; use crate::reservation_token::{ReservationParams, ReservationToken}; use crate::reserve_scope::{ReserveScope, ReservedLayers}; +use crate::scoped_budget::ScopedBudget; /// Build the token both entry points return. They differ in how they reach /// a committed [`ReservedLayers`], never in what they build from one. @@ -62,21 +63,34 @@ impl MemoryGovernor { scope.try_credit_global()?; { - let map = self - .database_budgets - .read() - .unwrap_or_else(|p| p.into_inner()); - if let Some(budget) = map.get(&db) { - match budget.try_reserve(size) { - Ok(arc) => scope.credit_database(arc), - Err(denied) => { - return Err(MemError::DatabaseBudgetExhausted { - db, - requested: size, - available: budget.available(), - limit: denied.limit, - }); - } + let budget = { + let map = self + .database_budgets + .read() + .unwrap_or_else(|p| p.into_inner()); + map.get(&db).cloned() + }; + let budget = match budget { + Some(b) => b, + None => { + let mut map = self + .database_budgets + .write() + .unwrap_or_else(|p| p.into_inner()); + map.entry(db) + .or_insert_with(|| ScopedBudget::new(None)) + .clone() + } + }; + match budget.try_reserve(size) { + Ok(arc) => scope.credit_database(arc), + Err(denied) => { + return Err(MemError::DatabaseBudgetExhausted { + db, + requested: size, + available: budget.available(), + limit: denied.limit, + }); } } } @@ -138,13 +152,26 @@ impl MemoryGovernor { scope.credit_global_unchecked(); { - let map = self - .database_budgets - .read() - .unwrap_or_else(|p| p.into_inner()); - if let Some(budget) = map.get(&db) { - scope.credit_database(budget.credit(size)); - } + let budget = { + let map = self + .database_budgets + .read() + .unwrap_or_else(|p| p.into_inner()); + map.get(&db).cloned() + }; + let budget = match budget { + Some(b) => b, + None => { + let mut map = self + .database_budgets + .write() + .unwrap_or_else(|p| p.into_inner()); + map.entry(db) + .or_insert_with(|| ScopedBudget::new(None)) + .clone() + } + }; + scope.credit_database(budget.credit(size)); } { @@ -493,4 +520,33 @@ mod tests { tok.err() ); } + + #[test] + fn uncapped_database_usage_tracks_allocation() { + let gov = MemoryGovernor::new(test_config()).unwrap(); + assert_eq!(gov.database_usage_bytes(db()), 0); + + let tok = gov + .try_reserve(db(), tenant(), EngineId::Vector, 512) + .unwrap(); + assert_eq!(gov.database_usage_bytes(db()), 512); + + drop(tok); + assert_eq!(gov.database_usage_bytes(db()), 0); + } + + #[test] + fn unbudgeted_database_tracks_8192_bytes_and_releases_on_drop() { + let gov = MemoryGovernor::new(test_config()).unwrap(); + let unbudgeted_db = DatabaseId::new(42); + assert_eq!(gov.database_usage_bytes(unbudgeted_db), 0); + + let token = gov + .try_reserve(unbudgeted_db, tenant(), EngineId::Vector, 8192) + .unwrap(); + assert_eq!(gov.database_usage_bytes(unbudgeted_db), 8192); + + drop(token); + assert_eq!(gov.database_usage_bytes(unbudgeted_db), 0); + } } diff --git a/nodedb-mem/src/scoped_budget.rs b/nodedb-mem/src/scoped_budget.rs index d756898ed..2384c7d64 100644 --- a/nodedb-mem/src/scoped_budget.rs +++ b/nodedb-mem/src/scoped_budget.rs @@ -11,7 +11,7 @@ use std::sync::atomic::{AtomicUsize, Ordering}; /// /// Quota changes mutate `limit` in place so live tokens keep decrementing the /// same counter. -#[derive(Debug)] +#[derive(Debug, Clone)] pub(crate) struct ScopedBudget { /// `None` means uncapped, still counted. pub(crate) limit: Option, @@ -25,7 +25,7 @@ pub(crate) struct BudgetDenied { } impl ScopedBudget { - fn new(limit: Option) -> Self { + pub(crate) fn new(limit: Option) -> Self { Self { limit, allocated: Arc::new(AtomicUsize::new(0)), diff --git a/nodedb/src/bootstrap/background_loops.rs b/nodedb/src/bootstrap/background_loops.rs index 83556fb1d..27073a470 100644 --- a/nodedb/src/bootstrap/background_loops.rs +++ b/nodedb/src/bootstrap/background_loops.rs @@ -145,40 +145,7 @@ pub fn spawn_background_loops( // Database metrics sampler (10-second interval). // Samples live per-database metrics across subsystems and updates // `DatabaseMetricsRegistry` gauges for Prometheus scraping. - { - let shared_sampler = Arc::clone(shared); - crate::control::shutdown::spawn_loop( - &shared.loop_registry, - &shared.shutdown, - "database_metrics_sampler", - crate::control::shutdown::ShutdownPhase::DrainingControlPlane, - move |mut shutdown| async move { - let mut tick = tokio::time::interval(Duration::from_secs(10)); - loop { - tokio::select! { - _ = shutdown.wait_cancelled() => break, - _ = tick.tick() => {} - } - if shutdown.is_cancelled() { - break; - } - let catalog = shared_sampler.credentials.catalog(); - let databases = match catalog.list_databases() { - Ok(d) => d, - Err(e) => { - tracing::warn!(error = %e, "database_metrics_sampler: catalog list error"); - continue; - } - }; - for db in databases { - sample_database_metrics(&shared_sampler, db.id, &db.name); - } - } - }, - ); - info!("database metrics sampler running"); - } - + crate::bootstrap::database_metrics_sampler::spawn_database_metrics_sampler(shared); // Wire stream delivery managers before Event Plane creation can admit // CREATE CHANGE STREAM delivery tasks. @@ -461,68 +428,3 @@ pub fn spawn_response_poller( }, ); } - -/// Sample live metrics for a single database and update its gauges in `shared.database_metrics`. -pub fn sample_database_metrics( - shared: &SharedState, - db_id: nodedb_types::DatabaseId, - db_name: &str, -) { - // 1. connections: active connections holding an admission permit for this database - // in `shared.admission_registry`. - let connections = shared - .admission_registry - .database_live_connections(db_id) - .map(u64::from) - .unwrap_or(0); - shared - .database_metrics - .set_connections(db_name, connections); - - // 2. memory: resident memory in bytes allocated by this database, read directly - // from the `database_budgets` map in `shared.governor`. - let memory_bytes = shared.governor.database_usage_bytes(db_id) as u64; - shared - .database_metrics - .set_memory_bytes(db_name, memory_bytes); - - // 3. storage: per-database storage usage in bytes from `shared.system_metrics` - // (the same source read by `SHOW DATABASE USAGE`). - let storage_bytes = shared - .system_metrics - .as_ref() - .map(|m| m.database_storage_bytes(db_name)) - .unwrap_or(0); - shared - .database_metrics - .set_storage_bytes(db_name, storage_bytes); - - // 4. bridge_queue_depth: sum of SPSC bridge virtual-queue depths for this database - // across all Data Plane cores in `shared.dispatcher`. - let bridge_queue_depth = match shared.dispatcher.lock() { - Ok(d) => d.virtual_queue_depth(db_id.as_u64()), - Err(p) => p.into_inner().virtual_queue_depth(db_id.as_u64()), - }; - shared - .database_metrics - .set_bridge_queue_depth(db_name, bridge_queue_depth); - - // 5. wal_latency_p99: P99 WAL group-commit fsync latency in microseconds from - // `shared.wal`'s commit latency histogram (falling back to `system_metrics`). - let wal_latency_p99 = { - let from_wal = shared.wal.commit_latency_p99_us(); - if from_wal > 0 { - from_wal - } else { - shared - .system_metrics - .as_ref() - .map(|s| s.wal_fsync_seconds.percentile(99.0)) - .unwrap_or(0) - } - }; - shared - .database_metrics - .set_wal_latency_p99(db_name, wal_latency_p99); -} - diff --git a/nodedb/src/bootstrap/database_metrics_sampler.rs b/nodedb/src/bootstrap/database_metrics_sampler.rs new file mode 100644 index 000000000..4b7530642 --- /dev/null +++ b/nodedb/src/bootstrap/database_metrics_sampler.rs @@ -0,0 +1,202 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! Background sampler for per-database Prometheus metrics. +//! +//! Samples live gauges every 10 seconds across memory, admission, storage, +//! virtual queue depths, and rolling-window WAL fsync latency. + +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; + +use tracing::info; + +use crate::control::metrics::histogram::HistogramSnapshot; +use crate::control::security::catalog::SystemCatalog; +use crate::control::state::SharedState; + +/// Spawn the 10-second periodic database metrics sampler. +pub fn spawn_database_metrics_sampler(shared: &Arc) { + let shared_sampler = Arc::clone(shared); + crate::control::shutdown::spawn_loop( + &shared.loop_registry, + &shared.shutdown, + "database_metrics_sampler", + crate::control::shutdown::ShutdownPhase::DrainingControlPlane, + move |mut shutdown| async move { + let mut tick = tokio::time::interval(Duration::from_secs(10)); + let mut last_fsync_snapshot: HistogramSnapshot = shared_sampler + .system_metrics + .as_ref() + .map(|s| s.wal_fsync_seconds.snapshot_counts()) + .unwrap_or_default(); + + loop { + tokio::select! { + _ = shutdown.wait_cancelled() => break, + _ = tick.tick() => {} + } + if shutdown.is_cancelled() { + break; + } + let catalog = shared_sampler.credentials.catalog(); + let databases = match catalog.list_databases() { + Ok(d) => d, + Err(e) => { + tracing::warn!(error = %e, "database_metrics_sampler: catalog list error"); + continue; + } + }; + + // Compute rolling P99 fsync latency across the 10s tick window. + let curr_fsync_snapshot = shared_sampler + .system_metrics + .as_ref() + .map(|s| s.wal_fsync_seconds.snapshot_counts()) + .unwrap_or_default(); + let wal_latency_p99 = + last_fsync_snapshot.delta_percentile(&curr_fsync_snapshot, 0.99); + last_fsync_snapshot = curr_fsync_snapshot; + + // Lock dispatcher ONCE per sampler tick to read virtual-queue depths for all databases. + let queue_depths: HashMap = { + let d = shared_sampler + .dispatcher + .lock() + .unwrap_or_else(|p| p.into_inner()); + databases + .iter() + .map(|db| (db.id.as_u64(), d.virtual_queue_depth(db.id.as_u64()))) + .collect() + }; + + for db in databases { + let depth = queue_depths.get(&db.id.as_u64()).copied().unwrap_or(0); + sample_database( + &shared_sampler, + catalog, + db.id, + &db.name, + depth, + wal_latency_p99, + ); + } + } + }, + ); + info!("database metrics sampler running"); +} + +/// Sample and record gauges for a single database (internal helper, not pub). +fn sample_database( + shared: &SharedState, + catalog: &SystemCatalog, + db_id: nodedb_types::DatabaseId, + db_name: &str, + bridge_queue_depth: u64, + wal_latency_p99: u64, +) { + // 1. connections: active connections holding an admission permit for this database + let connections = shared + .admission_registry + .database_live_connections(db_id) + .map(u64::from) + .unwrap_or(0); + shared + .database_metrics + .set_connections(db_name, connections); + + // 2. memory: resident memory in bytes allocated by this database in governor + let memory_bytes = shared.governor.database_usage_bytes(db_id) as u64; + shared + .database_metrics + .set_memory_bytes(db_name, memory_bytes); + + // 3. storage: real on-disk storage usage in bytes from system catalog tables + let storage_bytes = catalog.database_storage_bytes(db_id).unwrap_or(0); + shared + .database_metrics + .set_storage_bytes(db_name, storage_bytes); + + // 4. bridge_queue_depth: sum of SPSC bridge virtual-queue depths across cores + shared + .database_metrics + .set_bridge_queue_depth(db_name, bridge_queue_depth); + + // 5. wal_latency_p99: node-wide rolling 10s P99 group-commit fsync latency + shared + .database_metrics + .set_wal_latency_p99(db_name, wal_latency_p99); +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::bridge::dispatch::Dispatcher; + use crate::control::metrics::histogram::HistogramSnapshot; + use crate::control::state::SharedState; + use crate::wal::WalManager; + use nodedb_mem::engine::EngineId; + use nodedb_types::{DatabaseId, TenantId}; + + #[tokio::test] + async fn sampler_records_memory_connections_and_wal_fsync() { + let dir = tempfile::tempdir().expect("tempdir"); + let wal = Arc::new( + WalManager::open_for_testing(&dir.path().join("sampler-test.wal")).expect("open wal"), + ); + let (dispatcher, _data_sides) = Dispatcher::new(1, 64); + let shared = SharedState::new(dispatcher, wal).expect("construct state"); + + if let Some(metrics) = shared.system_metrics.as_ref() { + for _ in 0..90 { + metrics.record_wal_fsync(100); + } + for _ in 0..10 { + metrics.record_wal_fsync(500); + } + } + + let db_id = DatabaseId::new(42); + let db_name = "test_analytics"; + + // Reserve memory in governor for db_id. + let _token = shared + .governor + .try_reserve(db_id, TenantId::new(1), EngineId::Vector, 8192) + .expect("reserve memory"); + + // Acquire connection permit for db_id. + let _permit = shared + .admission_registry + .try_acquire_database(db_id) + .expect("acquire permit"); + + let catalog = shared.credentials.catalog(); + let curr_snapshot = shared + .system_metrics + .as_ref() + .map(|s| s.wal_fsync_seconds.snapshot_counts()) + .unwrap_or_default(); + let delta_p99 = HistogramSnapshot::default().delta_percentile(&curr_snapshot, 0.99); + + // Run sample_database. + sample_database(&shared, catalog, db_id, db_name, 5, delta_p99); + + // Verify gauge outputs in DatabaseMetricsRegistry. + let mut text = String::new(); + shared.database_metrics.render_prometheus(&mut text); + assert!(text.contains(&format!( + "nodedb_database_connections{{database=\"{db_name}\"}} 1" + ))); + assert!(text.contains(&format!( + "nodedb_database_memory_used_bytes{{database=\"{db_name}\"}} 8192" + ))); + assert!(text.contains(&format!( + "nodedb_database_bridge_queue_depth{{database=\"{db_name}\"}} 5" + ))); + assert!(text.contains(&format!( + "nodedb_database_wal_commit_latency_p99_us{{database=\"{db_name}\"}}" + ))); + } +} diff --git a/nodedb/src/bootstrap/mod.rs b/nodedb/src/bootstrap/mod.rs index 8faa98793..2ba898ca1 100644 --- a/nodedb/src/bootstrap/mod.rs +++ b/nodedb/src/bootstrap/mod.rs @@ -8,6 +8,7 @@ pub mod core_stall_monitor; pub mod credentials; pub mod data_group_recovery; pub mod data_plane; +pub mod database_metrics_sampler; pub mod diagnostics; pub mod index_registry_seed; pub mod listeners; diff --git a/nodedb/src/bridge/dispatch/dispatcher.rs b/nodedb/src/bridge/dispatch/dispatcher.rs index f92779699..f9af4038c 100644 --- a/nodedb/src/bridge/dispatch/dispatcher.rs +++ b/nodedb/src/bridge/dispatch/dispatcher.rs @@ -344,7 +344,6 @@ impl Dispatcher { .sum() } - /// Poll responses from all Data Plane cores. /// /// A core whose channel has been observed dead contributes a synthesized @@ -847,8 +846,47 @@ mod tests { #[test] fn virtual_queue_depth_reporting() { - let (dispatcher, _data_sides) = Dispatcher::new(2, 8); + let (mut dispatcher, data_sides) = Dispatcher::new(2, 2); assert_eq!(dispatcher.virtual_queue_depth(1), 0); + + // Fill physical rings on core 0 and core 1 so subsequent requests park in WFQ. + dispatcher + .dispatch_to_core(0, make_request_for_db(0, 1, 1)) + .unwrap(); + dispatcher + .dispatch_to_core(0, make_request_for_db(0, 1, 2)) + .unwrap(); + assert_eq!(data_sides[0].request_rx.len(), 2); + + dispatcher + .dispatch_to_core(1, make_request_for_db(1, 1, 10)) + .unwrap(); + dispatcher + .dispatch_to_core(1, make_request_for_db(1, 1, 11)) + .unwrap(); + assert_eq!(data_sides[1].request_rx.len(), 2); + + // Park 2 requests in core 0's WFQ for db 1. + dispatcher + .dispatch_to_core(0, make_request_for_db(0, 1, 3)) + .unwrap(); + dispatcher + .dispatch_to_core(0, make_request_for_db(0, 1, 4)) + .unwrap(); + + // Park 3 requests in core 1's WFQ for db 1. + dispatcher + .dispatch_to_core(1, make_request_for_db(1, 1, 12)) + .unwrap(); + dispatcher + .dispatch_to_core(1, make_request_for_db(1, 1, 13)) + .unwrap(); + dispatcher + .dispatch_to_core(1, make_request_for_db(1, 1, 14)) + .unwrap(); + + // Virtual queue depth must sum across all cores. + assert_eq!(dispatcher.virtual_queue_depth(1), 5); + assert_eq!(dispatcher.virtual_queue_depth(2), 0); } } - diff --git a/nodedb/src/control/catalog_entry/post_apply/database.rs b/nodedb/src/control/catalog_entry/post_apply/database.rs index fb9bff19f..45eb35714 100644 --- a/nodedb/src/control/catalog_entry/post_apply/database.rs +++ b/nodedb/src/control/catalog_entry/post_apply/database.rs @@ -21,12 +21,16 @@ pub fn put(descriptor: DatabaseDescriptor, shared: Arc) { } /// Post-apply for `DeleteDatabase` — release the quota caps of the dropped -/// scope, its tenants' caps included, and unregister from maintenance budget. +/// scope, its tenants' caps included, unregister from maintenance budget, +/// and remove any metrics handles from `DatabaseMetricsRegistry`. pub fn delete(db_id: u64, shared: Arc) { - shared - .maintenance_budget - .remove_database(DatabaseId::new(db_id)); - super::quota::release_database_scope(DatabaseId::new(db_id), &shared); + let db = DatabaseId::new(db_id); + if let Ok(Some(name)) = shared.credentials.catalog().get_database_name_by_id(db) { + shared.database_metrics.remove(&name); + } + shared.database_metrics.remove(&format!("db-{}", db_id)); + shared.maintenance_budget.remove_database(db); + super::quota::release_database_scope(db, &shared); } /// Post-apply for `PutDatabaseGrant`. diff --git a/nodedb/src/control/maintenance/budget.rs b/nodedb/src/control/maintenance/budget.rs index b42d9a3b0..1d5a92ec9 100644 --- a/nodedb/src/control/maintenance/budget.rs +++ b/nodedb/src/control/maintenance/budget.rs @@ -446,4 +446,3 @@ mod tests { assert!(out.contains("nodedb_database_maintenance_cpu_us_total")); } } - diff --git a/nodedb/src/control/metrics/database.rs b/nodedb/src/control/metrics/database.rs index d70790cbc..0e5c13ed6 100644 --- a/nodedb/src/control/metrics/database.rs +++ b/nodedb/src/control/metrics/database.rs @@ -27,10 +27,12 @@ pub struct DatabaseCounters { /// Current storage usage in bytes. pub storage_bytes: AtomicU64, /// Current active connection count. + /// Current active connection count. pub connections: AtomicU64, /// SPSC bridge virtual-queue depth snapshot. pub bridge_queue_depth: AtomicU64, - /// WAL commit latency P99 in microseconds (updated by WAL group-commit path). + /// WAL commit latency P99 in microseconds. Node-wide metric sampled from the node's + /// single WalManager; every database reports the same node-wide value. pub wal_commit_latency_p99_us: AtomicU64, /// Cumulative maintenance CPU-seconds consumed by this database. pub maintenance_cpu_seconds_total: AtomicU64, @@ -77,6 +79,12 @@ impl DatabaseMetricsRegistry { .clone() } + /// Remove all metrics for `db_name` from the registry (e.g. on DROP DATABASE). + pub fn remove(&self, db_name: &str) { + let mut w = self.counters.write().unwrap_or_else(|p| p.into_inner()); + w.remove(db_name); + } + /// Increment the QPS counter for `db_name` by 1. pub fn record_qps(&self, db_name: &str) { self.get_or_create(db_name) @@ -113,6 +121,9 @@ impl DatabaseMetricsRegistry { } /// Set the WAL commit latency P99 (microseconds) for `db_name`. + /// + /// Note: WAL latency is a node-wide metric (one WalManager per node); all databases + /// report the same node-wide latency. pub fn set_wal_latency_p99(&self, db_name: &str, us: u64) { self.get_or_create(db_name) .wal_commit_latency_p99_us @@ -139,7 +150,6 @@ impl DatabaseMetricsRegistry { /// Called from [`crate::control::maintenance::MaintenanceLease::drop`] /// when a maintenance lease finishes and records elapsed wall-clock CPU time. pub fn add_maintenance_cpu_secs(&self, db_name: &str, secs: f64) { - let us = if secs.is_finite() && secs > 0.0 { (secs * 1_000_000.0).round() as u64 } else { @@ -230,7 +240,7 @@ impl DatabaseMetricsRegistry { ); emit_gauge!( "nodedb_database_wal_commit_latency_p99_us", - "WAL commit latency P99 in microseconds per database", + "WAL commit latency P99 in microseconds (node-wide; shared across all databases)", wal_commit_latency_p99_us ); emit_counter!( @@ -349,12 +359,42 @@ mod tests { } #[test] - fn add_maintenance_cpu_secs_accumulates() { + fn remove_clears_database_metrics() { let reg = DatabaseMetricsRegistry::new(); - reg.add_maintenance_cpu_secs("db1", 1.5); - reg.add_maintenance_cpu_secs("db1", 0.5); - let c = reg.get_or_create("db1"); - assert_eq!(c.maintenance_cpu_seconds_total.load(Ordering::Relaxed), 2_000_000); + reg.set_connections("db_to_drop", 10); + assert_eq!( + reg.get_or_create("db_to_drop") + .connections + .load(Ordering::Relaxed), + 10 + ); + reg.remove("db_to_drop"); + let mut out = String::new(); + reg.render_prometheus(&mut out); + assert!(!out.contains(r#"database="db_to_drop""#)); } -} + #[test] + fn add_maintenance_cpu_secs_accumulates() { + let reg = Arc::new(DatabaseMetricsRegistry::new()); + let tracker = Arc::new( + crate::control::maintenance::MaintenanceBudgetTracker::with_metrics(Arc::clone(®)), + ); + let db = DatabaseId::new(42); + tracker.set_database_name(db, "analytics"); + tracker.set_cap(db, 100); + + { + let _lease = tracker.try_acquire(db, 1.0).expect("acquire lease"); + std::thread::sleep(std::time::Duration::from_millis(15)); + // Lease dropped here, calling add_maintenance_cpu_secs via Drop + } + + let c = reg.get_or_create("analytics"); + let recorded = c.maintenance_cpu_seconds_total.load(Ordering::Relaxed); + assert!( + recorded > 0, + "MaintenanceLease drop must increment maintenance_cpu_seconds_total, got {recorded}" + ); + } +} diff --git a/nodedb/src/control/metrics/histogram.rs b/nodedb/src/control/metrics/histogram.rs index 48a228efa..3cff8a517 100644 --- a/nodedb/src/control/metrics/histogram.rs +++ b/nodedb/src/control/metrics/histogram.rs @@ -168,6 +168,19 @@ impl AtomicHistogram { snap } + /// Create a point-in-time snapshot of bucket counts for rolling delta calculations. + pub fn snapshot_counts(&self) -> HistogramSnapshot { + HistogramSnapshot { + boundaries: self.boundaries, + buckets: self + .buckets + .iter() + .map(|b| b.load(Ordering::Relaxed)) + .collect(), + count: self.count.load(Ordering::Relaxed), + } + } + /// Merge another histogram's counts into this one. /// /// Both histograms must share the same bucket boundaries — if they do @@ -186,6 +199,71 @@ impl AtomicHistogram { } } +/// Point-in-time snapshot of bucket counts for rolling-window percentile calculations. +#[derive(Debug, Clone, Default)] +pub struct HistogramSnapshot { + /// Upper bounds in microseconds. + pub boundaries: &'static [u64], + /// Bucket counters at snapshot time. + pub buckets: Vec, + /// Total observations at snapshot time. + pub count: u64, +} + +impl HistogramSnapshot { + /// Compute percentile from the bucket count delta between `self` (earlier) and `newer`. + /// + /// `p` is a fraction between `0.0` and `1.0` (e.g. `0.99` for P99). + /// Returns `0` if no new observations occurred in the interval. + pub fn delta_percentile(&self, newer: &HistogramSnapshot, p: f64) -> u64 { + let is_empty_baseline = self.boundaries.is_empty() && self.count == 0; + if !is_empty_baseline && (self.boundaries != newer.boundaries || newer.count <= self.count) + { + return 0; + } + if newer.count == 0 { + return 0; + } + let delta_count = if is_empty_baseline { + newer.count + } else { + newer.count - self.count + }; + if delta_count == 0 { + return 0; + } + let target = (p * delta_count as f64) as u64; + let mut cumulative = 0u64; + let mut prev_boundary = 0u64; + + for (i, &boundary) in newer.boundaries.iter().enumerate() { + let older_bucket = if is_empty_baseline { + 0 + } else { + self.buckets.get(i).copied().unwrap_or(0) + }; + let newer_bucket = newer.buckets.get(i).copied().unwrap_or(0); + let bucket_delta = newer_bucket.saturating_sub(older_bucket); + cumulative += bucket_delta; + if cumulative >= target { + let bucket_start = prev_boundary; + let bucket_width = boundary - bucket_start; + if bucket_delta == 0 { + return boundary; + } + let fraction = if cumulative > target { + (bucket_delta - (cumulative - target)) as f64 / bucket_delta as f64 + } else { + 1.0 + }; + return bucket_start + (fraction * bucket_width as f64) as u64; + } + prev_boundary = boundary; + } + newer.boundaries.last().copied().unwrap_or(0) + } +} + impl Default for AtomicHistogram { fn default() -> Self { Self::new() @@ -229,6 +307,69 @@ mod tests { assert!((500..=1000).contains(&p50), "p50={p50}"); } + #[test] + fn percentile_p99_numerically_correct() { + // WAL buckets: [100, 500, 1000, 5000, 10000, 50000, 100000, 500000, 1000000] + let h = AtomicHistogram::with_buckets(WAL_FSYNC_BUCKETS_US); + // Observe 90 items at 80us (falls in <=100us bucket) + for _ in 0..90 { + h.observe(80); + } + // Observe 9 items at 400us (falls in <=500us bucket) + for _ in 0..9 { + h.observe(400); + } + // Observe 1 item at 800us (falls in <=1000us bucket) + h.observe(800); + + assert_eq!(h.count(), 100); + + // p50 is rank 50 (within first bucket 0..100us) + let p50 = h.percentile(0.50); + assert!(p50 <= 100, "expected p50 <= 100, got {p50}"); + + // p99 is rank 99 (90 in bucket0 + 9 in bucket1 = 99 -> top of <=500us bucket) + let p99 = h.percentile(0.99); + assert!( + (100..=500).contains(&p99), + "expected p99 in [100, 500], got {p99}" + ); + // Passing 99.0 would have returned 1_000_000 (the last boundary). + assert_ne!(p99, 1_000_000); + } + + #[test] + fn rolling_window_delta_percentile() { + let h = AtomicHistogram::with_buckets(WAL_FSYNC_BUCKETS_US); + + // Window 1: 100 observations at 80us (<=100us) + for _ in 0..100 { + h.observe(80); + } + let snap1 = h.snapshot_counts(); + + // Window 2: 90 observations at 80us, 10 observations at 4000us (<=5000us) + for _ in 0..90 { + h.observe(80); + } + for _ in 0..10 { + h.observe(4000); + } + let snap2 = h.snapshot_counts(); + + // Delta p99 in Window 2 alone (100 new observations: 90 at 80us, 10 at 4000us) + let delta_p99 = snap1.delta_percentile(&snap2, 0.99); + assert!( + (1000..=5000).contains(&delta_p99), + "expected delta p99 in [1000, 5000], got {delta_p99}" + ); + + // Window 3: No new observations + let snap3 = h.snapshot_counts(); + let delta_p99_empty = snap2.delta_percentile(&snap3, 0.99); + assert_eq!(delta_p99_empty, 0); + } + #[test] fn prometheus_output() { let h = AtomicHistogram::new(); diff --git a/nodedb/src/control/security/catalog/database.rs b/nodedb/src/control/security/catalog/database.rs index e4f15996c..6343843cd 100644 --- a/nodedb/src/control/security/catalog/database.rs +++ b/nodedb/src/control/security/catalog/database.rs @@ -223,6 +223,256 @@ impl SystemCatalog { .map_err(|e| catalog_err("list_databases read txn", e))?; list_databases_in(&txn) } + + /// Calculate the real storage size in bytes used by a database across all catalog tables. + /// + /// Sums the persisted key and value byte lengths for all entries belonging to `db_id` + /// across collections, surrogate indexes, topics, streams, materialized views, + /// policies, procedures, triggers, functions, and database descriptors. + pub fn database_storage_bytes(&self, db_id: DatabaseId) -> crate::Result { + let read_txn = self + .db + .begin_read() + .map_err(|e| catalog_err("read txn for database_storage_bytes", e))?; + + let id = db_id.as_u64(); + let mut total_bytes: u64 = 0; + + // 1. Collections: (database_id: u64, "{tenant_id}:{name}") -> msgpack + if let Ok(table) = read_txn.open_table(super::types::COLLECTIONS) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let (db, name) = k.value(); + if db == id { + total_bytes = total_bytes + .saturating_add(8 + name.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 2. Surrogate PKs: (database_id: u64, tenant_id: u64, collection, encoded_pk) -> u32 + if let Ok(table) = read_txn.open_table(super::types::SURROGATE_PK_V3) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, _v)) = entry { + let (db, _tid, coll, pk) = k.value(); + if db == id { + total_bytes = total_bytes + .saturating_add(16 + coll.len() as u64 + pk.len() as u64 + 4); + } + } + } + } + } + + // 3. Surrogate PK Rev: (database_id: u64, tenant_id: u64, collection, surrogate: u32) -> encoded pk + if let Ok(table) = read_txn.open_table(super::types::SURROGATE_PK_REV_V3) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let (db, _tid, coll, _surr) = k.value(); + if db == id { + total_bytes = total_bytes + .saturating_add(20 + coll.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 4. Topic messages: [database_id: be u64]... + if let Ok(table) = read_txn.open_table(super::types::TOPIC_MESSAGES) { + let be_bytes = id.to_be_bytes(); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_bytes = k.value(); + if k_bytes.starts_with(&be_bytes) { + total_bytes = total_bytes + .saturating_add(k_bytes.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 5. Ep Topics: "v2/{database_id}/{tenant_id}..." + if let Ok(table) = read_txn.open_table(super::types::TOPICS_EP) { + let prefix = format!("v2/{id}/"); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.starts_with(&prefix) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 6. Change streams: "v2/{database_id}/..." + if let Ok(table) = read_txn.open_table(super::types::CHANGE_STREAMS) { + let prefix = format!("v2/{id}/"); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.starts_with(&prefix) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 7. Consumer groups: "v2:{database_id}:..." + if let Ok(table) = read_txn.open_table(super::types::CONSUMER_GROUPS) { + let prefix = format!("v2:{id}:"); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.starts_with(&prefix) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 8. Streaming MVs: "v2:{database_id}:..." + if let Ok(table) = read_txn.open_table(super::types::STREAMING_MVS) { + let prefix = format!("v2:{id}:"); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.starts_with(&prefix) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 9. Continuous Aggregates: (database_id: u64, "{tenant_id}:{name}") + if let Ok(table) = read_txn.open_table(super::types::CONTINUOUS_AGGREGATES) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let (db, name) = k.value(); + if db == id { + total_bytes = total_bytes + .saturating_add(8 + name.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 10. Retention policies: (database_id: u64, "{tenant_id}:{policy_name}") + if let Ok(table) = read_txn.open_table(super::types::RETENTION_POLICIES) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let (db, name) = k.value(); + if db == id { + total_bytes = total_bytes + .saturating_add(8 + name.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 11. Alert rules: (database_id: u64, "{tenant_id}:{alert_name}") + if let Ok(table) = read_txn.open_table(super::types::ALERT_RULES) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let (db, name) = k.value(); + if db == id { + total_bytes = total_bytes + .saturating_add(8 + name.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 12. Synonym groups: "{database_id}:{tenant_id}:{group_name}" + if let Ok(table) = read_txn.open_table(super::types::SYNONYM_GROUPS) { + let prefix = format!("{id}:"); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.starts_with(&prefix) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 13. Column stats: "{database_id}:{tenant_id}:{collection}:{column}" + if let Ok(table) = read_txn.open_table(super::types::COLUMN_STATS) { + let prefix = format!("{id}:"); + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.starts_with(&prefix) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + + // 14. Database descriptor: `database_id (u64)` + if let Ok(table) = read_txn.open_table(DATABASES) { + if let Ok(Some(v)) = table.get(id) { + total_bytes = total_bytes.saturating_add(8 + v.value().len() as u64); + } + } + + // 15. Functions, triggers, procedures, arrays: "v2:{tenant_id}:{database_id}:{name}" or "\0v2:..." + let mid_pattern = format!(":{id}:"); + for table_def in [ + super::types::FUNCTIONS, + super::types::TRIGGERS, + super::types::PROCEDURES, + super::types::ARRAYS, + super::types::DEPENDENCIES, + ] { + if let Ok(table) = read_txn.open_table(table_def) { + if let Ok(range) = table.iter() { + for entry in range { + if let Ok((k, v)) = entry { + let k_str = k.value(); + if k_str.contains(&mid_pattern) { + total_bytes = total_bytes + .saturating_add(k_str.len() as u64 + v.value().len() as u64); + } + } + } + } + } + } + + Ok(total_bytes) + } } /// Body of [`SystemCatalog::list_databases`], over an already-open read @@ -349,4 +599,28 @@ mod tests { // second delete is a no-op cat.delete_database(DatabaseId::DEFAULT).unwrap(); } + + #[test] + fn database_storage_bytes_reflects_actual_usage() { + let (_dir, cat) = open_catalog(); + cat.bootstrap_default_database().unwrap(); + let initial_bytes = cat.database_storage_bytes(DatabaseId::DEFAULT).unwrap(); + assert!( + initial_bytes > 0, + "bootstrapped default database descriptor must contribute non-zero bytes" + ); + + let mut coll = super::super::collection::StoredCollection::new(1, "users", "admin"); + coll.database_id = DatabaseId::DEFAULT; + cat.put_collection(DatabaseId::DEFAULT, &coll).unwrap(); + + let after_coll_bytes = cat.database_storage_bytes(DatabaseId::DEFAULT).unwrap(); + assert!( + after_coll_bytes > initial_bytes, + "adding collection must increase database_storage_bytes: {after_coll_bytes} > {initial_bytes}" + ); + + // Unrelated database reports 0 bytes + assert_eq!(cat.database_storage_bytes(DatabaseId::new(999)).unwrap(), 0); + } } diff --git a/nodedb/src/control/server/admission/registry.rs b/nodedb/src/control/server/admission/registry.rs index 7cf365d11..c13551750 100644 --- a/nodedb/src/control/server/admission/registry.rs +++ b/nodedb/src/control/server/admission/registry.rs @@ -280,6 +280,17 @@ impl AdmissionRegistry { &self, db: DatabaseId, ) -> Result, AdmissionError> { + let has_entry = { + let map = self.db_semaphores.read().unwrap_or_else(|p| p.into_inner()); + map.contains_key(&db) + }; + if !has_entry { + let mut map = self + .db_semaphores + .write() + .unwrap_or_else(|p| p.into_inner()); + map.entry(db).or_insert_with(|| LimitEntry::new(None)); + } let map = self.db_semaphores.read().unwrap_or_else(|p| p.into_inner()); let permit = try_acquire_entry(map.get(&db), |limit| { AdmissionError::DatabaseCapExhausted { db, limit } @@ -372,10 +383,12 @@ mod tests { #[test] fn no_database_cap_allows_unlimited() { let reg = AdmissionRegistry::new(); - // No entry configured → Ok(None). - let r = reg.try_acquire_database(db(0)); - assert!(r.unwrap().is_none()); - assert_eq!(reg.database_live_connections(db(0)), None); + // An uncapped database admits and tracks live connections. + let p = reg.try_acquire_database(db(0)).unwrap(); + assert!(p.is_some()); + assert_eq!(reg.database_live_connections(db(0)), Some(1)); + drop(p); + assert_eq!(reg.database_live_connections(db(0)), Some(0)); } #[test] @@ -563,7 +576,7 @@ mod tests { reg.set_database_limit(db(5), 0); assert_eq!(reg.database_live_connections(db(5)), None); - assert!(reg.try_acquire_database(db(5)).unwrap().is_none()); + assert!(reg.try_acquire_database(db(5)).unwrap().is_some()); } #[test] diff --git a/nodedb/src/control/server/shared/ddl/neutral/observability.rs b/nodedb/src/control/server/shared/ddl/neutral/observability.rs index 94cfa2a4f..c7fa02d3d 100644 --- a/nodedb/src/control/server/shared/ddl/neutral/observability.rs +++ b/nodedb/src/control/server/shared/ddl/neutral/observability.rs @@ -171,19 +171,19 @@ pub fn show_metrics( if let Some(sys) = state.system_metrics.as_ref() { rows.push(( "wal_fsync_p50_us".into(), - sys.wal_fsync_seconds.percentile(50.0).to_string(), + sys.wal_fsync_seconds.percentile(0.5).to_string(), )); rows.push(( "wal_fsync_p99_us".into(), - sys.wal_fsync_seconds.percentile(99.0).to_string(), + sys.wal_fsync_seconds.percentile(0.99).to_string(), )); rows.push(( "query_latency_p50_us".into(), - sys.query_latency.percentile(50.0).to_string(), + sys.query_latency.percentile(0.5).to_string(), )); rows.push(( "query_latency_p99_us".into(), - sys.query_latency.percentile(99.0).to_string(), + sys.query_latency.percentile(0.99).to_string(), )); } diff --git a/nodedb/src/control/state/init.rs b/nodedb/src/control/state/init.rs index c108e7261..c3105469b 100644 --- a/nodedb/src/control/state/init.rs +++ b/nodedb/src/control/state/init.rs @@ -521,6 +521,9 @@ impl SharedState { startup: Arc::clone(&startup_gate), }); Self::wire_session_handle_audit(&state); + if let Some(metrics) = &state.system_metrics { + state.wal.set_metrics(Arc::clone(metrics)); + } state .maintenance_budget .set_metrics(Arc::clone(&state.database_metrics)); diff --git a/nodedb/src/control/state/init_prod/post_init.rs b/nodedb/src/control/state/init_prod/post_init.rs index 4eb5e6967..b788f27f6 100644 --- a/nodedb/src/control/state/init_prod/post_init.rs +++ b/nodedb/src/control/state/init_prod/post_init.rs @@ -33,6 +33,9 @@ pub(super) fn hydrate_caches(state: &Arc) { "boot: failed to populate idle_timeout_cache from catalog" ); } + if let Some(metrics) = &state.system_metrics { + state.wal.set_metrics(Arc::clone(metrics)); + } state .maintenance_budget .set_metrics(Arc::clone(&state.database_metrics)); diff --git a/nodedb/src/wal/manager/core.rs b/nodedb/src/wal/manager/core.rs index 5a6243154..c410842fa 100644 --- a/nodedb/src/wal/manager/core.rs +++ b/nodedb/src/wal/manager/core.rs @@ -45,8 +45,8 @@ pub struct WalManager { /// Wakes `wait_durable` followers when `durable_lsn` advances (or a leader's /// fsync fails, so they re-attempt and observe the same error). pub(super) durable_notify: tokio::sync::Notify, - /// WAL group-commit fsync latency distribution. - pub(super) commit_latency: crate::control::metrics::AtomicHistogram, + /// System metrics registry handle for recording group-commit fsync latency. + pub(super) metrics: std::sync::RwLock>>, } impl WalManager { @@ -65,11 +65,19 @@ impl WalManager { self.durable_lsn.load(std::sync::atomic::Ordering::Acquire) } - /// WAL group-commit fsync latency P99 in microseconds. - pub fn commit_latency_p99_us(&self) -> u64 { - self.commit_latency.percentile(99.0) + /// Wire or replace the system metrics handle used to record group-commit fsync latency. + pub fn set_metrics(&self, metrics: Arc) { + let mut w = self.metrics.write().unwrap_or_else(|p| p.into_inner()); + *w = Some(metrics); } + /// Read the optional system metrics handle. + pub fn metrics(&self) -> Option> { + self.metrics + .read() + .unwrap_or_else(|p| p.into_inner()) + .clone() + } /// Return the stable in-memory root for per-user CRDT signing keys. /// The root is persisted only as WAL-key-wrapped ciphertext and is @@ -173,20 +181,17 @@ impl WalManager { Ok(Self { wal: Arc::new(Mutex::new(wal)), - wal_dir, + wal_dir: wal_dir.to_path_buf(), encryption_ring: None, crdt_signing_root: None, audit_wal, durable_lsn: AtomicU64::new(0), commit_lock: tokio::sync::Mutex::new(()), durable_notify: tokio::sync::Notify::new(), - commit_latency: crate::control::metrics::AtomicHistogram::with_buckets( - crate::control::metrics::histogram::WAL_FSYNC_BUCKETS_US, - ), + metrics: std::sync::RwLock::new(None), }) } - /// Open without `O_DIRECT`. /// /// The in-process test harnesses put their data directories in tempdirs, diff --git a/nodedb/src/wal/manager/durable_commit.rs b/nodedb/src/wal/manager/durable_commit.rs index a5e8c6a25..52d916079 100644 --- a/nodedb/src/wal/manager/durable_commit.rs +++ b/nodedb/src/wal/manager/durable_commit.rs @@ -63,21 +63,24 @@ impl WalManager { // advances to exactly what the fsync made durable, never past // it. let wal = std::sync::Arc::clone(&self.wal); - let join = tokio::task::spawn_blocking(move || -> crate::Result<(u64, u64)> { - let start = std::time::Instant::now(); + let metrics = self.metrics(); + let join = tokio::task::spawn_blocking(move || -> crate::Result { let mut guard = wal.lock().unwrap_or_else(|p| p.into_inner()); + let start = std::time::Instant::now(); guard.sync().map_err(crate::Error::Wal)?; let elapsed_us = start.elapsed().as_micros() as u64; + if let Some(m) = &metrics { + m.record_wal_fsync(elapsed_us); + } // `next_lsn()` is the next LSN to assign; the highest LSN // this sync made durable is one below it. - Ok((guard.next_lsn().saturating_sub(1), elapsed_us)) + Ok(guard.next_lsn().saturating_sub(1)) }) .await; let outcome = match join { - Ok(Ok((durable_through, elapsed_us))) => { + Ok(Ok(durable_through)) => { self.durable_lsn.fetch_max(durable_through, AcqRel); - self.commit_latency.observe(elapsed_us); Ok(()) } // fsync error: do NOT advance `durable_lsn`. @@ -89,7 +92,6 @@ impl WalManager { }), }; - // Wake followers on every leader-exit path so a failed or // panicked fsync never strands them on `notified.await`. On // success `durable_lsn` was advanced first, so they observe @@ -116,6 +118,8 @@ impl WalManager { mod tests { use super::*; use crate::types::{DatabaseId, TenantId, VShardId}; + use std::sync::Arc; + use std::sync::atomic::Ordering; fn open_wal(dir: &std::path::Path) -> WalManager { WalManager::open_for_testing(&dir.join("test.wal")).expect("open wal") @@ -141,6 +145,8 @@ mod tests { async fn wait_durable_fast_path_when_already_durable() { let dir = tempfile::tempdir().expect("tempdir"); let wal = open_wal(dir.path()); + let metrics = Arc::new(crate::control::metrics::SystemMetrics::new()); + wal.set_metrics(Arc::clone(&metrics)); let lsn = wal .append_put( TenantId::new(1), @@ -150,12 +156,13 @@ mod tests { ) .expect("append"); wal.wait_durable(lsn).await.expect("first"); - assert!(wal.commit_latency.count() >= 1); + assert_eq!(metrics.wal_fsync_count.load(Ordering::Relaxed), 1); + assert_eq!(metrics.wal_fsync_seconds.count(), 1); // Second call is the fast path — no further fsync required. wal.wait_durable(lsn).await.expect("second"); + assert_eq!(metrics.wal_fsync_count.load(Ordering::Relaxed), 1); } - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn concurrent_waiters_coalesce() { let dir = tempfile::tempdir().expect("tempdir"); diff --git a/nodedb/tests/inproc/cases/quota_drop_cleanup.rs b/nodedb/tests/inproc/cases/quota_drop_cleanup.rs index 63a7b7d57..f834b3e61 100644 --- a/nodedb/tests/inproc/cases/quota_drop_cleanup.rs +++ b/nodedb/tests/inproc/cases/quota_drop_cleanup.rs @@ -245,11 +245,11 @@ async fn dropped_database_id_does_not_inherit_cap() { None, "the dead cap must not survive for a database reusing this id" ); + let permit = registry + .try_acquire_database(db) + .expect("an uncapped database admits"); assert!( - registry - .try_acquire_database(db) - .expect("an uncapped database admits") - .is_none(), - "an uncapped database hands out no permit" + permit.is_some(), + "an uncapped database hands out an unconstrained permit" ); } diff --git a/nodedb/tests/inproc/cases/quota_live_enforcement_apply.rs b/nodedb/tests/inproc/cases/quota_live_enforcement_apply.rs index 49bcfff32..29f3a3eb2 100644 --- a/nodedb/tests/inproc/cases/quota_live_enforcement_apply.rs +++ b/nodedb/tests/inproc/cases/quota_live_enforcement_apply.rs @@ -90,12 +90,12 @@ async fn alter_database_quota_zero_clears_connection_cap() { None, "zero drops the entry, so the database is uncapped again" ); + let permit = registry + .try_acquire_database(db) + .expect("an uncapped database admits"); assert!( - registry - .try_acquire_database(db) - .expect("an uncapped database admits") - .is_none(), - "an uncapped database hands out no permit" + permit.is_some(), + "an uncapped database hands out an unconstrained permit" ); }