NET·VI Networks Chapter 44 of 65

The parliament of Paxos

The island’s legislators keep wandering off to the market, messengers get lost, and still a law must be passed: one law, written down the same way in every book. This chapter retells Leslie Lamport’s parable, the one reviewers told him to rewrite without the Greeks, and lets you run the parliament and break it: remove the chair, call up a storm, bribe a general.

Further 70 minutes Networks Concurrency History
NET·VI

Networks

  1. 41 Networks
  2. 42 TCP/IP
  3. 43 The web
  4. 44 Distributed you are here

Builds on: 43 · Anatomy of this page 39 · Races

What you will take away

  • why “distributed” is hard: telling the dead from the slow, ordering events without a shared clock, agreeing despite lost messages and betrayal
  • how consensus works: quorums, the Paxos protocol, and Raft with its elected chair and its log; and what “the law is passed” means
  • what a system chooses when the network is cut in two (the CAP theorem), and how eventually consistent systems live

6How do billions of computers work together when none of them is in charge?

The last chapter ended with three copies of one account that had drifted apart: some messages got lost, and to the question “how much money is in the account?” the system had three answers. Resending messages, as TCP does, won’t help here: a copy can die along with a message or fall silent for a long time. How do you make a thousand unreliable machines behave like one reliable machine when none of them is in charge and there is nobody to appoint one? In the late 1980s Leslie Lamport set out to prove that this is impossible. The proof didn’t work out. An algorithm did, and Lamport told it as a parable about the parliament of a Greek island.

Report from the dig

What follows is the parable, retold. The main text is the island’s chronicle, and boxed comments translate it into the language of computers, as Marzullo’s did. The chronicle can also be played: you will send legislators to the market, call up a storm, bribe a general, and see what holds and what breaks.

An island where nobody is in charge

The parliament of Paxos passed laws: “no selling olives after sunset,” “goats graze on even days.” Each legislator kept his own book of laws, and the books had to agree: if one book had the olive law as its seventh law, then every other book had it as the seventh too, and not the goat law. There was no common book, and no chief legislator either. The legislators communicated by messenger.

The trouble was that the Paxons were busy people. A legislator could leave the chamber for the market at any moment and come back an hour later, or a week later. A messenger could lose his way, linger in a tavern or arrive twice. On the other hand, to the island’s credit, nobody lied: a legislator at the market said nothing, and a messenger who did arrive delivered his message word for word.

The list of troubles looks familiar: Chapter 42 dealt with lost packets, and Chapter 39 put the events of threads in order. But a distributed system has a trouble that neither of those had, and all the other troubles start from it.

Dead, or just thinking?

The legislator Alcaeus is waiting for an answer from Bias. Two minutes have passed with no answer. Did Bias go to the market? Did the messenger get lost? Or is Bias thinking over his answer and about to send it? From inside the chamber these cases look the same: in all three, Alcaeus sees one thing, silence. All he can do is agree to wait a set time and then count the silent one as gone. TCP’s timer in Chapter 42 answered the same question, how long to wait for an ack, except that there a mistake cost one needless resend, and here it costs burying the living. The Pathfinder’s watchdog timer in Chapter 39, which rebooted the computer when a high-priority job missed its deadline, was the same kind of agreement. Suppose Bias sends a messenger with “I’m here” every 100 ms, and Alcaeus decides how long to wait before declaring him gone.

Bias is alive the whole hour, yet a short timeout buries him dozens of times an hour. A long timeout never errs, but when death does come it notices only seconds later, and all that time the parliament waits for a dead man. There is no perfect timeout; it is always a trade between speed and false alarms. Our model is not to blame. In 1983 Michael Fischer, Nancy Lynch and Michael Paterson proved that if messages can be delayed for arbitrarily long, no deterministic algorithm can guarantee that the nodes will agree, even when only one of them may fail. The result goes by their initials, FLP; the journal version of the paper came out in 1985. It doesn’t forbid agreeing. It forbids guaranteeing agreement in finite time whatever the network does. That is why all the algorithms in this chapter are built the same way: whatever happens, they never pass two different laws, and they promise to pass some law only once the network calms down for a while.

