Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,8 @@ fn producer_plan(rows: &[&Row]) -> Vec<u8> {
rows: msgpack_array(rows),
filters: Vec::new(),
projection: Vec::new(),
computed_columns: Vec::new(),
window_functions: Vec::new(),
sort_keys: Vec::new(),
limit: None,
offset: 0,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,8 @@ fn provider_scan_plan(rows: &[&Row]) -> Vec<u8> {
rows: msgpack_array(rows),
filters: Vec::new(),
projection: Vec::new(),
computed_columns: Vec::new(),
window_functions: Vec::new(),
sort_keys: Vec::new(),
limit: None,
offset: 0,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,8 @@ fn provider_scan_plan(rows: &[Vec<u8>]) -> Vec<u8> {
rows: msgpack_array(rows),
filters: Vec::new(),
projection: Vec::new(),
computed_columns: Vec::new(),
window_functions: Vec::new(),
sort_keys: Vec::new(),
limit: None,
offset: 0,
Expand Down
2 changes: 2 additions & 0 deletions nodedb-physical/src/physical_plan/collection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,8 @@ impl PhysicalPlan {
PhysicalPlan::Query(QueryOp::Exchange(op)) => op.child.collection(),
// PostProcess: recurse into the materialized input plan.
PhysicalPlan::Query(QueryOp::PostProcess { input, .. }) => input.collection(),
// SetOp merges N branches; no single collection names the node.
PhysicalPlan::Query(QueryOp::SetOp { .. }) => None,
// ProviderScan is a catalog/constant source — no user collection.
PhysicalPlan::Query(QueryOp::ProviderScan { .. }) => None,
// KV ops carry their own collection (sorted-index-only ops → None).
Expand Down
2 changes: 2 additions & 0 deletions nodedb-physical/src/physical_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ pub mod plan;
pub mod query;
pub mod rls_write_check_accessor;
pub mod routing;
pub mod set_op;
pub mod sort_key;
pub mod spatial;
pub mod streaming;
Expand Down Expand Up @@ -51,6 +52,7 @@ pub use meta::MetaOp;
pub use plan::PhysicalPlan;
pub use query::{AggregateSpec, GroupKeySpec, JoinProjection, QueryOp};
pub use routing::plan_contains_cluster_partitioned_leaf;
pub use set_op::SetOpKind;
pub use sort_key::SortKeySpec;
pub use spatial::{SpatialOp, SpatialPredicate};
pub use text::TextOp;
Expand Down
47 changes: 39 additions & 8 deletions nodedb-physical/src/physical_plan/query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,14 @@ pub enum QueryOp {
/// Output column names to keep. Empty = emit all columns.
#[serde(default)]
projection: Vec<String>,
/// Serialized `Vec<ComputedColumn>` (MessagePack), same encoding as
/// `DocumentOp::Scan::computed_columns`. Empty = none.
#[serde(default)]
computed_columns: Vec<u8>,
/// Serialized `Vec<WindowFuncSpec>` (MessagePack), same encoding as
/// `DocumentOp::Scan::window_functions`. Empty = none.
#[serde(default)]
window_functions: Vec<u8>,
/// ORDER BY terms, each an expression. Empty = unordered.
#[serde(default)]
sort_keys: Vec<crate::physical_plan::SortKeySpec>,
Expand Down Expand Up @@ -118,6 +126,14 @@ pub enum QueryOp {
/// Output column names to keep. Empty = emit all columns.
#[serde(default)]
projection: Vec<String>,
/// Serialized `Vec<ComputedColumn>` (MessagePack), same encoding as
/// `DocumentOp::Scan::computed_columns`. Empty = none.
#[serde(default)]
computed_columns: Vec<u8>,
/// Serialized `Vec<WindowFuncSpec>` (MessagePack), same encoding as
/// `DocumentOp::Scan::window_functions`. Empty = none.
#[serde(default)]
window_functions: Vec<u8>,
/// ORDER BY terms, each an expression. Empty = unordered.
#[serde(default)]
sort_keys: Vec<crate::physical_plan::SortKeySpec>,
Expand All @@ -133,18 +149,33 @@ pub enum QueryOp {
distinct: bool,
},

/// Set operation over N materialized children. Coordinator-only: the
/// resolver materializes every child and merges the rows into one
/// `ProviderScan` before dispatch. A Data-Plane core never sees this node.
///
/// Lowered from a derived-table body that is `UNION [ALL]`,
/// `INTERSECT [ALL]`, or `EXCEPT [ALL]`, so the body is one relation for
/// an outer [`QueryOp::PostProcess`] or input-sourced [`QueryOp::Aggregate`].
SetOp {
/// Child relations in SQL order. Each sharded child is wrapped in
/// `Exchange{Gather}` by the converter so its gather runs once.
inputs: Vec<crate::physical_plan::PhysicalPlan>,
/// Which set operation merges the inputs.
op: crate::physical_plan::SetOpKind,
},

/// Aggregate: GROUP BY + aggregate functions.
Aggregate {
collection: QualifiedCollection,
/// Optional sub-plan whose decoded rows are aggregated instead of
/// scanning `collection` per-shard. `Some` currently means EXACTLY a
/// catalog source (a `ProviderScan` lowered by the converter): the
/// aggregate runs over the coordinator-materialized catalog rows and is
/// therefore coordinator-local (never broadcast — see
/// `is_sharded_source`). `None` = legacy path: scan the named
/// `collection` on every shard. `collection` stays populated in both
/// cases so downstream RLS / permission / classification continue to
/// read it; the executor simply prefers `input` when present.
/// scanning `collection` per-shard. `Some` = an input-sourced
/// aggregate over a materialized relation: a catalog `ProviderScan`,
/// or any derived-table body the coordinator materializes into a
/// `ProviderScan` before dispatch. Coordinator-local, never broadcast
/// (see `is_sharded_source`). `None` = scan the named `collection` on
/// every shard. `collection` stays populated in both cases so
/// downstream RLS / permission / classification continue to read it;
/// the executor prefers `input` when present.
#[serde(default)]
input: Option<Box<crate::physical_plan::PhysicalPlan>>,
group_by: Vec<GroupKeySpec>,
Expand Down
8 changes: 8 additions & 0 deletions nodedb-physical/src/physical_plan/routing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,11 @@ pub fn plan_contains_cluster_partitioned_leaf(plan: &PhysicalPlan) -> bool {
plan_contains_cluster_partitioned_leaf(input)
}

// Recurse through every SetOp branch for the same reason.
PhysicalPlan::Query(QueryOp::SetOp { inputs, .. }) => {
inputs.iter().any(plan_contains_cluster_partitioned_leaf)
}

// Recurse through lateral outer plans.
PhysicalPlan::Query(QueryOp::LateralTopK { outer_plan, .. })
| PhysicalPlan::Query(QueryOp::LateralLoop { outer_plan, .. }) => {
Expand Down Expand Up @@ -159,6 +164,9 @@ impl PhysicalPlan {
right_bitmap,
)
}
// Coordinator-local: the resolver materializes every branch and
// merges on the coordinator, so the node itself is never fanned out.
PhysicalPlan::Query(QueryOp::SetOp { .. }) => false,
_ => self.is_sharded_source_leaf(),
}
}
Expand Down
33 changes: 33 additions & 0 deletions nodedb-physical/src/physical_plan/set_op.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
// SPDX-License-Identifier: Apache-2.0

//! Set-operation kinds for [`crate::physical_plan::QueryOp::SetOp`].

/// Which SQL set operation a [`crate::physical_plan::QueryOp::SetOp`] node
/// applies over its materialized inputs. Coordinator-resolved, never reaches
/// a Data-Plane core.
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
serde::Serialize,
serde::Deserialize,
zerompk::ToMessagePack,
zerompk::FromMessagePack,
)]
#[msgpack(c_enum)]
pub enum SetOpKind {
/// `UNION ALL`: concatenate every input in order.
UnionAll,
/// `UNION`: concatenate, then drop duplicate rows.
UnionDistinct,
/// `INTERSECT`: rows present in every input, deduplicated.
Intersect,
/// `INTERSECT ALL`: rows present in every input, bag semantics.
IntersectAll,
/// `EXCEPT`: rows of the first input absent from the rest, deduplicated.
Except,
/// `EXCEPT ALL`: rows of the first input absent from the rest, bag semantics.
ExceptAll,
}
3 changes: 2 additions & 1 deletion nodedb-physical/src/physical_plan/streaming.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,9 @@ impl PhysicalPlan {
sort_keys,
offset,
distinct,
window_functions,
..
}) => sort_keys.is_empty() && *offset == 0 && !*distinct,
}) => sort_keys.is_empty() && *offset == 0 && !*distinct && window_functions.is_empty(),

