Three copies of the same binary, each told where the other two are, make a Queen cluster. They keep one replicated log between them, so a write is on a majority of disks before anyone hears yes, and any of them can serve any client. One node can be down for a crash or an upgrade while the other two go on committing. The quickest way to see it is the Compose file in the repository, which runs the whole cluster on a laptop.
Try it on a laptop
deploy/compose/three-node/ starts three nodes on a private network, each with its own volume.
Node N answers on localhost:N6632 (16632, 26632 and 36632), API and dashboard alike, which keeps
the cluster clear of a single node on 6632. Give it a token, start it, and ask each node what it
is:
cd deploy/compose/three-node # in a clone of github.com/queen-mq/queen
echo "QUEEN_RAFT_TOKEN=$(openssl rand -hex 32)" > .env
docker compose up -d --wait
for n in 1 2 3; do echo "queen-$n $(curl -s localhost:${n}6632/health | jq -r .raft.role)"; donequeen-1 follower
queen-2 follower
queen-3 leaderThe three nodes read the same peer list on empty volumes, voted, and node 3 won. --wait returns
once every node’s healthcheck passes, which on an Apple-silicon laptop, running the amd64 image
under emulation, took eight seconds. Now push through node 1, a follower, pop through node 2 and
ack through node 3:
curl -s -X POST localhost:16632/api/v1/push -H 'content-type: application/json' -d '{
"items": [{ "queue": "orders", "partition": "customer-42",
"transactionId": "order-1001-created",
"payload": { "orderId": 1001, "total": 42 } }]
}'
POP=$(curl -s "localhost:26632/api/v1/pop/queue/orders?batch=1&wait=true&timeout=5000")
echo "$POP" | jq -c '.messages[0] | {transactionId, partition, data}'
curl -s -X POST localhost:36632/api/v1/ack -H 'content-type: application/json' -d "$(echo "$POP" | jq '{
transactionId: .messages[0].transactionId, partitionId: .messages[0].partitionId,
leaseId: .leaseId, status: "completed" }')"[{"index":0,"message_id":"01a0fcd6-2243-7000-a59b-d7da302543ea","transaction_id":"order-1001-created","queueName":"orders","status":"queued","offset":0}]
{"transactionId":"order-1001-created","partition":"customer-42","data":{"orderId":1001,"total":42}}
[{"index":0,"transactionId":"order-1001-created","success":true,"error":null,"leaseReleased":true,"dlq":false,"noop":false}]Each node answered as the leader would have. Now stop the leader:
docker compose stop queen-3
for n in 1 2; do echo "queen-$n $(curl -s localhost:${n}6632/health | jq -r .raft.role)"; donequeen-1 follower
queen-2 leaderdocker compose stop sends SIGTERM, and a leader that receives one hands its leadership to the
most caught-up follower before it exits. Node 2 was leading within a second of the stop, and node
3’s log says how long the hand-off itself took:
INFO rsm: raft: leadership handed off to=2 ms=121Two nodes of three are a majority, so pushes, pops and acks go on through nodes 1 and 2 as before. Start node 3 again and it catches up from the new leader:
docker compose up -d --wait queen-3
curl -s localhost:36632/health | jq -c .raft{"role":"follower","leader":true,"term":3,"applied":13,"commit":13,"lag":0,"storageReady":true,"clusterVersion":3,"kinds":3,"apply":{"failure":null,"skipped":[]}}applied equal to commit, the same 13 the leader reports, means node 3 has applied everything
the cluster committed, the writes it missed included. docker compose down -v removes the
containers and the three volumes when you are done.
This is the whole file. The comments give the reason for each line that is not obvious:
# Three Queen nodes on one machine: a real Raft cluster for trying replication,
# leader failover and rolling restarts. Walkthrough in README.md next to this
# file and at https://queenmq.com/operate/cluster/.
#
# echo "QUEEN_RAFT_TOKEN=$(openssl rand -hex 32)" > .env
# docker compose up -d --wait
# curl -s localhost:16632/health
#
# Node N answers on localhost:N6632 (16632, 26632, 36632), so the cluster does
# not collide with a single node on 6632. Each node keeps its state in its own
# volume; `docker compose down -v` deletes them.
name: queen-cluster
x-node: &node
# QUEEN_VERSION in .env picks another release (README.md: rolling upgrade).
image: ghcr.io/queen-mq/queen:${QUEEN_VERSION:-latest}
# The images are published for linux/amd64 only. On an arm64 machine
# (Apple silicon) Docker runs them under emulation: slower, fine for a trial.
platform: linux/amd64
# A node that was away long enough to need a snapshot exits with code 75 to
# load it at start; the restart policy brings it back.
restart: unless-stopped
# On SIGTERM a node hands its ephemeral partitions and, if it leads, its
# leadership to a peer, then lets the long polls in flight finish (30 s by
# default). Docker's 10 s default would kill it halfway.
stop_grace_period: 60s
healthcheck:
# 200 once the node knows a leader and has caught up, 503 "settling" before.
test: ["CMD", "curl", "-fsS", "-o", "/dev/null", "http://localhost:6632/health"]
interval: 5s
timeout: 3s
retries: 3
start_period: 60s
x-env: &env
QUEEN_RAFT_REPLICATOR: openraft
# Every voter as id=raft_address/http_address, the same list on every node.
QUEEN_RAFT_PEERS: "1=queen-1:7400/queen-1:6632,2=queen-2:7400/queen-2:6632,3=queen-3:7400/queen-3:6632"
# The secret every Raft RPC carries. The raft port (7400) is not published, so
# only these three containers reach it.
QUEEN_RAFT_TOKEN: ${QUEEN_RAFT_TOKEN:?put QUEEN_RAFT_TOKEN in .env next to compose.yaml, see the header}
QUEEN_RAFT_DIR: /var/lib/queen/raft
services:
queen-1:
<<: *node
hostname: queen-1
environment:
<<: *env
QUEEN_RAFT_NODE_ID: "1"
# The broker port has no authentication of its own: published on loopback.
ports: ["127.0.0.1:16632:6632"]
volumes: ["queen-1:/var/lib/queen/raft"]
queen-2:
<<: *node
hostname: queen-2
environment:
<<: *env
QUEEN_RAFT_NODE_ID: "2"
ports: ["127.0.0.1:26632:6632"]
volumes: ["queen-2:/var/lib/queen/raft"]
queen-3:
<<: *node
hostname: queen-3
environment:
<<: *env
QUEEN_RAFT_NODE_ID: "3"
ports: ["127.0.0.1:36632:6632"]
volumes: ["queen-3:/var/lib/queen/raft"]
volumes:
queen-1:
queen-2:
queen-3:How it works
A cluster runs one Raft group over one log, with the openraft crate doing the protocol. The
leader puts every write into that log in one order and replicates it, and an entry commits once a
majority of the voters have written it and fsynced it. Each node then applies the same entries in
the same order to its own copy of the state: the queue logs, the cursors and leases, the KV store,
the timers. There is no second store beside the log, which is why a node that was away needs only
the entries it missed (or a snapshot, if it was away long) to be whole again.
deploy/compose/three-node/compose.yamlWe keep a tenant inside one Raft group because then a transaction is a single entry. An ack, three pushes and a KV write go into the log together and commit together, or none of them does, and no protocol between logs has to agree about it.
Followers carry their share of the work. A follower parses the requests of its own clients, prepares each one as a command, and sends the commands to the leader over a few long-lived streams on the raft port; it answers its client once the entry has committed. A pop is answered once this node has applied the claim, because the messages it hands out come from its own copy of the queue log. A read first waits until the node has applied everything the cluster had committed when the read began, so no client reads behind a write it was already told about, whichever node it asks. That is what made the walkthrough above work, and it is why a client can hold every node’s URL without caring which one leads:
const queen = new Queen({ urls: ['http://queen-1:6632', 'http://queen-2:6632', 'http://queen-3:6632'] })With several URLs the JavaScript client spreads its requests over the nodes and moves to the next one when a node fails with a network error or a 5xx.
A stop is planned work, so we made it cheap. On SIGTERM a node first hands its ephemeral partitions to their next owners. If it leads, it then lets the acks already in flight commit (at most 500 ms) and transfers leadership to the most caught-up voter, waiting up to 3 s for it to take over. Only then does it close its listener and finish the requests in flight, as a follower. A rolling restart therefore costs the cluster one transfer and a vote, about a tenth of a second in the run above.
What each failure costs
| Event | What the cluster does |
|---|---|
| A follower stops or crashes | The other two keep committing, and its clients move to another URL. When it starts again it catches up from the leader. |
| The leader gets SIGTERM | It hands leadership over first: under a second to a new leader in the run above. |
| The leader crashes (kill -9, OOM, a lost host) | The followers wait out the election timeout (1 to 2 s, QUEEN_RAFT_ELECTION_MS) and elect a new leader: 3.7 s in the Compose cluster under emulation, at most 4 s in the recovery rehearsals. |
| A leader hears from no quorum for 3 s | It hands leadership to the most caught-up voter, so a leader that can send but not receive does not hold every write hostage. |
| A majority is down | Nothing commits. Requests wait out their deadline and fail with 503, code no_leader or timeout, and Retry-After: 1, until a majority is back. |
A node was away more than 10 minutes (QUEEN_RAFT_PURGE_HOLD_S, 600) |
The leader stops holding its log for it. If the entries it missed are purged, the node receives a snapshot, exits with code 75 and loads it on its next start, and a restart policy brings it back. |
| A majority of the disks is lost for good | Only a forced recovery brings the cluster back, and it can lose acknowledged writes (recovery). |
In the recovery rehearsals of 2026-10-01 (three nodes in Docker under load, test/recovery/),
kill -9 of the leader or of a follower and SIGTERM of the leader lost none of 3,300 acknowledged
writes, and no leader change took more than 4 s. The guarantees behind those numbers are
Jepsen-tested.
Losing the majority deserves a look, because the node left standing does not notice. With nodes 1
and 3 of the Compose cluster stopped, node 2 went on answering /health with 200,
"status":"healthy" and "role":"leader", while a push to it waited out its 30 s deadline and
got this:
HTTP/1.1 503 Service Unavailable
retry-after: 1
{"error":"the request deadline elapsed","code":"timeout"}A write that fails like that, or with 503 outcome_unknown while the leader changes, may or may
not have applied, so send it again. A repeated ack of a lease it already released is answered as
the first one was, and a push repeated with the same transactionId is deduplicated
(dedup). The push above, sent again once the majority was back, answered
"status":"queued": it had never committed, and the retry wrote it once.
Three nodes or five
Run three voters, or five. A majority of three is two, so three nodes survive one failure; five survive two, at the price of a larger majority on every commit and two more followers for the leader to feed. Four nodes survive one failure like three while every write waits for three of them, so an even count only costs. Put each voter on its own machine, and in its own zone where you have zones: two voters on one host turn that host’s failure into a lost majority, which is also why the Compose cluster is a place to try things and nothing more.
More nodes make a cluster harder to kill, not faster, because one leader orders every write of a group. To spread write load, run several Raft groups per process (below) or more clusters. What one three-node cluster carries is on the benchmarks page.
Give every voter the same disk size. The disk gate (507 above 85% full,
QUEEN_RAFT_DISK_HIGH_PCT) only turns away a node’s own clients, and a smaller disk fills from
replication anyway.
Rolling upgrades
Upgrade one node at a time: stop it, start the new release on the same data directory, and wait
until its /health answers 200 before you touch the next. Healthy means it knows the leader and has
nearly caught up, within 1,000 entries of the leader’s commit or less than 2 s behind it
(QUEEN_RAFT_READY_LAG_ENTRIES, QUEEN_RAFT_READY_LAG_MS), so by the time you stop the next node
this one is back in the majority. In the Compose cluster, set QUEEN_VERSION in
.env to the release you move to and let --wait hold each step until the node is healthy:
for n in 1 2 3; do docker compose up -d --wait queen-$n; doneWe rehearsed that roll from 2.0.0-beta.5 to beta.6 while a writer pushed about 15 messages a
second round-robin over the three nodes, retrying a failed push on the next node with the same
transactionId. All 1,172 pushes were answered, and afterwards every node held all 1,172 messages,
each once. Six attempts met a node while it restarted and went on to the next one, and the slowest
push took 0.63 s.
A release that adds something to the log’s format has to roll without any old node meeting an
entry it cannot decode. So each node tells the leader, on every append it answers, the highest
format version it reads (kinds in /health), and the leader raises the cluster’s version
(clusterVersion) only once every member reports the new one. Until the last node runs the new
release, the leader goes on writing the old format and the feature that needs the new one stays
off; then the version rises and the feature turns on, with no flag for you to flip. A node whose
build reads less than the cluster’s version refuses to start on that data, and the leader refuses
to add it as a learner, so a downgrade past a format change stops at the first node. The gate
exists from 2.0.0-beta.2 on; the move to beta.1 was a stop-all upgrade.
Ephemeral queues in a cluster
Ephemeral queues keep their messages in memory and never in the log, and
still every node of a cluster serves them. Each (queue, partition) has exactly one owner,
picked by a highest-random-weight hash over the live members, so every node computes the same owner
without asking anyone, and a node that does not own a partition forwards the push, pop or ack to the
one that does. In the Compose cluster, six partitions pushed through node 1 landed two on node 2 and
four on node 3, and every message could be popped through any node.
When the membership changes, the partitions that move are handed over with their messages and each
consumer group’s position. A node stopping on SIGTERM tells its peers it is leaving and ships every
partition it owns to the next owner before it hands off leadership. Stopping node 3 logged
ephemeral rings handed over rings=4 delivered=4 lost=0, and its eight messages then popped
through node 1. When node 3 came back, nodes 1 and 2 shipped 13 partitions back to it, none lost.
The messages live in one node’s memory, so a crash (kill -9, OOM, a lost host) loses what that node owned. Of six messages spread over the cluster, a SIGKILL of node 3 left the two that node 2 owned. Two issues are still open in 2.0.0-beta.6. A hand-over to a node that is still restarting can drop the partition’s messages, which we have seen during rolling restarts. And through the proxy, a request forwarded to the owning node does not keep its tenant, so ephemeral queues behind the proxy are safe on a single node only, for now.
In production
The Compose file holds every setting a cluster needs, and production changes what surrounds it, starting with one machine per voter and a private network between them. Every node gets the single-node settings plus these:
| Variable | Default | What it does |
|---|---|---|
QUEEN_RAFT_REPLICATOR |
local |
openraft (or raft) makes the node a cluster member. A peer list without it fails the boot |
QUEEN_RAFT_NODE_ID |
1 |
A number from 1, different on each node. ordinal reads the trailing number of HOSTNAME plus one, so pod queen-0 is node 1 |
QUEEN_RAFT_PEERS |
unset | Every voter as id=raft_addr/http_addr, the same list on every node. Unset, the node is a single voter |
QUEEN_RAFT_LISTEN |
0.0.0.0: plus the node’s own raft port |
Where this node’s raft RPC server binds |
QUEEN_RAFT_TOKEN |
unset | The secret every raft RPC carries in x-queen-raft-token. Unset, anything that reaches the raft port can vote and append |
QUEEN_RAFT_DIR |
/var/lib/queen/raft |
This node’s data directory, never shared between nodes |
QUEEN_RAFT_JOIN |
off | The node joins a cluster that already exists, and serves 503 until the leader adds it |
The raft port (7400 by convention) carries votes, log entries, snapshots and the commands followers
send the leader, over plain HTTP with the token in a header. Keep it on a private network that only
the nodes reach (security). Port 6632 has no authentication
unless you turn on broker JWT, so it does not face the internet either. When you add the
proxy for tenants and API keys, give it its own port (QUEEN_PROXY_PORT) and
keep PORT, the address in the peer list, for the nodes and your internal clients.
Every node of a new cluster starts on an empty directory with the same peer list, and each one
initializes the cluster from it. After that first start the membership lives in the log, and the
list is read only for this node’s own listen address, the token and a group’s preferred leader.
A committed membership change therefore stands whatever QUEEN_RAFT_PEERS still says. Update the
list anyway, so the next new node sees the current members. The addresses are part of the
membership too: renaming hosts, ports or a namespace on a live cluster takes a forced recovery under
the new names.
Give a stopping node time to finish. The ephemeral hand-over can take up to 15 s, the leadership transfer up to 3 s, and a long-poll pop in flight up to its timeout (30 s by default), which is why the Compose file sets a 60 s stop grace period and the Kubernetes manifest a 90 s one.
Leases are timed on the leader’s clock. QUEEN_RAFT_MAX_CLOCK_SKEW_MS (500) is how long a lease
another node timed is held past its expiry, and it lengthens the pause before a new leader serves
consumption. Lease exclusivity across a leader change holds only while the nodes’ clocks agree
within that bound, so run NTP or chrony on every node.
/health answers 200 once a node knows a leader and has caught up
(monitoring), so use it for readiness. For alerts use a write probe
(recovery), because a node that lost its quorum keeps answering 200, as node 2 did
above.
Change the membership
QUEEN_RAFT_PEERS is read once, at the first boot. From then on the membership changes only
through these routes, which need the admin access level and are blocked by the proxy, so call them
on PORT:
| Route | Body, or what it answers |
|---|---|
GET /api/v1/system/raft/membership |
Answers the leader, the term, voters, learners, and each member’s lag and live |
POST /api/v1/system/raft/membership/learners |
{"id": 4, "raft": "queen-4:7400", "http": "queen-4:6632"} |
POST /api/v1/system/raft/membership/promote |
{"ids": [4]} |
PUT /api/v1/system/raft/membership/voters |
{"voters": [1, 2, 4]} |
DELETE /api/v1/system/raft/membership/members/:id |
none |
The leader checks every change before it writes it. One change runs at a time, a live majority has
to exist before and after (409 no_quorum otherwise), the last voter cannot be removed, and a
learner is promoted only within 1,000 entries of the leader (QUEEN_RAFT_PROMOTE_MAX_LAG;
409 learner_behind, or "force": true). Every change is idempotent, so a lost answer is safe to
repeat. To replace node 3 after its disk is gone:
curl -X DELETE http://queen-1:6632/api/v1/system/raft/membership/members/3
# wipe node 3's directory, then start it with QUEEN_RAFT_JOIN=true: it serves 503 until added
curl -X POST http://queen-1:6632/api/v1/system/raft/membership/learners \
-H 'content-type: application/json' -d '{"id":3,"raft":"queen-3:7400","http":"queen-3:6632"}'
curl -X POST http://queen-1:6632/api/v1/system/raft/membership/promote \
-H 'content-type: application/json' -d '{"ids":[3]}'Remove a member before you wipe its disk. A voter that comes back empty under its old id is never repaired: it looks healthy and its clients hang.
Raft groups
QUEEN_RAFT_GROUPS=3 (default 1, at most 64, the same on every node) runs that many Raft groups in
each node, each with its own log, leader and data directory (groups/g<N>/). A tenant lives in
exactly one group, picked by a stable hash of its id or pinned with
QUEEN_TENANT_GROUPS=<tenant>=<group>,..., so its transactions stay single entries. Group N listens
on the raft port plus N and prefers the member at position N mod n as its leader, so with as many
groups as nodes every node leads one and follows the others. It needs QUEEN_RAFT_CLIENT_OFFLOAD
on, which is the default.
Raft groups were not part of the Jepsen pass, which ran one group. A tenant cannot change group while it has data, because there is no migration yet, so pin a large tenant before it writes.
Recovery
When something looks wrong, start with a write probe rather than /health, because a node that
lost its quorum can keep answering 200:
curl -m 10 -X PUT http://queen-1:6632/api/v1/kv/ops/probe \
-H 'content-type: application/json' -d '{"value":"probe","ttlSeconds":60}'An answer within a second means the cluster commits; nothing after 10 s means it does not. From there, Recovery walks through every case we rehearsed on real clusters (a lost node, a full disk, a lost majority, a forced recovery, a poisoned entry, a rolling restart gone wrong), and Backups and restore covers snapshots of the data directory.
Limits
- Three containers on one host, as in the Compose file, all go down with that host.
- Losing a majority stops writes, and losing a majority of the disks for good can lose acknowledged data.
/healthcan stay 200 on a node that lost its quorum, and on an empty node that never caught up (F8 and F2 intest/recovery/FINDINGS.md). Alert on the write probe.- The raft port is plain HTTP with a shared token: keep it on a private network.
- Peer addresses are fixed at the first boot. Changing them on a live cluster is a forced recovery.
- Ephemeral messages live in their owner’s memory: a crash loses that node’s share, and a hand-over to a node that is still restarting can drop a partition’s messages. Behind the proxy they are safe on a single node only (2.0.0-beta.6).
- One leader orders every write of a group. Adding nodes to a group adds no write throughput.
Kubernetes turns this cluster into a StatefulSet, and monitoring covers what to watch once it runs.