Kind-3 term markers in the WAL stream, STREAM-relative vote frontiers, heartbeat-only divergence detection + quarantine, and fenced promote. Elections converge in 0.6–1.0s; zero acked-write loss across all kill points. Closes G5 (leaderless recovery) from the v0.9 wave.
28 KiB
m11p4 — Failure Detection, Election, Fencing (COMPLETE — 2026-06-12)
Phase spec and exit gate: docs/roadmap-to-cluster.md §4/m11p4. Closes ROADMAP gap G5 ("operator-driven failover; restarted node believes the static topology — split-brain seed", incident §1.4-1). Predecessors: phase-1.md (ship queue), phase-2.md (one replicated log), phase-3.md (quorum machinery — the prerequisite this phase's election rides on).
Already shipped in this phase (d0a52e4, 2026-06-11): catch-up timer retry
(re-arm-on-skip liveness) + TSEG segment version header + structured
"snapshot required" refusal. This document covers the remaining XL work.
Goal: "a machine died" is a non-event. Kill the leader under load → writes resume <10 s p99 with zero operator action and zero acknowledged loss; a restarted or partitioned ex-leader cannot accept a single write; election churn under flapping links is bounded.
Design (as adopted)
0. Open question #2 resolved: purpose-built election, no raft crate
The roadmap deferred "minimal purpose-built Raft over tidal-net vs. a
vetted raft crate" to this phase. Decision: purpose-built election-only
Raft (Raft §5.1–§5.4 election subset + pre-vote + check-quorum + leadership
transfer), for one structural reason: the replicated log already exists.
m11p2 made the WAL the stream; m11p3 built the quorum commit index over it.
A raft crate (openraft, raft-rs) owns its own log and storage contracts —
adopting one means either running a second log beside the WAL or contorting
the crate to delegate its log to ours. The state that actually needs
consensus is tiny — (term, leader) — and the safety-critical machinery
(durable votes, the up-to-date restriction, term fencing) maps directly onto
primitives we already have (flushed_seq(), the stream_baseline
persistence pattern, the commit-index epoch). Membership stays static this
phase (the topology file); m11p5 makes it dynamic through this phase's
metadata machinery.
The election state machine is pure and deterministic (no clocks, no I/O
— inputs carry Instants, randomness is an injected seed), so the entire
protocol is unit-testable without processes or sleeps.
1. The durable election state (kills restart amnesia)
A small per-node file, data_dir/election_state, written with the exact
stream_baseline discipline (write .tmp, sync_file_durable
[F_FULLFSYNC on macOS], atomic rename, sync_dir_durable), carrying a
magic + format version + checksum so corruption is detectable:
| Field | Type | Meaning |
|---|---|---|
current_term |
u64 | Highest term this node has seen. Strictly monotonic — a term is adopted only upward; no inbound message can ever lower it. |
voted_for |
Option | Vote cast in current_term (at most one). |
The WAL-tail term deliberately does NOT live in this file — see §1b.
Persistence is a fence, not bookkeeping: the file is fsynced BEFORE any externally visible action that depends on it — before a vote response is sent, before a higher term is acknowledged, before leadership is assumed. Votes and term changes are rare (not per-write), so the fsync cost is irrelevant. All election-state mutations serialize under one lock, and the boot-time persist completes before the HTTP routes serve a single request (design-review C6).
Boot rule (the incident §1.4-1 fix), four cases (design-review C1):
- File exists and reads clean → the node always boots as a follower, regardless of what the file or the topology says — leadership is learned from heartbeats or won by election, never assumed from a restart.
- File exists but is corrupt (bad magic/version/checksum) → refuse to boot with a FATAL error naming the file. A node that cannot prove its term knowledge must not guess at term 0 — guessing re-opens the §1.4-1 amnesia hole through the side door of a disk fault.
- File absent but the data dir is NOT fresh (a WAL exists) → the file was deleted out from under a node that has run before. Boot as a follower at term 0 with a loud ERROR (never assume topology leadership), and persist the recreated state immediately. Terms recover upward from the first heartbeat/ship the node hears.
- File absent and the data dir is fresh → a genuinely new node: boot with
the topology's term-0 roles exactly as today, persisting
{term: 0}before serving. Existing deployments boot byte-for-byte identically on first start.
Term-0 conflicting-leader guard (design-review C5): a term-0 leader that
receives another node's term-0 leadership assertion (heartbeat or ship
stamped term=0, leader=other) — possible only through mismatched topology
files, a deployment error — logs ERROR and starts an election. Term ≥ 1
leadership is unique by votes; the election converges the misdeployment
instead of letting two config-believing leaders run open-loop.
1b. WAL term markers (the up-to-date restriction, crash-atomic)
Raft's vote restriction compares (lastLogTerm, lastLogIndex). The first
design draft persisted a last_applied_term in the election file; the
adversarial review (C7) found the killer: that file and the WAL frontier
can never be made crash-atomic — a crash between the two persists yields a
node whose up-to-date claim lies in the unsafe direction (it can win an
election while missing committed entries).
Fix: the term lives in the log itself. On becoming leader of term T
(any T ≥ 1), the leader journals a TermMarker(T) record — a new small
WAL record kind beside m11p2's item/embedding blobs — as its first term-T
entry. It replicates to every follower through the normal stream like any
other record; followers fold it by recording the term (no storage effect).
The node's WAL-tail term is then the term of the last marker at or
below its durable frontier — recovered from the WAL on boot, advanced on
apply, and atomic with the frontier by construction (they are the same
fsync stream). No marker (a pre-election log) reads as term 0.
The vote comparison becomes (wal_tail_term, durable_frontier),
lexicographic, both read from the same WAL state. A crash at any point
yields a node that understates — never overstates — its log, which is the
safe direction (it may lose an election it deserved; it can never win one
it shouldn't, because every committed entry's quorum contains an
un-crashed holder whose refusal stands).
Wire compatibility falls out of "markers only exist at term ≥ 1": a p4 binary on a term-0 cluster never writes one, so a rolling m11p3→m11p4 upgrade ships zero unknown record kinds; the first election (only possible once a vote quorum of upgraded binaries exists) is the moment the format becomes live. Runbook: complete the binary upgrade before relying on auto-failover.
2. Term model: term 0 is the topology, terms ≥ 1 are elected
- Term 0 = the static topology file's leader. A fresh cluster boots leading/following at term 0 with no election, preserving every existing boot path and test.
- Any election or fenced transfer moves the cluster to term ≥ 1. From
then on the durable state is authoritative and the topology's
leaderfield is dead weight (kept for bootstrap only). - Proto compatibility falls out: proto3 defaults absent term fields to 0,
so a pre-m11p4 peer's ships/reports read as term-0 traffic. A node still
at term 0 interoperates with them; a node at term ≥ 1 rejects them
(
FAILED_PRECONDITION) — exactly the fencing semantics, applied to version skew.
3. Election protocol
Roles: Follower → PreCandidate → Candidate → Leader (and anything →
Follower on observing a higher term).
- Failure detector: the leader sends
Heartbeatto every peer eachheartbeat_interval(default 500 ms), now carrying(term, leader_region, stream_baseline). A follower that hears nothing valid for a randomizedelection_timeout∈ [1500 ms, 3000 ms) starts a pre-vote. - Pre-vote (disruption guard): asks "would you vote for me in term+1?"
without bumping any term. A voter refuses if it heard from a live leader
within
election_timeout_min(lease check) or if the candidate's log loses the up-to-date comparison. Only a pre-vote majority proceeds — a partitioned or flapping node cannot inflate terms or depose a healthy leader from the outside. - Vote: candidate persists
term+1, voted_for=self, then requests votes. A voter grants iffreq.term ≥ current_term, it hasn't voted for another candidate in that term (durable before reply), and the candidate is up-to-date:(wal_tail_term, durable_frontier) ≥the voter's, compared lexicographically — Raft's §5.4.1 restriction with both values derived from the voter's own WAL (§1b). Majority → leader. - Becoming leader: persist election state, then journal the
TermMarker(T)(§1b), then the promote mechanics verbatim: baseline = ownflushed_seq()(now ≥ the marker's seqno),persist_stream_baseline,ship_queue.activate_from(baseline, term)(which activates the commit index scoped to T and bumps its epoch), leader view updated, heartbeats start. - Step-down ordering (design-review C3): on any step-down — higher-term
contact, check-quorum loss, fenced transfer — the leadership view flips
first (
leader = Noneor the new leader), THEN the ship queue deactivates. The write path'sis_leader()reads the view, so once it flips no new write is accepted; the queue deactivation then fails in-flightack=quorumwaits with the existingDemotedpath. The residual window (a request that read the old view and stages after the flip) produces a leader-local-durable, never-shipped write — exactly the documentedack=leadercrash contract, surfaced at rejoin as a divergent suffix (§5). - Check-quorum / lease step-down: a leader that cannot reach a majority
of peers (heartbeat acks) within
leader_lease_mssteps down (stops accepting writes; quorum waits fail). The lease bound (design-review C2): the leader's countdown starts at its last successful quorum contact while a follower's election timer starts at the last heartbeat it received — the two anchors differ by up to one heartbeat interval plus RTT, NOT by wall-clock offset (all timers are monotonicInstantintervals; clock offsets cancel). The enforced constraint is thereforeleader_lease_ms + heartbeat_interval_ms < election_timeout_min_ms(validated at topology load), guaranteeing the deposed leader stops before any successor can be elected. This bounds the minority-partitioned-leader window forack=leaderwrites;ack=quorumwrites were never at risk (the commit index cannot reach majority from a minority partition). - No pre-vote deadlock (design-review C9): the lease refusal can only
outlive a real leader if something keeps sending heartbeats. A
stepped-down leader stops heartbeating by construction (only the Leader
role sends), so after a leader dies or steps down every follower's lease
window drains within one
election_timeoutand pre-votes start granting. Asymmetric partitions degrade exactly one node (it pre-votes, is refused by peers that still hear the leader, and never inflates a term — the pre-vote's whole point); the majority side keeps serving throughout. - Timers use
std::time::Instantexclusively — monotonic, immune to the chaos suites'TIDAL_HLC_SKEW_MSwall-clock skew injection (HLC remains data-plane-only for hard-negative LWW).
4. Fencing everywhere
Every replication RPC carries the sender's term; every receiver gates:
| RPC | Stamp | Stale-term reaction | Higher-term reaction |
|---|---|---|---|
ShipSegment |
term, leader_region |
reject FAILED_PRECONDITION before enqueue |
persist & adopt term, treat as leader contact |
StreamSegments request |
puller's term |
source at a higher term refuses the pull | source at a lower term steps down (it is deposed and must not serve its possibly-divergent tail — design-review C14) |
StreamSegments chunks |
same as ship (chunks are ShipSegmentRequests) |
puller drops the chunk + aborts the pull (a stale source can never inject history) | adopt |
Heartbeat |
term, leader_region, stream_baseline |
respond accepted=false, term (sender steps down) |
adopt, reset election timer, run the §5 term-join check |
ReportApplied |
reporter's term |
leader folds only if reporter term == the commit index's activation term (design-review C4/C8/C12: the gate reads the term the index was activated with, so a stale, future, or cross-leadership report can never advance it; heal's resume_from folds through the same gate — C13) |
leader steps down |
RequestVote / TimeoutNow |
term |
refuse, return own term | adopt |
| HTTP forward (internal-marked writes) | already rejected on a non-leader (NotLeader 503); the error body now names (term, leader-or-none, leader_http_addr) so the forwarder follows one redirect to the named leader or returns a retryable 503 "election in progress" (design-review C11) |
— | — |
A ship/report/vote response carrying a higher term is a step-down
signal to the sender — the deposed leader learns its term is stale from
its own outbound traffic even if it never hears a heartbeat. The
CommitIndex itself is term-scoped: activate(baseline, term) records the
term, update_peer folds only same-term reports, so every quorum input —
ship-ack hints, ReportApplied, heal resumes — passes one gate.
5. Rejoin, and the divergent-suffix boundary
The check runs exactly once per term, at term-join (the first valid leader contact carrying term T — heartbeat or ship — when adopting T), and compares the node's durable WAL frontier against the term's activation baseline B (persisted by the leader at activation and announced in every heartbeat — design-review C10: B is immutable for the term, so the check never races a moving target; a frontier that grows past B within term T is just normal replication, never re-checked):
- frontier ≤ B: clean rejoin. Catch-up via the existing
StreamSegmentspull (gap-triggered, timer-retried). The node's WAL-tail term advances to T when the term's marker record applies (§1b). - frontier > B with WAL-tail term < T: the node holds a divergent
suffix — writes that were leader-durable on a dead leader but never
quorum-committed, while the cluster elected past them. The node
quarantines: it refuses to apply or report (receiver halt latch — the
existing
died()health surface), logs ERROR naming the suffix range and the reseed runbook, setstidaldb_cluster_divergenceand stays fenced from the data plane. It still votes (its stale WAL-tail term makes it lose every up-to-date comparison against healthy peers, so it can never win). - frontier > B with WAL-tail term == T: impossible at term-join (T-term entries only exist after the leader's activation, and joining is what admits them) — and after joining, it is the normal shape of a follower tracking the leader. Restarts within a term re-join idempotently: the tail term already equals T, so no quarantine.
Why quarantine instead of Raft-style suffix truncation: tidalDB followers
apply segments on receipt (before quorum commit) — by the time
divergence is discovered, the suffix's effects are already folded into
storage (decayed counters), and unapplying signal arithmetic is not a thing.
Truncating the WAL alone would lie about state. The honest recovery is a
reseed — which is m11p5's snapshot+stream work, and exactly what the
already-shipped p4 increment's FAILED_PRECONDITION ("snapshot required")
path points at. Probability note: with RF=3 a quorum-committed suffix holder
is a member of every election majority (the vote restriction protects it),
so wrongful quarantine is impossible at n=3; the quarantined case is
strictly leader-ack-only data, which is the documented leader-ack contract.
6. Promote becomes a fenced transfer
POST /cluster/promote {region}:
- Leader alive (the maintenance drain): forwarded to the current
leader, which waits until the target's durable mark reaches the leader's
flushed frontier (bounded; 503 naming the lag on timeout), then sends
TimeoutNow— the target starts an immediate election (skipping pre-vote and timeout; the leader sanctioned it), wins on term+1, and the old leader steps down on the first higher-term contact. No drain-the-world: the catch-up wait IS the drain. - Leader dead/unreachable (the legacy failover drill): the target runs a normal pre-vote + election among survivors. With auto-election on, the operator usually never gets here — the verb remains as a manual override and a way to choose the successor.
- Mixed-version fallback: if a majority of peers answer
Unimplementedto vote RPCs (pre-m11p4 binaries), promote falls back to the legacy fan-out promote at term 0 with a loud WARN. Auto-election quietly waits until the fleet is upgraded (vote quorum impossible until then).
7. What goes where
| Crate | Work |
|---|---|
tidal |
replication/election.rs: pure ElectionState step machine (inputs → actions) + config; replication/election_durable.rs: the fsync'd (current_term, voted_for) file with magic/version/checksum + the four-case boot classification; the TermMarker WAL record kind + recovery/apply tracking of the WAL-tail term (§1b); CommitIndex.activate(baseline, term) + same-term fold gate; WalSegmentPayload.term/leader_region; ship queue stamps terms (shared Arc<AtomicU64>); receiver-side divergence latch. |
tidal-net |
Proto: term fields on ShipSegmentRequest/Response, StreamRequest, HeartbeatRequest/Response, AppliedReport; new RequestVote (with prevote/transfer flags) + TimeoutNow RPCs. Server: late-bound ElectionHooks (vote/heartbeat/observed-term) mirroring the AppliedSink pattern; term gates on ship/stream/report handlers. Client: PeerPool::{request_vote, timeout_now, heartbeat} — vote RPCs bypass the circuit breaker (an election must be able to probe a peer the breaker quarantined; vote traffic is rare and self-limiting). Config: election knob block. |
tidal-server |
The election driver (one tokio task per region state): ticks the state machine, executes actions (persist → send → transition), owns the heartbeat fan-out and failure detection. Boot path: durable-state detection → rejoin-as-follower. Leadership transitions route through the existing promote_local mechanics. leader becomes RwLock<Option<RegionId>> (leaderless windows are real and reported honestly — forwards during an election return retryable 503 "election in progress"). /cluster/status grows term, role; election metrics (tidaldb_cluster_term, _role, _elections_started_total, _leader_changes_total, _stepdowns_total, _votes_{granted,denied}_total, _divergence). Promote rework per §6. |
| tests | Pure state-machine property/unit tests (no clocks); tidal-net vote-RPC tests over real sockets; tier-3 cluster_election.rs: the three exit gates below. |
8. Configuration (topology election block, all serde-defaulted)
| Knob | Default | Constraint |
|---|---|---|
heartbeat_interval_ms |
300 | |
election_timeout_min_ms |
1500 | ≥ 3× heartbeat interval |
election_timeout_max_ms |
3000 | > min |
leader_lease_ms |
900 | leader_lease_ms + heartbeat_interval_ms < election_timeout_min_ms (the C2 bound — validated at topology load) |
auto_election |
true | false = detector + fencing only; elections fire only through /cluster/promote. After a full-cluster restart in this mode the cluster is deliberately leaderless until an operator promotes (design-review C16 — that is the posture the knob selects). |
Failover arithmetic: detect (≤3 s worst-case timeout) + pre-vote+vote (2 RTTs + 2 fsyncs) + activation (µs) ≈ 3–4 s typical, comfortably inside the 10 s p99 gate, without making the detector so twitchy that a fsync stall triggers an election (the 10–50 ms F_FULLFSYNC tail from p1 informs the 1500 ms floor). Heartbeats are async fan-out tasks on the transport's runtime with a per-call timeout well under the interval — N peers at 300 ms is negligible load and can never stall the data plane (design-review C15).
Exit gate (from the roadmap, restated as tests)
- Auto-failover, zero acked loss: SIGKILL the leader under
ack=quorumload with auto-election on → writes resume on the new leader <10 s p99, zero operator verbs, and the m11p3 ledger checker (frontier + content proofs) passes across repeated random kill-points. - Fencing under partition + restart: partition the old leader away mid-election, let it restart → it boots a follower (durable state), its ships/votes carry a stale term and are rejected, and it cannot accept a single write (every write attempt 503s) until it heals and rejoins.
- Bounded churn: flapping links (proxy sever/heal cycles) produce bounded elections (pre-vote absorbs the flaps; no term explosion, no livelock), and the cluster converges to exactly one leader per term.
Status
- Design adopted (this document; 16 confirmed adversarial-review holes folded in pre-implementation)
- Election core (state machine + durable state) —
tidal - Transport (proto, hooks, gates, client verbs) —
tidal-net - Driver + boot + fencing + promote rework —
tidal-server - Exit-gate suites green
- Docs (runbook §9, CHANGELOG, roadmap, status fields)
As built (deltas vs the adopted design)
- The vote-restriction frontier is STREAM-relative, not own-WAL-relative.
The gate run exposed it: a reseeded node's local WAL numbering diverges
from stream numbering after a baseline jump (observed live: stream-applied
38 vs own-WAL 22), so raw local frontiers are not comparable across nodes.
LogPosition.frontieris the node's durable position in the last joined term's stream numbering: ownflushed_seq()when it led that term,applied_seqno(shard of that term's leader)otherwise. The term marker therefore carriesleader_region, recovered with the term. - Divergence is judged against the leader's ELECTION-TIME position, only
from heartbeats. The adopted baseline-comparison was cross-numbering for
the same reason. The heartbeat carries
(prev_log_term, prev_log_seq)— the same pair the vote restriction compared — and a joining node is divergent iff its own position exceeds it lexicographically. A NEW-term ship arriving before the first heartbeat defers with a retryableJoinPending(UNAVAILABLE) instead of guessing;stream_baselineon the heartbeat is used only for the frontier jump. - "Topology era" ends at the first JOINED election, not at one's own
durable term. A failed campaign (an isolated operator override, the
chaos drills' deliberate split-brain) inflates a node's durable term with
no elected leadership existing; gating term-0 traffic, the legacy promote
and the check-quorum fallback on
joined_terminstead ofcurrent_termkeeps that world functional while elected clusters stay absolutely fenced. - Check-quorum step-down is gated on
auto_election— with elections off no successor can exist, so stepping down would only convert a follower outage into total write unavailability; the knob deliberately preserves the m11p3 availability posture. - The transfer's catch-up wait reads the commit index's peer mark (fed
by
ReportApplied), not the ship queue's acked frontier — the latter stalls ~30s behind an open circuit breaker after the target restarts. - The frontier-report gate at the handler references the commit index's activation term (the same value the sink's race-free in-lock gate uses), not the machine's possibly-candidacy-inflated term.
- A rejoining term-0 ex-leader's boot view is
None, neverSome(self-from-topology)— caught by the gate before it shipped: the adopted text only forbade assuming the role, but a self-pointing VIEW still opened the §1.4-1 write-acceptance hole with the queue parked.
Exit-gate evidence (local 3-process clusters, tier-3, 2026-06-12)
| Gate | Result |
|---|---|
| Auto-failover, zero acked loss | 5 rounds × random kill points under ack=quorum load, ZERO operator verbs: elections in 0.59–0.98 s (budget 10 s); frontier invariant (elected leader's election-time position ≥ every acked seq) + content invariant (every acked item searchable on the new leader) held every round, including rounds where the killed node rejoined via quarantine + reseed |
| Fencing under partition + restart | leader isolated (bidirectional TCP severs) → survivors elected in ~0.7 s; the ex-leader restarted while partitioned and was 0-for-~30 on write attempts across the fence window, then healed and rejoined the new term as a converged follower |
| Bounded churn | 4 sub-timeout flaps absorbed with no elections + 3 super-timeout flaps → converged to ONE leader on ONE term ≤ 15, serving quorum writes |
Full verification: 2,312 default-feature tests + all six tier-3 suites green;
workspace clippy clean (--all-targets --all-features, 0 warnings).
cluster_quorum/cluster_lifecycle pin auto_election: false (they validate
the manual drill, which the knob preserves); the quorum ledger drill gained
the promote-retry step (the election can refuse a stale max-applied sample —
the protocol now enforces the rule the runbook used to assign the operator).
Post-implementation review (seven-dimension, 2026-06-12)
Twenty-one findings, adversarially verified: 10 confirmed → all fixed, 11 refuted. The fixes that mattered:
joined_termboots at ZERO (three reviewers converged on it): it tracks terms whose JOIN CHECK passed, and the durable term is not evidence of that — initializing from it both fenced legitimate topology-era traffic after a failed candidacy and skipped the divergence check on restart. The fix surfaced a latent hole in the join check itself: the design'stail_term == joining term ⇒ cleanclause (§5, third bullet) was never implemented, so the corrected init would have false-quarantined every mid-term rejoiner. Both landed together; the failover gates re-proved them.- Torn-pair races:
WalTermMark's(term, seq, leader_region)and the node's activation(prev_log_term, prev_log_seq)were split atomics — a reader between stores could pick the new term with the previous marker's REGION and read a frontier in the wrong stream's numbering. Both are now single-lock values (written once per term, read a few times a second), and the activation pair publishes only AFTER the term marker is durable, so an aborted activation never exposes it. - Unit coverage added for
on_activation_failedandtransfer_to(the two uncovered machine transitions).
Notable refutations (recorded for the next reviewer): the "stale
wal_term_mark during join" claim — the mark and the frontier advance
through the same apply path, so a node's position is always internally
consistent; a behind position simply joins clean and catches up. And the
quarantine-metric-never-cleared claim — clearing rides m11p5's reseed
machinery by design.
Carried hazards (tracked, not regressions)
- A reseeded node converges on the live stream but does not hold pre-baseline
history; quorum marks it contributes cover ranges it genuinely holds, but
full-history re-replication is m11p5's snapshot transfer (same gap the
d0a52e4increment's "snapshot required" refusal points at). mp_rolling_upgrade_no_loss_no_stallis load-sensitive on saturated hosts (the post-restart heal races the 30 s circuit breaker inside its 60 s budget; widen via the harness env on slow runners). m11p8's self-driving heal removes the pattern.- Data below a term's activation baseline is unreachable through that term's stream for followers that were behind at the transfer — pre-existing m11p2 semantics, unchanged by this phase, resolved by p5 snapshots.