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
79 changes: 78 additions & 1 deletion nodedb-sql/src/ddl_ast/graph_parse/entry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,11 @@ pub fn try_parse(sql: &str) -> Option<Result<NodedbStatement, SqlError>> {

let toks = tokenizer::tokenize(trimmed);

let parsed = if upper.starts_with("GRAPH INSERT EDGE ") {
let parsed = if upper.starts_with("GRAPH INSERT EDGES ") {
variants::parse_insert_edges(&toks)
} else if upper.starts_with("GRAPH DELETE EDGES ") {
variants::parse_delete_edges(&toks)
} else if upper.starts_with("GRAPH INSERT EDGE ") {
variants::parse_insert_edge(&toks)
} else if upper.starts_with("GRAPH DELETE EDGE ") {
variants::parse_delete_edge(&toks)
Expand Down Expand Up @@ -114,6 +118,79 @@ mod tests {
/// A missing required clause is a malformed graph statement, not a
/// non-graph one. `None` here would send it to the SQL parser, which
/// reports only that `GRAPH` is not SQL.
#[test]
fn parse_graph_insert_edges_batch() {
let stmt =
parsed("GRAPH INSERT EDGES IN 'edges' VALUES ('a','b','CALLS'), ('c','d','IMPORTS')");
match stmt {
NodedbStatement::Graph(GraphStmt::GraphInsertEdges { collection, edges }) => {
assert_eq!(collection, "edges");
assert_eq!(edges.len(), 2);
assert_eq!(edges[0].src, "a");
assert_eq!(edges[0].dst, "b");
assert_eq!(edges[0].label, "CALLS");
assert_eq!(edges[1].src, "c");
assert_eq!(edges[1].label, "IMPORTS");
}
other => panic!("expected GraphInsertEdges, got {other:?}"),
}
}

#[test]
fn parse_graph_delete_edges_batch_accepts_bare_words() {
let stmt = parsed("GRAPH DELETE EDGES IN 'edges' VALUES (a, b, CALLS)");
match stmt {
NodedbStatement::Graph(GraphStmt::GraphDeleteEdges { collection, edges }) => {
assert_eq!(collection, "edges");
assert_eq!(edges.len(), 1);
assert_eq!(edges[0].src, "a");
assert_eq!(edges[0].dst, "b");
assert_eq!(edges[0].label, "CALLS");
}
other => panic!("expected GraphDeleteEdges, got {other:?}"),
}
}

#[test]
fn batch_edge_cap_is_enforced() {
let mut sql = String::from("GRAPH INSERT EDGES IN 'edges' VALUES ");
for i in 0..=variants::MAX_EDGES_PER_BATCH {
if i > 0 {
sql.push(',');
}
sql.push_str(&format!("('s{i}','d{i}','L')"));
}
let error = try_parse(&sql)
.expect("input is graph DSL")
.expect_err("a batch over the cap must not produce a statement");
assert!(
error.to_string().contains("at most 1000 edges"),
"the error must name the cap: {error}"
);
}

#[test]
fn batch_edge_rejects_malformed_triples() {
let error = try_parse("GRAPH INSERT EDGES IN 'edges' VALUES ('a','b')")
.expect("input is graph DSL")
.expect_err("a partial triple must not produce a statement");
assert!(
error.to_string().contains("triples"),
"the error must name the tuple shape: {error}"
);
}

#[test]
fn batch_edge_rejects_per_edge_properties() {
let error = try_parse("GRAPH INSERT EDGES IN 'edges' VALUES ('a','b','L') PROPERTIES '{}'")
.expect("input is graph DSL")
.expect_err("per-edge PROPERTIES is not supported in the batch form");
assert!(
error.to_string().contains("PROPERTIES"),
"the error must name the rejected clause: {error}"
);
}

#[test]
fn parse_graph_insert_edge_missing_collection_names_the_clause() {
let error = try_parse("GRAPH INSERT EDGE FROM 'a' TO 'b' TYPE 'l'")
Expand Down
86 changes: 84 additions & 2 deletions nodedb-sql/src/ddl_ast/graph_parse/variants.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,11 +12,12 @@ use super::{
super::statement::{GraphStmt, NodedbStatement},
fusion_params::{FusionParams, RAG_FUSION_KEYWORDS},
helpers::{
direction_after, extract_properties, missing_clause, quoted_after, quoted_list_after,
usize_after, usize_after_checked, word_after,
direction_after, extract_properties, find_keyword, missing_clause, quoted_after,
quoted_list_after, usize_after, usize_after_checked, word_after,
},
tokenizer::Tok,
};
use crate::ddl_ast::graph_types::GraphEdgeTuple;
use crate::error::SqlError;

pub(super) fn parse_insert_edge(toks: &[Tok<'_>]) -> Result<NodedbStatement, SqlError> {
Expand All @@ -36,6 +37,87 @@ pub(super) fn parse_insert_edge(toks: &[Tok<'_>]) -> Result<NodedbStatement, Sql
}))
}

/// Maximum number of edges one batch statement may carry. Keeps one
/// statement's worth of work bounded for the write-admission and lock paths.
pub const MAX_EDGES_PER_BATCH: usize = 1000;

pub(super) fn parse_insert_edges(toks: &[Tok<'_>]) -> Result<NodedbStatement, SqlError> {
let (collection, edges) = parse_edge_batch(toks, "GRAPH INSERT EDGES")?;
Ok(NodedbStatement::Graph(GraphStmt::GraphInsertEdges {
collection,
edges,
}))
}

pub(super) fn parse_delete_edges(toks: &[Tok<'_>]) -> Result<NodedbStatement, SqlError> {
let (collection, edges) = parse_edge_batch(toks, "GRAPH DELETE EDGES")?;
Ok(NodedbStatement::Graph(GraphStmt::GraphDeleteEdges {
collection,
edges,
}))
}

/// Shared body of the batch parsers.
///
/// Grammar: `IN '<collection>' VALUES ('<src>','<dst>','<label>')[, (...)]*`.
/// The tokenizer drops commas and parentheses, so the tuple structure is
/// recovered by consuming string tokens in threes; a count that is not a
/// multiple of three is malformed. Per-edge `PROPERTIES` stays on the
/// single-edge form until the physical `BatchEdge` can carry one.
fn parse_edge_batch(
toks: &[Tok<'_>],
stmt: &str,
) -> Result<(String, Vec<GraphEdgeTuple>), SqlError> {
let collection =
quoted_after(toks, "IN").ok_or_else(|| missing_clause(stmt, "IN <collection>"))?;
let Some(values_pos) = find_keyword(toks, "VALUES") else {
return Err(missing_clause(stmt, "VALUES ('<src>','<dst>','<label>')"));
};
let mut fields: Vec<String> = Vec::new();
for t in &toks[values_pos + 1..] {
match t {
Tok::Quoted(s) => fields.push(s.clone().into_owned()),
Tok::Word(w) => {
if w.eq_ignore_ascii_case("PROPERTIES") {
return Err(SqlError::Parse {
detail: format!(
"{stmt}: per-edge PROPERTIES is not supported in the batch form"
),
});
}
fields.push((*w).to_string());
}
Tok::Object(_) => {
return Err(SqlError::Parse {
detail: format!("{stmt}: object literals are not supported in the batch form"),
});
}
}
}
if fields.is_empty() || !fields.len().is_multiple_of(3) {
return Err(SqlError::Parse {
detail: format!("{stmt}: VALUES takes (src, dst, label) triples"),
});
}
let count = fields.len() / 3;
if count > MAX_EDGES_PER_BATCH {
return Err(SqlError::Parse {
detail: format!(
"{stmt}: at most {MAX_EDGES_PER_BATCH} edges per statement, got {count}"
),
});
}
let edges = fields
.chunks(3)
.map(|c| GraphEdgeTuple {
src: c[0].clone(),
dst: c[1].clone(),
label: c[2].clone(),
})
.collect();
Ok((collection, edges))
}

pub(super) fn parse_delete_edge(toks: &[Tok<'_>]) -> Result<NodedbStatement, SqlError> {
const STMT: &str = "GRAPH DELETE EDGE";
let collection =
Expand Down
9 changes: 9 additions & 0 deletions nodedb-sql/src/ddl_ast/graph_types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,3 +27,12 @@ pub enum GraphProperties {
/// expected to already be a JSON document.
Quoted(String),
}

/// One `(src, dst, label)` triple inside a batched edge statement
/// (`GRAPH INSERT EDGES` / `GRAPH DELETE EDGES`).
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GraphEdgeTuple {
pub src: String,
pub dst: String,
pub label: String,
}
2 changes: 1 addition & 1 deletion nodedb-sql/src/ddl_ast/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ pub use alter_ops::{
};
pub use collection_type::build_collection_type;
pub use graph_parse::{FusionParams, parse_search_using_fusion};
pub use graph_types::{GraphDirection, GraphProperties};
pub use graph_types::{GraphDirection, GraphEdgeTuple, GraphProperties};
pub use nodedb_types::QuotaSpec;
pub use nodedb_types::{MirrorMode, MirrorStatus};
pub use parse::parse;
Expand Down
2 changes: 1 addition & 1 deletion nodedb-sql/src/ddl_ast/statement/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,4 +17,4 @@ pub use wrapper::*;
// Cross-module re-exports preserved at the `statement::` path for
// callers that imported these types via the pre-split surface.
pub use super::alter_ops::{AlterCollectionOp, AlterRoleOp, AlterUserOp};
pub use super::graph_types::{GraphDirection, GraphProperties};
pub use super::graph_types::{GraphDirection, GraphEdgeTuple, GraphProperties};
18 changes: 17 additions & 1 deletion nodedb-sql/src/ddl_ast/statement/types/graph.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

//! Graph DDL/DML statements.

use crate::ddl_ast::graph_types::{GraphDirection, GraphProperties};
use crate::ddl_ast::graph_types::{GraphDirection, GraphEdgeTuple, GraphProperties};

#[derive(Debug, Clone, PartialEq)]
pub enum GraphStmt {
Expand All @@ -20,6 +20,22 @@ pub enum GraphStmt {
dst: String,
label: String,
},
/// Batched edge insert: one statement carries many property-less
/// `(src, dst, label)` triples in one round trip.
///
/// The batch is deliberately property-less for now: the physical
/// `BatchEdge` carries no property object, so per-edge `PROPERTIES`
/// stays on the single-edge form until `BatchEdge` grows one
///
GraphInsertEdges {
collection: String,
edges: Vec<GraphEdgeTuple>,
},
/// Batched edge delete, the delete-side twin of [`GraphInsertEdges`].
GraphDeleteEdges {
collection: String,
edges: Vec<GraphEdgeTuple>,
},
GraphSetLabels {
node_id: String,
labels: Vec<String>,
Expand Down
28 changes: 28 additions & 0 deletions nodedb/src/bootstrap/background_loops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -278,6 +278,34 @@ pub fn spawn_background_loops(
// Quota period rollover is lazy — see `QuotaManager::rollover_if_due`.
// Every reader/writer of quota usage rolls the scope's period over on
// access, computed exactly from `period_start` and `period_secs`, so
// Per-database metric sampler (10-second interval). Fills the gauges whose
// registry setters have no owning writer: connections from the live
// session registry, and bridge queue depth from the dispatch WFQ. Families
// without a source yet (memory, storage, WAL commit latency, maintenance
// CPU) are left alone rather than published as constants.
{
let shared_metrics = 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;
}
crate::control::metrics::sampler::sample_once(&shared_metrics);
}
},
);
}

// there is no background sweep to spawn here and no interval to couple
// a quota's `period_secs` to.

Expand Down
62 changes: 49 additions & 13 deletions nodedb/src/bridge/dispatch/dispatcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -203,12 +203,15 @@ impl Dispatcher {
}

// Enqueue into the WFQ — returns Err if total capacity is full.
channel
.wfq
.try_enqueue(database_id, request)
.map_err(|_| crate::Error::Dispatch {
detail: format!("core {core_id}: total WFQ capacity exhausted"),
})?;
channel.wfq.try_enqueue(database_id, request).map_err(|_| {
note_capacity_exhausted();
crate::Error::Dispatch {
detail: format!(
"dispatch_capacity: core {core_id} capacity {}",
self.per_core_capacity
),
}
})?;

// Update per-DB pressure.
channel.update_db_pressure(database_id);
Expand Down Expand Up @@ -284,12 +287,15 @@ 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::Dispatch {
detail: format!("core {core_id}: total WFQ capacity exhausted"),
})?;
channel.wfq.try_enqueue(database_id, request).map_err(|_| {
note_capacity_exhausted();
crate::Error::Dispatch {
detail: format!(
"dispatch_capacity: core {core_id} capacity {}",
self.per_core_capacity
),
}
})?;

channel.update_db_pressure(database_id);
channel.flush_wfq();
Expand Down Expand Up @@ -325,6 +331,18 @@ impl Dispatcher {
.unwrap_or(0)
}

/// Queued requests for `database_id` across every core's WFQ.
///
/// Read-only and lock-light: metrics sampling asks once per database per
/// interval, and the per-queue depth is the raw number the fairness
/// thresholds are computed from.
pub fn db_queue_depth(&self, database_id: u64) -> u64 {
self.cores
.iter()
.map(|core| core.wfq.depth_for(database_id) as u64)
.sum()
}

/// Per-database pressure state for the given core (used by metrics exporters).
///
/// Returns `PressureState::Normal` when no pressure has been recorded for
Expand Down Expand Up @@ -445,6 +463,24 @@ impl Dispatcher {
}
}

/// Cumulative count of refused dispatches because a core's WFQ had no room.
///
/// Process-global on purpose: the WFQ rejection happens below the per-database
/// metrics handles, and the operational question ("is dispatch capacity the
/// thing clients are hitting?") is answered by the total. Per-database
/// granularity is a follow-up once the sampler owns a metrics handle.
static DISPATCH_CAPACITY_EXHAUSTED_TOTAL: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);

fn note_capacity_exhausted() {
DISPATCH_CAPACITY_EXHAUSTED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}

/// Total refused dispatches since process start (Prometheus export).
pub fn dispatch_capacity_exhausted_total() -> u64 {
DISPATCH_CAPACITY_EXHAUSTED_TOTAL.load(std::sync::atomic::Ordering::Relaxed)
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down Expand Up @@ -637,7 +673,7 @@ mod tests {
let _ = dispatcher.db_pressure_on_core(0, 2);
}

// --- Dead-core request loss (GitHub #265) ---
// --- Dead-core request loss ---
//
// When a Data Plane core's consumer/producer is dropped (the core thread
// died), `Dispatcher` must synthesize an error `Response` for every
Expand Down
1 change: 1 addition & 0 deletions nodedb/src/bridge/dispatch/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,5 +7,6 @@ mod drain;
pub use core_channel::{CoreChannel, CoreChannelDataSide};
pub use dispatcher::{
BridgeRequest, BridgeResponse, DatabasePriorityResolver, DefaultPriorityResolver, Dispatcher,
dispatch_capacity_exhausted_total,
};
pub use drain::CorePending;
Loading
Loading