Your country

Tools that support it use your country for local currency, number formats, units and paper size. Your choice is saved only in this browser.

Type a name or a two-letter code. Use the up and down arrow keys to move through the countries, Enter to choose one and Escape to close.

System Design (High-Level Design)  Module 3 – Performance and reliability fundamentals

How distributed systems fail

Classify crash, omission, timing and Byzantine faults, spot partial, gray and metastable failures, and handle the unknown outcome of a call that timed out.

  • Intermediate
  • 25 minutes
  • Examples run with Python 3.14.8, Pyodide 314.0.7, Node.js 24.21.0 and quickjs 0.32.0
  • By MySmartCoPilot

What you will learn

  • Classify crash, omission, timing and Byzantine faults
  • Recognise partial, gray and metastable failures from their symptoms
  • Design for the failure modes that actually occur in production, starting with unknown outcomes

Before you start

On this page

On one computer, a program mostly either works or stops. A distributed system can do something worse: fail partially. Some machines work, some have crashed, some messages arrive and some are lost, and a call can end without any answer at all. That last case is the hardest to design for, because a call that timed out may have succeeded, failed, or still be running. This lesson names the kinds of faults that distributed systems are designed against (crash, omission, timing and Byzantine), shows how a timeout leaves an unknown outcome and how idempotency keys handle it, and describes two failure patterns that health checks and retries make worse rather than better: gray failures and metastable failures.

Partial failure, and the unknown outcome

Amazon’s Builders’ Library takes one request and reply between a client and a server and lists eight steps that can each go wrong independently: sending the request, delivering it, the server validating it, the server updating its state, the server sending the reply, delivering the reply, the client validating it, and the client updating its own state. Client, server and network fail independently of each other, and most of these failures look identical from the client’s side: nothing comes back.

A remote call can fail on the way in, inside the server, or on the way back, and a timeout cannot tell which.Client sends'charge order 7'Networkthe request may be lost or lateServermay crash before or aftercharging the cardNetworkthe reply may be lost or lateClientreply, or a timeout:was the card charged?

One remote call, and the places it can fail

Text description of the diagram

The diagram follows one call from top to bottom.

  1. The client sends a request to charge order 7.
  2. The network carries it, and may lose or delay it.
  3. The server may crash before charging the card, or after charging it but before replying.
  4. The network carries the reply, and may lose or delay it.
  5. The client receives the reply or, if nothing arrives in time, a timeout.

A timeout looks the same whether the request was lost, the server crashed before acting, or the card was charged and the reply was lost. The client cannot tell which, so it has to handle an unknown outcome.

The article’s conclusion is the rule this lesson is built around: when a call times out, its result is unknown. The request may never have arrived, the server may have crashed before acting, or it may have acted and lost its reply. The client has to handle all three at once. This is not a rare corner: the same article warns that bugs in how a service handles these cases can stay hidden in production for months before a rare combination of failures sets them off.

Four fault models

Designs and protocols state which faults they tolerate, in the terms of a system model. Martin Kleppmann’s Cambridge lecture notes describe a system model in three parts, network, nodes and timing, and the four classic kinds of fault map onto them. From the easiest to tolerate to the hardest:

  • Crash. A node stops, and may restart with only what it saved to disk. From outside it looks like silence, then a node that has forgotten its memory; a process killed for using too much memory is the everyday case. Defences: replicas, durable writes and restarts.
  • Omission. A message is lost on the way in or on the way out, for example by a congested switch. From outside it looks like silence. Defences: timeouts, retries and deduplication.
  • Timing. A correct answer arrives too late, after a garbage-collection pause or behind a slow disk: silence, then a late reply. Defences: timeouts, hedging, and protocols that stay correct when delays have no bound.
  • Byzantine, or arbitrary. A node does anything at all: wrong answers, lies, different stories to different peers, because of corrupted memory, a compromised machine or a participant who cheats. Its answers look valid. Defences: signatures, checksums and protocols that tolerate a minority of liars.

The notes distinguish crash-stop nodes, which never come back, from crash-recovery nodes, which may restart having lost their memory but kept their disk; most real services are crash-recovery. On the network side, a “fair-loss” link may lose, duplicate or reorder messages, but a message retried often enough eventually gets through, which is why retries plus deduplication can turn it into a reliable link. For timing, the notes explain why real systems are only partially synchronous: networks and nodes are usually predictable, then occasionally pause for garbage collection, page faults, congestion or rerouting, and no timeout can tell a slow node from a dead one.

