Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
3899eb1
fix(resp): report row counts on KV, CRDT, and vector writes
farhan-syah Sep 18, 2026
9e55c39
fix(resp): classify every client write as count-bearing on pgwire
farhan-syah Sep 18, 2026
c052b21
refactor(cluster): split Calvin scheduler dispatch by concern
farhan-syah Sep 19, 2026
027c196
fix(resp): report row counts on graph edge and label writes
farhan-syah Sep 19, 2026
8c4b4f7
fix(resp): fold multi-task statements into one command tag
farhan-syah Sep 19, 2026
94a888e
feat(array): stage Put/Delete writes through the transaction overlay
farhan-syah Sep 19, 2026
c253d91
feat(array): read-your-own-writes for distributed array reads/writes
farhan-syah Sep 19, 2026
5bc7aa3
fix(array): report real affected counts for distributed array writes
farhan-syah Sep 19, 2026
99de7f5
fix(vector): make vector-primary DML correct end to end
farhan-syah Sep 19, 2026
e4a130e
feat(vector): stage vector-primary writes through the transaction ove…
farhan-syah Sep 19, 2026
74370fb
feat(crdt): stage CRDT document writes through the transaction overlay
farhan-syah Sep 19, 2026
59131fe
feat(txn): stage TRUNCATE through the transaction overlay
farhan-syah Sep 19, 2026
994a758
fix(kv): stamp KV writes with a per-table epoch for cache invalidation
farhan-syah Sep 19, 2026
c9de634
fix(graph): keep GraphRAG fusion alive on an empty vector leg
farhan-syah Sep 19, 2026
24218b5
feat(sql): route TRUNCATE per engine through EngineRules
farhan-syah Sep 19, 2026
366056f
refactor(wal): split WAL append methods by record family
farhan-syah Sep 19, 2026
459376d
fix(columnar): keep flushed tombstones and rows in their own segment
farhan-syah Sep 19, 2026
4bd9be5
fix(spatial): cascade R-tree maintenance through columnar DML and undo
farhan-syah Sep 19, 2026
7c38909
feat(sql): route columnar and timeseries TRUNCATE through their engines
farhan-syah Sep 19, 2026
4909f6c
fix(native): report real DML counts and the command verb
farhan-syah Sep 20, 2026
ca3ca4e
feat(kv): stage predicate UPDATE/DELETE through the transaction overlay
farhan-syah Sep 20, 2026
6d9782a
refactor(tests): share the RESP client between crash and wire harnesses
farhan-syah Sep 20, 2026
5138f16
fix(kv): keep a bare-value KV row raw through every read-modify-write
farhan-syah Sep 20, 2026
20e1adb
refactor(tests): share the raw pgwire client across test harnesses
farhan-syah Sep 20, 2026
4df5001
fix(resp): stage and fold gateway-forwarded writes like local ones
farhan-syah Sep 20, 2026
00f2fb0
feat(sql): replicate MERGE, UPDATE FROM, and INSERT SELECT resolved w…
farhan-syah Sep 20, 2026
7fec24c
test(cluster): cover multi-shard tag fold, orchestrated writes, and p…
farhan-syah Sep 20, 2026
d5d3bb7
refactor(routing): extract plan metering info once per gateway task
farhan-syah Sep 20, 2026
b891fc8
fix(ci): update allowed dispatch path for cluster array module move
farhan-syah Sep 20, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
2 changes: 2 additions & 0 deletions nodedb-client/src/native/connection/response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ pub(super) fn response_to_query_result(resp: NativeResponse) -> NodeDbResult<Que
columns: resp.columns.unwrap_or_default(),
rows: resp.rows.unwrap_or_default(),
rows_affected: resp.rows_affected.unwrap_or(0),
command: resp.command,
})
}

Expand All @@ -69,6 +70,7 @@ mod tests {
columns: vec!["x".into()],
rows: vec![vec![nodedb_types::Value::Integer(42)]],
rows_affected: 0,
command: None,
},
0,
);
Expand Down
1 change: 1 addition & 0 deletions nodedb-client/src/remote/client/sql_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ impl NodeDbRemote {
columns,
rows,
rows_affected: 0,
command: None,
})
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,10 @@ use crate::common;
use std::sync::atomic::Ordering;
use std::time::Duration;

use nodedb::types::{DatabaseId, VShardId};
use nodedb_cluster::calvin::SEQUENCER_GROUP_ID;
use tokio_postgres::SimpleQueryMessage;

use super::vshard_names::distinct_vshard_collections;
use common::cluster_harness::{TestClusterNode, wait_for};

