diff --git a/docs/architecture.md b/docs/architecture.md index c132910af..28429b98f 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -204,6 +204,8 @@ Admission ordering for requests: Memory reservation flips the order (largest scope first to fail fast): global → database → tenant → engine. Failure at any layer aborts before the next is consulted. +A dispatch refused by the bridge — a full per-core weighted-fair queue, a per-database virtual queue at its fair share, or a tenant at its in-flight cap — surfaces as `SERVER_OVERLOAD` (`57P03`). That class is retryable: the client waits for capacity and resends the same request. `nodedb_dispatch_capacity_busy_total` counts every such refusal. + ### Memory Governor Located in `nodedb-mem/src/governor.rs`, the `MemoryGovernor` enforces hierarchical memory reservations: @@ -342,6 +344,7 @@ 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 +- **Dispatch capacity refusals**: bridge dispatcher, at every refusal — WFQ full, per-database suspension, tenant in-flight cap (`nodedb_dispatch_capacity_busy_total`, also the `dispatch_capacity_busy_total` row of `SHOW STATS`) All metrics are dimensionalized by database and tenant to enable per-customer tracking and alerting. diff --git a/nodedb/src/bridge/dispatch/dispatcher.rs b/nodedb/src/bridge/dispatch/dispatcher.rs index 35981fd38..9cf62c95c 100644 --- a/nodedb/src/bridge/dispatch/dispatcher.rs +++ b/nodedb/src/bridge/dispatch/dispatcher.rs @@ -8,6 +8,7 @@ use std::collections::{HashMap, HashSet}; use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; use nodedb_bridge::backpressure::{BackpressureConfig, BackpressureController, PressureState}; use nodedb_bridge::buffer::RingBuffer; @@ -244,3 +245,29 @@ impl Dispatcher { &self.router } } + +/// Dispatches refused because a capacity limit had no room, since process +/// start. +/// +/// Every capacity refusal carries [`crate::Error::DispatchCapacity`], which +/// classifies as `server_overload` and reaches a pgwire client as `57P03`. The +/// count is process-wide because the condition is the same one whichever core, +/// database, or tenant hit its limit. +/// +/// Rendered as `nodedb_dispatch_capacity_busy_total` on `/metrics` and as the +/// `dispatch_capacity_busy_total` row of `SHOW STATS`. +static DISPATCH_CAPACITY_BUSY_TOTAL: AtomicU64 = AtomicU64::new(0); + +/// Count one dispatch refused on a capacity limit. +/// +/// Called only from the dispatch paths that raise +/// [`crate::Error::DispatchCapacity`], so the operator-visible counter and the +/// client-visible class cannot drift apart. +pub(crate) fn note_capacity_busy() { + DISPATCH_CAPACITY_BUSY_TOTAL.fetch_add(1, Ordering::Relaxed); +} + +/// Dispatches refused on a capacity limit since process start. +pub fn dispatch_capacity_busy_total() -> u64 { + DISPATCH_CAPACITY_BUSY_TOTAL.load(Ordering::Relaxed) +} diff --git a/nodedb/src/bridge/dispatch/enqueue.rs b/nodedb/src/bridge/dispatch/enqueue.rs index ada1ef612..2dab2f17f 100644 --- a/nodedb/src/bridge/dispatch/enqueue.rs +++ b/nodedb/src/bridge/dispatch/enqueue.rs @@ -18,6 +18,18 @@ use crate::types::Lsn; use super::dispatcher::Dispatcher; use super::refusal::DispatchRefusal; +/// Refuse a dispatch on a capacity limit, counting it once. +/// +/// Every capacity refusal goes through here, so the counter the operator reads +/// and the `57P03` class the client receives always describe the same event. +fn capacity_refusal( + scope: DispatchCapacityScope, + request: envelope::Request, +) -> Box { + super::dispatcher::note_capacity_busy(); + DispatchRefusal::boxed(crate::Error::DispatchCapacity { scope }, request) +} + impl Dispatcher { /// Dispatch a request to the correct Data Plane core. /// @@ -54,10 +66,7 @@ impl Dispatcher { inflight, cap: self.max_per_tenant_inflight, }; - return Err(DispatchRefusal::boxed( - crate::Error::DispatchCapacity { scope }, - request, - )); + return Err(capacity_refusal(scope, request)); } } @@ -80,10 +89,7 @@ impl Dispatcher { database_id: request.database_id, core_id, }; - return Err(DispatchRefusal::boxed( - crate::Error::DispatchCapacity { scope }, - request, - )); + return Err(capacity_refusal(scope, request)); } // Enqueue into the WFQ. A full queue hands the request back. @@ -92,10 +98,7 @@ impl Dispatcher { core_id, capacity: self.per_core_capacity, }; - return Err(DispatchRefusal::boxed( - crate::Error::DispatchCapacity { scope }, - request, - )); + return Err(capacity_refusal(scope, request)); } self.commit_enqueued(core_id, database_id, tenant_id, req_id, wal_lsn); @@ -129,14 +132,16 @@ impl Dispatcher { let cls = self.priority_resolver.priority_for(database_id); channel.wfq.set_priority(database_id, cls); - channel.wfq.try_enqueue(database_id, request).map_err(|_| { - crate::Error::DispatchCapacity { - scope: DispatchCapacityScope::QueueFull { - core_id, - capacity: self.per_core_capacity, - }, - } - })?; + if let Err(_request) = channel.wfq.try_enqueue(database_id, request) { + let scope = DispatchCapacityScope::QueueFull { + core_id, + capacity: self.per_core_capacity, + }; + // This path reports a flat `Error` rather than a boxed refusal + // that hands the request back, so the count is noted directly. + super::dispatcher::note_capacity_busy(); + return Err(crate::Error::DispatchCapacity { scope }); + } self.commit_enqueued(core_id, database_id, tenant_id, req_id, wal_lsn); Ok(()) @@ -365,4 +370,76 @@ mod tests { let _ = dispatcher.db_pressure_on_core(0, 1); let _ = dispatcher.db_pressure_on_core(0, 2); } + + /// A capacity refusal counts, whichever limit refused it. + /// + /// The counter is process-wide, so the assertion is a lower bound: a test + /// running beside this one in the same binary may refuse too. + #[test] + fn capacity_refusals_are_counted_once_each() { + let before = crate::bridge::dispatch::dispatch_capacity_busy_total(); + + // A full weighted-fair queue. Distinct tenants and databases keep the + // tenant cap and the per-database suspension out of play. + let (mut dispatcher, _) = Dispatcher::new(1, 4); + let mut queue_full = false; + for i in 1..=64u64 { + let mut request = make_request_for_db(0, i, i); + request.tenant_id = TenantId::new(i); + match dispatcher.dispatch(request) { + Ok(()) => continue, + Err(crate::Error::DispatchCapacity { + scope: DispatchCapacityScope::QueueFull { .. }, + }) => { + queue_full = true; + break; + } + Err(other) => panic!("expected a queue-full refusal, got: {other}"), + } + } + assert!( + queue_full, + "a one-core queue of four must fill within 64 distinct requests" + ); + + // The per-tenant in-flight cap. + let (mut dispatcher, _) = Dispatcher::new(1, 4); + for i in 0..4u64 { + dispatcher + .dispatch(make_request_for_db(0, i + 1, i + 1)) + .unwrap(); + } + let refusal = dispatcher + .try_dispatch(make_request_for_db(0, 99, 99)) + .expect_err("the fifth request exceeds the tenant cap"); + assert!( + matches!( + refusal.error, + crate::Error::DispatchCapacity { + scope: DispatchCapacityScope::TenantInflight { .. } + } + ), + "the second refusal must be the tenant cap, got: {}", + refusal.error + ); + + let after = crate::bridge::dispatch::dispatch_capacity_busy_total(); + assert!( + after >= before + 2, + "two capacity refusals must count at least twice: {before} → {after}" + ); + } + + /// A terminal refusal is not a capacity refusal, so it carries no scope. + #[test] + fn terminal_refusal_is_not_a_capacity_refusal() { + let (mut dispatcher, _) = Dispatcher::new(1, 4); + let error = dispatcher + .dispatch_to_core(7, make_request(0)) + .expect_err("core 7 does not exist on a one-core dispatcher"); + assert!( + matches!(error, crate::Error::Dispatch { .. }), + "an out-of-range core is a terminal dispatch error, got: {error}" + ); + } } diff --git a/nodedb/src/bridge/dispatch/mod.rs b/nodedb/src/bridge/dispatch/mod.rs index 2bd5862d2..224d604fb 100644 --- a/nodedb/src/bridge/dispatch/mod.rs +++ b/nodedb/src/bridge/dispatch/mod.rs @@ -15,7 +15,7 @@ mod test_requests; pub use core_channel::{CoreChannel, CoreChannelDataSide}; pub use dispatcher::{ BridgeRequest, BridgeResponse, DATA_PLANE_QUEUE_CAPACITY, DatabasePriorityResolver, - DefaultPriorityResolver, Dispatcher, + DefaultPriorityResolver, Dispatcher, dispatch_capacity_busy_total, }; pub use drain::CorePending; pub use outcome_floor::{OutcomeFloor, ResendRefusal, StuckFloor, WriteWindow}; diff --git a/nodedb/src/control/server/http/routes/metrics.rs b/nodedb/src/control/server/http/routes/metrics.rs index 3afb56ad1..f09052136 100644 --- a/nodedb/src/control/server/http/routes/metrics.rs +++ b/nodedb/src/control/server/http/routes/metrics.rs @@ -50,6 +50,19 @@ pub async fn metrics( output.push_str("# TYPE nodedb_wal_next_lsn gauge\n"); output.push_str(&format!("nodedb_wal_next_lsn {wal_lsn}\n\n")); + // Dispatches refused because a capacity limit had no room. The same refusal + // reaches the client as the retryable `57P03` class, so a non-zero counter + // tells an operator that bulk writers are backing off rather than failing. + output.push_str( + "# HELP nodedb_dispatch_capacity_busy_total \ + Dispatches refused because a dispatch capacity limit had no room.\n", + ); + output.push_str("# TYPE nodedb_dispatch_capacity_busy_total counter\n"); + output.push_str(&format!( + "nodedb_dispatch_capacity_busy_total {}\n\n", + crate::bridge::dispatch::dispatch_capacity_busy_total() + )); + // Outcome floor: every engine watermark and WAL truncation stays at or // below it. let outcome_floor = &state.shared.outcome_floor; diff --git a/nodedb/src/control/server/shared/ddl/neutral/observability.rs b/nodedb/src/control/server/shared/ddl/neutral/observability.rs index 94cfa2a4f..67ae87431 100644 --- a/nodedb/src/control/server/shared/ddl/neutral/observability.rs +++ b/nodedb/src/control/server/shared/ddl/neutral/observability.rs @@ -141,6 +141,13 @@ fn server_stats_rows(state: &SharedState) -> Vec<(String, String)> { )); } + // Not a `SystemMetrics` counter: the dispatcher counts its own refusals, + // and the row keeps the SQL surface in step with `/metrics`. + rows.push(( + "dispatch_capacity_busy_total".into(), + crate::bridge::dispatch::dispatch_capacity_busy_total().to_string(), + )); + rows } diff --git a/nodedb/tests/wire/cases/pgwire_show_dispatch.rs b/nodedb/tests/wire/cases/pgwire_show_dispatch.rs index 6a596a9f9..a3d21aa56 100644 --- a/nodedb/tests/wire/cases/pgwire_show_dispatch.rs +++ b/nodedb/tests/wire/cases/pgwire_show_dispatch.rs @@ -460,3 +460,69 @@ async fn show_session_set_parameter_round_trips() { .expect("SHOW application_name must succeed"); assert_eq!(rows, vec!["mae8_bootstrap".to_string()]); } + +// ── Dispatch capacity counter ──────────────────────────────────────── + +/// Fetch the raw Prometheus text body from `/metrics`. +async fn fetch_metrics(http_port: u16) -> String { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let mut stream = tokio::net::TcpStream::connect(("127.0.0.1", http_port)) + .await + .expect("connect to /metrics"); + let req = b"GET /metrics HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n"; + stream.write_all(req).await.expect("write metrics request"); + let mut body = String::new(); + stream + .read_to_string(&mut body) + .await + .expect("read metrics response"); + body +} + +/// The dispatcher's refusal counter must be readable from SQL, so an operator +/// without HTTP access can still see dispatch pressure. A refusal reaches the +/// client as `57P03`, and this row is the operator's view of the same event. +#[tokio::test] +async fn show_stats_carries_the_dispatch_capacity_counter() { + let server = TestServer::start().await; + let rows = server + .query_named_rows("SHOW STATS") + .await + .expect("SHOW STATS must succeed"); + + let row = rows + .iter() + .find(|r| { + r.get("name") + .map(|name| name == "dispatch_capacity_busy_total") + .unwrap_or(false) + }) + .unwrap_or_else(|| { + panic!("SHOW STATS must carry dispatch_capacity_busy_total; got {rows:?}") + }); + + let value = row.get("value").expect("the counter row carries a value"); + value.parse::().unwrap_or_else(|_| { + panic!("dispatch_capacity_busy_total must be a decimal integer, got {value:?}") + }); +} + +/// `/metrics` must declare the counter and export it. The declaration is what +/// a Prometheus scrape uses to read it as a monotonic counter rather than a +/// gauge. +#[tokio::test] +async fn metrics_exposes_the_dispatch_capacity_counter() { + let server = TestServer::start().await; + let body = fetch_metrics(server.http_port).await; + + assert!( + body.contains("# TYPE nodedb_dispatch_capacity_busy_total counter"), + "/metrics must declare nodedb_dispatch_capacity_busy_total as a counter" + ); + assert!( + body.lines() + .any(|line| line.starts_with("nodedb_dispatch_capacity_busy_total ")), + "/metrics must export a nodedb_dispatch_capacity_busy_total sample" + ); +}