An ephemeral queue keeps its messages in memory and nowhere else. A push is a memory write and a wake-up, with no disk and no raft entry in the path, so a consumer parked on a long poll gets the message as soon as it has crossed the network. We built it for the traffic a durable log serves worst: thousands of short-lived inboxes, small payloads consumed at once, and a history nobody will ever read.
Request and reply
The requester mints an inbox and waits on it; the responder pushes each answer into whatever inbox the request names.
import { randomUUID } from 'node:crypto'
// Responder: every process in the group shares the requests.
const { messages } = await queen.ephemeral.pop('rpc-requests', { group: 'workers', batch: 10, wait: true })
for (const m of messages) {
await queen.ephemeral.push(m.payload.replyTo, [{ squared: m.payload.n ** 2 }])
}
await queen.ephemeral.ack('rpc-requests', messages, { group: 'workers' })
// Requester: push, then park on the inbox until the answer arrives.
const inbox = `rpc-inbox-${randomUUID()}`
await queen.ephemeral.push('rpc-requests', [{ n: 7, replyTo: inbox }])
const reply = await queen.ephemeral.pop(inbox, { wait: true, timeout: 5000, autoAck: true })
if (reply.messages.length === 0) throw new Error('no reply within 5 s')Over HTTP, a push and a long-poll pop look like this:
curl -X POST localhost:6632/api/v1/ephemeral/push -H 'content-type: application/json' \
-d '{"queue":"presence","partition":"room-7","messages":[{"payload":{"user":"ada","typing":true}}]}'
curl 'localhost:6632/api/v1/ephemeral/pop?queue=presence&partition=room-7&group=ui-1&wait=true&timeout=5000'Neither queue was declared anywhere. The first push or pop that names a queue creates it, and an
implicit queue that stays empty and unpolled for QUEEN_EPHEMERAL_IMPLICIT_IDLE_S (300 s) is
collected, so a requester can mint one inbox per request and leave nothing behind. A popped
message is { id, partition, payload, attempts }, and order holds within a partition.
Groups, acks and bounds
The pop’s group decides how a queue is consumed, as it does on a durable queue: consumers that
share a group compete for messages, a subscriber with a group of its own sees every message, and a
pop with no group reads in queue mode. However many groups read a message, it sits in memory once.
Without autoAck, a popped message is leased for 30 seconds. If the lease runs out, or you ack it
failed or retry, it comes back with attempts incremented, and after retryLimit attempts (5)
it is dropped and counted. That is at-least-once for as long as the message stays in memory, so a
consumer needs to tolerate a repeat exactly as it would on a durable queue. With autoAck: true
the message counts as delivered when it leaves the broker, which is at-most-once.
Declare a queue when the defaults don’t fit it:
await queen.ephemeral.configure('presence', {
maxLength: 1000, // up to QUEEN_EPHEMERAL_QUEUE_MAX_LENGTH (10,000, also the default)
maxBytes: 1024 * 1024, // up to QUEEN_EPHEMERAL_QUEUE_MAX_BYTES (16 MiB, also the default)
policy: 'dropOldest', // or 'reject', the default: the push answers 429 queue_full
ttlSeconds: 30, // drop anything older, consumed or not
leaseSeconds: 15, // default 30
retryLimit: 3, // default 5
})The declaration goes through the raft log, so every node applies it and a declared queue comes
back after a restart, configured and empty. The option list is closed and an unknown option is a
400, because a misspelled bound would otherwise give you a queue that looks configured and
behaves like the default. ttlSeconds drops a message once it is older than the limit whether or
not anyone consumed it, which is the right thing for a typing indicator or a cache invalidation
that a newer one has already made worthless. Above all of this, a node holds at most
QUEEN_EPHEMERAL_MAX_BYTES (256 MiB) of ephemeral messages and answers 503 ephemeral_unavailable
past it.
In a cluster
Every node serves every ephemeral queue. Each (queue, partition) has exactly one owner, picked by
a highest-random-weight hash over the cluster’s live members, and because every node computes the
same owner from the same membership, a node that receives a request for a partition it does not
own forwards it to the owner. There is no lease to take and no rebalance to wait for. Here is a
three-node cluster after one push per room, all six sent to queen-1; each node lists the
partitions it owns:
for n in queen-1 queen-2 queen-3; do
curl -s http://$n:6632/api/v1/ephemeral/queues/presence/depth | jq -c '[.partitions[]?.partition]'
done[]
["room-4","room-5"]
["room-1","room-2","room-3","room-6"]On that cluster, a push for room-1 and a pop for it reach queen-3 from whichever node they land on:
server/src/ephemeral.rs (hrw_pick, Ephemeral::route), server/src/handlers/ephemeral.rs (forward_to_owner)What happens to a partition’s contents depends on how its owner goes away:
| Event | The contents of the partitions that move |
|---|---|
The owner stops with SIGTERM (a deploy, a rolling restart, docker stop) |
Handed to the next owners before the process exits, with every group’s position. A leased message comes back as a redelivery with attempts + 1 |
| A node joins, or comes back after a restart | Handed over by the current owners to the node they now hash to |
| The owner crashes (kill -9, out of memory, the machine is lost) | Lost. After QUEEN_EPHEMERAL_MEMBER_TTL_MS (4000) the other nodes place its partitions elsewhere, empty |
| Every node stops | Every queue is empty. Declared queues come back configured |
Here are the first and third rows as they happen to queen-3:
server/src/handlers/ephemeral.rs (ephemeral_drain, ship, handle_ephemeral_adopt), server/src/main.rsWhen queen-3 above was stopped with docker stop, it logged the hand-over on its way out, and
the two messages a consumer had leased from room-1 came back from queen-2 with attempts: 2:
INFO shutdown: ephemeral rings handed over rings=4 delivered=4 lost=0The id of a delivered message names the node incarnation that delivered it, so an ack sent after
the move answers stale instead of acknowledging whatever the new owner keeps under that number.
A consumer that reconnects after a failover can flush its pending acks and read the answers as
information:
{"results":[{"id":"e:624f5af82b58223f:room-1:0","outcome":"stale"}]}We hand partitions over and do not replicate them because a second copy of every push would put a network round trip, or a raft entry, back into the path this class exists to keep short. The price is the crash case: a partition’s contents live on one node at a time, and a node that dies takes its share with it. If request/reply is your use, that is the failure the requester’s timeout already handles.
Limits
None of the guarantees apply here, and ephemeral queues are not part of
the Jepsen runs. There are no transactions, no dead-letter queue, no replay and no
subscriptionMode, and an ephemeral push cannot ride a /transaction. When a message is part of
a step that has to commit, use a durable queue.
A hand-over to a node that is still starting up can drop the partition’s contents. We have seen it
in rolling restarts and it is an open issue. A dropped hand-over is never silent: the sending node
logs hand-over failed; the partition's contents are dropped and counts it in
queen_ephemeral_wipes_total, although a node that drops on its way out takes that counter with
it, so alert on the log line.
Behind the proxy, ephemeral queues are safe on a single node only, for now. In a cluster, a request the receiving node forwards to the partition’s owner does not keep its tenant and lands in the broker’s default tenant. It is a bug we are fixing, and durable queues are not affected.
With the broker’s own JWT authentication on (JWT_ENABLED=true), the nodes’ hand-over calls are
refused with 401, so every move drops the contents it carries. Requests that clients send through
any node still reach the owner, because they carry the client’s own token.
While the nodes disagree about who is a member (a few seconds after a crash), a request can answer
503 owner_moved or 503 ephemeral_forward_failed; retry it. A handed-over message can also land
behind messages the new owner took during those seconds, so order can bend across a move.
The bounds count what one node holds: three nodes can hold up to three times a queue’s maxLength
between them. The queue listing and depth show the share of the node that answers, and a pop
with no partition serves the partitions the answering node owns (or routes on Default when it
owns none), so name the partition, or keep one partition per queue as inboxes do.
Next: consuming covers groups and acks in depth, and
monitoring lists the queen_ephemeral_* series.