Skip to content

Commit f679209

Browse files
authored
Merge pull request #197 from NodeDB-Lab/fix/issues-188-193
Fix issues #188#193
2 parents d75bfd4 + e642ba9 commit f679209

94 files changed

Lines changed: 5072 additions & 1013 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

Cargo.lock

Lines changed: 11 additions & 15 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -164,7 +164,7 @@ chrono = { version = "0.4", features = ["std"], default-features = false }
164164

165165
# MessagePack
166166
rmpv = "1"
167-
zerompk = { version = "0.5", features = ["std", "derive"] }
167+
zerompk = { version = "0.6", features = ["std", "derive"] }
168168

169169
# Bitmap
170170
roaring = "0.11"

nodedb-client/src/native/connection/mod.rs

Lines changed: 43 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -356,9 +356,10 @@ impl NativeConnection {
356356
op: OpCode,
357357
fields: TextFields,
358358
) -> NodeDbResult<NativeResponse> {
359+
let req_seq = self.next_seq();
359360
let req = NativeRequest {
360361
op,
361-
seq: self.next_seq(),
362+
seq: req_seq,
362363
fields: RequestFields::Text(fields),
363364
};
364365

@@ -374,9 +375,9 @@ impl NativeConnection {
374375
self.stream.flush().await.map_err(io_err)?;
375376

376377
let mut combined_rows: Vec<Vec<nodedb_types::Value>> = Vec::new();
377-
let mut final_resp: Option<NativeResponse> = None;
378+
let mut partial_columns: Option<Vec<String>> = None;
378379

379-
loop {
380+
let final_resp = loop {
380381
let mut len_buf = [0u8; FRAME_HEADER_LEN];
381382
self.stream.read_exact(&mut len_buf).await.map_err(io_err)?;
382383
let resp_len = u32::from_be_bytes(len_buf);
@@ -396,30 +397,51 @@ impl NativeConnection {
396397
NodeDbError::serialization("msgpack", format!("response decode: {e}"))
397398
})?;
398399

400+
if resp.seq != req_seq {
401+
// A fan-out query for a preceding request can leave stale
402+
// trailing frames on the wire after that request's terminal
403+
// frame was already returned to its caller. Discard any
404+
// frame that doesn't belong to this request rather than
405+
// misattributing it — never surface another request's rows
406+
// or status as this request's response.
407+
tracing::warn!(
408+
expected_seq = req_seq,
409+
got_seq = resp.seq,
410+
"native connection: discarding stale response frame"
411+
);
412+
continue;
413+
}
414+
399415
if resp.status == ResponseStatus::Partial {
416+
if partial_columns.is_none() {
417+
partial_columns = resp.columns;
418+
}
400419
if let Some(rows) = resp.rows {
401420
combined_rows.extend(rows);
402421
}
403-
if final_resp.is_none() {
404-
final_resp = Some(NativeResponse { rows: None, ..resp });
405-
}
406-
} else {
407-
if combined_rows.is_empty() {
408-
final_resp = Some(resp);
409-
} else {
410-
if let Some(ref rows) = resp.rows {
411-
combined_rows.extend(rows.iter().cloned());
412-
}
413-
let mut merged = final_resp.unwrap_or(resp);
414-
merged.rows = Some(combined_rows);
415-
merged.status = ResponseStatus::Ok;
416-
final_resp = Some(merged);
417-
}
418-
break;
422+
continue;
419423
}
420-
}
421424

422-
final_resp.ok_or_else(|| NodeDbError::internal("no final response received"))
425+
// The terminal frame owns status and all terminal metadata. In
426+
// particular, never turn a stream error into success merely
427+
// because earlier partial rows were received.
428+
if resp.status == ResponseStatus::Error {
429+
break resp;
430+
}
431+
let mut terminal = resp;
432+
if let Some(rows) = terminal.rows.take() {
433+
combined_rows.extend(rows);
434+
}
435+
if !combined_rows.is_empty() {
436+
terminal.rows = Some(combined_rows);
437+
}
438+
if terminal.columns.is_none() {
439+
terminal.columns = partial_columns;
440+
}
441+
break terminal;
442+
};
443+
444+
Ok(final_resp)
423445
}
424446
}
425447

nodedb-cluster-tests/tests/descriptor_versioning_cross_node.rs

Lines changed: 154 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ mod common;
1919

2020
use std::time::Duration;
2121

22-
use common::cluster_harness::{TestCluster, wait_for};
22+
use common::cluster_harness::{TestCluster, TestClusterNode, wait_for};
2323

2424
const TENANT: u64 = 1;
2525

@@ -144,6 +144,49 @@ async fn alter_collection_bumps_version_monotonically() {
144144
cluster.shutdown().await;
145145
}
146146

