Skip to content

[Epic] P2 — Cluster Consensus Safety (v0.6) #165

Description

@habibtalib

Epic — v0.6 production-hardening, phase P2: Cluster Consensus Safety

Goal: a restart or partition never produces two leaders, divergent committed entries, or silent state-machine divergence; linearizable reads are actually linearizable.

Every box below is checked against this repository's main.

Critical (filed)

High

  • Deposed/partitioned leader serves stale linearizable reads — no ReadIndex heartbeat, no leader lease. Partial. ReadIndex closes the safety hole: start_read_index issues its own quorum round and returns ReadIndexStatus, which separates "not yet confirmed" from "leadership lost". f991044, b1f52e9. Dead follower_read.rs deleted. Check-quorum has since landed: a leader tracks when a majority last acknowledged it and demotes itself once that gap reaches an election timeout, nodedb-raft/src/node/quorum_contact.rs, refreshed from rpc/append_entries.rs:145. Remaining: leader lease, a latency optimization over the ReadIndex round, not safety. Done in Make cluster consensus, crash recovery, and restore crash-safe #403. Leader lease: nodedb-raft/src/node/leader_lease.rs serves lease reads and refuses competing votes (request_vote.rs, pre_vote.rs) while the lease holds.
  • BoundedStaleness freshness measures time-since-last-apply, not lag-vs-leader — a badly-lagging follower reports "fresh". nodedb-raft/src/node/staleness.rs returns Behind when last_applied trails the leader's commit index. b1f52e9. Dead closed_timestamp.rs deleted.
  • Rolling upgrade blocked on any wire bump (MIN_WIRE_FORMAT_VERSION == WIRE_FORMAT_VERSION) → version hard-partition mid-upgrade. Not started here. nodedb-types/src/wire_version.rs:52 still reads MIN_WIRE_FORMAT_VERSION = WIRE_FORMAT_VERSION. Closed by design with Make cluster consensus, crash recovery, and restore crash-safe #403; no wire change. The floor stays pinned until 1.0: WIRE_BUILD_ID forces one build per cluster, and wire_version.rs documents why the floor stays pinned.
  • No fencing tokens anywhere in raft/cluster. Two layers land. Transport: cluster_epoch is stamped in the RPC header and enforced at decode, rpc_codec/header.rs. Catalog: descriptor-version fencing on every dispatch path, 6023d2a, 5445f8b, a18c74b — see Descriptor-version fencing.

Medium / Low

  • Scatter-gather has no shard-failure/timeout guard → hang or silently-partial aggregate. Four coordinators rebuilt on one shape: typed error, proof-carrying totals with private fields, BTreeSet responded-tracking, 30s deadline. b037bab, ae7718f, 516d263.
  • Descriptor-lease expiry uses raw cross-node wall clock, no skew bound / fencing token. Partial. The drain no longer judges another node's deadline by its own clock; an explicit DescriptorDrainEnd is the only thing that clears a drain. 1f67bb4. The skew bound the clamp derives from now exists: MAX_CLOCK_SKEW_NS (5s) with HlcClock::update_checked refusing a remote HLC further ahead, nodedb-types/src/hlc.rs:25, wired at metadata-apply ingress. bb7a949. Remaining: apply that clamp to lease expiry itself, and a SWIM-Dead to on_node_crash release hook. Done in Make cluster consensus, crash recovery, and restore crash-safe #403. A remote lease stays live until expires_at + MAX_CLOCK_SKEW_NS (lease_liveness/dead_holders.rs). A SWIM-Dead holder is recorded and its leases are released by raft_loop/lease_gc.rs. A dead DDL owner's prepare lease is reclaimed.
  • SWIM fast-restart rejoin can stick (restart at incarnation 0 refuted vs lingering Dead). The incarnation is persisted and bumped on restart, nodedb-cluster/src/bootstrap/start.rs:84 — a node that crashed while peers held Dead(A, N) restarts at N + 1, so its first Alive dominates every lingering rumour. A read error fails bootstrap rather than silently restarting at 0.
  • No pre-vote (partition-heal term inflation forces healthy leader to step down); no leadership transfer (TimeoutNow). Pre-vote grants require Follower or Candidate role, a higher term, the up-to-date log rule, and time-bounded leader stickiness. dd91eed. TimeoutNow, ae45049.
  • Snapshot GC sweep_orphans runs at startup only; decommission doesn't propagate cluster-wide ShutdownWatch. GC also runs every ~60s from the Raft tick loop, raft_loop/loop_core.rs:418. Decommission signals through consensus, 0f625a6, 7056676.

