Define consensus safety/liveness under explicit crash-fault and timing assumptions, then show why a minority partition must stop deciding rather than create two AtlasMart shard owners.
What Consensus Solves: Agreement Despite Failures Under Specific Assumptions
Consensus is not a magic availability feature. This lesson states the failure model, agreement/validity/liveness properties, majority intersection, and FLP-style liveness limit before any protocol mechanics.
Define consensus safety properties and liveness/termination assumptions without turning them into slogans.
Distinguish crash faults from Byzantine faults and explain what a simple majority crash-fault model does not tolerate.
Explain why a network partition can force a minority to stop deciding even though those nodes are still running.
Connect FLP-style asynchronous liveness limits to practical timeouts, partial synchrony, and leader-election behavior.
1. The problem: two healthy halves cannot both own the same shard
AtlasMart stores the authoritative assignment for payment shard
pay-17. Five control-plane nodes normally agree
that exactly one worker owns it. A network failure then
separates nodes A and B from nodes C, D, and E. Every process is
still alive. The dangerous temptation is to say, “each connected
group should keep operating and elect an owner.” That maximizes
local availability but can produce two simultaneously
valid-looking owners, which is a split brain.
Consensus is the problem of getting a group of participants to agree on a decision under a specified system model. A practical consensus protocol must state its assumptions: what kinds of faults occur, how messages behave, what persistent state survives crashes, and what timing conditions are needed for progress. Consensus does not remove partitions. Under a partition, a safe protocol may deliberately make one side unavailable for decisions.
2. Safety and liveness are different questions
Terminology varies across formalizations, but the core safety ideas are stable. Agreement means correct participants do not decide different values for the same consensus instance. Validity or integrity constrains what may be decided—for example, a decision must correspond to a value actually proposed, and a participant should not decide twice for one instance. Termination/liveness means correct participants eventually decide when the protocol's progress assumptions are satisfied.
| Property | Question it answers | Failure if violated |
|---|---|---|
| Agreement | Can two correct participants decide different values? | Two owners/terms/configurations become authoritative |
| Validity / integrity | Can the protocol invent or decide an invalid value? | A non-proposed or malformed decision becomes authoritative |
| Termination / liveness | Will a decision eventually happen under the stated timing/fault assumptions? | The system remains safe but unavailable/stuck |
| Durability | Does a decided value survive the required crash/restart model? | A node forgets promises/log state and can break protocol assumptions |
Safety is generally required even during severe delay. Liveness is conditional: if messages can be delayed forever and failures cannot be distinguished from delay, no deterministic asynchronous protocol can guarantee progress in all executions with even one crash failure—the classic Fischer-Lynch-Paterson (FLP) result. Practical systems escape the impossibility model by relying on additional timing assumptions, failure detectors, randomized mechanisms, or eventual periods of synchrony.
3. Crash faults are not Byzantine faults
A crash fault model assumes a failed process stops taking useful protocol steps (or becomes unreachable) but does not intentionally fabricate contradictory authenticated messages. A Byzantine fault allows arbitrary behavior: lying, equivocating, sending different claims to different peers, or being maliciously compromised. Raft and classic Paxos are normally taught as crash-fault-tolerant protocols; they are not Byzantine-fault-tolerant merely because they use quorums.
This difference matters for security. Mutual TLS, authentication, least privilege, and protection of persistent consensus state are not substitutes for a Byzantine consensus protocol, but they reduce the chance that an attacker can impersonate or alter members. Threat model and failure model must be documented separately.
4. AtlasMart lab: make the partition tradeoff visible
Python 3.13+ standard library only. The generated lab was verified with Python 3.13.5. No database server, Docker, cloud account, paid feature, credential, firewall change, clock manipulation, process killing, or destructive failure injection is required. All failures and partitions are deterministic in-memory simulations.
The deliberately broken rule lets any connected component appoint an owner. The corrected rule requires a majority of the fixed five-member crash-fault consensus group.
nodes = ["A", "B", "C", "D", "E"]
majority = len(nodes) // 2 + 1
print("CLUSTER")
print("nodes:", nodes, "majority:", majority)
# Network partition: A,B can talk to each other; C,D,E can talk to each other.
components = [{"A", "B"}, {"C", "D", "E"}]
print("partition components:", [sorted(c) for c in components])
print("\nBROKEN RULE: any connected component may appoint an owner")
broken_owners = []
for component in components:
candidate = sorted(component)[0]
broken_owners.append(candidate)
print(sorted(component), "appoints", candidate)
print("two active owners:", len(broken_owners) == 2, broken_owners)
print("\nMAJORITY RULE FOR ONE CRASH-FAULT CONSENSUS GROUP")
decisions = []
for component in components:
if len(component) >= majority:
decision = "owner=" + sorted(component)[0]
decisions.append(decision)
print(sorted(component), "can form quorum and decide", decision)
else:
print(sorted(component), "cannot decide; availability is sacrificed")
print("number of decisions:", len(decisions))
print("agreement in this run:", len(set(decisions)) <= 1)
print("\nFAILURE-MODEL REMINDER")
print("crash fault: node stops/responds late; protocol assumes it does not lie arbitrarily")
print("Byzantine behavior: arbitrary/malicious messages; this majority model does NOT tolerate it")
print("evidence proves this deterministic scenario only, not liveness under an unbounded asynchronous delay")
The broken rule produces two active owners. Under the majority rule, the two-node component cannot decide while the three-node component can. This demonstrates safety/availability behavior for this fixed deterministic partition only. It does not prove a full consensus protocol, tolerate Byzantine members, or prove liveness under unbounded message delay.
5. Majority is necessary context, not magic
For a fixed five-member group, any two majorities of size three intersect. That intersection is a useful safety ingredient because later decisions can be forced to encounter state from earlier quorums. But quorum intersection alone is not a complete protocol. Raft adds terms, voting restrictions, log matching, and commit rules. Paxos adds ballots, promises, and a rule for carrying forward previously accepted values. Reconfiguration also needs a safe membership protocol; simply changing the denominator from five to four or six can invalidate quorum reasoning if old and new configurations overlap incorrectly.
6. Deliberately wrong approach: “a timeout proves the node is dead”
A timeout is evidence of missing progress, not proof of death. If A is paused by a long garbage collection, saturated CPU, or network queue, another node may win a legitimate later term while A eventually resumes. The system remains safe only if stale authority is rejected through terms/epochs, leases with explicit assumptions, and often downstream fencing. Treating a timeout as a certificate of death is how stale leaders continue writing.
7. Production judgment and bridge
Use consensus for values that must have one authoritative sequence—membership configuration, shard ownership, lock/fencing generations, schema epochs, or replicated state-machine commands—when the latency and majority-availability costs are justified. Place members across failure domains deliberately; a five-node cluster spread badly across one rack is not five independent failures. Monitor election frequency, leader changes, quorum health, proposal/commit latency, persistent-log fsync latency, snapshot/compaction pressure, clock symptoms for lease-based features, and rejected stale terms.
The next lesson makes these abstractions concrete with Raft: terms, elections, replicated logs, commit index, conflict repair, and the reason an entry merely present on a leader is not yet committed.
Check your understanding
- What does agreement require?
- Why can a two-node minority of a five-node group be unavailable even when both nodes are healthy?
- What does FLP say at a high level?
- Does a majority crash-fault protocol tolerate a malicious member that equivocates?
- Why is quorum intersection not a complete consensus algorithm?
Review the answers
1. Correct participants must not decide different values for the same consensus instance.
2. Allowing it to decide independently could permit a disjoint decision and violate safety; it cannot form the configured majority.
3. In a fully asynchronous model with even one possible crash, a deterministic consensus protocol cannot guarantee termination for every execution, although safety can still be preserved.
4. No. Byzantine behavior is a different failure model requiring different mechanisms/assumptions.
5. The protocol must also constrain voting/proposals/log history so intersecting members carry forward the information needed for safety.
References
Foundational claims use primary research where practical. Product documentation is used only as a current implementation example and is not required for the mandatory labs.
- Fischer, Lynch & Paterson — Impossibility of Distributed Consensus with One Faulty Process — Primary liveness/impossibility result for deterministic consensus in a fully asynchronous model with one crash failure.
- Ongaro & Ousterhout — In Search of an Understandable Consensus Algorithm — Primary Raft reference; also states the crash-fault replicated-log problem and majority-based design.
- Lamport — Paxos Made Simple — Primary concise Paxos explanation and safety reasoning.
- Herlihy & Wing — Linearizability — Primary correctness condition used later when relating consensus decisions to real-time operation order.