Two instances of the same Autumn binary can find each other, agree on who is running, and share one distributed primitive — a cluster-wide counter — with no external coordination service. No Redis, no Postgres, no etcd, no ZooKeeper, and no new entry in your dependency tree: the substrate is compiled into the framework and switched on by a config section.
The properties, stated plainly up front: this is an opt-in subsystem
([cluster] enabled = false by default), it is eventually consistent
(nodes converge on a shared view; they never agree instantly), it is
convergent (CRDT) by construction rather than by timing, and its traffic is
authenticated (HMAC-SHA256) but unencrypted. This first slice is a
deliberate 2-node slice: designed, tested, and documented for exactly two
nodes and exactly one shared primitive.
This is not database sharding. Horizontal Sharding also uses the word "cluster" — a Redis-Cluster-style 16384-slot map routing tenants across
[[database.shards]]Postgres databases. That is data placement across databases. This guide is about application processes finding each other over the network. The two features share no configuration keys, no code, and nothing but the unfortunate word.
What you get
- A member view. Each node knows which nodes it currently believes are running, with their advertised address, status, and incarnation.
- A cluster-wide counter.
increment()on one node is observable viaget()on the other, typically within one push interval. - A health component.
cluster:membershipreports the local view on/actuator/health, and the substrate's own metrics land on/actuator/metrics.
What you do not get
- Not mutual exclusion, not leader election. For work that must run on one replica only, use Distributed Locks or the multi-replica scheduler. The cluster counter cannot fence anything.
- Not durable. Counter values live in process memory for the lifetime of the process (see Failure semantics).
- Not encrypted. Authenticated (HMAC) only — deploy on a trusted network.
- Not a replacement for the Redis or Postgres backends behind sessions, jobs, channels, or the scheduler. It is an additional option that needs no datastore, not a substitute for one.
- Not tested past two nodes. The protocol pushes full state to every known peer on every interval; three or more nodes is out of scope for this slice.
Enable it
# autumn.toml
[cluster]
enabled = true
secret = "…" # prefer the env var below
cluster_name = "autumn"
bind_addr = "0.0.0.0:7946"
advertise_addr = "10.0.1.7:7946"
seed_peers = ["10.0.1.8:7946"]
node_id = "web-1"
push_interval_ms = 500
suspicion_timeout_ms = 2500
Every key has an environment override, which is how you supply the secret and the per-instance addresses on a container platform:
AUTUMN_CLUSTER__ENABLED=true
AUTUMN_CLUSTER__SECRET=… # shared by every node
AUTUMN_CLUSTER__CLUSTER_NAME=autumn
AUTUMN_CLUSTER__BIND_ADDR=0.0.0.0:7946
AUTUMN_CLUSTER__ADVERTISE_ADDR=10.0.1.7:7946
AUTUMN_CLUSTER__SEED_PEERS=10.0.1.8:7946 # comma-separated
AUTUMN_CLUSTER__NODE_ID=web-1
AUTUMN_CLUSTER__PUSH_INTERVAL_MS=500
AUTUMN_CLUSTER__SUSPICION_TIMEOUT_MS=2500
Configuration keys
| Key | Type | Default | Meaning |
|---|---|---|---|
enabled | bool | false | Master switch. When false nothing binds, no task spawns, no extension installs. |
secret | secret string | none | Shared HMAC key, at least 16 bytes. A SecretString, never logged. |
cluster_name | string | "autumn" | Logical cluster name. Covered by the MAC, so two clusters cannot mix. |
bind_addr | socket addr | "127.0.0.1:0" | Where the cluster listener binds. Port 0 picks an ephemeral port. |
advertise_addr | socket addr | none | The address peers should dial. Defaults to the bound address. |
seed_peers | list of socket addrs | [] | Addresses to dial on startup. CSV in the env form. |
node_id | string | none | Stable node identity. Defaults to a per-boot random id. |
push_interval_ms | u64 | 500 | Base period of the state push, carrying ±20% per-node jitter. |
suspicion_timeout_ms | u64 | 2500 | Silence after which a peer is locally down and leaves the view. |
Validation
[cluster] is validated at boot and a bad section is a startup error, not a
warning:
enabled = truerequiressecret, and the secret must be at least 16 bytes. There is no lenient unauthenticated mode.suspicion_timeout_msmust be at least3 × push_interval_ms. Below that ratio a single delayed push evicts a healthy peer and the view flaps.push_interval_msmust be at least10.bind_addr,advertise_addr, and every entry ofseed_peersmust parse as aSocketAddr(host:port, IP literal — hostnames are not resolved). Onlybind_addrmay carry port0: there it means "any free port" and the node advertises the port it was actually given, while anadvertise_addror a seed on port0is an address no peer can ever dial.node_idandcluster_name, when set, must be non-empty, at most 64 bytes, and must not contain#(reserved as the separator in counter cell keys). Both travel in every frame and are covered by the MAC.- An unknown key under
[cluster]is a config error, like every other section.
Choosing addresses
bind_addr is a local bind; advertise_addr is what other nodes dial. They
differ whenever the bind is a wildcard or the node sits behind NAT:
- Binding
0.0.0.0:7946requires an explicitadvertise_addr—0.0.0.0is not a dialable address, and a peer that learns it can never reach you. - Both nodes must be dialable at the address they advertise. Replies never
travel back over an accepted inbound socket: each node opens its own
connection to the address it learned from the pushed state, so a node its
peer cannot dial — behind NAT, or on a container network that only publishes
one port — is a node its peer can accept frames from and never send frames
to. That is a one-way cluster, and the asymmetry runs opposite to
intuition: the reachable node converges on a two-member view, because it
keeps accepting the unreachable node's pushes and learning its record, while
the unreachable node receives nothing back and stays at a one-member view
forever — even though it is the one that dialed first.
seed_peersdecides who makes first contact, not who has to be reachable. - Leaving
bind_addrat the default ephemeral port is therefore fine only when the peer can reach the resolved port — same host, same pod network — because that resolved port is what gets advertised. Give each node a fixed port and a dialableadvertise_addron anything more complicated. seed_peersis a dial list, not a membership list. It is used to make first contact; after that each node learns peers (and their advertised addresses) from the pushed state. Seeding one direction is enough — node B pointed at node A produces a two-member view on both.
Using the counter
The cluster installs a ClusterHandle as an app extension when it is enabled,
so the lookup is an Option and your handler decides what a disabled cluster
means:
use autumn_web::cluster::ClusterHandle;
use autumn_web::extract::State;
use autumn_web::prelude::*;
#[post("/sightings")]
async fn record_sighting(State(state): State<AppState>) -> AutumnResult<Json<serde_json::Value>> {
let cluster = state
.extension::<ClusterHandle>()
.ok_or_else(|| AutumnError::internal_server_error_msg("cluster is not enabled"))?;
// Grow-only, cluster-wide. Applied locally right away, pushed to peers on
// the next interval.
let sightings = cluster.counter("boids_sighted");
sightings.increment();
Ok(Json(serde_json::json!({
"node": cluster.node_id(),
"members": cluster.members().len(),
"boids_sighted": sightings.get(),
})))
}
#[get("/sightings")]
async fn read_sightings(State(state): State<AppState>) -> AutumnResult<Json<serde_json::Value>> {
let cluster = state
.extension::<ClusterHandle>()
.ok_or_else(|| AutumnError::internal_server_error_msg("cluster is not enabled"))?;
// Read-only: no increment. This is the endpoint the two-terminal
// walkthrough polls on node B after incrementing on node A.
Ok(Json(serde_json::json!({
"node": cluster.node_id(),
"boids_sighted": cluster.counter("boids_sighted").get(),
})))
}
| Method | Returns | Notes |
|---|---|---|
ClusterHandle::node_id() | &str | This node's id — configured, or the per-boot random one. |
ClusterHandle::local_addr() | SocketAddr | The address actually bound (resolved port). |
ClusterHandle::members() | Vec<ClusterMemberInfo> | The local view, including this node. |
ClusterHandle::counter(name) | ClusterCounter | Cheap handle; call it per request. |
ClusterCounter::increment() | () | Adds 1 to this node's own entry. Synchronous, never fails. |
ClusterCounter::get() | u64 | Saturating sum of every entry this node has merged. |
ClusterMemberInfo carries { id, addr, status, incarnation }, where status
is Alive or Suspect — a member this node considers down is not in the view
at all.
get() may jump upward between two calls with no local increment in
between: that is a merge from a peer landing. It never moves downward. Treat it
as a lower bound on the true cluster-wide total, never as a limit to enforce.
Verify with two terminals
One binary, two processes, one shared secret. Both terminals need the same secret, so export it first in each:
export AUTUMN_CLUSTER__ENABLED=true
export AUTUMN_CLUSTER__SECRET="dev-only-cluster-secret-change-me"
Terminal A — the node that will be seeded against:
AUTUMN_SERVER__PORT=3000 \
AUTUMN_CLUSTER__NODE_ID=node-a \
AUTUMN_CLUSTER__BIND_ADDR=127.0.0.1:7946 \
cargo run
Terminal B — same binary, different ports, seeded at A:
AUTUMN_SERVER__PORT=3001 \
AUTUMN_CLUSTER__NODE_ID=node-b \
AUTUMN_CLUSTER__BIND_ADDR=127.0.0.1:7947 \
AUTUMN_CLUSTER__SEED_PEERS=127.0.0.1:7946 \
cargo run
Within a couple of push intervals both nodes report the same two members:
curl -s localhost:3000/actuator/health | jq '.components["cluster:membership"]'
curl -s localhost:3001/actuator/health | jq '.components["cluster:membership"]'
The details object renders only when health details are enabled
(health.detailed = true) — the dev profile's default, which is what cargo run uses here. With details off you still get the component's status.
{
"status": "UP",
"details": {
"node_id": "node-a",
"cluster": "autumn",
"local_addr": "127.0.0.1:7946",
"member_count": 2,
"members": [
{ "id": "node-a", "addr": "127.0.0.1:7946", "status": "alive", "incarnation": 1765430001204 },
{ "id": "node-b", "addr": "127.0.0.1:7947", "status": "alive", "incarnation": 1765430017861 }
]
}
}
Increment on A, read on B:
curl -s -X POST localhost:3000/sightings # {"boids_sighted":1,…} on node-a
curl -s localhost:3001/sightings # 1 on node-b within a push interval
Now stop node B and watch node A converge to a one-member view. Ctrl-C takes
the clean path (B sends a leave, A drops it inside a few hundred milliseconds);
kill -9 takes the failure path (A drops B after suspicion_timeout_ms).
Either way:
curl -s localhost:3000/actuator/health \
| jq '.components["cluster:membership"] | {status, count: .details.member_count}'
# { "status": "UP", "count": 1 }
A one-member view is healthy. The indicator registers in the HealthOnly
group (see Health Indicators), so it never gates
/ready and never reports DOWN because a peer went away — an orchestrator
cannot be tricked into restarting the survivor.
/actuator/metrics exposes the same facts as counters, which is what you
alert on:
| Metric | Kind | Notes |
|---|---|---|
autumn_cluster_members | gauge | Size of the local view. |
autumn_cluster_pushes_sent_total | counter | State pushes written. |
autumn_cluster_pushes_unsendable_total | counter | Outbound messages that could not be signed or framed, so the transport never saw them. Anything above zero is almost always a document past the 64 KiB frame cap — see State growth. |
autumn_cluster_pushes_received_total | counter | State pushes accepted from peers after verification. A leave is verified and applied without being counted here. |
autumn_cluster_merges_applied_total | counter | Merges that changed local state. |
autumn_cluster_frames_rejected_total | counter | Labelled reason — see the receive path below. Every label is published from boot, zeroes included, and reason="oversize" covers the connection-fatal step-1 rejections too. |
autumn_cluster_frames_dropped_total | counter | Outbound frames swallowed by a full or closed peer queue. Anti-entropy re-sends the state, so drops are lossy only to latency — with one exception, the departure, which travels on a lane of its own and is warned about when it cannot be handed over (Leave is advisory). |
A steady frames_rejected_total{reason="mac"} means somebody is talking to
your port with the wrong secret. Any pushes_unsendable_total at all means
this node has stopped gossiping and its peer is about to evict it — the push is
also the heartbeat, so that failure is otherwise invisible from the inside.
How it works
The wire format
One TCP listener per node. Connections are byte pipes and nothing more: connection state carries zero liveness meaning. A live connection with no recent state push is a silent peer; a dropped connection that reconnects before the suspicion timeout is a non-event.
Each frame is a 4-byte big-endian length prefix followed by exactly that many bytes of JSON:
+-------------------+--------------------------------------+
| u32 big-endian N | JSON envelope, exactly N bytes |
+-------------------+--------------------------------------+
N is capped at 65536 (64 KiB) and the cap is checked before any
allocation. A prefix of 0 or greater than the cap closes that connection
without reading further. The sender applies the same cap and refuses to emit an
oversized frame rather than sending something the peer must reject.
The envelope:
| Field | Type | Meaning |
|---|---|---|
v | u8 | Envelope version. 1 in this slice; any other value is dropped. |
key_id | u8 | Which signing key made the MAC. Always 0; reserved for key rotation. Anything else is dropped, and so is an envelope that omits the field. |
cluster | string | Sender's cluster_name. A mismatch with the local name is dropped. |
sender | string | Sender's node id — the authenticated identity. The source address is never an identity. |
incarnation | u64 | Sender's incarnation when the frame was produced. |
seq | u64 | Per-sender counter, incremented per frame, reset to 0 when the incarnation increases. |
payload | string | The inner message as a JSON string. Opaque until the MAC verifies. |
mac | string | Lowercase hex HMAC-SHA256, 64 characters. |
The MAC is computed over a length-delimited concatenation, so no field value can be shifted into another field:
signing_input = L(v) ‖ L(cluster) ‖ L(sender) ‖ L(incarnation) ‖ L(seq) ‖ L(payload)
L(x) = 8-byte big-endian byte length of x, followed by x
v → 1 byte, the u8 value
cluster → UTF-8 bytes
sender → UTF-8 bytes
incarnation → 8 bytes, big-endian
seq → 8 bytes, big-endian
payload → UTF-8 bytes of the payload string
key_id is deliberately outside the signing input: it selects the key
rather than being protected by it. Flipping it selects a key that does not
exist, the MAC then fails, and the frame is dropped — a tampered key_id can
cost an attacker a dropped frame, never a forged one.
Receive path
The order is normative. Every rejection logs and increments
frames_rejected_total with the named reason, and nothing ever panics. Steps
2-7 drop the offending frame and continue the read loop. Step 1 is the one
deliberate exception: a bad length prefix means the stream framing itself can
no longer be trusted — there is no way to know where the next frame starts —
so the connection closes and the peer re-dials with backoff. Closing a
connection is never an eviction; connection state carries zero liveness
meaning.
Step 4 is a second, quieter boundary: everything at or before it refuses unauthenticated bytes, and each such refusal is charged to the connection the frame arrived on — three of them and that connection is closed too (see Failure semantics). Rejections past step 4 are never charged, because the frame carried a valid MAC.
- Length prefix —
N == 0orN > 65536closes the connection with no allocation (reason="oversize"). - Envelope parse — malformed JSON, or an envelope missing a field, is
dropped (
reason="malformed"). Only the fixed envelope is parsed here. - Header checks —
v != 1(reason="version"),key_id != 0(reason="key_id"),cluster != cluster_name(reason="cluster"). - MAC — recompute over the signing input and compare in constant time.
A mismatch drops the frame (
reason="mac"). Nothing below this line runs on unauthenticated bytes. - Self-origin —
sender == local node idis dropped (reason="self_origin"). This is decided on the authenticated field, never on the peer address, so a seed list containing your own address and a reflected frame are both handled by the same rule. - Replay watermark — per sender, the node remembers the highest
incarnation it has accepted and the highest
seqaccepted at that incarnation. A frame is dropped (reason="replay") when itsincarnationis lower than the recorded one, or equal to it withseqnot greater than the recordedseq. A higherincarnationadopts the new incarnation and resets the sequence watermark. - Payload parse — only now is
payloaddeserialized. An unknown message type or malformed body is dropped (reason="payload"). - Merge — the message is applied, and receipt is recorded against the sender for the liveness overlay.
Messages
Two message types share the envelope. They are internally tagged so an unknown future variant is a clean drop rather than a parse failure:
{"type":"state_push","state":{
"members": {
"node-a": { "addr": "127.0.0.1:7946", "incarnation": 1765430001204, "status": "alive" },
"node-b": { "addr": "127.0.0.1:7947", "incarnation": 1765430017861, "status": "alive" }
},
"counters": {
"boids_sighted": { "node-a#1765430001204": 3, "node-b#1765430017861": 2 }
}
}}
{"type":"leave"}
state_push carries the entire replicated document and doubles as the
heartbeat — there is no separate ping. leave carries no fields at all: it
applies to the (sender, incarnation) pair in the authenticated envelope, so a
captured leave can never be replayed against a newer incarnation of that node.
Wire structs tolerate unknown fields, and the structs inside the payload
default missing ones, so a newer node can add fields without breaking an older
peer. (This is the opposite of the config layer, where an unknown key is an
error.) The envelope itself is the exception: every one of its fields is
required, and one that is absent is a step-2 malformed drop — an envelope is
what authenticates a frame, so there is nothing in it to be lenient about.
Membership
Membership is two things that are easy to confuse, so they are separated in the types:
Replicated status is part of the shared document and converges. A
MemberRecord is { addr, incarnation, status } where status is only
Alive or Left. Records merge pairwise:
- Higher
incarnationwins. - At equal incarnation,
LeftbeatsAlive— a leave is never undone by an in-flight older push at the same incarnation. - At equal incarnation and equal status, the lexicographically greater
addrwins. This tie-break exists only to keep the merge commutative; it should never fire in practice.
Incarnation is seeded at boot from the wall clock (Unix milliseconds,
read through the injected clock), so a restart — clean or crash — comes back
at a strictly higher incarnation with no persistence anywhere. Millisecond
granularity is load-bearing: no real process restarts within the same
millisecond, so in practice two boots do not mint equal incarnations — which is
what keeps refutation's exact-echo check sound (a record byte-equal to this
boot's own record really is an echo of this boot, rather than a leftover from
the previous one). That is what lets a node with a stable node_id rejoin
after a crash: its peer's replay watermark is keyed by (sender, incarnation),
and a fresh, higher incarnation starts a fresh sequence rather than colliding
with the dead boot's watermark.
That argument is a probability, not a guarantee. A clock that is frozen,
pre-epoch (the reading clamps to 0), or stepped backwards onto the exact
millisecond a previous boot read will mint an equal incarnation, and the two
boots then share a counter cell and a replay watermark until an operator
restarts one of them. Run the cluster on hosts with a working clock; the
millisecond seeding makes the collision a coincidence, not a routine
same-second restart.
Refutation keeps a live node from being buried, and covers the residual
restart cases (a restart within the same clock millisecond, or a clock that
stepped backwards). A node that receives any record about itself at an
incarnation greater than or equal to its own — whatever the status, Left or a
stale Alive — sets its incarnation to one above the highest it has seen for
itself, marks itself Alive, and pushes immediately. Because rule 1 outranks
rule 2, that refutation wins everywhere. The trigger works even when the
returning node is being replay-dropped by its peer, because the peer's own
pushes still reach it: the returning node's watermark table is fresh, so it
accepts the peer's state, sees its own stale record, and bumps past it —
provided the peer still holds a record about it at all. Once that record has
been pruned there is nothing left to refute, which is why pruning a member
forgets it whole: the peer drops the pruned node's replay watermark along
with its tombstone, so the returning node's next frame is judged as a fresh
sender at whatever incarnation it carries. Its record still waits out the
peer's recently-pruned note when it returns at an incarnation that peer already
collected — which only a backward clock step produces — so that rejoin is
adopted when the note lapses instead of on the first frame. A boot with a
working clock carries a higher incarnation and is adopted at once.
Local liveness is not replicated. Each node times the silence since it last accepted a frame from each peer, using the injected monotonic clock:
accepted a frame accepted a frame
│ │
▼ ▼
┌─────────┐ silence > 2×push ┌─────────┐ silence > suspicion ┌──────┐
│ Alive │────────────────────▶│ Suspect │──────────────────────▶│ Down │
└─────────┘ └─────────┘ └──────┘
in view in view not in view
Suspect at twice the push interval is a warning, not an eviction: with ±20%
jitter a healthy peer's gap never approaches it, but two missed pushes are
worth showing an operator. Down at suspicion_timeout_ms removes the member
from the view. Because validation forces suspicion_timeout_ms ≥ 3 × push_interval_ms, Suspect always strictly precedes Down.
The view is therefore local: replicated Alive records, minus peers this
node currently considers down, plus this node itself. Two nodes can briefly
disagree about the view — that is exactly what eventually consistent means
here, and it is why the counter's correctness does not depend on the view.
Leaving has a fast path and a correctness path, and only the second one is
a contract. On shutdown a node marks its own record Left at its current
incarnation and sends its final document — one last state_push, carrying
both the departure and every counter cell written since the previous push —
followed by the leave, together bounded to 250 ms so they always fit
inside the shutdown drain budget. The final push is what keeps an increment
that arrived after the last push round from dying with the process; the leave
that follows costs one frame and covers the peer that holds no record for this
node at all. Both frames take a per-peer departure lane that the writer reads
ahead of anything already queued, so a peer stalled long enough to fill its send
queue cannot swallow them (see Failure semantics). That is
the fast path, and it is best-effort.
The departure runs after in-flight HTTP requests have drained, not when the listener closes: a request served during shutdown can still increment a counter, and it must find a push loop still running. If the process is killed, the network eats the frame, or the peer is mid-reconnect, the peer still converges — by silence, at the suspicion timeout. Never build anything on a leave arriving.
Records for departed members stay in the document as tombstones so Left
propagates, and are pruned locally after ten suspicion timeouts — measured from
when this node first observed the departure, so a tombstone re-taught by a
lagging peer after it was pruned ages from that re-observation and expires
again rather than living forever.
A tombstoned member also stays in the push target set until its tombstone is
pruned. That is deliberate and costs a few frames aimed at a dead address: it
is the only channel by which a node that restarted with a lower incarnation
(a clock that stepped backwards) can hear its own Left record, refute it, and
rejoin — its pushes are being replay-dropped by the peer that holds the
tombstone, so it has to be told.
Pruning that tombstone therefore drops the departed node's replay watermark
with it. Once the record is gone there is no Left to hear and no address to
push to, so a watermark that outlived it would replay-drop the returning node's
every frame with no way left to argue — a permanent partition between two nodes
that both look healthy. A member is forgotten whole or not at all.
An Alive record needs an exit of its own, because only Left records are ever
pruned. A record nothing refreshes — a peer that vanished without a leave, a
member learned from a peer's document and never heard from, a record re-admitted
by a replayed frame — would otherwise stay in the document for the life of the
process. So a member this node has not heard from for a full tombstone window
(the same ten suspicion timeouts) is recorded as Left at that member's current
incarnation, and the ordinary tombstone lifecycle takes it from there. That is a
local decision written into the replicated document, and it says exactly what
the suspicion timeout already says — this node considers that member gone —
one whole tombstone window after the view said it. A live member recorded that
way (a partition that outlasted the window) has the standing answer: tombstones
stay push targets, so it hears the Left record about itself and refutes at a
higher incarnation, which wins everywhere by rule 1.
The counter
ClusterCounter is a grow-only CRDT counter, the simplest structure that makes
"observable on the other node" a property of the data rather than of the
network's timing:
- Each counter name maps to a map of
cell → u64, where a cell is keyed bynode id # incarnation— one cell per node per boot. A node writes only its own current cell, so no two writers ever share a cell, and a restarted node never resumes a cell an older boot of itself already populated. Without the incarnation in the key, a node restarting with a stablenode_idwould start its cell at zero while its peer remembers the old, higher value; the max-merge would then silently absorb every new increment until the new count overtook the old one. Keying cells by boot removes that failure by construction — no recovery step, no readiness gate. increment()adds 1 to this node's current cell, immediately and locally, then nudges the push task. Localget()reflects it at once. The nudge is rate-limited to one prompt push per quarter ofpush_interval_ms(never under 50 ms, never longer than the interval itself): the first write after a quiet gap still propagates promptly, and a handler incrementing on every request propagates at that floor instead of signing and sending the whole document once per request.- Merge is per-entry maximum. That is commutative, associative, and idempotent, so pushes may arrive out of order, be duplicated, or be dropped entirely, and the result is the same once any later push lands.
get()is the saturating sum of every entry.
There is no decrement. A grow-only counter is what makes the merge total and
the arithmetic auditable; a decrement map is a possible later addition and the
wire format leaves room for it.
Saturation, not wrapping. A per-node entry saturates at u64::MAX instead
of wrapping, and get() saturates the sum. An implausible count clamps; it
never wraps to zero and never panics.
Three rules, three mechanisms. A boid flock keeps formation with alignment, cohesion, and separation, and this substrate is built from the same three moves under different names. Alignment is
ClusterState::merge— each node steers its view toward the view of the neighbours it can see. Cohesion is the periodic state push to every known peer, which is the only thing keeping the flock from drifting apart. Separation is the ±20% jitter on each node's push interval, which stops identically configured nodes from flying in lockstep and colliding on the wire. The real flocking math, with actual vectors, lives inexamples/island-flock.
Failure semantics
Partition. Both sides keep accepting increments and keep serving their own view; neither side is fenced and no write is rejected. During the split each side reports a total missing the other side's contributions. When the partition heals, the next push in each direction merges both documents and both nodes converge on the full sum. Nothing is lost, because the counter only grows.
Replay. A captured frame replayed later is dropped by the per-sender
sequence watermark, and a replayed leave can only ever name the incarnation
it was signed at. If it names a stale incarnation, it loses the merge; if it
somehow lands, the live node refutes at a higher incarnation and re-appears.
Frames are authenticated but not timestamped, so replay protection is
per-sender-sequence, not clock-based.
The one replay a watermark cannot catch is a frame captured before a member
departed and replayed after that member's record was pruned: pruning forgets the
sender whole, watermark included, so the frame verifies again as a fresh sender.
Two bounds answer it, and neither needs the secret to have stayed secret. While
this node still remembers collecting that id — twice the tombstone window — the
record is refused whatever its status, Left or Alive. Once that memory
lapses and the record is re-learned, it is a member nothing refreshes, so a full
tombstone window of silence records it as Left and it prunes like any other
tombstone. Replaying one captured frame per departed id therefore costs bounded
churn over one window, never a document that ratchets upward toward the frame
cap.
Authenticated (HMAC), unencrypted. Every frame is authenticated with
HMAC-SHA256 over the shared secret, and unauthenticated bytes never reach a
payload parse. Frames are not encrypted: anyone who can read the wire can
read your node ids, advertised addresses, counter names, and counts, and anyone
who can reach the port can make you compute one HMAC per frame. Run the cluster
port on a trusted network — a private VPC subnet, a WireGuard or Tailscale
interface, a container network — and never expose it to the internet. Do not
put anything sensitive in a counter name. Key rotation is not in this slice:
key_id is reserved for it, and changing the secret today means restarting
both nodes.
What an unauthenticated connection can cost you. Reaching the port is enough to open a socket, so the listener bounds what a socket is worth before anyone proves they know the secret. Three bounds, and the third is what makes the first one hold:
- At most 128 inbound connections are held open at once; past that a new one is accepted and closed immediately.
- A connection that has not delivered a complete frame within four suspicion timeouts (at least 10 seconds) is closed.
- A connection whose frames keep being refused before they authenticate is
closed on the third such frame. The receive path's MAC check (step 4) is
the line:
oversize,malformed,version,key_idandmacare refusals of bytes anyone who can reach the port could have sent, and each one is charged to the connection it arrived on. Verdicts past the MAC —self_origin,replay,payload— are never charged, because the frame came from a holder of the secret. The count does not decay.
The idle deadline alone bounds a silent socket, not a talkative one: without
the third bound a stranger holds a slot for as long as it likes by sending one
well-framed garbage frame per idle window, and holds all 128 the same way — at
which point your actual peer is the one being refused at the cap. Three frames
rather than one is deliberate: a peer that is misconfigured, or caught halfway
through a secret rotation, gets far enough onto the wire to be diagnosed from
frames_rejected_total{reason="mac"} and the rejection log line before its
socket is taken away.
None of these can fire on a healthy peer — one peer needs one connection, a
peer silent past its suspicion timeout is already out of the view, and a peer
whose frames verify is never charged for one. Note what the third bound does
not cover: frames are authenticated but not encrypted, so somebody who can
capture a peer's frame off the wire can replay it to hold a connection open
(reason="replay") indefinitely. That is inside the trust boundary this design
already assumes, and it is the same reason none of this is a substitute for
keeping the port off the public internet.
Leave is advisory, suspicion is the contract. The 250 ms departure (final
state_push plus leave) is a latency optimization for clean shutdowns.
Correctness comes from the suspicion timeout, which handles the kill, the
crash, the pulled cable, and the lost leave identically. A shutdown that
overruns its drain budget is a kill as far as the peer is concerned.
Both departure frames travel on a small per-peer lane of their own, which the
peer's writer reads before anything queued: a farewell supersedes the state
pushes already queued for that peer, so a peer stalled long enough to fill its
64-deep send queue cannot swallow this node's last words along with them. That
is safe because every push carries the whole document — the queued frames are
older copies of the one the farewell carries, at lower sequence numbers, and the
peer's replay watermark discards whichever of them still arrive. When even that
lane cannot take the frame, the departure is counted in
autumn_cluster_frames_dropped_total and warned about on the way out rather
than being lost quietly: a dead writer still means the suspicion path.
Known residuals (slice 2). Four bounded consequences of the above, none of them correctness:
- Pushes still queued behind a farewell are written after it and rejected by the
survivor's replay watermark, so
frames_rejected_total{reason="replay"}can move by up to one queue's worth at a clean departure of a stalled peer. - The departure flush waits for every queued frame, not just the farewell, so a departing node whose queue had backed up normally spends its whole 250 ms budget and logs the "not fully flushed" line even though its farewell left first.
- The lane guarantees the farewell reaches the writer; it does not guarantee the writer gets to write it. A writer parked in a dial or a write when shutdown begins is cancelled at the end of the same 250 ms grace, and the peer converges on the suspicion timeout.
- A stable-id node that restarts with a backward clock and a different advertised address rejoins only after the tombstone and recently-pruned windows expire: the incumbent's refutation pushes go to the departed boot's old address, and the returning node's own pushes are dropped as replays before the payload carrying its new address is read. The delay is bounded — the windows expire and the node is then admitted as a fresh sender — and the clean answer, an authenticated recovery exchange addressed to the frame's source, is the next slice's work.
Counter volatility. Counters live in process memory for the lifetime of the process. There is no persistence. A restarted node starts a fresh cell (its cells are keyed by boot incarnation) and re-learns every older cell — including its own previous boots' — from its peer, so a rolling restart of one node at a time preserves the total and new increments count from the first request; but if every node is down at once, the count is gone. If you need a durable count, write it to the database (see Counter Caches) and use this for ephemeral, cluster-wide tallies.
State growth. The document holds one member record per distinct node id
ever seen, one counter cell per (node id, boot incarnation) that has
incremented each counter name, and each node keeps one replay-watermark row
per distinct sender it has accepted a frame from. Member records are pruned
after their tombstone window, and the pruned member's watermark row goes with
it — a member is forgotten whole, which is what lets it come back (see
Refutation above). A node does keep one small local note per pruned member —
the (node id, incarnation) it collected — for twice the tombstone window, and
refuses to re-adopt any record for that member at that incarnation or lower
while the note lasts; without it two peers pruning at different times re-teach each other the
same Left record indefinitely and departed ids never actually leave. The note
is local garbage collection, never gossiped, and never applied to a higher
incarnation, so a rejoining node is unaffected. An Alive record has an exit
too, and it is what bounds the member table: a member this node has not heard
from for a full tombstone window is recorded Left at its current incarnation
and then prunes on the ordinary schedule (see Membership above). That covers
every record no leave will ever arrive for — a peer that vanished, a member
learned from a document and never heard from, or one re-admitted by replaying a
captured pre-departure frame after its watermark was forgotten — so each such
record costs one window of churn rather than a permanent row. Every member
record either keeps being refreshed or leaves. Counter cells are not
pruned: dropping one would make
get() move downward. The consequence is operational: every boot that
increments a counter leaves one u64 cell behind, and a node with no
configured node_id additionally mints a fresh member id per restart. Set an
explicit, stable node_id on any long-lived deployment, and keep counter names
a small fixed set from your code rather than anything derived from user input.
Cell garbage collection is out of scope for this slice.
A document that outgrows the 64 KiB frame cap stops being sendable at all —
and since the state push is the heartbeat, the node then keeps merging its
peer's pushes and looks healthy to itself while the peer sees silence and
evicts it at the suspicion timeout. That is what
autumn_cluster_pushes_unsendable_total is for: it counts every message the
node could not frame, and a warning naming the serialized size against the cap
is logged (at most once a minute) the moment it starts.
A one-member view is healthy. cluster:membership reports UP with one
member, exactly as with two. It is a HealthOnly indicator: it never gates
/ready, and it never reports DOWN because a peer disappeared. A surviving
node keeps serving its counter and its traffic, and a liveness probe must never
be able to kill it for being alone.
Two nodes. Every push carries the full document to every known peer, there are no indirect probes, and there is no quorum anywhere. That is a sound design at two nodes and an unproven one beyond that. Treat larger fleets as future work, not as a supported configuration.