KnowSys

Failure Detection & Membership

Follow one worker, w3, whose heartbeats stop for a second. You'll see why a monitor can't tell a paused process from a dead one, how detectors turn silence into a verdict, how a whole cluster shares the watching, and what stops a wrong verdict from leaving two servers in charge.

⏱ 41 min read◆ BeginnerAssumes: a terminal and Python; TCP basics, replication and what a leader is help
Start reading

You run a small service called w3. Every 100 milliseconds it sends a tiny "I'm still here" message to a second program, the monitor, whose only job is to watch it. These messages are called heartbeats. One afternoon the heartbeats stop arriving. What should the monitor conclude?

The tempting answer is that w3 has crashed. But the monitor hasn't seen w3 crash. All it has seen is that nothing arrived, and plenty of healthy situations look exactly the same: w3 might be frozen for a moment while its runtime tidies up memory, the network between the two might have dropped some packets, or the monitor itself might have stalled and not been listening. From where the monitor sits, a dead server and a slow one produce the same silence.

Every group of servers that cooperate, a cluster, has to make this call about its members many times a second, because the answer decides who takes over the work. The part of the system that makes the call is the failure detector. This chapter follows w3 and asks one question: when all you can go on is silence, how do you decide a server is gone, and what keeps a wrong decision from doing damage? We'll start by watching a monitor get exactly this wrong, and each fix along the way will lead to the next problem, ending with two servers that both believe they're in charge.

01Watching w3 go quiet

1.1The simplest monitor

The simplest monitor has one rule: if no heartbeat has arrived for a set length of time, the timeout, declare the worker dead. We'll use 300 milliseconds, which is three missed heartbeats' worth.

To see what that rule does, we need a worker that is alive but occasionally stops. Many programs, Java and Go services among them, manage memory with a garbage collector, a piece of the runtime that every so often freezes the whole program while it cleans up. A long freeze like that is a GC pause, and they can last seconds. A stalled disk can freeze a process in the same way. The script below imitates one: the worker sends a heartbeat every 100 ms, and one second in, it sleeps for a full second. A monitor running in a second thread (a second line of execution inside the same program) applies the 300 ms rule every 20 ms.

Predict before you read on

The worker pauses for exactly 1 second, then carries on. The monitor's timeout is 300 ms. What will the monitor say?

Run a heartbeating worker and a 300 ms timeout monitor, and pause the worker for one second
python
Python
import threading, time
 
last_beat = time.monotonic()
start = time.monotonic()
stop = False
 
def say(msg): print(f"  t={time.monotonic() - start:4.1f}s  {msg}", flush=True)
 
def worker():                                  # a healthy service that sends a heartbeat every 100 ms
    global last_beat
    while not stop:
        last_beat = time.monotonic()
        time.sleep(0.1)
        if 1.0 <= time.monotonic() - start < 1.05:
            say("worker: entering a 1.0 s pause (a long GC, a stalled disk)")
            time.sleep(1.0)
            say("worker: pause over, still alive")
 
def monitor(timeout=0.3):                      # declares the worker dead after 300 ms of silence
    dead = False
    while not stop:
        silent = time.monotonic() - last_beat
        if silent > timeout and not dead:
            say(f"monitor: no heartbeat for {silent:.1f}s -> declares the worker DEAD"); dead = True
        elif silent <= timeout and dead:
            say("monitor: heartbeat is back -> the worker was never dead"); dead = False
        time.sleep(0.02)
 
threading.Thread(target=worker, daemon=True).start()
threading.Thread(target=monitor, daemon=True).start()
time.sleep(3)
stop = True
output
C++
  t= 1.0s  worker: entering a 1.0 s pause (a long GC, a stalled disk)
  t= 1.3s  monitor: no heartbeat for 0.3s -> declares the worker DEAD
  t= 2.0s  worker: pause over, still alive
  t= 2.1s  monitor: heartbeat is back -> the worker was never dead

Read the four lines in order. The pause starts at 1.0 s. The monitor gives its verdict at 1.3 s, exactly one timeout later, and the verdict is wrong: the worker wakes at 2.0 s fully alive, and the monitor takes the verdict back a moment later when the first heartbeat arrives. Run it again and the timestamps shift a little, because they depend on how the operating system schedules the threads. A second run can declare the worker dead at 1.2 s instead of 1.3 s.

Here is the same run from the monitor's side, with the worker's name in place and a failover added, because that's what a real monitor would do next:

w3's one-second pause, seen from the monitor
w3the workerNetworkheartbeats in flightMonitortimeout: 300 msWhat the monitor tells the systemits verdict about w3w3runningheartbeatevery 100 mssilence0.1 sw3: alivew4 startstakes over w3's jobheartbeat
Step 1. t = 0.9 s. w3 sends a heartbeat every 100 ms, and each one reaches the monitor long before its silence timer gets near 300 ms. The verdict is alive.
1 / 6

1.2What went wrong

Nothing was broken. The monitor followed its rule exactly, and the rule can't tell a worker that has died from one that has stopped for a moment. A verdict of "dead" about a worker that is alive is called a false positive. If the monitor had also promoted a replacement, as in the last frames above, we'd now have two workers doing one job. If w3's job was to be the leader, the one server allowed to accept writes, we'd now have two leaders, each taking writes the other knows nothing about.

Raising the timeout to a full second would have saved this particular run. It would also mean that a real crash takes a full second to notice, with the service down the whole time, so no value is right for everyone. Before looking for a better value, we should ask whether any rule could do better. Could a cleverer monitor look at the silence and tell which kind it is?

02Why silence can't be explained

On one machine, a crashed process is a fact: the kernel notices and tells the process's parent. Across a network there is no such signal. All a monitor ever has is messages that did or didn't arrive. From here on, a node is any member of the cluster, whether a whole server or one process on it, so w3 is a node.

2.1Four failures that look identical

Suppose the monitor hears nothing from w3. Here are the four things that could have happened, and what the monitor sees in each:

What happenedWhat the observer sees
The process crashedNo replies
The machine is up but the process is paused (a GC pause, memory being swapped out to disk, a stuck disk)No replies, then a burst of late ones
The network between you dropped or delayed packetsNo replies, while others may hear it fine
Your own process was paused, and the peer is fineNo replies arrived while you weren't looking

In the table, "you" are the monitor and the peer is the node it's watching, w3. The last row is the one that bites. A monitor that is itself frozen for three seconds will wake up, see three seconds of silence from everyone, and conclude that the whole cluster died.

