Skip to content

Day 6: Failure I, netsplits

Days 1 through 5 built the working cluster: discovery, handshake, mesh, membership, messaging. Today the network breaks. This is the 03_failure section of the companion repo and the hinge of the talk, the part of the abstract that promised “the drawbacks nobody warns you about”. Two bodies of material: how the runtime decides a node is gone, and what :global does when partitions form and heal. The second half includes the OTP 25 change, prevent_overlapping_partitions, that quietly invalidated most split-brain demos written before 2022, including possibly the one you were planning. Everything below is verified against current OTP docs and source as of today.

The runtime learns that a peer is gone in exactly two ways. Either the TCP connection closes, which the operating system reports immediately, or the connection goes silent, which costs a timeout. Every behavior in today’s material follows from that split.

Clean closes are the common case in production. kill -9, pod eviction, OOM kill, :init.stop/0, a BEAM crash: in all of these the kernel closes the socket and the peer gets nodedown within milliseconds. A partition is the other thing. Packets stop arriving and nothing tells you why. A pulled cable, a dropped BGP route, an overlay network hiccup, a conntrack table flush, and also, importantly, a machine that is alive but not scheduling: a VM paused for live migration, a container throttled to nothing, a BEAM stopped with SIGSTOP. From the outside these are all identical. The classic result that “crashed” and “slow” are indistinguishable in an asynchronous network is not theory trivia here, it is literally the branch in the code: socket closed versus tick timeout.

Day 4’s Cluster GenServer already subscribes to monitor_nodes. Extend it with the reason option, because the reason tells you which of the two worlds you are in:

def init(_) do
:ok = :net_kernel.monitor_nodes(true, %{nodedown_reason: true})
{:ok, %{}}
end
def handle_info({:nodedown, node, %{nodedown_reason: reason}}, state) do
Logger.warning("nodedown #{node}: #{inspect(reason)}")
{:noreply, handle_loss(node, reason, state)}
end

For the standard TCP distribution the reasons that matter are :connection_closed (the peer’s OS closed the socket: crash, halt, deploy), :disconnect (this node forcibly disconnected it, for example via :erlang.disconnect_node/1), and :net_tick_timeout (silence, the subject of the next section). The rarer ones are :connection_setup_failed, :send_net_tick_failed, :get_status_failed, :net_kernel_terminated, :no_network, and :shutdown. Logging the reason costs nothing and converts “nodes keep dropping” from a mystery into a diagnosis: deploys produce :connection_closed, network trouble and frozen VMs produce :net_tick_timeout. In a Kubernetes NETZlive context, pod restarts are clean closes and arrive fast; the slow painful cases are CNI trouble and nodes swapping under memory pressure, which arrive as tick timeouts a minute later.

2. The tick machinery and the 45 to 75 second window

Section titled “2. The tick machinery and the 45 to 75 second window”

The mechanism is small enough to describe completely. Two kernel parameters: net_ticktime, default 60 seconds, and net_tickintensity, default 4. The tick interval is their quotient, 15 seconds by default. Once per tick interval, each connection that had no outbound traffic during the last interval gets a tick, a small packet, so a healthy idle link carries four ticks per net_ticktime period. On the receiving side, a node is considered down if neither ticks nor payload have arrived for the last net_ticktime seconds, and this check also runs only once per tick interval. The kernel docs state the resulting detection window for defaults plainly: 45 < T < 75 seconds.

The window is a range rather than a number because both the sending and the checking are quantized to the 15 second interval, and the last packet you received lands at an arbitrary point inside one. Detection therefore spreads between 0.75 and 1.25 times net_ticktime. Worth internalizing rather than memorizing: raise net_tickintensity and the window narrows toward net_ticktime exactly, because the quantization shrinks.

Both endpoints run this detection independently with their own timers. Node A can declare B down at second 48 while B considers A alive until second 71. During that gap B sends to A, every send/2 returns :ok (Day 5: fire and forget), and every message vanishes. Asymmetric views during detection are routine, not exotic.

Operational constraints to state on stage. First, all communicating nodes must use the same net_ticktime; the docs are explicit, and a single mismatched node produces connections that drop for no visible reason. Changing it on a live cluster is possible with :net_kernel.set_net_ticktime/2, which runs a coordinated transition period, but every node must initiate the change before any node’s transition ends. Second, the temptation after your first partition is to lower the value for faster detection. RabbitMQ’s operational docs push in the opposite direction, toward raising it on unreliable networks, and they run more BEAM clusters than anyone. A short window converts every long GC, every CPU-starved scheduler, every noisy-neighbor hiccup into a false nodedown, and false nodedowns are expensive: :global prunes the node’s names, :pg prunes its members, singletons restart elsewhere, and then the whole thing reverses thirty seconds later on reconnect.