/// Observed sequencer-group leader id from a node's local Raft status, or `0`
Expand Down Expand Up @@ -61,23 +61,6 @@ async fn value_of(client: &tokio_postgres::Client, coll: &str, id: &str) -> Opti
})
}

/// Find two collection names whose vShard ids differ.
fn two_distinct_vshard_collections() -> (String, String) {
let mut first: Option<(String, u32)> = None;
for i in 0u32..512 {
let name = format!("calvin_e2e_{i}");
let vshard = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &name).as_u32();
if let Some((ref fname, fv)) = first {
if fv != vshard {
return (fname.clone(), name);
}
} else {
first = Some((name, vshard));
}
}
panic!("could not find two distinct-vshard collections in 512 tries");
}

/// Calvin multi-shard batch via pgwire `simple_query` COMMITS when sent as an
/// interactive `BEGIN ... COMMIT` block.
///
Expand Down Expand Up @@ -124,7 +107,7 @@ async fn calvin_multishard_write_in_explicit_block_commits() {
)
.await;

let (col_a, col_b) = two_distinct_vshard_collections();
let (col_a, col_b) = distinct_vshard_collections("calvin_e2e_0", "calvin_e2e");

// Create both collections on this (single) node.
node.client
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,7 @@
//!
//! 1. Two `document_schemaless` collections, each created `WITH
//! (bitemporal=true)`, are placed on DIFFERENT vShards
//! (`distinct_vshard_bitemporal_collections`, same technique as
//! `calvin_multi_shard_redo_restart.rs`).
//! (`vshard_names::distinct_vshard_collections`).
//! 2. `BEGIN; INSERT INTO <a>; INSERT INTO <b>; COMMIT` is sent as ONE
//! `simple_query` call, so both writes are buffered inside the block and,
//! on COMMIT, `classify_dispatch` sees writes on two vShards → MultiShard.
Expand All @@ -32,9 +31,9 @@ use crate::common;
use std::sync::atomic::Ordering;
use std::time::Duration;

use nodedb::types::{DatabaseId, VShardId};
use tokio_postgres::SimpleQueryMessage;

use super::vshard_names::distinct_vshard_collections;
use common::cluster_harness::{TestClusterNode, read_once_a_leader_exists, wait_for};

/// Observed sequencer-group leader id from a node's local Raft status, or `0`
Expand All @@ -61,26 +60,6 @@ fn admitted_total(node: &TestClusterNode) -> u64 {
.unwrap_or(0)
}

/// A `(coll_a, coll_b)` pair of bitemporal collection names whose vShard ids
/// differ, so a transaction writing to both is genuinely multi-shard.
/// Deterministic: `VShardId::from_collection_in_database` is a pure function
/// of the database id + collection name bytes. Same technique as
/// `calvin_multi_shard_redo_restart.rs::distinct_vshard_collections`.
fn distinct_vshard_bitemporal_collections() -> (String, String) {
let a_name = "bt_a".to_string();
let va = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &a_name).as_u32();
for i in 0u32..512 {
let b_name = format!("bt_b_{i}");
if VShardId::from_collection_in_database(DatabaseId::DEFAULT, &b_name).as_u32() != va {
return (a_name, b_name);
}
}
panic!(
"could not find a second bitemporal collection name on a distinct vShard \
from the first in 512 tries"
);
}

/// Current `value` for `id` in a bitemporal document collection, or `None` if
/// not visible.
///
Expand Down Expand Up @@ -160,7 +139,7 @@ async fn calvin_multi_shard_bitemporal_best_effort_commit_survives_wal_only_rest
)
.await;

let (coll_a, coll_b) = distinct_vshard_bitemporal_collections();
let (coll_a, coll_b) = distinct_vshard_collections("bt_a", "bt_b");

