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
3 changes: 3 additions & 0 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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.

Expand Down
27 changes: 27 additions & 0 deletions nodedb/src/bridge/dispatch/dispatcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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)
}
117 changes: 97 additions & 20 deletions nodedb/src/bridge/dispatch/enqueue.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<DispatchRefusal> {
super::dispatcher::note_capacity_busy();
DispatchRefusal::boxed(crate::Error::DispatchCapacity { scope }, request)
}

impl Dispatcher {
/// Dispatch a request to the correct Data Plane core.
///
Expand Down Expand Up @@ -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));
}
}

Expand All @@ -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.
Expand All @@ -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);
Expand Down Expand Up @@ -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(())
Expand Down Expand Up @@ -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}"
);
}
}
2 changes: 1 addition & 1 deletion nodedb/src/bridge/dispatch/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down
13 changes: 13 additions & 0 deletions nodedb/src/control/server/http/routes/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
7 changes: 7 additions & 0 deletions nodedb/src/control/server/shared/ddl/neutral/observability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
66 changes: 66 additions & 0 deletions nodedb/tests/wire/cases/pgwire_show_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<u64>().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"
);
}
Loading