Skip to main content

I Built Atlas, Then Started Breaking It on Purpose

What a distributed systems lab taught me about Raft, network partitions, split votes, and the uncomfortable reality of making computers agree with each other.

By Joseph Gitau Chege9 min read

There is something deeply unreasonable about distributed systems.

You take a perfectly functional computer, connect it to a few other perfectly functional computers, and suddenly nobody knows who's in charge.

One machine thinks it's the leader. Another thinks the leader is dead. A third is sitting somewhere wondering why everyone is shouting about terms and heartbeats.

And somehow, you're expected to make them agree on what happened.

I built Atlas because I wanted to understand that problem beyond reading about it, drawing architecture diagrams, and pretending I understood Raft because I could explain leader election in an interview.

I wanted to build the thing.

Then I wanted to break it.

Not just kill a process and call it chaos engineering. I wanted to introduce network partitions, create competing leadership elections, interrupt communication, and observe what happened when the system was no longer operating under the assumption that everything was fine.

Because honestly, a distributed system that works only when everything works isn't particularly interesting.

It's just a collection of computers having a good day.

So, what exactly is Atlas?

Atlas is my distributed systems laboratory, built around the Raft consensus algorithm.

The goal was to implement and experiment with the mechanisms that allow multiple nodes to maintain a consistent view of replicated state despite failures.

The project includes:

  • Raft leader election and term management
  • Heartbeats and follower coordination
  • Replicated log management
  • A replicated key-value store
  • Virtual nodes (vnodes)
  • Dead-letter queue (DLQ) handling
  • Failure injection and chaos testing

The interesting part wasn't implementing each feature independently. It was understanding how they interact when the network starts behaving like it has personal problems.

Atlas: conceptual architecture

      Client requests
        Read / Write
             │
             ▼
      ┌──────────────────┐
      │   Raft cluster   │
      │  ┌────────────┐  │
      │  │  Node A    │  │  Leader
      │  ├────────────┤  │
      │  │  Node B    │  │  Follower
      │  ├────────────┤  │
      │  │  Node C    │  │  Follower
      │  └────────────┘  │
      │ Election · HB ·  │
      │ Log replication · │
      │ Term tracking     │
      └──────────────────┘
             │
             ▼
   Replicated key-value state
   State machine and operation outcomes

The diagram is simplified. The actual implementation has to deal with message ordering, log consistency, node state transitions, and failures that don't politely announce themselves.

And that's where things get fun.

1. The first lie distributed systems tell you: the leader is alive

In a single-process application, checking whether something is running is relatively straightforward.

In a distributed system, determining whether another machine is alive is a completely different problem.

A node stops receiving heartbeats.

Is the leader dead?

Did the network fail?

Is the machine overloaded?

Did a packet get delayed?

Is the connection broken only in one direction?

Or is the leader perfectly healthy and simply unable to communicate with you?

You don't necessarily know.

Raft handles this through timeouts and elections. Followers expect heartbeats from the leader. If those heartbeats stop arriving within the election timeout, a follower can transition into a candidate and attempt to become leader.

That sounds simple enough until multiple nodes reach that conclusion at approximately the same time.

Which brings us to the next problem.

2. Split votes: democracy, but make it a networking problem

Imagine three nodes.

The leader disappears.

Both remaining followers independently decide that it's election time.

Both become candidates.

Both request votes.

Neither receives enough votes to establish leadership.

Congratulations. You have successfully built a cluster that cannot agree on who should be in charge.

This is a split vote.

Raft addresses this through randomized election timeouts, reducing the likelihood that candidates repeatedly initiate elections simultaneously.

But the important distinction is that randomization reduces the likelihood of repeated collisions. It doesn't eliminate the need to handle them correctly.

For Atlas, this meant treating elections as state transitions rather than simply assigning a leader variable and hoping the network cooperated.

A candidate needs to:

  • Increment its current term
  • Vote for itself
  • Request votes from other nodes
  • Become leader only after obtaining a majority
  • Return to follower state when observing a higher term

That last point matters considerably.

A node cannot simply cling to leadership because it was elected previously.

Leadership is conditional on the protocol's current state.

And terms are how Raft keeps track of that state across changing leadership.

3. Terms: because stale leadership is a real problem

One of the more interesting parts of implementing Raft is understanding that leadership isn't permanent.

A node might have been elected leader in term 4.

Then the cluster loses connectivity.

Another group of nodes establishes a new leader in term 5.

When communication resumes, the old leader cannot continue behaving as though nothing happened.

It must recognize the newer term and step down.

This is one of the reasons term tracking is so important.

Without it, an isolated node could continue operating under an outdated understanding of cluster leadership.

The distinction between having been elected leader and being the current leader is fundamental.

It also made me appreciate why distributed systems need explicit protocol state rather than relying on assumptions that are obvious in a single-process application.

In a normal application, you might ask:

In a distributed system, you have to ask:

Slightly more annoying.

Significantly more important.

4. Network partitions: the part where everything gets personal

This is where failure injection becomes genuinely useful.

A network partition separates nodes into groups that cannot communicate with one another.

