Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
64 commits
Select commit Hold shift + click to select a range
0804511
refactor(raft): split commit_resolve into a directory module
farhan-syah Sep 23, 2026
6e6e924
feat(bridge): defer sequenced dispatch on capacity refusal
farhan-syah Sep 23, 2026
dbdf41c
refactor(bridge): name the per-core dispatch queue capacity
farhan-syah Sep 23, 2026
69d0ed4
feat(scheduler): gate Calvin intake on dispatch and backlog capacity
farhan-syah Sep 23, 2026
1885c7a
fix(executor): never expire already-ordered work on deadline
farhan-syah Sep 23, 2026
3741142
fix(executor): run transaction sub-plans in the session's database
farhan-syah Sep 23, 2026
a17406c
feat(scheduler): halt a vShard's Calvin scheduler on apply failure
farhan-syah Sep 23, 2026
e24f5ba
feat(scheduler): retry a Calvin scheduler's sequencer entries until a…
farhan-syah Sep 23, 2026
b32966d
refactor(scheduler): split Calvin lock manager into a directory module
farhan-syah Sep 23, 2026
0642975
refactor(distributed_applier): split apply_loop into a directory module
farhan-syah Sep 23, 2026
138bd27
refactor(dispatch_utils): split the write funnel into a directory module
farhan-syah Sep 23, 2026
08486cd
refactor(pgwire): split handler dispatch into a directory module
farhan-syah Sep 23, 2026
5df1dc6
refactor(executor): split Calvin handlers into a directory module
farhan-syah Sep 23, 2026
712b1b9
feat(wal): apply each replicated proposal exactly once
farhan-syah Sep 23, 2026
8092926
feat(kv): fault KV counter atomics on non-numeric values with typed d…
farhan-syah Sep 23, 2026
02f0851
refactor(wal): replace the apply-key thread-local with an explicit ap…
farhan-syah Sep 23, 2026
2ac0841
feat(executor): fail-stop a core whose rollback state is unknown
farhan-syah Sep 23, 2026
185f7e2
feat(bridge): track a node-wide outcome floor for restart replay
farhan-syah Sep 24, 2026
c7b6223
fix(control): apply every permission-cache event on both consumer paths
farhan-syah Sep 24, 2026
1117e7d
feat(bridge): bound write windows to the outcome floor and fix partia…
farhan-syah Sep 24, 2026
9bee554
test(pgwire-harness): install the production gateway and a single-vot…
farhan-syah Sep 24, 2026
ce322d0
feat(crdt): persist dead-letter entries across restart
farhan-syah Sep 24, 2026
6ac1670
fix(bridge): cancel or hold minted records dropped before dispatch
farhan-syah Sep 24, 2026
c2a8510
refactor(pgwire): drop unused user_id param on dispatch_task_no_wal
farhan-syah Sep 24, 2026
77b05d6
refactor(cluster): split shuffle-push streaming into its own module
farhan-syah Sep 24, 2026
8a988de
feat(physical): support staged-redo transactions and KV counter atomics
farhan-syah Sep 24, 2026
7f0c312
refactor(executor): install transactions from redo, not sub-plan replay
farhan-syah Sep 24, 2026
3960fd6
feat(executor): stamp checkpoints with an outcome-floor and applied-a…
farhan-syah Sep 25, 2026
8bb83f0
feat(control): route every autocommit write through a durable path
farhan-syah Sep 25, 2026
96d31d4
feat(executor): stamp checkpoints with LSN-range replay stamps
farhan-syah Sep 25, 2026
2ccb891
feat(executor): report the outcome floor, not the watermark, for dura…
farhan-syah Sep 25, 2026
9acdde9
test(cluster): accept a rebalanced learner as a follower in a PK-read…
farhan-syah Sep 25, 2026
0a1331e
feat(control): fence reads and writes on stale authorization state
farhan-syah Sep 25, 2026
6572818
feat(security): refuse a role assignment or drop that breaks a role rule
farhan-syah Sep 26, 2026
ad7c0ae
feat(backup): back up and restore on a replicated consistent cut
farhan-syah Sep 26, 2026
aad2507
feat(raft): bound queued proposals by their caller deadline
farhan-syah Sep 26, 2026
40261d0
feat(events): tag restored rows so triggers do not fire twice
farhan-syah Sep 26, 2026
41850d3
feat(events): add REMOVE EVENT and read definitions off the Event Plane
farhan-syah Sep 26, 2026
1d7dd58
feat(events): stamp WAL row writes with their event source
farhan-syah Sep 26, 2026
eabd490
feat(query): decode a KV row back into insert-shaped body fields
farhan-syah Sep 26, 2026
8971b41
feat(sync): push KV writes from Lite through the idempotency gate
farhan-syah Sep 26, 2026
ffec08b
feat(startup): bind client sockets early, listen only once serving
farhan-syah Sep 26, 2026
9f09e0a
feat(query): fail scalar function faults instead of folding to NULL
farhan-syah Sep 26, 2026
0060af6
feat(vector): fail on dimension mismatch instead of panicking
farhan-syah Sep 26, 2026
021b5c2
fix(query): fail search legs on error instead of folding to empty
farhan-syah Sep 26, 2026
70174e5
feat(vector): back IVF-PQ with a live, trainable index
farhan-syah Sep 26, 2026
9f3a62c
feat(vector): queue HNSW builds on a bounded, non-blocking builder
farhan-syah Sep 26, 2026
ad0ec5d
feat(reindex): rebuild graph CSR and FTS indexes without blocking writes
farhan-syah Sep 27, 2026
56fca1d
feat(errors): keep a typed refusal's SQLSTATE class across every plane
farhan-syah Sep 27, 2026
9d4faaf
refactor(routing): key the vShard hash and surrogates on a canonical …
farhan-syah Sep 27, 2026
3c94bb4
feat(diagnostics): capture a descriptor lease the renewal loop drops
farhan-syah Sep 27, 2026
a457bd2
feat(errors): keep a Data-Plane refusal's typed code across the clust…
farhan-syah Sep 27, 2026
b42e8ab
feat(errors): cross a shard's typed error, not just its Data-Plane code
farhan-syah Sep 27, 2026
88069c9
feat(errors): unify HTTP status mapping and widen typed error crossing
farhan-syah Sep 27, 2026
0ed6745
feat(errors): exhaust every typed error match instead of a wildcard arm
farhan-syah Sep 27, 2026
60c9b8e
feat(errors): replace generic DDL error strings with typed propagation
farhan-syah Sep 27, 2026
234e84d
feat(errors): give every SQLSTATE a numeric code and cross DDL errors…
farhan-syah Sep 27, 2026
4289416
feat(backup): restore selectively by database and route by bare colle…
farhan-syah Sep 28, 2026
05f7784
refactor(wal): gate wasm32 exclusion at the module level
farhan-syah Sep 28, 2026
c5265a0
test(inproc): drop all-cores exchange sites from SystemTask allowlist
farhan-syah Sep 28, 2026
e8b3a41
refactor(cluster): move peer pre-registration into warm_peers
farhan-syah Sep 28, 2026
8bdf14c
refactor(control): pull local-dispatch helpers out of dispatch_utils
farhan-syah Sep 28, 2026
32caad3
refactor(errors): name the class-parity test's function-pointer tuples
farhan-syah Sep 28, 2026
15580a0
chore(ci): update authorized-dispatch allowlist for the dispatch modu…
farhan-syah Sep 28, 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: 1 addition & 1 deletion .config/nextest.toml
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,7 @@ slow-timeout = { period = "30s", terminate-after = 8 }
# hide the cause rather than fix it. The tail this costs is bounded — roughly a
# dozen crash/shutdown tests at about ten seconds each.
[[profile.default.overrides]]
filter = 'binary(wal_direct_io) | binary(ilp_client_address) | binary(crash_recovery) | binary(crash_recovery_overlays) | binary(crash_recovery_analytics) | binary(crash_resp_kv_write) | binary(crash_metadata_applier_wedge) | binary(crash_dropped_collection_reclaim) | binary(crash_mid_replay) | binary(crash_checkpoint_corruption) | binary(crash_checkpoint_truncate_window) | binary(crash_refused_write_not_resurrected) | binary(crash_replay_fail_stop) | binary(crash_core_stall) | test(/^cases::startup_failure::/) | test(/^cases::shutdown_in_flight::/) | test(/^cases::shutdown_budget::/) | test(/^cases::shutdown_abort_offender::/) | test(/^cases::shutdown_idempotent::/)'
filter = 'binary(wal_direct_io) | binary(ilp_client_address) | binary(crash_recovery) | binary(crash_recovery_overlays) | binary(crash_recovery_analytics) | binary(crash_resp_kv_write) | binary(crash_metadata_applier_wedge) | binary(crash_dropped_collection_reclaim) | binary(crash_purge_not_resurrected) | binary(crash_mid_replay) | binary(crash_checkpoint_corruption) | binary(crash_checkpoint_truncate_window) | binary(crash_refused_write_not_resurrected) | binary(crash_replay_fail_stop) | binary(crash_core_stall) | binary(crash_replay_stamp) | binary(crash_replay_stamp_calvin) | binary(calvin_hold_liveness) | binary(apply_pipeline_group_independence) | binary(crash_kv_atomic_autocommit) | test(/^cases::startup_failure::/) | test(/^cases::shutdown_in_flight::/) | test(/^cases::shutdown_budget::/) | test(/^cases::shutdown_abort_offender::/) | test(/^cases::shutdown_idempotent::/)'
test-group = 'server-process'
threads-required = 'num-test-threads'

