Day 4: Forming the cluster
Yesterday ended with two nodes that had survived EPMD, the cookie challenge, and the ‘N’ handshake, and were exchanging control messages over a live TCP connection. Today covers everything above that: how a set of nodes becomes a cluster, who decides to dial whom, how your application finds out the topology changed, and what actually runs in production at 3am during a rolling deploy. This is the 01_clustering section of the companion repo, and the point in the talk where you move from wire protocol to patterns.
Keep the thesis in frame. Everything today is about reach: discovery libraries and monitor messages extend your reach to other nodes, and none of it gives you agreement. Cluster membership on the BEAM is not a consensus value. It is each node’s local opinion, updated asynchronously by messages, so two nodes can hold different opinions about who is in the cluster at the same instant, and both are “right” from where they stand. Day 6’s netsplit material builds directly on this.
1. What “connect” actually does, and the mesh you get for free
Section titled “1. What “connect” actually does, and the mesh you get for free”Node.connect/1 is a thin wrapper over :net_kernel.connect_node/1. It triggers exactly the machinery from yesterday: EPMD PORT_PLEASE2 for the target’s dist port, TCP connect, ‘N’ handshake, cookie challenge, and then the connection is up and bidirectional.
In practice you rarely call it, because two mechanisms call it for you.
The first is auto-connect on first use. Distribution connections are established lazily. Send a message to {name, :"other@host"}, call :erpc.call/4, Node.monitor/2, Node.spawn/2, or even Node.ping/1 against a node you have never spoken to, and the runtime dials it for you, handshake and all, before delivering the signal. Addressing a process on another node is enough to create the link to that node. You can turn this off with the kernel parameter dist_auto_connect set to never (connections then only happen via explicit connect_node) or once (a lost connection is never automatically re-established). never is a surprisingly useful knob for netsplit demos, note it for Day 6.
The second is the transitive full mesh. When a visible node connects to another visible node, the global name server on each side exchanges its list of known nodes, and each node then connects to everyone the other knew. Connect A to B while B already knows C, and you shortly have A-B, B-C, and A-C. The default topology is a full mesh, n(n-1)/2 connections, which is why the folk wisdom caps comfortable cluster sizes somewhere in the dozens to low hundreds of nodes rather than thousands. You can suppress mesh formation with the kernel parameter connect_all set to false (historically the -connect_all false emulator flag; OTP 25 moved it to a kernel parameter, verify the exact spelling you want on the slide before stage), at the cost of global not working. Hidden nodes (--hidden) opt out of the mesh per node: they connect only where explicitly told, appear in Node.list(:hidden) but not Node.list(), and do not propagate. Observer CLI sessions and debug shells attached to production are the classic hidden-node use case.
The practical consequence for today: a discovery library does not need to connect every pair. It needs to get each node connected to one member of the mesh, and transitivity does the rest. This is why most libcluster strategies are so simple internally.
2. :net_kernel.monitor_nodes and its guarantees
Section titled “2. :net_kernel.monitor_nodes and its guarantees”The subscription API for topology changes:
:ok = :net_kernel.monitor_nodes(true)# => {:nodeup, node} and {:nodedown, node} to the calling process
:ok = :net_kernel.monitor_nodes(true, [:nodedown_reason, node_type: :visible])# => {:nodeup, node, info} and {:nodedown, node, info}, info is a keyword listWith options you get the 3-tuple form, and info mirrors the option format (pass a map of options, get a map back; pass a list, get a list). The useful options:
node_type::visible(default behavior),:hidden, or:all.:nodedown_reason: includes why the connection died. The documented reasons are:connection_setup_failed,:no_network,:net_kernel_terminated,:shutdown,:connection_closed,:disconnect,:net_tick_timeout,:send_net_tick_failed, and:get_status_failed.:connection_id: an identifier for the specific connection instance, letting you distinguish “same node, new connection” from “same connection”.
The reason atoms deserve a slide of their own, because they are your first diagnostic in production. :connection_closed means the TCP connection died, typically because the remote BEAM process exited or was killed, the common signature of a deploy. :net_tick_timeout means the connection was alive as far as TCP knew but the remote stopped answering heartbeats, which points to a partition, a frozen VM, or a severed network path. :disconnect means somebody called Node.disconnect/1, a clean administrative action. The NETZlive war stories almost certainly sort into these buckets. Tell the audience to log the reason, because it distinguishes a deploy from a netsplit, and the abstract promises exactly this kind of production detail.
Now the ordering guarantees, quoted from the kernel docs because they are precise and they matter:
- “nodeup messages are delivered before delivery of any signals from the remote node through the newly established connection.”
- “nodedown messages are delivered after all the signals from the remote node over the connection have been delivered.”
- “nodeup messages are delivered after the corresponding node appears in results from erlang:nodes().”
- “nodedown messages are delivered after the corresponding node has disappeared in results from erlang:nodes().”
And since OTP 23: a nodedown for a dying connection is guaranteed to arrive before the nodeup of a replacement connection to the same node. Before that fix, a fast reconnect could deliver nodeup then nodedown and your membership state would end up wrong while the node was actually fine.
These four guarantees mean the naive bootstrapping pattern is sound. Subscribe first, then read Node.list(), and you cannot miss a node or double-count one, because any nodeup arriving after your snapshot refers to a node that was not yet in it. Monitor messages are also delivered to the subscriber in the order the events happened per node. So the runtime hands you a linearized, per-observer-consistent event stream about topology, which is a well-engineered corner of the BEAM and worth a moment in the talk. What it does not hand you is cross-node agreement about that topology. Node A’s stream and node B’s stream can disagree about whether C is up right now.
One sharp edge: monitor messages go to the process that subscribed, and if that process dies, the subscription dies with it. Put the subscription in init/1 of a supervised GenServer so a restart re-subscribes, which brings us to the pattern.
3. The Cluster GenServer pattern
Section titled “3. The Cluster GenServer pattern”Jola’s “Elixir Cluster 101” post distills what most production Elixir shops end up writing: a small supervised GenServer that owns the cluster view. The shape:
defmodule NetzLive.Cluster do use GenServer require Logger
def start_link(opts), do: GenServer.start_link(__MODULE__, opts, name: __MODULE__)
def members, do: GenServer.call(__MODULE__, :members)
@impl true def init(_opts) do :ok = :net_kernel.monitor_nodes(true, [:nodedown_reason, node_type: :visible]) members = MapSet.new([Node.self() | Node.list()]) {:ok, members} end
@impl true def handle_info({:nodeup, node, _info}, members) do members = MapSet.put(members, node) Logger.info("nodeup #{node}, cluster size #{MapSet.size(members)}") :telemetry.execute([:netzlive, :cluster, :change], %{size: MapSet.size(members)}, %{event: :nodeup, node: node}) {:noreply, members} end
def handle_info({:nodedown, node, info}, members) do members = MapSet.delete(members, node) reason = Keyword.get(info, :nodedown_reason) Logger.warning("nodedown #{node} reason=#{inspect(reason)}, cluster size #{MapSet.size(members)}") :telemetry.execute([:netzlive, :cluster, :change], %{size: MapSet.size(members)}, %{event: :nodedown, node: node, reason: reason}) {:noreply, members} end
@impl true def handle_call(:members, _from, members), do: {:reply, MapSet.to_list(members), members}endAn expert audience will ask the obvious question: why a GenServer when Node.list() is free? Three honest answers. First, hooks. Node.list() tells you the state, not the transitions, and the transitions are where you attach behavior (rebalance a hash ring, re-run singleton election, emit metrics) and where you capture nodedown_reason, which is otherwise lost. Second, a stable synchronous read point. Everything that asks members/0 sees a view consistent with the event stream this one process has consumed, instead of racing Node.list() at arbitrary moments. Third, observability. Jola’s production advice is to turn membership changes into gauges and counters and alert when a node fails to rejoin within an expected window. Cluster size as a first-class metric is the cheapest netsplit detector you will ever deploy, and it makes a concrete slide: one Grafana line going from 6 to 5 and back within 90 seconds is a deploy, while going to 3 and staying there is an incident.
Say in the talk that this GenServer is each node’s local opinion of membership. Deploy it on six nodes and you have six opinions that are usually identical and occasionally not. The post is also clear about when this foundation suffices, and it matches the thesis: distributed use cases where some data loss is acceptable and temporary inconsistency is acceptable. When neither is acceptable, you need Day 8 and ra.
Rolling deploys are topology churn, on schedule
Section titled “Rolling deploys are topology churn, on schedule”Modern deployment means this machinery runs constantly, not just during incidents. A Kubernetes rolling update of a 6-node cluster is, from the BEAM’s perspective, six partial cluster failures and six joins in a few minutes: pod terminates (SIGTERM, BEAM shuts down, connections close, nodedown with :connection_closed everywhere), replacement pod starts with a new IP and therefore a new node name, discovery finds it, nodeup everywhere, repeat. Two design consequences:
- Anything keyed on node identity (HRW hash rings, singleton election,
:globalnames) re-shuffles on every deploy, not just on failures. Day 7’s Swoosh incident is exactly a bug in this churn path, so plant the seed today:nodedown/nodeuphandlers run many times per day in production, and any pathological behavior in them gets exercised at every deploy. - Code that assumes “the cluster” is a fixed set is wrong within a week of launch. The Cluster GenServer pattern treats membership as a stream of events, which is the correct mental model.
4. libcluster under the hood
Section titled “4. libcluster under the hood”libcluster (current version 3.5.0) is the standard answer to the question of who calls Node.connect/1. Its architecture is smaller than its reputation suggests, and walking through it gives the talk another “what is this library actually doing” moment, parallel to the EPMD unveiling.
You configure topologies, each naming a strategy and its config, and start a Cluster.Supervisor:
config :libcluster, topologies: [ netzlive: [ strategy: Cluster.Strategy.Kubernetes, config: [ mode: :ip, kubernetes_node_basename: "netzlive", kubernetes_selector: "app=netzlive", kubernetes_ip_lookup_mode: :pods, polling_interval: 5_000 ] ] ]
# application.exchildren = [ {Cluster.Supervisor, [Application.get_env(:libcluster, :topologies), [name: NetzLive.ClusterSupervisor]]}, ...]Each topology also accepts connect, disconnect, and list_nodes as MFA overrides. The defaults are essentially :net_kernel.connect_node/1, :erlang.disconnect_node/1, and Node.list/1. That is the whole contract: a strategy is a process that periodically (or reactively) computes “nodes that should exist”, diffs against “nodes I see”, and calls connect/disconnect. Writing a custom strategy means implementing the Cluster.Strategy behaviour, one start_link/1 that receives the topology state. There is no magic layer. libcluster is a discovery loop with pluggable discovery.
The shipped strategies, with the details that matter:
Epmd: a static host list in config (hosts: [:"a@127.0.0.1", :"b@127.0.0.1"]), connect to each. Discovery by configuration, fine for fixed fleets and demos.LocalEpmd: asks the local EPMD (yesterday’s NAMES request) for every node registered on this machine and connects to all of them. Good for local dev, since N iex sessions started with sname find each other, and a natural callback to Day 2 in the talk.ErlangHosts: uses the.hosts.erlangfile, the OTP-native mechanism that predates all of this.Gossip: multicast UDP heartbeats. Defaults: port 45892, multicast address 233.252.1.32, TTL 1. Each node periodically broadcasts its name and receivers try to connect. Asecretoption encrypts the gossip packets, which both authenticates (sort of) and lets multiple clusters share one network segment. There is alsobroadcast_only: truefor networks without multicast. The production caveat you must state: most cloud VPCs (AWS, GCP, Azure) block or do not route multicast, so Gossip is effectively a LAN and on-prem strategy. It is popular for demos because zero config feels magical, and it is also the strategy most likely to mysteriously fail at a venue. Do not build the live demo on it over conference WiFi; AP isolation kills it, and the travel-router backup plan exists for this reason.Kubernetes: polls the Kubernetes API (default every 5000 ms) with a label selector (kubernetes_selector: "app=myapp").kubernetes_ip_lookup_modechooses:endpoints(default) or:pods, and the difference matters in production. Endpoints are readiness-gated, so a pod joins the cluster only after it passes its readiness probe, while:podsmode sees pod IPs as soon as they exist, letting a node join the mesh before it serves traffic. Which one you want depends on whether cluster membership should precede readiness (it usually should, so handoff can happen before the old pod dies). Requires a ServiceAccount allowed to list endpoints or pods; token read from/var/run/secrets/kubernetes.io/serviceaccount. Themodeoption controls name shape::ip(default,basename@10.2.3.4),:hostname(StatefulSet FQDNs), or:dns(pod A records).kubernetes_node_basenamemust match the name your release boots with, which is the classic footgun: the strategy computesbasename@ip, and if yourrel/env.sh.eexsets a different RELEASE_NODE, connects fail silently forever.Kubernetes.DNS: same goal, but via DNS A records of a headless service instead of the API, so no RBAC needed.Rancher: metadata API of a platform you will never mention on stage.
For the talk, the taxonomy is more useful than the full list of seven: static list, local EPMD, LAN multicast, orchestrator API, DNS. Every strategy anywhere fits one of these five shapes.
5. DNSCluster: the Phoenix default, and what it does not do
Section titled “5. DNSCluster: the Phoenix default, and what it does not do”In early 2024 (1.7.11), phx.new started generating clustering support by default, with dns_cluster rather than libcluster. Current version 0.3.1 (released Oct 1, 2026, so check the changelog once more the week of the talk). The generated line:
{DNSCluster, query: Application.get_env(:my_app, :dns_cluster_query) || :ignore}:ignore means it starts as a no-op unless you set the query, so every Phoenix app since then carries a dormant clustering child spec. Set DNS_CLUSTER_QUERY=myapp.internal (Fly.io popularized this) and clustering turns on with zero code.
Mechanics: every interval (default 5000 ms) it resolves the query for A and AAAA records (default resource_types: [:a, :aaaa], SRV also supported), takes each IP, and connects to basename@ip where basename defaults to your own node’s basename. If you run as myapp@10.0.0.7 and DNS returns 10.0.0.8, it dials myapp@10.0.0.8. connect_timeout defaults to 10,000 ms. Multiple queries are allowed, and the tuple form {"otherapp", "otherapp.internal"} connects across basenames, so one BEAM cluster spanning two Phoenix apps is a one-liner.
What it requires: node names must be IP-based and match the records DNS returns, and in Kubernetes the service must be headless (clusterIP: None) so the query returns pod IPs rather than one virtual IP. What it does not do: no API queries, no RBAC, no multicast, no pluggable strategy behaviour, no disconnect logic to speak of. The hexdocs say it themselves: for advanced strategies, use libcluster.
This makes an easy comparison slide. DNSCluster is libcluster’s Kubernetes.DNS strategy extracted into a ~200 line dependency-free library and made the default. The ecosystem’s center of gravity moved from “install a clustering library” to “clustering is a built-in that you switch on with an env var”, which fits the thesis: distribution is baked into the runtime, and now even the discovery glue nearly is. For NETZlive-scale concerns (readiness gating, custom strategies, controlled disconnects), libcluster is still the right tool.
6. :peer: scripted multi-node without tmux gymnastics
Section titled “6. :peer: scripted multi-node without tmux gymnastics”For the companion repo and live demos, the tool is OTP’s peer module (OTP 25+). Its predecessor slave is deprecated and scheduled for removal (OTP 29 per the deprecations page, verify the removal release once more before stage). peer starts additional BEAM nodes as children of your process:
# run: elixir --name ctl@127.0.0.1 --cookie demo cluster_demo.exsnodes = for name <- [:n1, :n2, :n3] do {:ok, _pid, node} = :peer.start_link(%{ name: name, host: ~c"127.0.0.1", longnames: true, args: [~c"-setcookie", ~c"demo"] })
node end
# each peer connected to ctl; global's transitive mesh connects them# to each other within moments. Prove it:Process.sleep(500)for n <- nodes, do: IO.inspect({n, :erpc.call(n, Node, :list, [])})Note the charlists in host and args, it is an Erlang API. Three things make peer better than shelling out to elixir --name ...:
- Lifecycle coupling. With
start_link, the peer dies when your process dies, and losing the control connection kills the peer. No orphaned BEAM processes accumulating after a crashed demo, which matters a great deal on stage. peer:call/4,5gives you RPC over the control connection, independent of distribution. Combined with theconnection: :standard_iooption, you can start a peer with no distribution and no EPMD at all, controlled purely over stdin/stdout, with the peer’s console output relayed through your node. This is the clean way to demo a node that boots undistributed and then becomes distributed, and it connects back to the Day 2 EPMD-less material.- Deterministic startup.
wait_bootblocks until the peer is actually up, so scripts do not race the boot.
Also relevant from yesterday’s curiosity pile: the transitive mesh forming between peers that only ever connected to ctl is the global auto-mesh from section 1, observable in a five-line script. That is the cheapest available demonstration that the mesh is real.
For the repo’s justfile, a just cluster target running a script like the above, plus :erpc.call/4 to install the Cluster GenServer on each peer, replaces the three-terminal dance for everything except the demos where you want visible separate shells (netsplits on Day 6 will want visible shells).
7. Production notes to carry into the talk
Section titled “7. Production notes to carry into the talk”- Discovery is a retry loop, not a source of truth. Every strategy keeps trying to converge the connection set toward the discovered set. There is no registration authority, no membership service, no epoch numbers. If DNS is stale or the k8s API lags, your cluster view lags. The BEAM layer below (
monitor_nodes) is exact about connections, while the discovery layer above is approximate about intent. - The cookie is still the only gate. Every strategy assumes all nodes share the cookie, and connecting is authorizing (full RPC, as Day 7’s security block will cover). libcluster and DNSCluster widen reach automatically, and they widen the blast radius just as automatically.
- Name discipline accounts for most debugging. Nearly every “libcluster doesn’t work” issue is basename mismatch, sname vs name mismatch, or a hostname that does not resolve the same way everywhere. A 10-second
Node.ping/1from a remote shell beats an hour of log reading. - Watch the mesh size. Full mesh plus tick traffic (Day 6) is fine at NETZlive scale and most scales, but say the number: 100 nodes is 4,950 connections, each with heartbeats and a dist buffer. The audience running 3-pod deployments should know the ceiling exists and that they are nowhere near it.
References
Section titled “References”- jola.dev, “Elixir Cluster 101”, the Cluster GenServer pattern and the observability advice.
- Erlang kernel docs,
net_kernel,monitor_nodes/1,2, the delivery-ordering guarantees, and nodedown reasons. - libcluster README (v3.5.0), topology config and strategy list.
- libcluster
Cluster.Strategy.GossipandCluster.Strategy.Kubernetes, defaults quoted above. - DNSCluster hexdocs (v0.3.1), options, defaults, and limitations.
- Erlang stdlib docs,
peer, options, alternative connections, shutdown modes. - Erlang deprecations page,
slavedeprecation status. - ElixirForum, “New Phoenix project, why do we need DNSCluster?”, useful for the “what the ecosystem default means” framing.
Active block (45 to 60 min, repo section 01_clustering)
Section titled “Active block (45 to 60 min, repo section 01_clustering)”Build 01_clustering/cluster_demo.exs: a single script, run via just cluster, that (1) starts three :peer nodes with longnames and a shared cookie, (2) uses :erpc.call/4 to load and start the Cluster GenServer from this briefing on every node (put the module in a .ex file compiled with Code.compile_file/1 on each peer, or :erpc.call(n, Code, :eval_file, [path])), (3) sleeps, then prints each node’s members/0 to show four identical local opinions, (4) calls :peer.stop/1 on one peer and prints the nodedown_reason the survivors logged, and (5) starts a replacement peer and shows it meshing transitively after connecting only to ctl. Stretch goal: rerun step 4 but kill the peer OS process with System.cmd("kill", ["-9", os_pid]) (get the pid via :erpc.call(n, System, :pid, [])) and compare the nodedown reason against the clean stop. You should see the difference between a shutdown and a vanished node, which is the seed of Day 6.
Exit questions
Section titled “Exit questions”- What are the four delivery-ordering guarantees of
monitor_nodesmessages relative toNode.list()and to signals on the connection, and why do they make “subscribe, then snapshot” race-free? - A node in your 6-node cluster logs
{:nodedown, n, [nodedown_reason: :net_tick_timeout]}while another deploy-related disconnect logged:connection_closed. What different physical situations do these two reasons imply, and which one should page somebody? - DNSCluster resolved
myapp.internalto three IPs but the cluster never forms. Name the three most likely causes given how it derives node names, and what singleiexcommand you would run first to discriminate.
Say it out loud (60 seconds, to the Haarlem room)
Section titled “Say it out loud (60 seconds, to the Haarlem room)”“Explain how a BEAM cluster forms, from one Node.connect call to a full mesh, and why the cluster you end up with is each node’s local opinion rather than an agreed-upon fact. Cover what libcluster and the Phoenix DNSCluster default actually do (a retry loop around discovery and connect_node), and finish with why that is enough for reach and why it will never be enough for agreement.”