Prototypes QuorumP29 · Distributed database

The fuzzer

Each version of Raft runs on the same seeds. A run is eight simulated seconds of five clients and a fault every 100 to 400 ms: power cuts, election storms, split networks and lost messages. After each run the history is checked for linearizability, and the simulation watches that no term has two leaders and that every node applies the same entry at each index.

VersionSeedsPromisesWhat brokeFirst broken seed

How Quorum is built
  1. The whole cluster is one simulation in the page: five nodes, the network between them, their disks and five clients, on a clock that moves only when the next event is due. Every delay and every chance comes from the seed instead of the real time, so a seed gives the same run every time. That is how the clock can go back: the page builds the run again from its seed with the same actions.
  2. Each node runs Raft. A leader is elected by a majority of votes for a term. It appends each write to its log, sends it to the others, and applies it once a majority has it on disk. A node writes its vote and its entries to disk before it answers.
  3. A read goes through the leader too. Before it answers, the leader checks with a round of heartbeats that a majority still follows it, so an old leader cut off from the others cannot answer with old data.
  4. Each node’s disk holds a write-ahead log and a checkpoint. A record in the log is its length, a CRC-32 of its body and the body. Pulling the power keeps only what was flushed, and can cut the last record in the middle. On restart the node reads its checkpoint, then the log up to the first record whose checksum fails.
  5. The data is a B+tree. A checkpoint writes it as pages to a new file and renames the file, so a crash leaves either the old checkpoint or the new one whole. A follower too far behind gets the leader’s checkpoint instead of the entries the leader no longer keeps.
  6. Each client’s last request and its answer are kept in the tree, so a request that reaches the log twice is applied once.
  7. The checker places every operation at one moment between its call and its answer, so that each read returns the last value written. It is Wing and Gong’s search with a memo of states already tried, run for each key on its own. An operation that timed out may have happened at any moment after its call, or never.
  8. Four versions of Raft each have one known mistake. The fuzzer runs them on the same seeds and finds the seeds where they break a promise. While this prototype was built, the fuzzer found two mistakes in the correct version: a client that sent a request twice, and a follower that answered for entries it had received again before they were flushed. Both are fixed, and the correct version passes every seed the fuzzer runs.