Prototypes
Quorum
P29 · Distributed database
The cluster
The fuzzer
How it is built
Seed
New seed
Correct Raft
Random faults
0.02×
0.1×
0.5×
2×
Entries
Votes
Checkpoint
Answers
Committed entry
Not committed
Cut the leader’s power
Isolate the leader
Split 2 and 3
Heal the network
Restart every node
Lose 30 % of messages
Linearizable.
Waiting for operations.
History
Press the history to move the clock there.
Five clients
Stop when a promise breaks
Correct Raft
Reads from the leader’s own state
Answers before the disk has it
Does not write its vote to disk
Votes without comparing logs
How Quorum is built
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.
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.
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.
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.
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.
Each client’s last request and its answer are kept in the tree, so a request that reaches the log twice is applied once.
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.
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.