node.client
.simple_query(&format!(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,7 @@
//!
//! 1. Two `document_schemaless` collections, each created `WITH
//! (bitemporal=true)`, are placed on DIFFERENT vShards
//! (`distinct_vshard_bitemporal_collections`, same technique as
//! `calvin_multi_shard_redo_restart.rs`).
//! (`vshard_names::distinct_vshard_collections`).
//! 2. `BEGIN; INSERT INTO <a>; INSERT INTO <b>; COMMIT` is sent as ONE
//! `simple_query` call, so both writes are buffered inside the block and,
//! on COMMIT, `classify_dispatch` sees writes on two vShards → MultiShard
Expand All @@ -29,9 +28,9 @@ use crate::common;
use std::sync::atomic::Ordering;
use std::time::Duration;

use nodedb::types::{DatabaseId, VShardId};
use tokio_postgres::SimpleQueryMessage;

use super::vshard_names::distinct_vshard_collections;
use common::cluster_harness::{TestClusterNode, read_once_a_leader_exists, wait_for};

/// Observed sequencer-group leader id from a node's local Raft status, or `0`
Expand All @@ -58,26 +57,6 @@ fn admitted_total(node: &TestClusterNode) -> u64 {
.unwrap_or(0)
}

/// A `(coll_a, coll_b)` pair of bitemporal collection names whose vShard ids
/// differ, so a transaction writing to both is genuinely multi-shard.
/// Deterministic: `VShardId::from_collection_in_database` is a pure function
/// of the database id + collection name bytes. Same technique as
/// `calvin_multi_shard_redo_restart.rs::distinct_vshard_collections`.
fn distinct_vshard_bitemporal_collections() -> (String, String) {
let a_name = "bt_a".to_string();
let va = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &a_name).as_u32();
for i in 0u32..512 {
let b_name = format!("bt_b_{i}");
if VShardId::from_collection_in_database(DatabaseId::DEFAULT, &b_name).as_u32() != va {
return (a_name, b_name);
}
}
panic!(
"could not find a second bitemporal collection name on a distinct vShard \
from the first in 512 tries"
);
}

/// Current `value` for `id` in a bitemporal document collection, or `None` if
/// not visible.
///
Expand Down Expand Up @@ -156,7 +135,7 @@ async fn calvin_multi_shard_bitemporal_commit_survives_wal_only_restart() {
)
.await;

let (coll_a, coll_b) = distinct_vshard_bitemporal_collections();
let (coll_a, coll_b) = distinct_vshard_collections("bt_a", "bt_b");

node.client
.simple_query(&format!(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@
//! `vector_index_txn_restart.rs` for the single-shard analogue).
//!
//! 1. Two collections — a KV collection and a vector-indexed document
//! collection — are created on DIFFERENT vShards (`distinct_vshard_
//! collections`, same technique as `calvin_cluster_pgwire_e2e.rs`).
//! collection — are created on DIFFERENT vShards
//! (`vshard_names::distinct_vshard_collections`).
//! 2. `BEGIN; INSERT INTO <kv>; INSERT INTO <vecdocs>; COMMIT` is sent as ONE
//! `simple_query` call. tokio-postgres ships this as a single wire
//! message; the server buffers the two INSERTs during the transaction and,
Expand All @@ -25,9 +25,9 @@ use crate::common;
use std::sync::atomic::Ordering;
use std::time::Duration;

use nodedb::types::{DatabaseId, VShardId};
use tokio_postgres::SimpleQueryMessage;

use super::vshard_names::distinct_vshard_collections;
use common::cluster_harness::{TestClusterNode, read_once_a_leader_exists, wait_for};

/// Observed sequencer-group leader id from a node's local Raft status, or `0`
Expand All @@ -54,26 +54,6 @@ fn admitted_total(node: &TestClusterNode) -> u64 {
.unwrap_or(0)
}

/// A `(kv_name, vec_name)` pair of collection names whose vShard ids differ,
/// so a transaction writing to both is genuinely multi-shard. Deterministic:
/// `VShardId::from_collection_in_database` is a pure function of the database
/// id + collection name bytes. Same technique as
/// `calvin_cluster_pgwire_e2e.rs::two_distinct_vshard_collections`.
fn distinct_vshard_collections() -> (String, String) {
let kv_name = "cmr_kv".to_string();
let vkv = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &kv_name).as_u32();
for i in 0u32..512 {
let vec_name = format!("cmr_vecdocs_{i}");
if VShardId::from_collection_in_database(DatabaseId::DEFAULT, &vec_name).as_u32() != vkv {
return (kv_name, vec_name);
}
}
panic!(
"could not find a vector-doc collection name on a distinct vShard from \
the KV collection in 512 tries"
);
}

/// Single-row `col` value for `id` in a KV/document collection, or `None` if
/// not visible.
///
Expand Down Expand Up @@ -140,7 +120,7 @@ async fn calvin_multi_shard_write_in_explicit_block_commits_and_survives_restart
)
.await;

let (kv, vecdocs) = distinct_vshard_collections();
let (kv, vecdocs) = distinct_vshard_collections("cmr_kv", "cmr_vecdocs");

node.client
.simple_query(&format!(
Expand Down
Loading
Loading