?Why can't a smarter protocol tell them apart?

Because in an asynchronous network, one where a message can be delayed by any amount, a slow process and a crashed one produce exactly the same observations. Fischer, Lynch and Paterson proved a famous consequence in 1985, known as the FLP result. No deterministic protocol (one that makes no random choices) can guarantee consensus, meaning getting several nodes to agree on one value such as "who is the leader" (chapter 27 builds it), if even one process may crash. The reason is the one we've just seen: the protocol can never be sure whether to keep waiting for the silent node or to go on without it.

So real systems give up on certainty and add timing assumptions. Chandra and Toueg's unreliable failure detectors (1996) formalised the trade by naming two properties a detector can have:

  • Completeness: every crashed process is eventually suspected.
  • Accuracy: correct processes, the ones that haven't crashed, aren't suspected (or eventually stop being suspected).

Completeness is easy to get: suspect everyone who is quiet. Accuracy is the hard part, and every design in this chapter is a way of buying more of it.

2.2Gray failure: alive to one observer, dead to another

The table above treats a node as either up or down, but real failures are often partial. w3 might answer the monitor's health checks while timing out on every real request, or it might reach the monitor but not the other servers it works with. Huang et al. call this gray failure (HotOS 2017): the failure detector and the application disagree about whether a component is healthy.

The network fails partially too. When it splits so that some machines can reach each other but not the rest, that's a network partition. Bailis and Kingsbury's The Network is Reliable (ACM Queue, 2014) collects public postmortems of partitions that cut off only some machines, links that carried traffic in one direction but not the other, and switches that dropped some traffic and not other traffic. It's the best argument against assuming failures are clean.

No detector can be certain, then, and some will be wrong in both directions. The practical question is how to be wrong rarely and cheaply, and the first tool is the one our monitor already used: a timeout. We need to see what choosing its length costs.

03Heartbeats and timeouts

3.1The trade you can't avoid

w3's monitor uses the basic detector. Each node sends a heartbeat every interval, which we'll call Δ (100 ms for w3). The monitor declares a node failed if nothing arrives for a timeout T (300 ms for us). Every other design in this chapter is a refinement of that pair of numbers, and choosing T means choosing between two kinds of error:

Short timeoutLong timeout
Detection timeFastSlow: every real crash costs T of unavailability
False positivesFrequent: pauses and delays look like deathRare
Cost of each false positiveA failover, data copied to other nodes, maybe two leadersSame, but it happens less often

No setting is good on both rows, so choosing T is choosing which failure you'd rather have, and how often. That makes the timeout a design decision, to be priced like one, and the rest of the chapter is about making it well.

?Why not use TCP to tell you the peer is gone?

A natural thought is that the monitor already has a TCP connection to w3, and TCP knows when the other end disappears. It does, but slowly, because TCP's defaults are tuned to wait patiently rather than to detect quickly. A keepalive probe is an empty packet the kernel sends on an otherwise idle connection to check that the peer still answers. On Linux (kernel 6.10) the defaults are:

SettingDefaultWhat it means
tcp_keepalive_time7,200 sIdle time before the first keepalive probe (only if SO_KEEPALIVE is set)
tcp_keepalive_intvl75 sBetween probes
tcp_keepalive_probes9Probes before giving up: about 11 more minutes
tcp_retries215Retransmissions of unacknowledged data before the connection dies: at least 924.6 s, per the kernel docs

Unless a program turns on keepalive with the SO_KEEPALIVE socket option, an idle connection to a machine that lost power never notices at all. Even with it on, the first probe waits two hours, and nine unanswered probes take about eleven minutes more. A connection with unacknowledged data in flight fails faster, but "faster" is still over 15 minutes. Applications that care set their own heartbeats, or TCP_USER_TIMEOUT (a socket option that makes the kernel drop a connection once sent data has gone unacknowledged for that many milliseconds), or both.

3.2A real choice: etcd and Raft

Systems that elect a leader use the same mechanism at a small scale. Raft is a protocol for keeping copies of data on several nodes in step by electing one leader that everyone follows. Its leader sends heartbeats, and a follower that hears nothing for its election timeout assumes the leader is gone and starts an election for a new one.

etcd's tuning guide sets the defaults at a 100 ms heartbeat and a 1,000 ms election timeout, and gives two rules. The heartbeat should be around the round-trip time between members, meaning the time for a message to reach another machine and the reply to come back. The election timeout should be "at least 10 times the round-trip time" so it can absorb variance in the network.

The Raft paper adds one more trick: randomise each follower's election timeout (its example is 150–300 ms). Then one follower usually times out first and wins before the others start competing elections. The consensus chapter follows an election message by message.

RaftScope simulation of five Raft servers, all followers in term 1, each with a grey arc of a different length around it
Randomised election timeouts in a five-server simulation. Every server is a follower in term 1 and nobody has heard from a leader. The grey arc around each one is the time left on its election timer, and the arcs differ because each timeout was drawn at random. S4's is almost gone, so S4 will start the election and will most likely collect its votes before S3, the next closest, even begins.Screenshot: RaftScope by Diego Ongaro, © Stanford University, ISC licence

An election is cheap and safe to repeat, which is why etcd can afford a timeout as short as one second. For w3 we don't yet know whether 300 ms is too short or too long, because we haven't measured how late a heartbeat from a perfectly healthy process can be.

04How quiet a healthy process gets

4.1Scheduling jitter

To choose T for w3, we need the spread of gaps between heartbeats from a process that is fine. The experiment is a Python loop that sleeps for 100 ms and records the real time since the previous wake-up, 600 times per run. Nothing is wrong with the process. When a thread sleeps, the operating system decides when it gets the CPU back (chapter 06), so on a busy machine it may have to wait its turn. Small random variations in timing like this are called jitter.

The table below summarises each set of runs with three numbers. The median is the middle gap, with half the gaps shorter and half longer. The p99 is the gap that 99 out of 100 heartbeats stay under. The standard deviation is the typical wobble of a gap around the average. There were three runs each on two machines: an idle laptop, and a shared 4-CPU Linux container whose other tenants kept the load average, the number of threads wanting a CPU averaged over a minute, between 7 and 15, so two to four threads were queuing for each CPU.