In a distributed system you can’t tell a dead node from a slow one. Every decision that “he is dead” is a guess based on a timeout, and the algorithm must stay correct even when the guess is wrong: the “dead” may come back and start talking.

Clocks without a clockmaker

Every legislator had his own hourglass, and each ran as it pleased. So the island’s books said “after” rather than “at noon,” and for good reason: trust the clocks and you get nonsense. Here are three legislators, who are also three threads from Chapter 39, exchanging messengers through queues. Bias’s clock is 40 ms slow and Glaucus’s 25 fast. Each writes in a shared journal the time by his own clock.

By the wall clock, Bias received the messenger about olives before Alcaeus proposed the law. A clock running 40 ms slow put the effect ahead of its cause. Servers synchronize their clocks over the network, but a discrepancy of milliseconds always remains, while events on fast machines are microseconds apart. The second list is in a different order, and everything in it is where it belongs: every answer comes after its question.

Lamport proposed measuring time in a distributed system by causes and effects. Event $a$ happened before event $b$ (written $a \to b$) if $a$ could have influenced $b$. That happens in three cases: both events belong to one legislator and $a$ came first; $a$ is sending a messenger and $b$ is receiving him; or a chain of such steps leads from $a$ to $b$. If neither $a \to b$ nor $b \to a$, the events are concurrent: neither could know of the other, and asking which of them came first “in reality” makes no sense, as in relativity.

The happened-before order is kept by Lamport clocks, which are clocks in name only: each is a counter kept by one legislator. There are two rules. Before every event, add one to your counter, and when sending a messenger, give him the counter’s value to carry. On receiving a messenger, set your counter to the larger of your own value and the one he brought, and add one to that too. That’s how, in the cell above, Bias, whose counter stood at 0, received stamp 2 from Alcaeus and went to 3.

If $a \to b$, then the stamps satisfy $C(a) < C(b)$.

It is enough to check one step of the chain: after that the inequalities add up, as in $C(a) < C(c) < C(b)$. If $a$ and $b$ belong to one legislator and $a$ came first, then between them the counter only grew, and one more was added before $b$. If $a$ is a sending and $b$ the receiving, then the messenger brought $C(a)$, and the receiver set his counter to at least $C(a) + 1$.

The converse is false, and that’s worth remembering. In the cell, “Bias got ‘I’m for it’” and “Glaucus went to the market” share stamp 7, and if Glaucus had done a couple more things first, his trip to the market would have received a larger stamp than Bias’s receipt of the answer, though the two events are concurrent. A smaller stamp doesn’t mean “happened earlier”; it means “didn’t happen later.” But if you sort all events by stamp, and equal stamps by the legislator’s name, you get one common row in which no effect comes before its cause. That row is what Lamport proposed to use: if every legislator carries out the laws in that order, their books will agree. Draw your own history of messengers and see how the stamps behave.

Three legislators, three lines; time runs to the right. “Event”: tap a line, and something of his own happens to that legislator. “Messenger”: tap a sending event or a spot on a line, then the receiver’s line. “Compare”: pick two events, and the widget says which happened first or that they are concurrent, and highlights everything that could have influenced the one you picked. The switch shows Lamport clocks or vector clocks.
Vector clocks: telling “before” from “concurrent”

Lamport clocks can’t tell you whether two events are concurrent. Vector clocks can. Several authors came up with the idea independently in the early 1980s, and in 1988 Colin Fidge and Friedemann Mattern, also independently of each other, gave it a name and a rigorous theory. Here every legislator keeps a whole vector, with one number for each legislator on the island: how many of that legislator’s events he knows about. He increases his own number before every event, the messenger carries the whole vector, and the receiver takes the maximum position by position. Then $a \to b$ if and only if $a$’s vector is no greater than $b$’s in every position and smaller somewhere; and if it is greater in one position and smaller in another, the events are concurrent. The price is size: with a thousand nodes, a vector of a thousand numbers rides along with every message. The Dynamo database, which we come back to at the end of the chapter, kept similar version vectors for every object.

The Lamport counter with sending and receiving is the chapter’s first task, “The island’s clocks.” The catch is that the legislators’ journals arrive separately, and a messenger’s arrival in one journal can’t be processed until his departure in another one has been.

A bribed legislator

The Paxons were honest. Had a liar turned up among them, the island’s protocol would have had a much harder time, and Lamport had already worked out how much harder before Paxos, with a different cast.

Several divisions are besieging a city. The commander decides whether to attack or retreat and sends the order to his lieutenants by messenger. The lieutenants can talk among themselves. There may be a traitor among the generals, who can lie to anyone about anything, including telling different generals different things. All the loyal lieutenants must carry out the same order, and if the commander is loyal, it must be his. A failure in which a node sends contradictory lies is called Byzantine. A broken sensor, damaged memory or a hacked server behave this way.

Generals and a traitor. Choose how many generals there are, who the traitor is and what he says: tap the circles on the arrows that lead from the traitor, and the order on them flips. The council follows the algorithm: the lieutenants first get the commander’s order, then tell one another what they heard, and each decides by majority (on a tie, retreat). The verdict below says whether the council held.

With three generals the traitor can always wreck the council, and the reason is clear once you stand in the first lieutenant’s shoes. The commander tells him “attack,” and the second lieutenant says the commander told him “retreat.” Which of the two is lying? If the second lieutenant is, then the loyal commander wants an attack. If the commander is lying, then the second lieutenant may have honestly passed on what he was told. The first lieutenant sees the same thing in both cases and can’t tell them apart, yet the right actions differ. With four generals, each lieutenant has three votes, and one lie can’t outweigh the rest. The cell below checks this by trying every possible betrayal against the algorithm from the 1982 paper, OM(1).

With three generals, a traitorous lieutenant wrecks the council in two cases out of twelve: when the loyal commander ordered an attack and the traitor reported “retreat,” the loyal lieutenant gets a tie and retreats against orders. With four generals, not one council out of 32 is wrecked. Trying every case proves this for one algorithm only, while the 1980 paper proves it for any algorithm. The overall result: to withstand $f$ traitors with ordinary, unsigned messages, a system needs $3f + 1$ nodes, and no fewer will do. Byzantine fault tolerance is built where a mistake is costly or the participants can’t be trusted: according to public sources, it is used, for example, in the flight control systems of the Boeing 777 and 787. The parliament of Paxos, like most systems in data centers, assumes that nodes are honest and only crash, and so it pays less for failures: to survive $f$ legislators at the market, $2f + 1$ legislators are enough. The next section explains why that many.

How to pass one law

Start with the simplest case: the parliament must pass one law. Anyone may propose it, there may be several proposals, and in the end every book must name the same law. This is the problem of agreement, or consensus. It has three requirements. Only a proposed law can be passed: you can’t reach agreement by always passing “nothing.” Only one law is passed, never two different ones. And if a majority of the legislators are in the chamber and messengers get through, some law is passed sooner or later.

The first idea is a vote: a law passes if more than half the legislators vote for it. Half isn’t an arbitrary choice: any two majorities overlap, so every decision is known to at least one participant of any later decision.

If the parliament has $n$ legislators, and $A$ and $B$ are two groups with more than $n/2$ legislators each, then $A$ and $B$ have a member in common.

If they had no member in common, $A$ and $B$ together would contain $|A| + |B| > n/2 + n/2 = n$ different legislators, and there are only $n$.

A group whose agreement is enough for a decision to count is called a quorum. A majority is the simplest quorum, and on five legislators you can check the property from the theorem by trying every group.

Pairs can miss each other; triples always overlap. So a parliament of five has nothing to fear from two members leaving: the remaining three are a quorum. Out of $2f + 1$ legislators you can lose $f$ and no more: beyond that no quorum gathers, and the parliament freezes, though it makes no mistakes.

But a quorum alone isn’t enough. Suppose two legislators propose different laws at the same moment, one about olives and one about goats, and the messengers carry the proposals around in different orders. If each legislator votes for the first proposal he hears and never changes his vote, you can end up with two votes for olives, two for goats and one lost messenger, and no law, ever. If you allow changing your vote to the latest proposal heard, some law will get a majority sooner or later, but afterward another one may get a majority too. The Paxons solved this as follows.

A law is proposed by a priest, and each of his proposals is a ballot with a number of its own; different priests use different numbers. A ballot has two phases. In the first, the priest asks the legislators: “Do you promise not to take part in ballots numbered lower than mine? And what have you already voted for?” A legislator promises, unless he has already promised someone with a higher number, and names his last vote. Once the priest has promises from a quorum, he looks at those votes. If anyone in the quorum has already voted, the priest must propose the law from the highest-numbered ballot among them, and his own law is dropped. Only if nobody has voted does he propose his own. In the second phase the priest sends out “we are voting on such-and-such law in ballot number such-and-such,” and a legislator votes unless that would break his promise. When a quorum has voted for one law in one ballot, the law is passed.

An experiment. Two priests at once propose a tax on olives and a tax on goats; there are five legislators, every fifth messenger gets lost, and messengers take from 10 to 100 ms on the road. A priest who doesn’t see his law passed starts a new ballot with a higher number. We compare the Paxos protocol with the plain rule “vote for the highest-numbered proposal you have heard,” which has no first phase, and run the scenario two thousand times with different random delays and losses. The queue of events is a heap from Chapter 18.

Without promises the parliament passes both laws in nearly every run: first a majority votes for olives, then the goat priest, tired of waiting, starts a ballot with a higher number, and the majority votes again. Two books on the island will record different laws, and the mistake can never be undone. With promises there isn’t a single such case in two thousand: the second priest, collecting promises, learns of the votes for olives and proposes olives himself. The “none” column holds a zero too: with every fifth messenger lost and two rival priests, some law still got passed within five seconds. Nothing guarantees that, since FLP forbids it, and with bad luck the two priests can keep outbidding each other’s numbers forever. That’s why in practice a single priest is chosen.

The chair

The island can’t make do with one law. It needs a series of them, the seventh, the eighth, the hundredth, in the same order in every book. You could hold a synod for every place in the book, but in quiet times it’s wiser to choose one priest for a long stretch: a chair. While the chair is alive, he alone proposes laws, he has no rivals, and the first phase needn’t be repeated. If he goes to the market, a new one is chosen. That’s how “multi-Paxos” works, and so does the algorithm in more common use today.

Both algorithms aim at a shared book of laws, which Raft calls the log. If all the legislators carry out the same laws in the same order, starting from identical books, their books will agree. This is Lamport’s idea from the 1978 paper, state machine replication. Agree on the log, and determinism does the rest.

In Raft a legislator always plays one of three roles: follower, candidate or chair. Time is divided into numbered terms, and each term has at most one chair. The chair regularly sends out messengers, with new laws or with no more than “I’m here.” A follower who hasn’t heard from the chair for a long time declares a new term, votes for himself and asks the others for their votes: this is leader election, the choosing of the chair. Each legislator votes once per term, and only for a candidate whose book isn’t behind his own. Whoever gathers a majority is chair. And so that two legislators don’t stand at once and split the votes forever, each has his own random timeout, from 150 to 300 ms in the paper’s example. Usually someone wakes up first and gathers the votes while the others are still waiting.

The chair first adds a new law to his own book and then sends it to the others. Once the law is written in a majority of the books, the chair declares it passed, and from then on no later chair will erase it: a new chair needs the votes of a majority, majorities overlap, and nobody votes for a candidate whose book lags behind. Followers whose books have drifted from the chair’s, such as ones back from the market with laws that nobody passed, rewrite the part that differs from the chair’s book. Try it.

A parliament on Raft: five legislators, with messengers as dots on the lines. Next to the circle are the books: the number in a cell is the term in which the law was written, and filled cells are passed laws. Tap a legislator to send him to the market or bring him back. “New law” hands a proposal to the chair. “Storm” cuts two off from the other three. Try this: remove the chair; call up a storm while the chair is in the smaller part and propose a law; then bring everyone back.

If the chair ends up in the smaller part of the island, he keeps sending out laws, but he can’t gather a quorum, and none of his laws will pass: they hang unfilled in the books. Meanwhile the larger part elects a new chair with a higher term and passes laws without him. When the storm dies down, the old chair hears of the new term and becomes a follower, and his unpassed laws are erased. Nothing that was passed is lost.

Raft has one subtlety that its authors take apart in a figure of its own. A law from an old term can sit in a majority of books and still not be passed. Suppose the chair of term 2 managed to write a law only in his own book and one neighbor’s before he went to the market. In term 3 the chair was a legislator whom that law never reached; he wrote his own law in the same place in the book, also only in his own book, and also left. In term 4 the first one came back, became chair again and copied the term-2 law to a third legislator: now the law is in a majority of books. If the chair declares it passed and leaves at once, the election can be won by the owner of the term-3 law: his book counts as newer, because its last law is from term 3 and not 2, and three legislators will vote for him. The new chair will overwrite that place in every book, and the “passed” law will vanish. So the rule is stricter: by counting copies, the chair passes only laws of his own term, and older ones become passed along the way, when a later law of the current term is passed. This rule is the task “The chair’s quorum.”

All of this can be programmed with the threads and queues of Chapter 39, where each legislator is a thread and the messengers are messages in his queue. The chapter’s last task, “Electing the elder,” is an older and simpler way to choose a chair, the bully algorithm: the eldest of the living becomes chair, and the younger ones give way to him as soon as they hear from him.

The storm

The chair in the smaller part of the island, in the widget above, behaved modestly: he wrote laws down and sent them out but didn’t declare them passed, and a conscientious legislator in the smaller part, asked “what is the law on goats right now?” ought to answer “I don’t know, there’s no word from the others.” That is a choice, and it has a price: part of the island stops working while the storm lasts. You can choose otherwise: let each part pass laws on its own, and sort things out after the storm. Then both parts work, but the books diverge. There is no third option, and blaming Paxos or Raft for it is pointless: this is a theorem.

In 2000, at the Symposium on Principles of Distributed Computing (PODC), Eric Brewer of Berkeley presented a conjecture: a distributed system can’t at the same time provide consistency (every read sees the latest write), availability (every working node answers) and tolerance of network partitions. In 2002 Seth Gilbert and Nancy Lynch of MIT proved it, and it became the CAP theorem, named for the initials of consistency, availability and partition tolerance. It is often summed up as “pick two out of three,” and in 2012 Brewer himself explained why that is misleading. Nobody picks network partitions; they happen. The choice arises at the moment of a partition: refuse to answer, or answer and risk diverging. While the network is whole, you can have both.

Many systems choose availability in a storm. Let the copies diverge, as long as each half of the island keeps working; when the connection returns, the copies will be merged. This guarantee is called eventual consistency: if the changes stop, all copies eventually become identical, and until then they may answer differently. You have already used such a system: DNS from the last chapter. Change a site’s address, and for ten more minutes, until the records in caches expire, part of the world goes to the old one. Nobody waits for the whole world to agree, and nobody suffers for it.

The classic example is an online store’s shopping cart, from Amazon’s 2007 paper on its Dynamo storage system. The paper says plainly that customers must be able to view their cart and add items to it even if “disks are failing, network routes are flapping, or data centers are being destroyed by tornados.” The cart is kept in several copies, and during a storm each copy accepts changes. How do you merge them afterward?

“The last write wins” lost the wine: the customer put it in the cart, and it vanished. The union of the sets loses nothing, but it brings back the deleted cheese. A cart on Dynamo merges its copies by union, and the paper says so directly: an “add to cart” operation is never lost, but “deleted items can resurface.” For a cart that is sensible: the customer will take out the extra cheese himself, while vanished wine is a lost sale. For a bank account neither will do: there you choose consistency and a quorum, as in Raft, and accept refusing to answer while the storm lasts.

Billions of computers work together with nobody in charge because each layer of the network is built so that it depends on no one entirely. The packets of Chapter 41 pass through routers, each of which knows only where to send a packet next, and the routers build their routing tables by talking with their neighbors; Paul Baran’s network was designed to survive the loss of many nodes. On top of unreliable delivery, the TCP of Chapter 42 settles things between the two ends alone: sequence numbers, acks and a timer make a reliable stream without asking the network’s permission, and congestion control is millions of senders, each of which slows down by itself when it notices losses. The DNS of Chapter 43 is a tree in which every subtree is handed to its own owner, and its speed rests on caches with lifetimes. And where machines do need a single opinion, in a bank’s ledger, in GitHub’s database, in the store behind Kubernetes, the consensus algorithms of this chapter provide it: a majority decides, with no chief, and any two majorities overlap, so two contradictory decisions can’t both pass. The failure of any one machine is routine here: a timeout, a new election, a new term. The price of this independence is known too. You can’t tell the dead from the slow, and when the network is cut you can’t both always answer and always answer the same, so every system chooses which mistake it will pay for.

Tasks

Three tasks for three levels of the island: clocks, quorum and elections. The last one is solved with threads and queues, and the tests run it several times with random messenger delays: a solution that works “usually” won’t pass.

The legislators’ journals are collected in a dictionary: a legislator’s name → the list of his events in order. An event is a tuple: ("local",) is something of his own, ("send", m) means he sent messenger number m, and ("recv", m) means he received messenger m. A messenger’s number is any immutable value, different for every messenger. Write lamport(trace): a dictionary with the same names and, for each, the list of Lamport clock stamps of his events. Every clock starts at zero; one is added before every event, and receiving a messenger sets the clock to the larger of its own value and the sending’s stamp, plus one. A messenger may be sent and never received: he got lost. If the journals are impossible (a messenger is received but nobody sent him, he is received twice, two messengers share a number, or the legislators wait for one another in a circle), raise ValueError.

For the council from the clocks.py cell, {"Alcaeus": [("local",), ("send", "m1"), ("recv", "m4"), ("local",)], "Bias": [("recv", "m1"), ("send", "m2"), ("recv", "m3"), ("send", "m4")], "Glaucus": [("recv", "m2"), ("send", "m3"), ("local",)]}, the answer is {"Alcaeus": [1, 2, 9, 10], "Bias": [3, 4, 7, 8], "Glaucus": [5, 6, 7]}.

The starter goes through the journals one after another and fails with a KeyError as soon as a receiver comes before the sender in the dictionary: the messenger hasn’t been “sent” yet. You can move a legislator forward as long as his next event isn’t “received a messenger nobody has sent yet.” At such an event he stops and waits.

Keep for every legislator a pointer to his next event, a list of those who can move, and a dictionary “messenger number → who is waiting for him.” When a messenger is sent, put the legislator waiting for him back on the ready list. That way every event is processed exactly once, in linear time.

When nobody else can move and someone’s journal isn’t finished, that’s a waiting circle: ValueError. Nonexistent and twice-sent messengers are easiest to catch up front, in one pass over all the sends.

This is the topological sort from Chapter 19: events are the vertices, and edges run from each event to the next one of the same legislator and from every sending to the matching receiving. The result is the graph of the happened-before relation, and the Lamport stamps can be computed in any order consistent with it. A waiting circle is a cycle in the graph: such a journal can’t exist, because you can only receive what has already been sent.

