Lesson 032 · Phase 2, Mechanisms

Leader Election: Why Someone Has to Be in Charge

Picking one machine to be in charge is the easy half; the hard half is what the loser is still allowed to do.

19 min read

Lesson 32 · 32 published · 90 planned

On this page
The systems in this lessonUsed here: Marlow Books and Galewatch.

Made up for this course and reused from lesson to lesson so their numbers become familiar. None of them exist. All three

Marlow Books · A small online bookshop
Four people, one server and one Postgres database. About 40 requests a second on a normal day and ten times that in the week before Christmas. The one box that the early lessons stress until it breaks.
Galewatch · Telemetry for wind farms
Nine hundred turbines, a reading every two seconds, over links that drop for hours in bad weather and come back with a backlog. Dashboards that lag by seconds, reports that scan a year.

The job's log said it had moved all 133 turbines, verified all 133 and hit no errors. Every word of that was true, and a year of one turbine's readings was gone.

Galewatch, which collects a reading every two seconds from each turbine on the wind farms it monitors, spent that April night finishing the job lesson 031 priced. When Ardnave Point grew to three hundred turbines in March the farm outgrew two shards and the divisor in hash(serial) % 2 became a 3. March's hundred new turbines arrived empty; of the two hundred that already had a year behind them, a third kept their address and 133 did not. Their history did not move with them. 665 gigabytes, five gigabytes a turbine, had to be copied from the machine they were on to the machine the new arithmetic said they belonged on. At the two hundred megabytes a second lesson 013 measured on a restore drill, which lesson 031 flagged as the optimistic end, that is 3,325 seconds. Fifty five minutes of copying.

They wrote a mover, deployed it the way they deploy everything, to all four ingest boxes, and put lesson 007's lock in front of it: a claim row in the small Postgres that holds lesson 014's routing table, the one database in the fleet that is not a shard. First box to insert the row owns the move. Lesson 007 was blunt about what that leaves out, so they added the missing half: the claim expires sixty seconds after the last time its owner touches it, and the owner touches it between turbines, which is about every twenty five seconds. If the mover dies at turbine forty, another box picks the job up a minute later instead of nobody ever finishing it.

21:40 on the second Monday in April. Four boxes raced, box 2 won, and the copying started.

Turbine by turbine, twenty five seconds each, renewing the claim in between. At 22:24:35 box 2 finished the hundred and seventh turbine on its list and reached for the claim row.

The routing database was not there. Somebody had put a minor version upgrade of that Postgres on the April maintenance calendar for 22:24, eighty seconds of downtime on a machine nothing in the ingest path reads. Box 2's renewal was refused. It retried, backed off, retried again, and spent a minute doing it. Then it did the thing its author had thought about and chosen on purpose: rather than throw away forty four minutes of copying over a blip, it carried on.

At 22:25:10 the claim expired, sixty seconds after the renewal at 22:24:10, the last one that got an answer. Nothing happened, because nothing could reach the row to notice. Twenty five seconds after that, box 2 started the next turbine on its list.

At 22:25:40 Postgres came back. Box 3 polled, found an expired claim, took it, and began its own pass at the top of the list. The first turbine on the list was already gone from shard 1, so it copied nothing and moved on. So was the second. Nine seconds later it had walked through a hundred and seven turbines of somebody else's finished work and arrived at AP-0163, which is where box 2 was.

Both of them copied that turbine's year into shard 3. Box 3 got there first, reading rows shard 1 still had in memory from box 2's pass, counted what it had read out of shard 1 against what it had written into shard 3, got 15.77 million both ways, and at 22:26:19 deleted AP-0163 from shard 1, because that is what finishing a turbine means.

Two seconds later box 2 finished its own copy of the same turbine, slower than its usual twenty five seconds because two processes were now reading the same machine, and renewed its claim. The row belonged to box 3. So box 2 did the responsible thing. It was holding a half-finished turbine in shard 3 with no right to be there, so it tidied up: delete from readings where turbine_id = 'AP-0163' on shard 3.

A turbine sends 43,200 readings a day, so a year is 15.77 million rows and the five gigabytes lesson 013 costed. At 22:26:21 they were on neither machine.

The lease worked. Box 3 took over because box 2 had stopped renewing, which is the entire purpose of an expiry. Box 2 stopped the moment it learned it had lost, which is what a well behaved holder does. Nobody wrote a bug.

It surfaced two days later, which is faster than March and is nobody's credit. An engineer was still running the script they had written in March to compare all three shards, and a turbine whose history begins at twenty six minutes past ten on a Monday evening is loud in a way that twenty four thousand misfiled readings never was. Lesson 026 called Marlow's monitoring a museum of its own outages. Once in a while a museum piece earns its keep.

Which decisions cannot be made twice

Start before the election, because the question people skip is which decisions need one at all.

