13 min read

The 3 AM Nightmare: Why We Bet Big on Formal Verification for Our Distributed KV Store

Betting Big on Formal Verification for Distributed KV Stores

It’s 3:14 AM. Your phone buzzes. PagerDuty. The alert isn’t a simple CPU spike or a failed disk. It’s the dreaded split-brain scenario: two nodes, both convinced they are the sole leader of a shard, both accepting writes. Your key-value store—the beating heart of your company's infrastructure—is silently corrupting data at 50,000 writes per second.

You do the math. The last known good backup was six hours ago. The blast radius? Every user session, every financial transaction, every cached authorization token. You’re not just looking at a rollback; you’re looking at a forensic nightmare.

This isn't a hypothetical. It’s the ghost that haunts every distributed systems engineer. We’ve all been there—staring at a stack trace of a Raft or Paxos implementation, wondering how a protocol that looks so elegant on paper could fail so spectacularly in production.

Today, we’re going to talk about how we moved beyond testing our consensus protocols and started proving them. This is the story of how we integrated Formal Verification into the CI/CD pipeline of a large-scale, multi-tenant key-value store. We’re not talking about unit tests. We’re not talking about Jepsen (though we love those folks). We’re talking about mathematically proving that our system cannot exhibit a split-brain, even under adversarial network conditions and arbitrary hardware failures.

Let’s dive into the deep end. Put on your math hats and your Ops boots.


The Hype vs. The Reality: Why Now?

If you’ve been on Hacker News in the last two years, you’ve seen the buzzwords: TLA+, Coq, Isabelle, P, PlusCal. Amazon used TLA+ to find bugs in S3 and DynamoDB. Microsoft used it for Azure Cosmos DB. The "hype" is real, but the context is often lost in the Twitter threads.

The hype gained traction because distributed systems have hit a complexity ceiling. We can’t just “test our way” to correctness anymore. A modern KV store like etcd, TiKV, or CockroachDB has a state space so vast that a comprehensive integration test would take longer than the lifespan of the universe to run.

The actual technical substance behind the hype is this: Model Checking. Instead of running code, we build a mathematical model of the system—states, transitions, invariants—and then exhaustively explore every possible state within a bounded scope. If the invariant "Two leaders cannot exist in the same term" holds for all explored states, we gain a level of confidence that no amount of fuzzing can provide.

But here’s the catch: Formal verification doesn’t prove your code is correct. It proves your design is correct. The gap between the model and the implementation is where the bugs live. Our journey was about closing that gap in a large-scale KV store handling millions of ops per second.


The Architecture: Where Does Consensus Fit?

Before we get into the math, let’s ground ourselves in the architecture of a large-scale KV store. Imagine a cluster of 500 nodes across three regions. Data is sharded into Raft groups—each group is a consensus cluster of 3 or 5 nodes.

  • Shard A: Replicas in us-east-1, us-west-2, eu-west-1.
  • Shard B: Replicas in ap-south-1, ap-northeast-1, us-west-2.

Each shard runs a Raft instance. The Raft leader handles all writes for that shard. The followers replicate the log. A global metadata service (itself a Raft group) manages shard assignments.

The critical invariant is simple: For any given term, there is at most one leader per shard. If that breaks, data diverges. Game over.

