Skip to main content
Testing

Deterministic Simulation Testing: Fault Injection Methodology in TigerBeetle and FoundationDB

9 min read
LD
Lucio Durán
Engineering Manager & AI Solutions Architect
Also available in: Español, Italiano

What Makes Testing Distributed Systems So Hard

The fundamental problem is non-determinism. A distributed system has multiple sources of it:

  1. Thread scheduling: The OS decides which thread runs when
  2. Network timing: Packets arrive in unpredictable order with variable latency
  3. Disk I/O timing: Reads and writes complete at variable speeds
  4. Clock skew: System clocks on different machines drift independently
  5. Failure timing: Crashes, network partitions, and disk errors happen at arbitrary moments

A bug that requires a specific combination of these non-deterministic events might take billions of test runs to hit with random testing. And even if you hit it, you can't reproduce it because you don't control the OS scheduler, the network, or the disk.

Traditional integration tests try to handle this with sleeps and retries:

# This is how most people test distributed systems (don't do this)
def test_leader_election():
 cluster.start(3)
 time.sleep(5) # "wait for election to complete" 🤞
 leader = cluster.get_leader()
 assert leader is not None

 cluster.kill(leader)
 time.sleep(10) # "wait for re-election" 🤞🤞
 new_leader = cluster.get_leader()
 assert new_leader is not None
 assert new_leader != leader

This test passes 99% of the time and misses every interesting bug. The bugs live in the timing edges — what happens when the election timeout fires at the exact same instant as a heartbeat? What happens when a network partition heals during a leadership transfer? These tests can't explore that space.

The FoundationDB Methodology

FoundationDB pioneered simulation testing in production databases. Their key insight: if you control all sources of non-determinism, you can explore the entire state space systematically.

The architecture has three pillars:

1. Single-Threaded, Event-Driven Core

The entire database runs in a single thread (or multiple deterministic threads with explicit scheduling). There is no OS-level concurrency. All asynchronous operations are modeled as events in a priority queue:

Event Queue (ordered by simulated time):
 t=100ms: Node A receives AppendEntries from Node B
 t=102ms: Node C election timeout fires
 t=103ms: Disk write on Node A completes
 t=105ms: Node B heartbeat timer fires
 t=108ms: Network partition begins (injected fault)

The simulator pops events in order, executes them, and any resulting events go back into the queue. Because the queue is deterministic (same seed → same initial state → same event order), the execution is perfectly reproducible.

2. Abstracted I/O Layer

Every interaction with the outside world goes through an interface:

// FoundationDB's flow runtime abstracts everything
class INetwork {
 virtual Future<Void> connect(NetworkAddress addr) = 0;
 virtual Future<Void> send(Connection conn, Bytes data) = 0;
 virtual Future<Bytes> receive(Connection conn) = 0;
};

class IDisk {
 virtual Future<Void> write(FileHandle fh, Offset off, Bytes data) = 0;
 virtual Future<Bytes> read(FileHandle fh, Offset off, Size len) = 0;
 virtual Future<Void> sync(FileHandle fh) = 0;
};

class IClock {
 virtual double now() = 0;
 virtual Future<Void> delay(double seconds) = 0;
};

In production, these call real syscalls. In simulation, they're replaced with fakes that introduce controlled delays, failures, and reorderings.

3. Seed-Based Fault Injection

Given a random seed, the simulator decides:

  • When to fire timers (with jitter)
  • When to deliver or drop network packets
  • When to inject disk errors
  • When to simulate process crashes and restarts
  • When to create and heal network partitions
