From caac467124be6835e9992024f25584b2b39e6209 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Tue, 21 Jul 2026 04:29:35 +0000 Subject: [PATCH 01/10] fix issues 188 through 193 --- nodedb-client/src/native/connection/mod.rs | 46 +- .../tests/descriptor_versioning_cross_node.rs | 155 ++++- .../tests/single_node_calvin_graph_txn.rs | 221 ++++++- .../src/calvin/types/transaction.rs | 95 ++- nodedb-cluster/src/metadata_group/cache.rs | 7 + nodedb-cluster/src/metadata_group/entry.rs | 20 + nodedb-sql/src/catalog.rs | 68 +++ nodedb-sql/src/planner/catalog_expr_fold.rs | 223 +++++++ nodedb-sql/src/planner/catalog_fold.rs | 411 ++++++++----- nodedb-sql/src/planner/catalog_plan_shapes.rs | 120 ++++ .../src/planner/catalog_plan_validate.rs | 195 ++++++ nodedb-sql/src/planner/const_fold.rs | 3 + nodedb-sql/src/planner/mod.rs | 3 + nodedb-sql/src/planner/select/helpers.rs | 2 +- nodedb-sql/src/planner/select/mod.rs | 2 +- nodedb-sql/src/planner/select/select_stmt.rs | 6 +- .../src/native_harness/frames.rs | 64 +- .../control/catalog_entry/apply/collection.rs | 42 +- .../control/catalog_entry/descriptor_stamp.rs | 559 +++++++++++++++--- .../scheduler/driver/core/commit_redo.rs | 8 +- .../driver/core/commit_resolution_dispatch.rs | 53 ++ .../scheduler/driver/core/commit_resolve.rs | 115 ++-- .../scheduler/driver/core/completion_route.rs | 13 +- .../calvin/scheduler/driver/core/dispatch.rs | 68 ++- .../calvin/scheduler/driver/core/mod.rs | 1 + .../calvin/scheduler/driver/core/request.rs | 3 +- .../calvin/scheduler/driver/core/routing.rs | 84 ++- .../driver/core/write_version_record.rs | 6 +- .../cluster/calvin/scheduler/driver/types.rs | 3 + .../cluster/metadata_applier/catalog_ddl.rs | 21 +- .../cluster/metadata_applier/dispatch.rs | 47 ++ nodedb/src/control/gateway/version_set.rs | 10 +- nodedb/src/control/metadata_proposer.rs | 167 +++++- nodedb/src/control/planner/calvin/dispatch.rs | 11 +- .../calvin/tx_class/dependent_builder.rs | 19 +- .../planner/calvin/tx_class/static_builder.rs | 39 +- .../catalog_adapter/sql_catalog_impl.rs | 34 +- .../control/planner/context/query/planning.rs | 42 +- .../control/security/catalog/collections.rs | 38 ++ .../dispatch_utils/change_events/extract.rs | 42 +- .../dispatch_utils/change_events/mod.rs | 4 +- .../dispatch_utils/change_events/publish.rs | 2 +- .../src/control/server/dispatch_utils/mod.rs | 5 +- .../server/native/dispatch/transaction.rs | 76 ++- .../control/server/pgwire/handler/dispatch.rs | 4 +- .../shared/ddl/neutral/graph_ops/stats.rs | 36 +- .../server/shared/ddl/neutral/tenant/drop.rs | 10 + .../server/shared/ddl/neutral/user/drop.rs | 8 +- .../server/shared/ddl/neutral/user/mod.rs | 1 + .../shared/ddl/neutral/user/reassign_owned.rs | 11 +- .../shared/ddl/neutral/user/tenant_purge.rs | 114 ++++ .../control/server/shared/session/commit.rs | 197 ++---- .../server/shared/session/commit_calvin.rs | 49 +- .../server/shared/session/ddl_buffer.rs | 163 ++++- .../server/shared/session/leader_forward.rs | 4 +- .../server/shared/session/lifecycle.rs | 2 +- .../src/control/server/shared/session/mod.rs | 2 + .../server/shared/session/overlay_drop.rs | 56 ++ nodedb/src/control/state/fields.rs | 23 +- nodedb/src/control/state/init.rs | 4 + nodedb/src/control/state/init_prod/open.rs | 4 + nodedb/src/data/executor/core_loop/open.rs | 1 + nodedb/src/data/executor/core_loop/state.rs | 4 + .../executor/handlers/aggregate/cache_key.rs | 72 ++- .../data/executor/handlers/aggregate/exec.rs | 4 + .../handlers/aggregate/streaming/over_docs.rs | 2 + .../data/executor/handlers/control/calvin.rs | 18 +- .../executor/handlers/graph_edge_write.rs | 40 +- .../src/data/executor/handlers/graph_stats.rs | 11 +- .../executor/handlers/transaction/batch.rs | 18 +- .../overlay/graph_staged/txn_overlay.rs | 20 + .../handlers/transaction/resolve/entry.rs | 19 +- .../handlers/transaction/resolve/graph.rs | 38 +- .../executor/handlers/transaction/sub_plan.rs | 98 ++- .../handlers/transaction/sub_plan_write.rs | 31 +- .../handlers/transaction/undo/apply.rs | 52 +- .../handlers/transaction/undo/rollback.rs | 30 +- .../handlers/transaction/undo/tests.rs | 48 ++ .../data/executor/wal_replay_redo_graph.rs | 110 +++- .../src/engine/graph/edge_store/stats/read.rs | 84 +-- .../engine/graph/edge_store/stats/table.rs | 104 +++- .../engine/graph/edge_store/temporal/write.rs | 115 +++- nodedb/tests/engine_surface_graph_stats.rs | 39 ++ .../executor_tests/test_aggregate_aliases.rs | 201 +++++++ nodedb/tests/native_gateway_txn_overlay.rs | 127 ++++ nodedb/tests/pg_catalog_regclass.rs | 138 +++++ nodedb/tests/tenant_drop_owned_objects.rs | 171 ++++++ 87 files changed, 4520 insertions(+), 916 deletions(-) create mode 100644 nodedb-sql/src/planner/catalog_expr_fold.rs create mode 100644 nodedb-sql/src/planner/catalog_plan_shapes.rs create mode 100644 nodedb-sql/src/planner/catalog_plan_validate.rs create mode 100644 nodedb/src/control/cluster/calvin/scheduler/driver/core/commit_resolution_dispatch.rs create mode 100644 nodedb/src/control/server/shared/ddl/neutral/user/tenant_purge.rs create mode 100644 nodedb/src/control/server/shared/session/overlay_drop.rs create mode 100644 nodedb/tests/pg_catalog_regclass.rs create mode 100644 nodedb/tests/tenant_drop_owned_objects.rs diff --git a/nodedb-client/src/native/connection/mod.rs b/nodedb-client/src/native/connection/mod.rs index e6db85703..d610c3c31 100644 --- a/nodedb-client/src/native/connection/mod.rs +++ b/nodedb-client/src/native/connection/mod.rs @@ -374,9 +374,9 @@ impl NativeConnection { self.stream.flush().await.map_err(io_err)?; let mut combined_rows: Vec> = Vec::new(); - let mut final_resp: Option = None; + let mut partial_columns: Option> = None; - loop { + let final_resp = loop { let mut len_buf = [0u8; FRAME_HEADER_LEN]; self.stream.read_exact(&mut len_buf).await.map_err(io_err)?; let resp_len = u32::from_be_bytes(len_buf); @@ -397,29 +397,35 @@ impl NativeConnection { })?; if resp.status == ResponseStatus::Partial { + if partial_columns.is_none() { + partial_columns = resp.columns; + } if let Some(rows) = resp.rows { combined_rows.extend(rows); } - if final_resp.is_none() { - final_resp = Some(NativeResponse { rows: None, ..resp }); - } - } else { - if combined_rows.is_empty() { - final_resp = Some(resp); - } else { - if let Some(ref rows) = resp.rows { - combined_rows.extend(rows.iter().cloned()); - } - let mut merged = final_resp.unwrap_or(resp); - merged.rows = Some(combined_rows); - merged.status = ResponseStatus::Ok; - final_resp = Some(merged); - } - break; + continue; } - } - final_resp.ok_or_else(|| NodeDbError::internal("no final response received")) + // The terminal frame owns status and all terminal metadata. In + // particular, never turn a stream error into success merely + // because earlier partial rows were received. + if resp.status == ResponseStatus::Error { + break resp; + } + let mut terminal = resp; + if let Some(rows) = terminal.rows.take() { + combined_rows.extend(rows); + } + if !combined_rows.is_empty() { + terminal.rows = Some(combined_rows); + } + if terminal.columns.is_none() { + terminal.columns = partial_columns; + } + break terminal; + }; + + Ok(final_resp) } } diff --git a/nodedb-cluster-tests/tests/descriptor_versioning_cross_node.rs b/nodedb-cluster-tests/tests/descriptor_versioning_cross_node.rs index 0389adb8d..7d317db80 100644 --- a/nodedb-cluster-tests/tests/descriptor_versioning_cross_node.rs +++ b/nodedb-cluster-tests/tests/descriptor_versioning_cross_node.rs @@ -19,7 +19,7 @@ mod common; use std::time::Duration; -use common::cluster_harness::{TestCluster, wait_for}; +use common::cluster_harness::{TestCluster, TestClusterNode, wait_for}; const TENANT: u64 = 1; @@ -144,6 +144,49 @@ async fn alter_collection_bumps_version_monotonically() { cluster.shutdown().await; } +#[tokio::test(flavor = "multi_thread", worker_threads = 6)] +async fn concurrent_cross_node_updates_allocate_distinct_descriptor_versions() { + let cluster = TestCluster::spawn_three().await.expect("3-node cluster"); + cluster + .exec_ddl_on_any_leader("CREATE USER owner_a WITH PASSWORD 'pw' ROLE READWRITE") + .await + .expect("create owner_a"); + cluster + .exec_ddl_on_any_leader("CREATE USER owner_b WITH PASSWORD 'pw' ROLE READWRITE") + .await + .expect("create owner_b"); + cluster + .exec_ddl_on_any_leader("CREATE COLLECTION concurrently_owned") + .await + .expect("create collection"); + + let update_a = cluster.nodes[0] + .client + .simple_query("ALTER COLLECTION concurrently_owned OWNER TO owner_a"); + let update_b = cluster.nodes[1] + .client + .simple_query("ALTER COLLECTION concurrently_owned OWNER TO owner_b"); + let (result_a, result_b) = tokio::join!(update_a, update_b); + result_a.expect("node 0 concurrent owner update"); + result_b.expect("node 1 concurrent owner update"); + + wait_for( + "both concurrent updates apply as distinct versions on every node", + Duration::from_secs(15), + Duration::from_millis(50), + || { + cluster.nodes.iter().all(|node| { + node.collection_descriptor(TENANT, "concurrently_owned") + .map(|stamp| stamp.0) + == Some(3) + }) + }, + ) + .await; + + cluster.shutdown().await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 6)] async fn distinct_collections_get_independent_versions() { let cluster = TestCluster::spawn_three().await.expect("3-node cluster"); @@ -187,3 +230,113 @@ async fn distinct_collections_get_independent_versions() { cluster.shutdown().await; } + +#[tokio::test(flavor = "multi_thread", worker_threads = 6)] +async fn historical_descriptor_entries_replay_without_regressing_the_latest_version() { + let data_dir = tempfile::tempdir().expect("tempdir"); + let data_path = data_dir.path().to_path_buf(); + let node = TestClusterNode::spawn_single_node_calvin_on_path(4, data_path.clone()) + .await + .expect("spawn single-node metadata group"); + wait_for( + "single-node sequencer leader elected", + Duration::from_secs(10), + Duration::from_millis(50), + || node.sequencer_leader() == node.node_id, + ) + .await; + wait_for( + "single-node metadata leader elected", + Duration::from_secs(10), + Duration::from_millis(50), + || node.shared.is_metadata_leader(), + ) + .await; + + node.client + .simple_query( + "CREATE COLLECTION replay_graph (id TEXT PRIMARY KEY, name TEXT) \ + WITH (engine='document_strict')", + ) + .await + .expect("create graph-bearing collection"); + wait_for( + "collection descriptor reaches version 1", + Duration::from_secs(10), + Duration::from_millis(50), + || { + node.collection_descriptor(TENANT, "replay_graph") + .map(|v| v.0) + == Some(1) + }, + ) + .await; + + node.client + .simple_query( + "GRAPH INSERT EDGE IN replay_graph FROM 'a' TO 'b' \ + TYPE 'knows' PROPERTIES '{}'", + ) + .await + .expect("insert edge and mark collection edge-bearing"); + wait_for( + "edge-bearing descriptor reaches version 2", + Duration::from_secs(10), + Duration::from_millis(50), + || { + node.collection_descriptor(TENANT, "replay_graph") + .map(|v| v.0) + == Some(2) + }, + ) + .await; + + node.graceful_shutdown_wal_only().await; + let node = TestClusterNode::spawn_single_node_calvin_on_path(4, data_path) + .await + .expect("restart against the persisted catalog and full metadata log"); + wait_for( + "single-node sequencer leader re-elected after restart", + Duration::from_secs(10), + Duration::from_millis(50), + || node.sequencer_leader() == node.node_id, + ) + .await; + wait_for( + "single-node metadata leader re-elected after restart", + Duration::from_secs(10), + Duration::from_millis(50), + || node.shared.is_metadata_leader(), + ) + .await; + + wait_for( + "latest collection descriptor remains visible after replay", + Duration::from_secs(10), + Duration::from_millis(50), + || { + node.collection_descriptor(TENANT, "replay_graph") + .map(|v| v.0) + == Some(2) + }, + ) + .await; + + node.client + .simple_query( + "CREATE COLLECTION ddl_after_descriptor_replay \ + (id TEXT PRIMARY KEY) WITH (engine='document_strict')", + ) + .await + .expect( + "historical metadata replay must advance its watermark so later DDL remains usable", + ); + assert_eq!( + node.collection_descriptor(TENANT, "replay_graph") + .map(|version| version.0), + Some(2), + "replaying historical version 1 must not overwrite the persisted latest version 2" + ); + + node.shutdown().await; +} diff --git a/nodedb-cluster-tests/tests/single_node_calvin_graph_txn.rs b/nodedb-cluster-tests/tests/single_node_calvin_graph_txn.rs index 5272044a0..a56ceff65 100644 --- a/nodedb-cluster-tests/tests/single_node_calvin_graph_txn.rs +++ b/nodedb-cluster-tests/tests/single_node_calvin_graph_txn.rs @@ -51,16 +51,30 @@ fn sequencer_leader(node: &TestClusterNode) -> u64 { /// between them is genuinely cross-shard. Deterministic: `VShardId::from_key` is /// a pure function of the key bytes, and it is how `insert_edge` homes each /// endpoint. Same key-picking approach as the sibling `single_node_calvin` suite. -fn distinct_vshard_node_keys() -> (String, String) { +fn distinct_core_node_keys(num_cores: u32) -> (String, String) { let dst = "sncgtx_hub".to_string(); let vdst = VShardId::from_key(dst.as_bytes()).as_u32(); for i in 0u32..4096 { let src = format!("sncgtx_src_{i}"); - if VShardId::from_key(src.as_bytes()).as_u32() != vdst { + let vsrc = VShardId::from_key(src.as_bytes()).as_u32(); + if vsrc != vdst && vsrc % num_cores != vdst % num_cores { return (src, dst); } } - panic!("could not find a node key on a distinct vShard from the hub in 4096 tries"); + panic!("could not find node keys on distinct vShards and cores in 4096 tries"); +} + +fn another_source_on_distinct_core(num_cores: u32, first: &str, dst: &str) -> String { + let first_core = VShardId::from_key(first.as_bytes()).as_u32() % num_cores; + let dst_core = VShardId::from_key(dst.as_bytes()).as_u32() % num_cores; + for i in 0u32..4096 { + let candidate = format!("sncgtx_other_src_{i}"); + let core = VShardId::from_key(candidate.as_bytes()).as_u32() % num_cores; + if candidate != first && core != first_core && core != dst_core { + return candidate; + } + } + panic!("could not find a second source on a distinct core"); } /// Run `GRAPH NEIGHBORS OF '' LABEL '