Skip to content

Ephemeral queues

In-memory queues for request/reply, presence and signals. Every node serves every partition, a graceful stop hands the contents to the next owner, and a crash loses what the crashed node held.

Updated View as Markdown

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:

Three Queen nodes serving one ephemeral queue. A producer pushes to presence/room-1 through queen-1 and a consumer pops room-1 through queen-2. Neither node owns room-1, so both forward the request to queen-3, which owns room-1 (with room-2, room-3 and room-6) and keeps its messages in memory. queen-2 owns room-4 and room-5, queen-1 owns none of the six rooms. The owner of a (queue, partition) is a highest-random-weight hash over the live members, which every node computes the same way.producerpush room-1consumerpop room-1, long pollqueen-1owns none of the roomsqueen-2owns room-4, room-5queen-3owns room-1, 2, 3, 6holds them in memoryany nodeany nodeforwardforwardowner = highest-random-weight hashof (queue, partition)over the live members
Clients can use any node. The node that receives a request hashes the partition, finds the owner and forwards to it, and the owner answers through the same hop. A long-poll pop waits at the owner. Source: 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:

A sequence on three nodes. queen-3 receives SIGTERM and from then on owns nothing. It tells queen-2 and queen-1 that it is leaving, and both leave it out of the placement hash for 30 seconds. queen-3 then sends each of them the partitions they now own, with the messages and every consumer group's position; each answers adopted. queen-3 logs the hand-over, hands off raft leadership and exits. In the alternative where queen-3 crashes instead (kill -9, out of memory), its partitions' messages are lost with the process, and after 4 seconds without hearing from it the other nodes re-hash and its partitions start again on them, empty.queen-1queen-2queen-3stoppingSIGTERM: owns nothing nowleavingleavingqueen-3 left out of the hash for 30 sthe partitions queen-2 now ownsmessages, group positionsadoptedthe partitions queen-1 now ownsmessages, group positionsadoptedall handed overraft hand-off, exitor a crash: kill -9, out of memoryits messages die with it4 s without a word from queen-3: re-hash,its partitions start again here, empty
A stop moves the contents, a crash loses them. The hand-over runs before the node gives up raft leadership, while its peers still answer and it still serves. Source: server/src/handlers/ephemeral.rs (ephemeral_drain, ship, handle_ephemeral_adopt), server/src/main.rs

When 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=0

The 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.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close