Give every entity its own partition, keep its state in one KV entry, and commit each transition as
one transaction. Any number of stateless workers in one consumer group then run your entities as
state machines, without locks, an outbox table or a scheduler. This page builds an order that goes
from created to paid to shipped, and cancels itself if nobody pays within 30 minutes.
The machine
Every event of an order lands in that order’s partition, in the order it was pushed. A worker reads the order’s state, decides, and commits one transaction:
| Event | State in KV | The one transaction |
|---|---|---|
order_created |
none | ack, put created, schedule payment_timeout in 30 minutes into this partition |
payment_succeeded |
created |
ack, put paid, cancel the timer, push ship_order to shipping |
payment_timeout |
created |
ack, put cancelled, push order_cancelled to notifications |
payment_succeeded |
cancelled |
ack, push refund to refunds (paid after the deadline) |
shipped |
paid |
ack, put shipped |
| anything else | any | ack only: a duplicate or a late event changes nothing |
examples/apps/js/saga.mjsEvents enter like any message. The partition is the order id, and the transactionId makes a
producer’s retry harmless, because the broker drops a second push with the same id inside the
queue’s dedup window (an hour by default):
await queen.queue('orders').partition('ord-1042').push([
{ transactionId: 'ord-1042:created', data: { type: 'order_created', total: 4200 } },
])The worker
The table becomes a pure function, and the worker turns its answer into one transaction:
import { Queen } from 'queen-mq'
const queen = new Queen({ url: 'http://localhost:6632' })
const ORDERS = 'orders' // one partition per order: the input
const STATE = 'order-state' // KV namespace: the state
const timeoutKey = (orderId) => `payment-timeout:${orderId}`
// Pure: no I/O. Given the state and the event, say what this step does.
function transition(state, event) {
switch (`${state?.status ?? 'none'}/${event.type}`) {
case 'none/order_created':
return { state: { status: 'created', total: event.total }, startTimeout: '30m' }
case 'created/payment_succeeded':
return { state: { ...state, status: 'paid' }, cancelTimeout: true,
emit: [{ queue: 'shipping', data: { type: 'ship_order' } }] }
case 'created/payment_timeout':
return { state: { ...state, status: 'cancelled' },
emit: [{ queue: 'notifications', data: { type: 'order_cancelled' } }] }
case 'cancelled/payment_succeeded':
return { state, emit: [{ queue: 'refunds', data: { type: 'refund', amount: event.amount } }] }
case 'paid/shipped':
return { state: { ...state, status: 'shipped' } }
default:
return { state } // a duplicate or a late event: ack it, change nothing
}
}
await queen.queue(ORDERS)
.group('order-machine')
.subscriptionMode('all') // a new group starts at the tail unless told otherwise
.autoAck(false) // the ack rides the transaction below
.concurrency(10)
.each()
.consume(async (msg) => {
const orderId = msg.partition // the partition is the entity
const row = await queen.kv.get(STATE, orderId) // { found, value, version }
const state = row.found ? row.value : null
const step = transition(state, msg.data)
const tx = queen.transaction()
.ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
if (step.state !== state) {
tx.kv.put(STATE, orderId, step.state, { ttl: '90d' })
}
for (const e of step.emit ?? []) {
tx.queue(e.queue).partition(orderId).push([{ data: { orderId, ...e.data } }])
}
if (step.startTimeout) {
tx.timer(ORDERS).key(timeoutKey(orderId)).partition(orderId)
.delay(step.startTimeout).payload({ type: 'payment_timeout' }).schedule()
}
if (step.cancelTimeout) {
tx.timer(ORDERS).key(timeoutKey(orderId)).cancel()
}
try {
await tx.commit()
} catch (err) {
if (err.reason === 'rejected_ack') return // lease lost: nothing was written
await queen.ack(msg, 'failed', { group: msg.consumerGroup, error: err.message })
}
})clients/client-js/client-v2/builders/TransactionBuilder.js, examples/apps/js/saga.mjsA few details in that code matter. The ack names consumerGroup so the code reads the same over
HTTP, where an ack without one acks in queue mode, outside any group; the JS SDK takes it from the
message when you leave it out. Every KV put carries an expiry (ttl,
ttlSeconds or until) or forever: true, so state never outlives its use by accident. And the
timer’s .partition(orderId) is what sends the timeout into the order’s own partition; without it
the timer fires into Default.
A commit that fails for another reason is nacked, so the event comes back. Each nack spends one
retry, and once retryLimit retries are spent the next failure files the event in the dead-letter
queue (it has then been tried retryLimit + 1 times), which is on for every queue by default. See transactions,
KV and timers for each piece on its own.
Why it needs no locks
While a worker holds the lease on an order’s partition for order-machine, no other worker in that
group receives that order’s next event. A step’s read with kv.get and its write in the
transaction therefore cannot interleave with another step of the same order: the next step’s read
happens after this commit, and a single-key KV read is linearizable, so it sees the commit.
A worker can still stall past the queue’s leaseTime, and then the partition goes to another
worker. The stalled worker’s commit carries its old lease, so the broker refuses the ack with
rejected_ack and rolls back the KV write, the pushes and the timer with it. The new worker reads
the old state and runs the step once. A compare-and-set on the state alone cannot give you this,
because a version that still matches succeeds even for a worker that no longer owns the event. If
you want the assumption checked as well, add expect: row.found ? row.version : 0 and
required: true to the put: a second writer then gets kv_precondition back and the value stays
as it was.
The timeout needs no special handling either. The timer fires into the order’s partition, so it
takes its place in the order’s sequence like any other event. If the payment lands first, its step
cancels the timer; if the timer had already fired, the cancel answers absent without failing the
transaction, and payment_timeout arrives in state paid, where the table says to ack it and do
nothing. If the timeout lands first, the order is cancelled and the payment behind it is refunded.
Every worker sees one sequence of events per order, so every worker decides the same way.
That is what we mean by no outbox and no scheduler. The events a step emits are pushed in the same commit as its state, so they exist exactly when the state changed, and the timeout is a timer in that same commit, so it exists exactly when the order is waiting for a payment.
A full program
The booking saga below runs the same pattern with three states (held, confirmed, expired),
end to end, in the examples suite. Its timer fires into a separate expiries queue, so its
compensator reads the state again before it acts, which is the check the order above gets for free
by firing into the entity’s own partition.
//
// A booking saga: a room is held, paid for, and released by a timer when the
// payment never comes.
//
// The release is the part that usually breaks. Kept as a setTimeout in a
// worker, it dies with the worker, and a deploy in the middle of a hold leaves
// the room held forever. Here the release is a timer in the broker, scheduled
// in the same transaction that holds the room: the saga's state, the timer,
// the payment request and the ack of the booking commit as one log entry. If
// the room is held, its release exists; if anything failed, none of it
// happened.
//
// bookings
// └── group "reserver" ONE transaction: state + timer + push + ack
// ├── payments (one partition per booking)
// │ └── group "payer" confirm + cancel the timer + ack
// └── expiries (the timer delivers here when the hold runs out)
// └── group "compensator" reads the state before it releases
//
// Run it:
// QUEEN_URL=http://localhost:6632 node saga.mjs
import { Queen } from 'queen-mq'
const QUEEN_URL = process.env.QUEEN_URL || 'http://localhost:6632'
const RUN = Date.now().toString(36)
// Fresh queues and a fresh KV namespace per run. Saga entries outlive the queues
// (they expire with their TTL), so a second run in the same namespace would
// find every booking already held.
const BOOKINGS = `app-js-saga-bookings-${RUN}`
const PAYMENTS = `app-js-saga-payments-${RUN}`
const EXPIRIES = `app-js-saga-expiries-${RUN}`
const NS = `app-js-saga-${RUN}`
// How long a room stays held before it is released. In production this is
// minutes; it is the only number that changes. It has to outlast the reserving
// and paying phases below, or a release would fire before the payment that
// cancels it and the run would measure a race.
const HOLD_MS = 10_000
// Each phase ends on a count of messages, with this deadline behind it so a
// stall fails the run instead of hanging it.
const PHASE_MS = 20_000
// Timers fire on the leader's next tick after they are due (every 50 ms), so
// the compensation phase only needs the hold plus a margin.
const TIMER_DEADLINE_MS = HOLD_MS + 20_000
// Four bookings, five submissions: B-2 is submitted twice, which is what a
// redelivery looks like from the reserver's side.
const BOOKINGS_IN = [
{ bookingId: 'B-1', room: '101', cents: 24000 },
{ bookingId: 'B-2', room: '102', cents: 31000 },
{ bookingId: 'B-2', room: '102', cents: 31000 },
{ bookingId: 'B-3', room: '103', cents: 18000 },
{ bookingId: 'B-4', room: '104', cents: 27000 },
]
const BOOKING_IDS = ['B-1', 'B-2', 'B-3', 'B-4']
// B-3's card is declined, so its saga never reaches "confirmed" and the timer
// is what gives the room back.
const DECLINED = 'B-3'
// B-4 pays, but its cancel is skipped on purpose. That is the cancel that comes
// too late, made reproducible: the release is delivered for a booking that is
// already confirmed, and the compensator has to refuse it.
const CANCEL_SKIPPED = 'B-4'
let checks = 0
const assert = (condition, description) => {
if (!condition) throw new Error(description)
checks++
console.log(` ok: ${description}`)
}
// The saga's state key derives from the booking id, which is also the
// partition key of the payments queue. That is what makes the payer's
// read-then-write safe below.
const sagaKey = (bookingId) => `saga:${bookingId}`
const reserveDecisions = []
const paymentsRequested = []
const compensationsDelivered = []
const roomsReleased = []
const compensationsRefused = []
let gatesLost = 0
const queen = new Queen({ url: QUEEN_URL, handleSignals: false })
// A consumer for one phase: acks ride the transactions, so autoAck is off, and
// the phase ends after `limit` messages or after `idleMillis` with none.
const phase = (queue, group, limit, idleMillis) => queen
.queue(queue)
.group(group)
.subscriptionMode('all')
.autoAck(false)
.each()
.limit(limit)
.timeoutMillis(1000)
.idleMillis(idleMillis)
try {
console.log(`broker ${QUEEN_URL}`)
for (const q of [BOOKINGS, PAYMENTS, EXPIRIES]) {
await queen.queue(q).config({ leaseTime: 30, retryLimit: 3 }).create()
}
console.log('\nsubmitting bookings')
for (const [i, booking] of BOOKINGS_IN.entries()) {
await queen.queue(BOOKINGS).push({
// One id per submission, so the duplicate of B-2 is stored and reaches
// the reserver, where the gate has to catch it.
transactionId: `submit-${i}-${booking.bookingId}`,
data: booking,
})
}
console.log(` ${BOOKINGS_IN.length} submissions for ${BOOKING_IDS.length} bookings`)
// --------------------------------------------------------------- reserving
//
// The transaction is the whole point of the example: four things commit
// together, so there is no order between them to get wrong. Written as four
// calls, a crash between the timer and the push leaves a release for a
// payment that was never asked for, and a crash the other way round leaves a
// hold with no release.
console.log('\nreserving')
await phase(BOOKINGS, 'reserver', BOOKINGS_IN.length, PHASE_MS).consume(async (msg) => {
const { bookingId, room, cents } = msg.data
const res = await queen
.transaction()
// 1. The gate and the first state, in one KV entry. required: true makes
// the putIfAbsent a gate: if the booking is already held, the whole
// transaction rolls back and nothing below happens.
.kv.putIfAbsent(NS, sagaKey(bookingId), { step: 'held', room, cents }, { ttl: '1h', required: true })
// 2. The release. From this commit on it is a record in the broker's
// replicated state, independent of this process. The key is ours,
// which is what lets the payer cancel it by name.
.timer(EXPIRIES).key(bookingId).delayMs(HOLD_MS).payload({ bookingId, room }).schedule()
// 3. The payment request, in the booking's own partition. It needs no
// transactionId: it can only commit together with the saga entry, so
// the gate above is its idempotency key.
.queue(PAYMENTS).partition(bookingId).push({ data: { bookingId, cents } })
// 4. The ack, with this delivery's lease. If the lease ran out, the ack is
// refused and the other three are refused with it.
.ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
.commit()
reserveDecisions.push(bookingId)
// A lost gate comes back as a value (HTTP 200, success: false, reason
// 'kv_precondition'), because a duplicate is a normal outcome and not an
// error to retry. Nothing was written, so the message is acked on its own.
if (res.success === false && res.reason === 'kv_precondition') {
gatesLost++
await queen.ack(msg, 'completed', { group: msg.consumerGroup })
console.log(` ${bookingId}: already held, nothing written (${res.kvReason})`)
return
}
console.log(` ${bookingId}: room ${room} held, release armed for ${HOLD_MS} ms`)
})
assert(
reserveDecisions.length === BOOKINGS_IN.length,
`the reserver decided every submission (${BOOKINGS_IN.length}, got ${reserveDecisions.length})`
)
assert(gatesLost === 1, 'the duplicate submission of B-2 lost the gate, once')
// Pending timers can be listed: the releases are records, not promises.
const armed = await queen.timer(EXPIRIES).list({ limit: 50 })
console.log(` timers armed: ${armed.rows.map(r => r.timerKey).sort().join(', ')}`)
assert(
armed.rows.length === BOOKING_IDS.length,
`one release per booking, none for the duplicate (${BOOKING_IDS.length}, got ${armed.rows.length})`
)
// ------------------------------------------------------------------ paying
//
// A settled payment confirms the saga and cancels the release in one commit.
// A declined card leaves the state alone and lets the timer do its work.
console.log('\npaying')
await phase(PAYMENTS, 'payer', BOOKING_IDS.length, PHASE_MS).consume(async (msg) => {
const { bookingId } = msg.data
paymentsRequested.push(bookingId)
// A read now and a write in the transaction below. That is safe here
// because the key derives from the partition key: every message about this
// booking is in one partition, and a partition is held by one worker of
// the group at a time.
const state = await queen.kv.get(NS, sagaKey(bookingId))
if (bookingId === DECLINED) {
// A declined card is a business outcome, not a failed delivery: the
// message is done. The room stays held, and nothing in this process is
// responsible for giving it back.
await queen.ack(msg, 'completed', { group: msg.consumerGroup })
console.log(` ${bookingId}: card declined, hold left to expire`)
return
}
const tx = queen
.transaction()
// expect makes the "one worker per booking" assumption checkable: if it
// ever fails, two consumers were serving one partition.
.kv.put(NS, sagaKey(bookingId), { ...state.value, step: 'confirmed' }, {
ttl: '1h',
expect: state.version,
required: true,
})
if (bookingId !== CANCEL_SKIPPED) {
// The cancel rides the transaction: the booking is confirmed and its
// release cancelled, or neither happens.
tx.timer(EXPIRIES).key(bookingId).cancel()
}
const res = await tx.ack(msg, 'completed', { consumerGroup: msg.consumerGroup }).commit()
if (res.success === false) throw new Error(`${bookingId}: confirmation lost its fence (${res.kvReason})`)
console.log(
` ${bookingId}: paid and confirmed, ` +
(bookingId === CANCEL_SKIPPED ? 'release deliberately NOT cancelled' : 'release cancelled')
)
})
assert(
paymentsRequested.length === BOOKING_IDS.length && new Set(paymentsRequested).size === BOOKING_IDS.length,
`every booking was asked to pay once, B-2 included (${paymentsRequested.join(', ')})`
)
// A cancelled timer is gone before it fires. peek answers { found: false }
// with HTTP 200 for a timer that does not exist.
const peeked = Object.fromEntries(
await Promise.all(BOOKING_IDS.map(async id => [id, (await queen.timer(EXPIRIES).key(id).peek()).found]))
)
assert(peeked['B-1'] === false, 'the release cancelled with the confirmation is gone')
assert(peeked[DECLINED] === true, `${DECLINED} was never confirmed, so its release is still armed`)
assert(peeked[CANCEL_SKIPPED] === true, `${CANCEL_SKIPPED} is confirmed and its release is still armed, on purpose`)
// ------------------------------------------------------------ compensating
//
// What the timers deliver, and the consumer that must not trust them. A
// release message asks a question: is this booking still only held? A fired
// timer leaves nothing behind, so a cancel that arrives a moment too late
// answers 'absent' and the release is delivered anyway. The saga's state
// decides.
//
// This message arrives on another queue, in a partition unrelated to the
// payments partition, so nothing serialises the compensator with the payer.
// Here expect is what stops a release computed from a stale read from
// overwriting a confirmation that landed in between.
console.log('\ncompensating')
const compensate = (limit, idleMillis) =>
phase(EXPIRIES, 'compensator', limit, idleMillis).consume(async (msg) => {
const { bookingId, room } = msg.data
compensationsDelivered.push(bookingId)
const state = await queen.kv.get(NS, sagaKey(bookingId))
if (!state.found || state.value.step !== 'held') {
compensationsRefused.push(bookingId)
await queen.ack(msg, 'completed', { group: msg.consumerGroup })
console.log(` ${bookingId}: state is ${state.value?.step ?? 'gone'}, release refused`)
return
}
const res = await queen
.transaction()
.kv.put(NS, sagaKey(bookingId), { ...state.value, step: 'expired' }, {
ttl: '1h',
expect: state.version,
required: true,
})
.ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
.commit()
if (res.success === false) {
// Confirmed between the read and the commit: nothing was written.
compensationsRefused.push(bookingId)
await queen.ack(msg, 'completed', { group: msg.consumerGroup })
console.log(` ${bookingId}: confirmed in the meantime, release refused`)
return
}
roomsReleased.push(room)
console.log(` ${bookingId}: hold expired, room ${room} released`)
})
// Two releases were left armed, so two messages have to arrive.
await compensate(2, TIMER_DEADLINE_MS)
assert(
compensationsDelivered.length === 2,
`both armed releases were delivered (got ${compensationsDelivered.join(', ') || 'none'})`
)
// A second, short pass with room for more: the only way to show that the
// cancelled releases never arrive is to wait for them and see nothing.
await compensate(2, 4000)
// ---------------------------------------------------------------- checking
console.log('\nchecking')
assert(
compensationsDelivered.length === 2 && !compensationsDelivered.includes('B-1') && !compensationsDelivered.includes('B-2'),
'no cancelled release was ever delivered'
)
assert(
roomsReleased.join(',') === '103',
`exactly one room went back on sale, the declined one (got ${roomsReleased.join(', ') || 'none'})`
)
assert(
compensationsRefused.join(',') === CANCEL_SKIPPED,
`the late release for ${CANCEL_SKIPPED} was refused by the compensator`
)
const states = await queen.kv.getMany(NS, BOOKING_IDS.map(sagaKey))
const step = Object.fromEntries(states.rows.map(r => [r.key.replace('saga:', ''), r.value.step]))
assert(
states.rows.length === BOOKING_IDS.length && states.missing.length === 0,
`every booking has exactly one saga entry (${states.rows.length})`
)
assert(
step['B-1'] === 'confirmed' && step['B-2'] === 'confirmed' && step[CANCEL_SKIPPED] === 'confirmed',
'B-1, B-2 and B-4 ended confirmed'
)
assert(step[DECLINED] === 'expired', `${DECLINED} was released by its timer, with no process waiting for it`)
console.log(`\n final: ${Object.entries(step).map(([k, v]) => `${k}=${v}`).sort().join(', ')}`)
console.log(`\nPASS: ${checks} checks`)
} catch (err) {
console.error(`\nFAIL: ${err.message}`)
process.exitCode = 1
} finally {
// Clean up in every case. Saga entries are in KV, and a pending timer is
// stored by its queue and key: deleting the queues removes neither, and a
// timer that fires into a deleted queue creates the queue again.
try {
for (const bookingId of BOOKING_IDS) {
await queen.timer(EXPIRIES).key(bookingId).cancel()
await queen.kv.delete(NS, sagaKey(bookingId))
}
for (const q of [BOOKINGS, PAYMENTS, EXPIRIES]) await queen.queue(q).delete()
} catch (err) {
console.error(` (cleanup incomplete: ${err.message})`)
}
await queen.close()
}#
# A booking saga: a room is held, paid for, and released by a timer when the
# payment never comes.
#
# The release is the part that usually breaks. Kept as a sleeping task in a
# worker, it dies with the worker, and a deploy in the middle of a hold leaves
# the room held forever. Here the release is a timer in the broker, scheduled
# in the same transaction that holds the room: the saga's state, the timer,
# the payment request and the ack of the booking commit as one log entry. If
# the room is held, its release exists; if anything failed, none of it
# happened.
#
# bookings
# `-- group "reserver" ONE transaction: state + timer + push + ack
# |-- payments (one partition per booking)
# | `-- group "payer" confirm + cancel the timer + ack
# `-- expiries (the timer delivers here when the hold runs out)
# `-- group "compensator" reads the state before it releases
#
# Run it:
# QUEEN_URL=http://localhost:6632 python3 saga.py
import asyncio
import os
import sys
import time
from datetime import timedelta
from queen import Queen
QUEEN_URL = os.environ.get("QUEEN_URL", "http://localhost:6632")
# Fresh queues and a fresh KV namespace per run. Saga entries outlive the
# queues (they expire with their TTL), so a second run in the same namespace
# would find every booking already held.
RUN = f"{int(time.time() * 1000):x}"
BOOKINGS = f"app-py-saga-bookings-{RUN}"
PAYMENTS = f"app-py-saga-payments-{RUN}"
EXPIRIES = f"app-py-saga-expiries-{RUN}"
NS = f"app-py-saga-{RUN}"
# How long a room stays held before it is released. In production this is
# minutes; it is the only number that changes. It has to outlast the reserving
# and paying phases below, or a release would fire before the payment that
# cancels it and the run would measure a race.
HOLD_MS = 10000
# Each phase ends on a count of messages, with this deadline behind it so a
# stall fails the run instead of hanging it.
PHASE_MS = 20000
# Timers fire on the leader's next tick after they are due (every 50 ms), so
# the compensation phase only needs the hold plus a margin.
TIMER_DEADLINE_MS = HOLD_MS + 20000
# Four bookings, five submissions: B-2 is submitted twice, which is what a
# redelivery looks like from the reserver's side.
BOOKINGS_IN = [
{"bookingId": "B-1", "room": "101", "cents": 24000},
{"bookingId": "B-2", "room": "102", "cents": 31000},
{"bookingId": "B-2", "room": "102", "cents": 31000},
{"bookingId": "B-3", "room": "103", "cents": 18000},
{"bookingId": "B-4", "room": "104", "cents": 27000},
]
BOOKING_IDS = ["B-1", "B-2", "B-3", "B-4"]
# B-3's card is declined, so its saga never reaches "confirmed" and the timer
# is what gives the room back.
DECLINED = "B-3"
# B-4 pays, but its cancel is skipped on purpose. That is the cancel that comes
# too late, made reproducible: the release is delivered for a booking that is
# already confirmed, and the compensator has to refuse it.
CANCEL_SKIPPED = "B-4"
CHECKS = 0
def check(condition: bool, description: str) -> None:
"""Record one verified fact, or abort the run.
This raises instead of using the `assert` statement, because `python3 -O`
removes `assert` and the checks are the whole point of the program.
"""
global CHECKS
if not condition:
raise AssertionError(description)
CHECKS += 1
print(f" ok: {description}")
def saga_key(booking_id: str) -> str:
"""The saga's state key derives from the booking id, which is also the
partition key of the payments queue. That is what makes the payer's
read-then-write safe below."""
return f"saga:{booking_id}"
async def main() -> int:
global CHECKS
queen = Queen(url=QUEEN_URL)
reserve_decisions: list = []
payments_requested: list = []
compensations_delivered: list = []
rooms_released: list = []
compensations_refused: list = []
gates_lost = 0
verdict, failed = "", False
# A consumer for one phase: acks ride the transactions, so auto_ack is off,
# and the phase ends after `limit` messages or after `idle_millis` with
# none. Long polls end after a second, so the idle bound is checked at
# least once a second.
def phase(queue: str, group: str, limit: int, idle_millis: int):
return (
queen.queue(queue)
.group(group)
.subscription_mode("all")
.auto_ack(False)
.each()
.limit(limit)
.timeout_millis(1000)
.idle_millis(idle_millis)
)
try:
print(f"broker {QUEEN_URL}")
for queue in (BOOKINGS, PAYMENTS, EXPIRIES):
await queen.queue(queue).config({"lease_time": 30, "retry_limit": 3}).create()
print("\nsubmitting bookings")
for index, booking in enumerate(BOOKINGS_IN):
await queen.queue(BOOKINGS).push(
{
# One id per submission, so the duplicate of B-2 is stored
# and reaches the reserver, where the gate has to catch it.
"transactionId": f"submit-{index}-{booking['bookingId']}",
"data": booking,
}
)
print(f" {len(BOOKINGS_IN)} submissions for {len(BOOKING_IDS)} bookings")
# ---------------------------------------------------------- reserving
#
# The transaction is the whole point of the example: four things commit
# together, so there is no order between them to get wrong. Written as
# four calls, a crash between the timer and the push leaves a release
# for a payment that was never asked for, and a crash the other way
# round leaves a hold with no release.
print("\nreserving")
async def reserve(msg) -> None:
nonlocal gates_lost
booking = msg["data"]
booking_id, room, cents = booking["bookingId"], booking["room"], booking["cents"]
group = msg.get("consumerGroup")
tx = queen.transaction()
# 1. The gate and the first state, in one KV entry. required=True
# makes the put_if_absent a gate: if the booking is already
# held, the whole transaction rolls back and nothing below
# happens.
tx.kv.put_if_absent(
NS,
saga_key(booking_id),
{"step": "held", "room": room, "cents": cents},
# The Python client takes a timedelta where the JavaScript one
# takes "1h"; both become ttlSeconds on the wire.
ttl=timedelta(hours=1),
required=True,
)
# 2. The release. From this commit on it is a record in the
# broker's replicated state, independent of this process. The
# key is ours, which is what lets the payer cancel it by name.
tx.timer(EXPIRIES).key(booking_id).after_ms(HOLD_MS).payload(
{"bookingId": booking_id, "room": room}
).schedule()
# 3. The payment request, in the booking's own partition. It needs
# no transactionId: it can only commit together with the saga
# entry, so the gate above is its idempotency key.
tx.queue(PAYMENTS).partition(booking_id).push({"data": {"bookingId": booking_id, "cents": cents}})
# 4. The ack, with this delivery's lease. If the lease ran out, the
# ack is refused and the other three are refused with it.
res = await tx.ack(msg, "completed", {"consumer_group": group}).commit()
reserve_decisions.append(booking_id)
# A lost gate comes back as a value (HTTP 200, success False,
# reason "kv_precondition"), because a duplicate is a normal
# outcome that must not be retried. Nothing was written, so the
# message is acked on its own.
if res.get("success") is False and res.get("reason") == "kv_precondition":
gates_lost += 1
await queen.ack(msg, "completed", {"group": group})
print(f" {booking_id}: already held, nothing written ({res.get('kvReason')})")
return
print(f" {booking_id}: room {room} held, release armed for {HOLD_MS} ms")
await phase(BOOKINGS, "reserver", len(BOOKINGS_IN), PHASE_MS).consume(reserve)
check(
len(reserve_decisions) == len(BOOKINGS_IN),
f"the reserver decided every submission ({len(BOOKINGS_IN)}, got {len(reserve_decisions)})",
)
check(gates_lost == 1, "the duplicate submission of B-2 lost the gate, once")
# Pending timers can be listed, because each release is a record in the
# broker.
armed = await queen.timers.list(EXPIRIES, limit=50)
print(f" timers armed: {', '.join(sorted(row['timerKey'] for row in armed['rows']))}")
check(
len(armed["rows"]) == len(BOOKING_IDS),
f"one release per booking, none for the duplicate ({len(BOOKING_IDS)}, got {len(armed['rows'])})",
)
# ------------------------------------------------------------- paying
#
# A settled payment confirms the saga and cancels the release in one
# commit. A declined card leaves the state alone and lets the timer do
# its work.
print("\npaying")
async def pay(msg) -> None:
booking_id = msg["data"]["bookingId"]
group = msg.get("consumerGroup")
payments_requested.append(booking_id)
# A read now and a write in the transaction below. That is safe
# here because the key derives from the partition key: every
# message about this booking is in one partition, and a partition
# is held by one worker of the group at a time.
state = await queen.kv.get(NS, saga_key(booking_id))
if booking_id == DECLINED:
# A declined card is a business outcome, and the message is
# done. The room stays held, and nothing in this process is
# responsible for giving it back.
await queen.ack(msg, "completed", {"group": group})
print(f" {booking_id}: card declined, hold left to expire")
return
tx = queen.transaction()
# expect makes the "one worker per booking" assumption checkable:
# if it ever fails, two consumers were serving one partition.
tx.kv.put(
NS,
saga_key(booking_id),
{**state["value"], "step": "confirmed"},
ttl=timedelta(hours=1),
expect=state["version"],
required=True,
)
if booking_id != CANCEL_SKIPPED:
# The cancel rides the transaction: the booking is confirmed
# and its release cancelled, or neither happens.
tx.timer(EXPIRIES).key(booking_id).cancel()
res = await tx.ack(msg, "completed", {"consumer_group": group}).commit()
if res.get("success") is False:
raise AssertionError(f"{booking_id}: confirmation lost its fence ({res.get('kvReason')})")
tail = "release deliberately NOT cancelled" if booking_id == CANCEL_SKIPPED else "release cancelled"
print(f" {booking_id}: paid and confirmed, {tail}")
await phase(PAYMENTS, "payer", len(BOOKING_IDS), PHASE_MS).consume(pay)
check(
len(payments_requested) == len(BOOKING_IDS),
f"every booking was asked to pay, B-2 included ({', '.join(payments_requested)})",
)
check(len(set(payments_requested)) == len(payments_requested), "no booking was asked to pay twice")
# A cancelled timer is gone before it fires. peek answers
# {"found": False} with HTTP 200 for a timer that does not exist.
peeked = {b: (await queen.timers.peek(EXPIRIES, b))["found"] for b in BOOKING_IDS}
check(peeked["B-1"] is False, "the release cancelled with the confirmation is gone")
check(peeked[DECLINED] is True, f"{DECLINED} was never confirmed, so its release is still armed")
check(
peeked[CANCEL_SKIPPED] is True,
f"{CANCEL_SKIPPED} is confirmed and its release is still armed, on purpose",
)
# ------------------------------------------------------- compensating
#
# What the timers deliver, and the consumer that must not trust them. A
# release message asks a question: is this booking still only held? A
# fired timer leaves nothing behind, so a cancel that arrives a moment
# too late answers "absent" and the release is delivered anyway. The
# saga's state decides.
#
# This message arrives on another queue, in a partition unrelated to
# the payments partition, so nothing serialises the compensator with
# the payer. Here expect is what stops a release computed from a stale
# read from overwriting a confirmation that landed in between.
print("\ncompensating")
async def compensate_one(msg) -> None:
booking_id, room = msg["data"]["bookingId"], msg["data"]["room"]
group = msg.get("consumerGroup")
compensations_delivered.append(booking_id)
state = await queen.kv.get(NS, saga_key(booking_id))
if not state["found"] or state["value"]["step"] != "held":
compensations_refused.append(booking_id)
await queen.ack(msg, "completed", {"group": group})
step = state["value"]["step"] if state["found"] else "gone"
print(f" {booking_id}: state is {step}, release refused")
return
res = await (
queen.transaction()
.kv.put(
NS,
saga_key(booking_id),
{**state["value"], "step": "expired"},
ttl=timedelta(hours=1),
expect=state["version"],
required=True,
)
.ack(msg, "completed", {"consumer_group": group})
.commit()
)
if res.get("success") is False:
# Confirmed between the read and the commit: nothing was
# written.
compensations_refused.append(booking_id)
await queen.ack(msg, "completed", {"group": group})
print(f" {booking_id}: confirmed in the meantime, release refused")
return
rooms_released.append(room)
print(f" {booking_id}: hold expired, room {room} released")
def compensator(limit: int, idle_millis: int):
return phase(EXPIRIES, "compensator", limit, idle_millis).consume(compensate_one)
# Two releases were left armed, so two messages have to arrive.
await compensator(2, TIMER_DEADLINE_MS)
check(
len(compensations_delivered) == 2,
f"both armed releases were delivered (got {', '.join(compensations_delivered) or 'none'})",
)
# A second, short pass with room for more: the only way to show that the
# cancelled releases never arrive is to wait for them and see nothing.
await compensator(2, 4000)
# ----------------------------------------------------------- checking
print("\nchecking")
check(
len(compensations_delivered) == 2,
f"nothing else arrived on the second pass (got {', '.join(compensations_delivered)})",
)
check(
"B-1" not in compensations_delivered and "B-2" not in compensations_delivered,
"no cancelled release was ever delivered",
)
check(
rooms_released == ["103"],
f"exactly one room went back on sale, the declined one (got {', '.join(rooms_released) or 'none'})",
)
check(
compensations_refused == [CANCEL_SKIPPED],
f"the late release for {CANCEL_SKIPPED} was refused by the compensator",
)
states = await queen.kv.get_many(NS, [saga_key(b) for b in BOOKING_IDS])
check(
len(states["rows"]) == len(BOOKING_IDS) and len(states["missing"]) == 0,
f"every booking has exactly one saga entry ({len(states['rows'])})",
)
step = {row["key"].replace("saga:", ""): row["value"]["step"] for row in states["rows"]}
check(
step["B-1"] == "confirmed" and step["B-2"] == "confirmed",
"B-1 and B-2 ended confirmed",
)
check(
step[CANCEL_SKIPPED] == "confirmed",
f"{CANCEL_SKIPPED} is still confirmed after its late release was delivered",
)
check(
step[DECLINED] == "expired",
f"{DECLINED} was released by its timer, with no process waiting for it",
)
print("\n final: " + ", ".join(f"{k}={v}" for k, v in sorted(step.items())))
verdict = f"\nPASS: {CHECKS} checks"
except Exception as err: # noqa: BLE001 - the program's verdict is its exit code
verdict, failed = f"\nFAIL: {err}", True
finally:
# Clean up in every case. Saga entries are in KV, and a pending timer
# is stored by its queue and key: deleting the queues removes neither,
# and a timer that fires into a deleted queue creates the queue again.
# Best effort: a cleanup that raised would replace the real verdict
# with its own.
try:
for booking_id in BOOKING_IDS:
await queen.timers.cancel(EXPIRIES, booking_id)
await queen.kv.delete(NS, saga_key(booking_id))
for queue in (BOOKINGS, PAYMENTS, EXPIRIES):
await queen.queue(queue).delete()
except Exception as err: # noqa: BLE001 - the run's verdict outranks this
print(f" (cleanup incomplete: {err})")
await queen.close()
sys.stdout.flush()
print(verdict, file=sys.stderr if failed else sys.stdout)
return 1 if failed else 0
if __name__ == "__main__":
sys.exit(asyncio.run(main()))//
// A booking saga: a room is held, paid for, and released by a timer when the
// payment never comes.
//
// The release is the part that usually breaks. Kept as a time.AfterFunc in a
// worker, it dies with the worker, and a deploy in the middle of a hold leaves
// the room held forever. Here the release is a timer in the broker, scheduled
// in the same transaction that holds the room: the saga's state, the timer,
// the payment request and the ack of the booking commit as one log entry. If
// the room is held, its release exists; if anything failed, none of it
// happened.
//
// bookings
// `-- group "reserver" ONE transaction: state + timer + push + ack
// |-- payments (one partition per booking)
// | `-- group "payer" confirm + cancel the timer + ack
// `-- expiries (the timer delivers here when the hold runs out)
// `-- group "compensator" reads the state before it releases
//
// Run it:
//
// QUEEN_URL=http://localhost:6632 GOWORK=off go run ./saga
package main
import (
"context"
"encoding/json"
"fmt"
"os"
"sort"
"strconv"
"strings"
"time"
queen "github.com/smartpricing/queen/clients/client-go/v2"
)
// Fresh queues and a fresh KV namespace per run. Saga entries outlive the queues
// (they expire with their TTL), so a second run in the same namespace would
// find every booking already held.
var runID = strconv.FormatInt(time.Now().UnixMilli(), 36)
var (
bookingsQueue = "app-go-saga-bookings-" + runID
paymentsQueue = "app-go-saga-payments-" + runID
expiriesQueue = "app-go-saga-expiries-" + runID
ns = "app-go-saga-" + runID
)
const (
reserverGroup = "reserver"
payerGroup = "payer"
compensatorGroup = "compensator"
// How long a room stays held before it is released. In production this is
// minutes; it is the only number that changes. It has to outlast the
// reserving and paying phases below, or a release would fire before the
// payment that cancels it and the run would measure a race.
hold = 10 * time.Second
// Each phase ends on a count of messages, with this deadline behind it so a
// stall fails the run instead of hanging it.
phaseMillis = 20000
// Timers fire on the leader's next tick after they are due (every 50 ms),
// so the compensation phase only needs the hold plus a margin.
timerDeadlineMillis = int((hold + 20*time.Second) / time.Millisecond)
// B-3's card is declined, so its saga never reaches "confirmed" and the
// timer is what gives the room back.
declined = "B-3"
// B-4 pays, but its cancel is skipped on purpose. That is the cancel that
// comes too late, made reproducible: the release is delivered for a booking
// that is already confirmed, and the compensator has to refuse it.
cancelSkipped = "B-4"
)
// Four bookings, five submissions: B-2 is submitted twice, which is what a
// redelivery looks like from the reserver's side.
type booking struct {
BookingID string `json:"bookingId"`
Room string `json:"room"`
Cents int `json:"cents"`
}
var bookingsIn = []booking{
{BookingID: "B-1", Room: "101", Cents: 24000},
{BookingID: "B-2", Room: "102", Cents: 31000},
{BookingID: "B-2", Room: "102", Cents: 31000},
{BookingID: "B-3", Room: "103", Cents: 18000},
{BookingID: "B-4", Room: "104", Cents: 27000},
}
var bookingIDs = []string{"B-1", "B-2", "B-3", "B-4"}
// sagaState is what the saga's KV entry holds. It is a struct because the value
// comes back as raw JSON and this program reads a field of it on every hop:
// with a map, a typo in a key would read as "the saga is not held" and the run
// would pass for the wrong reason.
type sagaState struct {
Step string `json:"step"`
Room string `json:"room"`
Cents int `json:"cents"`
}
// The saga's state key derives from the booking id, which is also the
// partition key of the payments queue. That is what makes the payer's
// read-then-write safe below.
func sagaKey(bookingID string) string { return "saga:" + bookingID }
var checks int
// assert is the whole test framework here. Go has no exceptions, so a failed
// check is an error that unwinds run() and is printed once, at the bottom.
func assert(condition bool, description string) error {
if !condition {
return fmt.Errorf("%s", description)
}
checks++
fmt.Printf(" ok: %s\n", description)
return nil
}
func main() {
if err := run(); err != nil {
fmt.Fprintf(os.Stderr, "\nFAIL: %v\n", err)
os.Exit(1)
}
fmt.Printf("\nPASS: %d checks\n", checks)
}
func run() error {
brokerURL := os.Getenv("QUEEN_URL")
if brokerURL == "" {
brokerURL = "http://localhost:6632"
}
// Every call in the Go client takes a context, and it is the only deadline
// there is. This one bounds the whole program, so a broker that stops
// answering fails the run when it expires.
ctx, cancel := context.WithTimeout(context.Background(), 300*time.Second)
defer cancel()
client, err := queen.New(brokerURL)
if err != nil {
return fmt.Errorf("create client: %w", err)
}
defer client.Close(context.Background())
// Clean up in every case, which is why this is deferred. Saga entries live
// in KV, and a pending timer is stored by its queue and key: deleting the
// queues removes neither, and a timer that fires into a deleted queue
// creates the queue again. It is best effort, on a context of its own, so a
// cleanup problem is reported without replacing the verdict.
defer func() {
cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cleanupCancel()
for _, bookingID := range bookingIDs {
if _, err := client.Timers().Cancel(cleanupCtx, expiriesQueue, bookingID); err != nil {
fmt.Fprintf(os.Stderr, " (cleanup incomplete: cancel %s: %v)\n", bookingID, err)
}
if _, err := client.KV().Delete(cleanupCtx, ns, sagaKey(bookingID)); err != nil {
fmt.Fprintf(os.Stderr, " (cleanup incomplete: delete %s: %v)\n", sagaKey(bookingID), err)
}
}
for _, q := range []string{bookingsQueue, paymentsQueue, expiriesQueue} {
if _, err := client.Queue(q).Delete().Execute(cleanupCtx); err != nil {
fmt.Fprintf(os.Stderr, " (cleanup incomplete: delete %s: %v)\n", q, err)
}
}
}()
// A consumer for one phase: acks ride the transactions, so AutoAck is off,
// and the phase ends after limit messages or after idleMillis with none.
// Concurrency stays at the default of one, so each handler below runs on a
// single goroutine, its bookkeeping needs no lock, and Limit (which this
// client counts per worker) is the count for the whole phase.
phase := func(queue, group string, limit, idleMillis int) *queen.QueueBuilder {
return client.Queue(queue).
Group(group).
SubscriptionMode(queen.SubscriptionModeAll).
AutoAck(false).
Each().
Limit(limit).
TimeoutMillis(1000).
IdleMillis(idleMillis)
}
fmt.Printf("broker %s\n", brokerURL)
for _, q := range []string{bookingsQueue, paymentsQueue, expiriesQueue} {
if _, err := client.Queue(q).
Config(queen.QueueConfig{LeaseTime: 30, RetryLimit: 3}).
Create().Execute(ctx); err != nil {
return fmt.Errorf("create %s: %w", q, err)
}
}
fmt.Println("\nsubmitting bookings")
for i, b := range bookingsIn {
if _, err := client.Queue(bookingsQueue).
Push(b).
// One id per submission, so the duplicate of B-2 is stored and
// reaches the reserver, where the gate has to catch it.
TransactionID(fmt.Sprintf("submit-%d-%s", i, b.BookingID)).
Execute(ctx); err != nil {
return fmt.Errorf("submit %s: %w", b.BookingID, err)
}
}
fmt.Printf(" %d submissions for %d bookings\n", len(bookingsIn), len(bookingIDs))
// -------------------------------------------------------------- reserving
//
// The transaction is the whole point of the example: four things commit
// together, so there is no order between them to get wrong. Written as four
// calls, a crash between the timer and the push leaves a release for a
// payment that was never asked for, and a crash the other way round leaves
// a hold with no release.
fmt.Println("\nreserving")
var reserveDecisions []string
gatesLost := 0
err = phase(bookingsQueue, reserverGroup, len(bookingsIn), phaseMillis).
Consume(ctx, func(ctx context.Context, msg *queen.Message) error {
bookingID, _ := msg.Data["bookingId"].(string)
room, _ := msg.Data["room"].(string)
cents, ok := msg.Data["cents"].(float64)
if !ok {
return fmt.Errorf("submission %s carries no numeric cents", msg.TransactionID)
}
res, err := client.Transaction().
// 1. The gate and the first state, in one KV entry. Required
// makes the putIfAbsent a gate: if the booking is already
// held, the whole transaction rolls back and nothing below
// happens.
KV(queen.KVPutIfAbsentOp(
ns,
sagaKey(bookingID),
sagaState{Step: "held", Room: room, Cents: int(cents)},
// Every KV write carries an expiry, and the zero value of
// queen.Expiry is refused. queen.Forever() exists, but an
// example that runs in CI must not be able to leave an
// entry behind for good.
queen.TTL(time.Hour),
queen.KVWriteOptions{Required: true},
)).
// 2. The release. From this commit on it is a record in the
// broker's replicated state, independent of this process.
// The key is ours, which is what lets the payer cancel it by
// name.
Timers(queen.ScheduleTimerOp(queen.TimerSchedule{
Queue: expiriesQueue,
TimerKey: bookingID,
Delay: hold,
Payload: map[string]interface{}{"bookingId": bookingID, "room": room},
})).
// 3. The payment request, in the booking's own partition. It
// needs no transaction id of its own (the client mints one):
// it can only commit together with the saga entry, so the gate
// above is its idempotency key.
Queue(paymentsQueue).
Partition(bookingID).
Push(map[string]interface{}{"bookingId": bookingID, "cents": int(cents)}).
// 4. The ack, with this delivery's lease. If the lease ran out,
// the ack is refused and the other three are refused with it.
Ack(msg, "completed", queen.AckOptions{ConsumerGroup: reserverGroup}).
Commit(ctx)
reserveDecisions = append(reserveDecisions, bookingID)
// A lost gate comes back as a value (err is nil, Success is false,
// Reason is "kv_precondition"), because a duplicate is a normal
// outcome and not an error to retry. Every other failed commit is
// an error. Nothing was written, so the message is acked on its
// own.
if res.IsKVPrecondition() {
gatesLost++
if _, err := client.Ack(ctx, msg, true, queen.AckOptions{ConsumerGroup: reserverGroup}); err != nil {
return fmt.Errorf("ack the duplicate submission: %w", err)
}
fmt.Printf(" %s: already held, nothing written (%s)\n", bookingID, res.KVReason)
return nil
}
if err != nil {
return fmt.Errorf("reserve %s: %w", bookingID, err)
}
fmt.Printf(" %s: room %s held, release armed for %s\n", bookingID, room, hold)
return nil
}).
Execute(ctx)
if err != nil {
return fmt.Errorf("reserving: %w", err)
}
if err := assert(
len(reserveDecisions) == len(bookingsIn),
fmt.Sprintf("the reserver decided every submission (%d, got %d)", len(bookingsIn), len(reserveDecisions)),
); err != nil {
return err
}
if err := assert(gatesLost == 1, "the duplicate submission of B-2 lost the gate, once"); err != nil {
return err
}
// Pending timers can be listed, because each release is a record in the
// broker.
armed, err := client.Timers().List(ctx, expiriesQueue, queen.TimerListOptions{Limit: 50})
if err != nil {
return fmt.Errorf("list the armed releases: %w", err)
}
armedKeys := make([]string, 0, len(armed.Rows))
for _, row := range armed.Rows {
armedKeys = append(armedKeys, row.TimerKey)
}
sort.Strings(armedKeys)
fmt.Printf(" timers armed: %s\n", strings.Join(armedKeys, ", "))
if err := assert(
len(armedKeys) == len(bookingIDs),
fmt.Sprintf("one release per booking, none for the duplicate (%d, got %d)", len(bookingIDs), len(armedKeys)),
); err != nil {
return err
}
// ----------------------------------------------------------------- paying
//
// A settled payment confirms the saga and cancels the release in one
// commit. A declined card leaves the state alone and lets the timer do its
// work.
fmt.Println("\npaying")
var paymentsRequested []string
err = phase(paymentsQueue, payerGroup, len(bookingIDs), phaseMillis).
Consume(ctx, func(ctx context.Context, msg *queen.Message) error {
bookingID, _ := msg.Data["bookingId"].(string)
paymentsRequested = append(paymentsRequested, bookingID)
// A read now and a write in the transaction below. That is safe
// here because the key derives from the partition key: every
// message about this booking is in one partition, and a partition
// is held by one worker of the group at a time.
state, version, err := readSaga(ctx, client, bookingID)
if err != nil {
return err
}
if bookingID == declined {
// A declined card is a business outcome, not a failed delivery:
// the message is done. The room stays held, and nothing in this
// process is responsible for giving it back.
if _, err := client.Ack(ctx, msg, true, queen.AckOptions{ConsumerGroup: payerGroup}); err != nil {
return fmt.Errorf("ack the declined payment: %w", err)
}
fmt.Printf(" %s: card declined, hold left to expire\n", bookingID)
return nil
}
state.Step = "confirmed"
tx := client.Transaction().
// Expect makes the "one worker per booking" assumption
// checkable: if it ever fails, two consumers were serving one
// partition.
KV(queen.KVPutOp(ns, sagaKey(bookingID), state, queen.TTL(time.Hour), queen.KVWriteOptions{
Expect: queen.Expect(version),
Required: true,
}))
if bookingID != cancelSkipped {
// The cancel rides the transaction: the booking is confirmed
// and its release cancelled, or neither happens.
tx = tx.Timers(queen.CancelTimerOp(expiriesQueue, bookingID))
}
res, err := tx.Ack(msg, "completed", queen.AckOptions{ConsumerGroup: payerGroup}).Commit(ctx)
if res.IsKVPrecondition() {
return fmt.Errorf("%s: confirmation lost its fence (%s)", bookingID, res.KVReason)
}
if err != nil {
return fmt.Errorf("confirm %s: %w", bookingID, err)
}
tail := "release cancelled"
if bookingID == cancelSkipped {
tail = "release deliberately NOT cancelled"
}
fmt.Printf(" %s: paid and confirmed, %s\n", bookingID, tail)
return nil
}).
Execute(ctx)
if err != nil {
return fmt.Errorf("paying: %w", err)
}
if err := assert(
len(paymentsRequested) == len(bookingIDs),
fmt.Sprintf("every booking was asked to pay once, B-2 included (%s)", orNone(strings.Join(paymentsRequested, ", "))),
); err != nil {
return err
}
if err := assert(distinct(paymentsRequested), "no booking was asked to pay twice"); err != nil {
return err
}
// A cancelled timer is gone before it fires. Peek answers Found: false
// with HTTP 200 for a timer that does not exist.
peeked := map[string]bool{}
for _, bookingID := range bookingIDs {
info, err := client.Timers().Peek(ctx, expiriesQueue, bookingID)
if err != nil {
return fmt.Errorf("peek %s: %w", bookingID, err)
}
peeked[bookingID] = info.Found
}
if err := assert(!peeked["B-1"], "the release cancelled with the confirmation is gone"); err != nil {
return err
}
if err := assert(peeked[declined], declined+" was never confirmed, so its release is still armed"); err != nil {
return err
}
if err := assert(peeked[cancelSkipped], cancelSkipped+" is confirmed and its release is still armed, on purpose"); err != nil {
return err
}
// ------------------------------------------------------------ compensating
//
// What the timers deliver, and the consumer that must not trust them. A
// release message asks a question: is this booking still only held? A
// fired timer leaves nothing behind, so a cancel that arrives a moment too
// late answers "absent" and the release is delivered anyway. The saga's
// state decides.
//
// This message arrives on another queue, in a partition unrelated to the
// payments partition, so nothing serialises the compensator with the
// payer. Here Expect is what stops a release computed from a stale read
// from overwriting a confirmation that landed in between.
fmt.Println("\ncompensating")
var compensationsDelivered []string
var roomsReleased []string
var compensationsRefused []string
compensate := func(limit, idleMillis int) error {
return phase(expiriesQueue, compensatorGroup, limit, idleMillis).
Consume(ctx, func(ctx context.Context, msg *queen.Message) error {
bookingID, _ := msg.Data["bookingId"].(string)
room, _ := msg.Data["room"].(string)
compensationsDelivered = append(compensationsDelivered, bookingID)
state, version, err := readSaga(ctx, client, bookingID)
if err != nil {
return err
}
if state.Step != "held" {
compensationsRefused = append(compensationsRefused, bookingID)
if _, err := client.Ack(ctx, msg, true, queen.AckOptions{ConsumerGroup: compensatorGroup}); err != nil {
return fmt.Errorf("ack the refused release: %w", err)
}
step := state.Step
if step == "" {
step = "gone"
}
fmt.Printf(" %s: state is %s, release refused\n", bookingID, step)
return nil
}
state.Step = "expired"
res, err := client.Transaction().
KV(queen.KVPutOp(ns, sagaKey(bookingID), state, queen.TTL(time.Hour), queen.KVWriteOptions{
Expect: queen.Expect(version),
Required: true,
})).
Ack(msg, "completed", queen.AckOptions{ConsumerGroup: compensatorGroup}).
Commit(ctx)
if res.IsKVPrecondition() {
// Confirmed between the read and the commit: nothing was
// written.
compensationsRefused = append(compensationsRefused, bookingID)
if _, err := client.Ack(ctx, msg, true, queen.AckOptions{ConsumerGroup: compensatorGroup}); err != nil {
return fmt.Errorf("ack the fenced release: %w", err)
}
fmt.Printf(" %s: confirmed in the meantime, release refused\n", bookingID)
return nil
}
if err != nil {
return fmt.Errorf("release %s: %w", bookingID, err)
}
roomsReleased = append(roomsReleased, room)
fmt.Printf(" %s: hold expired, room %s released\n", bookingID, room)
return nil
}).
Execute(ctx)
}
// Two releases were left armed, so two messages have to arrive.
if err := compensate(2, timerDeadlineMillis); err != nil {
return fmt.Errorf("compensating: %w", err)
}
if err := assert(
len(compensationsDelivered) == 2,
fmt.Sprintf("both armed releases were delivered (got %s)", orNone(strings.Join(compensationsDelivered, ", "))),
); err != nil {
return err
}
// A second, short pass with room for more: the only way to show that the
// cancelled releases never arrive is to wait for them and see nothing.
if err := compensate(2, 4000); err != nil {
return fmt.Errorf("second compensation pass: %w", err)
}
// --------------------------------------------------------------- checking
fmt.Println("\nchecking")
if err := assert(
len(compensationsDelivered) == 2,
fmt.Sprintf("nothing arrived on the second pass, still 2 releases (got %d)", len(compensationsDelivered)),
); err != nil {
return err
}
if err := assert(
!contains(compensationsDelivered, "B-1") && !contains(compensationsDelivered, "B-2"),
"no cancelled release was ever delivered",
); err != nil {
return err
}
if err := assert(
len(roomsReleased) == 1 && roomsReleased[0] == "103",
fmt.Sprintf("exactly one room went back on sale, the declined one (got %s)",
orNone(strings.Join(roomsReleased, ", "))),
); err != nil {
return err
}
if err := assert(
len(compensationsRefused) == 1 && compensationsRefused[0] == cancelSkipped,
"the late release for "+cancelSkipped+" was refused by the compensator",
); err != nil {
return err
}
keys := make([]string, 0, len(bookingIDs))
for _, bookingID := range bookingIDs {
keys = append(keys, sagaKey(bookingID))
}
states, err := client.KV().GetMany(ctx, ns, keys)
if err != nil {
return fmt.Errorf("read the saga entries: %w", err)
}
if err := assert(
len(states.Rows) == len(bookingIDs) && len(states.Missing) == 0,
fmt.Sprintf("every booking has exactly one saga entry (%d)", len(states.Rows)),
); err != nil {
return err
}
step := map[string]string{}
for _, row := range states.Rows {
var st sagaState
if err := json.Unmarshal(row.Value, &st); err != nil {
return fmt.Errorf("decode saga entry %s: %w", row.Key, err)
}
step[strings.TrimPrefix(row.Key, "saga:")] = st.Step
}
if err := assert(
step["B-1"] == "confirmed" && step["B-2"] == "confirmed",
"B-1 and B-2 ended confirmed",
); err != nil {
return err
}
if err := assert(
step[cancelSkipped] == "confirmed",
cancelSkipped+" ended confirmed, although its release was delivered",
); err != nil {
return err
}
if err := assert(
step[declined] == "expired",
declined+" was released by its timer, with no process waiting for it",
); err != nil {
return err
}
final := make([]string, 0, len(step))
for k, v := range step {
final = append(final, k+"="+v)
}
sort.Strings(final)
fmt.Printf("\n final: %s\n", strings.Join(final, ", "))
return nil
}
// readSaga reads one saga entry and its version. The version is what a later
// write passes back as Expect, so the two always travel together.
//
// An entry past its expiry reads as absent, and an absent entry comes back as
// the zero sagaState, whose Step is the empty string and so never "held".
func readSaga(ctx context.Context, client *queen.Queen, bookingID string) (sagaState, int64, error) {
var state sagaState
entry, err := client.KV().Get(ctx, ns, sagaKey(bookingID))
if err != nil {
return state, 0, fmt.Errorf("read the saga state of %s: %w", bookingID, err)
}
if !entry.Found {
return state, 0, nil
}
if err := json.Unmarshal(entry.Value, &state); err != nil {
return state, 0, fmt.Errorf("decode the saga state of %s: %w", bookingID, err)
}
return state, entry.Version, nil
}
func contains(values []string, want string) bool {
for _, v := range values {
if v == want {
return true
}
}
return false
}
func distinct(values []string) bool {
seen := map[string]bool{}
for _, v := range values {
if seen[v] {
return false
}
seen[v] = true
}
return true
}
func orNone(s string) string {
if s == "" {
return "none"
}
return s
}#!/usr/bin/env bash
#
# A booking saga, with nothing but curl: a room is held, paid for, and released
# by a timer when the payment never comes.
#
# The release is the part that usually breaks. Kept as a sleep in a worker, it
# dies with the worker, and a deploy in the middle of a hold leaves the room
# held forever. Here the release is a timer in the broker, scheduled in the same
# transaction that holds the room: the saga's state, the timer, the payment
# request and the ack of the booking commit as one log entry. If the room is
# held, its release exists; if anything failed, none of it happened.
#
# bookings
# `-- group "reserver" ONE transaction: state + timer + push + ack
# |-- payments (one partition per booking)
# | `-- group "payer" confirm + cancel the timer + ack
# `-- expiries (the timer delivers here when the hold runs out)
# `-- group "compensator" reads the state before it releases
#
# There is no client library here and none is needed, and this file is worth
# reading even if you use one. `kv` and `timers` are keys of the ROOT of the
# transaction request, beside `operations` and never elements of it, and a timer
# payload travels base64. Everything an SDK hides is written out.
#
# Run it:
# QUEEN_URL=http://localhost:6632 bash saga.sh
set -euo pipefail
QUEEN_URL="${QUEEN_URL:-http://localhost:6632}"
# Fresh queues and a fresh KV namespace per run, so runs never share state. The
# namespace matters most: the saga's state is KV entries, which outlive the
# queues (they expire with their TTL), so a second run in the same namespace
# would find every booking already held. $$ is the process id, which keeps two
# runs in the same second apart.
RUN="$(date +%s)-$$"
BOOKINGS="app-http-saga-bookings-$RUN"
PAYMENTS="app-http-saga-payments-$RUN"
EXPIRIES="app-http-saga-expiries-$RUN"
NS="app-http-saga-$RUN"
# A group's cursor lives on the queue, and the queue names are already unique
# per run, so these need no suffix.
RESERVER=app-http-reserver
PAYER=app-http-payer
COMPENSATOR=app-http-compensator
# Four bookings, five submissions: B-2 is submitted twice, which is what a
# redelivery looks like from the reserver's side, and the reason the
# transaction opens with a gate.
BOOKING_IDS="B-1 B-2 B-3 B-4"
SUBMISSIONS="B-1 B-2 B-2 B-3 B-4"
SUBMISSION_COUNT=5
BOOKING_COUNT=4
# B-3's card is declined, so its saga never reaches "confirmed" and the timer is
# what gives the room back.
DECLINED=B-3
# B-4 pays, but its cancel is skipped on purpose. That is the cancel that comes
# too late, made reproducible: a cancel that arrives after the fire answers
# `absent`, which may mean the release was already delivered. So the release
# for a confirmed booking has to be refused by the compensator that receives
# it.
CANCEL_SKIPPED=B-4
# How long a room stays held before it is released. In production this is
# minutes; it is the only number that changes. It has to outlast the reserving
# and paying phases below, or a release would fire before the payment that
# cancels it and the run would measure a race.
HOLD_MS=10000
# Each phase ends on a count of messages, with a deadline behind it so a stall
# fails the run instead of hanging it. Timers fire on the leader's next tick
# after they are due (every 50 ms), so the compensation phase only needs the
# hold plus a margin.
PHASE_MS=20000
TIMER_DEADLINE_MS=$((HOLD_MS + 20000))
# Every pop long-polls for this many milliseconds and no longer.
POLL_MS=1000
command -v jq >/dev/null 2>&1 || { echo "FAIL: jq is not installed"; exit 1; }
CHECKS=0
TMP="$(mktemp -d)"
# One line per delivery handled: "<group> <booking> <outcome>".
OBSERVED="$TMP/observed"
: > "$OBSERVED"
# One line per room actually put back on sale. This is the release's external
# effect, and the reason the program exists.
RELEASED="$TMP/released"
: > "$RELEASED"
fail() { FAILURE="$*"; exit 1; }
# check <actual> <expected> <description>
check() {
[ "$1" = "$2" ] || fail "$3 (expected [$2], got [$1])"
CHECKS=$((CHECKS + 1))
echo " ok: $3"
}
# A millisecond clock. GNU date spells it %3N; BSD date (macOS) has no %N and
# leaves the unconverted tail in the output, so a probe for anything that is not
# a digit tells the two apart, and perl, whose Time::HiRes is core, is the
# fallback.
if [ -z "$(date +%s%3N 2>/dev/null | tr -d '0-9')" ]; then
now_ms() { date +%s%3N; }
else
command -v perl >/dev/null 2>&1 \
|| { echo "FAIL: need GNU date or perl for a millisecond clock"; exit 1; }
now_ms() { perl -MTime::HiRes -e 'printf "%d", Time::HiRes::time() * 1000'; }
fi
# Sets $STATUS to the HTTP status code and writes the response body to $OUT.
# There is no --fail: Queen reports outcomes in the body and several of the
# interesting ones arrive as 200, so read the status, then the body.
OUT="$TMP/body"
request() {
local method="$1" path="$2" body="${3:-}"
if [ -n "$body" ]; then
STATUS="$(curl -sS -o "$OUT" -w '%{http_code}' \
-X "$method" "$QUEEN_URL$path" \
-H 'content-type: application/json' -d "$body")"
else
STATUS="$(curl -sS -o "$OUT" -w '%{http_code}' -X "$method" "$QUEEN_URL$path")"
fi
}
# saga_key <booking>: the state key. It derives from the booking id, which is
# also the partition key of the payments queue, and that derivation is what
# makes the payer's read-then-write safe further down.
saga_key() { printf 'saga:%s' "$1"; }
# kv_get <key>: prints the whole result as JSON, {"found":false,...} when
# absent. `found` is a field of its own because null is a legal stored value:
# absence is never inferred from the value being empty. Everything goes through
# POST /api/v1/kv, including single-key reads, which keeps the key out of access
# logs, proxy samples and tracing spans.
kv_get() {
local body
body="$(jq -cn --arg ns "$NS" --arg key "$1" \
'{operations: [{op: "get", ns: $ns, key: $key}]}')"
request POST /api/v1/kv "$body"
[ "$STATUS" = 200 ] || fail "kv get returned HTTP $STATUS"
jq -c '.results[0]' "$OUT"
}
# One exit path for everything. A failed check calls fail(), which records the
# reason and exits 1; any other command that fails under `set -e` arrives here
# too, with its own status. FAIL is printed exactly once, and only on failure.
#
# It also cleans up, in every case, a failed run included. Deleting the queues
# is not enough: the saga's state is KV entries, and a pending timer is stored
# by its queue and key, so deleting the queues removes neither, and a timer that
# fires into a deleted queue creates the queue again. So the timers are
# cancelled and the entries deleted first. The cleanup is best effort, with
# `|| true` throughout, so a cleanup that fails cannot overwrite the verdict.
cleanup() {
local status=$?
purge || true
rm -rf "$TMP"
if [ "$status" -ne 0 ]; then
echo
echo "FAIL: ${FAILURE:-a command exited with status $status}"
fi
exit "$status"
}
purge() {
local booking keys body
for booking in $BOOKING_IDS; do
# The cancel route, DELETE /api/v1/timers/:queue/*timerKey. Neither a quota
# nor an operator's switch ever blocks it: a pending timer fires whatever
# the quota says, so its cancel always has to get through.
request DELETE "/api/v1/timers/$EXPIRIES/$booking" || true
done
keys="$(printf '%s\n' $BOOKING_IDS | jq -R 'sub("^"; "saga:")' | jq -sc .)" || return 0
body="$(jq -cn --arg ns "$NS" --argjson keys "$keys" \
'{operations: [$keys[] | {op: "delete", ns: $ns, key: .}]}')" || return 0
request POST /api/v1/kv "$body" || true
request DELETE "/api/v1/resources/queues/$BOOKINGS" || true
request DELETE "/api/v1/resources/queues/$PAYMENTS" || true
request DELETE "/api/v1/resources/queues/$EXPIRIES" || true
}
trap cleanup EXIT
echo "broker $QUEEN_URL"
# Every broker serves /api/v1/kv and /api/v1/timers: there is no flag that turns
# them on. What can still refuse is an operator's runtime kill switch (503) or a
# quota (403), so probe once here and name that. Otherwise the first real call
# fails with something that reads like a bug.
request POST /api/v1/kv '{"operations":[{"op":"get","ns":"probe","key":"probe"}]}'
[ "$STATUS" = 200 ] \
|| fail "the kv probe returned HTTP $STATUS: $(cat "$OUT") (503 is an operator's kill switch, 403 a quota; GET /api/v1/system/kv-timers shows the switches)"
request GET "/api/v1/timers/$EXPIRIES?limit=1"
[ "$STATUS" = 200 ] \
|| fail "the timers probe returned HTTP $STATUS: $(cat "$OUT") (503 is an operator's kill switch, 403 a quota; GET /api/v1/system/kv-timers shows the switches)"
# /configure merges: an option not named here keeps whatever the queue already
# has. These three names are unique per run, so what is not named lands on its
# default.
for queue in "$BOOKINGS" "$PAYMENTS" "$EXPIRIES"; do
body="$(jq -n --arg queue "$queue" '{queue: $queue, options: {leaseTime: 30, retryLimit: 3}}')"
request POST /api/v1/configure "$body"
[ "$STATUS" = 200 ] || fail "configure of $queue returned HTTP $STATUS"
done
check "$(jq -r .configured "$OUT")" true 'three queues exist, each with a 30 second lease'
# ---------------------------------------------------------------------- queuing
echo
echo "submitting bookings"
index=0
room=101
cents=24000
for booking in $SUBMISSIONS; do
# One transaction id per submission, so the duplicate of B-2 is stored and
# reaches the reserver, where the gate has to catch it. The duplicate carries
# the same room and price, being the same booking submitted twice.
case "$booking" in
B-1) room=101; cents=24000 ;;
B-2) room=102; cents=31000 ;;
B-3) room=103; cents=18000 ;;
B-4) room=104; cents=27000 ;;
esac
body="$(jq -n --arg queue "$BOOKINGS" --arg booking "$booking" --arg room "$room" \
--argjson cents "$cents" --argjson i "$index" \
'{items: [{queue: $queue, transactionId: ("submit-" + ($i|tostring) + "-" + $booking),
payload: {bookingId: $booking, room: $room, cents: $cents}}]}')"
request POST /api/v1/push "$body"
[ "$STATUS" = 201 ] || fail "push of $booking returned HTTP $STATUS"
# HTTP 201 is not proof the message was stored: an item the broker refused
# comes back "error" inside a 201. The per-item status is the only answer.
[ "$(jq -r '.[0].status' "$OUT")" = queued ] \
|| fail "push of $booking came back $(jq -r '.[0].status' "$OUT")"
index=$((index + 1))
done
echo " $SUBMISSION_COUNT submissions for $BOOKING_COUNT bookings"
# ---------------------------------------------------------------------------
# handle_reserve: one delivery from the bookings queue.
#
# The transaction is the whole point of the example: four things commit
# together, so there is no order between them to get wrong. Written as four
# calls, a crash between the timer and the push leaves a release for a payment
# that was never asked for, and a crash the other way round leaves a hold with
# no release.
# ---------------------------------------------------------------------------
handle_reserve() {
local booking room cents txn partition lease body payload why ack_body
booking="$(jq -r '.messages[0].data.bookingId' "$TMP/pop")"
room="$(jq -r '.messages[0].data.room' "$TMP/pop")"
cents="$(jq -r '.messages[0].data.cents' "$TMP/pop")"
txn="$(jq -r '.messages[0].transactionId' "$TMP/pop")"
partition="$(jq -r '.messages[0].partitionId' "$TMP/pop")"
# The lease minted for THIS pop. It is what says the worker still owns the
# message, and it is the reason the acknowledgement below can refuse.
lease="$(jq -r '.leaseId' "$TMP/pop")"
# `kv` and `timers` are keys of the ROOT of this body, beside `operations`.
# They are separate top-level fields so that no client can send them under
# one key by accident.
#
# kv: the gate and the first state, in one KV entry. required:true makes
# the putIfAbsent a gate: if the booking is already held, the whole
# transaction rolls back and nothing else in it happens. Without it
# a lost race would come back applied:false while the payment and the
# timer went out anyway. ttlSeconds is mandatory on every KV write
# (`forever: true` is the only alternative), so an entry nobody
# deletes still goes away.
# timers: the release. From this commit on it is a record in the broker's
# replicated state, independent of this process. The key is ours,
# which is what lets the payer cancel it by name. The payload is
# base64, and delayMs is milliseconds from now: an absolute instant
# is not expressible, because the broker computes the due time on
# its own clock.
# push: the payment request, in the booking's own partition. It carries no
# transactionId: it can only commit together with the saga's state,
# so the gate is its idempotency key, and the broker mints an id.
# With an id of our own, the duplicate of B-2 would be refused as a
# duplicate push (reason "duplicate") before the gate could answer.
# ack: with this delivery's lease. If the lease ran out, the ack is
# refused and the other three are refused with it.
payload="$(jq -rn --arg booking "$booking" --arg room "$room" \
'{bookingId: $booking, room: $room} | tojson | @base64')"
body="$(jq -cn --arg ns "$NS" --arg key "$(saga_key "$booking")" \
--arg booking "$booking" --arg room "$room" --argjson cents "$cents" \
--arg payments "$PAYMENTS" --arg expiries "$EXPIRIES" --argjson hold "$HOLD_MS" \
--arg payload "$payload" \
--arg txn "$txn" --arg pid "$partition" --arg grp "$RESERVER" --arg lease "$lease" '
{operations: [{type: "push",
items: [{queue: $payments, partition: $booking,
payload: {bookingId: $booking, cents: $cents}}]},
{type: "ack", transactionId: $txn, partitionId: $pid,
consumerGroup: $grp, leaseId: $lease, status: "completed"}],
kv: [{op: "putIfAbsent", ns: $ns, key: $key,
value: {step: "held", room: $room, cents: $cents},
ttlSeconds: 3600, required: true}],
timers: [{op: "schedule", queue: $expiries, timerKey: $booking,
delayMs: $hold, txn: ("hold-" + $booking), payload: $payload}]}')"
request POST /api/v1/transaction "$body"
[ "$STATUS" = 200 ] || fail "the reserving transaction for $booking returned HTTP $STATUS: $(cat "$OUT")"
# A lost gate comes back as a value: HTTP 200 with success:false and reason
# "kv_precondition". A duplicate is a normal outcome, so it is handled here
# and kept out of the error path, where a retry would be the reflex.
if [ "$(jq -r '.success' "$OUT")" != true ]; then
[ "$(jq -r '.reason' "$OUT")" = kv_precondition ] \
|| fail "the reserving transaction for $booking failed: $(jq -r '.error' "$OUT")"
# Read the verdict before the next call: $OUT is one file and the ack below
# overwrites it. `kvReason` says why the precondition failed: here
# `exists`, the booking's state was already there.
why="$(jq -r '.kvReason' "$OUT")"
# Nothing was written: no second payment, no second timer, no second state.
# The message still has to leave the cursor, so it is acknowledged on its
# own.
ack_body="$(jq -cn --arg txn "$txn" --arg pid "$partition" --arg grp "$RESERVER" --arg lease "$lease" \
'{transactionId: $txn, partitionId: $pid, consumerGroup: $grp, leaseId: $lease, status: "completed"}')"
request POST /api/v1/ack "$ack_body"
[ "$STATUS" = 200 ] || fail "ack returned HTTP $STATUS"
printf '%s %s rolled-back\n' "$RESERVER" "$booking" >> "$OBSERVED"
echo " $booking: already held, nothing written ($why)"
return 0
fi
printf '%s %s held\n' "$RESERVER" "$booking" >> "$OBSERVED"
echo " $booking: room $room held, release armed for $HOLD_MS ms"
}
# ---------------------------------------------------------------------------
# handle_pay: one delivery from the payments queue.
#
# A settled payment confirms the saga and cancels the release in one commit. A
# declined card leaves the state alone and lets the timer do its work.
# ---------------------------------------------------------------------------
handle_pay() {
local booking txn partition lease state version value body timers ack_body
booking="$(jq -r '.messages[0].data.bookingId' "$TMP/pop")"
txn="$(jq -r '.messages[0].transactionId' "$TMP/pop")"
partition="$(jq -r '.messages[0].partitionId' "$TMP/pop")"
lease="$(jq -r '.leaseId' "$TMP/pop")"
# A read now and a write in the transaction below. That is safe here because
# the key derives from the partition key: every message about this booking is
# in one partition, and a partition is held by one worker of the group at a
# time. Where a key does not derive from the partition key this shape is a
# race, which is the compensator's situation further down.
state="$(kv_get "$(saga_key "$booking")")"
version="$(printf '%s' "$state" | jq -r '.version')"
value="$(printf '%s' "$state" | jq -c '.value')"
if [ "$booking" = "$DECLINED" ]; then
# A declined card is a business outcome, not a failed delivery: the message
# is done. The room stays held, and nothing in this process is responsible
# for giving it back.
ack_body="$(jq -cn --arg txn "$txn" --arg pid "$partition" --arg grp "$PAYER" --arg lease "$lease" \
'{transactionId: $txn, partitionId: $pid, consumerGroup: $grp, leaseId: $lease, status: "completed"}')"
request POST /api/v1/ack "$ack_body"
[ "$STATUS" = 200 ] || fail "ack returned HTTP $STATUS"
printf '%s %s declined\n' "$PAYER" "$booking" >> "$OBSERVED"
echo " $booking: card declined, hold left to expire"
return 0
fi
# The cancel rides the transaction: the booking is confirmed and its release
# cancelled, or neither happens. Inside a transaction a cancel travels in the
# timers array and shares the transaction's fate.
if [ "$booking" = "$CANCEL_SKIPPED" ]; then
timers='[]'
else
timers="$(jq -cn --arg expiries "$EXPIRIES" --arg booking "$booking" \
'[{op: "cancel", queue: $expiries, timerKey: $booking}]')"
fi
# `expect` makes the "one worker per booking" assumption checkable: if it ever
# fails, two consumers were serving one partition, and the run says so.
body="$(jq -cn --arg ns "$NS" --arg key "$(saga_key "$booking")" \
--argjson value "$value" --argjson version "$version" --argjson timers "$timers" \
--arg txn "$txn" --arg pid "$partition" --arg grp "$PAYER" --arg lease "$lease" '
{operations: [{type: "ack", transactionId: $txn, partitionId: $pid,
consumerGroup: $grp, leaseId: $lease, status: "completed"}],
kv: [{op: "put", ns: $ns, key: $key, value: ($value + {step: "confirmed"}),
ttlSeconds: 3600, expect: $version, required: true}],
timers: $timers}')"
request POST /api/v1/transaction "$body"
[ "$STATUS" = 200 ] || fail "the confirming transaction for $booking returned HTTP $STATUS: $(cat "$OUT")"
[ "$(jq -r '.success' "$OUT")" = true ] \
|| fail "$booking: confirmation lost its fence ($(jq -r '.kvReason' "$OUT"))"
printf '%s %s confirmed\n' "$PAYER" "$booking" >> "$OBSERVED"
if [ "$booking" = "$CANCEL_SKIPPED" ]; then
echo " $booking: paid and confirmed, release deliberately NOT cancelled"
else
echo " $booking: paid and confirmed, release cancelled"
fi
}
# ---------------------------------------------------------------------------
# handle_compensate: one delivery from the expiries queue, which is to say one
# message a timer produced.
#
# A release message asks a question: is this booking still only held? A fired
# timer leaves nothing behind, so a cancel that arrives a moment too late
# answers `absent` and the release is delivered anyway. The saga's state
# decides, and it is read first.
#
# This message arrives on another queue, in a partition unrelated to the
# payments partition, so nothing serialises the compensator with the payer.
# Here `expect` is what stops a release computed from a stale read from
# overwriting a confirmation that landed in between.
# ---------------------------------------------------------------------------
handle_compensate() {
local booking room txn partition lease state step version value body ack_body
booking="$(jq -r '.messages[0].data.bookingId' "$TMP/pop")"
room="$(jq -r '.messages[0].data.room' "$TMP/pop")"
txn="$(jq -r '.messages[0].transactionId' "$TMP/pop")"
partition="$(jq -r '.messages[0].partitionId' "$TMP/pop")"
lease="$(jq -r '.leaseId' "$TMP/pop")"
state="$(kv_get "$(saga_key "$booking")")"
step="$(printf '%s' "$state" | jq -r '.value.step // "gone"')"
version="$(printf '%s' "$state" | jq -r '.version')"
value="$(printf '%s' "$state" | jq -c '.value')"
ack_body="$(jq -cn --arg txn "$txn" --arg pid "$partition" --arg grp "$COMPENSATOR" --arg lease "$lease" \
'{transactionId: $txn, partitionId: $pid, consumerGroup: $grp, leaseId: $lease, status: "completed"}')"
if [ "$step" != held ]; then
# The booking was confirmed before this fired. Releasing the room now would
# put a sold room back on sale.
request POST /api/v1/ack "$ack_body"
[ "$STATUS" = 200 ] || fail "ack returned HTTP $STATUS"
printf '%s %s refused\n' "$COMPENSATOR" "$booking" >> "$OBSERVED"
echo " $booking: state is $step, release refused"
return 0
fi
body="$(jq -cn --arg ns "$NS" --arg key "$(saga_key "$booking")" \
--argjson value "$value" --argjson version "$version" \
--arg txn "$txn" --arg pid "$partition" --arg grp "$COMPENSATOR" --arg lease "$lease" '
{operations: [{type: "ack", transactionId: $txn, partitionId: $pid,
consumerGroup: $grp, leaseId: $lease, status: "completed"}],
kv: [{op: "put", ns: $ns, key: $key, value: ($value + {step: "expired"}),
ttlSeconds: 3600, expect: $version, required: true}]}')"
request POST /api/v1/transaction "$body"
[ "$STATUS" = 200 ] || fail "the compensating transaction for $booking returned HTTP $STATUS: $(cat "$OUT")"
if [ "$(jq -r '.success' "$OUT")" != true ]; then
# Confirmed between the read and the commit. The fence held, nothing was
# written, and the room stays sold.
[ "$(jq -r '.reason' "$OUT")" = kv_precondition ] \
|| fail "the compensating transaction for $booking failed: $(jq -r '.error' "$OUT")"
request POST /api/v1/ack "$ack_body"
[ "$STATUS" = 200 ] || fail "ack returned HTTP $STATUS"
printf '%s %s refused\n' "$COMPENSATOR" "$booking" >> "$OBSERVED"
echo " $booking: confirmed in the meantime, release refused"
return 0
fi
printf '%s\n' "$room" >> "$RELEASED"
printf '%s %s released\n' "$COMPENSATOR" "$booking" >> "$OBSERVED"
echo " $booking: hold expired, room $room released"
}
# ---------------------------------------------------------------------------
# drain <queue> <group> <handler> <deliveries> <deadline_ms>: pop and handle
# until this group has handled that many deliveries, or the deadline passes. The
# count is the bound and the deadline is the net; neither is a wait for silence.
# ---------------------------------------------------------------------------
drain() {
local queue="$1" group="$2" handler="$3" wanted="$4" budget="$5" deadline
deadline=$(( $(now_ms) + budget ))
while [ "$(grep -c "^$group " "$OBSERVED" || true)" -lt "$wanted" ]; do
[ "$(now_ms)" -lt "$deadline" ] || break
# subscriptionMode=all is what makes a group created now read what was
# pushed before it existed: a new cursor is seeded at the TAIL unless you
# say otherwise. batch=1 keeps one message in flight.
request GET "/api/v1/pop/queue/$queue?consumerGroup=$group&subscriptionMode=all&batch=1&wait=true&timeout=$POLL_MS"
# 204 is an empty pop, with no body at all. Go round again until the
# deadline.
[ "$STATUS" != 204 ] || continue
[ "$STATUS" = 200 ] || fail "pop returned HTTP $STATUS"
cp "$OUT" "$TMP/pop"
"$handler"
done
}
# -------------------------------------------------------------------- reserving
echo
echo "reserving"
drain "$BOOKINGS" "$RESERVER" handle_reserve "$SUBMISSION_COUNT" "$PHASE_MS"
check "$(grep -c "^$RESERVER " "$OBSERVED" || true)" "$SUBMISSION_COUNT" \
'the reserver decided every submission'
check "$(grep -c "^$RESERVER .* rolled-back$" "$OBSERVED" || true)" 1 \
'the duplicate submission of B-2 lost the gate, once'
# Pending timers can be listed: each release is a record in the broker.
request GET "/api/v1/timers/$EXPIRIES?limit=50"
[ "$STATUS" = 200 ] || fail "listing timers returned HTTP $STATUS"
echo " timers armed: $(jq -r '[.rows[].timerKey] | sort | join(", ")' "$OUT")"
check "$(jq -r '.rows | length' "$OUT")" "$BOOKING_COUNT" \
'one release per booking, none for the duplicate'
# ----------------------------------------------------------------------- paying
echo
echo "paying"
drain "$PAYMENTS" "$PAYER" handle_pay "$BOOKING_COUNT" "$PHASE_MS"
check "$(grep -c "^$PAYER " "$OBSERVED" || true)" "$BOOKING_COUNT" \
'every booking was asked to pay once, B-2 included'
check "$(awk -v g="$PAYER" '$1 == g {print $2}' "$OBSERVED" | sort -u | wc -l | tr -d ' ')" \
"$BOOKING_COUNT" 'no booking was asked to pay twice'
# A cancelled timer is gone before it fires. A peek is how you ask, and it
# answers {"found":false} with HTTP 200 for a timer that does not exist.
request GET "/api/v1/timers/$EXPIRIES/B-1"
[ "$STATUS" = 200 ] || fail "peek returned HTTP $STATUS"
check "$(jq -r '.found' "$OUT")" false \
'the release cancelled with the confirmation is gone'
request GET "/api/v1/timers/$EXPIRIES/$DECLINED"
check "$(jq -r '.found' "$OUT")" true \
"$DECLINED was never confirmed, so its release is still armed"
request GET "/api/v1/timers/$EXPIRIES/$CANCEL_SKIPPED"
check "$(jq -r '.found' "$OUT")" true \
"$CANCEL_SKIPPED is confirmed and its release is still armed, on purpose"
# ----------------------------------------------------------------- compensating
echo
echo "compensating"
# Two releases were left armed, so two messages have to arrive: that is the
# count, and TIMER_DEADLINE_MS is the deadline behind it.
drain "$EXPIRIES" "$COMPENSATOR" handle_compensate 2 "$TIMER_DEADLINE_MS"
check "$(grep -c "^$COMPENSATOR " "$OBSERVED" || true)" 2 \
'both armed releases were delivered'
# Then a short second pass with room for two more. The only way to show that
# the cancelled releases never arrive is to wait for them and see nothing: the
# first pass would have stopped at two, whatever those two were.
drain "$EXPIRIES" "$COMPENSATOR" handle_compensate 4 4000
# --------------------------------------------------------------------- checking
echo
echo "checking"
check "$(grep -c "^$COMPENSATOR " "$OBSERVED" || true)" 2 \
'nothing else arrived on a second pass: still 2 releases'
check "$(grep -c "^$COMPENSATOR B-1 \|^$COMPENSATOR B-2 " "$OBSERVED" || true)" 0 \
'no cancelled release was ever delivered'
check "$(cat "$RELEASED" | tr -d ' \n')" 103 \
'exactly one room went back on sale, the declined one'
check "$(grep -c "^$COMPENSATOR $CANCEL_SKIPPED refused$" "$OBSERVED" || true)" 1 \
"the late release for $CANCEL_SKIPPED was refused by the compensator"
# The saga's state is readable: "what happened to this booking" has an answer
# in the broker. getMany reports `missing` explicitly, so an absent key is named
# in the answer and never inferred by difference.
keys="$(printf '%s\n' $BOOKING_IDS | jq -R 'sub("^"; "saga:")' | jq -sc .)"
body="$(jq -cn --arg ns "$NS" --argjson keys "$keys" \
'{operations: [{op: "getMany", ns: $ns, keys: $keys}]}')"
request POST /api/v1/kv "$body"
[ "$STATUS" = 200 ] || fail "kv getMany returned HTTP $STATUS"
check "$(jq -r '.results[0].rows | length' "$OUT")" "$BOOKING_COUNT" \
'every booking has exactly one saga entry'
check "$(jq -r '.results[0].missing | length' "$OUT")" 0 \
'no booking is missing its saga entry'
check "$(jq -r '[.results[0].rows[] | select(.value.step == "confirmed") | .key] | sort | join(",")' "$OUT")" \
"saga:B-1,saga:B-2,saga:$CANCEL_SKIPPED" \
"B-1, B-2 and $CANCEL_SKIPPED ended confirmed"
check "$(jq -r --arg key "saga:$DECLINED" '.results[0].rows[] | select(.key == $key) | .value.step' "$OUT")" \
expired "$DECLINED was released by its timer, with no process waiting for it"
echo
echo " final: $(jq -r '[.results[0].rows[] | (.key | sub("^saga:"; "")) + "=" + .value.step] | sort | join(", ")' "$OUT")"
echo
echo "PASS: $CHECKS checks"Limits
Queen does not know your states. The table is your code: Queen orders the events and commits each step, and a wrong transition commits as faithfully as a right one.
A lease fences the workers of one group from each other and does nothing between groups, so other groups may read the same queue for projections or an audit trail, but only one group should write the machine’s KV keys.
A hot entity is sequential, because one order runs one step at a time. Throughput comes from running many orders at once.
A side effect your handler performs itself, an email or a call to another API, can repeat on redelivery. Emit it as an event and handle it with a marker, as in charge a card once.
The state is one KV value, 64 KiB by default (QUEEN_KV_MAX_VALUE_BYTES). POST /api/v1/timers
refuses a timer more than 90 days out (QUEEN_TIMERS_MAX_HORIZON_S), but in 2.0.0-beta.6 a timer
scheduled inside a transaction is not checked against that horizon, so keep your own delays
sensible.
The transaction, the queue and timer firing are all covered by Jepsen: the P26 campaign on 2.0.0-beta.3 ran them under kills, partitions, clock faults and power loss, 65 tests, all valid (guarantees, Jepsen).