For demos the calculus flips, because no stage survives a 75 second wait. Start every peer with -kernel net_ticktime 4: one second tick interval, detection in 3 to 5 seconds. Narrate that production defaults are 15 times slower, since the number is the point. Kubernetes people in the audience are calibrated to liveness probes tuned in seconds, and the BEAM’s refusal to guess faster than silence allows will surprise them.

One more production wrinkle: ticks share the TCP connection with data. Before distribution fragmentation (OTP 22, Day 3), serializing one multi-hundred-megabyte term could hold the link long enough to blow the tick deadline and get a perfectly healthy node declared down. Fragmentation interleaves traffic now, so that specific failure is mostly historical, but a connection deep in busy_dist_port backpressure still delays everything on it, and an overloaded link can look like a partitioned one from the outside.

Nothing pauses. This deserves an unhurried minute on stage because it is where BEAM distribution diverges hardest from systems that fence or quorum by default. Each side of a partition is a complete, functioning cluster that happens to be smaller than it was. Requests get served, writes get accepted, supervisors supervise. nodedown is a local opinion about reachability, not a cluster-wide event, and there is no built-in notion of which side is “the real cluster”. Day 5’s registry table said every tool picks a flavor of wrongness during churn; a partition is that table’s worst case, held open for minutes.

Per layer: :global reacts to a nodedown by deleting every name owned by a pid on the departed node. That is deliberate, the replicated table only ever contains reachable pids. The consequence is that :grid_coordinator is now unregistered on your side, whichever supervisor owns the singleton notices and restarts it, registration succeeds because the name is locally free, and both sides now run one coordinator, each correctly registered by everything its node can see. Split brain on the BEAM is not a malfunction; it is the specified behavior of a system that replicates opinions without requiring agreement. :pg meanwhile drops the unreachable node’s members and keeps serving local reads, so Phoenix.PubSub broadcasts reach exactly the subscribers on your side and nobody is told about the rest.

Now the distinction the whole second half of today turns on: clean partitions versus overlapping ones. Take three nodes where only the a-to-c link fails. Node a sees {a, b}, c sees {b, c}, and b sees everyone. These two partitions overlap at b. For :global this shape is poison. Its write path (Day 5) promises a name is registered on all nodes or none, and its cluster-wide locks are taken across “all nodes” as each node sees them. With b a member of both components, b receives two incompatible streams of registrations and lock operations that each side believes are fully replicated. The docs say what comes of it without dramatizing: overlapping partitions “might cause the internal state of global to become inconsistent”, and the inconsistency can persist after the partitions rejoin. Names that resolve on some nodes and not others, locks that never release. This exact failure class is what OTP 25 decided to kill, in section 5.

Reconnection is not an event you schedule. With the default dist_auto_connect, any traffic re-establishes the connection: a send, a :net_adm.ping/1, an :erpc call. libcluster re-lists the node on its next poll and reconnects it (Day 4). For demos this is a trap worth naming: the partition heals the moment anything touches the other node, including your own IEx probing.

On reconnect, :global runs its name exchange. The two nodes swap name tables, and for every name registered to different pids on the two sides, the resolve function configured at registration time runs, once per clashing name. Three built-ins ship with the module. random_exit_name/3, the default for register_name/2 and re_register_name/2, keeps one pid and kills the other. random_notify_name/3 keeps one and sends the loser {:global_name_conflict, name} to handle as it sees fit. notify_all_name/3 unregisters both and sends both {:global_name_conflict, name, other_pid}, delegating the whole decision to application code.

The default deserves a close reading, because the implementation is sharper than the docs. In global.erl, random_exit_name/3 is not random at all: it orders the two pids with minmax/2, keeps the smaller by term order, logs one info-level line, global: Name conflict terminating {Name, Pid}, and executes exit(Max, kill). Reason :kill is untrappable. The losing process gets no message, no terminate callback, no chance to flush or hand off state. For a stateless worker that is the right call and nobody mourns. For a coordinator that spent the two minutes of the partition accumulating divergent state, the heal silently destroys one side’s history, selected by pid ordering, and every process still holding the dead pid finds out via monitors or not at all.

Two further edges. If a resolver crashes, or returns anything other than one of the two pids, :global unregisters the name on both sides: both processes keep running, neither is told, and callers get :undefined from whereis_name. That is a failure mode to test, not just mention. And resolution restores exactly one property, name uniqueness. It merges nothing. Whatever the killed pid knew is gone, and reconciling what the two processes did during the split is entirely your application’s problem. That sentence is the load-bearing transition of your talk: :global can give you back a unique name after a partition, and nothing in the runtime can give you back a single history. Raft exists to make the second thing possible, which is act 3.

