Raft is the consensus algorithm most teams actually implement. It was designed to be understandable, and it is: a leader per term, a replicated log, a majority rule for commits. Understanding the algorithm is not the hard part. The hard part is knowing that your implementation, with its own storage layer, network code, batching and timers, still has the properties the paper proves. Many production Raft bugs were not misreadings of the paper; they were correct-looking code that lost a write on crash, or a timing interaction nobody had tested.

This article is about that gap. It recaps the algorithm briefly, then turns the paper's five safety properties into a checker you can run, shows how to build a deterministic simulator that explores crashes and partitions reproducibly, and walks through the classes of bugs that real implementations have shipped. For the full rules and message handlers, see Raft consensus in depth; for why each rule exists, see Raft consensus intuition.

Advertisement

The algorithm in one page

A Raft cluster of N servers, usually 3 or 5, tolerates the failure of a minority. Time is divided into numbered terms. Each server is a follower, a candidate or a leader. A follower that hears nothing from a leader within a randomised election timeout increments its term, votes for itself and sends RequestVote to everyone. A server grants at most one vote per term, and only to a candidate whose log is at least as up to date as its own, comparing the last entry's term first and then its index. A candidate with votes from a majority becomes leader.

The leader appends client commands to its log and replicates them with AppendEntries, which carries the index and term of the entry just before the new ones. A follower rejects the call unless its log has a matching entry there, and the leader backs up until the logs agree, then overwrites the follower's conflicting suffix. An entry is committed once it is stored on a majority and it belongs to the leader's current term; earlier-term entries become committed indirectly when a current-term entry after them commits. Servers apply committed entries to their state machine in index order. Every server must persist its current term, its vote and its log to stable storage before answering any RPC. Replication mechanics are covered in Raft log replication.

The five safety properties, as code

The Raft paper states five properties that hold at all times: Election Safety (at most one leader per term), Leader Append-Only (a leader never overwrites or deletes its own entries), Log Matching (if two logs have an entry with the same index and term, the logs are identical up to that index), Leader Completeness (a committed entry is present in the logs of all leaders of later terms) and State Machine Safety (if a server has applied an entry at an index, no other server applies a different entry at that index).

These are the specification of your implementation. Each one can be checked mechanically from snapshots of every node's state. The checker below keeps history across steps, because several properties are about all time, not about one moment: two leaders in term 7 is a violation even if they never coexist.

from collections import defaultdict

class SafetyViolation(Exception):
    pass

class RaftChecker:
    """Checks Raft's safety properties over snapshots of every node after each step."""

    def __init__(self):
        self.leaders = {}                     # term -> node id ever seen as leader
        self.committed = {}                   # index -> (term, command) once committed
        self.commit_term = {}                 # index -> term of the node that first showed it committed
        self.applied = defaultdict(dict)      # node -> {index: command}

    def check(self, nodes):
        for n in nodes:
            # Election Safety: at most one leader per term, ever.
            if n.role == "leader":
                prev = self.leaders.setdefault(n.term, n.id)
                if prev != n.id:
                    raise SafetyViolation(f"two leaders in term {n.term}: {prev}, {n.id}")

        # Log Matching: same index and term implies identical prefixes.
        for a in nodes:
            for b in nodes:
                if a.id >= b.id:
                    continue
                for i in range(min(len(a.log), len(b.log)) - 1, -1, -1):
                    if a.log[i].term == b.log[i].term:
                        if a.log[:i + 1] != b.log[:i + 1]:
                            raise SafetyViolation(f"log matching broken at {i}: {a.id}/{b.id}")
                        break

        for n in nodes:
            # Committed entries never change (part of State Machine Safety).
            for i in range(n.commit_index):
                e = (n.log[i].term, n.log[i].command)
                if self.committed.setdefault(i, e) != e:
                    raise SafetyViolation(f"committed index {i} changed on {n.id}")
                self.commit_term.setdefault(i, n.term)
            # Leader Completeness: leaders of LATER terms hold every committed entry.
            # A stale leader in an older term may legitimately lack them.
            if n.role == "leader":
                for i, e in self.committed.items():
                    if n.term <= self.commit_term[i]:
                        continue
                    if i >= len(n.log) or (n.log[i].term, n.log[i].command) != e:
                        raise SafetyViolation(f"leader {n.id} term {n.term} lacks index {i}")
            # State Machine Safety: no two nodes apply different commands at an index.
            for i, cmd in n.applied_log():
                for other, seen in self.applied.items():
                    if i in seen and seen[i] != cmd:
                        raise SafetyViolation(f"{n.id} and {other} applied different entries at {i}")
                self.applied[n.id][i] = cmd