Galewatch's four ingest boxes have never needed a leader and never will. They take readings and write them down. If two of them handled the same reading you would waste a few microseconds of CPU and an insert, and lesson 019 would make even that harmless. Four boxes doing the same work is a cost. Four boxes doing contradictory work is an outage.

That is the test, and it is cheap to apply: a decision needs one owner when its second copy contradicts the first rather than repeating it. Picking that owner out of a group of machines all willing to do the job is leader election, and it is worth knowing how few things need it.

Deleting a turbine from a shard needs it. Readings for AP-0163 keep arriving all night, so a second delete from readings where turbine_id = 'AP-0163' takes rows the first one never saw, which is lesson 019's rule exactly: a statement is safe to repeat only when its own effect falsifies its condition, and that where clause never stops being true. Changing the divisor from 2 to 3 needs it, and lesson 031 is a whole lesson on what two answers to that question cost. Promoting a replica needs it.

Emailing a publisher their monthly CSV does not, which is why lesson 019 spent a section on that job rather than reaching for a lock. A great deal of what gets a lease wanted lesson 019 instead, and lesson 019 is cheaper and has no clock in it. Spend ten minutes on which one you want before you spend a week on the other.

You already have a leader election and it is a unique constraint

Marlow Books, the four person online bookshop in this course, has run two leader elections. One of them has never picked the wrong winner. The other was deleted twelve days after it shipped, having signed the shop's customers out for an hour and fifty minutes. Lesson 007 built the one that works, and it is the ten lines that lesson called the whole mechanism.

Every box wakes at 2 am on the first Sunday, every box tries to insert one row into job_runs keyed on the job name and the date, and exactly one insert comes back with a row. Lesson 007 named the shape: whoever writes the row first runs the job. No heartbeat, no expiry, nothing to renew, and no possibility of two payouts, because there is exactly one thing deciding and it is the database that the payout was going to need anyway.

Look at what that election does not have to know. It never asks whether the other box is alive. It asks who got here first, and a unique constraint answers that correctly on the first try, every time, forever. An election that asks who was first needs no clock. An election that asks who is alive needs one, and the clock is where all of this goes wrong.

The prices are real and lesson 007 printed them. You get at-most-once: if the winner dies halfway through, the row sits there, no other box will ever take the job, and nobody is paid, which is lesson 018's eleven days of silence. The leader is exactly as available as that one Postgres. And you elect per run rather than per role, which is perfect for a cron and useless for anything that has to stay in charge for an hour.

Lesson 007 asked its reader where the insert stops being enough, and the answer is the hook. A row in the routing database saying box 2 owns the move is read by nobody on the three machines where the damage happens. For Marlow's payout the arbiter and the victim are the same Postgres. For Galewatch's mover they sat on four different machines.

An election is only as safe as the thing that checks its answer.

A lease is a promise to stop

The mechanism the mover added has a name. A lease is a claim with an expiry on it: you hold the role until a stated time, and you extend it by asking again before that time arrives. Lesson 029 described the same object from the other side, a signed token whose contents are a decision frozen for the whole lifetime you chose. Lesson 007 named this one when it said that turning the payout lock into at-least-once means a finished_at and an expiry, and that a job a second box can retake has to be safe to run twice. Lesson 019 then showed you cannot escape it, only make the work behind it survivable.

Here is the part that is almost never built. A lease has two halves. The lock service enforces the first one: it will not hand the role to a second holder while the current grant is live, and it is good at that, because it is one place answering one question. The second half is that the old holder stops at the expiry.

Nothing enforces that at all.

The old holder's stop time is enforced by the old holder's own code reading the old holder's own clock, and a process that is paused, swapping, blocked on a socket or waiting out a stop-the-world collection has neither. A lease is a promise to stop, and the promise is kept by the machine that is losing it. Every system that has ever had two leaders had this bug, including the one in the hook, which had the clock and the code and still kept copying for seventy one seconds after its authority ran out.

Now price what the mover was asking to get through. Fifty five minutes of copying, one renewal between each of 133 turbines, so a hundred and thirty two times that night, box 2 had to get a round trip to a Postgres and back inside sixty seconds or stop being the mover while believing it still was. That is not a bet about the mover's code. It is a bet about the worst thing that will happen to a process, a network and a database in the next hour, taken a hundred and thirty two times, and the thing that eventually happened was on somebody's maintenance calendar.

The two obvious dials make it worse in opposite directions, which is the useful part. A ten second lease cannot contain a twenty five second copy at all, so the renewal has to move onto its own thread, and a renewal on its own thread keeps the lease alive while the work it is vouching for is wedged on a socket: the lock service hears a healthy heartbeat from a mover that has not copied a row in four minutes. Go the other way, ten minutes, and every renewal is comfortable and you have also decided that a mover which dies at turbine forty blocks the next one for ten minutes. When both ends of a dial are wrong, the dial is not where the fix lives.