147+
#[tokio::test(flavor = "multi_thread", worker_threads = 6)]
148+
async fn concurrent_cross_node_updates_allocate_distinct_descriptor_versions() {
149+
let cluster = TestCluster::spawn_three().await.expect("3-node cluster");
150+
cluster
151+
.exec_ddl_on_any_leader("CREATE USER owner_a WITH PASSWORD 'pw' ROLE READWRITE")
152+
.await
153+
.expect("create owner_a");
154+
cluster
155+
.exec_ddl_on_any_leader("CREATE USER owner_b WITH PASSWORD 'pw' ROLE READWRITE")
156+
.await
157+
.expect("create owner_b");
158+
cluster
159+
.exec_ddl_on_any_leader("CREATE COLLECTION concurrently_owned")
160+
.await
161+
.expect("create collection");
162+
163+
let update_a = cluster.nodes[0]
164+
.client
165+
.simple_query("ALTER COLLECTION concurrently_owned OWNER TO owner_a");
166+
let update_b = cluster.nodes[1]
167+
.client
168+
.simple_query("ALTER COLLECTION concurrently_owned OWNER TO owner_b");
169+
let (result_a, result_b) = tokio::join!(update_a, update_b);
170+
result_a.expect("node 0 concurrent owner update");
171+
result_b.expect("node 1 concurrent owner update");
172+
173+
wait_for(
174+
"both concurrent updates apply as distinct versions on every node",
175+
Duration::from_secs(15),
176+
Duration::from_millis(50),
177+
|| {
178+
cluster.nodes.iter().all(|node| {
179+
node.collection_descriptor(TENANT, "concurrently_owned")
180+
.map(|stamp| stamp.0)
181+
== Some(3)
182+
})
183+
},
184+
)
185+
.await;
186+
187+
cluster.shutdown().await;
188+
}
189+
147190
#[tokio::test(flavor = "multi_thread", worker_threads = 6)]
148191
async fn distinct_collections_get_independent_versions() {
149192
let cluster = TestCluster::spawn_three().await.expect("3-node cluster");
@@ -187,3 +230,113 @@ async fn distinct_collections_get_independent_versions() {
187230

188231
cluster.shutdown().await;
189232
}
233+
234+
#[tokio::test(flavor = "multi_thread", worker_threads = 6)]
235+
async fn historical_descriptor_entries_replay_without_regressing_the_latest_version() {
236+
let data_dir = tempfile::tempdir().expect("tempdir");
237+
let data_path = data_dir.path().to_path_buf();
238+
let node = TestClusterNode::spawn_single_node_calvin_on_path(4, data_path.clone())
239+
.await
240+
.expect("spawn single-node metadata group");
241+
wait_for(
242+
"single-node sequencer leader elected",
243+
Duration::from_secs(10),
244+
Duration::from_millis(50),
245+
|| node.sequencer_leader() == node.node_id,
246+
)
247+
.await;
248+
wait_for(
249+
"single-node metadata leader elected",
250+
Duration::from_secs(10),
251+
Duration::from_millis(50),
252+
|| node.shared.is_metadata_leader(),
253+
)
254+
.await;
255+
256+
node.client
257+
.simple_query(
258+
"CREATE COLLECTION replay_graph (id TEXT PRIMARY KEY, name TEXT) \
259+
WITH (engine='document_strict')",
260+
)
261+
.await
262+
.expect("create graph-bearing collection");
263+
wait_for(
264+
"collection descriptor reaches version 1",
265+
Duration::from_secs(10),
266+
Duration::from_millis(50),
267+
|| {
268+
node.collection_descriptor(TENANT, "replay_graph")
269+
.map(|v| v.0)
270+
== Some(1)
271+
},
272+
)
273+
.await;
274+
275+
node.client
276+
.simple_query(
277+
"GRAPH INSERT EDGE IN replay_graph FROM 'a' TO 'b' \
278+
TYPE 'knows' PROPERTIES '{}'",
279+
)
280+
.await
281+
.expect("insert edge and mark collection edge-bearing");
282+
wait_for(
283+
"edge-bearing descriptor reaches version 2",
284+
Duration::from_secs(10),
285+
Duration::from_millis(50),
286+
|| {
287+
node.collection_descriptor(TENANT, "replay_graph")
288+
.map(|v| v.0)
289+
== Some(2)
290+
},
291+
)
292+
.await;
293+
294+
node.graceful_shutdown_wal_only().await;
295+
let node = TestClusterNode::spawn_single_node_calvin_on_path(4, data_path)
296+
.await
297+
.expect("restart against the persisted catalog and full metadata log");
298+
wait_for(
299+
"single-node sequencer leader re-elected after restart",
300+
Duration::from_secs(10),
301+
Duration::from_millis(50),
302+
|| node.sequencer_leader() == node.node_id,
303+
)
304+
.await;
305+
wait_for(
306+
"single-node metadata leader re-elected after restart",
307+
Duration::from_secs(10),
308+
Duration::from_millis(50),
309+
|| node.shared.is_metadata_leader(),
310+
)
311+
.await;
312+
313+
wait_for(
314+
"latest collection descriptor remains visible after replay",
315+
Duration::from_secs(10),
316+
Duration::from_millis(50),
317+
|| {
318+
node.collection_descriptor(TENANT, "replay_graph")
319+
.map(|v| v.0)
320+
== Some(2)
321+
},
322+
)
323+
.await;
324+
325+
node.client
326+
.simple_query(
327+
"CREATE COLLECTION ddl_after_descriptor_replay \
328+
(id TEXT PRIMARY KEY) WITH (engine='document_strict')",
329+
)
330+
.await
331+
.expect(
332+
"historical metadata replay must advance its watermark so later DDL remains usable",
333+
);
334+
assert_eq!(
335+
node.collection_descriptor(TENANT, "replay_graph")
336+
.map(|version| version.0),
337+
Some(2),
338+
"replaying historical version 1 must not overwrite the persisted latest version 2"
339+
);
340+
341+
node.shutdown().await;
342+
}

0 commit comments

Comments
 (0)