MachineMedian gapp99Worst gap per runStd deviation
macOS, idle laptop103.3 ms105.1 ms105.1 · 105.1 · 106.9 ms1.7 ms
Linux, shared container101.0 ms104.1–117.4 ms112.4 · 146.4 · 132.4 ms0.8–2.6 ms

The laptop's heartbeats are regular to within a few milliseconds. On the busy shared box the worst gap in 600 heartbeats was up to 46% late (146 ms instead of 100), with nothing wrong with the process, because scheduling delay alone does that. Even so, the 300 ms timeout from section 1 has room to spare against these numbers, so jitter isn't what fooled our monitor. What fooled it was a pause, and pauses come from somewhere else.

4.2Garbage collection

A garbage collector does much worse. The next test runs a heartbeat thread sleeping 100 ms inside a Java virtual machine (JVM) whose main thread churns through roughly 800 MB of live objects in a 1.6 GB heap, using the serial collector, for 30 seconds. The heap is the memory the JVM hands out to a program's objects, and the serial collector is the simplest one, which stops every thread while it works. Here is what the heartbeat thread saw:

RunFull collectionsLongest GC pauseMedian heartbeat gapWorst heartbeat gap
1342,781 ms609 ms2,884 ms
251660 ms584 ms936 ms
350769 ms592 ms943 ms

The median column shows how constantly this JVM was being stopped: a heartbeat meant for every 100 ms arrived about every 600 ms. That's a deliberately unhappy JVM, close to its heap limit. It is also a realistic one, and it makes the point: a live process with a 100 ms heartbeat went silent for almost three seconds. A detector with a one-second timeout would have declared it dead, and if it had been the leader, it would have woken up still believing it was.

Look at what the two experiments say together. A healthy process is usually regular to within a few milliseconds, and now and then it goes quiet for hundreds or thousands. A single T has to cover the rare long gap, which makes it far too loose for the regular case. And a fixed-timeout detector throws information away. It sees every one of w3's heartbeats arrive, so it could easily learn how regular w3 normally is, yet all it reports is one bit, alive or dead.

05A suspicion level instead of a verdict

A fixed timeout treats a node whose heartbeats arrive like clockwork the same as one whose heartbeats always wander, and it gives the application nothing but a yes or a no. We'd rather have a number that says how suspicious the current silence is, so each part of the system can decide how suspicious is too suspicious for what it's about to do.

5.1The phi accrual detector

Picture w3's heartbeats arriving every 101 ms, give or take 2. A silence of 105 ms is ordinary. A silence of 150 ms is odd. A silence of one second would be astonishing, if w3 were healthy. The detector can put a number on that surprise, using the gaps it has seen.

Hayashibara, Défago, Yared and Katayama's φ accrual failure detector (SRDS 2004) does that. It keeps a window of recent heartbeat inter-arrival times (the gaps between consecutive heartbeats), fits a distribution to them, and asks: given how long it's been since the last heartbeat, how unlikely is it that a healthy node would be this quiet?

"Fits a distribution" needs a word of explanation. To answer that question the detector needs a model of how likely each length of gap is. Two systems we'll look at, Akka (a toolkit for building clustered services on the JVM) and Cassandra (a distributed database), choose differently. Akka assumes gaps follow a normal distribution, the familiar bell curve centred on the average gap with a width given by the standard deviation. Cassandra uses a simpler exponential model, in which the chance of waiting past a time t depends only on t divided by the average gap.

With a model in hand, the detector takes the current silence and asks the model for the probability that a healthy node's next heartbeat would come even later than this. Call that P. If P is large, the silence is ordinary. If it's tiny, the silence is suspicious. φ is that probability on a log scale:

φ = −log10(P(a heartbeat arrives later than now))

So φ = 1 means P is 1 in 10, φ = 2 means 1 in 100, and so on. The paper puts it plainly: if you suspect a node when φ ≥ 1, "the likeliness that we will make a mistake … is about 10%. The likeliness is about 1% with Φ = 2, 0.1% with Φ = 3, and so on." A threshold of 8, the default in both Cassandra and Akka, means a one-in-a-hundred-million chance, if the model is right.

A bell curve with the area under its right tail, beyond a point z, shaded and labelled T(z)
The bell curve of heartbeat gaps Akka assumes, with the mean gap at the peak. Put the current silence at z: the shaded area to its right, labelled T(z), is the chance that a healthy node would have stayed quiet this long, the P in the formula. As the silence grows, z moves right, the shaded area shrinks towards nothing, and φ climbs.Image: Inductiveload, public domain, via Wikimedia Commons (axis label removed)
One φ evaluation
●
♥
Heartbeats
arrival times
▤
Window
last 1,000 intervals
∿
Model
mean, std dev
⏱
Now
time since last
φ
φ
−log10 P(later)
⚑
Threshold
suspect if φ > 8
Step 1. Each heartbeat's arrival time is recorded. The detector cares about the gaps between arrivals, not the arrivals themselves.
1 / 6

The last step is the point of the whole design. The detector reports how suspicious it is and each consumer decides. A load balancer that merely stops sending new requests to w3 can act at a low φ, because a mistake costs a little capacity. A failover that promotes a replacement should wait for a high one.

5.2Akka's version, and why it floors the deviation

Here is how Akka computes φ. The first function gathers the inputs, and the second turns a delay into a number:

akka-remote/src/main/scala/akka/remote/PhiAccrualFailureDetector.scala
akka/akka @ v2.6.20 ↗
scala
  private def phi(timestamp: Long): Double = {
    val oldState = state.get
    val oldTimestamp = oldState.timestamp
 
    if (oldTimestamp.isEmpty) 0.0 // treat unmanaged connections, e.g. with zero heartbeats, as healthy connections
    else {
      val timeDiff = timestamp - oldTimestamp.get
 
      val history = oldState.history
      val mean = history.mean
      val stdDeviation = ensureValidStdDeviation(history.stdDeviation)
 
      phi(timeDiff, mean + acceptableHeartbeatPauseMillis, stdDeviation)
    }
  }
  /* ... */
  private[akka] def phi(timeDiff: Long, mean: Double, stdDeviation: Double): Double = {
    val y = (timeDiff - mean) / stdDeviation
    val e = math.exp(-y * (1.5976 + 0.070566 * y * y))
    if (timeDiff > mean)
      -math.log10(e / (1.0 + e))
    else
      -math.log10(1.0 - 1.0 / (1.0 + e))
  }