State-machine divergence

  • A committed entry could apply twice. One Raft batch carried the same committed index twice, and apply_committed forwarded both copies. A PK-carrying write absorbs the repeat as an upsert and hides it; an append-shaped write (timeseries/columnar/spatial insert) gains a row. Each committed index is now handed off at most once, and the delivery watermark is claimed in one critical section — reading it and advancing it under separate lock acquisitions let two concurrent deliveries claim the same range. ea4f73f. Reproduced 3-in-6 before, 10/10 clean after.
  • Why the Raft layer emits a batch containing a repeated committed index is unexplained. The applier enforces its own contract, which is the right boundary and makes the duplicate unreachable, but the upstream duplication may be a defect in log delivery. Start at nodedb-cluster raft loop batch construction. Done in Make cluster consensus, crash recovery, and restore crash-safe #403. Cause: queue_committed_entries re-queued from last_applied+1 on each commit advance in one Ready window. It now resumes from the last queued index (nodedb-raft/src/node/internal.rs). A repeat across batches is dropped at apply_committed.rs by the <= last_applied filter.
  • Raft readiness was bounded by a fixed 30s wall clock, so a node with a large metadata replay was killed mid-boot while healthy. The gate now bounds on lack of progress: the deadline resets whenever applied_index advances, and a stuck group still fails, naming the index it stuck at. raft readiness timeout is hardcoded to 30s — large replay wedges the node #253, 0e1f6a3.
  • HlcClock::update had no caller at any node boundary and no skew bound, so cross-node HLC comparisons carried an unbounded offset and one far-future observation moved a node permanently. update_checked refuses a remote HLC more than MAX_CLOCK_SKEW_NS ahead; the one remote observation the metadata group carries is folded at apply, and every other Hlc on a MetadataEntry is a future deadline that must not be. HlcClock::update is unwired — no cross-node HLC merge and no clock-skew bound anywhere #264, bb7a949.

Descriptor-version fencing

A descriptor lease does not prove a plan is current. force_refresh_lease never compares the requested version against the catalog. A plan stamped just before a DDL commits still acquires its lease afterwards, at the version that DDL superseded.

Two dispatch paths carried a plan across that window.

Path Window Commit
Interactive transaction Statement time to COMMIT. Bounded only by how long the client holds the block open. a18c74b
DEFINE EVENT THEN action Plan to lease acquisition, then across each task's WAL append and SPSC round trip. 5445f8b

gateway::version_check::check_descriptor_holds owns the comparison. Holds are deduplicated by scope identity, then grouped by their own (database, tenant). A collection name is never compared against another tenant's catalog. Non-collection descriptor kinds are skipped. A stale hold aborts COMMIT with SQLSTATE 40001 before any durable write.

The event-action fence re-runs per dispatch. One check before the loop clears task 0, then lets a later task run against a version the catalog has left behind.

Correction. Two shipped module comments claimed a timed-out drain "is force-ended and the DDL proceeds anyway". It does not. The timeout proposes DescriptorDrainEnd to unblock new acquires, then drain_for_ddl returns Err and the DDL aborts. The fences are still required. Their stated reason was wrong. Both comments now name the real window above, and the real carve-out: a mixed-version cluster skips the drain outright.

Deferred actions

Auditing the trigger fence surfaced three defects in the Event Plane retry path. All three were confirmed in code before any edit.

  • The retry re-fired triggers that had already succeeded. dispatcher/single.rs called fire_for_operation with no trigger filter. The failed trigger's name was stored and never read.
  • A failing trigger cancelled its siblings. fire_common.rs returned Err on the first failure, so triggers queued behind it never ran on any path.
  • DEFINE EVENT failures had no retry path. consumer_helpers.rs passed no queue.
  • Retry unit is one action, not one source write. Records are keyed for idempotent enqueue and backed by a per-core redb store that opens lazily. 02f8a20.
  • Actions commit as one transaction, control/system_txn/. Partial dispatch became unrepresentable, so that variant and its per-dispatch fence are deleted — the COMMIT fence covers them.
  • DLQ is readable and replayable. SHOW TRIGGER DLQ and REQUEUE TRIGGER DLQ <entry_id>. It was durable but write-only, with no payload to replay. ddd8d09.
  • Event-action writes no longer cascade without bound. They dispatched as EventSource::User, and the event-definition path gated on is_data_event() but never on the source. An action writing to the collection it watches re-fired itself forever.
  • Trigger bodies are procedural blocks, not transactions. A body that fails mid-block has partially applied, and its retry repeats what landed. Apply the system_txn treatment event actions already get. Done in Make cluster consensus, crash recovery, and restore crash-safe #403. Trigger bodies run as one atomic block inside a system txn (control/trigger/fire_common.rs). A failed body discards its staged writes. BEFORE, INSTEAD OF, and sync AFTER bodies join the statement txn.

Test infrastructure

  • Both crash_resp_kv_write tests join the load-sensitive nextest group. Each fails the harness 20s readiness poll under the full suite and passes in ~5s in isolation. The child server is starved of cores by the unit-test fan-out, not slow. 805ebde.
  • calvin_multishard_write_in_explicit_block_commits failed 3 of 3 tries on "no sequencer leader elected yet". Its gate waited for sequencer_metrics to be set, which proves only that start_raft wired the service. It now waits on raft_status_fn reporting a SEQUENCER_GROUP_ID leader. de48fba.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:cluster-raftRaft, replication, consensus safetytype:epicTracking issue spanning multiple sub-issues

Type

No type

Projects

Milestone

No milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions