Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -342,6 +342,16 @@ Counters are incremented at the sites where resources are consumed:
- **WAL latency**: Group commit fsync completion
- **Maintenance CPU**: Lease acquisition / release
- **Replication lag**: Raft follower log application
- **Graph edge writes**: graph edge put handler, per applied edge in a single put or a batch put (`nodedb_graph_edges_written_total`, also the `graph_edges_written_total` row of `SHOW STATS`)
- **Graph edge deletes**: graph edge delete handler, per live edge tombstoned in a single delete or a batch delete (`nodedb_graph_edges_deleted_total`, also the `graph_edges_deleted_total` row of `SHOW STATS`) — WIP, see below

`nodedb_graph_edges_written_total` counts applied edge versions, so a put that rewrites a live edge increments it while the `nodedb_graph_edges` gauge stays flat. `nodedb_graph_edges_deleted_total` counts live edges removed, so a delete of an edge that was already absent writes a tombstone and increments neither counter — matching the affected count that statement reports.

A cross-vShard edge is dual-homed: the same `EdgePut` or `EdgeDelete` runs on both endpoint homes. The write counter is recorded under `owns_logical_edge_stats` — the same ownership boundary `SHOW GRAPH STATS` counts logical edges with — so the source home counts the edge version once and the destination home adds nothing.

`nodedb_graph_edges_deleted_total` is **WIP and over-reports on a multi-core cluster**. It counts each endpoint home that removes a live row. On one core the two dual-home participants share a store, so only the first finds the edge and the counter reads one; on two or more cores each home keeps its own copy, both find a live row, and one logical delete reads as two. Counting only the source home instead reads zero on one core, because the destination home performs the removal there. The deciding fact is the Control Plane's `single_home`, which the plan does not carry; distinguishing the cases needs a new `GraphOp::EdgeDelete` field threaded from `graph_ops/edge.rs`. Do not rely on this counter on a multi-core deployment until that lands.

Both counters record client activity, so neither counts a WAL-replayed edge. Boot replay re-enters the same handlers to rebuild engine state, and the metrics are attached before replay runs; counting there would report pre-restart writes as new work on every restart. The handlers skip both counters while `CoreLoop::boot_replaying_wal()` holds, which covers the whole boot pass and is never set on a serving core. An online committed-redo apply reaches the same handlers and does count, because it is a real write.

All metrics are dimensionalized by database and tenant to enable per-customer tracking and alerting.

Expand Down
14 changes: 14 additions & 0 deletions nodedb/src/control/metrics/prometheus/engines.rs
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,20 @@ impl SystemMetrics {
"Graph edges stored",
self.graph_edges.load(Ordering::Relaxed),
);
counter(
out,
"nodedb_graph_edges_written_total",
"Edge versions applied by a graph edge write",
self.graph_edges_written.load(Ordering::Relaxed),
);
counter(
out,
"nodedb_graph_edges_deleted_total",
"Live edges tombstoned by a graph edge delete (WIP: counted per \
endpoint home, so a cross-shard delete on a multi-core cluster \
reports one logical removal per home)",
self.graph_edges_deleted.load(Ordering::Relaxed),
);

// ── Document engine ──
counter(
Expand Down
8 changes: 8 additions & 0 deletions nodedb/src/control/metrics/system/fields.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,14 @@ pub struct SystemMetrics {
pub graph_traversals: AtomicU64,
pub graph_nodes: AtomicU64,
pub graph_edges: AtomicU64,
/// Edge versions applied by a graph edge write (single and batch puts).
/// Distinct from [`Self::graph_edges`], a gauge of live edges: a put that
/// rewrites an edge leaves the gauge flat and still counts here.
pub graph_edges_written: AtomicU64,
/// Live edges tombstoned by a graph edge delete (single and batch deletes).
/// A delete of an edge that was already absent writes a tombstone but
/// removes no live edge, so it does not count.
pub graph_edges_deleted: AtomicU64,

pub document_inserts: AtomicU64,
pub document_reads: AtomicU64,
Expand Down
17 changes: 17 additions & 0 deletions nodedb/src/control/metrics/system/record.rs
Original file line number Diff line number Diff line change
Expand Up @@ -217,6 +217,23 @@ impl SystemMetrics {
self.graph_edges.store(edges, Ordering::Relaxed);
}

/// One edge version was applied by a graph edge write.
///
/// Batch handlers call this once per edge as it is applied, never once per
/// batch: a batch that fails midway has already applied the edges before
/// the failure, and those writes are counted.
pub fn record_graph_edge_written(&self) {
self.graph_edges_written.fetch_add(1, Ordering::Relaxed);
}

/// One live edge was tombstoned by a graph edge delete.
///
/// Callers pass the pre-image result: a delete whose pre-image was absent
/// removes no live edge, so it does not count.
pub fn record_graph_edge_deleted(&self) {
self.graph_edges_deleted.fetch_add(1, Ordering::Relaxed);
}

// ── Document engine ──

pub fn record_document_insert(&self) {
Expand Down
84 changes: 84 additions & 0 deletions nodedb/src/control/server/shared/ddl/neutral/observability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,14 @@ fn server_stats_rows(state: &SharedState) -> Vec<(String, String)> {
"queries_graph".into(),
sys.queries_graph.load(Ordering::Relaxed).to_string(),
));
rows.push((
"graph_edges_written_total".into(),
sys.graph_edges_written.load(Ordering::Relaxed).to_string(),
));
rows.push((
"graph_edges_deleted_total".into(),
sys.graph_edges_deleted.load(Ordering::Relaxed).to_string(),
));
rows.push((
"queries_document".into(),
sys.queries_document.load(Ordering::Relaxed).to_string(),
Expand Down Expand Up @@ -254,3 +262,79 @@ pub fn show_memory(
rows,
))])
}