Leader Append-Only needs one more piece of history, each leader's log at its previous check, and is left out for space. Run the checker after every simulated step: a transient violation that later heals is still a bug, and easiest to diagnose where it occurs.

Advertisement

Deterministic simulation

Unit tests check the cases you thought of. The failures that matter in consensus come from interleavings you did not think of: a vote reply that arrives after a crash and restart, an AppendEntries delayed past a new election. Deterministic simulation explores those interleavings by running the real Raft code inside a single-threaded event loop where everything nondeterministic, meaning the network, the disk, the clock and randomness, is a fake controlled by one seeded random number generator.

The payoff is reproducibility. A failure prints its seed, and rerunning that seed replays the exact same sequence of events, so a one-in-a-million interleaving becomes a debuggable test case. The design requires discipline in the implementation: no direct calls to the system clock, no threads the scheduler does not own, and no iteration over unordered maps where the order affects behaviour.

A deterministic Raft simulator: one seed decides every interleavingSeeded schedulerPRNG(seed), event queuenode 1real Raft codenode 2real Raft codenode 3real Raft codeFake networkdrop, delay, reorder, partitionFake diskfsync boundary, crash loses unsyncedFake clocktimers fire only when scheduleddeliversendpersistset timerevents go back into the queueInvariant checker after every stepelection safety, log matching, leader completeness, state machine safetyOn violationprint seed + event trace; rerun = exact replaysnapshot statesEverything nondeterministic is behind an interface the scheduler controls
Figure 1. The Raft code is unchanged; its network, disk and clock are fakes driven by one seeded scheduler, and the checker runs after every event.
import heapq, random

def simulate(seed, steps=200_000):
    rng = random.Random(seed)
    net, disk, clock = FakeNetwork(rng), FakeDisk(), FakeClock()
    nodes = [RaftNode(i, net, disk.for_node(i), clock, rng) for i in range(3)]
    checker, queue, trace = RaftChecker(), [], []
    for n in nodes:
        n.start(queue)                        # schedules initial election timers

    for _ in range(steps):
        if not queue:
            break
        at, seq, event = heapq.heappop(queue)
        clock.now = at
        trace.append(event)
        fault = rng.random()
        if fault < 0.002:
            victim = rng.choice(nodes)
            disk.crash(victim.id)             # drops writes not yet fsynced
            victim.restart(queue)             # rebuilds state from disk only
        elif fault < 0.004:
            net.partition(rng.sample([n.id for n in nodes], rng.randint(1, 2)))
        elif fault < 0.006:
            net.heal()
        event.fire(queue)                     # may send, persist, set timers
        try:
            checker.check(nodes)
        except SafetyViolation as v:
            raise SafetyViolation(f"seed={seed}: {v}; last events: {trace[-20:]}")

for seed in range(10_000):                    # in CI; nightly runs use more seeds
    simulate(seed)

The fake disk is the component most often done badly. It must model the gap between a write and a durable write: on a simulated crash, everything written since the last fsync is lost, and optionally a torn final record is left behind. Without that, the simulator cannot find the most common class of real Raft bug, persistence that happens too late.

Which faults to inject

  • Message faults: drop, duplicate, delay and reorder each message independently. Raft must tolerate all four; duplicates and reordering catch code that assumes replies match the latest request.
  • Crash and restart with loss of unsynced writes, including crashes between persisting a vote and sending the reply.
  • Partitions, both clean splits and partial ones, where A reaches B and B reaches C but A cannot reach C, and asymmetric links where messages flow only one way.
  • Timing: slow disks that make fsync take longer than the election timeout, paused processes (garbage collection, VM pauses) and clock jumps, which matter as soon as you add leader leases.
  • Membership changes and snapshots concurrent with all of the above, because these code paths run rarely in production and are the least tested.

Beyond invariants: linearizability and model checking

The five properties are about the log. Clients care about something else: that the system behaves as if each operation took effect at one instant between its call and its return. An implementation can keep a perfect log and still serve stale reads, for example when a deposed leader that does not yet know it was replaced answers reads from local state. To catch this, record every client operation's start time, end time, input and output during a simulation or a real-cluster test, then run a linearizability checker such as Porcupine or Knossos over the history. This is how Jepsen tests find problems in real systems, including its analyses of etcd.

Model checking attacks the algorithm rather than the code. Diego Ongaro's dissertation includes a TLA+ specification of Raft, and a model checker such as TLC explores every reachable state of a small configuration, for example three servers, two terms and two client values, and checks invariants exhaustively within those bounds. It is the right tool when you change the protocol itself, such as adding a read optimisation or a new membership scheme. It does not test your code; the simulator does that.

