Skip to content

Day 1: The thesis

Seventeen days to Haarlem. Today covers the claim the whole talk hangs off: reach versus agreement, what the transparency layer actually promises, and the layer cake from epmd up to libcluster. The next nine days work through the parts. This one fixes the frame.

Distribution on the BEAM gives you reach, not agreement.

Reach means any process can address, signal, link to, monitor, and spawn any other process on any connected node, with the same primitives it uses locally. send/2 works on a remote pid. A pid carries its node with it. All of that lives in the runtime rather than in a library, and it is a lot of machinery.

Agreement means the cluster deciding one thing together, with every node acting on that single decision: one leader, one owner of a name, one committed value at index 42, one membership list. The BEAM ships none of this. Node.list/0 is a local opinion, not a membership record. :global looks like agreement, and Day 9 covers how it fails. If you need agreement you build it (Raft, Days 11 and 12), import it (ra), or outsource it (Oban’s Postgres coordination, Day 13).

Put in mechanism terms, distributed Erlang is a transparency layer over signal passing. It extends the process model across machines and is explicit about what it cannot know. It is not a message bus (no durability, no replay), not an RPC framework (though :erpc sits on top), and not a coordination service (there is no consensus anywhere in the stack).

In production terms, every distributed Erlang outage in your talk is a case of reach being used where the problem needed agreement: the :global split-brain, the Swoosh Process.link cascade, the rolling deploy that behaves like a netsplit. That is also why the caveats section works as the bridge into Raft rather than as a disclaimer at the end.

For the Haarlem audience, who mostly cross service boundaries over HTTP, there is a useful inversion. They work without transparency or agreement and compensate with infrastructure: service discovery, gateways, queues, meshes. The BEAM gives them transparency for free. The trap is assuming agreement came with it. It did not, and could not have: that is FLP and the two generals problem, not a BEAM limitation.

The phrase in your abstract, “not a library bolted on top”, needs to be made concrete, because the audience will ask.

The distribution channel carries the BEAM’s full signal vocabulary, not just messages. The control-message set of the wire protocol (Day 3 goes byte by byte) includes SEND and REG_SEND, but also LINK, UNLINK, EXIT, EXIT2, MONITOR_P, DEMONITOR_P, MONITOR_P_EXIT, GROUP_LEADER, and, since OTP 23, the SPAWN_REQUEST and SPAWN_REPLY pair that Node.spawn/2 and :erpc use. Exit propagation, supervision semantics, and monitor delivery cross the machine boundary natively. HTTP systems have no equivalent of “if that process dies, deliver me an exit signal with its reason”, because over HTTP a process is not an addressable thing. On the BEAM it is, cluster-wide.

Identity crosses the boundary too. Pids, references, and ports encode their origin node. node(pid) works on any pid. A ref created on node a and sent to node b still matches in a receive on b, and still demonitors correctly when it travels back. This is why GenServer calls work across nodes without modification: the {ref, reply} plumbing does not care where the server runs.

A ten-second demo that shows this on stage:

# terminal 1
# iex --sname a --cookie talk
# terminal 2
# iex --sname b --cookie talk
Node.connect(:"b@yourhost")
pid =
Node.spawn(:"b@yourhost", fn ->
receive do
{:ping, from} -> send(from, {:pong, node()})
end
end)
send(pid, {:ping, self()})
flush()
#=> {:pong, :"b@yourhost"}

And the group leader demo, which surprises even people who have run clusters for years:

Node.spawn(:"b@yourhost", fn -> IO.puts("where does this print?") end)
# It prints in THIS shell, on node a.

The spawned process runs on b, but its group leader is your shell process on a, so its IO arrives as a message over the dist channel and prints in your terminal. GROUP_LEADER is a control message, and IO is itself just messages. Works well as an opening beat: the runtime did not forward stdout, it routed an IO protocol message to a process on another machine through the same mechanism as everything else.

One cost note to keep handy: a local send copies the message into the receiver’s heap (large refc binaries are shared by reference instead). A remote send serializes the whole term to External Term Format, binaries included, and copies it again on arrival. The semantics are transparent. The cost is not.

