Jepsen-Style Verification and Failure Injection
LESSON
Jepsen-Style Verification and Failure Injection
By the end of this lesson, you will be able to...
Turn a distributed-system promise into operations, a history, and a checkable invariant.
Choose a fault schedule that challenges the assumptions behind that promise.
Interpret a counterexample or a clean run without claiming more certainty than the test earned.
Idea in one sentence: Failure injection becomes verification only when a recorded client history is checked against an explicit model of what the system was allowed to do.
Core Insight
Atlas has a controller lease for a critical rollout job. In staging, the team kills the leader process. Another controller takes over, dashboards become green, and the test passes.
At first, that seems like evidence of safety: the cluster recovered. But recovery answers a liveness question—did useful work resume? It does not answer the safety question: did a paused former leader perform an action after a newer leader had authority?
The difference becomes visible only in the client-visible history. Imagine that controller A pauses, its lease expires, controller B becomes leader, and then A wakes up with stale local state. If both A and B complete a protected action, a healthy dashboard cannot tell us whether the system maintained exclusive authority.
Jepsen-style verification makes that story inspectable. A workload invokes operations. A fault process disrupts the environment. Clients record invocations and results. A checker asks whether the completed results can fit the stated contract. The fault is pressure; the history is evidence.
The Production Symptom
The user-facing symptom is not necessarily an outage. Atlas might deploy two conflicting rollout plans, remove the wrong instance, or record a successful command from an actor that no longer owns the role.
Logs alone are hard to interpret because each component saw only part of the event. Controller A believes its earlier lease is still valid. Controller B sees that it acquired a newer lease. The protected deployment API may see two requests without knowing which one should have authority.
The initial test model is therefore too weak:
inject fault -> cluster elects a leader -> inspect health -> pass
It works for checking whether replacement occurs. It becomes insufficient when the promise is “only the current leader may change deployment state.” The missing element is a history-level invariant that can reject an impossible combination of successful results.
Write the Claim Before the Fault
For this lesson, Atlas gives each successful leadership acquisition an increasing fencing epoch. The deployment API persists the highest epoch it has accepted for this role.
The contract is:
Once the deployment API has accepted an action carrying epoch
8, it must reject every later action carrying an earlier epoch such as7.
This is deliberately narrower than “the entire cluster is correct.” It names the protected resource, the observable successful action, and the stale-actor boundary.
Plain meaning:
An old leader may still wake up, but it must not make the protected system go backward.
In Atlas:
acquire(A) -> epoch 7
acquire(B) -> epoch 8
apply(A, epoch 7) after B's accepted epoch 8 -> must fail
Technical name: this is a fencing invariant. It checks the enforcement point, not merely the coordination-store record. A lease can remove A's ownership record; fencing tests whether the target rejects A when A keeps running anyway.
For other systems, choose a different claim. A linearizable register needs an ordering of reads and writes consistent with real-time precedence. A queue may need “no acknowledged message is lost.” A lock may need “no two successful holders overlap.” Do not test a guarantee the API never promises.
The Test Pieces
The following four pieces form a useful Jepsen-style test:
| Piece | Atlas example | Why it matters |
|---|---|---|
| Workload | Clients acquire, renew, apply a fenced rollout, read current epoch, and release | Exercises the claimed contract |
| History recorder | Records invocation, completion, client, epoch, result, and operation ID | Lets a checker reconstruct what users observed |
| Fault process | Pauses a leader, partitions it from the store, delays a renewal, restarts a store node | Challenges the assumptions that make leadership safe |
| Checker | Rejects a successful stale epoch after a newer epoch was accepted | Turns the history into evidence |
The test does not need access to every protocol message. Jepsen describes this as opaque-box testing: real binaries run in a real cluster, while the test checks observable behavior. That is valuable because production users also observe results, not internal intentions.
Worked Trace: Finding a Stale Leader
The following is a synthetic trace. The times and identifiers are illustrative, not a report from a real cluster.
| Time | Client observation | Store or target state | What the checker records |
|---|---|---|---|
| 12:00:00 | A acquires leadership | lease owner A, epoch 7 | acquire(A) -> ok(7) |
| 12:00:01 | A applies rollout r-41 |
target accepts epoch 7 | apply(A, 7, r-41) -> ok |
| 12:00:02 | Fault process pauses A | A sends no renewals | fault: pause(A) |
| 12:00:12 | Lease expires; B acquires leadership | lease owner B, epoch 8 | acquire(B) -> ok(8) |
| 12:00:13 | B applies rollout r-42 |
target stores highest epoch 8 | apply(B, 8, r-42) -> ok |
| 12:00:15 | A resumes and retries | A still has epoch 7 locally | fault: resume(A) |
| 12:00:16 | A applies rollback r-40 |
target returns ok |
apply(A, 7, r-40) -> ok |
The last result is a counterexample. Once the target accepted epoch 8, it also accepted stale epoch 7. The violation is not “A sent a request.” A may send stale requests in any distributed system. The violation is that the protected target reported success.
Now compare a repaired system. The final row becomes:
apply(A, 7, r-40) -> rejected_stale_epoch
The old process still exists, the pause still happened, and the cluster may still have imperfect clocks. The design earns safety at the target boundary by checking a durable monotonic value.
So far, the history has made the stronger model concrete: do not infer exclusive authority from a clean failover; test whether stale authority is rejected after the failover.
From a Trace to a Checker
For a small fencing test, the checker can maintain one value:
highest_accepted_epoch = 0
for each successful apply(epoch):
if epoch < highest_accepted_epoch:
report violation
highest_accepted_epoch = max(highest_accepted_epoch, epoch)
This teaching model assumes completed successful actions are ordered by the target's durable acceptance record. A real checker must define how it handles concurrent and unknown operations. A timeout may mean an operation did not reach the target, or that it reached the target and its reply was lost. Do not silently treat every timeout as failure.
For linearizability, the model is different. Herlihy and Wing define a history as invocation and response events and ask whether completed operations can correspond to a legal sequential execution while respecting non-overlapping real-time order. The checker has to reason about concurrent intervals, not merely sort by client timestamps.
This is why an invariant should be as small as possible. A tiny model with clear operation semantics is easier to challenge and trust than a vague statement that “the system handled chaos.”
Choose Faults That Attack the Claim
The goal is not maximum drama. Choose faults that make the invariant's assumptions questionable:
| Claim | Useful workload | Faults that create pressure | Evidence to keep |
|---|---|---|---|
| Fencing rejects stale leaders | Acquire and apply with epochs | Pause owner, delay renewals, partition owner from store | Accepted epochs and target responses |
| Linearizable register | Concurrent read, write, compare-and-swap | Partition, leader failover, dropped replies | Invocation/response intervals and values |
| No lost acknowledged writes | Write, acknowledge, restart, read | Process crashes, disk disruption, dropped acknowledgements | Acknowledged write IDs and recovered state |
| Watch resync is correct | Read snapshot, watch, reconcile | Disconnect, compaction, rapid updates | Snapshot revision, event range, final cache |
A client pause is especially valuable for lease-based coordination because it separates two ideas that are easy to conflate: the coordination service can transfer authority, while the old client can remain alive and confused. A fault schedule should expose that gap.
Trade-offs and Limits of This Test
The trade-off is that a sharp history checker gives stronger evidence than a broad recovery demo, but it costs careful operation semantics, reliable recording, representative fault injection, and time to minimize failures. A tiny fencing model can prove that one target rejects stale epochs; it does not establish linearizability for every API, prove the coordination protocol correct, or cover every production workload.
The boundary is visible when the checker cannot explain an operation outcome. A timeout, partial log, or undocumented client retry rule should become an explicit unknown case or a repair to the instrumentation—not an invisible assumption that makes the history look clean.
Reading Results Honestly
When the checker finds a violation, preserve the smallest history that demonstrates it, the fault schedule, system version and configuration, random seed if available, and the raw evidence needed to reproduce it. The immediate triage question is: did the implementation violate the stated contract, or did the test assert a stronger contract than the system documented?
Both answers are useful. A stale read in a mode documented as eventually consistent is a test-design mismatch. A stale action accepted after an advertised fencing contract is an implementation or integration problem. The history makes that distinction discussable.
When a run is clean, report the scope honestly:
Under this workload, fault schedule, configuration, and checker, no violation was found.
That is evidence, not a proof. Jepsen's own analysis material notes that these tests are nondeterministic and can find errors but cannot prove correctness. Increase confidence with varied seeds, repeated runs, focused regressions, and new workloads when the product or client contract changes.
Mitigation and Prevention
When a counterexample appears, do not begin by adding more random faults. First make the contract and enforcement point explicit. For the Atlas case:
- Make the target persist and compare fencing epochs on every leader-only action.
- Return a distinct stale-epoch error so the old leader stops and resynchronizes.
- Add the counterexample as a deterministic regression test.
- Re-run the focused workload with the original pause and partition schedule.
The prevention is not “never pause.” It is designing operations so a pause cannot turn old authority into accepted work.
Signals to Watch
An operational verification run should expose its own evidence:
- number of completed, failed, and unknown operations;
- fault start and end times, including which client or node was affected;
- checker result and a minimized counterexample when it fails;
- highest accepted fencing epoch and rejected-stale count;
- watch cancellations, reconnects, or compaction events when the contract includes resync;
- test seed, software version, configuration, and topology.
The signal that a result is weak is a history with missing operation outcomes or a checker whose model does not state what timeout, retry, and concurrency mean. Repair the instrumentation or the model before declaring the system safe.
Check Your Understanding
Check: A test pauses the leader, another client becomes leader, and dashboards show one healthy leader after recovery. Has the test established that stale leader actions were impossible?
Think first, then reveal.
Answer: No. It established that leadership recovered. To establish the safety claim, clients must attempt and record protected actions, and a checker must reject a history where an earlier epoch succeeds after a newer one.
Check: A client times out while applying epoch 8. Should the checker always record that operation as failed?
Think first, then reveal.
Answer: No. The result may be unknown: the target could have accepted the action and the reply could have been lost. The model must say how unknown operations are treated rather than silently erasing uncertainty.
Practice: Design a Failure Test for a Watch Client
A controller reads a snapshot at revision 500, watches configuration changes, and rebuilds its local cache after a disconnection. The team claims, “after recovery, the cache eventually equals the authoritative state and never applies an update older than its snapshot.”
Specify a focused test. Name the operations to record, a fault schedule, the invariant, and the evidence that would make a failure reproducible.
Model answer: Record snapshot reads with their revision and contents, watch starts and cancellations, every received event with revision, cache updates, and a final authoritative read. Disconnect the controller while writing revisions 501 through 520, optionally compact history before it reconnects, then restore connectivity. Check that the controller takes a new snapshot when it cannot continue from its old revision; its final cache must equal the final authoritative state, and it must not apply an event at or before the revision represented by its replacement snapshot. Preserve the snapshot revisions, event sequence, final state, disconnect/compaction times, client logs, and random seed.
Connections
The previous lesson showed that coordination clients need snapshot-and-watch recovery and fencing at the protected boundary. This lesson turns those design requirements into histories and invariants that can be challenged under pause, partition, and retry.
The next capstone asks where consensus belongs in a control plane. Its answer is incomplete until every authority claim has a corresponding failure test and checker.
Resources
- [DOC] Jepsen Analyses — Focus: Opaque-box, generative tests that construct concurrent histories and check them against models under partial failure.
- [SOURCE] Jepsen Open-Source Framework — Focus: The framework's workload, nemesis/fault, and checker structure.
- [PAPER] Linearizability: A Correctness Condition for Concurrent Objects — Focus: Invocation/response histories and real-time ordering for a linearizability claim.
- [DOC] etcd API — Focus: A concrete coordination API whose revisions, watches, and leases can become testable client contracts.
Key Takeaways
- A fault test is verification only when it checks a recorded history against a named contract.
- Fencing tests whether a protected target rejects stale actors; successful failover alone does not establish that safety property.
- The best fault schedule attacks the assumptions behind one narrow invariant instead of creating unstructured chaos.
- A clean run is bounded evidence under its workload and fault model, never a general proof of correctness.
- Counterexamples, unknown outcomes, seeds, and fault timelines are operational artifacts, not debugging leftovers.