← Back to Writing

Leader Election Was the Easy Part: What Building Raft Actually Required

A cluster electing a leader makes a satisfying demo. It does not yet make a useful key-value store.

In QuorumKV, election was the point where the longer list of correctness obligations became visible: what survives a crash, what a read means, what happens when a client retries, how a lagging node catches up, and how the voter set changes without briefly creating two clusters.

Persistence changes the protocol

Raft's term, vote, and log are not ordinary in-memory fields. A node must persist the relevant state before sending a response that depends on it. Reordering those operations can produce a node that promises one history over the network and remembers another after restart.

QuorumKV uses checksummed, length-delimited WAL records and explicit replay behavior. The tests kill real subprocesses around meaningful persistence points, restart them from disk, and check the state machine that emerges. Message-drop tests alone did not exercise the durable transition I wanted to test.

A leader can still serve a stale read

An isolated former leader may not immediately know it lost authority. Returning its local value would violate linearizability even though no write path ran.

GET therefore uses a quorum-confirmed ReadIndex path. Before serving the read, the leader confirms authority with a current-term quorum and waits until its state machine has applied through the required commit index. QuorumKV deliberately does not claim lease reads; time-based authority would introduce a clock assumption the design does not need.

Client retry identity belongs in replicated state

A PUT can commit while its response is lost. Retrying the same logical request must not apply it twice. QuorumKV replicates client request identity with the command, so deduplication becomes state-machine state and survives leadership changes and snapshots.

This is stronger than remembering recent request IDs in the leader process. The leader is exactly the process most likely to change during the retry.

The log cannot grow forever

Snapshotting compacts applied history, but it creates a second catch-up protocol. A far-behind follower may need chunked InstallSnapshot transfer rather than normal AppendEntries. The installed snapshot must restore both application data and replicated request identity before log replication resumes.

Membership is not a configuration-file edit

Changing a peer list locally can create two groups that each believe they have a majority under a different configuration. QuorumKV represents membership changes in replicated state and uses joint consensus so commitment temporarily requires majorities from both voter sets.

Real crashes told me more than dropped messages

The failover test starts three QuorumKV binaries as separate OS processes and writes x=1 through the cluster. It finds the current leader and sends that process SIGKILL—no graceful shutdown and no chance to clean up memory first. After the surviving majority elects a replacement, the test reads x and writes y=2 through the new leader.

Then it restarts the old leader from the same data directory. I poll that node directly until its own applied index catches up; querying through the normal client would be weaker evidence because a follower could redirect the read to the new leader. Finally, the cluster returns both values.

put x=1 → kill leader with SIGKILL
surviving majority → elect replacement → get x → put y=2
restart old leader from same disk → wait for local catch-up → get x and y

This test exercises real process loss, persisted state, a leadership change, new work after failover, and follower catch-up. It does not prove every Raft execution correct, which is why the repository keeps the formal-proof claim out.

Backpressure belongs in the protocol boundary

An unbounded replication queue can turn a slow follower into memory growth on the leader. The transport and replication workers bound frames, batches, and pending work so overload becomes explicit rather than hidden resource exhaustion.

The limits of the evidence

QuorumKV's tests start real node processes, move traffic over TCP, kill leaders, restart nodes, and exercise persistence failures. That is stronger evidence than an in-memory election test, but it is not a proof of Raft correctness. The implementation has a bounded feature set, and the repository does not claim formal verification. I treat each failure test as evidence for a particular invariant, not as permission to generalize beyond what the suite covers.

What building Raft came to mean

Building QuorumKV changed my definition of “Raft implementation.” Election and replication are the recognizable center. Useful consistency lives in the surrounding details: persistence ordering, read authority, replicated client identity, compaction, configuration changes, and crash evidence.

Related project: QuorumKV case study →