Fix Calvin dispatch capacity, typed errors across planes, and backup/restart gaps - #392
Merged
Merged
Conversation
commit_resolve.rs held apply_tail, verdict, and vote handling in one 736-line file, over the per-file size limit. It becomes a commit_resolve/ directory with apply_tail.rs, verdict.rs, and vote.rs, each keeping its existing logic. The driver core tests also duplicated build_test_scheduler and make_sequenced_txn across catch_up.rs, process.rs, and scheduler.rs. They move into a shared test_support.rs so each test module imports the fixtures instead of redefining them.
The bridge dispatcher had no way to signal "not enqueued, retry later" distinct from a terminal failure. It gains `DispatchCapacity`, a new error carrying which limit refused the request (tenant in-flight cap, a suspended per-database virtual queue, or a full per-core queue), and notifies waiters once an abandoned or completed request frees a slot. The Calvin scheduler is the first caller that cannot tolerate a terminal refusal for sequenced work: every replica must apply a sequenced txn, so a capacity refusal now parks the request in a FIFO (`deferred.rs`) instead of aborting it. The txn keeps its locks and its pending entry; the run loop re-sends parked requests once capacity frees, refreshing the deadline and group-leadership flag on each resend. Catch-up replay stops at the first deferred dispatch and re-arms itself at that Raft index instead of replaying the whole range again. `dispatcher.rs` splits into `enqueue.rs` (admission and slot accounting), `refusal.rs` (the refusal outcome type), and `response_poll.rs` (response draining), with shared test fixtures moved into `test_requests.rs`. The new error variant is wired through the data-plane wire format, the classify table (mapped to the retryable server-overload class), and every gateway error map (HTTP, pgwire, RESP, native).
Replace the magic 1024 literal used at every Dispatcher::new call site with a documented DATA_PLANE_QUEUE_CAPACITY constant, so downstream code (scheduler backpressure tuning) can reference the same bound.
Stop the run loop from reading new sequenced input once a dispatch is deferred at capacity, or once the in-flight backlog (pending, blocked, dependent-barrier txns) sits at the dispatcher queue bound while some of it can only drain on executor responses. A backlog of blocked txns alone keeps intake open, since the reservation release they wait on arrives as input. Cap each catch-up drain to a bounded window of committed sequencer log entries instead of the whole armed range, and resume from the first unreplayed index on the next drain, so a closed intake gate cannot make one drain read unbounded history. Expose the gate state, backlog depth, and closure reasons as scheduler metrics.
Request::deadline was read directly by DeadlineCheck and ExecutionTask::is_expired, so Calvin applies, replicated applies, replay, clone, and checkpoint work could be dropped as DeadlineExceeded once their envelope deadline passed. That work is already ordered and every replica must run it to completion, or replicas diverge. Add Request::execution_deadline, which returns None for Admission::Exempt(ExemptReason::AlreadyOrdered) and Some(deadline) otherwise. Route every deadline check through it instead of the raw field, and carry the parent's admission (not just its deadline) into transaction sub-plan tasks and exec_tx_passthrough so a sub-plan inherits its parent's exemption.
exec_tx_passthrough hardcoded DatabaseId::DEFAULT and vShard 0 for every sub-plan request, so a predicate UPDATE/DELETE issued inside a transaction against a non-default session database silently applied to the default database's collection instead. Extract the sub-plan request-building shared by exec_tx_passthrough and build_dummy_task_at into a SubRequestScope that copies the parent's database, vShard, deadline, and admission, and use it from both call sites.
Every replica must mark a sequenced txn applied only after it applied on this replica, or after an abort every replica reaches identically. A replica-local infrastructure error (dispatch refusal, a disconnected executor response, a failed resolve/flush/stage, a failed identity bind, or a failed WAL append) was neither, but the scheduler had no way to stop rather than silently diverge from its peers. Add a HaltLatch to each Scheduler that records the first such failure, holds the stuck txn's locks and pending entry in place, closes intake via IntakeClosure::ApplyHalted, and stops the deferred re-send and catch-up drain, while continuing to route responses, verdicts, and promotions for every other in-flight txn. The first cause wins and is recorded as a CalvinApplyHalt on the node-wide CalvinApplyHaltMarker (exposed through SequencerHaltMarker::apply_halt), which /healthz and the native status report read the same way they already read a halted sequencer or a wedged metadata applier. A halt during node shutdown is attributed to draining, logs at info instead of error, and sets no node marker. Expose the new gauge reason and step labels as scheduler metrics.
…pplied propose_sequencer_entry sent a vote, completion ack, OLLP mismatch, or routing-failure signal to the sequencer group once and dropped it on any failure. A refused proposal or a leader change that dropped an appended entry then left every participant waiting for a verdict or ack forever, since nothing proposed it again. Add an owed-entries table that keeps each proposal until this node's completion registry shows its effect, and a stall-tick sweep that proposes every remaining owed entry again. Read progress from a new CalvinCompletionRegistry::participant_progress so the sweep can tell an applied entry from an outstanding one. Route every proposal through a new SequencerProposer seam instead of locking MultiRaft directly: RaftSequencerProposer appends locally on the sequencer leader and forwards to it otherwise, reusing the data-group forward RPC. Extend DataProposeRequest's target from a bare vshard_id to a ProposeTarget enum (VShard or Sequencer) so the same RPC and raft-loop handler carry both proposal kinds, and give each node one shared proposer so its forward concurrency limit bounds the whole node rather than one scheduler. Expose the retry counts as a per-kind Prometheus counter, and split the scheduler's flow-metric rendering into its own module alongside it.
manager.rs held the lock table, acquire, wound-wait, release, try_acquire, and introspection logic in one file. It becomes a manager/ directory with types.rs, acquire.rs, wound_wait.rs, release.rs, try_acquire.rs, and introspection.rs, each keeping its existing logic.
apply_loop.rs held the per-batch driver, Array/Calvin/write dispatch paths, applied-floor bookkeeping, and shared helpers in one file. It becomes an apply_loop/ directory with driver.rs, array_dispatch.rs, calvin_read_result.rs, write_dispatch.rs, bookkeeping.rs, and helpers.rs, each keeping its existing logic.
funnel.rs held write admission, WAL redo append, Data-Plane dispatch, and response collection/post-apply steps in one file. It becomes a funnel/ directory with admission.rs, wal_append.rs, dispatch.rs, driver.rs, and response.rs, each keeping its existing logic.
dispatch.rs held the dispatch entry points, per-task routing decision, Raft-replicated and local Data Plane submission paths, and the authorization helper in one file. It becomes a dispatch/ directory with entry.rs, routing.rs, replicated.rs, local.rs, and authorize.rs, each keeping its existing logic.
calvin.rs held the static-set, passive-participant, active-participant, flush, and discard handler logic in one file. It becomes a calvin/ directory with static_stage.rs, active_passive.rs, flush.rs, discard.rs, and shared.rs, each keeping its existing logic.
A proposal re-proposed after a leader change can commit at two Raft log indexes. Applying both copies double-counts every non-idempotent effect: a materialized-sum fold, a columnar append, a timeseries ingest. RecordHeader gains an `apply_key` field (replacing the unused `reserved` bytes) carrying the idempotency key of the proposal whose apply appended the record, durable in the same write and covered by the CRC. A new `ProposalApplied` record marks an apply that writes no record of its own. `ProposalLedger` tracks applied keys per data group, bounded in memory and recovered from WAL replay on restart, so the apply loop can recognize and skip a duplicate copy and answer its waiter with the first copy's outcome. Building on this, committed transactions now replicate as a single `TransactionRedo` record carrying post-images, applied identically on every replica via `ApplyTransactionRedo` instead of re-resolving each sub-plan per node. Materialized-sum resolution travels with the redo as `RedoSumTargets` so replicas fold source writes into their targets without re-running the join. Columnar row updates now carry the old row's cross-engine surrogate forward to the replacement row, and the vector, columnar, and array engines gain rollback/truncate primitives so a partially-applied write can be withdrawn cleanly. The `resolve/columnar.rs` transaction resolver is renamed to `resolve/timeseries.rs` to match the timeseries-specific resolution it now holds, with the columnar resolution split out separately. A new cluster test verifies a proposal committed twice applies its effect once.
…etail INCR/INCRBY/DECR/INCRBYFLOAT previously computed on whatever bytes were stored, silently producing garbage on a non-numeric or out-of-range value. CounterFault now names the four ways a counter atomic computes no value (not an integer, not a float, integer overflow, non-finite float), each mapped to the RESP client text and SQLSTATE a caller expects. A new float_text module adds INCRBYFLOAT's decimal exactly, matching Redis's trimmed-decimal output for values that fit a Decimal, falling back to f64 outside that range. KvCounterShape travels with a counter atomic from planning through every path that computes its value: the live handler, transaction staging, resolve, and WAL replay. It tells the Data Plane, which does not see the catalog, whether the key belongs to a raw single-value collection or a typed row (and if typed, which numeric column moves and what template the fresh row uses). Structured error details now travel across the wire: ErrorPayload and NativeResponse carry an optional ErrorDetails, and NodeDbError::from_wire_with_details rebuilds a typed error from a carried detail when it matches the error's category, falling back to the code-only reconstruction otherwise. NodeDbError::kv_counter_fault turns a CounterFault into the OVERFLOW or TYPE_MISMATCH error a caller already handles. Documentation and RESP/native tests cover the new fault text and SQLSTATEs, a fresh key's stored shape, and a bare raw value read through every read-modify-write.
…pender WalManager::with_apply_key scoped a proposal's idempotency key to the calling thread for the duration of an append closure, so a reader of an append call site could not see which key its records carried, and an append inside the wrong scope silently picked up the ambient key. WalManager::appender(apply_key) now returns a WalAppender handle whose append_* methods carry that key explicitly. Every WAL append site, across WAL dispatch, the write funnel, write abort, Calvin recovery and commit resolution, the surrogate appender, checkpoint and collection-tombstone paths, and bitemporal purge, takes an appender up front instead of pairing a WalManager reference with a closure. NO_APPLY_KEY names the key of a record no replicated proposal owns.
A rollback that fails part way, or a committed write whose post-install work (a memtable flush, an artifact republish) fails afterward, leaves a core holding state that neither matches its WAL nor can be undone. Continuing to serve reads and writes from that core would answer against state no replica or restart reproduces. CoreFailStop latches a core the first time either cause is observed: it logs an ERROR, files a diagnostic report, and refuses every queued and future request with RetryableRefusal until a restart rebuilds the core from the WAL. The latch is exposed as the nodedb_data_plane_core_fail_stopped gauge, folded into the /healthz readiness body and the native STATUS command, and gates opportunistic checkpointing and maintenance so a stopped core publishes no artifact of its unknown state. Reaching that guarantee required each engine's rollback to reverse exactly what it wrote instead of approximating it: - The graph CSR index gains an interning-aware restore module: exact edge and weight reversal, and withdrawal of an interned node or label only when it is the newest entry and nothing still refers to it, refused otherwise via a new GraphError::WithdrawRefused. - The edge store's temporal writer gains a revert module mirroring each bitemporal write with its exact undo, and node-edge cascade delete now runs through one cascade helper shared by point-delete and periodic sweep so a failed store write leaves the CSR and edge store still agreeing. - The IVF vector index gains roll_back_to, withdrawing every vector added after a mark and dropping training state the mark predates. - The transaction undo pipeline is restructured accordingly: undo/ apply.rs shrinks to dispatch, with edge and vector undo split into their own modules, and each engine's undo entry now carries the detail a failed reversal reports. The redo-apply and WAL-replay paths pass through the exact-restore data these rollbacks need, and SeriesCatalog gains Clone/PartialEq so a timeseries rollback can snapshot and compare its catalog state.
Introduce OutcomeFloor: the highest WAL LSN at or below which every record dispatched to a Data Plane core has a final outcome (applied, or refused with a durable WriteAborted marker). A write opens a window before it mints its LSN and settles it once the outcome is final; the dispatcher opens a window for every accepted request that carries a WAL LSN and settles it on final response, dead core, or abandoned drain. Every request pushed onto a core's ring now carries the floor read at enqueue time, and each core keeps the highest floor it has read via a new AppliedPrefix tracker, surfaced in checkpoint snapshot logging. Wire the funnel to open a window per control-plane write and abort undispatched writes whose dispatch was refused before settling. BridgeRequest and its constructors change shape throughout the test suite to carry the new field via BridgeRequest::unfloored() where no floor tracking is under test.
handle_permission_event used a non-blocking try_write and silently dropped the update on contention with no later repair, since each grant or edge event is the only update for its row. Make it async and wait for the write lock instead. Also route the WAL-catchup path's lone re-dispatched event through the same Normal-mode batch handling as the events that follow it, so a grant or hierarchy row consumed there reaches the permission cache too rather than only the trigger and CDC side effects.
…l-refusal replay Introduces an outcome-floor mechanism that tracks in-flight write windows (minted WAL records not yet resolved) so restart replay and checkpoint truncation never advance past a write whose outcome is still undecided. Every mint site opens a window, and settles or holds it on every exit path including panics, deferrals, and Calvin scheduler halts; a dropped window that never settled or held is reported as a diagnostic leak. Handlers that commit rows incrementally (columnar ingest, bulk update/ delete, CRDT apply, graph edge writes, vector writes, timeseries ingest) now track whether any row landed before a later failure, and answer with a refusal code that keeps their WAL records instead of one that claims nothing applied — otherwise the Control Plane would cancel records for a write that partially landed, and recovery would silently drop it. The same correction applies to transaction sub-plans with no engine-specific undo handling and to Calvin's redo-record bookkeeping. Adds `ErrorCode::ExpiredBeforeExecution` for requests whose deadline passed before a core started them, distinct from `DeadlineExceeded` for a task that ran partway; both surface as the same query-cancelled error to clients. Replaces the WAL appender's apply-key thread-local with an explicit appender parameter, and extracts gateway response-shaping into a dedicated module so remote and local dispatch share one response shape.
…er read gate The harness built `SharedState` without installing the gateway or a Raft read gate, so remote-routed dispatch and linearizable reads behaved differently under test than in production. Every server-start path now calls the same `install_gateway` production boot runs, and a routed harness installs a `SingleVoterReadGate` that answers leadership and read index the way a group with one voter does, since the harness runs no Raft loop. A routed harness also derives its `node_id` from the routing table's sole leader instead of leaving it unset.
Bind each rejected delta's dead-letter entry to the log position of the record that produced it, and store it in the sparse engine's redb table. A tenant's CRDT engine restores its entries from storage when created, so restart replay — which never reaches a rejected record — no longer loses them. Re-applying the same record binds to the same position, keeping one entry per record instead of accumulating duplicates. Route every rejection path (snapshot import, sync apply, local apply, transaction batch, WAL replay) through the same store-then-respond sequence, and report a store failure as part of the refusal instead of silently dropping it. Purge stored entries when a collection is dropped or a tenant is purged, reusing the bounded reclaim retry now extracted into its own module.
A caller dropped while its write waits for admission, or before a dispatch result arrives, left a set of minted WAL records with no path to settle their outcome-floor window: neither the writer's own cancel logic nor the eventual response ever ran. Give MintedRecords a Drop impl that cancels the records in place with a WriteAborted marker and settles the window when no core has been given the request yet, or holds the window and reports it when a core may already hold the records. Track which record a recording WAL appender wrote through RecordedAppend (LSN plus the tenant/vshard/database that placed it) instead of the bare LSN, so the drop path can append its own cancel markers without the caller re-deriving that context. Mark records as sent once a core has been handed the request, at every call site that dispatches to one, so a later drop knows to hold rather than cancel. Fold the health endpoint's outcome-floor-stuck body into its own function so it can be exercised directly by tests, and add a failpoint on the write-aborted WAL append to exercise the failure path where a cancel marker cannot be written.
Its one caller always passed None.
send.rs mixed the pooled single-attempt RPC path with the outbound streaming-shuffle helpers. Move send_shuffle_push, open_shuffle_push_stream, and ShufflePushStream into a dedicated shuffle_push module and re-export it from transport::client, leaving send.rs to the RPC send path plus the connection-target verification it now shares via a widened visibility.
Origin no longer installs a committed transaction by re-executing its sub-plans as a batch; it installs the transaction's redo record instead. Update MetaOp::TransactionBatch, CalvinExecute and StageWrite doc comments to describe the staged/redo-apply flow, and derive Eq on the document sum-target keys the redo path now compares. Add SortedIndexRead/SortedIndexSpec so a sorted-index read inside an explicit transaction can be answered from a transaction-local tree built from base rows with the transaction's staged writes folded in, routed as a new KvOp::SortedIndexTxnRead. Move KV counter/atomic-op computation (integer and float-text arithmetic, counter fault classification) out of the nodedb binary crate into a new nodedb-physical::kv_atomic module shared across engine code, pulling in rust_decimal for its float-text parsing. Add DataPlaneErrorCode::SyncRejected to carry a sync frame's rejection provenance across the wire, and box ClusterError::StreamTerminal's error payload now that TypedClusterError has grown larger variants. Add the active_sql_transaction SQLSTATE and rename the transaction-batch fail point to match the new Calvin overlay-stage code path.
A committed transaction used to re-execute its buffered plans as a 'sub-plan batch' at COMMIT — the sole durable apply path, with per-engine sub_plan_* modules translating each staged plan back into engine calls. Replace that with the staged/redo-apply pipeline already used for non-transactional writes: staged plans resolve into a redo record once, the record installs on every path (autocommit, explicit COMMIT, Calvin), and CalvinExecute now stages under lock instead of executing eagerly, resolving to redo only once the global verdict commits. Delete the sub_plan, sub_plan_doc, sub_plan_kv*, sub_plan_write, sub_request, batch/batch_crdt/batch_irreversible, index_write_values and resolve/vector_direct modules; extend redo_apply, resolve and stage_write to cover the cases they used to handle (document batches, timeseries ILP staging, columnar insert undo). Track write-outcome ownership end to end: add a bridge ClosedLsns range set and OutcomeFloor::ResendRefusal so a resend of an LSN whose outcome is already final, or one minted outside a live write window, is refused instead of replayed twice. Give Calvin's overlay-stage handlers their own calvin_reply module (stage, images, target, flush_read, reply) and split CoreLoop's Calvin and maintenance state into calvin_fence/calvin_state/ maintenance_state modules. Move KV counter and float-text arithmetic out of engine/kv into the new nodedb_physical::kv_atomic module, deleting engine_atomic_compute.rs and float_text.rs; give the sorted-index engine its own index module. Add transactional index DDL and sorted-index reads: index DDL statements inside an explicit transaction buffer their catalog entries and defer their engine side effects (secondary-index backfill, KV index build/drop, sorted-index tree build/drop, FTS analyzer bind) to COMMIT via a new DeferredDdlEffect/deferred_effects module, and a sorted-index read inside a transaction answers from a transaction-local tree via the new KvOp::SortedIndexTxnRead and kv/sorted_txn handler. Expose both as native opcodes (index_ddl_op, sorted_read_op) so the wire protocol reaches the same catalog and engine path as the SQL statements. Add ErrorCode::SyncRejected carrying the same rejection provenance as the cluster wire type. Update transaction, Calvin, WAL replication, undo, and crash/inproc/ native/wire test suites for the staged-redo install path, and add a single-core commit-plans test driver in nodedb-test-support.
…bove set A checkpoint used to record only the highest LSN applied. LSNs are node-global and reach a core out of mint order, so a record with a lower LSN can still be in flight when a higher one applies. A max-applied stamp then claims the lower record and restart replay skips it, losing the write. Replace the single LSN with a ReplayStamp: a `prefix` (the core's outcome floor when the artifact was written, below which every record has a final outcome) plus `applied_above`, the disjoint LSN ranges applied above that floor. Replay skips a record exactly when the stamp says so. Add LsnRanges, a BTreeMap-backed disjoint-range set that only merges touching LSNs and never bridges a gap, to track applied LSNs above the floor. Wire the stamp through every checkpoint format (columnar, KV, graph-label, sparse-vector, spatial, sync-hwm, vector) and through wal_replay_all and replay_floors so restart replay gates on it. Add an async fail-point action, WaitForFile, so a test can park one in-flight write at the funnel gate without blocking its core, plus a Control-Plane fail_gate module that calls it after a logged write's WAL append and before dispatch. Add a crash-replay test proving a sorted-index write parked at the gate while later writes are checkpointed still applies after a crash and restart.
KV counter atomics, sorted-index registration, rate-gate counters, and weighted-pick's audit write used to dispatch straight to the Data Plane, bypassing the write-admission funnel and Raft. That write has no WAL record and no replica ever sees it: a crash loses it, and a follower's KV_INCR answers with a stale value. Add a durable-write module (`dispatch_utils::durable_write`) that proposes a write through Raft when this node runs a proposer and the plan is replicable, and otherwise dispatches it through the funnel's `AppendHere` route under the write-admission guard. Route every planned and hand-built autocommit write through it: `dispatch_authorized_durable_write` for tasks that still need the clone-write and authorization gates, `dispatch_durable_autocommit_write` for one already built as an `AutocommitWrite`, and `dispatch_authorized_task_by_class` for a transport that dispatches reads and writes through one call site. Add `refuse_unlogged_write` at the read-dispatch boundary so a write that reaches a route with no WAL append is rejected before it applies, rather than silently losing durability. Extend the in-transaction staging gate's `InTxnRoute` with an `Autocommit` variant so a write outside a transaction block (or one a transaction cannot buffer) is distinguished from a read at the gate, and every caller that matched `InTxnRoute::Read` for both cases now routes the write side through the durable path. Publish the change-feed event for a replicated write from the proposing node in both the gateway and the executor's replicated-entry path, since a replica applies with `ChangeFeedOwner::Unowned`. Harden `weighted_pick`'s row scan and audit write to surface a refusal or a decode failure as an error instead of silently returning an empty result or logging and continuing. Add a `CrashHarness::standalone` boot mode (no Raft proposer, so autocommit writes take the local funnel route) and cover the new path with crash, in-process, native, and cluster-replication tests for KV counter atomics.
Move the applied-prefix stamp type to types::replay_stamp so vector, array, and timeseries checkpoints can all publish an exact record of what they hold, instead of a single high-water LSN. A single LSN wrongly claims out-of-order records that mint below it but arrive after, so restart replay skips them and the write is lost. Vector and array checkpoint manifests now carry a ReplayStamp and gate recovery on it. Timeseries collections get an analogous TsReplayStamp that names applied rows and effective truncates per partition and per collection directory, written before each flush and before retention removes a partition, so replay never re-applies or drops a record. Timeseries ingest also groups a committed record's sub-records so a memtable flush never lands mid-record, and retention enforcement moves into its own dispatch handler.
Add KvPush/KvPushAck wire messages so Lite forwards KV writes to Origin, and RowPushReject so Lite can refuse an Origin row push it cannot apply. Thread SyncProvenance through KvOp::Put/Delete so a pushed write clears the sync idempotency gate instead of bypassing it, and introduce SyncHold (Duplicate/Fenced/Gap) to distinguish a held-back frame from an outright rejection across the bridge and cluster RPC codec.
Boot now binds every protocol socket before it waits on cluster readiness, so a port conflict fails boot before anything is exposed, then opens the sockets for accept only after entering a new terminal Serving startup phase. A bound-but-not-listening socket refuses connections at once instead of leaving a client to wait out the rest of boot in the kernel's accept queue. The HTTP listener is the exception: it starts serving early so orchestrator probes can watch startup, gated by a new middleware that lets through only the health and metrics routes until Serving. Reaching authorization readiness during cluster-ready no longer waits on the bounded lease alone: a node that leads the metadata group as its only voter now qualifies through a pinned lease with no expiry, since a second voter would end that leadership before it could vote. Lease status reporting and Prometheus rendering account for this new SoleVoter state.
Give EvalError three new variants alongside DivisionByZero: UnknownFunction, VectorDimensionMismatch, ArgumentType, and InvalidJsonPath. A call to an unregistered name, a vector distance over mismatched dimensions or a non-vector operand, and a document function given a malformed JSONPath now surface as typed errors instead of silently evaluating to NULL. The new ErrorCode::DataException (SQLSTATE 22000) and NodeDbError::data_exception carry these through the bridge, cluster RPC codec, plan error map, and pgwire/native/HTTP error rendering, alongside the constant-fold path and every enforcement/executor site that used to hardcode DivisionByZero on any evaluator error. Add vector_distance/vector_cosine_distance/vector_neg_inner_product, doc_get/doc_exists/doc_array_contains/nav document accessors, and ndb_chunk_text as real per-row evaluators, plus make_array for non-literal ARRAY[...] elements. Rewrite sql_like_match as a tokenizing matcher shared by the scan-filter and scalar like/ilike paths, adding escape-character support. Add planner::search_scope::refuse_row_scoped_search_functions so an index-owned search function (bm25_score, text_match, rrf_score, sparse_score, ...) used outside a search plan is refused at plan time with SqlError::SearchFunctionOutsideSearch instead of resolving to a NULL score column.
Replace assert_eq! panics across the vector engine (flat index, HNSW, IVF-PQ, Vamana, PQ/SQ8 quantization, matryoshka, multivec, NAViX, SIEVE) with typed VectorError results. Split DimensionMismatch (bad caller input) from a new StoredDimensionMismatch (corrupt or foreign stored data) so the two are classified differently: an input error maps to SQLSTATE 22000 (data exception) and never takes the core down, while a stored-data mismatch keeps fail-stop handling. Add InvalidFilterBitmap for pre-filter bitmaps that fail to decode and InvalidInput for index build/train calls with unusable parameters. Split IVF-PQ search out of vector_search_exec.rs into its own handler module. Add a generated_always SQLSTATE and a shared constraint_sqlstate mapping so RejectedConstraint kinds (not_null, unique, generated_always, fk_missing, rls_policy, permission_denied) each keep their own error class instead of collapsing to unique_violation.
FTS query analysis, corpus stats, facet counting, and staged spatial row decoding used to swallow storage/analyzer errors with ok()/ unwrap_or_default(), silently degrading hybrid, triple, graph-RAG, and facet searches to partial results. They now propagate the underlying error through the search response instead.
The IVF-PQ index type used to build once from a static snapshot and serve reads only. It now buffers inserted vectors, searched exactly, until it holds max(ivf_cells, pq_k) live vectors, trains its k-means cells and PQ codebooks on them, and from then on supports insert, soft-delete, and search with cell-probe plus exact rerank. - nodedb-vector/src/ivf splits into mod.rs, index.rs, kmeans.rs, params.rs, search.rs, and checkpoint.rs; each entry carries a caller-assigned id so a collection can move its vectors under the ids they already have. - nodedb-vector/src/quantize/pq_kmeans.rs carries the k-means implementation shared by PQ codebook training and IVF cell training. - VectorCollection tracks IVF training state (collection/ivf_mode.rs) and settles (trains or reseeds) a filled index on transaction commit instead of only sealing HNSW builds. - The executor wires insert, delete, undo, WAL replay, snapshot restore, compaction, and reindex through the new lifecycle, and vector_settle.rs replaces vector_search_ivf.rs as the settle path. - VectorIndexStats reports IVF training state (threshold, trained, cell count, nprobe) through SHOW VECTOR INDEX. - FTS and spatial search legs used to swallow errors with ok()/ unwrap_or_default() and fold to empty rather than fail; they now propagate the underlying error through the search response.
The HNSW builder thread used unbounded channels and blocking sends, so a core could stall on a full completion queue and a failed insert only logged and skipped a vector, silently shifting every later node's id. The builder now takes bounded request/completion queues, builds vectors in order and fails the whole build on the first insert error, and reports success or failure back through a new VectorBuildQueue that backlogs requests the queue can't take yet and drains completions once per tick. Sealing a growing segment, settling after boot, and REINDEX CONCURRENTLY all route through this one queue and its collection.rs::build install path, replacing REINDEX's separate rebuild-on-a-plain-thread codepath and the old blocking send on seal/settle. VectorCollection tracks per-key completed/failed build counts and a configurable seal threshold (VectorTuning), surfaced through SHOW VECTOR INDEX and Prometheus. Fix a duplicate-node bug: a vector write over an already-indexed row left the row's original node live after a subsequent delete, since the delete looked up the surrogate's originally recorded node instead of the node a later write rebound it to. INSERT also stopped double-indexing fields a vector index already covers via the document write. Add ErrorCode::BadRequest so the Data Plane can report a malformed request (SQLSTATE 42601) the same way the Control Plane already does, and route CollectionDeactivated/CalvinSerializationConflict/SourceFrozen/ FeatureNotSupported/CrossCollectionNotColocated through the existing NotFound/ConflictRetry/Unsupported classes instead of falling through to Internal.
REINDEX and the vector settle/seal path used to rebuild the graph CSR adjacency index and the FTS inverted index inline, blocking the core for the duration of the build. Both now build a shadow copy off the core while it keeps serving reads and writes, journal writes made during the build, replay the journal onto the shadow copy at cutover, and swap it in atomically. A write that diverges between the live and shadow index, or that overflows the journal's byte bound, discards the shadow copy and leaves the live index untouched so REINDEX can be retried. - nodedb-graph gains csr::rebuild (seed, journal, install) and new GraphError variants for a rebuild already running, a superseded rebuild, a journal overflow, a replay divergence, and an invalid snapshot. - nodedb::engine::sparse::inverted gains the matching rebuild_snapshot/rebuild_journal/rebuild_install modules for the FTS posting index. - The executor's reindex handler is split into a control/reindex/ module (dispatch, csr, fts, pending holds, waiters) replacing the old single-file reindex/reindex_apply handlers, and MetaOp::RebuildIndex's contract is updated to describe the non-blocking build, deadline-based wait, and cutover semantics. - Point and bulk-DML updates, snapshot restore, and CONVERT now route text-index maintenance through the same rebuild-aware path (update_reindex_text, snapshot restore/text.rs) instead of updating the FTS index inline, and CONVERT reports a schema mismatch as a data-exception naming the row and column instead of a generic internal error. - Diagnostics record index-rebuild lifecycle events (diag/context/recording index_rebuild) alongside the existing vector-build diagnostics.
Data-Plane refusals now carry a typed cause end to end: a phase error such as MOVE_TENANT_SNAPSHOT_FAILED keeps its own code while the refusal that caused it rides alongside as ErrorCausePayload, rebuilt client-side as NodeDbError::cause. DDL apply, system dispatch, and the native/pgwire/HTTP gateways now render a Data-Plane ErrorCode's own SQLSTATE and class consistently instead of folding refusals to Internal, and native error codes are referenced by their named constants instead of magic numbers. Adds the ProgramLimitExceeded error variant for statements that exceed a server size or depth limit.
…collection key Introduce CollectionKey, the only (database_id, bare_name) pair the vShard hash and surrogate allocator accept. Every caller that used to hash a raw string, or DatabaseId::from_collection_in_database, now builds a CollectionKey via from_bare or from_qualified so a database-qualified name can never reach the hash and route rows to the wrong vShard or bind surrogates under the wrong key. Thread the type through nodedb-cluster routing, the Calvin sequencer and its diagnostics, physical layer surrogate assignment, control plane bind/dispatch paths, and the cluster and inproc test suites. Add a diagnostic capture for a Calvin batch whose participant set cannot be derived because a key-set collection name is missing its database qualifier.
Escalate the renewal loop's re-acquire and release failures from a warn log to error plus a faultbox Capture, grouped by step and error class so a retry every tick files one report instead of storming. An unrenewed lease expires under a holder that still plans against it; an unreleased one blocks every DDL drain on its descriptor.
…er and DDL paths Array cluster writes and reads, backup restore, merge/update-from-join resolve passes, calvin pre-execution scans, materialized-sum reconnaissance reads, and several DDL neutral functions (rate gate, weighted pick, KV index, collection purge) now propagate a refused Data-Plane response as its own typed ErrorCode instead of folding it into a generic Internal or Storage error. A shard's Data-Plane verdict crosses the cluster wire as a new VShardRefusal RPC frame carrying the coded ErrorCode verbatim, so the coordinator rebuilds ClusterError::DataPlane and renders the same SQLSTATE a single-node execution would. Such a refusal never counts against the array shard's circuit breaker, since it comes from a shard that answered. HTTP DDL errors now derive their status from the SQLSTATE through a shared sqlstate-to-HTTP status table instead of a two-way if/else, and carry their typed cause in the response body via a new HttpError.cause field, matching native and pgwire. RetryableSchemaChanged now classifies as a write conflict and renders 40001 (serialization failure, retried by drivers) instead of an internal 503/XX000. SessionTokenExpired now classifies as auth_expired and renders 401/28000 (invalid authorization) instead of a bad request, on native, pgwire, and HTTP alike.
Generalize VShardRefusal to carry any ClusterError as a ShardErrorWire mirror instead of only a Data-Plane verdict. WrongOwner, Raft redirects, Codec, and Transport errors now cross the wire typed and rebuild on the other side; an error with no wire mirror falls back to RemoteUntyped, keeping its message instead of collapsing to a closed stream.
Route every gateway error's HTTP status through the pgwire SQLSTATE table (GatewayErrorMap::to_http / sqlstate_to_http), so HTTP, native, and pgwire answer one class for one error; drop the now-redundant remote-code-to-HTTP table. ApiError::from derives its status the same way, and shape_error_to_api and the '3D000' database-not-found case follow the same table instead of their own if/else. Array cluster execution rebuilds a shard's typed error instead of only its Data-Plane code: WrongOwner reroutes to NotLeader/NoLeader, a local dispatch or channel timeout crosses as the typed DeadlineExceeded verdict, and a cross-shard write reports its target's refusal reason instead of a generic 'unexpected RPC response type'. A commit abort on a dispatch or DDL-propose error keeps its own class via the new SystemTxnError::CommitFailed, and a Calvin cancel/timeout now renders the deadline class instead of a bare internal error. The rate-gate DDL functions propagate a dispatch error or coded refusal with its own class instead of silently reading it as an absent key or zero usage, and add BACKUP_TENANT_MISMATCH / BACKUP_KEY_MISMATCH to the numeric-code-to-SQLSTATE table.
Replace the last catch-all arms in the error-crossing paths with exhaustive matches over crate::Error and bridge::envelope::ErrorCode, so a new variant fails to compile here instead of silently falling into Internal or the wrong compensation/retry class. Covers the cluster propose path (new propose_error module), CRDT delta compensation hints (new compensation module), sync refusal/retry classification, gateway error mapping, array cluster execution, backup restore, DDL neutral handlers, and the recursive-value/point/transaction executor handlers. Split error_classify.rs into a directory: public.rs keeps the Error-to-NodeDbError table, unclassified.rs keeps the is-unclassified-failure predicate. Add numeric_sqlstate.rs so a numeric ErrorCode crossing from a remote node maps back to the same SQLSTATE class its local variant would have chosen. Extend nodedb-cluster's wire/circuit-breaker/vshard/data_propose types and nodedb-types' sqlstate table with the codes this now threads through.
Give DdlError first-class helpers for typed sources: from_error_in_context carries a crate::Error's own SQLSTATE/code behind a message prefix, internal marks a fault with no typed source (codec, storage, broken invariant) as XX000, and in_context prefixes a built DdlError's message without disturbing its class. Replace every neutral DDL handler's ddl_err(format!(...)) call site with one of these, so a catalog or propose failure keeps its real class (deadlock retry, privilege denial, not-found) instead of collapsing to a generic internal error. Add three error codes this now threads through: TRANSACTION_ROLLBACK for a Calvin participant abort with no read-set validated (retriable, unlike a write conflict), ACTIVE_SQL_TRANSACTION for a statement that cannot run inside an explicit transaction block (folds in NotInTransactionBlock, CrdtApplyForbiddenInTransaction, and CrossShardInExplicitTransaction), and DEPENDENT_OBJECTS_EXIST for a drop or revoke blocked by dependents. Give DROP ROLE a typed RoleInUse error carrying RoleDependents (Users or ChildRoles) in place of a formatted BadRequest, so it classifies as DEPENDENT_OBJECTS_EXIST/2BP01 like other dependent-object refusals. Extract the DROP USER reassignment path's OwnerKind enum into its own owner_kind module, shared unchanged by reassign_owned.rs.
… typed Generate ErrorCode::ALL from one macro-driven list instead of a hand-kept const block, so the code table cannot miss a constant. Complete code_for_sqlstate so every named SQLSTATE constant, not just the ones with one classification, maps to its numeric code, and fold the shared value, integrity, and data-exception SQLSTATE classes into their default codes. Type QUERY_CANCELED, BACKUP_KEY_MISMATCH, and the new AUTH_TOKEN_EXPIRED as AmbiguousSqlstate so a caller can no longer compare a bare &str against a code whose class depends on the call site. Add AUTHENTICATION_FAILED and INVALID_PASSWORD codes for credential failures, and DUPLICATE_OBJECT, SERVER_REJECTED_ESTABLISHMENT, and PROTOCOL_VIOLATION SQLSTATE constants. Add Error::Ddl, a typed variant that boxes a DdlError so a DDL step that fails at COMMIT after its statement already returned reports through the same SQLSTATE/code/detail/cause path as its autocommit form, instead of collapsing to a generic Internal error. Wire the new variant into every exhaustive match over crate::Error across the cluster, gateway, and sync layers. Add static_sqlstate to intern a DdlError's SQLSTATE as a &'static str once, since the error renderers take &'static str but DdlError only ever carries a bounded, server-controlled set of SQLSTATEs.
…ction
Data-Plane collections outside the default database are stored under a
database-qualified name ("{database_id}/{name}"), but WAL tombstones,
vShard routing, and backup/restore sections keyed on it inconsistently:
some paths compared the qualified name against the bare catalog name a
tombstone or route recorded, so a purge in a named database silently
failed to shadow its own WAL tail, and a snapshot section could route to
the wrong vShard. Introduce CollectionKey/QualifiedCollection as the
single source of truth for the two name forms and thread them through
TombstoneSet (nodedb-wal), vshard_of_stored (nodedb/backup/snapshot_keys),
WAL replay dispatch, catalog-entry tombstone application, and security
catalog tombstones, so every consumer converts through the same rules.
Add DatabaseTarget/RestoredName (backup/restore/target.rs) to resolve a
backed-up section's source-qualified collection name to its destination
name during restore, and DatabaseDataSection (backup_envelope) to wrap
each database's snapshot as its own envelope section instead of merging
every database into one flat section. Add backup/metadata.rs and
backup/node_snapshot.rs to carry per-database metadata through the backup
path, and restore/databases.rs plus restore/orchestrate/database.rs to
let a restore target one or more specific databases instead of always
restoring the whole cluster.
Add a regression test (crash_purge_not_resurrected.rs) that kills and
reopens a node after DROP COLLECTION ... PURGE in a named database, and
a cluster/wire pair (cluster_backup_restore_databases.rs,
sql_backup_restore_databases.rs) covering the new selective-database
restore path. Remove the now-redundant single-database backup shape test
superseded by the multi-database cases.
Move the wasm32/native split from per-function `#[cfg]` branches inside writer, double-write, and record code to `#[cfg(not(target_arch = "wasm32"))]` on the module declarations and re-exports in lib.rs and record/mod.rs. wasm32 now builds only crypto, error, secure_mem, and the record header constants; the writer, readers, segments, double-write buffer, replay, and recovery compile solely on native targets instead of each carrying a dead wasm32 arm. Drop the now-unreachable `WalError::Unsupported` variant along with the wasm32 fallback paths in `pwrite_all`, `DoubleWriteBuffer::recover_record`, and `scan_max_seq` that only ever returned it or scanned the buffer by hand. `writer/flush.rs` and `double_write/raw_io.rs` keep a single unix `pwrite` loop instead of a `#[cfg(unix)]`-gated one, and `write_error` drops its own now-redundant cfg. Update the crate- and module-level docs to describe the new split, and remove the wasm32 cfg on the `mmap_reader_madvise` test case now that the module it exercises is native-only anyway.
Backup capture, the cluster snapshot builder, PURGE TENANT, and the MOVE TENANT snapshot step no longer authorize through a SystemTask call site; they fan out through the all-cores exchange instead.
register_peers_from_topology idempotently registers a fan-out target's address with the transport from the live cluster topology. Move it from server/exchange/resolve/peers (a shuffle-join/aggregate helper module) to cluster/warm_peers/register, alongside the peer pre-warming it complements, and update every caller — the Calvin sequencer proposer, reservation and submit paths, auth lease leadership, and surrogate exchange resolve — to import it from its new home.
Move collect_under_deadline/collect_bounded_response, reject_data_plane_error, and dispatch_local_read out of server/dispatch_utils into a new top-level control/local_dispatch module. They dispatch to this node's own Data Plane (bounded response collection, error-response-to-typed-Err conversion, and read-only internal scans), a distinct concern from dispatch_utils' remaining cross-node/authorized-write helpers, and are used well outside the server tree (Calvin planner reconciliation and submit paths, the permission tree reloader, system_txn). Update every caller to import from the new module.
Introduce HopEncoder and ClassCase type aliases for the (name, fn(...) -> ...) and (fn() -> Error, sqlstate, code) tuples the class-parity tests build, replacing the inline tuple types repeated at each array declaration.
…le split
handler/dispatch.rs was split into handler/dispatch/{replicated,routing}.rs;
point check_authorized_dispatch.py's ALLOWED_REFERENCES entries for
into_physical_task and propose_replicated_entry at their new file.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Cross-core Calvin transactions no longer stall when dispatch capacity runs out. Every typed error keeps its SQLSTATE class on every plane and across nodes. Backup, purge, and restart keep every committed row.
What changed
nodedb/src/bridge/dispatch/,nodedb/src/error/dispatch_capacity.rsnodedb/src/control/cluster/calvin/,nodedb-cluster/src/calvin/sequencer/nodedb/src/bridge/,nodedb/src/control/server/dispatch_utils/minted/nodedb/src/wal/manager/,nodedb/src/control/distributed_applier/nodedb/src/data/executor/*_checkpoint/,nodedb/src/data/executor/core_loop/checkpoint_floors/nodedb/src/data/executor/nodedb/src/control/server/dispatch_utils/Error::Ddl.nodedb-types/src/error/,nodedb/src/control/gateway/error_map/nodedb/src/control/gateway/error_map/class_parity.rsCollectionKey { database_id, bare_name }is the only input to the vShard hash and surrogates.nodedb-types/src/id/collection_key.rs,nodedb/src/control/surrogate/Serving.nodedb/src/main_boot/listeners.rs,nodedb/src/main_boot/gates.rsnodedb/src/control/security/auth_fence/,nodedb/src/control/security/permission_tree/nodedb/src/control/security/auth_lease/nodedb/src/control/backup/,nodedb/src/control/server/exchange/all_cores/REMOVE EVENT.nodedb/src/event/,nodedb-wal/src/record/nodedb-vector/src/hnsw/build.rs,nodedb-vector/src/ivf/nodedb/src/data/executor/handlers/control/reindex/nodedb-query/src/functions/nodedb/src/engine/kv/,nodedb/src/control/server/sync/nodedb/src/engine/crdt/nodedb-cluster/src/raft_loop/register_peers_from_topologyand the local-dispatch helpers belowcontrol::server.nodedb/src/control/cluster/warm_peers/register.rs,nodedb/src/control/local_dispatch/nodedb-walon wasm32 builds onlycrypto,secure_mem,error, andrecord::header.WalError::Unsupportedis removed.nodedb-wal/src/lib.rsWhy
XX000for typed refusalsLinked issues
How to check it
cargo nextest run --workspace --exclude nodedb-cluster-tests --all-features --no-fail-fastcargo nextest run -p nodedb-cluster-tests --all-features --no-fail-fastcargo clippy --workspace --all-targets --all-features -- -D warnings.github/workflows/static-gates.yml