The actual guarantees are few enough to recite, and you should be able to. This is the reference manual’s processes and distribution chapters, compressed.

Ordering. Signals from process P to process Q are delivered in the order sent, if they are delivered at all. The guarantee holds per sender-receiver pair and nowhere else: no ordering across pairs, no global order, no causal order. If P sends to Q and also to R, and R forwards to Q, then Q can see the forwarded signal before the direct one.

Delivery. Local sends to a live process are not lost. Remote sends are fire-and-forget past the dist boundary. send/2 returns the message in every case, and if the connection dies while your message sits in the dist buffer or in flight, the message is dropped with no error and no NACK. There is no acknowledgment anywhere in the dist protocol. The way you find out something went wrong is a monitor:

ref = Process.monitor(pid)
:rpc.cast(node(pid), :erlang, :halt, [])
receive do
{:DOWN, ^ref, :process, ^pid, reason} -> reason
end
#=> :noconnection

:noconnection, not :killed. The reason means the runtime no longer knows anything about that process, only that its node is unreachable. Links behave the same way, delivering an exit signal with reason :noconnection. A large part of production distributed Elixir is deciding what your system does when it receives that signal.

Blocking. Transparency also leaks in the other direction. Local send/2 never blocks. A remote send can suspend the calling process: when the dist buffer toward a node fills past the busy limit, senders to that node are descheduled until it drains (busy_dist_port, tunable with +zdbbl, default 1024 KB, on the Day 13 verify list). A slow or partitioned peer can therefore stall your hot path through an ordinary-looking send, including GenServer casts to remote names. For people who learned that send never blocks, this is usually the most surprising fact in the talk, and it belongs in the drawbacks section.

Failure detection. “Down” means “silent past the tick window”, roughly a minute at defaults (net_ticktime; the exact 45 to 75 second window is Day 9 material). The runtime cannot tell a dead node from a slow link, because nothing can. The two generals problem is the direct reason nodedown is a local judgment based on a timeout, not a fact about the other machine.

Membership. Every node runs its own net_kernel and keeps its own connection list. Node.list/0 on a and on b can disagree, and during a partition they will, for the whole detection window and after it. The default is a full mesh maintained transitively: when a connects to b, they exchange node lists and connect to everything the other knows (-connect_all true, the default; hidden nodes and -dist_auto_connect never opt out). A mesh of n nodes holds n(n-1)/2 connections, which is why very large meshes are rare. That number is on the Day 13 cheat sheet.

That is the whole contract: pair-wise ordering, best-effort delivery, explicit failure signals, timeout-based suspicion, and per-node membership opinions. The next nine days cover either the machinery that implements this contract or the patterns for living inside it.

Each of these gets its own day later. Today they matter as instances of the same mistake.

First, :global registration. It does locked multi-node coordination at registration time, and it works until a partition. Then each side registers its own “singleton”, because each side’s :global can only coordinate with the nodes it can reach. On heal, :global’s resolver discovers the conflict and resolves it by killing one of the claimants. Two singletons ran concurrently for the duration of the split, and reconciling whatever they both did is left to you. Day 9, including how OTP 25’s prevent_overlapping_partitions changed this behavior.

Second, the Swoosh incident from jola.dev, already in your notes. During rolling deploys, the loser of a :global registration linked to the winner. Deploys churned nodes, the links propagated exits, and the cascade took the service down. Links were the wrong tool because the actual requirement was agreement about which process runs and an orderly handover. Day 10.

Third, the quieter one. Phoenix.PubSub broadcast is at-most-once with no cross-publisher ordering. That is fine for cache invalidation, presence hints, and “something changed, go refetch”. It fails silently for anything that must arrive or must arrive in order. Teams promote PubSub from notifications to transport without noticing they crossed that line. Day 8.

For the talk’s structure: each caveat is the same story with different nouns. The team had reach, assumed agreement, and the network eventually disagreed. That sets up Raft as the thing that actually provides agreement, priced in majorities, quorum loss, and write latency. Oban is the counterweight: a team that already operates Postgres already operates a strongly consistent coordination point, with backups and an on-call rotation, so coordinating through it instead of through distributed Erlang is often the cheaper engineering decision. Day 13.