The resource keeps the score

Fencing is the fix, and it is a word this course has leaned on without ever opening. Lesson 011 named the four things a failover needs, a death certificate, a promotion, a repoint and a fence, and said the database supplies one of them. Lesson 024 added that a promotion with no repoint and no fence is a second opinion. This is the fence.

Every time the lock service grants the role, it attaches a number, and that number only ever goes up. The mover carries it on every statement it sends. Each shard remembers the highest number it has already obeyed and refuses anything lower. The standard name for the number is a fencing token.

Which in Galewatch's case is one row on each shard and one statement, run inside the same transaction as the work:

create table mover_fence (token bigint not null);
insert into mover_fence values (0);

-- box 2 holds token 17, so every batch it sends begins with
update mover_fence set token = 17 where token <= 17;
-- zero rows means a newer holder has already been here:
-- roll the transaction back and stop

Said in words, because the audio skips that: the shard stores the highest token it has served, the mover's own token has to be at least that high for its work to commit, and a mover whose token is older gets zero rows and a rollback instead of a write.

Now run the hook through it. Box 3 is granted token 18 at 22:25:44 and its first statement on shard 3 moves the stored token from 17 to 18. At 22:26:21 box 2 arrives with the tidy-up delete, carrying 17, and shard 3 tells it no. Box 2 logs an error at 22:26:21 instead of removing a year of readings, and the worst thing that happened that night is a wasted minute on one turbine and a half-copied duplicate somebody clears up on Tuesday. Box 2's forty four minutes were never at risk either: box 3 found those turbines already gone from shard 1 and skipped them in nine seconds, which is the thing box 2's author had not worked out before deciding to carry on.

Be careful about why this works, because it looks exactly like the thing lesson 031 told you not to do. Lesson 031's stale ingest box computed its own answer from its own stale input, so its claim agreed with itself and a comparison against it matched. The fencing token is different in one way that decides everything: it is minted by the only party allowed to mint it, and a stale holder cannot invent a newer one. And notice what the shard is being asked. Not "are you the leader", which it has no way to know. A fence does not know who the leader is and does not need to. It knows the newest authority it has already served, which is a fact about its own memory, and that is enough to refuse everything older.

Lesson 031's shard constraint and this are the same move one radius apart. The owner stops taking the caller's word for the thing the caller is wrong about.

The honest limit is that fencing needs a resource that can keep score. Three Postgres machines can. An object store with a conditional write can. An email provider, a card processor and a bare file on a mounted volume cannot, and no amount of cleverness in the leader fixes that, which is where you fall back to lesson 019 and make the work safe to do twice, or accept at-most-once and say so out loud. Lesson 035 owns the sharp edges of the lock service itself.

Nobody at GitHub needed to be wrong

On 21 October 2018, GitHub was replacing failing optical network equipment and the work cut the link between their US East Coast network hub and their primary East Coast data centre for forty three seconds. Their automated failover promoted a database cluster on the West Coast. Writes the East Coast primary had already taken never made it west. GitHub chose the data over the uptime and ran degraded for twenty four hours and eleven minutes. Lesson 011 published that story and lesson 024 divided it: about two thousand to one, because recovery reconciles decisions rather than packets.

What today adds is that the promotion was not the mistake. The automation did what it had been configured to do, on evidence that was true: the primary had stopped answering. Nothing in the picking has to have gone wrong for the next twenty four hours to happen.

The usual way a group makes that call with no central referee is a majority, and there is one reason it works that needs no proof: two groups that each hold more than half of the same set of machines must share a member, and that member will not vote twice. A majority is how "am I allowed to act" becomes a question a machine can answer without talking to the machine it is replacing, which is the only version of the question that survives a partition. Lesson 033 is the machinery.

So the picking was defensible, and the writes still landing on the East Coast primary during those forty three seconds were landing on a machine nothing had taken out of service. A correct election plus an unfenced resource is two leaders.

Lesson 024 ended on the failure detector it could not repair: from where Marlow's watcher stood, a machine that will answer in a moment and a machine that is never going to answer again send exactly the same thing down the wire, which is nothing. That is still true and it is not getting fixed. You cannot make the detector right. You can make a wrong answer harmless, and the two things that do it are a majority, so that only one side is permitted to decide, and a fence, so that a loser who has not found out yet cannot act on it.

The window where nobody is in charge

Everything above is paid for with an interval during which the role is vacant, and the interval is not a tuning detail. It is the lease.

You cannot detect death faster than your expiry, because until the expiry passes, silence and slowness are the same observation. So the role is vacant from the expiry until the next candidate looks: on the April Monday that was 22:25:10 to 22:25:44, thirty four seconds, against a worst case of the full sixty plus however often box 3 polls.

