etcd, Consul, CockroachDB, TiKV and Kafka's KRaft mode all rest on the same small idea: a group of machines agree on one ordered log of commands, and each machine applies that log to its own copy of the state. If a majority is alive and can talk, the group keeps working. If not, it stops instead of diverging.
Raft is the algorithm most of them use, and it was designed to be understood. Figure 2 of the paper fits the core rules on one page. Getting those rules right in code, under crashes, lost messages and partitions, is a much longer story, and that's this project. You'll finish with a KV store that three or five processes serve together, plus a test harness that tries hard to break it.
01Why build this
Consensus sits under the control plane of most modern infrastructure. Once you've implemented it:
- Split-brain stops being a vague fear. You'll know exactly which rule keeps two leaders from both accepting writes in the same term.
- Quorum sizes make sense. You'll see why five nodes tolerate two failures, and why a fourth node adds cost without adding tolerance.
- "Linearizable" gets a precise meaning. You'll learn the hard way that a leader serving reads from memory can return stale data.
- etcd and ZooKeeper outages read differently. Leader churn, slow disks causing missed heartbeats, and snapshots too large to send are all things you'll have debugged.
It's also the best project here for learning to test concurrent code. The bugs are rare, timing-dependent and silent, and finding them forces you to build real tooling.
02What you're building
A single write, from client to committed, in the finished system:
?Why wait for a majority, and not for everyone?
Because any two majorities of the same cluster overlap in at least one node. A committed entry lives on a majority, and a new leader needs votes from a majority, so at least one voter has the entry. Raft's election rule then refuses leadership to any candidate whose log is behind that voter's. Waiting for every node would add nothing to safety and would let one slow node stall the cluster.
03Before you start
| You need | Why | Where to get it |
|---|---|---|
| The extended Raft paper, especially Figure 2 | The rules, stated exactly | raft.github.io |
| A language with good concurrency and RPC | Timers, goroutines or tasks, and message passing | Go is what MIT 6.5840 uses; Rust with tokio works well |
| A simulated network you control | You need to drop, delay, reorder and partition messages on demand | Build it yourself in milestone 1 |
| Some background on replication and failure detection | Why timeouts can't tell slow from dead | Chapter 27, Chapter 30 |
| Patience with rare bugs | Some failures appear once in hundreds of runs | A script that runs your tests in a loop overnight |
04The roadmap
Eight milestones, roughly in the order of the MIT 6.5840 labs. Keep every earlier test passing as you add later ones.
A network you can break
1 weekendBefore any Raft code, write the harness: nodes in one process, sending messages through a router that can drop, delay, duplicate and reorder them, and cut the cluster into groups. Seed its randomness so a failing run can be replayed.
It feels like a detour. Probably every serious Raft bug you'll find this month gets found here, so it's time well spent.
Leader election
1–2 weekendsEvery node starts as a follower with an election timer. If it hears nothing from a leader before the timer fires, it becomes a candidate, increments its term, votes for itself and asks the others. A node grants at most one vote per term. Whoever wins sends empty AppendEntries as heartbeats.
Randomise the timeout so candidates don't keep splitting the vote. Also, any message carrying a higher term makes the receiver step down to follower immediately. Most early bugs are a missing check for that.
Log replication
2 weekendsThe leader keeps a nextIndex and matchIndex per follower. AppendEntries
carries the index and term of the entry just before the new ones, and a follower
rejects it if its log doesn't match there. On rejection the leader backs up and
tries again. A follower deletes conflicting entries and takes the leader's.
Advance the commit index when a majority has an entry, and only for entries from the leader's current term. That last condition is the subtle case in Figure 8 of the paper. Read it twice before you skip it.
Persistence and restarts
1 weekendThree fields must survive a crash: the current term, who this node voted for
in it, and the log. Write them to disk and fsync before replying to any RPC
that changed them. A node that forgets its vote can vote twice in one term, and
two leaders follow.
Backing up one entry at a time, as the paper first describes, is slow after a long outage. Add the faster version it sketches, where the follower reports the conflicting term and its first index, so the leader can skip a whole term per round trip.
A key-value store on top
1–2 weekendsEach node runs a KV map and applies committed log entries to it in order. A client sends requests to the leader; followers reply with a hint of who the leader is. Only after the entry commits and is applied does the leader reply.
A client that times out will retry, and the first attempt may already be in the
log. Tag each request with a client ID and sequence number, and have the state
machine skip duplicates. Without that, Append runs twice.
Linearizable reads
1 weekendA leader that has been partitioned off doesn't know it yet, and will happily answer reads from its stale map. The simplest fix is to put reads through the log like writes. That's correct, and slow.
Then try read index: record the commit index, confirm leadership with a round of heartbeats, wait until that index is applied, then answer. Leases skip the heartbeat round but depend on bounded clock drift. Know which assumption you've made.
Snapshots and log compaction
1–2 weekendsLeft alone, the log grows forever. When it passes a threshold, have the KV layer serialise its state, including the duplicate table, and let Raft discard every entry up to that index. Now log indexes no longer start at zero, and every array access needs an offset.
A follower too far behind asks for entries the leader no longer has. Send it the snapshot with InstallSnapshot instead. Getting indexes right after truncation is where most off-by-one bugs in this project live.
Real processes, real failures
1–2 weekendsSwap the simulated network for real RPC and run each node as its own process.
Use iptables rules or a proxy between nodes to create partitions, and
kill -9 to crash them.
Record every client operation with its start time, end time and result. Feed the history to a linearizability checker like Porcupine. If it rejects a history, the bug is real, and you have the exact interleaving that shows it.
05Traps that catch everyone
| Symptom | Cause | Fix |
|---|---|---|
| Two leaders in the same term | A vote wasn't persisted before replying, or votedFor wasn't reset on a new term | Persist before replying; reset votedFor only when the term increases |
| Elections never settle | Timeouts not randomised, or the timer resets on every message instead of only valid ones | Randomise per election; reset only on a vote granted or AppendEntries from the current leader |
| A committed entry disappears after a leader change | Committing an old-term entry by counting replicas | Only count replicas for entries from the current term (Figure 8) |
| Deadlock under load | Holding the node's lock while sending an RPC | Release the lock before any network call, then recheck the term after |
| A retried Append applies twice | No duplicate detection in the state machine | Client ID plus sequence number, stored in the snapshot too |
| Stale reads after a partition heals | Reads served from the leader's memory without checking leadership | Route reads through the log, or use read index |
| Index out of range after snapshots | Log positions mixed up with absolute indexes | Wrap the log in a type that converts between the two in one place |
06Stretch goals
- Membership changes. Add and remove nodes safely, one at a time, the way Ongaro's thesis describes.
- Sharding. Split the keyspace across several Raft groups with a configuration service that moves shards between them, as in the last 6.5840 lab. Chapter 29 covers the partitioning side.
- Pre-vote and leadership transfer. Stop a node that rejoins after a partition from forcing a pointless election.
- Batching and pipelining. Send many entries per AppendEntries and several in flight at once, and measure throughput before and after.
- Deterministic simulation. Run the whole cluster in one thread with a fake clock, so every bug replays exactly from a seed.
07References worth your time
The extended Raft paper. Figure 2 is the specification you're implementing; keep it printed next to you.
The course labs build Raft and a sharded KV store in Go with a test harness that injects failures. Probably the best structure for this project.
Written by a 6.824 teaching assistant from the bugs students actually hit. Read it before milestone 2 and again when you're stuck.
The PhD thesis behind Raft, covering membership changes, log compaction and client interaction in more depth than the paper.
Anish Athalye's fast linearizability checker in Go, used by the 6.5840 test suite. Plug your recorded histories into it in milestone 8.
Kyle Kingsbury's analyses of real databases under partitions, and Maelstrom, a workbench for testing your own distributed system against the same kind of faults.
A production Raft in Go, used by etcd and others. Read it after milestone 7 to see how a real implementation separates the protocol from I/O.