#[cfg(test)]
mod tests {
use super::*;
use crate::control::metrics::SystemMetrics;
use std::sync::Arc;

/// Read one `(name, value)` row from a rendered row set.
fn row_value(rows: &[(String, String)], name: &str) -> Option<String> {
rows.iter()
.find(|(key, _)| key == name)
.map(|(_, value)| value.clone())
}

/// The graph write counters render on the Prometheus surface as counters
/// with their own help text, next to the live-edge gauge they must not be
/// confused with.
#[test]
fn the_graph_write_counters_render_on_the_prometheus_surface() {
let metrics = SystemMetrics::new();
metrics.record_graph_edge_written();
metrics.record_graph_edge_deleted();
let output = metrics.to_prometheus();

assert!(
output.contains("# TYPE nodedb_graph_edges_written_total counter"),
"the write counter renders as a counter: {output}"
);
assert!(
output.contains("nodedb_graph_edges_written_total 1"),
"the write counter carries its value: {output}"
);
assert!(
output.contains("# TYPE nodedb_graph_edges_deleted_total counter"),
"the delete counter renders as a counter: {output}"
);
assert!(
output.contains("nodedb_graph_edges_deleted_total 1"),
"the delete counter carries its value: {output}"
);
}

/// The row builder reads the counters off `SystemMetrics`, the same
/// instance the Data Plane records through, so `SHOW STATS` cannot drift
/// from `/metrics`. Rows carry the value as decimal text like every other
/// counter in the set.
#[test]
fn the_stats_rows_read_the_graph_write_counters_from_metrics() {
let directory = tempfile::tempdir().expect("tempdir");
let wal = Arc::new(
crate::wal::WalManager::open_for_testing(&directory.path().join("obs-stats.wal"))
.expect("open WAL"),
);
let (dispatcher, _data_sides) = crate::bridge::dispatch::Dispatcher::new(1, 64);
let state = SharedState::new(dispatcher, wal).expect("shared state");
let metrics = state
.system_metrics
.as_ref()
.expect("system metrics are wired")
.clone();
metrics.record_graph_edge_written();
metrics.record_graph_edge_written();
metrics.record_graph_edge_deleted();

let rows = server_stats_rows(&state);

assert_eq!(
row_value(&rows, "graph_edges_written_total").as_deref(),
Some("2"),
);
assert_eq!(
row_value(&rows, "graph_edges_deleted_total").as_deref(),
Some("1"),
);
}
}
1 change: 1 addition & 0 deletions nodedb/src/data/executor/core_loop/open.rs
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,7 @@ impl CoreLoop {
redo_apply:
crate::data::executor::handlers::transaction::redo_apply::RedoApplyState::new(),
fail_stop: super::fail_stop::CoreFailStop::default(),
boot_replaying_wal: false,
})
}
}
5 changes: 5 additions & 0 deletions nodedb/src/data/executor/core_loop/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -491,4 +491,9 @@ pub struct CoreLoop {
crate::data::executor::handlers::transaction::redo_apply::RedoApplyState,
/// Set once this core's state is unknown. It then refuses every request.
pub(in crate::data::executor) fail_stop: super::fail_stop::CoreFailStop,

/// True while this core rebuilds engine state from the WAL at boot.
/// Boot replay re-enters the write handlers, so a handler that counts
/// client activity asks this first. See [`Self::boot_replaying_wal`].
pub(in crate::data::executor) boot_replaying_wal: bool,
}
Loading
Loading