For example, consider a five-node cluster:

  ┌─────────────────────┐   ┌──────────────┐
  │  Majority partition │   │   Minority   │
  │      3 nodes        │   │    2 nodes   │
  │                     │   │              │
  │ Can form a quorum   │   │ Cannot       │
  │ and elect a leader. │   │ independently│
  │                     │   │ establish a  │
  │                     │   │ majority.    │
  └─────────────────────┘   └──────────────┘
        no communication

The partition prevents communication between the two groups. This illustrates expected Raft quorum behavior, not a claim about a particular Atlas test run.

The majority side can continue making progress, assuming the relevant nodes remain healthy and can communicate.

The minority side cannot independently commit new entries because it lacks a majority.

This is an important trade-off.

The system sacrifices availability for certain operations in the minority partition to preserve its consistency guarantees.

And this is where I started appreciating the difference between building something that responds to requests and building something that responds correctly.

Returning a successful HTTP response is easy.

Returning a successful response when the system has actually committed the operation according to its consistency rules is a different conversation.

The uncomfortable question

What happens when the partition heals?

The system must reconcile its state.

Nodes need to recognize the current term, restore communication, and bring their logs into alignment through the Raft replication process.

A stale leader must not continue acting as leader.

Uncommitted entries must not be mistaken for committed state.

And the replicated state machine must not apply operations in an inconsistent order.

These are precisely the sorts of properties I wanted Atlas's failure-injection tests to exercise.

5. Failure injection: stop trusting the happy path

Most application testing follows a fairly predictable pattern.

Send a request.

Check the response.

Verify the database.

Move on.

That is useful, but it doesn't tell you much about how a distributed system behaves when its assumptions collapse.

With Atlas, I wanted to test the conditions that ordinary functional tests tend to avoid.

Failure scenarios I designed Atlas to exercise

The following are test scenarios and expected protocol properties, not a claim that every scenario produced a verified result in a particular run.

  • Network partition — Separate nodes into communicating groups and examine quorum behavior, leadership, and recovery.
  • Split vote — Trigger competing elections and observe whether the cluster can establish leadership through subsequent election rounds.
  • Delayed or interrupted communication — Exercise timeout behavior and test whether stale protocol messages are handled correctly.
  • Leader failure — Observe election recovery and whether the remaining nodes can establish a new leader.
  • Partition recovery — Restore connectivity and inspect term convergence, log reconciliation, and state-machine consistency.

The goal wasn't to make the cluster survive every possible failure.

That's not a realistic engineering objective.

The goal was to understand which guarantees the protocol provides, under which assumptions, and how the implementation behaves when those assumptions are violated.

A distributed system doesn't become reliable because you added retries and a health-check endpoint.

Reliability comes from defining what correctness means and testing whether the system preserves it under failure.

6. The replicated key-value store made everything more interesting

Implementing consensus in isolation is useful, but I wanted Atlas to do something with it.

That's where the replicated key-value store comes in.

A write operation isn't merely a mutation on whichever node happens to receive the request.

It needs to participate in the replicated log and follow the protocol's commitment rules before being applied to the state machine.

This introduces a distinction that is easy to overlook:

An operation being received is not the same as an operation being committed.

And an operation being present in a node's local log is not necessarily proof that it has been committed by the cluster.

That distinction matters when handling failures.

If a client times out while waiting for a write, the client may not know whether the operation committed.

If it blindly retries, the application may need additional mechanisms to prevent duplicate business effects.

This is one of the reasons the Atlas work connects directly to the integration problems I encounter in enterprise software.

The underlying question is the same:

7. What I took away from building Atlas

I went into Atlas wanting to understand Raft.

I came out with a much greater appreciation for the difference between implementing a protocol and building a system that behaves correctly under imperfect conditions.

A few things stuck with me.

Failure is not an exceptional state in distributed systems

It's part of the operating environment.

Nodes crash. Networks partition. Messages arrive late. Processes disagree about what they know.

Your architecture has to account for those conditions from the beginning.

Correctness is more important than appearing available

A system that accepts writes from two isolated groups and later discovers conflicting histories has a much bigger problem than one that temporarily refuses writes because it cannot establish a quorum.

Observability is part of the engineering work

When something goes wrong in a distributed system, you need to understand what each node believed was happening.

Terms, roles, election transitions, log indexes, commit indexes, and replication progress aren't just debugging details.

They're the evidence you need to understand the system.

Reading the algorithm isn't the same as implementing it

The Raft paper explains the protocol.

Writing the implementation forces you to confront the details.

What happens when a message arrives after a timeout? What if a node receives a higher term while performing another operation? What does a follower do when its log conflicts with the leader's?

Those questions become much harder to ignore when you're responsible for the code.

The part I still find fascinating

The funny thing about Atlas is that the more I learned about distributed systems, the less interested I became in pretending that failures could be eliminated.

They can't.

What we can do is design systems that understand their failure modes, preserve their guarantees where possible, and recover predictably when conditions improve.

That's a much more interesting engineering problem than making a demo work.

And frankly, it's a lot more fun to deliberately break five computers than to spend another afternoon wondering why an API returned a 500.

Although, knowing my luck, the API will probably return a 500 anyway.

Start a conversation

Reading beats doing — but only up to a point.

Where a post ends is where the real environment begins. If any of this sounds like your situation, the first conversation is free.

or email joseph.gitau.c@gmail.com