Leaders, Followers and Failover
Introductions, exercises and summaries stay visible.
34.0 What this chapter gives you#
- A leader is not simply the machine with the newest-looking screen. It is a participant currently authorized by a defined protocol to coordinate particular work.
- You will trace a role change through failure suspicion, selection, fencing, client routing and restoration of redundancy. You will also learn why a missing reply leaves an operation’s outcome uncertain.
- The examples use fictional nodes A, B and C and a small authority model. They do not install failover software, promote a real PostgreSQL server or establish production high availability.
34.1 Write authority#
34.1.1 PLAIN — in simple words#
- A single-leader design gives one participant responsibility for coordinating writes to a defined set of data. Other participants may copy those writes or serve eligible reads.
- This is a rule about authority, not a statement that one computer is permanently special. A safe design can transfer that role without allowing incompatible authorities to operate at the same time.
- The scope matters. Different partitions may have different leaders, and a machine coordinating database writes may have no authority over a payment service or a delivery system.
34.1.2 PLAIN — a picture in your head#
- Mira appoints one counter supervisor to approve reservations from a shared stock allocation. Other staff forward requests to that supervisor rather than independently promising the same last item.
- If the supervisor changes, the old approval badge must stop working. Giving a new person a badge while the old badge remains valid creates two competing authorities.
- Where the comparison breaks: a database protocol cannot rely on everyone noticing a human announcement. Delayed messages and paused processes can continue carrying obsolete authority.
34.1.3 PLAIN — a worked example#
- In a separate teaching system, node A holds authority generation 7
for resource
RESERVE-BR-A. Its writes include that resource identity and generation. - Node B may be a follower for the same resource while leading an unrelated reporting task. Calling B “a follower” without naming the resource obscures its actual responsibilities.
- A role transition to generation 8 must define when generation 7 ceases to authorize changes and how the affected resource enforces that decision.
- Merely changing a dashboard label from A to B does not alter what an old connection to A can still do.
34.1.4 PLAIN — what is really happening inside#
- Authority may come from a consensus group, a carefully managed failover controller or another explicit coordination mechanism. The storage and request paths must cooperate with it.
- Clients also need a way to find the current service endpoint and recover from stale routing. Connection pools may retain old sessions after a network address changes.
- The design must connect selection, enforcement and discovery. Implementing only one of these leaves gaps between who was chosen and who can actually modify data.
34.1.5 TECHNICAL — the engineer’s version#
- Separate control-plane authority from data-plane enforcement. The former chooses a valid owner or epoch; the latter rejects operations that do not satisfy the current authorization condition.
- PostgreSQL 17 provides standby promotion mechanisms but does not itself supply the complete external failure-detection and notification system described in its failover documentation. [S150]
- A single-leader label is not a complete safety proof. Review the failure model, durable state, resource boundaries and stale-client behaviour of the actual implementation.
34.1.6 WORDS — remember these#
Leader: the current coordinator for specified work — a role granted under a protocol with a defined authority scope. Control plane: the machinery choosing and configuring behaviour — coordination mechanisms such as role assignment and membership management. Data plane: where requests produce actual effects — the operational path that must enforce the control plane’s decisions.
34.2 Failure suspicion#
34.2.1 PLAIN — in simple words#
- When a machine stops replying, it may be dead, slow, disconnected or unable to reach the observer. The observer has evidence of missing communication, not complete knowledge of the machine’s state.
- Waiting forever prevents useful recovery. Acting too quickly can replace a healthy but temporarily unreachable leader. A failover design must handle both risks explicitly.
- A timeout is therefore a trigger for a protocol, not proof that the previous authority has ceased to exist.
34.2.2 PLAIN — a picture in your head#
- Dev calls Mira and hears nothing. Her phone may be off, the mobile network may be down, or she may be busy with a customer.
- Dev cannot safely infer that every approval she previously held has disappeared. A separate rule is needed to transfer responsibility and invalidate obsolete approvals.
- Where the comparison breaks: systems use timers, persistent terms and communication protocols rather than human judgement. The analogy explains uncertainty, not a recommended failure detector.
34.2.3 PLAIN — a worked example#
- Suppose a fictional controller expects a heartbeat every second and suspects failure after three missed intervals. At time 3 seconds it may initiate investigation or an election.
- In one schedule, A crashed at time 0. In another, A is still running but its messages to the controller are delayed for 5 seconds. The observer sees the same missing heartbeats during the first 3 seconds.
- Promoting B at time 3 without fencing A can allow both machines to accept incompatible writes between times 3 and 5, or longer if the network split persists.
- These timers are invented for reasoning. They are not suggested production settings or a measured failure-detection guarantee.
34.2.4 PLAIN — what is really happening inside#
- Failure detectors trade responsiveness against false suspicion under the timing conditions they assume. A process pause can resemble a network outage from another participant’s perspective.
- A witness or third participant may help a protocol decide which side can proceed, but only if its role and failure assumptions are correctly integrated. An arbitrary extra ping target is not a proof of safety.
- Good incident records preserve observations from multiple paths, the timer policy and the authority transition. They do not retroactively turn an initial suspicion into a fact that was known at the time.
34.2.5 TECHNICAL — the engineer’s version#
- Distinguish crash failure from communication failure and timing uncertainty. Raft uses election timeouts to initiate a new term, while safety depends on voting and log rules rather than the timeout proving that a leader died. [S154]
- PostgreSQL’s failover guidance explicitly discusses heartbeat mechanisms, possible witnesses and the need to prevent an old primary from continuing as primary. [S150]
- Record the actual detection-to-service-restoration interval, including selection, promotion, routing and readiness. A three-second suspicion threshold does not imply a three-second recovery time.
34.2.6 WORDS — remember these#
Heartbeat: a repeated liveness signal — periodic communication used to detect absence or delay under a configured policy. Failure suspicion: evidence that a participant may be unavailable — an observation requiring protocol handling rather than certain knowledge of a crash. False suspicion: a healthy participant is treated as failed — a classification possible under delays, partitions or process pauses.
34.3 Election and promotion#
34.3.1 PLAIN — in simple words#
- Selection chooses a participant to take responsibility. Promotion changes a replica’s operational role. These steps must agree with the data history and the authority mechanism.
- Choosing the first machine that answers a ping is not enough. The candidate must have suitable data, compatible configuration and a legitimate path to ownership.
- A system may intentionally stop rather than promote a candidate that would violate its acknowledged-data or authority requirements.
34.3.2 PLAIN — a picture in your head#
- Two assistants have different prefixes of the order notebook. One is available immediately but is missing the latest approved reservations; another has the necessary history but needs a moment to finish applying it.
- The shop must decide what evidence a replacement needs, not merely who can reach the counter fastest.
- Where the comparison breaks: consensus election rules are mathematical protocol conditions. Human preference for the most confident assistant is not an equivalent selection rule.
34.3.3 PLAIN — a worked example#
- Suppose A acknowledged a synthetic write at sequence 41, B has durable history through 41 and C has only through 40. A is now unreachable.
- Promoting C can lose the acknowledged change unless the system first obtains the missing history or its documented policy accepts that loss. The fact that C is running does not fill the gap.
- If B was the required synchronous destination for that acknowledgement, the failover procedure must preserve the relevant durability assumptions and select a valid survivor. An unrelated asynchronous copy is not automatically equivalent.
- In a Raft-style protocol, the election’s log-freshness rule is specific: compare the last entry’s term first, then its index when terms are equal. “Largest row count wins” is not the rule. [S154]
34.3.4 PLAIN — what is really happening inside#
- A candidate may need recovery replay before serving traffic. A process accepting TCP connections does not establish application readiness or complete recovery.
- The transition must also handle routing, credentials, replication slots, monitoring and downstream consumers. These are practical components of service continuity, not optional cosmetic tasks.
- Once a follower becomes primary, the system may temporarily have fewer healthy copies. Restoring redundancy is part of finishing the incident, even after requests start succeeding again.
34.3.5 TECHNICAL — the engineer’s version#
- Failover safety depends on the acknowledged commit set, eligible candidates and the mechanism preventing competing primaries. PostgreSQL streaming replication and Raft are different mechanisms; do not attribute Raft’s election rules to a plain PostgreSQL primary/standby pair. [S147] [S150] [S154]
- Readiness checks should cover recovery completion, intended role, required schema, protected write access and the ability to perform a bounded representative operation.
- The procedure must record any accepted data-loss window. Fast role switching with undisclosed loss is not equivalent to a verified no-loss transition.
34.3.6 WORDS — remember these#
Election: selecting a coordinator under protocol rules — an authority decision with explicit membership and eligibility conditions. Candidate eligibility: the conditions a replacement must satisfy — requirements covering authority, history and readiness rather than mere reachability. Degraded redundancy: fewer healthy copies than intended — an operational state that may persist after the service resumes.
34.4 Fencing stale writers#
34.4.1 PLAIN — in simple words#
- A paused old leader can wake up believing it still owns the job. Fencing prevents its obsolete authority from producing effects.
- A common pattern attaches a changing generation to authority. The protected resource checks that generation and rejects obsolete operations.
- The check must occur where the effect is accepted. A generation number written only in a log message provides evidence, not enforcement.
34.4.2 PLAIN — a picture in your head#
- The stock room changes its approval register from badge generation 7 to generation 8. The door attendant checks the register for every withdrawal.
- An old supervisor returning with badge 7 is refused even if the badge once worked. Merely asking the supervisor to remember to stop would be weaker.
- Where the comparison breaks: digital enforcement must be atomic with the protected change and must validate the actual authority source. A caller must not gain ownership merely by inventing a larger number.
34.4.3 PLAIN — a worked example#
- Our toy protected resource stores current generation 7 and value 5. A valid generation-7 request changes the value to 4.
- A separate trusted authority transition atomically advances the resource to generation 8 before new-owner traffic is admitted. A delayed generation-7 write proposing value 99 is rejected; the value remains 4.
- A valid generation-8 write changes it to 3. The model checks generation equality, not “accept any number larger than the last number.” Untrusted callers cannot advance authority through an ordinary write.
- This is a deterministic teaching model. Authentication, consensus-based generation allocation, durable fencing across hosts and actual device isolation are outside it.
34.4.4 PLAIN — what is really happening inside#
- If enforcement relies on remembering the highest generation ever seen, the exact acquisition and first-write protocol matters. A stale write arriving before the newer generation is registered can still be accepted unless the larger design prevents it.
- Physical or infrastructure fencing can instead remove the old node’s ability to affect the resource. That requires supported, authorized operational controls and verification, not an assumption that a failed ping means isolation.
- Every relevant path must be covered. Fencing database writes does not automatically fence an old worker’s email, object-storage update or command to another service.
34.4.5 TECHNICAL — the engineer’s version#
- Chubby’s original design describes sequencers containing lock generation information and recipient-side validation to reject obsolete lock-protected requests. The recipient’s participation is essential. [S163]
- PostgreSQL’s failover documentation requires a mechanism preventing the old primary from continuing as primary after replacement. The exact fencing implementation is deployment-specific. [S150]
- A correct design specifies epoch allocation, durable acceptance state, replay rules, atomic effect checks and all downstream resources. Do not describe a standalone monotonic integer as a complete distributed lock.
34.4.6 WORDS — remember these#
Fencing: preventing obsolete authority from acting — enforcement at a protected resource or infrastructure boundary that rejects stale owners. Epoch: a generation of authority or configuration — a value distinguishing successive valid ownership periods under a defined protocol. Sequencer: evidence describing a lock generation — a token that a recipient validates before accepting a protected operation.
34.5 Lost acknowledgements#
34.5.1 PLAIN — in simple words#
- A request can succeed even when the caller never receives the reply. The network can fail after the change is committed but before the success message arrives.
- The caller must therefore distinguish “known failed” from “outcome unknown.” Treating every timeout as failure can repeat an already completed operation.
- A stable request identity and stored outcome help the caller reconcile what happened across a retry or leader change.
34.5.2 PLAIN — a picture in your head#
- Mira approves a reservation and files the record, but the phone call drops before Dev hears the answer.
- Dev should ask about that same request number. Sending a brand-new request for the same goods can reserve them a second time.
- Where the comparison breaks: duplicate detection must be part of the durable operation, not an informal memory of the last caller. Its retention and authority boundaries must survive the relevant failures.
34.5.3 PLAIN — a worked example#
- Request
IDEM-Areserves three units from five, leaving two. The request record, reservation and notification intent commit together, as in Chapter 26. - The success reply is lost. After routing changes, the caller retries
IDEM-Awith the same validated semantic payload. A survivor containing that committed history returns the stored outcome without reserving again. - If an incorrectly selected replacement lacks the request record and stock change, the same retry can be treated as new. Local idempotency does not repair a failover procedure that discarded its evidence.
- Reusing
IDEM-Awith quantity four remains a conflict, not a convenient way to edit the original request.
34.5.4 PLAIN — what is really happening inside#
- Request identity, business mutation and durable outcome must share an appropriate transaction boundary. Otherwise one can survive without the others.
- Recovery also has to preserve enough request history for the declared retry horizon. Deleting deduplication evidence changes which old requests can safely be recognized.
- A client may need to query an authoritative operation-status endpoint rather than blindly repeat an external action. The endpoint itself must enforce tenant and permission boundaries.
34.5.5 TECHNICAL — the engineer’s version#
- The Raft paper explicitly separates log replication from client duplicate suppression: a retried command may execute again without request identifiers and retained results in the state machine. [S154]
- Idempotency contracts must bind the request key to caller scope and semantic intent. A timeout after a possible commit requires reconciliation under that contract, not automatic assumption of rollback. [S126]
- End-to-end guarantees include failover history selection and external-effect handling. A database commit record does not establish that a remote notification was delivered exactly once.
34.5.6 WORDS — remember these#
Unknown outcome: the caller lacks enough evidence to classify completion — a state distinct from confirmed success or confirmed rollback. Operation reconciliation: checking a prior request’s authoritative result — resolving uncertainty through stable identity and retained outcome evidence. Retry horizon: how long old requests remain recognizable — the interval over which a system retains the evidence needed by its idempotency contract.
34.6 Safe return of a former leader#
34.6.1 PLAIN — in simple words#
- Bringing an old primary back is not the same as turning on another read-only copy. Its data may have diverged from the accepted history after the replacement took over.
- Do not reconnect it as writable merely because its files are intact. First establish its role, compare histories and choose a supported resynchronization path.
- Restoring redundancy is complete only when the returning participant follows the intended authority and can be observed doing so.
34.6.2 PLAIN — a picture in your head#
- An old supervisor returns with a notebook containing some private entries made while disconnected. The new supervisor has a different continuation of the shared notebook.
- Stapling the two endings together arbitrarily can create contradictory reservations. The shop must determine the accepted history and preserve disputed evidence before rebuilding the old copy.
- Where the comparison breaks: database histories include recovery positions and timelines. They cannot generally be merged by comparing visible rows or choosing the newest file timestamp.
34.6.3 PLAIN — a worked example#
- A and B share entries through 40. While isolated, A has an unaccepted local entry 41A; B becomes valid authority and commits a different entry 41B.
- The matching sequence number does not make the entries interchangeable. The old node must not serve its 41A continuation as though it were the current history.
- A supported rebuild or rewind procedure may return A as a follower of B, subject to its prerequisites and preservation requirements. A later integrity check alone cannot decide which business history was authorized.
- The example is a history diagram, not permission to discard real data or run a destructive repair command.
34.6.4 PLAIN — what is really happening inside#
- Preserve the old state when investigation or reconciliation requires it. Establish the accepted authority and acknowledge any unresolved client outcomes before resynchronization.
- Use the database’s supported procedure for the exact version and configuration. Recovery must finish, and the node must be configured as the intended follower before it receives ordinary traffic.
- Finally verify replay progress, role, read-routing behaviour, monitoring and the number of healthy copies. Record what was tested rather than declaring the cluster healthy from one successful connection.
34.6.5 TECHNICAL — the engineer’s version#
- PostgreSQL
pg_rewindcan help resynchronize divergent clusters with a common history, but it has documented prerequisites involving WAL information and configuration. It is not a general merge tool. [S151] - Recovery replay and standby configuration remain necessary after rewind. The documentation warns that failures during the operation can leave the target unusable and that incorrect restart configuration can allow renewed divergence. [S151]
- A failover runbook should include a safe-return phase, evidence preservation and a restoration-of-redundancy check. This book supplies reasoning templates, not a verified procedure for an unseen production cluster.
34.6.6 WORDS — remember these#
Divergent history: copies contain incompatible continuations — a situation requiring authority and recovery analysis rather than simple row copying. Reinitialization: rebuilding a participant from valid source state — a controlled procedure restoring a compatible replica and its role. Failover runbook: the documented role-transition procedure — steps, prerequisites, checks and failure dispositions for transferring service authority.
34.97 Practice and worked answers#
- Three heartbeats are missing. What is established? Communication has not met the observation policy. A crash is possible but not proven.
- Why is a dashboard role change insufficient? Existing processes, connections and downstream resources may still accept the old authority unless enforcement changes.
- A has committed 41, B contains 41 and C contains 40. Is C automatically a safe replacement? No. The acknowledged-history and promotion requirements determine eligibility.
- A write carries generation 999. Should a resource accept it because 999 is large? No. Generation allocation and authorization belong to the authority protocol, not an untrusted request field.
- A caller times out after a possible commit. What should it preserve? The stable request identity, semantic payload and evidence needed to query or retry the same operation safely.
- Can local idempotency survive a failover that lost its request table? Not automatically. The selected survivor must retain the relevant committed request and business state.
- Two copies have different entries numbered 41. May they be merged by number? No. Positions are meaningful within histories and authority rules, not as universally interchangeable labels.
- When is the incident finished? After the required service and redundancy checks, unresolved outcomes are recorded and the tested scope is explicit—not merely when a replacement accepts a connection.
34.98 Common wrong ideas#
- Wrong: timeout means dead. Right: it means the observer lacks a timely reply.
- Wrong: the fastest responder should always become leader. Right: eligibility includes history and legitimate authority.
- Wrong: promotion automatically fences the old primary. Right: stale authority needs an enforced prevention mechanism.
- Wrong: a larger token is self-authenticating. Right: the token must originate from and be checked against valid authority.
- Wrong: no reply means no write. Right: a committed operation can lose its acknowledgement.
- Wrong: consensus automatically removes repeated business requests. Right: client identity and duplicate handling are additional state-machine responsibilities.
- Wrong: an old primary can safely restart with its previous configuration. Right: role, history and recovery must be reconciled first.
- Wrong: failover ends when traffic returns. Right: redundancy and safe return also require verification.
34.99 Chapter summary in 20 lines#
- Leadership is scoped authority, not a permanent machine identity.
- Selection, enforcement and discovery are separate responsibilities.
- A timeout creates suspicion rather than certainty about a crash.
- Process pauses and network partitions can resemble failure.
- Failure detection initiates a protocol instead of proving old authority vanished.
- Replacement candidates need suitable history and configuration.
- Acknowledged-data guarantees constrain failover selection.
- PostgreSQL standby promotion is not itself a Raft election.
- Fencing prevents obsolete owners from producing effects.
- Enforcement must occur at every relevant protected resource.
- Authority generations require legitimate allocation and validation.
- A generation number alone is not a distributed lock.
- Missing acknowledgements leave some outcomes unknown.
- Stable request identities support safe reconciliation and replay.
- Idempotency evidence must survive the promised failover and retry horizon.
- External effects need their own outcome and duplicate policy.
- A former leader may hold a divergent continuation.
- Supported resynchronization is not arbitrary history merging.
- Recovery and standby configuration must precede return to service.
- Complete failover evidence includes restored redundancy and remaining limits.