Byzantine faults are the general case, named after Lamport, Shostak and Pease’s paper on generals who must agree on a plan while some of them are traitors. Their result is that with ordinary messages, agreement is possible only if more than two thirds of the generals are loyal; with unforgeable signed messages it is possible for any number. Kleppmann notes that a bug shared by every node cannot be outvoted by such a protocol, because every node has it, so the term is usually kept for deliberate deviation. Services inside one organisation assume crash-recovery faults; blockchains and other systems shared between parties who do not trust each other must tolerate Byzantine ones.

A timeout does not tell you what happened

The program below pays for 10,000 orders through a fake payment server. The network loses 2% of requests before they reach the server and 3% of replies after the server has charged the card; the client sees the same timeout in both cases. Three clients face the same network luck:

Paying for 10,000 orders over a lossy network, three ways Python · unknown_outcome.py
"""A payment call that times out has an unknown outcome: the charge may or may not have happened.

A fake payment server charges a card when a request reaches it. The network loses 2% of requests on the way in
and 3% of replies on the way out; in both cases the client only sees a timeout. Three client strategies pay
for the same 10,000 orders with the same network luck. The rates are assumptions for the sketch.
"""

import random


class PaymentServer:
    def __init__(self, deduplicate):
        self.deduplicate = deduplicate
        self.charges = []  # every charge made: the order it was for
        self.done = {}  # idempotency key -> charge id, when deduplicating

    def charge(self, order, key):
        if self.deduplicate and key in self.done:
            return self.done[key]  # a retry: the same answer, no new charge
        self.charges.append(order)
        charge_id = f"ch-{len(self.charges)}"
        if self.deduplicate:
            self.done[key] = charge_id
        return charge_id


def pay(server, order, network, retries):
    """Try up to 1 + retries times; return True when a reply confirmed the payment."""
    key = f"order-{order}"  # the client picks one key per order and reuses it on every retry
    for _ in range(1 + retries):
        if network.random() < 0.02:  # the request is lost: the server never sees it
            continue
        server.charge(order, key)
        if network.random() < 0.03:  # the charge happened, but the reply is lost
            continue
        return True
    return False


print(f"{'strategy':28}{'charges':>8}{'charged 2+ times':>18}{'unconfirmed':>13}{'charged, but told no':>22}")
for name, retries, deduplicate in [
    ("no retry", 0, False),
    ("retry up to 3 times", 3, False),
    ("retry with idempotency key", 3, True),
]:
    server, network = PaymentServer(deduplicate), random.Random(9)
    confirmed = {order: pay(server, order, network, retries) for order in range(10_000)}
    per_order = {}
    for order in server.charges:
        per_order[order] = per_order.get(order, 0) + 1
    twice = sum(1 for n in per_order.values() if n > 1)
    unconfirmed = sum(1 for ok in confirmed.values() if not ok)
    charged_unconfirmed = sum(1 for order, ok in confirmed.items() if not ok and order in per_order)
    print(f"{name:28}{len(server.charges):>8,}{twice:>18,}{unconfirmed:>13,}{charged_unconfirmed:>22,}")

Output

strategy                     charges  charged 2+ times  unconfirmed  charged, but told no
no retry                       9,802                 0          503                   305
retry up to 3 times           10,325               308            0                     0
retry with idempotency key    10,000                 0            0                     0

Recorded with Python 3.14.8 on macOS 26 arm64. To run it yourself: mise exec python@3.14.8 -- python3 unknown_outcome.py

  • No retry leaves 503 orders unconfirmed. Of those, 305 were in fact charged: their replies were lost. Those customers paid and were told the payment failed.
  • Retry up to 3 times confirms every order, but charges 308 customers twice or more, because a retry after a lost reply makes a second charge.
  • Retry with an idempotency key makes exactly 10,000 charges. The client picks one key per order and sends it with every attempt; the server records the result under the key, and a repeat gets the stored result instead of a new charge.

An idempotency key turns “do this” into “do this once”, which is what makes retries safe. When even the retries run out, the key still lets a later process look the payment up and settle it, which is called reconciliation: the client, or a background job, asks the server what happened to order 7 rather than guessing. The next lessons of this module build retries, timeouts and backoff on top of this.

UUID Generator Generate random UUIDs to use as idempotency keys, one per logical operation.

