Five glowing amber orbs on a dark surface, three connected by thin gold threads in a quorum, two sitting isolated and unconnected Tech
AI-generated, Working Theory
Tech · ◉ Evergreen

They all have to agree — and any of them can disappear.

by · ·5 min·Working Theory

Quorums, leaders, logs, and elections — how a cluster of machines agrees on one truth when any of them might be dead, slow, or lying, explained in plain English.

Here is a problem that sounds trivial until you try it. You have five copies of the same service, running on five machines, so that if one dies the others carry on. A write comes in — say, the balance is now $40. You need all five copies to end up believing the same thing, in the same order, forever. That’s it. That’s the whole job.

Now add the two facts that make it hard. Machines crash without warning, mid-sentence. And the network doesn’t just drop messages — it delays them, so you can never tell “he’s dead” apart from “he’s slow and the answer is coming.” You are trying to get a room full of people to agree on one number, in the dark, where anyone might have quietly left, and a reply you’re waiting on might arrive in a millisecond or never.

The naive move is “ask everyone, proceed when all five say yes.” It fails the instant one machine is down for maintenance — now you can never write anything, which is the opposite of what the replicas were for. So you weaken it: proceed when a majority says yes. Three of five. This one change is the seed of the whole field, and it works because of a small piece of arithmetic: any two majorities of the same group must overlap in at least one member. If a decision needed three votes, and a later decision also needs three, at least one machine was in both rooms and remembers the first. Majorities can’t contradict each other without someone noticing. That overlap is the load-bearing wall.

leader f1 f2 f3 f4 ✓ ✓ ✓ commit at 3 of 5 = majority → durable minority — cannot commit
A write is durable once a majority has it. The stranded minority can't commit anything that would contradict the majority — the two sets always overlap. Original diagram · Working Theory

The second idea keeps things sane: pick one leader. Instead of every machine proposing values and arguing, the group elects a single machine whose job is to receive writes, stamp them into an ordered list — a log — and push that log to the others. Order stops being a debate because one machine decides it. A write is committed once the leader has copied it to a majority; at that point it survives any single crash, because a majority remembers it and the next leader will be forced to inherit it.

Then the leader dies — because everything dies. The others notice they’ve stopped hearing from it (a timeout, that same “dead or just slow?” guess), and they hold an election. Whoever wins gets a new term — think of it as a numbered era — and the term number is how the system fences out zombies. If the old leader was only slow and comes back mid-write, its messages carry a stale term, everyone else has moved on, and its writes are refused. There is only ever one current era, and stale eras can’t commit.

That’s essentially the whole shape of Raft, the algorithm Diego Ongaro and John Ousterhout published in 2014 with the explicit goal of being understandable — because its predecessor, Leslie Lamport’s Paxos, was correct and famously almost impossible to hold in your head. Same core guarantees, deliberately teachable decomposition: leader election, log replication, and the safety rules that keep a new leader from erasing committed history. If you’ve used etcd, Consul, CockroachDB, or the coordination layer under a dozen other systems, you’ve been standing on this.

One humbling footnote worth carrying: in 1985 Fischer, Lynch, and Paterson proved that in a truly asynchronous network — no bound on message delay — no deterministic algorithm can guarantee consensus if even one process can fail. Not “it’s hard.” Impossible. So how does any of this run in production? By cheating honestly: real systems use timeouts and a little randomness to make progress in practice, trading the theoretical guarantee of always eventually deciding for the practical one of deciding, almost always, quickly. The impossibility result isn’t a wall you break; it’s a reminder of which corner you agreed to give up.

The reason any of this belongs in a builder’s head, not just an academic’s: consensus is expensive, so spend it narrowly. Every committed decision costs a round trip to a majority. That’s the right price for the handful of facts that must be singular and ordered — who is the leader, what the cluster config is, the ordered log of transactions. It is a ruinous price to pay for your high-volume data. The mature pattern is a small, consensus-guarded core deciding the few things that truly cannot disagree, with the bulk of the system kept deliberately outside that core, allowed to be eventually consistent. Knowing which facts have to go through the expensive room — and refusing to send the rest — is most of the engineering.

Sources

  • Leslie Lamport, "Paxos Made Simple" (2001)
  • Diego Ongaro & John Ousterhout, "In Search of an Understandable Consensus Algorithm" (2014)
  • Fischer, Lynch & Paterson, FLP impossibility result (1985)

Liked this? Get the next one in Working Theory.

Going weekly in August (it's in beta now). One genuinely interesting read on building, the brain, and the science most people missed.

Subscribe →
Got a reaction, a counter-example, or something I missed? Reply by email — I read everything.
◉ join in

Where have you hit this — in a product you use, or one you're building?

Threads open here soon. For now, the conversation lives two clicks away — discuss on GitHub, or just reply by email. I read and answer everything.