The history first, because audience OTP versions vary. OTP-17843 introduced the fix in OTP 24.3 (also backported as far as 22.3.4.25 and 23.3.4.12), disabled by default, with notice given that OTP 25 would flip it. OTP 25.0 flipped it: the release notes list OTP-17911 both as a highlight and as a potential incompatibility, and the kernel docs state the parameter is enabled by default. On OTP 25, 26, 27, and 28, which is every version your Haarlem audience runs in production, this behavior is on unless someone deliberately turned it off.

What it does, from the algorithm comment in global.erl. When a node loses a connection to another node, it multicasts {lost_connection, Me, OtherNode} to every node it knows. A receiver that has not seen the message before re-multicasts it, so the information reaches everyone even if links are dying mid-gossip, then sends {remove_connection, me} to OtherNode and deletes OtherNode from its own view of the cluster. OtherNode, on receiving remove_connection, drops the connection and prunes in return. The result, enforced within one gossip round, is that every surviving partition is fully connected. No overlaps, no node like b from section 3 straddling two components, and therefore no inconsistent :global state.

The cost is stated with unusual candor in the same source comment: this takes down more connections than the minimum needed. If a single link between a and b dies, every other node drops both a and b. The comment explains why: receivers must decide instantly and identically, without knowing whether the lost node halted or merely lost one link, so they evict both endpoints. Walk the four-node case because it is the one that surprises people. Nodes a, b, c, d, fully meshed, and only the a-to-b link fails. c and d each receive lost_connection, each disconnects from both a and b, and the final partitions are {c, d}, {a}, {b}. One flaky link turned into two whole-node evictions. In the three-node version, cutting one link atomizes the cluster into three singletons.

When this happens you get a log line per eviction, verbatim from the source: 'global' at node c@host requested disconnect from node a@host in order to prevent overlapping partitions, at warning level. The docs describe these warnings as harmless, and mechanically they are, the cleanup rather than the outage. On a dashboard during a rolling deploy they read like an incident, and the question “why is global disconnecting my healthy nodes” is now a recurring one in Elixir forums. Expect it in your Q&A.

What this did to demos, which is the reason the curriculum flags it. Essentially every split-brain demo written before 2022, including most blog posts you will find today, produces the split by cutting one link: :erlang.disconnect_node(:b@host) from a, leaving both connected to c, then showing the asymmetric views. Under the default flag, a manual disconnect_node is indistinguishable from a failed link. The same nodedown fires, the same lost_connection gossip goes out, and the demo cluster disintegrates completely while warning logs scroll. If a three-node demo mysteriously ends with three isolated nodes, nothing is broken; the runtime did exactly what OTP 25 says it should.

Equally important is what the flag does not change. A symmetric full split, two halves with every cross-link dead, the thing a real datacenter partition produces, already yields fully connected partitions. The flag has nothing to trim; each half evicts the other, which was already gone. Each half then prunes names, restarts singletons, and registers the same :global name, and heal still runs the resolver and random_exit_name still kills. Split brain is fully alive under prevent_overlapping_partitions. What the flag eliminated is the overlapping partition shape and, incidentally, the single-link-cut method of producing a split on stage. Say that precisely in the talk, because “OTP 25 fixed netsplits” is a misreading the room will otherwise walk out with.

Operational notes worth one slide. The fix only works if every node has it enabled, and the docs strongly advise against disabling it; :global’s correctness now assumes it, together with connect_all (the Day 5 forward reference, now cashed in). For taking a node out of a cluster gracefully there is :global.disconnect/0 since OTP 25.1, which removes only the calling node with no eviction cascade. A halted node is also harmless: everyone loses the connection to it at once, the gossip only names the dead node, and nothing collateral happens. The pathological input is precisely one link down with both endpoints alive. That shape is what transient CNI trouble during rolling deploys produces, which is why clusters under libcluster on Kubernetes can see eviction-and-reconnect churn: a blip evicts two healthy nodes, libcluster re-adds them on its next poll, :global re-runs its exchange, repeat on the next blip. Not a livelock, since evictions only trigger on losses, but noisy. Verify the current state of libcluster issue threads before claiming more than that on stage.

RabbitMQ is the production data point. The rabbitmq-server startup script on main passes -kernel prevent_overlapping_partitions false, checked against the repo today. RabbitMQ shipped its own partition handling years before this flag existed and never wanted :global evicting nodes out from under it. The sequel is better for your purposes: as of RabbitMQ 4.3 the Mnesia-era partition-handling strategies, pause_minority and friends, are removed entirely, because the metadata store moved to Khepri, which is Raft built on ra, the same library in your act 3. The flagship BEAM clustering product spent a decade tuning partition heuristics and concluded the answer was consensus. Hard to buy a better endorsement of the talk’s arc.