Gray failures: healthy to the checker, broken for users

A gray failure is a component that is failing for the people who use it while the system’s own failure detectors see it as healthy. Huang and colleagues, writing from experience with large cloud systems, argue that this differential observability is the defining trait of gray failures, and that subtle faults of this kind, rather than clean crashes, cause many of the worst outages: a crash is detected and handled, a half-working component is not. Their example is a server whose request-handling module is stuck while its heartbeat module still runs, so every check that relies on heartbeats reports it as healthy.

A checkout service whose health checks stay green while it fails JavaScript · gray_failure.mjs
// A gray failure: the health checks stay green while users fail. A checkout service has a shallow health check
// (is the process up?) and a deep one (can it reach its database?). From minute 4, the path to one of its three
// pricing shards breaks, so checkouts for customers on that shard fail. The workload is an assumption.

function generator(seed) {
  let a = seed >>> 0;
  return () => {
    a = (a + 0x6d2b79f5) >>> 0;
    let t = Math.imul(a ^ (a >>> 15), a | 1);
    t ^= t + Math.imul(t ^ (t >>> 7), t | 61);
    return ((t ^ (t >>> 14)) >>> 0) / 4294967296;
  };
}
const random = generator(17);

const service = {
  brokenShard: null, // which pricing shard is unreachable, if any
  shallowHealth: () => "200 ok", // the process answers
  deepHealth: () => "200 ok", // the database answers SELECT 1
  checkout(customer) {
    const shard = customer % 3; // customers are spread over three pricing shards
    if (shard === this.brokenShard) return "503 pricing unavailable";
    return random() < 0.001 ? "500 error" : "200 paid";
  },
};

console.log("minute  shallow check  deep check  checkouts  succeeded  alarm on success rate < 99%");
for (let minute = 1; minute <= 8; minute++) {
  if (minute === 4) service.brokenShard = 2;
  let ok = 0;
  const requests = 600;
  for (let i = 0; i < requests; i++) {
    const customer = Math.floor(random() * 90000);
    if (service.checkout(customer).startsWith("200")) ok += 1;
  }
  const rate = ok / requests;
  const alarm = rate < 0.99 ? "FIRING" : "-";
  console.log(
    `${String(minute).padStart(6)}  ${service.shallowHealth().padEnd(13)}  ${service.deepHealth().padEnd(10)}  ${String(requests).padStart(9)}  ${(100 * rate).toFixed(1).padStart(8)}%  ${alarm}`,
  );
}

Output

minute  shallow check  deep check  checkouts  succeeded  alarm on success rate < 99%
     1  200 ok         200 ok            600     100.0%  -
     2  200 ok         200 ok            600     100.0%  -
     3  200 ok         200 ok            600      99.7%  -
     4  200 ok         200 ok            600      67.3%  FIRING
     5  200 ok         200 ok            600      65.5%  FIRING
     6  200 ok         200 ok            600      68.8%  FIRING
     7  200 ok         200 ok            600      67.3%  FIRING
     8  200 ok         200 ok            600      69.5%  FIRING

Recorded with Node.js 24.21.0 on macOS 26 arm64. To run it yourself: mise exec node@24.21.0 -- node gray_failure.mjs

From minute 4, one of the service’s three pricing shards is unreachable, so about a third of checkouts fail. The shallow health check only asks whether the process answers, and the deep one only asks whether the database answers, so both stay at “200 ok” while the success rate drops to between 65% and 70%. Only the measurement of what users actually get, the share of successful checkouts, raises the alarm. The paper’s advice is to move from a single failure detector, such as a heartbeat, to several signals of health, and to measure from where the applications stand; in practice that means alerting on user-facing success rates, the SLIs of the previous lesson, and health checks that exercise the paths users depend on.

Metastable failures: when the fix keeps the failure going

Bronson, Aghayev, Charapko and Zhu describe a pattern behind widespread outages at large internet companies, some lasting hours: a metastable failure. A system runs in a vulnerable state, efficient and healthy, but with a hidden feedback loop. A trigger pushes it into a bad state, and a sustaining effect keeps it there even after the trigger is gone. Leaving the state takes a strong push, such as cutting the load sharply or changing the policy that feeds the loop.

