skip to content

questions

5

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?

level: juniorimportance: must knowfreq 65%

answer

  1. crash-stop = dies forever
  2. crash-recovery = dies then reboots, loses non-persisted state
  3. omission = network drops messages
  4. Byzantine = arbitrary/malicious behavior
  5. 2f+1 vs 3f+1 quorum math

basics

~20 s

A 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 s

These 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

for a junior

Can name the four models and give a one-line description of each without confusing which is more permissive than which.

for a middle

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.

for a senior

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.

for a principal

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

context

open as a page

A replicated key-value store has N=5 replicas. A team configures writes to require acknowledgment from W=3 replicas and reads to query R=3 replicas. Why does this W+R>N configuration let the system mask node failures while still returning consistent reads, and what happens if the team instead sets W=2, R=2?

level: middleimportance: must knowfreq 62%

basics

~20 s

With W+R>N, any read set and any write set must share at least one replica, so a read always sees the latest write even if some nodes are down or slow. If W+R is not greater than N, like W=2,R=2 on N=5, reads and writes might miss each other entirely and return stale data.

open as a page

A team is deciding between synchronous and asynchronous replication for a database that must tolerate a primary node failure. What availability and durability guarantees does each give, and what specific failure mode does asynchronous replication risk that synchronous replication avoids?

level: middleimportance: must knowfreq 68%

basics

~20 s

Synchronous replication waits for a copy to confirm the write before telling the client it succeeded, so no acknowledged write is ever lost, but it's slower and can stall if a replica is unreachable. Asynchronous replication confirms instantly and copies data in the background, so it's fast, but a crash right after can lose the last few writes.

open as a page

A distributed lock service grants a lease to a client believed to hold exclusive access to a shared resource, such as a storage volume. The client experiences a long GC pause, its lease expires, and a second client acquires the lease and starts writing. The first client then resumes and, unaware its lease expired, also writes. How does fencing prevent this from corrupting the shared resource, and why isn't simply checking the lease is still valid on the client side sufficient?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Fencing gives each lease a rising number; the resource itself refuses any write whose number is older than the newest it has seen. That way even a confused old client that thinks it still owns the lock physically can't write, because the resource, not the confused client, enforces the check.

open as a page

The FLP impossibility result states that no deterministic consensus protocol can guarantee both safety and termination in a fully asynchronous system where even a single process may crash. Given that systems like Raft and Paxos are used in production to reach consensus every day, how do real systems reconcile their existence with this theoretical impossibility?

level: principalimportance: should knowfreq 35%

basics

~20 s

FLP proves that, in theory, a perfectly asynchronous network can always delay messages just enough to stall a consensus algorithm forever. Real systems dodge this by accepting that consensus might occasionally stall for a while, never violating correctness, rather than promising it will always finish quickly - and by using timeouts that work well enough in practice even though they aren't a formal guarantee.

open as a page