Skip to content

Life of a step

One push, one pop, one ack and one transaction followed from the HTTP edge to the answer: where dedup, leases and refusals happen, and when you are answered.

Updated View as Markdown

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}}]}'
  1. At the edge, the request takes admission permits sized by its Content-Length before 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 for QUEEN_RAFT_ADMIT_HOLD_MS (15 s, jittered), the push gets 429 overloaded with a Retry-After of 1 to 5 s. A node above QUEEN_RAFT_DISK_HIGH_PCT (85%) answers 507 storage_full, and a name too long for a store key 413.
  2. 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 of QUEEN_RAFT_FWD_STREAMS (4) long-lived streams, and it takes its own fair share of the leader’s admission budget there.
  3. 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.
  4. 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.
  5. 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.
  6. The answer is 201 when 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.
A push sent to a follower of a three-node cluster. The follower admits it, mints the message ids, hashes the transactionIds and packs one frame per message, then forwards only the prepared command to the leader over a long-lived stream. The leader plans it (request id, dedup, offsets), writes it to its queue logs and fsyncs them, and sends the entry to both followers, which write and fsync it too. Once a majority has it the entry is committed. The answering follower learns the commit index from the leader and answers the client 201.clientfollowernode 2leadernode 1followernode 3POST /api/v1/pushadmit, mint ids,hash, pack framesprepared commandplan: request id,dedup, offsetswrite, fsyncappendappendwrite, fsyncdonemajority: committedcommit index201
The edge work, parsing, ids and hashing, stays on the node the client called; the leader receives a prepared command and spends its time planning, ordering and writing. Source: server/src/rsm/facade/real.rs, server/src/rsm/planner/, server/src/rsm/qlog/mod.rs

Pop

curl 'localhost:6632/api/v1/pop/queue/orders?consumerGroup=billing&batch=10&wait=true&timeout=30000'
  1. 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.
  2. 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.
  3. 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 answers 204.
  4. 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=1 answers before that, and a leader change can then lose the lease.
  5. 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"}'
  1. The engine checks a leaseId against the live lease: the same holder, not expired. An ack without a leaseId skips the check and still advances the cursor.
  2. The acked hashes are resolved inside the leased batch and the cursor moves. A failed status charges the retry budget; past it, the message goes to the DLQ in the same checkpoint as the cursor row.
  3. The answer comes, as for a pop, once the checkpoint holding the new cursor row has committed.
  4. 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.
A pop and its ack through the consumption engine on the leader. The worker pops for group billing; the engine claims the oldest claimable partition and leases its batch. The lease is a change to the group's cursor row, logged in the engine's next checkpoint, at most 5 ms later, and the pop is answered with the messages and a leaseId once that checkpoint has committed. After the worker's step, its ack is checked against the live lease, moves the cursor, rides the next checkpoint, and is answered once that checkpoint commits.workerconsumption engineleaderraft logmajority fsyncpop, group billingclaim oldest, leasecheckpoint, within 5 mscommitted200: messages, leaseIdthe worker runs its stepack completed, leaseIdcheck lease, move cursornext checkpointcommitted200
Pops and acks never go through the planner. They change cursor rows in the engine's memory, and a few checkpoint commands every 5 ms make those changes durable, however many acks they carry. Source: server/src/rsm/consume/pop.rs, server/src/rsm/consume/ack.rs

Transaction

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}]}'
  1. The receiving node validates the body (kv and timers are top-level arrays, never operations) and the pushes are packed like a push. A body of the wrong shape is a 400.
  2. 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 for QUEEN_CONSUME_TXN_TTL_MS, 30 s, at most), and hands the planner the cursor rows to write.
  3. 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 lost required KV precondition (kv_precondition). A rollback restores the planner’s view, and nothing is logged.
  4. 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).
  5. After the entry applies on the answering node, the answer is 200 with success: true, its KV reads rendered from that node’s state. A rollback is also 200, with success: false and a reason.

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

Navigation

Type to search…

↑↓ navigate↵ selectEsc close