Their first example is retries. A web application normally sends 280 queries a second to a database that answers quickly below 300 queries a second and much more slowly above it, and each user request retries a query that has not answered within a second. A 10-second network outage triggers a surge of queued requests and retries; the database slows, so queries time out and are retried, and the load stays at 560 queries a second, twice the real demand. The network has recovered, yet the system stays down until the load drops or the retry policy changes. The paper notes that outages like this are often blamed on the trigger, when the real root cause is the sustaining effect, and that the features creating the loop, retries, failover and caching, are usually ones added for efficiency or reliability.

The defences are the subject of the rest of this module: retries with backoff and a budget, so that retries cannot double the load; circuit breakers that stop calling a failing dependency; load shedding, so that an overloaded service refuses work quickly instead of slowly failing all of it; and capacity tests that look for the load at which the system becomes vulnerable.

Design for the failures that happen

  • Put a timeout on every remote call, and treat a timeout as an unknown outcome, not a failure.
  • Make every operation that changes state idempotent, with a key chosen by the caller, so retries are safe.
  • Reconcile what timed out: look outcomes up by key rather than guessing.
  • Alert on what users get (success rates and latency at the edge), not only on health checks.
  • Make health checks exercise real dependencies, and keep the action they trigger proportionate: a check that removes every server when a shared dependency fails turns a partial failure into a total one.
  • Give retries, failovers and other recovery actions a budget, so they cannot become the sustaining effect.
  • Assume crash-recovery faults inside your organisation, and design for Byzantine faults only where participants do not trust each other.

Practice

Exercise · Medium · Python

A payment server that charges once per idempotency key

Complete PaymentServer.charge(key, amount) in payments.py, so that a client may safely retry a payment whose reply it never received.

- The first request with a new idempotency key makes a charge: it appends amount to self.charges and returns a new charge id, "ch-1" for the first charge the server ever makes, "ch-2" for the second, and so on. - A request that repeats a key with the same amount is a retry: it makes no new charge and returns the id of the original charge. - A request that repeats a key with a different amount is a client bug, not a retry: raise ValueError and make no charge.

For example, charge("order-7", 500) twice returns "ch-1" both times and leaves charges as [500].

Starter code · payments.py

class PaymentServer:
    """A payment server that charges at most once per idempotency key."""

    def __init__(self):
        self.charges = []  # every amount actually charged, in order
        self.seen = {}  # idempotency key -> (charge id, amount)

    def charge(self, key, amount):
        """Charge `amount` once for `key` and return the charge id; a retry returns the same id."""
        # Replace this line with your code.
        return ""
The sample tests · test_payments.py
from payments import PaymentServer


def test_first_charge():
    """a new key makes a charge and returns a new id"""
    server = PaymentServer()
    assert server.charge("order-7", 500) == "ch-1"
    assert server.charge("order-8", 120) == "ch-2"
    assert server.charges == [500, 120]


def test_retry_charges_once():
    """a repeated key with the same amount returns the first id and charges nothing"""
    server = PaymentServer()
    first = server.charge("order-7", 500)
    assert server.charge("order-7", 500) == first
    assert server.charge("order-7", 500) == first
    assert server.charges == [500]


def test_mismatched_amount():
    """a repeated key with another amount is refused and charges nothing"""
    server = PaymentServer()
    server.charge("order-7", 500)
    try:
        server.charge("order-7", 900)
        refused = False
    except ValueError:
        refused = True
    assert refused, "a different amount for the same key must raise ValueError"
    assert server.charges == [500]


def test_ids_count_real_charges():
    """ids count the charges made, so retries do not use up numbers"""
    server = PaymentServer()
    server.charge("a", 1)
    server.charge("a", 1)
    assert server.charge("b", 2) == "ch-2"
A hint

Keep a dictionary from each key to what it produced, for example self.seen[key] = (charge_id, amount). Look the key up before charging: if it is there, compare the amounts and either return the stored id or raise; if it is not, make the charge, store the result and return the new id. The id can come from the number of charges made so far.

The sample tests run on this device, in your browser (Pyodide): nothing is sent to mysmartcopilot.com. The first run downloads Python (about 13.5 MB), which is kept for the next runs. A check in your browser is feedback for you, not proof that the code is right for every input.

Check yourself