Disabling it for a demo. The reliable place is the VM arguments, because kernel parameters are read at boot:

# Day 4's :peer harness, extended for failure demos
{:ok, pid, node} =
:peer.start_link(%{
name: :a,
host: ~c"127.0.0.1",
args: [
~c"-kernel", ~c"prevent_overlapping_partitions", ~c"false",
~c"-kernel", ~c"net_ticktime", ~c"4",
~c"-setcookie", ~c"netzsplit"
]
})

For a release, the line goes in vm.args.src: -kernel prevent_overlapping_partitions false. Compile-time Elixir config (config :kernel, ...) does flow into sys.config and generally reaches the kernel app, but runtime.exs executes after the kernel is already up, so kernel flags set there do nothing, silently. Flag that subtlety or skip the config route entirely; vm.args is unambiguous. Remember it has to be every node, your ctl node included, or you get mixed behavior that is worse to debug than either mode.

Order the section as the day ordered it: detection, then the open split, then heal, then the OTP 25 change. The flag material lands best last because it recontextualizes whatever netsplit intuition the audience brought in, and because it sets up your demo honestly: you will either disable the flag and say so, or use a full symmetric split, and explaining why is itself teaching the material.

Put one number on a slide: 45 to 75 seconds. It is the gap between the audience’s Kubernetes-calibrated intuition and the BEAM’s defaults, and everything else in the section hangs off it.

The heal-and-kill sequence is demoable in under two minutes with net_ticktime 4 peers and makes the abstract point physical: two coordinators, two histories, one exit(pid, :kill), one survivor chosen by pid comparison. Show the info log line. Then say out loud that nothing merged, and let that fact hand the microphone to Raft.

Active block (45 to 60 min, repo section 03_failure)

Section titled “Active block (45 to 60 min, repo section 03_failure)”

Build 03_failure/netsplit_lab.exs on the Day 4 :peer harness, all peers started with net_ticktime 4. (1) Subscribe to monitor_nodes with nodedown_reason on every peer via :erpc and print each event with node, reason, and timestamp. (2) Clean close: :peer.stop/1 one peer and confirm everyone logs :connection_closed immediately. (3) Silence: restart it, fetch its OS pid with :erpc.call(node, :os, :getpid, []), System.cmd("kill", ["-STOP", pid]), watch :net_tick_timeout arrive within 3 to 5 seconds, then SIGCONT and watch the reconnect. (4) The cascade: four peers with the flag at its default, cut exactly one link with :erlang.disconnect_node/1, and capture the eviction warnings and the final {c, d}, {a}, {b} partition map for a slide. (5) Rerun step 4 with prevent_overlapping_partitions false on all peers: the same cut now leaves the overlapping shape, so print Node.list/0 from all four and show c and d still seeing both a and b. (6) Heal and kill: two peers, full disconnect, register a counter GenServer as {:global, :grid_coordinator} on each side, bump both counters differently, Node.connect/1, and verify one process is dead, which survivor won, and that its counter alone remains. Stretch: give the singleton a custom resolver that crashes, heal again, and confirm the name is unregistered on both sides while both processes run.

  1. A node is SIGSTOPed at t=0 in a default-configured cluster. Give the earliest and latest moment a peer declares it down, derive both bounds from the tick interval arithmetic, and name the nodedown_reason delivered. Contrast with kill -9 on the same node.
  2. Five nodes, full mesh, the c-to-d link drops with prevent_overlapping_partitions at its default. Trace the messages (lost_connection, remove_connection), state who disconnects whom and the final partitions, and explain why the algorithm evicts both endpoints instead of one.
  3. Both sides of a healed partition had :grid_coordinator registered. List what happens on reconnect in order, through the resolver, and state precisely what random_exit_name does: how the survivor is selected, what exit reason the loser receives, and what the loser can do about it.

Say it out loud (60 seconds, to the Haarlem room)

Section titled “Say it out loud (60 seconds, to the Haarlem room)”

“A netsplit on the BEAM starts as silence, and with default settings silence takes 45 to 75 seconds to become a verdict, because detection is ticks checked every fifteen seconds. During the split nothing stops: both sides keep serving, nodedown is a local opinion, :global frees the missing side’s names, supervisors restart your singleton, and now it exists twice, each copy correct by everything its node can see. When the network heals, the default resolver restores uniqueness by picking a survivor and killing the other process with an untrappable exit, and it merges nothing. Since OTP 25 the runtime also force-evicts nodes so partitions stay fully connected, which keeps :global’s state consistent and breaks every old demo that cut a single link. What none of this machinery gives you is a shared history of what happened during the split. Getting that back requires consensus, and that is where this talk goes next.”