Make cluster consensus, crash recovery, and restore crash-safe - #403
Merged
Merged
Conversation
This was referenced Oct 1, 2026
farhan-syah
force-pushed
the
fix/cluster-recovery
branch
from
October 1, 2026 11:40
8da09ef to
cde126f
Compare
The metadata log replays from its start on every boot. A delete proposed against one incarnation of a name must not remove a later incarnation of that name after a create-drop-recreate sequence. Stamp every delete with the modification_hlc (and descriptor version, where applicable) of the row it targeted at propose time, and refuse to apply it once the row has moved past that clock. Carries the fence through purge, sequence, trigger, function, procedure, materialized view, continuous aggregate, synonym group, topic, and vector index params deletes, across apply, post-apply dispatch, DDL compensation/reversal, and cluster replication. Adds a committed-only topic lookup and a compensation module that reverses finalized DDL under the same fencing rule when buffered DML fails to dispatch.
Raft and consensus - Leaders serve linearizable reads from a lease, and refuse competing votes while the lease holds; check-quorum tracks per-peer contact. - Voting, snapshot install, and membership changes persist before they are acknowledged, and group mount and unmount are explicit steps. - Leader hints, leader balancing, and routing persistence follow elections, and redirects name a hosted group. Calvin sequencing and snapshots - Multi-part transactions are collected, completed, and garbage collected by the sequencer, with a bounded entry size. - Sequencer state is captured in snapshots and restored on a joiner, including after the log has been compacted. - Schedulers gate on snapshot install and metadata catch-up. Crash recovery and replay - Metadata replay fences deletes to the incarnation they targeted and restarts from a durable floor. - Redo images, write-set journals, and hash chains replay after a crash, and snapshot install is staged so an interrupted install recovers. CDC and the trigger lane - Change events carry transaction boundaries and replay from the WAL. - Consumers resume across leader changes, and sinks have a single owner. - Trigger firing runs in a dedicated lane that survives failover. Placement and routing - Placement reconciliation, leader preference, and membership-aware routing track live peers; descriptor leases reclaim from dead holders. - Graph, array, and timeseries reads route to the owning home. Backup, restore, and PITR - Backups are cut consistently per database, scheduled, verified, and stored locally or remotely. - Restore reissues arrays and graph edges, protects the surrogate floor, resolves bind conflicts, and retries safely. - Point-in-time restore uses restore points and WAL time anchors. Tests - Cluster, in-process, native, and wire suites cover each area above, with matching harness support and nextest config.
Update the architecture, query language, real-time, and security docs and the changelog, and adjust the CI dispatch and reconstructed-SQL checks to match the code changes.
An empty data directory previously resolved the incarnation file against the process working directory, minting an incarnation in an unrelated location. Return a storage error instead.
The presence check and the old-expiry read now share a single entry metadata lookup.
retry_through_drain re-runs an operation refused with a retryable schema change until it succeeds, fails terminally, or a caller-supplied budget elapses. Each wait ends early when a drain ends on this node.
An ILP batch now projects its inferred time column apart from the tag and field columns. The catalog merge drops it for a timeseries collection whose declared time key names one of its fields, so a steady ingest stream no longer proposes a new descriptor version on every flush. The batch flush also takes its write leases through retry_through_drain, so a flush that meets a descriptor drain waits it out instead of dropping the connection.
farhan-syah
force-pushed
the
fix/cluster-recovery
branch
from
October 1, 2026 18:08
9fc96f8 to
8583388
Compare
Removing a vshard's sender now drops its armed catch-up, and a new scheduler registering a sender starts with none. The compaction floor counts only vshards with a registered sender, so an arm made by an exiting scheduler can no longer hold sequencer log compaction down.
…aker A receiver whose handler fails now answers with a typed refusal frame instead of dropping the stream. The sender treats any answer, a refusal included, as proof the link is up: it closes the peer's circuit breaker and does not resend the request. Only a failed connect, handshake, stream open, write, read timeout or lost connection counts as a failure. The breaker gains explicit admission, so an open circuit lets a single recovery probe through. Health pings use that probe path, and a ping refused by this node's own open circuit is no longer recorded as a ping failure. Per-peer Raft batches end on a link failure but move past a refusal, so one group a peer does not host yet no longer holds back another group's heartbeat.
Replacing a vshard's output sender no longer drops the catch-up armed for it. Only a vshard that leaves this node drops its catch-up, so the sequencer compaction floor keeps covering the replay a new scheduler still needs.
Raft: a successful AppendEntries response now reports the last entry the follower shares with the leader and holds durably, and the leader uses it as the follower's match index. Quorum contact tracking and the leader lease are reworked around this, and the contested-election marker is removed. Transport: concurrent sends to one peer share a single dial, and a connection is evicted only when it failed. Control plane: the auth lease, leased reads, the calvin scheduler, gateway routing, and graph dispatch read leaders through a new live-leaders snapshot that takes the Raft status before the routing guard. The lock order between MultiRaft and the routing table is documented to prevent a nested-read deadlock.
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.
A restart or partition no longer loses acknowledged state, and a restore no longer drops committed rows. Backups cut consistently per database, and point-in-time restore runs from restore points and WAL time anchors.
What changed
nodedb-raft/src/node/leader_lease.rs,nodedb-raft/src/node/quorum_contact.rsnodedb-cluster/src/raft_loop/,nodedb-cluster/src/group_disk/nodedb-cluster/src/raft_loop/nodedb-cluster/src/calvin/,nodedb-cluster/src/calvin/sequencer/nodedb-cluster/src/calvin/sequencer/state_machine/,nodedb/src/control/cluster/calvin_snapshot/,nodedb/src/control/state/calvin_bases.rsnodedb-cluster/src/multi_raft/proposals.rs,nodedb-raft/src/node/peer_contact.rsmodification_hlcof the row it targeted. Refuse the delete once the row has moved past that clock. Restart replay from a durable floor.nodedb/src/control/catalog_entry/incarnation/,nodedb/src/control/cluster/metadata_applier/nodedb/src/wal/redo/,nodedb/src/control/cluster/snapshot_install/nodedb/src/event/cdc/,nodedb/src/control/change_stream/nodedb/src/event/trigger/lane/nodedb/src/control/lease/,nodedb/src/control/cluster/nodedb/src/control/backup/cut_capture/,nodedb/src/control/backup/schedule/,nodedb/src/control/backup/verify/nodedb/src/control/backup/restore/,nodedb/src/ctl/restore/nodedb/src/control/pitr/,nodedb/src/wal/archiver/,nodedb-wal/src/record/restore_point.rsscripts/ci/check_authorized_dispatch.py,scripts/ci/check_reconstructed_sql.pynodedb-cluster-tests/,nodedb/tests/Why
Linked issues
wire_version.rsdocuments.Closes #165
Closes #166
Known limits
INSERT ... SELECT, upsert, CRDT, columnar, timeseries, spatial, array, and cross-home edge writes. A vShard with no running scheduler also skips it. Point writes and single-shard predicate writes take the locks.nodedb restore. SQLRESTORE DATABASErestores a logical backup and has no time target.RESTORE DATABASEchecks row counts and digests against the backup.nodedb restorechecks segment CRCs only, with no row tally.BACKUP DATABASEwrites a full logical envelope each time. Incremental storage comes from PITR base snapshots.last_applied.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