diff --git a/Cargo.lock b/Cargo.lock index fbf8b8688..c47be165a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6262,9 +6262,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.43" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "log", @@ -7020,7 +7020,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.4.3", + "getrandom 0.3.4", "once_cell", "rustix", "windows-sys 0.61.2", diff --git a/nodedb-cluster-tests/tests/common_suite/cases/shuffle_aggregate_cross_node.rs b/nodedb-cluster-tests/tests/common_suite/cases/shuffle_aggregate_cross_node.rs index d534d18ba..e56ec23b5 100644 --- a/nodedb-cluster-tests/tests/common_suite/cases/shuffle_aggregate_cross_node.rs +++ b/nodedb-cluster-tests/tests/common_suite/cases/shuffle_aggregate_cross_node.rs @@ -111,6 +111,8 @@ fn producer_plan(rows: &[&Row]) -> Vec { limit: None, offset: 0, distinct: false, + computed_columns: Vec::new(), + window_functions: Vec::new(), }); let plan = PhysicalPlan::Query(QueryOp::PartialAggregateState { collection: QualifiedCollection::new(nodedb_types::id::DatabaseId::DEFAULT, ""), diff --git a/nodedb-cluster-tests/tests/common_suite/cases/shuffle_consume_cross_node.rs b/nodedb-cluster-tests/tests/common_suite/cases/shuffle_consume_cross_node.rs index 205772766..4b41f71a0 100644 --- a/nodedb-cluster-tests/tests/common_suite/cases/shuffle_consume_cross_node.rs +++ b/nodedb-cluster-tests/tests/common_suite/cases/shuffle_consume_cross_node.rs @@ -108,6 +108,8 @@ fn provider_scan_plan(rows: &[&Row]) -> Vec { limit: None, offset: 0, distinct: false, + computed_columns: Vec::new(), + window_functions: Vec::new(), }); plan_wire::encode(&plan).expect("encode provider scan plan") } diff --git a/nodedb-cluster-tests/tests/common_suite/cases/shuffle_produce_cross_node.rs b/nodedb-cluster-tests/tests/common_suite/cases/shuffle_produce_cross_node.rs index 3f0602c28..ee3640bec 100644 --- a/nodedb-cluster-tests/tests/common_suite/cases/shuffle_produce_cross_node.rs +++ b/nodedb-cluster-tests/tests/common_suite/cases/shuffle_produce_cross_node.rs @@ -112,6 +112,8 @@ fn provider_scan_plan(rows: &[Vec]) -> Vec { limit: None, offset: 0, distinct: false, + computed_columns: Vec::new(), + window_functions: Vec::new(), }); plan_wire::encode(&plan).expect("encode provider scan plan") } diff --git a/nodedb-physical/src/physical_plan/kv/op.rs b/nodedb-physical/src/physical_plan/kv/op.rs index 1fe9126c7..17e01d800 100644 --- a/nodedb-physical/src/physical_plan/kv/op.rs +++ b/nodedb-physical/src/physical_plan/kv/op.rs @@ -152,6 +152,14 @@ pub enum KvOp { /// See `Get::surrogate_ceiling`; drops entries above the ceiling. #[serde(default)] surrogate_ceiling: Option, + /// Output column names (same format as DocumentOp::Scan). Empty = + /// return the whole row document. + #[serde(default)] + projection: Vec, + /// Serialized `Vec` applied per row after the scan + /// (same format as DocumentOp::Scan). Empty = none. + #[serde(default)] + computed_columns: Vec, }, /// Set or update TTL on an existing key. diff --git a/nodedb-physical/src/physical_plan/query.rs b/nodedb-physical/src/physical_plan/query.rs index 96f087d3c..5de8629f4 100644 --- a/nodedb-physical/src/physical_plan/query.rs +++ b/nodedb-physical/src/physical_plan/query.rs @@ -78,6 +78,18 @@ pub enum QueryOp { /// Output column names to keep. Empty = emit all columns. #[serde(default)] projection: Vec, + /// Serialized `Vec` applied per row after + /// projection-name extraction (same wire format as the engine + /// scans). Expression projections over materialized rows (derived + /// tables, constant subqueries) need these — name-only projection + /// would silently drop them. + #[serde(default)] + computed_columns: Vec, + /// Serialized `Vec` evaluated per partition after + /// computed columns (window over derived-table rows — issue #295 + /// Gap 3). Empty = no window functions. + #[serde(default)] + window_functions: Vec, /// ORDER BY terms, each an expression. Empty = unordered. #[serde(default)] sort_keys: Vec, @@ -118,6 +130,14 @@ pub enum QueryOp { /// Output column names to keep. Empty = emit all columns. #[serde(default)] projection: Vec, + /// Serialized `Vec` applied per row (see + /// `ProviderScan::computed_columns`). + #[serde(default)] + computed_columns: Vec, + /// Serialized `Vec` (see + /// `ProviderScan::window_functions`). + #[serde(default)] + window_functions: Vec, /// ORDER BY terms, each an expression. Empty = unordered. #[serde(default)] sort_keys: Vec, diff --git a/nodedb-query/src/window/eval.rs b/nodedb-query/src/window/eval.rs index cf6085ad6..b97791e6d 100644 --- a/nodedb-query/src/window/eval.rs +++ b/nodedb-query/src/window/eval.rs @@ -32,6 +32,13 @@ pub fn evaluate_window_functions( let partitions = build_partitions(rows, &spec.partition_by)?; for partition_indices in &partitions { + // Evaluate the partition in the window's own ORDER BY — the + // query's sort, when present, orders by different keys. Without + // this, row_number numbered by arrival order and the rank family + // compared peers in the wrong sequence. + let ordered = + super::helpers::ordered_partition_indices(rows, partition_indices, &spec.order_by)?; + let partition_indices = &ordered; match spec.func_name.as_str() { "row_number" => apply_row_number(rows, partition_indices, &spec.alias), "rank" => apply_rank(rows, partition_indices, &spec.alias, &spec.order_by)?, @@ -134,6 +141,84 @@ mod tests { assert_eq!(rows[4].1["rn"], json!(2)); } + /// `row_number() OVER (ORDER BY ...)` must number by the window's own + /// ordering, not by arrival order: arrival is salary 100, 120, 90; ORDER + /// BY salary DESC must give 120 → 1, 100 → 2, 90 → 3. + #[test] + fn row_number_follows_the_window_ordering() { + let mut rows = make_rows(); + let spec = WindowFuncSpec { + alias: "rn".into(), + func_name: "row_number".into(), + args: vec![], + partition_by: vec![], + order_by: vec![(SqlExpr::Column("salary".into()), false)], + frame: WindowFrame::default(), + }; + evaluate_window_functions(&mut rows, &[spec]).unwrap(); + assert_eq!(rows[1].1["rn"], json!(1), "120 is the highest salary"); + assert_eq!(rows[4].1["rn"], json!(2), "110 is second"); + assert_eq!(rows[0].1["rn"], json!(3), "100 is third"); + assert_eq!(rows[2].1["rn"], json!(4), "90 is fourth"); + assert_eq!(rows[3].1["rn"], json!(5), "80 is fifth"); + } + + /// An ORDER BY key that divides by zero fails the query instead of being + /// skipped along with the ordering. + #[test] + fn row_number_ordering_error_propagates() { + let mut rows = make_rows(); + let spec = WindowFuncSpec { + alias: "rn".into(), + func_name: "row_number".into(), + args: vec![], + partition_by: vec![], + order_by: vec![( + SqlExpr::BinaryOp { + left: Box::new(SqlExpr::Column("salary".into())), + op: crate::expr::BinaryOp::Div, + right: Box::new(SqlExpr::Literal(nodedb_types::Value::Integer(0))), + }, + true, + )], + frame: WindowFrame::default(), + }; + let err = evaluate_window_functions(&mut rows, &[spec]) + .expect_err("ordering by a division by zero must fail"); + assert!( + matches!(err, crate::expr::EvalError::DivisionByZero), + "unexpected error: {err:?}" + ); + } + + /// A singleton partition still evaluates its ordering key: the result + /// cannot change, but the error must be raised like anywhere else. + #[test] + fn singleton_partition_ordering_error_propagates() { + let mut rows = vec![("1".to_string(), json!({"n": 1}))]; + let spec = WindowFuncSpec { + alias: "rn".into(), + func_name: "row_number".into(), + args: vec![], + partition_by: vec![], + order_by: vec![( + SqlExpr::BinaryOp { + left: Box::new(SqlExpr::Column("n".into())), + op: crate::expr::BinaryOp::Div, + right: Box::new(SqlExpr::Literal(nodedb_types::Value::Integer(0))), + }, + true, + )], + frame: WindowFrame::default(), + }; + let err = evaluate_window_functions(&mut rows, &[spec]) + .expect_err("a division by zero in the ordering key must fail"); + assert!( + matches!(err, crate::expr::EvalError::DivisionByZero), + "unexpected error: {err:?}" + ); + } + #[test] fn running_sum() { let mut rows = make_rows(); @@ -146,11 +231,14 @@ mod tests { frame: WindowFrame::default(), }; evaluate_window_functions(&mut rows, &[spec]).unwrap(); - assert_eq!(rows[0].1["running_total"], json!(100.0)); - assert_eq!(rows[1].1["running_total"], json!(220.0)); - assert_eq!(rows[2].1["running_total"], json!(310.0)); - assert_eq!(rows[3].1["running_total"], json!(80.0)); - assert_eq!(rows[4].1["running_total"], json!(190.0)); + // Partition `eng` in salary order is 90, 100, 120 and `sales` is 80, + // 110. Arrival order (100, 120, 90 / 80, 110) is NOT the window's + // order, so the running total must follow the sorted sequence. + assert_eq!(rows[2].1["running_total"], json!(90.0), "eng 90 first"); + assert_eq!(rows[0].1["running_total"], json!(190.0), "eng +100"); + assert_eq!(rows[1].1["running_total"], json!(310.0), "eng +120"); + assert_eq!(rows[3].1["running_total"], json!(80.0), "sales 80 first"); + assert_eq!(rows[4].1["running_total"], json!(190.0), "sales +110"); } #[test] diff --git a/nodedb-query/src/window/helpers.rs b/nodedb-query/src/window/helpers.rs index 3bd7b36fb..cf61b1e98 100644 --- a/nodedb-query/src/window/helpers.rs +++ b/nodedb-query/src/window/helpers.rs @@ -82,3 +82,91 @@ pub(super) fn order_keys_equal( } Ok(true) } + +/// Indices of one partition, ordered by a window spec's ORDER BY keys. +/// +/// Window functions receive the sorted result set, but a window's own ORDER +/// BY is a different ordering from the query's — nothing upstream sorts by +/// it. Evaluating a partition in arrival order silently assigned +/// row numbers by arrival and compared peers in the wrong sequence, so every +/// ranking and running frame was wrong whenever the two orders differed. +/// +/// Keys are evaluated once per row; a division/modulo-by-zero propagates as +/// `Err(EvalError::DivisionByZero)`. The sort is stable, so rows with equal +/// keys keep their arrival order. +pub(super) fn ordered_partition_indices( + rows: &[(String, serde_json::Value)], + indices: &[usize], + order_by: &[(SqlExpr, bool)], +) -> Result, crate::expr::EvalError> { + if order_by.is_empty() { + return Ok(indices.to_vec()); + } + // Keys are evaluated even for a one-row partition: the ordering cannot + // change the result, but a key that divides by zero is still a statement + // error in PostgreSQL, and skipping the evaluation would answer where + // the same expression anywhere else fails. + let mut keyed: Vec<(usize, Vec)> = Vec::with_capacity(indices.len()); + for &i in indices { + let mut keys = Vec::with_capacity(order_by.len()); + for (expr, _) in order_by { + keys.push(eval_expr_on_json(expr, &rows[i].1)?); + } + keyed.push((i, keys)); + } + keyed.sort_by(|(_, a), (_, b)| compare_key_values(a, b, order_by)); + Ok(keyed.into_iter().map(|(i, _)| i).collect()) +} + +/// Compare two pre-evaluated ORDER BY key vectors under their direction flags. +fn compare_key_values( + a: &[serde_json::Value], + b: &[serde_json::Value], + order_by: &[(SqlExpr, bool)], +) -> std::cmp::Ordering { + for ((va, vb), (_expr, ascending)) in a.iter().zip(b.iter()).zip(order_by.iter()) { + let ord = compare_with_direction(va, vb, *ascending); + if ord != std::cmp::Ordering::Equal { + return ord; + } + } + std::cmp::Ordering::Equal +} + +/// Compare two key values, applying the window ORDER BY null convention: +/// NULLS LAST under ASC, NULLS FIRST under DESC. +fn compare_with_direction( + a: &serde_json::Value, + b: &serde_json::Value, + ascending: bool, +) -> std::cmp::Ordering { + use serde_json::Value as J; + use std::cmp::Ordering; + + let base = match (a, b) { + (J::Null, J::Null) => Ordering::Equal, + (J::Null, _) => Ordering::Greater, + (_, J::Null) => Ordering::Less, + _ => json_rank(a).cmp(&json_rank(b)).then_with(|| { + if let (Some(x), Some(y)) = (as_f64(a), as_f64(b)) { + x.partial_cmp(&y).unwrap_or(Ordering::Equal) + } else { + a.to_string().cmp(&b.to_string()) + } + }), + }; + if ascending { base } else { base.reverse() } +} + +/// Rank JSON kinds so mixed-type keys still compare deterministically. +fn json_rank(v: &serde_json::Value) -> u8 { + use serde_json::Value as J; + match v { + J::Null => 0, + J::Bool(_) => 1, + J::Number(_) => 2, + J::String(_) => 3, + J::Array(_) => 4, + J::Object(_) => 5, + } +} diff --git a/nodedb-sql/src/ddl_ast/graph_parse/cursor.rs b/nodedb-sql/src/ddl_ast/graph_parse/cursor.rs new file mode 100644 index 000000000..2fac116f9 --- /dev/null +++ b/nodedb-sql/src/ddl_ast/graph_parse/cursor.rs @@ -0,0 +1,285 @@ +// SPDX-License-Identifier: Apache-2.0 + +//! Consume-tracking cursor over a graph DSL token list. +//! +//! The module's parse is seek-based: a clause reader finds its keyword +//! wherever it appears and ignores every token between. That is fine while +//! every token belongs to some clause. It is wrong the moment a token belongs +//! to none: `GRAPH TRAVERSE FROM 1 DEPTS 3 IN g` finds no `DEPTH`, defaults +//! the depth, and leaves `DEPTS 3` unread — the statement answers a different +//! question than it asked, with no error. +//! +//! The cursor records which tokens each clause consumed. After the statement +//! is built, [`Cursor::finish`] refuses the first token no clause claimed. +//! A typo is then a parse error that names the token, and a new clause cannot +//! be added without deciding what it consumes. + +use super::tokenizer::Tok; +use crate::error::SqlError; + +/// A token list plus the set of tokens a clause has claimed. +pub(super) struct Cursor<'a> { + toks: Vec>, + used: Vec, +} + +impl<'a> Cursor<'a> { + /// Build a cursor over `toks`. The first `prefix_len` tokens are the + /// command words (`GRAPH TRAVERSE`), which the dispatcher matched and no + /// clause will claim. + pub(super) fn new(toks: Vec>, prefix_len: usize) -> Self { + let mut used = vec![false; toks.len()]; + for slot in used.iter_mut().take(prefix_len.min(toks.len())) { + *slot = true; + } + Self { toks, used } + } + + fn is_keyword(tok: &Tok<'_>, keyword: &str) -> bool { + matches!(tok, Tok::Word(w) if w.eq_ignore_ascii_case(keyword)) + } + + /// Position of `keyword`, claiming it. A keyword already claimed by an + /// earlier clause still matches: `IN` in one statement is one clause, but + /// the seek-based readers may be called in any order. + fn find(&mut self, keyword: &str) -> Option { + let pos = self + .toks + .iter() + .position(|tok| Self::is_keyword(tok, keyword))?; + self.used[pos] = true; + Some(pos) + } + + /// Claim the value token at `pos` when it is a word or a quoted literal. + /// An object literal is claimed by the callers that accept one. + fn claim_text(&mut self, pos: usize) -> Option { + let value = match self.toks.get(pos)? { + Tok::Quoted(s) => s.clone().into_owned(), + Tok::Word(w) => (*w).to_string(), + Tok::Object(_) => return None, + }; + self.used[pos] = true; + Some(value) + } + + /// The word or quoted literal after `keyword`. + pub(super) fn quoted_after(&mut self, keyword: &str) -> Option { + let pos = self.find(keyword)?; + self.claim_text(pos + 1) + } + + /// Every consecutive word or quoted literal after `keyword`, up to the + /// first token that is neither. Used by `AS