A distributed system fails in parts. One replica stops answering while the others carry on, a switch drops traffic in one direction only, a garbage-collection pause makes a leader look dead for twelve seconds and then it wakes up still believing it leads. Each of those cases has code written to handle it, and most of that code runs so rarely that nobody knows whether it works. Chaos engineering runs it on purpose: state a hypothesis about how the system behaves under a specific fault, inject that fault in a controlled way, measure, and either gain confidence or find the bug before an outage does.
This page takes the distributed-systems view: which faults to inject, which tool injects each one and at which layer, and how to check the property that matters most and is checked least, which is correctness rather than availability. The general practice is covered in the chaos engineering guide and rehearsals with people in the loop in game days. The telemetry side, with steady state written as SLI queries and abort guards driven by error-budget burn, is in chaos engineering and observability. The worked example below partitions the leader of a three-node replicated store and checks that no acknowledged write is lost.
An experiment harness in one picture
Principles, restated for partial failure
The method has four steps: define steady state as a measurable output, hypothesise that it holds in both a control and an experimental group, introduce variables that reflect real-world events, and try to disprove the hypothesis. In a distributed system the interesting faults are partial and the interesting outcomes are quiet.
Partial means the fault affects some nodes, some links or one direction of a link. Crashing every replica at once tests your backup restore, not your replication protocol. Quiet means the outcome is often not an error rate. A split brain that lets two leaders accept writes for four seconds may produce no errors at all; the damage is two divergent histories, one of which is thrown away when the partition heals. If your steady state is an error rate below 0.1%, the experiment passes and the data is gone.
So every experiment needs two kinds of hypothesis. An availability hypothesis says what clients see: with one replica partitioned, p99 write latency stays under 300 ms after a recovery window of at most 15 seconds. A safety hypothesis says what must never happen, however bad availability gets: no acknowledged write disappears, no read returns a value older than one a completed read already returned, no two nodes act as leader for the same term. Availability can be checked from metrics. Safety needs a record of operations, and most chaos programmes skip it.
The fault model
Start from what actually goes wrong in production and map each fault to a way of injecting it and to the code it exercises. The table is the minimum fault model for a replicated service.
| Fault | Real-world cause | Inject with | What it exercises |
|---|---|---|---|
| Process crash | OOM kill, host failure, bad deploy | SIGKILL, pod delete | restart recovery, log replay, leader election |
| Process pause | GC pause, VM steal time, live migration | SIGSTOP then SIGCONT | leases, fencing tokens, writes from a stale leader |
| Full partition | switch failure, security-group change | iptables DROP both ways, Chaos Mesh partition | quorum loss, election, minority-side behaviour |
| One-way partition | asymmetric ACL, NIC fault | iptables DROP in one direction | failure detectors that disagree, leadership flapping |
| Latency and jitter | congested or cross-region path | tc netem delay, Toxiproxy latency | timeouts, retries, hedged requests, election timeouts |
| Packet loss | bad optics, overflowing buffers | tc netem loss | TCP retransmission stalls, tail latency |
| Bandwidth cap | saturated link, noisy neighbour | Toxiproxy bandwidth, Chaos Mesh bandwidth | replication lag, backpressure |
| Clock skew or jump | NTP misconfiguration, VM resume | libfaketime, or stepping the clock on an isolated host | lease expiry, TTLs, timestamp ordering, certificate checks |
| Disk full or slow | log growth, failing drive | fill the volume; a device-mapper delay target | fsync stalls, write-ahead log behaviour at ENOSPC |
| Dependency errors | downstream bug, throttling | mesh fault abort, Toxiproxy reset_peer or timeout | circuit breakers, fallbacks, retry amplification |
Begin with crash, pause, partition and latency. Pause is the one people forget: a paused process is not dead, and when it resumes it acts on stale state. Slow-but-alive components are the subject of gray failure; latency injection tests your detectors against them.
Injecting faults: the layer decides what breaks
The layer an injector works at decides what the fault touches. Kernel-level tools act on packets. tc netem on the root qdisc shapes egress only, so apply it on both hosts to delay both directions. iptables DROP makes a partition that looks like a real one, with timeouts; REJECT gives fast errors instead. Kernel faults hit everything on the path, including heartbeats and your own SSH session.
# Latency: 100 ms +/- 20 ms, normally distributed, on traffic leaving eth0
tc qdisc add dev eth0 root netem delay 100ms 20ms distribution normal
# Switch to 5% loss instead, then heal
tc qdisc change dev eth0 root netem loss 5%
tc qdisc del dev eth0 root
# Partition this host from 10.0.1.12 in both directions, then heal
iptables -I INPUT -s 10.0.1.12 -j DROP
iptables -I OUTPUT -d 10.0.1.12 -j DROP
iptables -D INPUT -s 10.0.1.12 -j DROP
iptables -D OUTPUT -d 10.0.1.12 -j DROP
# A 12-second "GC pause"
kill -STOP "$PID"; sleep 12; kill -CONT "$PID"Proxy-level injection, as with Toxiproxy, places a TCP proxy in front of one dependency and adds toxics such as latency, bandwidth, timeout and reset_peer to that connection only, through an HTTP API on port 8474. It is precise and safe, but the client must connect through the proxy, so it suits integration tests. Mesh-level injection delays or aborts a percentage of HTTP requests in the sidecar; it is good for testing fallbacks but never touches heartbeats or replication on ports the mesh does not carry. Orchestrated injection, such as Chaos Mesh on Kubernetes, wraps kernel tools in a resource with selectors, a mode such as one or fixed-percent, a direction and a duration after which the fault is removed.
apiVersion: chaos-mesh.org/v1alpha1
kind: NetworkChaos
metadata:
name: isolate-kv-leader
spec:
action: partition
mode: one
selector:
namespaces: [kv]
labelSelectors:
role: leader
direction: both
target:
mode: all
selector:
namespaces: [kv]
labelSelectors:
app: kv
duration: "30s"
Recording a history and checking safety
A safety hypothesis is checked against a history: every operation a client invoked, when it started and finished, and what the client was told. The outcome has three values. An acknowledged write must survive. A rejected write must not appear. A write that timed out is unknown: it may have been applied, may be applied later, or never. Treat timeouts as failures and you report false bugs; treat them as successes and you miss real ones.
The harness below is deliberately small; client stands for your store's client wrapper. It runs a write workload, injects a fault, always heals in a finally block, waits for the cluster to settle, and checks that each key's final value is its last acknowledged write, an overlapping acknowledged write, or a write with an unknown outcome.
import random, threading, time, uuid
history, lock = [], threading.Lock()
def record(**ev):
with lock:
history.append(ev)
def writer(client, stop):
while not stop.is_set():
key, val = f"k{random.randrange(50)}", uuid.uuid4().hex
t0 = time.monotonic()
try:
client.put(key, val, timeout=1.0)
outcome = "ok" # acknowledged: must survive
except TimeoutError:
outcome = "unknown" # may or may not have been applied
except Exception:
outcome = "fail" # rejected, per the client contract
record(key=key, value=val, t0=t0, t1=time.monotonic(), outcome=outcome)
def lost_acks(reader):
"""Final value of each key must be its last ack, or a write that could have landed later."""
lost = []
for key in {e["key"] for e in history}:
ops = [e for e in history if e["key"] == key]
acks = [e for e in ops if e["outcome"] == "ok"]
if not acks:
continue
last = max(acks, key=lambda e: e["t1"])
allowed = {last["value"]}
allowed |= {e["value"] for e in ops if e["outcome"] == "unknown"}
allowed |= {e["value"] for e in ops if e["outcome"] == "ok" and e["t1"] > last["t0"]}
got = reader.get(key, linearizable=True)
if got not in allowed:
lost.append((key, last["value"], got))
return lost
def run(cluster, inject, heal, fault_s=30, settle_s=20, workers=8):
stop = threading.Event()
threads = [threading.Thread(target=writer, args=(cluster.client(), stop))
for _ in range(workers)]
for t in threads:
t.start()
time.sleep(10) # baseline with no fault
inject()
try:
time.sleep(fault_s)
finally:
heal() # heal even if the harness is interrupted
time.sleep(settle_s) # let election and catch-up finish
stop.set()
for t in threads:
t.join()
return lost_acks(cluster.client())This catches lost acknowledged writes and nothing subtler. Linearizability and transactional isolation need a proper checker, such as Jepsen's Knossos and Elle or the Porcupine library in Go; see what Jepsen taught us. Log operation ids on the server too, so a violation can be traced to its request.
Worked example: partitioning a Raft leader
Take a three-node key-value store that replicates with Raft. Node A leads, B and C follow, and clients connect to any node; followers forward writes to the leader. The experiment isolates A from B and C in both directions for 30 seconds, using the NetworkChaos resource above. Clients can still reach all three nodes. The numbers below are illustrative of a typical run, not measurements of a particular product.
Hypotheses. Availability: B and C elect a new leader and writes through them resume within 5 seconds, given a 1-second election timeout and client retries; writes sent to A fail or time out. Safety: no acknowledged write is lost, reads issued as linearizable never return stale data, and at most one node acts as leader in any term.
What happened. B won the election about 1.4 seconds in, and writes through B and C recovered within 3 seconds. Writes sent to A timed out and were recorded as unknown; on healing, A stepped down and truncated its uncommitted entries. The lost-ack check passed. Two bugs appeared anyway. First, A kept answering reads from its local state for the whole 30 seconds, because the service served leader reads without confirming leadership; a client reading from A saw values that B had already overwritten, which is a stale read on an operation documented as linearizable. Second, the client library retried timed-out writes against B with no request id, so an increment applied on the old leader before the partition was applied a second time.
Fixes and rerun. Leader reads now confirm leadership with a heartbeat round to a majority, the read-index approach from the Raft thesis, and retries carry an idempotency key. The rerun passed and joined the nightly suite. Neither bug changed the error rate: an availability-only experiment would have passed.
Blast radius and abort controls
Blast radius is the share of users, data and capacity an experiment can hurt. Grow it in steps: one process under Toxiproxy, then a staging cluster, then one production instance or shard, then a zone. Production matters because staging differs in traffic mix, data volume and configuration, and many outage bugs live in those differences.
Three controls keep production experiments safe. A heal timer owned by the injector, so the fault ends even if the controller crashes; Chaos Mesh's duration field does this, and for hand-rolled faults a dead-man switch that heals unless renewed does the same. An abort guard that heals when an SLO burns too fast. And one command, known to the on-call engineer, that removes every active fault. Run experiments in working hours with the owning team present, and never during a deploy, because two changes at once make the result unreadable.
Failure modes
- The fault outlives the experiment. An iptables rule or netem qdisc survives a crashed script. Keep heal commands idempotent, run them in
finally, and sweep hosts for leftover rules after every run. - The injector cuts its own control path. netem on the interface that carries SSH or the agent connection can make the host unreachable. Exclude management traffic, or inject from an agent that reaches the host another way.
- The fault misses the traffic. A mesh delay on a port replication does not use tests nothing. Measure that the fault took effect before reading results.
- No load, no finding. A partition during idle time exercises elections but not the write path.
- Clock faults break more than the target. Skewing a host clock can invalidate TLS certificates for everything on that host. Prefer per-process tools such as libfaketime.
Trade-offs
Kernel faults are realistic and blunt, proxies are precise and need client changes, meshes speak HTTP but miss everything else. Production experiments find the most and carry real risk. Random instance killing keeps teams honest about restarts but rarely finds protocol bugs, which need targeted faults and a checker. A small fixed suite of targeted experiments run on every release usually returns more than a large random one run occasionally.
What to do next
- Write down the fault model for one service: which of crash, pause, partition, one-way partition, latency, loss and clock skew it must survive.
- For each fault, write an availability hypothesis with numbers and a safety hypothesis stated as something that must never happen.
- Build a harness that records every operation with an id, start and end times, and an ok, fail or unknown outcome.
- Run the first experiment locally with Toxiproxy or netem, with a workload running and a checker for lost acknowledged writes.
- Move the passing experiments to staging under an injector with an owned heal timer and a single command that removes every fault.
- Add an abort guard tied to SLO burn before running anything in production, and start with one instance or shard.
- Turn every experiment that found a bug into a regression test that runs on each release.