A worker’s step in Queen is three calls at most: it pops a batch, does its work, and commits a transaction that acks the input and writes the output. This page follows each of those through a node, from the HTTP edge to the answer, so you know exactly how far a write had got when you were told about it, and what a status code means about the state behind it. The parts it names are on architecture.
Push
curl -X POST localhost:6632/api/v1/push -H 'content-type: application/json' -d '{
"items": [{"queue": "orders", "partition": "customer-42",
"transactionId": "order-1001", "payload": {"total": 30}}]}'- At the edge, the request takes admission permits sized by its
Content-Lengthbefore its body is read, so a push that has to wait holds only its connection and its bytes stay in the socket. When the budget (QUEEN_RAFT_ADMIT_MAX_MB, at least 128 MiB) stays full forQUEEN_RAFT_ADMIT_HOLD_MS(15 s, jittered), the push gets429 overloadedwith aRetry-Afterof 1 to 5 s. A node aboveQUEEN_RAFT_DISK_HIGH_PCT(85%) answers507 storage_full, and a name too long for a store key413. - The receiving node mints a request id and a deadline, gives each item a message id, hashes
its
transactionId(or the message id, when there is none) with xxh3-128, packs one frame per message and encrypts the payload if the queue is encrypted. On a follower, only this prepared command crosses to the leader, over one ofQUEEN_RAFT_FWD_STREAMS(4) long-lived streams, and it takes its own fair share of the leader’s admission budget there. - In the planner, a retry under a request id the leader has already answered gets the recorded
outcome, for
QUEEN_RAFT_REQUEST_ID_WINDOW_S(60 s). The first push that names a queue or a partition creates it in the same entry as its first message, which is why there is no partition count to plan and no topic to create first. - For dedup, on a partition the planner knows, a ring of bloom filters answers “certainly new” or “maybe” for each hash. A “maybe” is checked exactly against the committed per-append hash rows and the entries still in flight, inside the queue’s dedup window, so the filter can save a read but never decide a duplicate. A duplicate gets its original offset and adds nothing; the survivors get one gapless run of offsets.
- The log writer puts the payloads into their queue’s log file and an entry record into every log the entry touches, then fsyncs each touched log once. Each follower does the same before it acknowledges the append.
- The answer is
201when the entry commits (QUEEN_RAFT_ANSWER_AT_COMMIT, on by default): on one node once the fsync returns, in a cluster once a majority has fsynced it. A push with a duplicate item waits for the local apply, because its answer carries the original message id, read from the log. Every node then applies the entry, and on the leader apply arms the partition for every consumer group that holds it and wakes the long polls waiting there.
server/src/rsm/facade/real.rs, server/src/rsm/planner/, server/src/rsm/qlog/mod.rsPop
curl 'localhost:6632/api/v1/pop/queue/orders?consumerGroup=billing&batch=10&wait=true&timeout=30000'- The pop goes to the consumption engine on the leader. The engine serves only while
its node leads and has heard from a quorum within
QUEEN_CONSUME_LEADER_LEASE_MS(400 ms); otherwise the pop is taken to the new leader. - The engine claims from the group’s claimable partitions, oldest first. A partition becomes
claimable when an append to it applies, or when a lease on it is released, expires or a delay
ends. Each claimed partition gets a lease for the queue’s
leaseTime. The engine stops claiming a margin before the caller’s deadline (250 ms at most, and never more than a quarter of the time left), so it never leases messages nobody will receive. - When there is nothing to claim and
wait=true, the pop is held on the leader until a partition of the group becomes claimable or the deadline comes, and then it answers204. - A lease is a change to the group’s cursor row, so it has to be durable before it counts. Every
QUEEN_CONSUME_CHECKPOINT_MS(5 ms) the engine logs the changed rows as a checkpoint, and the pop is answered only once the checkpoint holding its lease has committed.QUEEN_CONSUME_FAST=1answers before that, and a leader change can then lose the lease. - The node that answers reads the messages from its own queue logs once it has applied up to the leader’s index. A follower that cannot catch up before the deadline hands the leases back at once, so the partition is not frozen for a whole lease.
Ack
curl -X POST localhost:6632/api/v1/ack -H 'content-type: application/json' -d '{
"transactionId": "order-1001", "partitionId": "17", "leaseId": "<from the pop>",
"consumerGroup": "billing", "status": "completed"}'- The engine checks a
leaseIdagainst the live lease: the same holder, not expired. An ack without aleaseIdskips the check and still advances the cursor. - The acked hashes are resolved inside the leased batch and the cursor moves. A
failedstatus charges the retry budget; past it, the message goes to the DLQ in the same checkpoint as the cursor row. - The answer comes, as for a pop, once the checkpoint holding the new cursor row has committed.
- An ack sent again after it already released its lease (because its reply was lost) is answered as the first one was, on any leader: the released lease is kept in the cursor row for exactly this.
server/src/rsm/consume/pop.rs, server/src/rsm/consume/ack.rsTransaction
curl -X POST localhost:6632/api/v1/transaction -H 'content-type: application/json' -d '{
"operations": [
{"type": "ack", "transactionId": "order-1001", "partitionId": "17",
"consumerGroup": "billing", "leaseId": "<from the pop>", "status": "completed"},
{"type": "push", "items": [{"queue": "receipts", "partition": "customer-42",
"transactionId": "receipt-1001", "payload": {"total": 30}}]}],
"kv": [{"op": "put", "ns": "orders", "key": "1001",
"value": {"status": "paid"}, "ttlSeconds": 2592000}]}'- The receiving node validates the body (
kvandtimersare top-level arrays, never operations) and the pushes are packed like a push. A body of the wrong shape is a400. - On the leader, before anything is planned, the engine validates every ack on shadow
copies of the cursors. An ack whose lease has expired or belongs to another worker, or whose
message is already acked, rolls the whole transaction back (
rejected_ack). Otherwise the engine reserves those partitions, so no pop, ack, expiry or checkpoint touches them until the entry lands (or forQUEEN_CONSUME_TXN_TTL_MS, 30 s, at most), and hands the planner the cursor rows to write. - The planner takes the pushes first, then positions, then KV, then timers, then the engine’s
rows. A duplicate push rolls the whole transaction back (
duplicate), and so does a lostrequiredKV precondition (kv_precondition). A rollback restores the planner’s view, and nothing is logged. - Everything goes into one entry, written into every queue log it touches. Crash recovery replays it only if every copy is present, so a crash can never leave half a transaction behind (see storage).
- After the entry applies on the answering node, the answer is
200withsuccess: true, its KV reads rendered from that node’s state. A rollback is also200, withsuccess: falseand areason.
After a leader change
The node serving a call retries it under the same request id, against the new leader, until its deadline. That is safe for everything the planner logs (pushes, transactions, KV, timers): if the entry committed, the new leader returns the recorded outcome, and if it did not, the command is planned for the first time.
Acks and nacks are served from the engine’s memory and leave no request-id record. If the leader
loses its lead while the checkpoint carrying an ack is on its way to the log, or a forwarded ack is
lost after it went out, the answer is 503 outcome_unknown: the ack may or may not have
applied. Send the same ack again; if it applied, the repeat is answered as a success. A
transaction is not repeat-safe through its acks (a retry finds them settled and rolls back), so
give its pushes a deterministic transactionId: a retry of a committed transaction then rolls back
as a duplicate, and the first commit stands.
A graceful stop of the leader waits up to 500 ms for the acks it already answered to commit, then hands leadership to a caught-up peer. A new leader takes over only after applying an entry of its own term, and it loads leases from the committed cursor rows, so a leased batch is not delivered again before its lease expires.
When you are answered
| Command | Leader answers | Follower answers |
|---|---|---|
| Push, no duplicate | when the entry commits | when it learns the entry committed in the leader’s term |
| Push with a duplicate | after the entry applies locally | after it applies the entry |
| Pop | after the lease’s checkpoint commits | after it applies up to the leader’s index |
| Ack, nack, renew | after the cursor row’s checkpoint commits | when the leader’s answer arrives |
| Transaction, KV | after the entry applies locally | after it applies the entry |
Limits
Delivery is at-least-once: a lease that expires before its ack lands puts the batch back. A
503 outcome_unknown or 503 timeout does not say whether the change applied, so make retries
harmless with a transactionId on every push and a once gate on every side effect (see
dedup). Lease exclusivity across a leader change needs node clocks within
QUEEN_RAFT_MAX_CLOCK_SKEW_MS (500 ms) of each other. And a transaction is one call, with no
BEGIN and COMMIT across requests (see transactions).