Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
617b6e2
test(sql): cover computed columns and errors over derived tables
farhan-syah Sep 17, 2026
8180eb2
refactor(planner): split sql_plan_convert/expr.rs by concern
farhan-syah Sep 17, 2026
f6bc126
feat(query): evaluate window functions and computed columns in Provid…
farhan-syah Sep 17, 2026
1699789
fix(query): order window partitions by each spec's own ORDER BY
farhan-syah Sep 17, 2026
2cbe9bf
feat(query): evaluate window functions over subquery and CTE tails
farhan-syah Sep 17, 2026
0916580
feat(query): aggregate over derived-table and union bodies
farhan-syah Sep 17, 2026
a1a3199
feat(query): evaluate query tails over set-operation bodies
farhan-syah Sep 17, 2026
c22ed89
feat(query): evaluate a SELECT-list sequence accessor per output row
farhan-syah Sep 17, 2026
39077bd
feat(query): evaluate projections and computed columns over kv scans
farhan-syah Sep 18, 2026
744a39a
feat(query): stamp Control-Plane computed columns across response sha…
farhan-syah Sep 18, 2026
759f31d
feat(sql): resolve RETURNING items as full expressions, not just columns
farhan-syah Sep 18, 2026
2a75594
feat(query): support RETURNING on KV UPDATE and DELETE
farhan-syah Sep 18, 2026
7f093f0
feat(kv): make SQL UPDATE a no-op against an absent key
farhan-syah Sep 18, 2026
7193e22
fix(resp): surface rejected KV writes as errors, not fallback replies
farhan-syah Sep 18, 2026
5168e18
docs(sql): document sequence DDL, functions, and evaluation rules
farhan-syah Sep 18, 2026
041b348
Merge remote-tracking branch 'origin/main' into feat/sequence-row-scope
farhan-syah Sep 18, 2026
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
65 changes: 63 additions & 2 deletions docs/query-language.md
Original file line number Diff line number Diff line change
Expand Up @@ -239,19 +239,80 @@ TRUNCATE users;

### RETURNING

Both `UPDATE` and `DELETE` support a `RETURNING` clause to read back affected rows in the same statement:
`INSERT`, `UPSERT`, `UPDATE`, `DELETE`, and `MERGE` accept a `RETURNING` clause that reads back the affected rows in the same statement:

```sql
-- INSERT RETURNING: returns the stored row
INSERT INTO users (id, name) VALUES ('u2', 'Bo') RETURNING id, name;

-- UPDATE RETURNING: returns the post-update image
UPDATE users SET role = 'admin' WHERE id = 'u1' RETURNING id, role;
UPDATE orders SET status = 'shipped' WHERE id = 'o1' RETURNING *;

-- DELETE RETURNING: returns the pre-delete image
DELETE FROM users WHERE id = 'u1' RETURNING id, name;
DELETE FROM orders WHERE status = 'cancelled' RETURNING *;

-- Expressions: evaluated per returned row against the stored image
UPDATE orders SET qty = qty + 1 WHERE id = 'o1' RETURNING id, qty * price AS total;
INSERT INTO events (id, kind) VALUES ('e1', 'click') RETURNING id, nextval('event_seq') AS n;
```

`RETURNING *` expands to all columns. Named columns are returned as bare values — arithmetic expressions in `RETURNING` are not supported. Works in both simple-query and extended-query (prepared statement) protocols.
`RETURNING *` expands to all columns. Every item is a scalar expression over the target collection: a bare column, a column under an alias (`col AS name`), arithmetic, a function call, or a sequence accessor (`nextval`, `currval`, `setval`). The Data Plane returns the base columns an expression reads, and the Control Plane evaluates the expression once per returned row. A sequence accessor advances once per row, in row order. The clause works in both the simple-query and extended-query (prepared statement) protocols, and `Describe` announces an expression under its alias.

### Sequences

A sequence is a named `bigint` counter, independent of any collection.

```sql
CREATE SEQUENCE event_seq START WITH 1 INCREMENT BY 1 MINVALUE 1 CYCLE CACHE 20;
DROP SEQUENCE event_seq;
DROP SEQUENCE IF EXISTS event_seq;
SHOW SEQUENCES;
DESCRIBE SEQUENCE event_seq;
ALTER SEQUENCE event_seq RESTART WITH 100;
```

`CREATE SEQUENCE [IF NOT EXISTS] <name>` accepts these options, in any order:

| Option | Effect |
|---|---|
| `START [WITH] n` | First value `nextval` returns |
| `INCREMENT [BY] n` | Step between successive values |
| `MINVALUE n` | Lower bound |
| `MAXVALUE n` | Upper bound |
| `CYCLE` / `NO CYCLE` | Wrap to the bound instead of erroring at exhaustion |
| `CACHE n` | Values a node pre-allocates per round-trip to the registry |
| `FORMAT 'template'` | Render template applied to the numeric value |
| `RESET period` | Period after which the counter restarts |
| `GAP_FREE` | Accepted and stored; allocation is not yet serialized or rolled back per transaction |
| `SCOPE name` | Named allocation scope |

`DROP SEQUENCE [IF EXISTS] <name>` removes it. `SHOW SEQUENCES` lists every sequence. `DESCRIBE SEQUENCE <name>` reports one sequence's current state. `ALTER SEQUENCE <name> RESTART [WITH n]` and `ALTER SEQUENCE <name> FORMAT '<template>'` change it in place.

Three functions read and move a sequence:

| Function | Effect |
|---|---|
| `nextval('s')` | Advances `s` and returns the new value; records it as this session's `currval` |
| `currval('s')` | The last value THIS session obtained from `nextval('s')` |
| `setval('s', n)` | Positions `s` so the next `nextval` returns `n + increment` |