The second function measures the delay in standard deviations (y) and converts that to a tail probability with a formula that approximates the bell curve. Two lines in the first function carry lessons from production. acceptableHeartbeatPauseMillis is added to the mean, so the detector treats pauses of that length as normal. And ensureValidStdDeviation never lets the deviation drop below min-std-deviation. Why bother? Imagine w3 running on the busy shared Linux box from section 4, and try the numbers before reading the answer.

Predict before you read on

A heartbeat sender on the busy Linux box has a mean gap of 101.1 ms and a standard deviation of 2.4 ms. With no floor and no acceptable pause, at what gap does φ reach 8?

You can reproduce those numbers with Akka's formula copied into a few lines of Python. The script scores five heartbeat gaps twice, once with the raw 2.4 ms deviation and once with the deviation floored at 100 ms. Then it scores longer silences using Akka's real defaults: a heartbeat once a second, 3 s of acceptable pause added to the mean, and the 100 ms floor.

Score the same heartbeat gaps with and without Akka's deviation floor
python
Python
import math
 
def phi(gap, mean, std):
    # Akka's formula (PhiAccrualFailureDetector.scala, v2.6.20), gaps in ms
    y = (gap - mean) / std
    e = math.exp(-y * (1.5976 + 0.070566 * y * y))
    if gap > mean:
        return -math.log10(e / (1.0 + e))
    return -math.log10(1.0 - 1.0 / (1.0 + e))
 
mean = 101.1                 # a busy Linux box: heartbeats every ~101 ms ...
print("gap    raw (std 2.4)   std floored to 100")
for gap in (105, 112, 114, 132, 146):
    print(f"{gap} ms  {phi(gap, mean, 2.4):12.1f}  {phi(gap, mean, 100):14.2f}")
 
print()
print("Akka defaults: mean 1000 ms + 3000 ms acceptable pause, std 100 ms")
for gap in (4000, 4400, 4500, 4600):
    print(f"silence {gap} ms  phi = {phi(gap, 1000 + 3000, 100):.1f}")
output
C++
gap    raw (std 2.4)   std floored to 100
105 ms           1.3            0.31
112 ms           6.0            0.34
114 ms           8.5            0.35
132 ms          74.3            0.42
146 ms         213.7            0.49
 
Akka defaults: mean 1000 ms + 3000 ms acceptable pause, std 100 ms
silence 4000 ms  phi = 0.3
silence 4400 ms  phi = 4.7
silence 4500 ms  phi = 7.3
silence 4600 ms  phi = 10.8

The first table shows the problem. With the raw standard deviation of 2.4 ms, φ climbs past 8 at 114 ms and reaches 214 at 146 ms, though that gap is ordinary jitter on a busy machine. With the floor at 100 ms, the same gaps all score under 0.5. The second table shows what Akka's defaults do to a heartbeat sent once a second: φ crosses 8 at about 4.5 seconds of silence. That's a timeout in disguise, but one that stretches by itself if the heartbeats become irregular.

5.3Cassandra's version, with a guard for its own pauses

Cassandra's nodes watch each other's heartbeats too (they arrive through a mechanism called gossip, which section 7 explains). Its detector keeps the name φ and drops the normal distribution:

src/java/org/apache/cassandra/gms/FailureDetector.java
apache/cassandra @ cassandra-4.1.0 ↗
java
    public void interpret(InetAddressAndPort ep)
    {
        /* ... */
        long now = preciseTime.now();
        long diff = now - lastInterpret;
        lastInterpret = now;
        if (diff > MAX_LOCAL_PAUSE_IN_NANOS)
        {
            logger.warn("Not marking nodes down due to local pause of {}ns > {}ns", diff, MAX_LOCAL_PAUSE_IN_NANOS);
            lastPause = now;
            return;
        }
        /* ... */
        double phi = hbWnd.phi(now);
        /* ... */
        if (PHI_FACTOR * phi > getPhiConvictThreshold())
        {
            /* ... */
                listener.convict(ep, phi);
        }
    }
 
    // see CASSANDRA-2597 for an explanation of the math at work here.
    double phi(long tnow)
    {
        assert arrivalIntervals.mean() > 0 && tLast > 0; // should not be called before any samples arrive
        long t = tnow - tLast;
        lastReportedPhi = t / mean();
        return lastReportedPhi;
    }

phi here is the elapsed time divided by the mean interval, and PHI_FACTOR is 1 / ln(10). The two together give the exponential model's answer. In that model the chance of a heartbeat arriving later than t is e^(−t/mean), and taking −log10 of that gives (t/mean) / ln(10). So a node is convicted when its silence exceeds 8 × ln(10) ≈ 18.4 times its mean heartbeat interval. With gossip running once a second (Gossiper.intervalInMillis = 1000), that's around 18 seconds.

?Why check for a local pause first?

Because of the last row of the table in section 2.1. If interpret itself hasn't run for more than 5 seconds (DEFAULT_MAX_PAUSE), the silence it's about to judge is probably its own fault: a GC pause or a stalled VM on this node. Cassandra logs "Not marking nodes down due to local pause" and refuses to convict anyone until the pause is 5 seconds behind it. The check costs a few lines, and a detector without it turns its own pause into a verdict that the whole cluster has died.

With a suspicion level in hand, the next question is what to do at each level. Declaring w3 suspect and acting on it are two different decisions, and the best-known system shows how to keep them apart.

06Suspect first, act later

6.1Kubernetes: a chain of timeouts

Kubernetes runs your programs in containers (chapter 11) grouped into pods, spread over machines it calls nodes. Usually you don't create pods one by one. You create a Deployment, which says "keep this many copies of this pod running", and Kubernetes creates or replaces pods to match. Let w3 now be a whole machine in such a cluster, running a pod called web-1 that belongs to a Deployment.

The cluster keeps its state in a central service, the API server. On each node an agent called the kubelet keeps a small record there, the node's Lease, and renews it every 10 seconds. The renewal is the node's heartbeat, and it's lighter than updating the node's whole status. Inside the cluster, the node controller (in full, the node lifecycle controller) checks every 5 seconds which nodes have stopped renewing. Step through what happens when w3 loses power:

From a dead node to a rescheduled pod, with default settings
Node w3a machine in the clusterLease for w3stored in the API serverNode controllerchecks every 5 sNode w5a healthy machineSchedulerfinds a home for pods that need oneweb-1pod, runningkubeletrenews the Leasereneweda few seconds agow3: Readyweb-2pod, runningtaintunreachable:NoExecuteweb-1 copyneeds a node
Step 1. Normally the kubelet on w3 renews the node's Lease every 10 seconds, and the controller sees a fresh heartbeat each time it looks. Pod web-1 runs on w3.
1 / 6