The map for the rest of the curriculum, bottom up:

6 libcluster, dns_cluster, Phoenix.PubSub, Horde, syn, ra ecosystem
5 :global :pg :erpc / :rpc :net_adm Node OTP primitives
4 net_kernel: mesh maintenance, ticks, nodeup/nodedown membership opinions
3 control messages + ETF payloads, atom cache, fragmentation signal vocabulary
2 handshake: 'N', capability flags, cookie challenge-response identity + capabilities
1 carrier: inet_tcp_dist | inet_tls_dist | custom bytes
0 epmd on 4369: node name -> port rendezvous discovery

For each layer, know what it does and, just as useful on stage, what it does not do.

Layer 0, EPMD. A small daemon on a fixed port that maps node names to the dynamically chosen dist ports. The departure board at Tulln Hauptbahnhof: it tells you the platform, it does not drive the train. It does not authenticate, does not route traffic, and does not know about clusters. Deregistration is the TCP connection dropping, nothing more. It is also replaceable (-start_epmd false, custom epmd_module). Day 2.

Layer 1, the carrier. Plain TCP by default via inet_tcp_dist, and pluggable: -proto_dist inet_tls swaps in TLS distribution without touching anything above it. That this is a swappable module is worth a line in the talk, because it shows distribution is a protocol with defined seams rather than something fused into the VM.

Layer 2, the handshake. Names exchanged, capability flags negotiated (which is how old and new nodes interoperate), then the cookie challenge-response: each side proves knowledge of the shared cookie by MD5-hashing it with a challenge. Worth one line today and a full treatment later: this authenticates, it does not encrypt, and possession of the cookie is remote code execution on every node. The cookie is a shared password, not a security boundary. Days 3 and 10.

Layer 3, the wire protocol. The control-message vocabulary plus ETF payloads, with the atom cache so repeated atoms cost one byte, fragmentation so a 100 MB term does not head-of-line-block everything else (since OTP 22), and tick heartbeats feeding the failure detector. This layer has no acks, and no flow control beyond TCP’s plus the busy dist buffer. Day 3, with Wireshark.

Layer 4, net_kernel. The process that owns connections, maintains the transitive full mesh, runs ticks, and publishes nodeup and nodedown to subscribers of :net_kernel.monitor_nodes/1. Membership lives here, as per-node opinion. Day 7 builds the Cluster GenServer on exactly this API.

Layer 5, OTP primitives. Node, :erpc, :net_adm, :pg, :global. All ordinary OTP code built on signals. :pg is local ETS plus cross-node gossip. :global is a locking protocol among connected nodes. Nothing below changes when you use them. Day 8.

Layer 6, the ecosystem. The fact that reliably surprises people: libcluster does not do distribution. It schedules Node.connect/1 calls. A strategy (gossip, Kubernetes DNS, EC2 tags) produces a list of names, libcluster connects to them, and everything else happens in layers 0 through 4. The same goes for dns_cluster, which the Phoenix generator now ships by default (it polls DNS and connects; the exact generator version is on the Day 7 verify list). Phoenix.PubSub, Horde, and syn are patterns over :pg and friends. ra is the one library in the list that adds something the lower layers lack, namely agreement. That restates the thesis as a diagram: six layers of reach, with agreement only appearing when you import it at the top.

One more reason to teach it bottom-up: everything above layer 2 is processes sending ETF-encoded signals over a socket you can watch in Wireshark. The stack is thin and inspectable end to end, which is what makes your Section 1 packet-capture material possible and gives the audience a reason to trust the rest.

Your target listener compensates for missing transparency with infrastructure. A short mapping early in the talk connects their world to the internals you are about to show.

Service discovery (Consul, k8s DNS) corresponds to EPMD plus a libcluster strategy. RPC frameworks (gRPC, REST clients) correspond to send/2 and :erpc.call/4. Fan-out buses used purely for broadcast (Redis pub/sub, SNS) correspond to :pg and Phoenix.PubSub. Health checks and liveness probes correspond to ticks, monitors, and nodedown. Sidecars doing mTLS correspond to -proto_dist inet_tls.

