Learn Labs
9. The Trouble with Distributed Systems

9.7 Formal Methods and Randomized Testing

Because of concurrency, partial failures, and network delays, THERE ARE A HUGE NUMBER OF POTENTIAL STATES. We need to guarantee that the properties hold in EVERY possible state and that we haven't forgotten any edge cases.

Model checkingFault injectionDeterministic simulation (DST)
Runs a model, not your code.Runs the real system and breaks it for real.Runs your actual code, with all nondeterminism replaced by mocks.
TLA+, Gallina, FizzBeeChaos Monkey, JepsenFoundationDB/Flow, TigerBeetle, FrostDB, MadSim, Antithesis
Systematic state exploration.Realistic, but coarse control.Systematic and replayable.

7.1 Model checking

Model checkers verify that INVARIANTS HOLD ACROSS ALL OF AN ALGORITHM'S STATES by systematically trying all the things that could happen. Specifications are written in a purpose-built language (TLA+, Gallina, FizzBee) that lets you focus on behavior without worrying about implementation details.

The honest limitations:

Model checking CAN'T ACTUALLY PROVE that invariants hold for every possible state, since most real-world algorithms have AN INFINITE STATE SPACE. True verification would require a formal proof, which is typically more difficult than running a model checker. Instead, model checkers encourage you to REDUCE the model to an approximation that can be fully verified, or to LIMIT the execution to an upper bound (e.g. a maximum number of messages). Any bugs occurring only with longer executions would then NOT BE FOUND.

They also DON'T RUN YOUR ACTUAL CODE — which makes state-space exploration tractable but RISKS THAT YOUR SPECIFICATION AND YOUR IMPLEMENTATION GO OUT OF SYNC.

Track record: CockroachDB, TiDB, Kafka and many others use model specifications. Using TLA+, researchers demonstrated the potential for DATA LOSS in viewstamped replication caused by AMBIGUITY IN THE PROSE DESCRIPTION of the algorithm.

7.2 Fault injection

Inject faults into a running system's environment and see how it behaves: network failures, machine crashes, disk corruption, paused processes.

Mechanics: deploy the system alongside fault injection COORDINATORS (deciding what faults to execute and when) and SCRIPTS (injecting failures into individual nodes). Tools: kill to pause or kill a process, umount to unmount a disk, firewall settings to disrupt network connections.

The myriad of tools required make fault injection tests CUMBERSOME TO WRITE. It's common to adopt a framework like JEPSEN, which comes with integrations for various operating systems and many prebuilt fault injectors. JEPSEN HAS BEEN REMARKABLY EFFECTIVE AT FINDING CRITICAL BUGS IN MANY WIDELY USED SYSTEMS.

(Production fault injection = chaos engineering, popularized by Netflix's Chaos Monkey.)

7.3 Deterministic simulation testing (DST)

DST uses a similar state-space exploration process to a model checker, BUT IT TESTS YOUR ACTUAL CODE, NOT A MODEL.

Network communication, I/O, and clock timing are all replaced with MOCKS that allow the simulator to CONTROL THE EXACT ORDER in which things happen. This lets it explore MANY MORE SITUATIONS than handwritten tests or fault injection.

If a test fails, IT CAN BE RERUN, since the simulator knows the exact order of operations that triggered the failure — in contrast to fault injection, which does NOT have such fine-grained control.

Three strategies for making code deterministic:

LevelHowExamples
ApplicationBuilt from the ground up to execute deterministically. FoundationDB uses Flow, an async communication library providing a point to inject a deterministic network simulation. TigerBeetle models state as a state machine with all mutations in a SINGLE EVENT LOOP, plus mock deterministic clocksFoundationDB, TigerBeetle
RuntimeA single-threaded runtime forces all asynchronous code to run sequentially. FrostDB patches Go's runtime to execute goroutines sequentially. Rust's MadSim provides deterministic implementations of Tokio's async API, Amazon's S3 library, Kafka's Rust library — applications swap in deterministic libraries WITHOUT CHANGING THEIR CODEFrostDB, MadSim
MachineA custom HYPERVISOR replaces normally nondeterministic operations with deterministic ones — everything from clocks to network and storage. Then developers run their entire distributed system in containers within the hypervisor and get a COMPLETELY DETERMINISTIC DISTRIBUTED SYSTEMAntithesis

Two bonus advantages beyond replayability:

  • Antithesis BRANCHES a test execution into multiple subexecutions when it discovers less common behavior, exploring many paths
  • Mocked clocks let tests run FASTER THAN WALL CLOCK TIME. TigerBeetle's time abstraction simulates network latency and timeouts WITHOUT ACTUALLY TAKING THE FULL LENGTH OF TIME — exploring more code paths faster

7.4 The power of determinism — the book's synthesis

NONDETERMINISM IS AT THE CORE OF ALL THE DISTRIBUTED SYSTEMS CHALLENGES in this chapter: concurrency, network delay, process pauses, clock jumps, and crashes all happen in unpredictable ways that vary from one run to the next.

Conversely, IF YOU CAN MAKE A SYSTEM DETERMINISTIC, THAT CAN HUGELY SIMPLIFY THINGS. Making things deterministic is a simple but powerful idea that ARISES AGAIN AND AGAIN in distributed system design.

Where determinism has already appeared in this book:

ChapterMechanismWhat determinism buys
Ch 3Event sourcingdeterministically Replay a log of events to reconstruct derived materialized views
Ch 5Workflow enginesdurable execution requires workflow definitions to be deterministic
Ch 6Statement-based replicationreplicating by re-executing statements requires determinism (NOW(), RAND() break it)
Ch 8Serial execution with stored proceduresVoltDB replicates by running the Same stored procedure on each replica
Ch 9Deterministic simulation testingreplay a failing execution exactly
Ch 10State machine replicationreplicate by independently executing the same sequence of deterministic transactions

⚠️ But making code FULLY deterministic requires care. Even once you have removed all concurrency and replaced I/O, network, clocks, and RNGs with deterministic simulations, elements of nondeterminism may remain: in some languages, THE ORDER IN WHICH YOU ITERATE OVER A HASH TABLE may be nondeterministic. WHETHER YOU RUN INTO A RESOURCE LIMIT (memory allocation failure, stack overflow) is also nondeterministic.


On this page