diff --git a/Cargo.lock b/Cargo.lock index 148c6975c..6b01ae303 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2165,7 +2165,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -5591,7 +5591,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -6218,7 +6218,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -6286,7 +6286,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -6937,10 +6937,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.2", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -8100,7 +8100,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -8597,18 +8597,18 @@ dependencies = [ [[package]] name = "zerompk" -version = "0.5.0" +version = "0.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6ed7119bccd8297686328e73233a8c839645e742e16c2b8e94a48bda9f118729" +checksum = "147801aff255b90229e978171d3911c19b07214e9c808e6d2155b90094748dbf" dependencies = [ "zerompk_derive", ] [[package]] name = "zerompk_derive" -version = "0.5.0" +version = "0.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4cec68b6d21b4d29083199a2b2f42fb23712cfdf3bdb221b86c2485c6b6c4248" +checksum = "0a314420de412a9e7337bb7fc1ab7797902b1e6cd41cc485d240f60116e4184d" dependencies = [ "proc-macro2", "quote", @@ -8709,7 +8709,3 @@ dependencies = [ "cc", "pkg-config", ] - -[[patch.unused]] -name = "pagedb" -version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 284e6a090..c6c1bd143 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -164,7 +164,7 @@ chrono = { version = "0.4", features = ["std"], default-features = false } # MessagePack rmpv = "1" -zerompk = { version = "0.5", features = ["std", "derive"] } +zerompk = { version = "0.6", features = ["std", "derive"] } # Bitmap roaring = "0.11" diff --git a/nodedb-client/src/native/connection/mod.rs b/nodedb-client/src/native/connection/mod.rs index e6db85703..05cf17858 100644 --- a/nodedb-client/src/native/connection/mod.rs +++ b/nodedb-client/src/native/connection/mod.rs @@ -356,9 +356,10 @@ impl NativeConnection { op: OpCode, fields: TextFields, ) -> NodeDbResult { + let req_seq = self.next_seq(); let req = NativeRequest { op, - seq: self.next_seq(), + seq: req_seq, fields: RequestFields::Text(fields), }; @@ -374,9 +375,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); @@ -396,30 +397,51 @@ impl NativeConnection { NodeDbError::serialization("msgpack", format!("response decode: {e}")) })?; + if resp.seq != req_seq { + // A fan-out query for a preceding request can leave stale + // trailing frames on the wire after that request's terminal + // frame was already returned to its caller. Discard any + // frame that doesn't belong to this request rather than + // misattributing it — never surface another request's rows + // or status as this request's response. + tracing::warn!( + expected_seq = req_seq, + got_seq = resp.seq, + "native connection: discarding stale response frame" + ); + continue; + } + 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 '