Bug classes real implementations shipped

  • Committing by counting old-term replicas. The Raft paper's Figure 8 shows how a leader that marks a previous-term entry committed because a majority stores it can see that entry overwritten by a later leader. The fix is the current-term rule from the algorithm section.
  • Persisting after replying. A server that sends its vote or its AppendEntries success before the vote or entries are on stable storage can, after a crash, vote twice in one term or forget an acknowledged entry. Two leaders in one term follows. Asynchronous write paths and batched fsyncs are where this hides.
  • Single-server membership changes. In 2015 Ongaro reported on the raft-dev list a safety bug in the single-server membership change algorithm from his dissertation: with changes made by leaders of different terms, two disjoint majorities could form. The published fix is that a new leader must commit an entry from its own term before it starts a configuration change.
  • Liveness under partial partitions. Raft's safety holds under any partition, but its liveness can collapse. In a partial partition, a server cut off from the leader but not from others keeps timing out and raising the term, deposing a healthy leader again and again. Cloudflare's November 2020 control-plane outage involved an etcd cluster in this kind of partial failure; Jensen, Howard and Mortier examined Raft's behaviour in that scenario in a HAOC 2021 paper. PreVote and CheckQuorum are the standard defences, discussed in Raft consensus in depth; test them under partial partitions, not only clean ones.
  • Snapshot and log truncation races. Installing a snapshot while entries after it are still being appended, or truncating the log before a snapshot is durable, can lose committed entries on restart.

Worked example: reading a failing seed

Suppose seed 4417 fails with "two leaders in term 5: 2, 1". The trace's last events, condensed, read:

  1. Term 4: node 1 is leader. The simulator crashes node 1.
  2. Node 2 times out, starts term 5, and sends RequestVote to nodes 1 and 3.
  3. Node 3 records votedFor=2 for term 5 in memory, sends its grant, and queues the write to disk.
  4. Node 2 becomes leader of term 5 with votes from 2 and 3.
  5. The simulator crashes node 3 before its fsync. On restart node 3 reads term 4 and no vote.
  6. Node 1 restarts from disk as a follower in term 4, times out before node 2's heartbeat reaches it, starts term 5 and asks node 3 for a vote. Node 3 has no record of voting in term 5 and grants it.
  7. Node 1 becomes leader of term 5. The checker sees leaders 2 and 1 in term 5.

The cause is step 3: node 3 sent its vote before the vote was durable, so after the crash it could vote again in the same term. The fix is to make the reply wait for the fsync, then rerun seed 4417 to confirm it passes, and keep it as a regression test. Unit tests of the RequestVote handler would never find this; it needs a crash at one precise point.

Trade-offs between testing approaches

ApproachFindsMissesCost
Unit testsHandler logic errorsInterleavings, crashesLow
Deterministic simulationCrash and timing bugs in your code, reproduciblyBugs in real OS, disk and network behaviourMedium: requires injectable I/O
Jepsen-style real-cluster testsLinearizability violations end to endRare interleavings; failures are hard to replayMedium to high
TLA+ model checkingProtocol design errors, exhaustively in small boundsImplementation bugsSpecialist skill

Beyond unit tests, deterministic simulation is the best first investment: it finds bugs in the code you ship. For comparisons with the algorithm Raft replaced, see Paxos.

What to do next

  1. Write the five safety properties as an executable checker over node snapshots, keeping history across steps.
  2. Put the network, disk, clock and randomness behind interfaces and build a seeded single-threaded simulator.
  3. Make the fake disk lose unsynced writes on crash, then audit every RPC reply to confirm it follows the fsync it depends on.
  4. Inject message faults, crashes, partial and asymmetric partitions, slow fsync and membership changes; run thousands of seeds in CI.
  5. Record client histories and check them for linearizability, especially reads after leader changes.
  6. Before changing the protocol, model it in TLA+ and check the same invariants in small bounds.
  7. Keep every failing seed as a regression test.
Key takeaway: Raft is easy to understand and hard to implement correctly, because the bugs that matter live in persistence order, crash timing and partial failures rather than in the rules themselves. Turn the paper's five safety properties into a checker that runs after every step, run the real implementation inside a deterministic simulator whose disk loses unsynced writes, inject message, crash, partition and timing faults across thousands of seeds, and check client histories for linearizability. Use TLA+ for protocol changes. Each failing seed is an exact replay, which makes rare consensus bugs ordinary debugging work.