Then show the reverse column, because conceding it is what keeps the talk credible: Kafka gives durability, replay, per-partition ordering, and consumer backpressure, and distributed Erlang gives none of those, by design, because it is a signal layer rather than a log. The planted Q&A question (“when would you not reach for distributed Erlang”) is answered from this slide: when you need agreement, durability, or replay, pick the tool that has it, and sometimes that tool is Postgres with Oban on top.

Reading notes for today’s two references

Section titled “Reading notes for today’s two references”

The Distribunomicon (Learn You Some Erlang) states today’s thesis without ever using the word “agreement”. Read it for three things. Its treatment of the fallacies of distributed computing (Deutsch’s list: the network is reliable, latency is zero, bandwidth is infinite, the network is secure, topology does not change, there is one administrator, transport cost is zero, the network is homogeneous) applied specifically to Erlang, which gives you quotable material. Its blunt framing of cookies as a convenience for grouping nodes rather than a security mechanism. And its “my other cap theorem” section, a usable 90-second CAP answer if a question drags you there, though CAP should stay out of the prepared material since the reach/agreement frame does the same work with less baggage.

The Distributed Erlang chapter of the system documentation is short, authoritative, and current. Read it for the precise definitions: what an alive node is, how -connect_all false and hidden nodes change mesh behavior, the -dist_auto_connect never flag, and the pointers into net_kernel and global. It is also the right citation for any slide that claims something about default behavior, since blog posts go stale and this does not.

Current OTP is the 29 line (29.0 released May 2026, 29.1 is out). Checked today: the OTP 29 highlights contain nothing distribution-level, no changes to epmd, the dist protocol, net_kernel, global, or pg, so everything in this curriculum describes OTP 26 through 29 semantics unless a day flags otherwise. The two version-sensitive stories in the talk remain prevent_overlapping_partitions (OTP 25 era, Day 9) and the ra API (Day 12). Your audience will mostly run OTP 27 and 28 in production, and nothing in the thesis material differs across 26 to 29.

Flagged today as “verify before stage”: the 1024 KB +zdbbl default (Day 13 check), and the exact Phoenix generator version that began shipping dns_cluster (Day 7 check).

In the companion repo, create 02_processes/day01_reach.exs as the seed of the opening demo. Three acts, scripted with :peer so it runs as one file.

Act 1, transparency: start a peer node, spawn the ping process on it, show the send/receive round-trip, node(pid), and the IO.puts group-leader behavior printing on the origin node.

Act 2, the failure signal: Process.monitor the remote pid, then :rpc.cast(peer, :erlang, :halt, []) and pattern match the {:DOWN, _, _, _, :noconnection}. Print the reason prominently, since the reason is the point.

Act 3, fire-and-forget: start the origin with --erl "-kernel dist_auto_connect never", connect explicitly, then :erlang.disconnect_node/1 and send to the now-unreachable remote pid. Observe send return the message with no error and nothing arriving. Add a one-line comment above each act stating which contract clause it demonstrates.

Done when the file runs clean twice in a row with elixir --sname origin day01_reach.exs and each act prints a labeled result.

  1. State the complete ordering and delivery contract for remote sends, and name the only mechanism by which you learn a remote send did not arrive.
  2. Walk the layer cake from port 4369 to libcluster, one sentence per layer on what it adds, and one on what it deliberately does not do.
  3. Under what condition does send/2 to a remote pid block the caller, and why does that matter for the location-transparency story in production?

Sixty seconds, Haarlem voice: “You’ve all built systems where services talk over HTTP. The BEAM makes that boundary disappear: a process on another machine is addressed, monitored, linked, and spawned exactly like a local one, because distribution is in the runtime, not in a library. But here’s the thing I need you to hold onto for the next forty minutes…” Finish it from memory, landing on reach versus agreement and why the disappearing boundary is both the best feature and the trap.