`currval` before this session ever called `nextval` on that sequence fails with SQLSTATE `55000` (`object_not_in_prerequisite_state`). Any accessor naming an unknown sequence fails with SQLSTATE `42704` (`undefined_object`).

A sequence accessor is allowed only where the plan controls how many times it runs:

| Context | Evaluates |
|---|---|
| FROM-less `SELECT nextval('s')` | Once, at plan time |
| `VALUES` list | Once, at plan time |
| Column `DEFAULT nextval('s')` | Once per inserted row |
| SELECT list of a top-level SELECT over a relation | Once per output row, in output order, after `ORDER BY`/`LIMIT` — on the Control Plane, after the Data Plane returns the rows. `LIMIT n` consumes exactly `n` values |
| `RETURNING` | Once per returned row |

`WHERE`, `ORDER BY`, `GROUP BY`, `HAVING`, `JOIN ... ON`, `UPDATE ... SET`, an aggregate or window argument, an `INSERT ... SELECT` source, and any nested subquery all refuse a sequence accessor with SQLSTATE `0A000` (`feature_not_supported`) — those clauses run on the per-row evaluator, which holds no sequence state.

A plan whose SELECT list holds a sequence accessor is never cached: each execution must call the accessor again, not replay a cached row. `EXPLAIN` and `Describe` plan the statement without executing it, so neither advances a sequence.

## DDL

Expand Down
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
8 changes: 5 additions & 3 deletions nodedb-physical/src/physical_plan/document/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -237,10 +237,12 @@ pub enum ReturningColumns {
Named(Vec<ReturningItem>),
}

/// Parsed representation of a RETURNING clause carried through the bridge.
/// The Data-Plane projection of a RETURNING clause carried through the bridge.
///
/// Produced by the Control Plane's `strip_returning()` and injected into
/// `PointUpdate`, `BulkUpdate`, `PointDelete`, and `BulkDelete` variants
/// Derived on the Control Plane from the resolved clause: every stored column
/// the clause names or an expression in it reads, by bare name. The Control
/// Plane evaluates expressions and applies display names after the rows
/// return. Injected into the DML plan variants that carry a `returning` slot
/// before crossing the SPSC bridge.
#[derive(
Debug,
Expand Down
42 changes: 41 additions & 1 deletion nodedb-physical/src/physical_plan/kv/op.rs
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,13 @@ pub enum KvOp {
/// predicate is attached. Only `RlsWriteCheck::Predicate` makes the
/// handler read the pre-image at all.
rls_write_check: RlsWriteCheck,
/// When `Some`, return the STORED pre-image of every removed row
/// (row as `SELECT` showed it, `key` included).
#[serde(default)]
returning: Option<ReturningSpec>,
/// See `Put::rls_filters`.
#[serde(default)]
rls_filters: Vec<u8>,
},

/// Cursor-based scan with optional filter predicate.
Expand All @@ -143,6 +150,13 @@ pub enum KvOp {
count: usize,
/// Optional filter predicates (same format as DocumentScan filters).
filters: Vec<u8>,
/// Output column names to keep. Empty = emit the full row.
#[serde(default)]
projection: Vec<String>,
/// Serialized `Vec<ComputedColumn>` (MessagePack), same encoding as
/// `DocumentOp::Scan::computed_columns`. Empty = none.
#[serde(default)]
computed_columns: Vec<u8>,
/// Optional glob pattern for key matching (e.g., "user:*").
match_pattern: Option<String>,
/// ORDER BY terms, each an expression, applied to the scan result
Expand Down Expand Up @@ -260,14 +274,26 @@ pub enum KvOp {
FieldSet {
collection: QualifiedCollection,
key: Vec<u8>,
/// Field name → new value (JSON-encoded bytes).
/// Field name → new value (msgpack-encoded bytes; empty = NULL).
updates: Vec<(String, Vec<u8>)>,
/// Content-addressed identity on `(collection, key)`, threaded to the
/// write-back so a field merge keeps the row's original surrogate.
surrogate: Surrogate,
/// Update only an existing row. `true` for SQL UPDATE (an absent key
/// is a no-op); `false` for the RESP hash-set family, which creates
/// the row.
#[serde(default)]
if_present: bool,
/// Write policy evaluated against the merged body, which exists only
/// after the stored row is read and updates applied.
rls_write_check: RlsWriteCheck,
/// When `Some`, return the STORED post-image (merged row as `SELECT`
/// shows it, `key` included). Never the caller's submitted updates.
#[serde(default)]
returning: Option<ReturningSpec>,
/// See `Put::rls_filters`.
#[serde(default)]
rls_filters: Vec<u8>,
},

/// Truncate: delete ALL entries in a KV collection.
Expand Down Expand Up @@ -490,6 +516,13 @@ pub enum KvOp {
/// matched row's post-image once the assignments have been applied, or
/// the reason no predicate is attached.
rls_write_check: RlsWriteCheck,
/// When `Some`, return one row per matched key — the STORED
/// post-image of each, in scan order — projected per spec.
#[serde(default)]
returning: Option<ReturningSpec>,
/// See `Put::rls_filters`.
#[serde(default)]
rls_filters: Vec<u8>,
},

/// Delete every row matching `filters`.
Expand All @@ -503,5 +536,12 @@ pub enum KvOp {
/// Compiled row-level-security WRITE predicate, evaluated against the
/// pre-image of every row this removes.
rls_write_check: RlsWriteCheck,
/// When `Some`, return one row per removed key — the STORED
/// pre-image of each, in scan order — projected per spec.
#[serde(default)]
returning: Option<ReturningSpec>,
/// See `Put::rls_filters`.
#[serde(default)]
rls_filters: Vec<u8>,
},
}
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
Loading
Loading