// Every other Document / Kv / Columnar / Timeseries op, plus all
// other engines and query ops, are not unordered-streamable.
Expand Down
40 changes: 35 additions & 5 deletions nodedb-query/src/window/eval.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,8 @@ use super::spec::WindowFuncSpec;
///
/// `rows` is the sorted result set. Each row is a `(doc_id, serde_json::Value)`.
/// The same rows are mutated in place with window columns appended to each
/// document.
/// document. The row array keeps its input order; each spec's partitions are
/// ordered by that spec's own ORDER BY, independent of the row array order.
///
/// Unknown window function names must be rejected by the planner before
/// reaching this dispatcher; an unrecognised name here is an internal bug
Expand All @@ -29,7 +30,7 @@ pub fn evaluate_window_functions(
specs: &[WindowFuncSpec],
) -> Result<(), crate::expr::EvalError> {
for spec in specs {
let partitions = build_partitions(rows, &spec.partition_by)?;
let partitions = build_partitions(rows, &spec.partition_by, &spec.order_by)?;

for partition_indices in &partitions {
match spec.func_name.as_str() {
Expand Down Expand Up @@ -146,9 +147,11 @@ 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));
// The frame runs in salary order within each dept, not in row
// arrival order: eng = Carol(90) → Alice(100) → Bob(120).
assert_eq!(rows[0].1["running_total"], json!(190.0));
assert_eq!(rows[1].1["running_total"], json!(310.0));
assert_eq!(rows[2].1["running_total"], json!(90.0));
assert_eq!(rows[3].1["running_total"], json!(80.0));
assert_eq!(rows[4].1["running_total"], json!(190.0));
}
Expand Down Expand Up @@ -262,6 +265,33 @@ mod tests {
assert_eq!(rows[4].1["nv"], json!(2));
}

#[test]
fn rank_orders_by_spec_order_by_not_row_arrival_order() {
// Rows arrive as Alice(100), Bob(120), Carol(90) within dept "eng" —
// not sorted by salary. RANK() OVER (ORDER BY salary DESC) must rank
// by salary, and the row array order must stay unchanged.
let mut rows = make_rows();
let spec = WindowFuncSpec {
alias: "rnk".into(),
func_name: "rank".into(),
args: vec![],
partition_by: vec![SqlExpr::Column("dept".into())],
order_by: vec![(SqlExpr::Column("salary".into()), false)],
frame: WindowFrame::default(),
};
evaluate_window_functions(&mut rows, &[spec]).unwrap();
assert_eq!(rows[0].1["name"], json!("Alice"));
assert_eq!(rows[1].1["name"], json!("Bob"));
assert_eq!(rows[2].1["name"], json!("Carol"));
assert_eq!(rows[3].1["name"], json!("Dave"));
assert_eq!(rows[4].1["name"], json!("Eve"));
assert_eq!(rows[0].1["rnk"], json!(2)); // Alice, salary 100
assert_eq!(rows[1].1["rnk"], json!(1)); // Bob, salary 120
assert_eq!(rows[2].1["rnk"], json!(3)); // Carol, salary 90
assert_eq!(rows[3].1["rnk"], json!(2)); // Dave, salary 80
assert_eq!(rows[4].1["rnk"], json!(1)); // Eve, salary 110
}

#[test]
#[should_panic(expected = "should have been rejected at planning time")]
fn unknown_function_panics_at_evaluator() {
Expand Down
115 changes: 96 additions & 19 deletions nodedb-query/src/window/helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,35 +6,112 @@ use std::collections::HashMap;

use crate::expr::types::SqlExpr;

/// Group row indices by partition key, preserving first-seen partition order.
/// Group row indices by partition key, preserving first-seen partition order,
/// then sort each partition's indices by the spec's ORDER BY.
///
/// A division/modulo-by-zero in a PARTITION BY expression propagates as
/// `Err(EvalError::DivisionByZero)` rather than being folded to NULL.
/// The returned index lists are ordered by `order_by`; the `rows` array
/// itself keeps its input order — only the per-partition index lists move.
///
/// A division/modulo-by-zero in a PARTITION BY or ORDER BY expression
/// propagates as `Err(EvalError::DivisionByZero)` rather than being folded to
/// NULL.
pub(super) fn build_partitions(
rows: &[(String, serde_json::Value)],
partition_by: &[SqlExpr],
order_by: &[(SqlExpr, bool)],
) -> Result<Vec<Vec<usize>>, crate::expr::EvalError> {
if partition_by.is_empty() {
return Ok(vec![(0..rows.len()).collect()]);
}
let mut partitions = if partition_by.is_empty() {
vec![(0..rows.len()).collect::<Vec<usize>>()]
} else {
let mut groups: HashMap<String, Vec<usize>> = HashMap::new();
let mut order = Vec::new();

for (i, (_id, doc)) in rows.iter().enumerate() {
let key: String = partition_by
.iter()
.map(|expr| eval_expr_on_json(expr, doc).map(|v| v.to_string()))
.collect::<Result<Vec<_>, _>>()?
.join("\x00");
let entry = groups.entry(key.clone()).or_default();
if entry.is_empty() {
order.push(key);
}
entry.push(i);
}

order.iter().filter_map(|k| groups.remove(k)).collect()
};

let mut groups: HashMap<String, Vec<usize>> = HashMap::new();
let mut order = Vec::new();
if !order_by.is_empty() {
let mut keys: Vec<Vec<serde_json::Value>> = Vec::with_capacity(rows.len());
for (_id, doc) in rows.iter() {
keys.push(
order_by
.iter()
.map(|(expr, _)| eval_expr_on_json(expr, doc))
.collect::<Result<Vec<_>, _>>()?,
);
}

for (i, (_id, doc)) in rows.iter().enumerate() {
let key: String = partition_by
.iter()
.map(|expr| eval_expr_on_json(expr, doc).map(|v| v.to_string()))
.collect::<Result<Vec<_>, _>>()?
.join("\x00");
let entry = groups.entry(key.clone()).or_default();
if entry.is_empty() {
order.push(key);
for partition in &mut partitions {
partition.sort_by(|&a, &b| compare_order_keys(&keys[a], &keys[b], order_by));
}
entry.push(i);
}

Ok(order.iter().filter_map(|k| groups.remove(k)).collect())
Ok(partitions)
}

/// Decide NULL placement for one ORDER BY column, shared by every window
/// evaluator's `compare_order_keys`.
///
/// NULL placement follows PostgreSQL's default: ASC places NULLs last, DESC
/// places NULLs first. A window spec carries no explicit NULLS FIRST/LAST
/// override, so this default is fixed by direction alone. Returns `None`
/// when neither value is NULL, leaving the non-null comparison to the
/// caller.
pub(super) fn null_order(
a_null: bool,
b_null: bool,
ascending: bool,
) -> Option<std::cmp::Ordering> {
use std::cmp::Ordering;
let nulls_first = !ascending;
match (a_null, b_null) {
(true, true) => Some(Ordering::Equal),
(true, false) => Some(if nulls_first {
Ordering::Less
} else {
Ordering::Greater
}),
(false, true) => Some(if nulls_first {
Ordering::Greater
} else {
Ordering::Less
}),
(false, false) => None,
}
}

/// Compare two rows' pre-evaluated ORDER BY keys.
fn compare_order_keys(
a: &[serde_json::Value],
b: &[serde_json::Value],
order_by: &[(SqlExpr, bool)],
) -> std::cmp::Ordering {
use std::cmp::Ordering;
for (idx, (_, ascending)) in order_by.iter().enumerate() {
let (Some(va), Some(vb)) = (a.get(idx), b.get(idx)) else {
continue;
};
let ord = null_order(va.is_null(), vb.is_null(), *ascending).unwrap_or_else(|| {
let c = crate::json_expr::compare_json(va, vb);
if *ascending { c } else { c.reverse() }
});
if ord != Ordering::Equal {
return ord;
}
}
Ordering::Equal
}

pub(super) fn set_window_col(row: &mut serde_json::Value, alias: &str, val: serde_json::Value) {
Expand Down
1 change: 1 addition & 0 deletions nodedb-query/src/window/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ pub mod running;
pub mod spec;
pub mod value_agg;
pub mod value_eval;
pub mod value_partition;

pub use eval::evaluate_window_functions;
pub use spec::{FrameBound, WindowFrame, WindowFuncSpec};
Expand Down
Loading
Loading