Skip to content

fix(obsv): wire live per-database metrics sampler (fixes #386) - #386

Open
AubaidAhmedSaiyed wants to merge 4 commits into
NodeDB-Lab:mainfrom
AubaidAhmedSaiyed:fix-sampler
Open

AubaidAhmedSaiyed wants to merge 4 commits into
NodeDB-Lab:mainfrom
AubaidAhmedSaiyed:fix-sampler

Conversation

@AubaidAhmedSaiyed

@AubaidAhmedSaiyed AubaidAhmedSaiyed commented Sep 27, 2026 •

Copy link
Copy Markdown

fix(obsv): wire live per-database metrics sampler (fixes #375)


Fixes #375

What this does and why

Six per-database Prometheus gauges (defined in DatabaseMetricsRegistry, control/metrics/database.rs) had setter methods but no callers — nothing outside the registry ever invoked them, so they stayed frozen at their initial values, silently misleading anyone watching /metrics or building alerts/dashboards on them. In particular, bridge_queue_depth staying at zero meant bulk loaders had no real signal to back off before the write-fair-queue started rejecting them.

This PR adds a periodic sampler, database_metrics_sampler, spawned in background_loops.rs on a 10-second interval — following the same pattern already used by the existing mirror_lag_monitor loop. Each tick it lists active databases from the catalog and pushes real current values into five of the six gauges:

Gauge Source
set_connections AdmissionRegistry::database_live_connections
set_memory_bytes MemoryGovernor::database_usage_bytes (new accessor)
set_storage_bytes SystemMetrics::database_storage_bytes
set_bridge_queue_depth Dispatcher::virtual_queue_depth (new accessor)
set_wal_latency_p99 WalManager::commit_latency_p99_us (new histogram)

Three new accessor/measurement points were added where no existing one exposed per-database data:

  • MemoryGovernor::database_usage_bytes (nodedb-mem/src/governor/metrics.rs)
  • Dispatcher::virtual_queue_depth (nodedb/src/bridge/dispatch/dispatcher.rs)
  • A commit_latency AtomicHistogram in WalManager, sampled during the group-commit fsync in wait_durable (nodedb/src/wal/manager/durable_commit.rs)

The sixth setter, add_maintenance_cpu_secs, is intentionally left unfilled — see Tradeoffs below.

How to test

1. Unit tests — includes new coverage for each accessor (e.g. governor::metrics::tests::database_usage_bytes_tracks_allocation, a dispatcher test for virtual_queue_depth):

cargo test -p nodedb -p nodedb-mem -p nodedb-bridge

2. Caller-verification script (from the issue) — confirms all five gauges now have a real caller outside the registry file:

python3 -c "
import pathlib, re, subprocess
src = pathlib.Path('nodedb/src/control/metrics/database.rs').read_text()
for m in re.findall(r'pub fn (\w+)', src):
    out = subprocess.run(['git','grep','-rn',f'.{m}(','nodedb/src'],
                         capture_output=True, text=True).stdout
    sites = [l for l in out.splitlines() if 'control/metrics/database.rs' not in l]
    print(f'{m:32} callers_outside={len(sites)}')
"

set_memory_bytes, set_storage_bytes, set_connections, set_bridge_queue_depth, and set_wal_latency_p99 all now report callers_outside=1.

3. Manual verification — run the server locally, create a database, generate some load, and curl /metrics twice ~10s apart. The five gauges should show non-zero, changing values instead of remaining at 0.

Tradeoffs / alternatives considered

  • Pull-based sampler vs. push-on-event. I used a 10-second poll (per the issue's suggestion) rather than having each subsystem push its own metric update on every state change, the way set_mirror_lag_ms is called directly from the mirror observer today. A push model would be fresher and cheaper per-update, but touches more call sites across unrelated subsystems for a first pass. Open to converting bridge_queue_depth specifically to push-based in a follow-up, since that's the one that matters most for real-time backpressure decisions.

  • WAL commit latency measured in wait_durable, not at the raw fsync syscall. This times the full group-commit wait from the caller's perspective, consistent with how wal_fsync_latency_us is already measured elsewhere in the codebase, rather than isolating just the fsync duration.

Performance

This change touches the WAL group-commit path (wait_durable) to record one histogram observation per commit, and the bridge dispatcher to sum queue depths once per sampler tick (every 10 seconds, not per-request).

I have not run a formal before/after benchmark for this PR. The WAL-path addition is a single atomic histogram observe() call per group commit, which should sit well under the noise floor of an fsync itself, but I haven't measured it directly and don't want to state a number I don't have. Happy to run a benchmark against whichever suite maintainers consider authoritative before merge.

Files changed

  • nodedb-mem/src/governor/metrics.rs — added database_usage_bytes + unit test
  • nodedb/src/bridge/dispatch/dispatcher.rs — added virtual_queue_depth + unit test
  • nodedb/src/wal/manager/core.rs — added commit_latency histogram + commit_latency_p99_us
  • nodedb/src/wal/manager/durable_commit.rs — timed group-commit fsyncs into commit_latency
  • nodedb/src/bootstrap/background_loops.rs — added database_metrics_sampler loop
  • nodedb/src/control/metrics/database.rs — documented add_maintenance_cpu_secs, added unit tests for the five gauge setters

@AubaidAhmedSaiyed AubaidAhmedSaiyed changed the title obsv: wire live per-database metrics sampler (fixes #375) fix(obsv): wire live per-database metrics sampler (fixes #375) Sep 27, 2026
…dration and DDL post-apply hooks) and forwards elapsed CPU time to DatabaseMetricsRegistry on drop, completing the sixth metric from NodeDB-Lab#375.

@farhan-syah farhan-syah left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for picking up #375. The issue is real: none of the six per-database setters has a caller today.

This PR does not close it yet. Four of the five wired gauges still export a constant or near-constant value. That is the bug the issue reports.

What landed correctly

  • Dispatcher::virtual_queue_depth sums wfq.depth_for across cores. It uses the same database_id.as_u64() key as the dispatch path. The bridge queue depth gauge is real.
  • The sampler spawns through spawn_loop with a shutdown phase. This matches the mirror_lag_monitor pattern.

Blockers

# Defect Inline
1 percentile(99.0) pins WAL p99 at 1_000_000 µs after the first fsync wal/manager/core.rs:70
2 Storage source is a stub that only ever holds 0 background_loops.rs:494
3 Memory reads 0 for every database without a quota governor/metrics.rs:18
4 Connections read 0 for every database without a connection cap background_loops.rs:475
5 A second WAL fsync histogram duplicates SystemMetrics::wal_fsync_seconds wal/manager/core.rs:49
6 rustfmt --check fails on all six files; all six pass on main body only
7 bootstrap/background_loops.rs grows to 528 lines, over the 500-line file limit background_loops.rs:466

On 6: run cargo fmt --all. The causes are doubled blank lines, a blank line after the add_maintenance_cpu_secs opening brace, and extra blank lines at end of file.

Claims that do not hold

  • "consistent with how wal_fsync_latency_us is already measured elsewhere." SystemMetrics::record_wal_fsync has zero callers. Nothing measures WAL fsync latency today. nodedb_wal_fsync_seconds is also dead.
  • "callers_outside=1" as proof of a fix. A caller exists for each setter. Items 1 to 4 show that most callers pass a constant.
  • cargo test. This repo runs tests with cargo nextest run. cargo test ignores the nextest test groups and can hang.

Same bug on main

control/server/shared/ddl/neutral/observability.rs:174-186 passes 50.0 and 99.0 to percentile for WAL fsync and query latency. Those outputs are wrong today. Fix them in this PR, since it is the same bug class.

Required end state

  • Every gauge reads a source that covers all databases, with no quota or cap needed.
    • The governor accounts memory per database without a budget.
    • The admission registry counts live connections without a cap.
    • Per-database storage accounting exists and feeds set_storage_bytes.
  • WAL latency uses one histogram. The group-commit path calls record_wal_fsync. The sampler reads p99 over the window since the previous tick.
  • add_maintenance_cpu_secs is either wired to maintenance tasks or no longer exported.
  • DROP DATABASE removes the database's series from DatabaseMetricsRegistry.
  • Tests prove each gauge moves: a queued-item depth test, a p99 value test, and a sampler test over a real reservation, connection, and fsync.

Comment thread nodedb/src/wal/manager/core.rs Outdated

/// WAL group-commit fsync latency P99 in microseconds.
pub fn commit_latency_p99_us(&self) -> u64 {
self.commit_latency.percentile(99.0)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. AtomicHistogram::percentile takes a fraction, not a percent. per_vshard.rs calls it with 0.99.

With 99.0, the target rank is 99× the count. The loop never reaches it and returns the last bucket boundary. After the first fsync, this gauge reads 1_000_000 µs forever.

Pass 0.99. Add a test that observes known values and asserts the p99.

Comment thread nodedb/src/wal/manager/core.rs Outdated
/// 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,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. This duplicates SystemMetrics::wal_fsync_seconds. That histogram already exists and is exported as nodedb_wal_fsync_seconds. Its writer record_wal_fsync has no callers.

Remove this field. Route the fsync duration into record_wal_fsync, so one histogram feeds both the node metric and the per-database gauge.

let wal = std::sync::Arc::clone(&self.wal);
let join = tokio::task::spawn_blocking(move || -> crate::Result<u64> {
let join = tokio::task::spawn_blocking(move || -> crate::Result<(u64, u64)> {
let start = std::time::Instant::now();

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The timer starts before wal.lock(). Appender lock contention counts as fsync time.

Start the timer after the lock is taken.

shared
.system_metrics
.as_ref()
.map(|s| s.wal_fsync_seconds.percentile(99.0))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This fallback reads wal_fsync_seconds, which nothing writes. It also has the same 99.0 bug.

Remove the fallback once the WAL path feeds record_wal_fsync.

The p99 is also cumulative since process start and never decays. A bulk loader needs recent latency. Compute p99 from the bucket-count delta since the previous tick.

There is one WAL per node, so every database gets the same value. Document the gauge as a node-wide value repeated per database.

let storage_bytes = shared
.system_metrics
.as_ref()
.map(|m| m.database_storage_bytes(db_name))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. SystemMetrics::database_storage_bytes is a stub. Its only writer is set_database_storage_bytes(name, 0) at CREATE DATABASE. This gauge copies a constant 0.

The issue asks for storage from the owning subsystem. Per-database storage accounting needs to exist first. Build it in this PR.


// 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() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This locks the dispatcher mutex once per database per tick. The dispatch hot path contends on the same lock.

Take the lock once per tick and read every database's depth under it.

break;
}
let catalog = shared_sampler.credentials.catalog();
let databases = match catalog.list_databases() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The sampler visits only live databases. DatabaseMetricsRegistry has no remove function.

After DROP DATABASE, the database's last values stay in /metrics forever. Add removal on the drop path.

Comment thread nodedb/src/control/metrics/database.rs Outdated
/// 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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A rustdoc note does not reach anyone reading /metrics. The series still exports a constant 0.

The note also describes future work ("pending a dedicated … source"). Comments describe the code as it is.

Either wire this counter to the maintenance tasks, or stop exporting it.

}

#[test]
fn virtual_queue_depth_reporting() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test only asserts 0 on an empty dispatcher. A function that always returns 0 passes it.

Queue items for a database on more than one core and assert the summed depth.

Comment thread nodedb/src/control/metrics/database.rs Outdated
}

#[test]
fn add_maintenance_cpu_secs_accumulates() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This tests an existing setter that the PR does not wire. It does not cover the change.

Remove it, or replace it with a test of the sampler.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

yeah i am just doing that
i actually missed out

@AubaidAhmedSaiyed

Copy link
Copy Markdown
Author

please check

@farhan-syah farhan-syah left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the fast turnaround. Most of round one landed. Two round-one blockers are still open. The new maintenance wiring adds two more.

Round one

Point Status
percentile(99.0) pins WAL p99 Fixed. observability.rs is fixed too.
Duplicate WAL fsync histogram Fixed. The group-commit path feeds record_wal_fsync.
Timer includes lock wait Fixed.
Cumulative p99, node-wide value Fixed. The p99 is computed per 10 s window, and the gauge doc says it is node-wide.
Memory 0 without a quota Fixed. A new race is noted inline at reserve.rs:71.
Connections 0 without a cap Fixed.
rustfmt --check Fixed.
background_loops.rs over 500 lines Fixed. The sampler has its own file.
Dispatcher lock per database Fixed.
Queue-depth test proves nothing Fixed. The test now sums across two cores.
Storage source is a stub Open. The new source measures the wrong thing. See catalog/database.rs:232.
Stale series after DROP DATABASE Open. The removal never matches. See post_apply/database.rs:28.
add_maintenance_cpu_secs constant 0 Wired, but the counter is now written from the Data Plane. See budget.rs:304.

Blockers

# Defect Inline
1 Storage sums catalog metadata, not database data, with a full catalog scan per database every tick catalog/database.rs:232
2 MaintenanceLease::drop runs on the Data Plane and writes a Control Plane registry budget.rs:304
3 Drop cleanup looks up the name after apply deletes it post_apply/database.rs:28
4 db-{id} is a guessed label that matches no real series budget.rs:118, post_apply/database.rs:31

Commit history: rebuild it

The branch needs a rebuild before merge. main keeps every PR commit, so these messages land as they are.

  • Issue and PR numbers. (fixes #375), from #375, and address PR #375 review put tracker numbers in commit messages. This repo keeps commit messages free of tracker numbers. The last one also names the wrong number, since this PR is #386. Put Fixes #375 in the PR body only.
  • Merge commit. Merge branch 'upstream/main' into fix-sampler adds a merge commit inside the PR. Rebase onto main instead.
  • Fix-up commits. "address … review" is a fix-up of the first commit. Squash the review fix-ups into the commits they fix.
  • Message shape. Each message states what changed and why, in the type(scope): summary form used on main.

Rebase onto current main, rebuild the history as described, and force-push the branch.

Should-fix

  • MaintenanceBudgetTracker::try_acquire_named and try_acquire_with_metrics have no callers. Remove them.
  • set_cap_named repeats the cap arithmetic in set_cap. Call one from the other.
  • init.rs:399 and post_init.rs:37 repeat the same wiring block. Move it into one function that both call.
  • try_acquire_database still documents Ok(None) for a database with no entry. That case no longer occurs. Update the doc.
  • governor/metrics.rs:17 still says 0 is returned without a scoped budget. Every database is now tracked. Update the doc.
  • DatabaseCounters repeats the doc line /// Current active connection count.
  • reserve.rs repeats the get-or-insert block in both entry points. Extract one helper.

/// 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<u64> {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. This sums rows in the system catalog: DDL descriptors, surrogate maps, and topic rows. That is catalog metadata. Collection data lives in the engines, not here. A database with 1 TB of documents reports a few kilobytes.

The cost is also unbounded. SURROGATE_PK_V3 and SURROGATE_PK_REV_V3 hold one row per user row across every database. The sampler scans them in full, once per database, every 10 s.

Every if let Ok(...) also drops read errors without a word, so a corrupt table reads as a smaller number.

Required: read storage from the owners of the data. That means per-database byte counts kept by the engines and the WAL, updated on write and compaction. The sampler reads the counter and never scans.

.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);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

unwrap_or(0) turns a catalog read error into a 0 gauge. Nothing in the output shows that the read failed.

Log the error and skip the storage gauge for that tick. Do not publish 0.

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);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. Leases are acquired and dropped on the Data Plane, in the compaction handlers (compact/budget.rs, compact/segments.rs, compact/runner.rs). This Drop now calls DatabaseMetricsRegistry::add_maintenance_cpu_secs, which runs get_or_create. That takes the registry RwLock, a write lock for a new name, and it allocates. The Control Plane holds the same lock while it renders /metrics. A Data Plane core can stall behind a scrape.

Only the SPSC bridge and the Event Bus carry data across planes.

Required: keep a cumulative per-database total inside the tracker, next to the sliding window. The Control Plane sampler reads that total each tick and writes the registry. With that change the tracker needs no names, no registry handle, and no MaintenanceLease::db_name.

if db == DatabaseId::DEFAULT {
return "default".to_string();
}
format!("db-{}", db.as_u64())

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. db-{id} is a guess at a label. No series uses it, so maintenance time lands under a name no dashboard queries.

This goes away when the sampler owns the name lookup (see the comment on line 304). The sampler already has real names from list_databases.

pub fn delete(db_id: u64, shared: Arc<SharedState>) {
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) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. Post-apply runs after apply::database::delete, which already removed the descriptor. get_database_name_by_id returns None here every time. The removal by real name never runs.

Line 31 then removes db-{id}, which matches no series. The dropped database's series stay in /metrics.

Required: remove the series by the name the registry actually uses. One option is a registry keyed by DatabaseId, with the name as a label. Another is to carry the name in the DeleteDatabase entry. Add a test that runs DROP DATABASE through the real apply path and checks that the series is gone.

.database_budgets
.read()
.unwrap_or_else(|p| p.into_inner());
map.get(&db).cloned()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The read lock is now released before try_reserve. The budget is a clone, so its limit is a snapshot. Before this change, a quota change under the write lock excluded any reservation in progress. Now a reservation can pass against a quota that was lowered after the snapshot.

Hold the read lock across try_reserve for the existing-entry path, as before. Take the write lock only to insert a missing entry. Then drop Clone from ScopedBudget, because it breaks the "mutate limit in place" rule its own doc states.

"nodedb_database_bridge_queue_depth{{database=\"{db_name}\"}} 5"
)));
assert!(text.contains(&format!(
"nodedb_database_wal_commit_latency_p99_us{{database=\"{db_name}\"}}"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This asserts that the p99 line exists, never its value. A constant 0 or 1_000_000 still passes.

The fsyncs recorded above give a known distribution. Assert the p99 value falls in the bucket they produce.

(100..=500).contains(&p99),
"expected p99 in [100, 500], got {p99}"
);
// Passing 99.0 would have returned 1_000_000 (the last boundary).

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This comment describes the old bug, and the assert under it tests the old call site. Once merged, neither makes sense.

Remove the comment and the assert_ne!. The range assert above already pins the value.

}

#[test]
fn add_maintenance_cpu_secs_accumulates() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This duplicates lease_drop_increments_prometheus_metrics in budget.rs. It also tests the tracker from the metrics module.

Remove it. The registry's own behavior is covered by remove_clears_database_metrics and gauge_setters_update_counters.

}

#[test]
fn unbudgeted_database_tracks_8192_bytes_and_releases_on_drop() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is uncapped_database_usage_tracks_allocation again with a different size and id. Remove one of them.

@AubaidAhmedSaiyed

Copy link
Copy Markdown
Author

sure

@AubaidAhmedSaiyed AubaidAhmedSaiyed changed the title fix(obsv): wire live per-database metrics sampler (fixes #375) fix(obsv): wire live per-database metrics sampler (fixes #386) Sep 30, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

obsv: six per-database metric setters have no writer; add a sampler

2 participants