Expand Down
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

15 changes: 13 additions & 2 deletions docs/kv.md
Original file line number Diff line number Diff line change
Expand Up @@ -155,10 +155,21 @@ SELECT KV_GETSET('session_token', 'player-123', 'new-token-xyz');

**RESP (Redis) equivalents:** `INCR`, `DECR`, `INCRBY`, `DECRBY`, `INCRBYFLOAT`, `GETSET` — all work over the RESP protocol.

**Value shapes:**

- Raw value (a single `value` column, or RESP `SET`): the value is a byte string. `INCR`/`INCRBY`/`DECR`/`DECRBY` read it as decimal integer text and store the result as decimal text. `INCRBYFLOAT` and `KV_INCR_FLOAT` read decimal text (plain or exponent form, such as `5.0e3`), add exactly, and store the trimmed decimal text: `"0.1"` plus `0.2` stores `"0.3"`, `"3.0"` plus `0` stores `"3"`. Exact addition covers 28 significant digits below 7.9e28. A value outside that range adds in 64-bit float. `SET k 5` then `INCR k` leaves `"6"`.
- Absent key: the counter starts at 0 and is stored as decimal text.
- Typed row (several columns): the first numeric column in key order moves. Every other column stays.

**Error handling:**

- `TYPE_MISMATCH` (SQLSTATE 42846) — INCR on a non-numeric value
- `OVERFLOW` (SQLSTATE 22003) — i64 overflow on INCR
| Condition | RESP reply | SQLSTATE |
|---|---|---|
| Raw value is not a decimal integer in the i64 range | `ERR value is not an integer or out of range` | `22P02` |
| Raw value is not a decimal float | `ERR value is not a valid float` | `22P02` |
| Integer result leaves the i64 range | `ERR increment or decrement would overflow` | `22003` |
| Float result is NaN or infinite | `ERR increment would produce NaN or Infinity` | `22003` |
| Typed row has no numeric column | `WRONGTYPE ...` | `42846` |

