In distributed systems, engineers talk about crash-stop, crash-recovery, omission, and Byzantine failure models. What does each model assume about how a faulty node can misbehave, and why does the choice of model matter when designing a fault-tolerant system?
answer
- crash-stop = dies forever
- crash-recovery = dies then reboots, loses non-persisted state
- omission = network drops messages
- Byzantine = arbitrary/malicious behavior
- 2f+1 vs 3f+1 quorum math
basics
~20 sA failure model is an assumption about how something can break. Crash-stop: a node dies and never comes back. Crash-recovery: it dies but can restart. Omission: messages get lost. Byzantine: a node can lie or act maliciously. Stronger assumptions need more expensive protection.
solid answer
~50 sThese are increasingly permissive assumptions about what a faulty component can do. Crash-stop assumes a node halts forever and sends nothing more - simplest to reason about. Crash-recovery allows a node to halt and later restart, possibly with amnesia about in-flight state unless it persisted it - this is what most real systems (databases, Raft/Paxos deployments) actually assume, which is why nodes must fsync state before acknowledging. Omission failures mean messages can be dropped in transit, effectively 'the network can lose packets.' Byzantine failures are the most general: a node can send arbitrary, contradictory, or malicious messages to different peers. Real systems pick a model based on threat surface and cost: internal replicated databases assume crash-recovery because operators control the hardware; blockchains and multi-organization systems assume Byzantine because no single party is trusted. The model chosen determines the fault-tolerance math (n>2f for crash faults vs n>3f for Byzantine).
go deeper
Can name the four models and give a one-line description of each without confusing which is more permissive than which.
Can explain why crash-recovery, not crash-stop, is the realistic model for restart-prone production systems, and connects durable state (e.g., fsync before ack) to correctness under that model.
Can derive or explain the 2f+1 vs 3f+1 quorum math and justify why a given system (internal Raft cluster vs a blockchain) picked its failure model based on trust boundaries and cost.
Can reason about failure-model composition in real architectures - e.g., assuming crash-recovery for owned replicas but treating an untrusted third-party integration as Byzantine - and weigh the operational cost of over-provisioning beyond the actual threat surface.
## What a failure model is A **failure model** is a formal statement of the ways a component in a distributed system is permitted to deviate from correct behavior, and every fault-tolerance algorithm is designed and proven correct only under an assumed model - pick the wrong one and the protection you built provides no actual guarantee. The four models line up on a spectrum from most restrictive (easiest to protect against) to least restrictive (hardest, costliest to protect against). ## Crash-stop **Crash-stop** (also called fail-stop) is the simplest and most optimistic model: a faulty process halts at some point and from then on sends and receives nothing, forever. Other processes may or may not detect the halt (that's the job of a failure detector, covered elsewhere), but they never receive a message from a 'dead' node again. This model underlies textbook consensus proofs and is easy to reason about, but it's unrealistic on its own - real processes restart after: - crashes - OOM kills - container restarts - power blips ## Crash-recovery **Crash-recovery** relaxes that: a process can crash (stop sending/responding) and later resume, but when it resumes it may have lost all memory that wasn't durably persisted before the crash. This is the model that actually matches production infrastructure - a Raft or Paxos node that crashes and restarts must reload its persisted log and term/vote state from disk before it can safely rejoin, precisely because crash-recovery permits amnesia of anything not fsynced. Getting this model wrong is a classic real bug: a node that recovers and re-votes for a leader in a term it already voted in (because the vote wasn't flushed to disk) can cause **split-brain**, which is why protocols like Raft mandate persisting `currentTerm` and `votedFor` before responding to RPCs. ## Omission **Omission** failures generalize the picture to the network layer: a process is correct, but messages it sends or receives can be silently dropped (send-omission, receive-omission, or both). This is essentially 'the network is asynchronous and lossy' and is the default assumption for anything running across machines. Omission is, behaviorally, a subset of what crash-recovery can look like, because a crashed node's silence is indistinguishable, over an asynchronous network, from a node whose messages are all being dropped - this indistinguishability is exactly what makes failure detection fundamentally hard. ## Byzantine **Byzantine** failure is the most permissive and expensive model: a faulty node can do anything at all - - send different, contradictory values to different peers - forge messages - stay silent selectively - actively try to subvert the protocol This isn't hypothetical malice only; it also captures: - bit-flip memory corruption - buggy code computing wrong values - a compromised machine under attacker control Byzantine fault tolerance (BFT) underlies PBFT, Tendermint, and blockchain consensus, because in those settings no single party is trusted and nodes are run by mutually distrusting operators. ## The redundancy math the model drives The model chosen directly drives the redundancy math for masking failures via quorums. | Fault model | Replicas required | Why | |---|---|---| | Crash faults | tolerating f failures among n replicas typically requires n ≥ 2f+1 (majority quorums, as in Raft/Paxos) | a non-faulty majority is enough to out-vote silence | | Byzantine faults | tolerating f faulty nodes requires n ≥ 3f+1 | you need enough correct replicas to both outvote the liars and distinguish a genuinely correct minority report from a fabricated one sent by faulty nodes | This is a common interview trap: candidates who quote 'n > 2f' for a BFT system have conflated the two models. ## What production actually assumes In practice, most cloud infrastructure (etcd, ZooKeeper, internal replica sets, cloud-native SQL like CockroachDB/Spanner) is engineered under crash-recovery plus omission, because operators trust their own hardware and the primary risk is hardware failure, process crashes, and restarts - not malicious replicas. Byzantine tolerance is reserved for cross-organizational or adversarial settings, because it costs roughly 50% more replicas and heavier per-message overhead for the same fault budget, which is wasted expense if the actual threat model is 'servers occasionally crash,' not 'servers occasionally lie.'
- Why does Raft require persisting currentTerm and votedFor to disk before replying to a RequestVote RPC?Because Raft assumes the crash-recovery model, where a node can crash and restart with amnesia for anything not durably stored. If votedFor weren't flushed before the reply, a restarted node could vote again in the same term it already voted in, letting two leaders be elected in the same term and causing split-brain writes.
- Why does Byzantine fault tolerance need n ≥ 3f+1 replicas instead of the 2f+1 used for crash faults?With crash faults, a silent node simply doesn't vote, so a majority of 2f+1 is enough to make progress and know any decision reflects the honest majority. With Byzantine faults, faulty nodes can actively vote for wrong or contradictory values, so you need enough correct nodes left over after subtracting the maximum faulty set to still outnumber a second, disjoint set of faulty nodes lying differently to different observers - hence 3f+1.
- Is TCP's guaranteed, ordered delivery enough to rule out omission failures at the application layer?No. TCP guarantees ordering and retransmission within an established connection, but connections themselves can drop, time out, or never establish during a partition, and the application still experiences a message never arriving. Omission at the distributed-algorithm level is about the end-to-end guarantee an application protocol can rely on, not the transport's internal retry behavior.
Think of a courier service: crash-stop is a courier who quits and never comes back; crash-recovery is one who goes on unannounced leave and returns having forgotten anything not written in their notebook; omission is a courier who sometimes loses packages in transit; Byzantine is a courier who might deliver forged or tampered packages on purpose.
saying these in an interview costs you the question
- Says Byzantine tolerance also needs n>2f, conflating it with crash fault tolerance
- Thinks 'omission' and 'crash-stop' are the same thing
- Believes crash-recovery nodes always retain full state after restart
- Can't explain why the model choice changes the replica count needed
- Assumes Byzantine faults only mean malicious attackers, missing memory corruption or bugs as a cause