~ $ cd teaching/raft && cat README.md

Raft, and a recipe for concurrency

In the autumn of 2013 the distributed systems laboratory at the UOC set its students a consensus algorithm to implement: Raft, which was then a draft going round, a year away from being presented. I taught that laboratory. The assignment took three months, and my implementation is public: one Java class over the course's skeleton, dated October 2013, which was the reference the students' work was compared against, and was given back to them as the solution.

Raft was designed to be understandable, and it is. The hard part of the assignment is somewhere else. A server is doing four things at once — timing out, asking for votes, answering other servers' requests, replicating its log — every one of them reads and writes the same few fields, and between any two lines the network may hand it a message that makes it a different kind of server. That is a lot to ask of someone writing their first concurrent program.

So the students were given three hints, and the implementation follows them.

One: the timeouts are always on

The obvious way to write a timeout is to arm it, and to cancel it and arm it again whenever a message says it is not needed yet. Every one of those is a chance for a message and a timer to cross. Here the timers are started once, and they tick for as long as the server runs; each tick asks one question instead — has a leader been heard from since the last one? — and returns if it has:

timerQueue.schedule(electionTimeoutTask, electionTimeout, electionTimeout);
…
private void electionTimeout() {
    // abort timeout if leader has been seen
    if (seenLeader.getAndSet(false)) { return; }

It does a little work for nothing, and it has no race left to lose.

Two: a recipe for the shared state

The Java habit is to make the whole method synchronized and be safe. It is not always safe — a synchronized method is a monitor, and a student who knows mutexes and semaphores does not yet know what a monitor will do — and it makes long stretches of the program wait on each other. So the implementation follows a recipe simple enough to be followed by someone who cannot yet reason about interleavings, and still concurrent:

something changed1 · inside the guardcheck who you arecopy what you need2 · outside the guardcompute, wait, talkto the network3 · inside the guard againcheck nothing changedonly then writegive up quietlythe next timeoutwill try again
  1. One guard for all the state. Not a lock per field: one. Take it, check that you are still what you think you are — only leaders send heartbeats — and copy everything you are about to need into local variables that nothing else can touch.
  2. Let go before doing anything slow. Never hold the guard across the network. The remote call runs with the copies, on another thread, for as long as it takes, and the server goes on answering everyone else.
  3. Take the guard again, and trust nothing. The answer arrives in a world that has moved. Am I still the leader? Is it still the same term? Is this follower's index still where I left it? If any answer is no, drop the result and return. There is nothing to undo, because nothing was written.

Here it is in the leader's heartbeat, shortened but in the code's own words and with its own comments:

synchronized (GUARD) {
    // only leaders perform heartbeats
    if (state != RaftState.LEADER) return;

    // gather common info (from iteration to iteration may become rotten)
    term = persistentState.getCurrentTerm();
    prevLogIndex = nextIndexes.get(otherServer) - 1;
    prevLogTerm = persistentState.getTerm(prevLogIndex);
    entries = prevLogIndex > -1 ? persistentState.getLogEntries(prevLogIndex+1) : new ArrayList<LogEntry>();
    commitIndex = this.commitIndex;
}

// send the message (and listen the answer) in concurrent
executorQueue.execute(new Runnable() {
    public void run() {
        AppendEntriesResponse response = RMIsd.getInstance()
            .appendEntries(otherServer, term, leaderId, prevLogIndex, prevLogTerm, entries, commitIndex);

        // execute inside the guard, any sent data could be changed and must be reevaluated
        synchronized (GUARD) {
            // still leader?
            if (state != RaftState.LEADER) return;
            // term changed?
            if (term != persistentState.getCurrentTerm()) return;
            // prevLogIndex changed?
            if (nextIndexes.get(otherServer) - 1 != prevLogIndex) return;

            // … only now is anything written
        }
    }
});

That is the whole of it. The two comments that matter are the code's own: gather common info (from iteration to iteration may become rotten) going in, and execute inside the guard, any sent data could be changed and must be reevaluated coming back.

Three: the answer is handled where the question is asked

A server sends many requests and gets their answers back in any order. The usual Java of the time had a handler for each kind of answer, somewhere else, matching answers to questions; every attempt written that way was hard to follow and full of races. The idea that worked came from JavaScript: each request goes to an executor as a small anonymous task that sends it, waits, and deals with its own answer, with everything it needs already in its local variables. Here is the vote, as the election asks for it:

// request votes
for (final Host otherHost : otherServers) {
    executorQueue.execute(new Runnable() {
        public void run() {
            RequestVoteResponse response = RMIsd.getInstance()
                .requestVote(otherHost, term, candidateId, lastLogIndex, lastLogTerm);
            // not now or not me
            if (response.getTerm() != term || !response.isVoteGranted()) { checkReceivedTerm(term); return; }

            // who wons?
            int votes = voteCount.incrementAndGet();
            if (votes == minimumVoteCount) {
                synchronized (GUARD) {
                    if (persistentState.getCurrentTerm() != term || state != RaftState.CANDIDATE) {
                        // Ops! Something changed while network RPC go and come
                        return;
                    }
                    // I'm the leader
                    state = RaftState.LEADER;
                    // …
                }
            } // else greater values ignored to avoid become leader to often
        }
    });
}

It is the second hint again, from the other side: the task carries the copies it was given, and when the answer arrives it takes the guard and checks that the world it was asked in is still there.

Why the recipe works

It removes the two things a beginner gets wrong. There is one lock, so there is no order of locks to get wrong and no deadlock. And no lock is held while waiting, so nothing stalls behind a slow server. What is left is the one real difficulty, stale data, and the recipe turns it from something to reason about into something to check: a list of ifs at the top of step three.

It costs something. Work is sometimes thrown away, and it leans on Raft being built the same way — terms and indices are exactly the version numbers step three needs. But that is not a coincidence to apologise for. Optimistic concurrency, compare-and-swap, a database's UPDATE … WHERE version = ?: read, work outside, write only if nothing moved. It is the pattern most concurrent code that works turns out to have.

Making parallel machines usable by people who are not parallel programmers was what my PhD was about. This was the same problem, with students in place of scientists.

~/teaching/raft $

~/teaching/raft $