## Sorted Indexes (Leaderboards)

Expand Down
10 changes: 8 additions & 2 deletions nodedb-client/src/native/connection/response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,10 +40,16 @@ fn error_frame_to_typed(
if payload.ndb_code == 0 {
return NodeDbError::internal(payload.message.clone());
}
NodeDbError::from_wire(
let error = NodeDbError::from_wire_with_details(
nodedb_types::error::ErrorCode(payload.ndb_code),
payload.message.clone(),
)
payload.details.clone(),
);
// The typed cause, when the server sent one, becomes the error's cause.
match &payload.cause {
Some(cause) => error.with_cause(cause.to_error()),
None => error,
}
}

pub(super) fn response_to_query_result(resp: NativeResponse) -> NodeDbResult<QueryResult> {
Expand Down
4 changes: 4 additions & 0 deletions nodedb-cluster-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,10 @@ description = "3-node integration test suite for NodeDB. Heavy; runs in its own
[lib]
path = "src/lib.rs"

[features]
# Arms the in-process fail points the cluster tests park writes at.
failpoints = ["nodedb/failpoints", "nodedb-types/failpoints"]

[dependencies]

[dev-dependencies]
Expand Down
34 changes: 19 additions & 15 deletions nodedb-cluster-tests/tests/cluster_common/calvin_test_node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@
//! for shutdown.
//! - `add_vshard_sender(vshard_id, sender)` — wire a per-vshard channel
//! into the state machine so tests can assert fan-out.
//! - `try_recv_txn(rx)` — read the next sequenced txn a fan-out channel
//! holds.

#![allow(dead_code)] // Not every test file uses every helper.

Expand Down Expand Up @@ -130,24 +132,15 @@ impl CalvinTestNode {

/// Register a per-vshard output sender so the state machine can
/// fan out sequenced transactions to the receiving test code.
pub fn add_vshard_sender(&self, vshard_id: u32, sender: mpsc::Sender<SequencedTxn>) {
// The state machine now fans out `SchedulerInput`; these tests assert on
// the sequenced-txn stream, so adapt: forward only `Txn` payloads to the
// caller's `SequencedTxn` channel (reservation inputs don't occur here).
let (adapt_tx, mut adapt_rx) = mpsc::channel::<SchedulerInput>(512);
tokio::spawn(async move {
while let Some(input) = adapt_rx.recv().await {
if let SchedulerInput::Txn(txn) = input
&& sender.send(txn).await.is_err()
{
break;
}
}
});
///
/// The state machine sends into `sender` itself, inside `apply`, before
/// it advances `last_applied_epoch`. A test that waits for the epoch can
/// then read the fan-out with `try_recv_txn` and no further wait.
pub fn add_vshard_sender(&self, vshard_id: u32, sender: mpsc::Sender<SchedulerInput>) {
self.state_machine
.lock()
.unwrap_or_else(|p| p.into_inner())
.set_vshard_sender(vshard_id, adapt_tx);
.set_vshard_sender(vshard_id, sender);
}

/// Start the sequencer service epoch-ticker task on this node.
Expand Down Expand Up @@ -417,3 +410,14 @@ pub async fn wait_for_sequencer_leader(
tokio::time::sleep(step).await;
}
}

/// The next sequenced txn `rx` holds, skipping any other scheduler input.
/// `None` when the channel holds no txn.
pub fn try_recv_txn(rx: &mut mpsc::Receiver<SchedulerInput>) -> Option<SequencedTxn> {
while let Ok(input) = rx.try_recv() {
if let SchedulerInput::Txn(txn) = input {
return Some(txn);
}
}
None
}
2 changes: 1 addition & 1 deletion nodedb-cluster-tests/tests/cluster_common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,6 @@ pub mod rebalancer;
pub mod test_node;

pub use calvin_test_node::{
CalvinApplier, CalvinTestNode, spawn_with_sequencer, wait_for_sequencer_leader,
CalvinApplier, CalvinTestNode, spawn_with_sequencer, try_recv_txn, wait_for_sequencer_leader,
};
pub use test_node::{NoopApplier, TestNode, test_transport, wait_for};
Original file line number Diff line number Diff line change
Expand Up @@ -19,22 +19,21 @@ use std::time::Duration;

use nodedb_cluster::calvin::{
sequencer::{SequencerConfig, new_inbox},
types::{EngineKeySet, ReadWriteSet, SequencedTxn, SortedVec, TxClass, VersionedReadSet},
};
use nodedb_types::{
TenantId,
id::{DatabaseId, VShardId},
types::{EngineKeySet, ReadWriteSet, SchedulerInput, SortedVec, TxClass, VersionedReadSet},
};
use nodedb_types::{TenantId, id::DatabaseId};
use tokio::sync::mpsc;

use super::cluster_common::{spawn_with_sequencer, wait_for_sequencer_leader};
use super::cluster_common::{spawn_with_sequencer, try_recv_txn, wait_for_sequencer_leader};

/// Find two collection names that hash to distinct vshards.
fn two_distinct_collections() -> (String, String) {
let mut first: Option<(String, u32)> = None;
for i in 0u32..512 {
let name = format!("col_{i}");
let vshard = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &name).as_u32();
let vshard = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &name)
.vshard()
.as_u32();
if let Some((ref fname, fv)) = first {
if fv != vshard {
return (fname.clone(), name);
Expand All @@ -48,8 +47,12 @@ fn two_distinct_collections() -> (String, String) {

fn make_multishard_txclass() -> (TxClass, u32, u32) {
let (col_a, col_b) = two_distinct_collections();
let va = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_a).as_u32();
let vb = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_b).as_u32();
let va = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_a)
.vshard()
.as_u32();
let vb = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_b)
.vshard()
.as_u32();
let write_set = ReadWriteSet::new(vec![
EngineKeySet::Document {
collection: col_a,
Expand Down Expand Up @@ -92,11 +95,15 @@ async fn sequencer_normal_path_commit_on_all_replicas() {

// Wire per-vshard receivers on every node.
let (tx_a, col_b_name) = two_distinct_collections();
let va = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &tx_a).as_u32();
let vb = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_b_name).as_u32();

let mut vshard_rxs_a: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut vshard_rxs_b: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let va = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &tx_a)
.vshard()
.as_u32();
let vb = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_b_name)
.vshard()
.as_u32();

let mut vshard_rxs_a: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
let mut vshard_rxs_b: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
for node in &nodes {
let (tx_a_ch, rx_a) = mpsc::channel(64);
let (tx_b_ch, rx_b) = mpsc::channel(64);
Expand Down Expand Up @@ -149,8 +156,8 @@ async fn sequencer_normal_path_commit_on_all_replicas() {
.zip(vshard_rxs_b.iter_mut())
.enumerate()
{
let got_a = rx_a.try_recv().is_ok();
let got_b = rx_b.try_recv().is_ok();
let got_a = try_recv_txn(rx_a).is_some();
let got_b = try_recv_txn(rx_b).is_some();
assert!(
got_a || got_b,
"node {}: neither vshard receiver got the txn fan-out",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,23 +43,22 @@ use std::time::Duration;

use nodedb_cluster::calvin::{
sequencer::{SequencerConfig, new_inbox},
types::{EngineKeySet, ReadWriteSet, SequencedTxn, SortedVec, TxClass, VersionedReadSet},
};
use nodedb_types::{
TenantId,
id::{DatabaseId, VShardId},
types::{EngineKeySet, ReadWriteSet, SchedulerInput, SortedVec, TxClass, VersionedReadSet},
};
use nodedb_types::{TenantId, id::DatabaseId};
use tokio::sync::mpsc;

use super::cluster_common::{spawn_with_sequencer, wait_for_sequencer_leader};
use super::cluster_common::{spawn_with_sequencer, try_recv_txn, wait_for_sequencer_leader};

// ── Helpers ──────────────────────────────────────────────────────────────────

fn two_distinct_collections() -> (String, String) {
let mut first: Option<(String, u32)> = None;
for i in 0u32..512 {
let name = format!("col_{i}");
let vshard = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &name).as_u32();
let vshard = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &name)
.vshard()
.as_u32();
if let Some((ref fname, fv)) = first {
if fv != vshard {
return (fname.clone(), name);
Expand Down Expand Up @@ -130,11 +129,15 @@ async fn scheduler_catchup_via_raft_log_replay() {

// Wire per-vshard receivers on every node so we can verify fan-out.
let (col_a, col_b) = two_distinct_collections();
let va = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_a).as_u32();
let vb = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_b).as_u32();
let va = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_a)
.vshard()
.as_u32();
let vb = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_b)
.vshard()
.as_u32();

let mut vshard_rxs_a: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut vshard_rxs_b: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut vshard_rxs_a: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
let mut vshard_rxs_b: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
for node in &nodes {
let (tx_a, rx_a) = mpsc::channel(128);
let (tx_b, rx_b) = mpsc::channel(128);
Expand Down Expand Up @@ -246,7 +249,7 @@ async fn scheduler_catchup_via_raft_log_replay() {
// Drain whatever arrived — we care that the routing worked, not the count.
let mut total_received = 0usize;
for rx in vshard_rxs_a.iter_mut().chain(vshard_rxs_b.iter_mut()) {
while rx.try_recv().is_ok() {
while try_recv_txn(rx).is_some() {
total_received += 1;
}
}
Expand Down
39 changes: 22 additions & 17 deletions nodedb-cluster-tests/tests/cluster_suite/cases/calvin_e2e_ollp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,15 +37,12 @@ use std::time::Duration;

use nodedb_cluster::calvin::{
sequencer::{SequencerConfig, new_inbox},
types::{EngineKeySet, ReadWriteSet, SequencedTxn, SortedVec, TxClass, VersionedReadSet},
};
use nodedb_types::{
TenantId,
id::{DatabaseId, VShardId},
types::{EngineKeySet, ReadWriteSet, SchedulerInput, SortedVec, TxClass, VersionedReadSet},
};
use nodedb_types::{TenantId, id::DatabaseId};
use tokio::sync::mpsc;

use super::cluster_common::{spawn_with_sequencer, wait_for_sequencer_leader};
use super::cluster_common::{spawn_with_sequencer, try_recv_txn, wait_for_sequencer_leader};

/// Find two collection names that hash to distinct vshards.
///
Expand All @@ -54,7 +51,9 @@ fn two_distinct_vshard_collections() -> (String, String) {
let mut first: Option<(String, u32)> = None;
for i in 0u32..512 {
let name = format!("ollp_col_{i}");
let vshard = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &name).as_u32();
let vshard = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &name)
.vshard()
.as_u32();
if let Some((ref fname, fv)) = first {
if fv != vshard {
return (fname.clone(), name);
Expand Down Expand Up @@ -104,14 +103,17 @@ fn make_ollp_tx_class(
.expect("valid multi-vshard OLLP TxClass")
}

/// Assert: one specific vshard channel received at least one SequencedTxn.
/// Assert: one specific vshard channel received at least one sequenced txn.
///
/// The state machine sends into the channel inside `apply`, before it
/// advances the epoch the caller waited for, so the txn is already there.
fn assert_fan_out_received(
rx: &mut mpsc::Receiver<SequencedTxn>,
rx: &mut mpsc::Receiver<SchedulerInput>,
vshard_id: u32,
replica_idx: usize,
) {
assert!(
rx.try_recv().is_ok(),
try_recv_txn(rx).is_some(),
"replica {replica_idx}: vshard {vshard_id} fan-out channel received no txn"
);
}
Expand All @@ -134,13 +136,16 @@ async fn ollp_bulk_update_txclass_admitted_and_fanned_out() {

// Find two collections that hash to distinct vshards.
let (col_static, col_ollp) = two_distinct_vshard_collections();
let vs_static =
VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_static).as_u32();
let vs_ollp = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &col_ollp).as_u32();
let vs_static = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_static)
.vshard()
.as_u32();
let vs_ollp = nodedb_types::CollectionKey::from_bare(DatabaseId::DEFAULT, &col_ollp)
.vshard()
.as_u32();

// Wire per-vshard fan-out receivers on every replica.
let mut rxs_static: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut rxs_ollp: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut rxs_static: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
let mut rxs_ollp: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
for node in &nodes {
let (tx_s, rx_s) = mpsc::channel(64);
let (tx_o, rx_o) = mpsc::channel(64);
Expand Down Expand Up @@ -191,8 +196,8 @@ async fn ollp_bulk_update_txclass_admitted_and_fanned_out() {

// Simulate an OLLP retry: concurrent insert added surrogate 4 to col_ollp.
// Re-wire fresh fan-out receivers and re-submit with the corrected set.
let mut retry_rxs_static: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut retry_rxs_ollp: Vec<mpsc::Receiver<SequencedTxn>> = Vec::new();
let mut retry_rxs_static: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
let mut retry_rxs_ollp: Vec<mpsc::Receiver<SchedulerInput>> = Vec::new();
for node in &nodes {
let (tx_s, rx_s) = mpsc::channel(64);
let (tx_o, rx_o) = mpsc::channel(64);
Expand Down
Loading
Loading