8 questions about this lesson. Every answer and why it is right is on the page, behind “Show the answer”. Your score stays in this browser.

  1. Question 1 of 8 A cache node runs out of memory, is killed by the operating system and restarts a minute later with an empty cache. Which kind of fault is this?

    Choose one answer.

    Show the answer to question 1

    Answer: A crash fault, in the crash-recovery model

    The node stopped and came back, losing what it held in memory, which is the crash-recovery model. Anything that depended on its memory, such as the cached data, has to be rebuilt.

  2. Question 2 of 8 A switch between two racks drops about 1% of packets. Some calls never get a reply and time out; retried, they succeed. Which kind of fault is this?

    Choose one answer.

    Show the answer to question 2

    Answer: An omission fault, because messages are lost

    Messages that should arrive do not. The switch keeps working, so nothing crashed, and what does arrive is correct, so nothing is arbitrary. Retries with idempotent requests turn such lossy links into reliable ones.

  3. Question 3 of 8 A replica answers every query correctly, but long garbage-collection pauses make some replies arrive after the client's 2-second timeout. Which kind of fault is this?

    Choose one answer.

    Show the answer to question 3

    Answer: A timing fault, because correct answers arrive too late

    The answers are right but late. The client cannot tell a paused node from a dead one, so its timeout treats the call as failed even though the replica will reply, and may already have acted.

  4. Question 4 of 8 A server with a faulty memory chip returns corrupted account balances while reporting success. Which kind of fault is this?

    Choose one answer.

    Show the answer to question 4

    Answer: An arbitrary (Byzantine) fault, because the server gives wrong answers that look right

    The node does not stop or go silent: it says something false as if it were true. Systems that assume only crashes do not defend against this, which is why storage systems add checksums to detect corruption.

  5. Question 5 of 8 Every host passes its health check, but 20% of users' uploads fail because one storage path is broken. What is this called?

    Choose one answer.

    Show the answer to question 5

    Answer: A gray failure, because the failure detector sees health while users see failure

    The health check and the users observe different things, which the gray-failure paper calls differential observability. Measuring the success rate of real uploads, or a check that exercises the broken path, would have caught it.

  6. Question 6 of 8 A 10-second network blip makes clients retry. The retries keep the database overloaded for an hour after the network recovered, until retries are switched off. What is this?

    Choose one answer.

    Show the answer to question 6

    Answer: A metastable failure, because a feedback loop keeps the system down after the trigger is gone

    The trigger, the blip, lasted seconds; the sustaining effect, the extra load from retries, kept the database overloaded. A metastable failure ends only with a strong push, such as cutting the load or the retries.

  7. Question 7 of 8 A call to charge a card times out. What does the client know?

    Choose one answer.

    Show the answer to question 7

    Answer: Nothing certain: the charge may have failed, succeeded with a lost reply, or still be running

    A timeout means no reply arrived in time, not that nothing happened. The safe response is to retry with the same idempotency key, or to look up the payment's status by that key, so the customer is charged once.

  8. Question 8 of 8 In a ledger shared by several companies, one participant sends different lists of transactions to different peers. Which fault model must the protocol tolerate?

    Choose one answer.

    Show the answer to question 8

    Answer: Byzantine faults, because a participant may deviate from the protocol on purpose

    When participants do not trust each other, one may lie, and differently to different peers. Lamport, Shostak and Pease showed that with plain messages, agreement then needs more than two thirds of the participants to be honest. Services inside one company usually assume crash-recovery faults instead.

Interview questions

Warm-up (fresher to mid level): what is a partial failure, and why does it make distributed systems harder than one machine? A partial failure is one in which some parts of a system fail while others keep working: one node crashes, one network link drops messages, one dependency slows down. On a single machine a failure usually stops the whole program, which is easy to notice. In a distributed system the working parts must carry on around the broken ones, and a caller often cannot tell which happened: a timed-out call may have failed, succeeded with a lost reply, or still be running. So every remote call needs a timeout, and every operation that changes state needs to be safe to retry.

Key takeaways

  • A remote call can fail at eight points between client and server; a timeout means its outcome is unknown.
  • Crash, omission, timing and Byzantine faults are increasingly hard to tolerate; internal services assume crash-recovery, multi-party systems Byzantine.
  • Idempotency keys make retries safe: without them, retries after lost replies charge customers twice; without retries, customers are charged and told it failed.
  • Gray failures pass health checks while users fail: alert on user-facing success rates.
  • Metastable failures persist after their trigger because a feedback loop, often retries, sustains them; budget every recovery action.

References

Related tools

Report a problem with this lesson

Quick answers and tool search

Type to search tools or to get a quick answer, for example 18% of 2500. Use the up and down arrow keys to move through the results, Enter to choose, and Escape to close.