What half a minute of vacancy costs is something you get to decide in advance.

During those thirty four seconds Galewatch ingested every reading it was sent. Twelve hundred turbines, the fleet lesson 031 counted, at one reading every two seconds is six hundred a second, and not one of them waited on anything, because the ingest boxes ask nobody where a reading goes. Lesson 031's formula needs no agreement. What stopped was the rebalance. Keep the leader off the request path, and a leaderless minute costs you a decision rather than a service.

Marlow is the other arrangement, and it is why the founder has no auto-promotion today. There the leader is the request path: box A is the thing every write goes to, so the vacancy is the shop being shut. Lesson 024 told that story. Fifty lines of Python, a thirty second constant, a ninety second network fault inside one rack, and an hour and fifty minutes of quoting prices it would not honour, after which the fifty lines were deleted. What Marlow has instead is a switch with a human on it: lesson 025 measured the rebuild at forty one minutes and put a promotion rehearsal in the calendar for the first Saturday of every month, and that one the founder actually does.

That is the right call for four people, and it is not where it should stay. Lesson 011's four things are still the list, and two of them are now today's lesson: a decision somebody is allowed to make, and a fence so the loser cannot write. Both need a third machine, because with two there is no majority either side can hold on its own. And a fence on a Postgres primary is real work: fencing it means taking away its ability to accept a write at all, which is its network address, its storage or its power, not a flag you set from the box you cannot reach. The monthly rehearsal promotes the replica and has never once had to stop the old primary, because in a rehearsal the old primary is cooperating. The reason small shops fail over by hand is that the manual version has a human majority of one, and that is not the embarrassing answer people think it is.

Recap

A decision needs one owner when its second copy contradicts the first rather than repeating it. Two boxes writing the same reading is waste. Two boxes deleting the same turbine is an incident. Ask which one you have before you add a lease, because lesson 019's idempotency is cheaper and has no clock in it.

An election that asks who was first needs no clock; an election that asks who is alive needs one. Marlow's payout lock is a unique constraint on the job name and the date, it is ten lines, and it has never been wrong. It buys at-most-once and it stops being enough the moment the leader's work lands somewhere that row does not guard.

An election is only as safe as the thing that checks its answer. Marlow's arbiter and victim were the same Postgres. Galewatch's claim row sat in the routing database and the damage happened on three shards that had never heard of it.

A lease is a promise to stop, and the promise is kept by the machine that is losing it. The lock service enforces the start. The stop is enforced by the old holder's own code and its own clock, which a paused process does not have. Fifty five minutes of copying under a sixty second lease is a bet taken a hundred and thirty two times, and shortening the lease only moves the renewal onto a thread that cannot tell whether the work is still moving.

The resource keeps the score. A fence is a number that only goes up, issued by the only thing allowed to issue it, carried on every statement and compared by the resource against the highest it has already obeyed. A fence does not know who the leader is and does not need to. It works where the resource can remember, and nowhere else.

You cannot make the detector right; you can make a wrong answer harmless. A majority decides who is permitted to act, a fence stops the one who is not, and GitHub's forty three seconds in October 2018 became a day because the promotion was not the error and nothing fenced the primary it replaced.

Check your understanding

  1. Your leader is thirty minutes into an hour of work when its connection to the lock service fails. Give the two things it can do next, say which you would ship, and then say what has to be true elsewhere for that choice to stop mattering.

  2. Marlow's founder wants automatic failover again, some months after deleting the fifty lines. Name what has to exist before you would agree, in the order you would build it, and say which item lesson 011 listed that the database still will not supply.

  3. A colleague argues for a five second lease so that failover is quick. Lesson 002 measured Galewatch's ingestion at a p99.9 of 2,100 milliseconds. Use that number to make their case, then make the case against, and name the measurement you do not have.

  4. Lesson 019 split Marlow's payout so that everything except the send is safe to run twice, and left the send itself in its third pile. You now want a second box able to take the job over mid-run. Say whether a lease helps the send at all, what you would fence and what you cannot, and what a publisher sees when you get it wrong.

  5. Take the April timeline and add fencing tokens. Write the log line box 2 produces at 22:26:21, say what state shard 3 and shard 1 are in at 22:26:30, and name the cleanup somebody still has to do on the Tuesday.

  6. A team runs three service instances and elects a leader to own an in-memory cache refresh that every instance could do independently. They have a lock service, a lease and a fencing token. Argue for deleting all three, then give the one fact about the refresh that would change your mind.

Next lesson

033 Consensus, Intuitively: Raft Without the Proofs. Today leaned on a majority twice without saying how a group of machines actually reaches one; next lesson opens that up, and electing a leader and keeping a log in step turn out to be the same mechanism done once.

Finished reading?

Marking a lesson done keeps your place on the course index. It is stored only in this browser.

Tip: use the ← and → keys to move between lessons.