A Raft chair knows how many laws of his book each legislator already holds. Write commit_index(match, log, term, commit). match is a list: for every legislator, the chair included, the number of the last law he is known to have (laws are numbered from 1; zero means none). log is the chair’s book: log[k - 1] is the term in which law number k was written. term is the chair’s current term, and commit the number of the last law already passed. Return the new number of the last passed law: the largest k that a strict majority of legislators have and that was written in the current term. If there is no such k greater than commit, return commit: what was passed stays passed.

For example, commit_index([7, 7, 5, 3, 7], [1, 1, 2, 2, 3, 3, 3], 3, 3) is 7: three legislators out of five have the seventh law, and it is from term 3.

The starter has two bugs, and one of them is in the majority test. Of four legislators a majority is three, while 4 // 2 is two. A strict majority means more than half: have * 2 > n.

The second bug is the subtlety from the section on the chair: a law from an old term can’t be passed by counting, even if a majority has it. Check log[k - 1] == term. Once a suitable law of the current term is found, everything before it is passed along with it, so returning its number is enough.

Checking every k by counting over all legislators is $O(n \cdot L)$: with two thousand legislators and two hundred thousand laws, that’s hundreds of millions of steps. Sort match in descending order: the element at index n // 2 is the largest number a majority has. From there it’s enough to walk down to the first law of the current term.

After sorting in descending order, the legislators at indices 0 through n // 2, all n // 2 + 1 of them, a majority, have a number no less than best, while a number greater than best is no longer held by a majority. This is the median, and it can even be found in linear time on average, with the quickselect of Chapter 21. The term check is what the task was written for: without it the algorithm is right almost always, and that is why such a bug can live in a program for years, until three failures happen in a row.

On the island, the elder is the senior one present, the one with the highest number. Nobody knows in advance who is present. Write elect(me, ids, inbox, send, timeout); it will be run in a separate thread for every legislator present, all at the same time. me is your own number, ids the numbers of all the legislators on the island, including those at the market. inbox is a queue.Queue of messengers to you; each messenger is a tuple (kind, sender). send(to, kind) sends legislator to the messenger (kind, me); the kind is "ELECTION", "OK" or "LEADER". A messenger takes from 1 to 30 milliseconds, and messengers to absent legislators vanish. timeout is how many seconds it is reasonable to wait for an answer (0.15 in the tests). The function must return the elder’s number, and everyone present must return the same one: the highest number among those present.

This is the bully algorithm, described by Hector Garcia-Molina in 1982. Send "ELECTION" to everyone senior to you. If nobody answers "OK" within timeout, there is nobody senior, and you are the elder: send "LEADER" to everyone and return your own number.

While you wait, answer the others: to an "ELECTION" from someone junior, reply "OK", or the junior will decide there is nobody above him. If you get "LEADER", return the sender’s number. Waiting with a limit is convenient this way: inbox.get(timeout=remaining) raises queue.Empty if nothing arrived in that time.

If a senior answered "OK", wait for his "LEADER", for longer, say 4 * timeout, and keep answering juniors. If it doesn’t come, the senior has gone off to the market in the meantime, and you start the election over.

Everything rests on the timeout, and so on the assumption that messengers travel faster than timeout: the bully algorithm is correct in a synchronous system, where delays are bounded. If a messenger to a living senior is delayed longer, the junior will declare himself the elder, and the island will have two of them. Raft guards against this with term numbers and a quorum: a chair without a majority can’t pass anything. But the bully algorithm is simple and fast, and small systems on a reliable local network still use it.

What next

The parliament of Paxos has what it wanted: the books of laws at every end of the island agree, and neither failures nor storms can break that. Now open one of those books. It holds a hundred thousand entries in a row, and the island’s citizens come to it with questions.

Every question is a new loop over the whole book, written from scratch. A third question means a third loop, and when the book grows to a hundred million entries and the questions come a hundred a second, there will be no time to go through the whole book for each one. The nodes keep the data. How do you store and query millions of records so that a question like “how many laws about goats after the year 300?” can simply be asked, and the machine decides how to look for the answer? Databases have been working on that for half a century, and the next chapter is set in an archive where visitors come to you, the new archivist, with their questions.