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()),

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.

[P2] Enforce the batch bound before copying the complete input

You copy every field into fields before checking the 1,000-edge limit. Oversized input consumes allocation and parsing work proportional to the full request before rejection. The tokenizer also materializes the complete token vector first.

Enforce limits while structurally parsing tuples, before collecting their owned strings. Keep a shared bound at the typed-plan admission boundary for callers that bypass SQL. Avoid duplicate full-input collections.

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)

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.

[P1] Preserve tuple structure before constructing edge identities

The tokenizer discards parentheses and commas. Grouping its flattened fields into threes cannot check the original tuples. VALUES ('a','b'), ('c','d','L','M') passes and creates (a,b,c) and (d,L,M). The parser also accepts missing delimiters and an unterminated final quoted value.

Correct the structural parser or tokenizer boundary. Check clause order, tuple arity, delimiters, quote closure, and trailing input before constructing edges. Reject malformed input before dispatch. Add negative cases for both insert and delete.

.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
17 changes: 17 additions & 0 deletions nodedb/src/control/planner/calvin/tx_class/static_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,23 @@ fn build_static_tx_class_impl(
homes.push(VShardId::from_key(dst_id.as_bytes()).as_u32());
continue;
}
// A batched edge write derives one participant home per edge, so a
// single batch task routes to every home its edges touch.
if let PhysicalPlan::Graph(
GraphOp::EdgePutBatch { edges } | GraphOp::EdgeDeleteBatch { edges },
) = &task.plan
{
for edge in edges {
edge_pairs
.entry(edge.collection.to_string())
.or_default()
.push((edge.src_surrogate.as_u32(), edge.dst_surrogate.as_u32()));
let homes = edge_homes.entry(edge.collection.to_string()).or_default();
homes.push(VShardId::from_key(edge.src_id.as_bytes()).as_u32());
homes.push(VShardId::from_key(edge.dst_id.as_bytes()).as_u32());
}
continue;
}
// KV and Vector writes carry their own key representation.
match &task.plan {
PhysicalPlan::Kv(op) => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use crate::control::state::SharedState;
use crate::types::DatabaseId;

use super::super::super::result::{DdlError, DdlResult};
use super::{algo, edge, rag_fusion, stats, traverse};
use super::{algo, edge, edge_batch, rag_fusion, stats, traverse};

/// Dispatch a parsed graph-overlay variant to its handler.
///
Expand Down Expand Up @@ -68,6 +68,14 @@ pub async fn dispatch_graph(
)
.await,
),
NodedbStatement::Graph(GraphStmt::GraphInsertEdges { collection, edges }) => Some(
edge_batch::insert_edges(state, identity, database_id, collection, edges, txn_ctx)
.await,
),
NodedbStatement::Graph(GraphStmt::GraphDeleteEdges { collection, edges }) => Some(
edge_batch::delete_edges(state, identity, database_id, collection, edges, txn_ctx)
.await,
),
NodedbStatement::Graph(GraphStmt::GraphSetLabels {
node_id,
labels,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,6 @@
//! Each function receives already-parsed typed fields; handlers never touch
//! `&str` parse paths.

use nodedb_sql::ddl_ast::GraphProperties;

use crate::bridge::envelope::PhysicalPlan;
use crate::control::planner::calvin::{build_static_tx_class, submit_calvin_routed};
use crate::control::security::identity::AuthenticatedIdentity;
Expand All @@ -16,8 +14,10 @@ use crate::control::server::shared::sql::staging_predicates::require_affected_co
use crate::control::server::surrogate_exchange::assign_surrogate_routed;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TraceId, VShardId};

use nodedb_physical::physical_plan::GraphOp;
use nodedb_physical::physical_task::{PhysicalTask, PostSetOp};
use nodedb_sql::ddl_ast::GraphProperties;

use super::super::super::result::{DdlError, DdlResult};
use super::edge_parse::{properties_to_json, validate_edge_label};
Expand Down
Loading
Loading