Now, let’s talk about the failure modes that keep us up at night:

  1. Network Partitions: A leader in us-east-1 is partitioned from the followers in us-west-2 and eu-west-1. The followers elect a new leader. The old leader, unaware, still thinks it’s the leader. Can it accept writes? (Spoiler: In Raft, it shouldn't, but a bug in the implementation of the election timeout or the log matching property can cause this.)
  2. Byzantine Failures: A node doesn’t just crash; it sends malformed messages or lies. (Raft is not Byzantine fault-tolerant, but we need to ensure our implementation doesn’t accidentally become vulnerable to simple message corruption.)
  3. Clock Skew: The election timeout is based on a heartbeat interval. If clocks drift, a follower might start an election too early or too late.
  4. Log Truncation Bugs: The leader crashes after appending an entry but before replicating it. The new leader must not commit that entry. A bug in the log matching property can cause committed entries to be lost.

These are not theoretical. We’ve seen all of these in production, albeit rarely. The cost of a single incident is millions of dollars in lost revenue and trust.


Enter Formal Verification: The Toolchain

We chose TLA+ (Temporal Logic of Actions) as our specification language. Why? It’s designed for concurrent and distributed systems. It’s backed by a powerful model checker (TLC) and has a rich ecosystem. Plus, it’s used by Amazon, Microsoft, and many others—so we’re not alone.

But we didn’t just write a TLA+ spec and call it a day. We built a pipeline that integrates formal verification into our development lifecycle.

Step 1: Modeling the Protocol

We started by writing a high-level TLA+ spec of our Raft variant. Here’s a simplified snippet of the Election module:

VARIABLES currentTerm, state, votedFor, log, commitIndex

\* Invariant: At most one leader per term
AtMostOneLeader ==
    \A t \in Term:
        Cardinality({s \in Server: state[s] = "Leader" /\ currentTerm[s] = t}) <= 1

\* Election action
StartElection(s) ==
    /\ state[s] = "Follower"
    /\ currentTerm' = [currentTerm EXCEPT ![s] = @ + 1]
    /\ votedFor' = [votedFor EXCEPT ![s] = s]
    /\ state' = [state EXCEPT ![s] = "Candidate"]
    /\ UNCHANGED <<log, commitIndex>>

This is a state machine. TLC will explore every possible interleaving of StartElection, RequestVote, AppendEntries, and Crash actions. If AtMostOneLeader is violated, TLC gives us a counterexample—a trace of states leading to the bug.

Step 2: Bounded Model Checking

The state space is infinite if we allow unbounded terms, logs, and message queues. So we bound the model:

  • 3 servers
  • 2 terms
  • 3 log entries
  • 2 messages in flight

This is not a proof for all possible systems, but it’s a proof for a representative subset. If a bug exists in the unbounded system, it often manifests in the bounded system. We run this in CI on every PR.

Step 3: The Gap Between Spec and Code

This is where most teams fail. You have a beautiful TLA+ spec, but your Go/C++/Rust code is a mess of channels, mutexes, and syscalls. How do you ensure the code matches the spec?

We used a combination of trace validation and runtime verification.

Trace Validation: We instrumented our Raft implementation to emit a trace of events (e.g., StartElection, VoteGranted, AppendEntries). We then wrote a TLA+ spec that consumes these traces and checks if they are valid according to the protocol. If the code emits an event that the spec says is impossible, we have a bug.

// Go code emitting a trace event
func (r *Raft) StartElection() {
    r.trace.Emit("StartElection", r.id, r.currentTerm)
    // ... actual logic
}

The TLA+ trace validator then checks:

TraceInvariant ==
    \A i \in 1..Len(trace):
        LET event == trace[i] IN
        \/ event.type = "StartElection" /\ IsValidStartElection(event)
        \/ event.type = "VoteGranted" /\ IsValidVoteGranted(event)
        \* etc.

This catches implementation bugs that deviate from the spec. It’s not a full proof, but it’s a powerful safety net.

Step 4: Fuzzing the Model

We also use fuzzing to generate random traces and feed them to the TLA+ validator. This is similar to Jepsen, but instead of checking for linearizability, we check for conformance to the formal spec. We found a bug where a follower would grant a vote to a candidate with a shorter log—a violation of the Log Matching property. This was in a code path that was only triggered when a specific network partition healed in a specific order. Traditional testing would have never caught it.


The Infrastructure: Compute Scale and Cost

Formal verification is computationally expensive. TLC is a model checker that explores state spaces. For our 3-server, 2-term model, the state space is around 10^8 states. TLC can explore about 10^5 states per second on a single core. That’s 1000 seconds—about 17 minutes.

We run this in CI. 17 minutes is too long for a PR check. So we parallelized.

  • Distributed TLC: We run TLC on a Kubernetes cluster with 32 cores. The state space is partitioned by the initial state. We use the -workers flag to spawn multiple TLC instances. This brings the time down to ~2 minutes.
  • Incremental Checking: We cache the state space for unchanged modules. If a PR only changes the logging module, we don’t re-check the election module.
  • Nightly Deep Dives: For the full model (5 servers, 3 terms, 5 log entries), we run a nightly job on a 128-core machine. This takes ~8 hours and explores 10^12 states. We’ve found bugs that only appear with 5 servers—specifically, a liveness bug where a leader could be elected but never commit an entry due to a subtle interaction between the heartbeat and the log replication.

The cost? A few thousand dollars a month in cloud compute. The ROI? A single prevented outage saves millions.


The Engineering Curiosities: What We Learned

1. The Model is Not the Code (But It’s Close)

We initially thought we could generate code from the TLA+ spec. We tried PlusCal (a TLA+ algorithm language) and then used a transpiler to Go. It was a disaster. The generated code was unreadable, slow, and didn’t handle real-world concerns like batching, pipelining, and backpressure.

Instead, we treated the TLA+ spec as a reference oracle. We wrote the Go code by hand, but we validated every state transition against the spec. This gave us the best of both worlds: performance and correctness.

2. The Invariant is the Product

The most valuable artifact we produced was not the code, but the invariants. We defined invariants for:

  • Election Safety: At most one leader per term.
  • Log Matching: If two logs contain an entry with the same index and term, then the logs are identical up to that index.
  • Leader Completeness: If a log entry is committed in a given term, then that entry will be present in the logs of the leaders for all higher-numbered terms.
  • State Machine Safety: If a server has applied a log entry at a given index, no other server will ever apply a different log entry for the same index.

These invariants are now part of our SLO documentation. When we onboard new engineers, we don’t just tell them to read the Raft paper. We tell them to read our TLA+ invariants. It’s a contract.

3. The Bug That Almost Shipped

We found a bug in our pre-vote implementation. Pre-vote is an optimization to prevent a partitioned server from disrupting the cluster by incrementing its term. Our TLA+ model showed that in a specific scenario—a server that was partitioned, then reconnected, then immediately started a pre-vote—the pre-vote could be granted by a follower that had already voted for a different candidate in the same term. This would cause the original candidate to lose the election, leading to a liveness violation (no leader elected for an extended period).

This bug was in a code path that was only triggered when the network partition healed in a specific order. Traditional testing would have never caught it. We fixed it by adding a check in the pre-vote handler to ensure the follower hasn’t already voted in the current term.

4. The Cost of Correctness

Formal verification is not free. It requires:

  • Expertise: You need engineers who understand temporal logic and model checking.
  • Time: Writing a good spec takes weeks, not days.
  • Tooling: TLC is not a friendly tool. The error messages are cryptic. We built a custom visualization tool that renders the counterexample trace as a sequence diagram.
  • Cultural Shift: You need to convince your team that spending a week on a spec is worth it. The payoff is invisible until it prevents a catastrophic bug.

We invested in a formal methods team of 3 engineers. They don’t write production code. They write specs, build tooling, and train other engineers. It’s a significant investment, but it’s paid off in ways we can’t fully quantify.


The Code: A Glimpse into Our CI Pipeline

Here’s a simplified version of our GitHub Actions workflow for formal verification:

name: Formal Verification

on:
  pull_request:
    paths:
      - 'raft/**'
      - 'specs/**'

jobs:
  tlc:
    runs-on: ubuntu-latest
    container:
      image: our-tlaplus-image:latest
    steps:
      - uses: actions/checkout@v2
      - name: Run TLC on Election Spec
        run: |
          tlc -workers 8 -config election.cfg election.tla
      - name: Run Trace Validator
        run: |
          ./trace-validator --spec election.tla --trace ./traces/election_trace.ndjson
      - name: Upload Counterexample
        if: failure()
        uses: actions/upload-artifact@v2
        with:
          name: counterexample
          path: ./counterexample/

When a PR fails the formal verification check, it blocks the merge. The developer gets a counterexample trace and a link to a visualization. This is not a suggestion; it’s a hard gate.


The Future: Beyond Raft

We’re now looking at formal verification for other parts of the KV store:

  • Transaction Isolation: We’re modeling our MVCC (Multi-Version Concurrency Control) implementation in TLA+ to prove that our snapshot isolation level never allows dirty reads or lost updates.
  • Shard Rebalancing: We’re modeling the algorithm that moves shards between nodes. The invariant is that no shard is ever served by two nodes at the same time, and no data is lost during a move.
  • Byzantine Fault Tolerance: We’re exploring PBFT and HotStuff for a future version of the KV store that needs to tolerate malicious nodes. The state space is larger, but the invariants are similar.

We’re also experimenting with Coq for proving the correctness of our CRDT (Conflict-free Replicated Data Type) implementations. Coq is a proof assistant, not a model checker. It allows us to write machine-checked proofs of mathematical properties. It’s a different beast, but the goal is the same: mathematical certainty.


The Takeaway: You Don’t Need to Be a Mathematician

You don’t need a PhD in formal methods to start. You can begin with a simple TLA+ spec for a critical component. You can use PlusCal to write algorithms in a more familiar syntax. You can use P (from Microsoft) which is more like a programming language. You can even use Alloy for lightweight modeling.

The key is to start small. Pick one invariant—like "no two leaders"—and model it. You’ll be surprised at what you find.

Formal verification is not a silver bullet. It’s a tool. It’s a powerful tool that, when combined with rigorous testing, fuzzing, and observability, can give you a level of confidence that your distributed system will behave correctly even in the face of adversarial conditions.

So, the next time you’re paged at 3 AM, you might still have to fix a bug. But it won’t be a split-brain. It won’t be a data corruption. It will be something mundane—like a disk filling up. And you’ll smile, because you know that the core of your system is mathematically sound.

Now, go write some specs. Your future self will thank you.


This article is based on our experiences building a multi-tenant KV store that handles over 10 million operations per second across 500 nodes. The TLA+ specs and tooling are open-sourced under the Apache 2.0 license. Check out our GitHub repo for more details.


More to explore

Keep diving in