Adding up the delays gives the time from power loss to eviction. The controller's 5-second check adds up to one more interval in the worst case, because the grace period may run out just after it looked.

Lease renewal intervalkubelet default10 s
Grace period before NotReady / Unknown--node-monitor-grace-period50 s
Default unreachable tolerationtolerationSeconds300 s
Controller check interval (worst case added)--node-monitor-period≤ 5 s
Time from power loss to eviction, defaults≈ 5 min 55 s

?Why wait five minutes when the node was declared unreachable at 50 seconds?

Because declaring a node unreachable is cheap and reversible, and evicting its pods isn't. A node that was merely cut off from the API server is still running its pods, so evicting them early starts a second copy of every pod while the first may still be serving, which is the duplicate w4 from section 1 at the scale of a whole machine. Kubernetes separates suspecting from acting, and gives acting a much longer fuse. Workloads that need faster failover set a smaller tolerationSeconds on their own pods.

It also refuses to act on a mass failure. If more than 55% of nodes in a zone are unhealthy (--unhealthy-zone-threshold), the controller slows evictions sharply, on the theory that the problem is more likely the controller's own connectivity than half the zone dying at once (nodes docs).

In Kubernetes, one controller does the watching and forms one opinion about each node. That worked for w3's lone monitor too, but it leaves two problems. The monitor's own view might be the broken one, as in row four of our table. And with thousands of nodes, you can't have every node watch every other.

07Many watchers: gossip and SWIM

Heartbeats answer "is that node alive?" for one monitor. A cluster needs every member to agree, roughly, on the whole membership list, the list of which nodes are in the cluster and alive. The naive way to get there is for every node to send heartbeats to every other node. With n nodes that's about n² messages per interval: with a thousand nodes and one heartbeat a second, a million messages a second, and the load grows with the square of the cluster.

7.1Gossip: spreading news like an epidemic

Think of how a rumour spreads in an office. Nobody tells everyone. Each person tells one or two others, who each tell one or two more, and soon the whole office knows. A gossip protocol does the same: each node periodically picks a few random peers and exchanges what it knows with them.

Suppose w1 probes w3, gets no answer, and starts suspecting it. Here's how that suspicion spreads through a cluster of eight nodes if each node that knows the rumour tells one peer per round:

A rumour about w3 spreads through eight nodes
Haven't heard yetstill think w3 is fineHave heard the rumourknow that w3 is suspectedw1w2w3w4w5w6w7w8
Step 1. w1 probed w3 and got no answer, so w1 suspects it. Nobody else knows. The other seven nodes still think w3 is fine.
1 / 4

In this idealised picture the number of nodes that know doubles each round, so eight nodes need three rounds and a thousand would need about ten, because 2 to the power 10 is about 1,000. That's what O(log n) rounds means, and meanwhile each node sends a constant number of messages per round. Real gossip picks peers at random, so some messages land on nodes that already know and it takes a few more rounds than the ideal, but the growth stays logarithmic.

Cassandra gossips once a second to a random live peer (and sometimes to a seed, one of a few well-known nodes that new members contact first, or to a node it currently believes is unreachable). Each node's state carries a heartbeat counter that the node increments itself, so a peer knows a node is alive when that counter keeps going up, even if it heard the news second- or third-hand. That counter is what feeds the φ detector from section 5.3.

So in Cassandra, gossip does two jobs at once. It spreads news, and it is also how nodes find out who is alive, since the heartbeat counters travel the same second-hand way. That works, but it makes detection only as fresh as the latest rumour, and nobody ever checks w3 directly. A cleaner design would split finding failures from spreading the news.

7.2SWIM: probe, ask for help, then suspect

Das, Gupta and Motivala's SWIM (DSN 2002; the name stands for Scalable Weakly-consistent Infection-style process group Membership) does exactly that, separating detecting failures from disseminating (spreading) membership changes. Detection is a direct check: each period, a node pings one other node and waits for a reply, an ack. Dissemination rides along on those same ping and ack messages, which is called piggybacking, and it spreads news the way we just watched.

One more piece of state makes it work. Every node keeps an incarnation number about itself, a counter that only that node may increase. Any rumour about a node is tagged with the node's incarnation number, and a message with a higher number beats one with a lower number. Follow one probe period, with M as the prober and N as its target (in our story, M is w1 and N is w3):

One SWIM protocol period: M probes N
M (prober)N (target)K1, K2, K3Everyoneping(no ack)ping-req(N)pingsuspect N (inc 7)alive N (inc 8)dead N
Step 1. Once per protocol period, M picks one member, N, and pings it directly. (memberlist: every 1 s, with a 500 ms probe timeout.)
1 / 7

The indirect probe handles the case where the path between w1 and w3 is the broken part, and w3 is fine. The suspect step handles the case where w3 is merely slow or paused, as in section 1: the rumour reaches w3, which can answer for itself, and its higher incarnation number wins. Only a node that stays silent through all of that is declared dead.

?Why does the load stay constant as the cluster grows?

Because each member sends one probe per period, plus at most k indirect requests when a probe fails, whatever the group size. The SWIM paper shows the expected time until some member first detects a crash is at most e/(e−1), about 1.58 protocol periods, also independent of group size. Picking probe targets round-robin from a shuffled list, instead of at random, bounds the worst case: every member is probed by each other member within 2n − 1 periods.

7.3memberlist: SWIM in production

HashiCorp's memberlist is the SWIM implementation inside Consul (which tracks which services run where), Nomad (which schedules jobs onto machines) and Serf (a standalone membership tool). Its defaults for a local network, from config.go at v0.5.1:

SettingDefaultMeaning
ProbeInterval1 sOne probe per period
ProbeTimeout500 msBefore falling back to indirect probes
IndirectChecks3Helpers asked to ping the target
SuspicionMult4Scales the suspicion timeout
SuspicionMaxTimeoutMult6Upper bound, as a multiple of the minimum
GossipInterval / GossipNodes200 ms / 3Dissemination of updates
AwarenessMaxMultiplier8How far a struggling node slows its own probing

The suspicion timeout grows with the cluster, because rumours take longer to reach a node in a big cluster, and a suspect needs time to hear and refute its own suspicion. From util.go, the minimum is SuspicionMult × max(1, log10(n)) × ProbeInterval:

Cluster sizeMinimum suspicion timeoutMaximum (6×)
10 nodes4 s24 s
100 nodes8 s48 s
1,000 nodes12 s72 s
10,000 nodes16 s96 s

(The comment in config.go says the maximum reaches "120 seconds" for 10,000 nodes; the formula in util.go gives 96.)

7.4Lifeguard: when the prober is the problem

SWIM has the same weakness Cassandra guards against with its local-pause check: a node that is slow to process messages misses acks, suspects healthy peers, and spreads false rumours. Dadgar, Phillips and Currey's Lifeguard extensions to memberlist fix that in two ways.

First, local health awareness. A node that has to refute suspicions about itself, or whose probes keep failing, raises its own "awareness" score and stretches its probe interval, up to 8× by default. A struggling node becomes a more patient judge.

Second, suspicion that accelerates with independent confirmations. The suspicion timer starts at the maximum and shrinks as other members independently confirm the suspicion:

Go
// remainingSuspicionTime takes the state variables of the suspicion timer and
// calculates the remaining time to wait before considering a node dead. The
// return value can be negative, so be prepared to fire the timer immediately in
// that case.
func remainingSuspicionTime(n, k int32, elapsed time.Duration, min, max time.Duration) time.Duration {
	frac := math.Log(float64(n)+1.0) / math.Log(float64(k)+1.0)
	raw := max.Seconds() - frac*(max.Seconds()-min.Seconds())
	timeout := time.Duration(math.Floor(1000.0*raw)) * time.Millisecond
	if timeout < min {
		timeout = min
	}
 
	// We have to take into account the amount of time that has passed so
	// far, so we get the right overall timeout.
	return timeout - elapsed
}

Here n is the number of independent confirmations received so far, and k is the number wanted, set to SuspicionMult − 2, so 2 by default. With none, a suspect in a 100-node cluster gets the full 48 seconds to refute. One independent confirmation brings the timer down to about 23 seconds, and the second brings it to the 8-second minimum. One confused node alone can't get w3 removed quickly. Several agreeing nodes can.

Many witnesses, a suspicion that the suspect can answer, and a detector that distrusts itself all make a wrong verdict rarer. None of them makes it impossible, as section 2 told us. So the last question is the one section 1 left hanging: when a verdict turns out wrong, what stops the replacement and the original from both acting as the leader?

08When the verdict is wrong

A false positive in a load balancer costs a little capacity. A false positive in a leader election costs much more: the old leader is still alive, still believes it's the leader, and now there are two. That state is called split brain, and it's how w3 and w4 ended up in section 1.

8.1What a 43-second partition did at GitHub

On 21 October 2018, routine maintenance to replace failing optical equipment cut connectivity between GitHub's US East Coast network hub and its primary East Coast data center. Connectivity came back in 43 seconds. In that window, their MySQL failover tool, Orchestrator, promoted primaries on the West Coast. (A primary is the database copy that accepts writes. The others, the replicas, copy it.)

When the network healed, both coasts held writes the other didn't have: several seconds of East Coast writes that hadn't replicated west, and, by the time the team came to rebuild the topology, nearly 40 minutes of application writes on the West Coast primaries. Failing back wasn't safe, and the incident ran to 24 hours and 11 minutes of degraded service while they reconciled.

?Why was 43 seconds enough to cause a day of damage?

Because the failover's timeout was shorter than the partition, and the action it triggered wasn't reversible. The detector was right that the East Coast was unreachable from where it stood. It was wrong about what that meant, and promotion across regions turned a 43-second network event into divergent histories.

8.2The tools that stop two leaders

Systems stack several defences, each closing a different hole. Here they are on w3, the leader that was wrongly declared dead, and w4, the replacement.

A quorum is a majority of the members, such as three out of five. Two groups can't both hold a majority, so if a leader must keep in touch with a majority to act, at most one side of a partition can have a working leader. If w3 is cut off with a minority, it has to stop leading, and the majority can elect w4. A lease is a promise of leadership that runs out after a fixed time unless it's renewed. If w3 can't renew, its lease lapses, and w4 takes over only once it has waited the lease out. An epoch (Raft calls it a term) is a number that goes up by one with every new leader, so w4 leads in a higher epoch than w3, and followers reject messages from a leader with an older one. A fencing token applies the same idea at the storage, and it gets a scene of its own below. STONITH stands for "shoot the other node in the head", and it means what it says: cut the power to w3 before promoting w4.

ToolHow it worksExamples
Majority quorumA leader must be able to reach a majority. At most one side of a partition has oneRaft, ZooKeeper, etcd, Redis Cluster's master votes
LeasesLeadership expires after a fixed time unless renewed; a new leader waits it outetcd leases, Chubby, ZooKeeper sessions
Epochs / termsEach new leader gets a higher number; followers reject older onesRaft terms, Kafka leader epochs, Redis configEpoch
Fencing tokensThe storage rejects writes carrying an older tokenZooKeeper zxid, etcd revisions used as tokens
STONITHPower off or isolate the old node before promotingPacemaker fencing, cloud API "stop instance"

In the table, Chubby is Google's lock service, ZooKeeper and etcd are open-source services of the same kind, a zxid is ZooKeeper's ever-increasing transaction number, and Pacemaker is a Linux tool for managing high-availability clusters.

An open server rack with a blue power strip running down the left side and a red one down the right, each server plugged into both
What STONITH acts on. Each server in this rack takes power from two power distribution units, a blue A feed and a red B feed, so it survives losing either one. A fencing agent that works by cutting power has to switch off the node's outlets on both, or the node it believes is dead keeps running.Photo: Schleifenbauer, CC BY-SA 4.0, via Wikimedia Commons

?Why isn't a lease enough on its own?

Because a lease is a promise about time, and the process holding it can lose track of time. A leader that's paused for longer than its lease wakes up, doesn't know how long it was out, and writes as if it still holds the lease. That's w3 from section 1 again, with the pause stretched past the lease. The JVM in section 4 paused for 2.9 seconds. Martin Kleppmann's analysis of distributed locking notes HBase GC pauses that lasted minutes, and a GitHub incident where packets were delayed about 90 seconds.

The fix is to make the thing being protected do the checking, instead of trusting the lock holder. Every time the lock service grants the lease, it hands out a fencing token, a number that goes up with each grant. The storage remembers the highest token it has accepted and rejects anything lower:

A fencing token stops a paused leader's write
w3client 1Lock servicegrants leases and tokensw4client 2Storageremembers the highest token it has acceptedtoken 33lease heldlast granted33highest seennone yetw3pausedtoken 34lease heldwritetoken 33grant 33
Step 1. w3 takes the lease. The lock service gives it token 33, a number that rises with every grant. Storage hasn't seen any token yet.
1 / 6

We now have the pieces: detectors that give a suspicion level, a separate step for acting, many witnesses, and fencing for the actions that must have one owner. What's left is to put numbers on them for a real system.

09Choosing a timeout on purpose

9.1Defaults across real systems

Here's how the systems in this chapter answer "how long until a silent node is declared failed, and what happens then":

SystemHeartbeatDeclared failed afterWhat happens then
etcd / Raft100 ms1 s election timeoutNew election
Akka Cluster1 sφ > 8 ≈ 4.5 s of silenceMarked unreachable; removing it from the cluster is a separate decision
memberlist (Consul)1 s probe8–48 s suspicion at 100 nodesMarked dead, removed from the service list
Cassandra1 s gossipφ > 8 ≈ 18 × mean intervalMarked down; no data moves
Redis Clustercluster-node-timeout, 15 s by defaultMajority of masters within 2 × timeoutReplica promoted
ZooKeeper sessionsClient pingsNegotiated, 2–20 × tickTimeThe session's ephemeral nodes (records tied to it) are deleted
Kubernetes nodes10 s lease50 sTainted; pods evicted 300 s later

The answers run from one second to about six minutes, and each is reasonable for what it triggers. etcd's timeout is short because an election is cheap and safe. A Kubernetes eviction is slow because rescheduling a node's worth of pods is expensive and may duplicate them.

9.2Watching it on a real cluster

Each question the chapter raised has a command that answers it on a running system.

Shell
# What will TCP do about a dead peer? (section 3.1)
sysctl net.ipv4.tcp_keepalive_time net.ipv4.tcp_keepalive_intvl \
       net.ipv4.tcp_keepalive_probes net.ipv4.tcp_retries2
ss -tno                              # per-connection timers, including keepalive
 
# How is Kubernetes tracking this node's heartbeat? (section 6)
kubectl get lease -n kube-node-lease
kubectl describe node NODE           # Conditions, and any unreachable taint
 
# What does the cluster believe about its members? (section 7)
consul members                       # alive / failed, from memberlist
redis-cli cluster nodes              # fail? (PFAIL) and fail flags
 
# Is a node convicting others because of its own pause? (section 5.3)
grep "Not marking nodes down due to local pause" system.log     # Cassandra

9.3Rules that hold up

  1. Measure the healthy gap distribution on your real machines under real load, including GC. Size the timeout from the worst gap, because the p99 leaves out exactly the rare pauses that fool detectors.
  2. Price both errors. What does a missed crash cost per second? What does one false positive cost: a failover, a cold cache, a rebalance, a possible second leader?
  3. Split suspect from act. Suspect early and cheaply (stop routing reads there); act late and expensively (promote, evict, move data).
  4. Require independent witnesses for expensive actions: a quorum, indirect probes, confirmations from other nodes.
  5. Guard against your own pauses. If the detector itself was stalled, don't convict anyone.
  6. Rate-limit actions. If many nodes look dead at once, the observer is the likeliest suspect.
  7. Fence anything that must have one owner.

9.4What you trade for what

Each tool in the chapter buys something and charges for it somewhere else. In the tables below, flapping means a node being marked down and back up again and again, the pattern a too-short timeout produces.

You getYou payWhen the bill arrives
A short timeout: crashes found fastHealthy nodes convicted during pauses and jitterAs flapping, needless failovers, a possible second leader
A long timeout: few false alarmsEvery real crash costs the whole timeoutAs minutes of downtime on a real failure
A φ detector: one dial per consumerA model that can be overconfident in the tailAs false convictions if the deviation isn't floored
Gossip and SWIM: constant load per nodeNews takes several periods to arrive, and suspicion takes seconds to clearAs slow detection in large clusters
Fencing tokens: stale leaders can't writeThe storage has to check themAs silent corruption if it doesn't

9.5Symptom, cause, fix

SymptomLikely causeFix
Nodes flap between up and downTimeout below the healthy gap distributionLengthen it, raise φ threshold, add acceptable pause
Many nodes marked down at once, then backThe observer paused or lost its networkLocal-pause guard; Lifeguard-style awareness; rate-limit actions
Failover takes minutesLong chain of timeouts (e.g. Kubernetes defaults)Shorter tolerationSeconds for that workload; app-level health checks
Two primaries after a network blipFailover without quorum or fencingMajority quorum, fencing tokens, STONITH
Dead peer's TCP connections stay open for hoursKernel keepalive defaultsApp heartbeats, TCP_USER_TIMEOUT, shorter keepalive on those sockets
Healthy node repeatedly suspected by one peer onlyAsymmetric path or gray failureIndirect probes; require multiple witnesses
Not marking nodes down due to local pause in Cassandra logsGC or VM stall on that nodeFix the pause; the detector is protecting you

10Summary

  1. A dead node and a slow one look the same. Across a network, failure detection is always a guess made from silence, and a monitor that guesses "dead" about a live worker has made a false positive.
  2. Every timeout trades detection time for false positives. There's no setting that's good at both, so pick by what each error costs.
  3. TCP won't tell you in time. Without keepalive an idle dead peer is never noticed, and with it the first probe waits two hours.
  4. Healthy processes go quiet for surprisingly long. A 100 ms heartbeat ran up to 46% late on a shared box, and a JVM busy with garbage collection went silent for 2.9 s.
  5. Accrual detectors output a suspicion level. φ is a log10 mistake probability, and each consumer can pick its own threshold.
  6. Fitted distributions are overconfident in the tail. That's why Akka floors the standard deviation and adds an acceptable pause.
  7. Detectors must distrust themselves. Cassandra's local-pause guard stops a stalled node from convicting everyone.
  8. Separate suspecting from acting. Kubernetes marks a node unreachable at 50 s and evicts its pods only after 300 more.
  9. SWIM keeps per-node load constant. Direct probe, indirect probes through k helpers, then suspicion that the suspect can refute with a higher incarnation, spread by gossip in about log n rounds.
  10. Lifeguard adds local health and witnesses. A struggling node probes more slowly, and a suspicion shrinks only as independent members confirm it.
  11. A wrong verdict becomes split brain without quorums and fencing. Leases alone fail when the holder pauses; storage that checks tokens doesn't.