# Pseudocode for the simulation fault injector
class FaultInjector:
 def __init__(self, seed):
 self.rng = Random(seed)

 def should_drop_packet(self) -> bool:
 return self.rng.random() < 0.01 # 1% packet loss

 def network_delay_ms(self) -> float:
 return self.rng.exponential(scale=5.0) # avg 5ms delay

 def should_inject_partition(self, time) -> bool:
 # Partitions happen ~once per simulated hour
 return self.rng.random() < (1.0 / 3600000.0) * time_step_ms

 def should_crash_node(self, node_id) -> bool:
 return self.rng.random() < 0.0001 # rare but devastating

 def disk_write_should_fail(self) -> bool:
 return self.rng.random() < 0.001 # occasional disk errors

TigerBeetle's Approach: VOPR

TigerBeetle (a financial transactions database written in Zig) took FoundationDB's ideas and built VOPR — the Viewstamped Operation Replication simulator. What makes TigerBeetle's approach distinctive is how deeply simulation is integrated into the codebase.

In TigerBeetle, the I/O abstraction isn't an afterthought — it's the foundation. Every syscall goes through IO, which has two implementations:

// Production I/O — real io_uring syscalls
const IO = if (is_simulation)
 @import("io_simulation.zig").IO
else
 @import("io_uring.zig").IO;

The simulation IO replaces io_uring with a deterministic event loop:

// Simplified simulation I/O
pub const IO = struct {
 prng: std.rand.DefaultPrng,
 event_queue: PriorityQueue(Event),
 tick: u64,

 pub fn submit_write(self: *IO, buffer: []const u8, offset: u64) Completion {
 // In simulation, writes complete after a random delay
 const delay = self.prng.random_int(u64) % 10 + 1; // 1-10 ticks
 const completion = Completion{
 .result = if (self.should_inject_fault())
 error.DiskError
 else
 .success,
 };
 self.event_queue.insert(.{
 .tick = self.tick + delay,
 .completion = completion,
 });
 return completion;
 }

 fn should_inject_fault(self: *IO) bool {
 return self.prng.random_int(u64) % 1000 == 0; // 0.1% failure rate
 }
};

The VOPR runner executes the entire TigerBeetle cluster — multiple replicas, clients, and the network — in a single process:

// VOPR runs a full cluster simulation
pub fn run(seed: u64) !void {
 var prng = std.rand.DefaultPrng.init(seed);

 // Create simulated cluster
 var cluster = Cluster.init(&prng, .{
 .replica_count = 6,
 .standby_count = 2,
 .client_count = 3,
 });

 // Run simulation for N ticks
 var tick: u64 = 0;
 while (tick < 10_000_000) : (tick += 1) {
 // Maybe inject a fault
 if (prng.random_int(u64) % 100 == 0) {
 const fault = random_fault(&prng);
 cluster.inject_fault(fault);
 }

 // Advance all components by one tick
 cluster.tick();

 // Check invariants after every tick
 cluster.check_invariants();
 }
}

The check_invariants() call is key. After every single tick of the simulation — not just at the end — TigerBeetle verifies that all safety properties hold: no double-spending, no lost transactions, no committed state divergence. This is what catches the subtle interleaving bugs that only manifest mid-execution.

Building Your Own Simulation Harness

This methodology can be applied to smaller-scale projects as well. Here is the practical architecture:

Step 1: Define Your Determinism Boundary

List every source of non-determinism in your system:

Sources of non-determinism:
 ✅ System clock → wrap with injectable Clock interface
 ✅ Random numbers → use seeded PRNG everywhere
 ✅ Network I/O → wrap with injectable Network interface
 ✅ Disk I/O → wrap with injectable Storage interface
 ✅ Thread scheduling → use single-threaded event loop in simulation
 ❌ Memory allocation order → usually deterministic, but watch out for hash maps
 ❌ Signal handling → disable in simulation

Step 2: Build the Abstract I/O Layer

// Go example: abstract I/O for a Raft implementation
type Clock interface {
 Now() time.Time
 After(d time.Duration) <-chan time.Time
}

type Network interface {
 Send(to NodeID, msg Message) error
 Receive() <-chan Message
}

type Storage interface {
 Append(entries []LogEntry) error
 ReadRange(start, end uint64) ([]LogEntry, error)
 Sync() error
}

// Production implementations use real I/O
type RealClock struct{}
func (c *RealClock) Now() time.Time { return time.Now() }

// Simulation implementations are deterministic
type SimClock struct {
 current time.Time
}
func (c *SimClock) Now() time.Time { return c.current }
func (c *SimClock) Advance(d time.Duration) { c.current = c.current.Add(d) }

Step 3: The Simulation Runner

type Simulator struct {
 seed int64
 rng *rand.Rand
 clock *SimClock
 network *SimNetwork
 nodes []*RaftNode
 events *PriorityQueue
}

func (s *Simulator) Run(ticks int) error {
 for i := 0; i < ticks; i++ {
 // Advance clock
 s.clock.Advance(time.Millisecond)

 // Maybe inject faults
 s.maybeInjectFault()

 // Deliver pending network messages (with possible reordering)
 s.network.DeliverPending(s.rng)

 // Tick all nodes
 for _, node := range s.nodes {
 node.Tick()
 }

 // Check invariants
 if err := s.checkInvariants(); err != nil {
 return fmt.Errorf("invariant violated at tick %d (seed %d): %w",
 i, s.seed, err)
 }
 }
 return nil
}

func (s *Simulator) checkInvariants() error {
 // All committed entries must be identical across nodes
 committed := s.nodes[0].CommittedLog()
 for _, node := range s.nodes[1:] {
 nodeCommitted := node.CommittedLog()
 minLen := min(len(committed), len(nodeCommitted))
 for j := 0; j < minLen; j++ {
 if committed[j] != nodeCommitted[j] {
 return fmt.Errorf("log divergence at index %d: node0=%v node%d=%v",
 j, committed[j], node.ID, nodeCommitted[j])
 }
 }
 }
 return nil
}

Step 4: CI Integration

# Run 500 seeds on every PR, different range per CI shard
simulation-test:
 strategy:
 matrix:
 seed-range: ["0-100", "100-200", "200-300", "300-400", "400-500"]
 steps:
 - run: |
 START=$(echo ${{ matrix.seed-range }} | cut -d- -f1)
 END=$(echo ${{ matrix.seed-range }} | cut -d- -f2)
 for seed in $(seq $START $END); do
 echo "Running simulation with seed $seed"
 ./bin/simulation-test --seed=$seed --ticks=1000000
 done

Bugs This Methodology Has Found

In practice, simulation testing consistently finds bugs in these categories:

Stale leader writes: A node believes it's the leader but was deposed. It accepts a write that gets lost when the real leader's log overwrites it. Only happens when the deposition message is delayed by exactly the right amount.

Split-brain during partition healing: When a network partition heals, two sub-clusters have divergent state. The reconciliation protocol has a window where both think they have quorum. Simulation catches this by healing partitions at adversarial moments.

Disk write reordering: On real hardware, write(A); write(B); doesn't guarantee A is persisted before B. Simulation injects write reordering to find assumptions about persistence order.

Timer race conditions: When two timers fire at the same simulated instant, the order matters. Simulation explores both orderings for each seed.

A simulation harness can find complex interleaving bugs in seconds that would take weeks to discover through manual debugging. The contrast between weeks of manual investigation and seconds of simulation execution underscores the methodology's value.

For systems that need to be correct under failure — databases, consensus protocols, transaction processors, distributed locks — deterministic simulation testing is not optional. It is the only methodology that systematically finds the bugs that matter most: the ones that corrupt data.

deterministic-simulationtestingdistributed-systemstigerbeetlefoundationdbfault-injectioncorrectnesszig

Tools mentioned in this article

AWSTry AWS
DigitalOceanTry DigitalOcean
Disclosure: Some links in this article are affiliate links. If you sign up through them, I may earn a commission at no extra cost to you. I only recommend tools I personally use and trust.
Seguime