11Build this

A failure detector you can break.

  • Write a heartbeat sender and a monitor over UDP. Log every inter-arrival gap, then run it on a loaded machine and inside a JVM or Go process under allocation pressure. Plot the gap distribution and mark where a 3-interval timeout would have convicted.
  • Implement φ with both the normal and exponential models over the same logged gaps. For thresholds 1 through 12, count false convictions and the delay to detect a real kill -9.
  • Build a five-node SWIM simulation with message loss and one node that processes messages slowly. Count false "dead" verdicts with and without indirect probes, and with and without a Lifeguard-style awareness score.
  • Add a fencing check to a toy key-value store and show a paused lock holder's write being rejected.

12Interview questions

beginnerWhy can't a distributed system reliably detect that a node has crashed?›

It only sees messages, and a crashed node, a paused one and one behind a broken link all produce the same silence. Any detector either waits, and is slow on real crashes, or guesses, and sometimes convicts a healthy node. FLP makes this formal for asynchronous systems: no deterministic protocol can guarantee consensus if one process may crash, because it can't know whether to keep waiting.

beginnerWhat's the trade-off in choosing a heartbeat timeout?›

A short timeout detects crashes quickly and produces more false positives from GC pauses, scheduling delay and network jitter. A long one is accurate and slow, because every real crash costs the whole timeout in downtime. Pick it from the worst gap a healthy process produces under load, and from what each kind of error costs.

intermediateHow does the phi accrual failure detector work?›

It records heartbeat inter-arrival times in a sliding window, fits a distribution (normal in Akka, exponential in Cassandra), and computes φ = −log10 of the probability that a heartbeat would still arrive given the time since the last one. φ = 1 means about a 10% chance of a mistake, φ = 2 about 1%. Applications choose thresholds; 8 is the common default.

intermediateWalk through SWIM's failure detection.›

Each period a member pings one random (or round-robin) peer. If no ack arrives within the probe timeout, it asks k other members to ping that peer. If none succeed, it marks the peer suspect and gossips that with the peer's incarnation number. The peer can refute by gossiping alive with a higher incarnation. If the suspicion timeout passes without refutation, the peer is declared dead. Updates piggyback on probe messages, so load per member is constant.

intermediateA Kubernetes node loses power. How long until its pods run elsewhere, and why?›

About six minutes with defaults. The kubelet renews a Lease every 10 s; after 50 s without one (node-monitor-grace-period) the node goes Unknown and gets an unreachable:NoExecute taint. Pods tolerate that taint for 300 s by default, then are evicted and rescheduled. The long toleration exists because a merely partitioned node may still be running its pods.

deepWhy does Akka enforce a minimum standard deviation in its phi detector?›

Because a very regular heartbeat gives a tiny standard deviation, and a normal distribution with a tiny deviation calls any small delay astronomically unlikely. On a busy shared Linux box the gaps have a mean of about 101 ms and a deviation of about 2.4 ms. With those, φ would cross 8 at 114 ms, and a real 146 ms scheduling hiccup would score about 214. The floor (100 ms by default) plus an acceptable heartbeat pause (3 s) keeps ordinary jitter from convicting healthy nodes.

deepA leader holds a 10-second lease. Why can it still corrupt data, and what fixes it?›

It can pause (GC, VM stall) for longer than the lease, wake without knowing time passed, and write while a new leader also writes. Clocks and leases can't fix that from the holder's side. The fix is at the resource: every grant carries a monotonically increasing fencing token, and storage rejects any write with a token lower than the highest it has seen. Consensus-based stores provide suitable tokens (ZooKeeper zxid, etcd revision).

deepWhat went wrong in GitHub's October 2018 incident, in failure-detection terms?›

A 43-second partition exceeded the failover tool's detection threshold, so Orchestrator promoted West Coast primaries while East Coast ones still held unreplicated writes. The detector was locally correct and globally wrong, and the action (cross-region promotion) wasn't reversible. Both sides then had writes the other lacked, and recovery took over 24 hours. The lessons: price the action, not just the detection, and don't let a short partition trigger an irreversible topology change.

13Go deeper

check yourself
φ = 3. What's the probability the detector is wrong, according to the paper's model?›

About 0.1%: each unit of φ is a factor of ten.

Cassandra with phi_convict_threshold 8 and 1-second gossip. Roughly how long a silence convicts a node?›

About 18 seconds: phi is t / mean scaled by 1/ln 10, so conviction needs t above 8 × ln 10 ≈ 18.4 mean intervals.

In SWIM, who can increment a node's incarnation number?›

Only the node itself, when it refutes a suspicion. That's why its alive message always overrides the suspect rumour.

Why does Kubernetes slow evictions when more than 55% of a zone looks unhealthy?›

Because a mass failure is more likely a problem with the controller's own view, such as a partition, than half the zone dying at once. Evicting everything would make it worse.

Das, Gupta and Motivala, SWIM (DSN 2002)

Randomised probing, indirect probes, suspicion and infection-style dissemination, with the analysis of detection time and load. PDF

Hayashibara et al., The φ Accrual Failure Detector (SRDS 2004)

The accrual idea and the φ scale, with measurements against fixed-timeout detectors. DOI

Dadgar, Phillips and Currey, Lifeguard

What went wrong with SWIM in HashiCorp's production clusters and the local health extensions that fixed it. arXiv 1707.00788

Chandra and Toueg, Unreliable Failure Detectors (JACM 1996)

Completeness, accuracy, and the failure-detector classes that make consensus solvable. PDF

GitHub: October 21 post-incident analysis

A 43-second partition, a cross-region failover and 24 hours of recovery, told in detail. github.blog

Kleppmann, How to do distributed locking

Why leases need fencing tokens, with the pause and delay examples. martin.kleppmann.com

Consensus: Raft, Paxos & Leases

Elections, terms and leases, the machinery that turns a failure verdict into one new leader. Chapter 27.

Replication & Consistency Models

What a failover can lose, depending on how replicas were kept in step. Chapter 28.

Partitioning & Rebalancing

Why moving a dead node's partitions is expensive, and why systems wait before doing it. Chapter 29.

Contention, Queueing & Tail Latency

Where the long tail of heartbeat gaps comes from: scheduling, queueing and pauses. Chapter 16.