Skip to content

Charge a card once

Make a redelivered order harmless: write a charge marker in the same transaction as the ack, gate it with putIfAbsent, and let the lease fence a worker that ran out of time.

Updated View as Markdown

A consumer that charges cards will see some orders twice: a lease runs out while the card network is slow, a node restarts, someone replays the queue. Make the second delivery harmless by writing a marker in the same transaction as the ack, so an order is either charged and acknowledged or neither. That gives you exactly-once effects inside Queen. The charge itself is a call to another system, and the last section is honest about what that leaves open.

The fix

The marker lives in Queen’s KV store and commits with the ack:

await queen.queue('orders')
  .group('charger')
  .subscriptionMode('all')
  .batch(1)                                         // one order per lease, see below
  .autoAck(false)
  .each()
  .consume(async (msg) => {
    const key = `charge:${msg.data.orderId}`
    const marker = await queen.kv.get('charges', key)
    if (marker.found) {                               // a redelivery: already charged
      await queen.ack(msg, 'completed', { group: msg.consumerGroup })
      return
    }

    try {
      const chargeId = await chargeCard(msg.data)     // the external effect
      const res = await queen.transaction()
        .kv.putIfAbsent('charges', key, { chargeId }, { ttl: '30d', required: true })
        .ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
        .commit()
      if (res.success === false) {                    // kv_precondition: another delivery won
        await queen.ack(msg, 'completed', { group: msg.consumerGroup })
      }
    } catch (err) {
      if (err.reason === 'rejected_ack') return       // lease lost: nothing was written
      await queen.ack(msg, 'failed', { group: msg.consumerGroup, error: err.message })
    }
  })

Two details in that consumer matter more than they look. .batch(1) keeps each charge inside its own lease: a failed ack releases the partition’s lease, so with a bigger batch the orders popped after a failure would be charged, refused at commit with rejected_ack, and charged again when they came back. And the charge sits inside the try so the handler chooses each ending itself: a lost lease (rejected_ack) needs no ack at all, and any other failure is acked as failed with its error. An error that escaped would be nacked by the SDK, and the consumer would keep going. The complete program, with the ledger that proves it, is the charge each order once example.

The usual worker keeps its “already charged” flag in a second system: it charges, writes the flag to a database, then acks the broker. The flag and the ack commit separately, so they can disagree, and nothing stops a worker whose lease has expired from writing its flag after another worker took the message. Here the marker exists exactly when the order is acknowledged, because they are one entry in the log. A worker that dies before the commit leaves no marker behind, and the redelivery runs the handler again.

A sequence between a worker of the charger group, Queen and the card provider. The worker pops one order with its lease, reads the marker charge:1042 from KV and finds nothing, then asks the provider to charge the card with the order id as idempotency key and gets a chargeId back. It commits one transaction that puts the marker with putIfAbsent and required: true and acks the order under the lease. Queen writes both as one log entry if the lease still holds and no marker exists, and answers success. In the alternative where the lease ran out during the charge, the commit is refused with rejected_ack: neither the marker nor the ack is written, and the order goes to the next worker.workergroup chargerQueencard providerpop orders, batch 1order 1042 + leasekv.get charge:1042not foundcharge the cardorder id as idempotency keychargeIdtransactionputIfAbsent marker + acklease held, no marker yet(required: true): one log entrysuccessthe same commit, after the lease ran outrejected_ack: no marker, no ack;the order goes to the next worker
The marker and the ack are one log entry, so the marker exists exactly when the order is acknowledged. The provider's call is outside that entry, which is why the order id goes with it as an idempotency key. Source: examples/apps/js/exactly-once.mjs, clients/client-js/client-v2/builders/TransactionBuilder.js

required: true is what turns putIfAbsent into a gate. Without it, losing the race to an existing marker is only a verdict in the results and the ack still commits. With it, the whole transaction rolls back and commit() returns success: false with reason: 'kv_precondition', over HTTP 200, since for this code a duplicate is an expected outcome that the caller checks for. The gate covers the case the read at the top cannot: two deliveries of one order in flight at once (the order pushed twice, or read by two groups). Only one of them commits its marker, and with it everything else in its transaction.

The lease is the fence. If the worker’s lease expired while it was charging, the broker refuses the ack, and the marker write with it, and commit() throws with reason: 'rejected_ack'. A compare-and-set cannot do this on its own, because a version that still matches succeeds for a worker that no longer owns the message.

once is the same gate under a shorter name inside a transaction: .once('charges', key, { ttl: '30d', value: { chargeId } }) is putIfAbsent with required: true. The SDKs send the pop’s lease with every transaction; over HTTP, put the pop’s leaseId on the ack operation yourself, as the curl program below does, because an ack without a lease is not fenced. Transactions and dedup cover both mechanisms in full.

The verified programs

Each program charges five orders. The first attempt at one of them fails before it charges anything, and must leave no trace. Then a second consumer group replays the whole queue, finds a marker for every order, and must charge nothing, so the ledger still holds exactly five charges.

examples/apps/js/exactly-once.mjsjs
//
// Charging each order exactly once, when orders can be delivered twice.
//
// Any consumer that charges cards sees some orders twice: a lease runs out
// while the card network is slow, a node restarts, someone replays the queue.
// The usual guard is an "already charged" flag in a database next to the
// broker, written after the charge and before the ack. The flag and the ack
// then commit separately, and in the gap between them the work can happen
// twice.
//
// Here the flag is a KV entry in the broker itself, written in the same
// transaction as the ack. The order is marked and acknowledged together, or
// neither happens, and a redelivery finds the marker and charges nothing.
//
//   orders
//     ├── group "charger"  marker + ack in ONE transaction
//     └── group "replay"   reads every order again and must charge nothing
//
// Run it:
//   QUEEN_URL=http://localhost:6632 node exactly-once.mjs

import { Queen } from 'queen-mq'

const QUEEN_URL = process.env.QUEEN_URL || 'http://localhost:6632'
const RUN = Date.now().toString(36)

// A fresh queue and a fresh KV namespace per run. The namespace needs it more
// than the queue does: markers outlive the queue that produced them (they
// expire with their TTL), so a second run in the same namespace would find
// every order already charged and pass without charging anything.
const ORDERS = `app-js-exactly-once-${RUN}`
const NS = `app-js-exactly-once-${RUN}`

// Five orders. The first attempt at ORD-3 fails before it charges anything, the
// failure that must leave no trace at all.
const ORDER_IDS = ['ORD-1', 'ORD-2', 'ORD-3', 'ORD-4', 'ORD-5']
const CRASHING_ORDER = 'ORD-3'

// The charging phase ends after six deliveries, five orders plus the retry of
// ORD-3, and the deadline is there so that a stall fails the run instead of
// hanging it.
const CHARGE_DELIVERIES = ORDER_IDS.length + 1
const PHASE_MS = 30000

let checks = 0
const assert = (condition, description) => {
  if (!condition) throw new Error(description)
  checks++
  console.log(`  ok: ${description}`)
}

// The external effect. Each call is a charge on a real card, and the ledger is
// what the customers' statements would show.
const ledger = []
const chargeCard = (order) => {
  const chargeId = `ch_${order.orderId}_${ledger.length}`
  ledger.push({ orderId: order.orderId, chargeId, cents: order.cents })
  return chargeId
}

const attempts = new Map()
const markerKey = (orderId) => `charge:${orderId}`

const queen = new Queen({ url: QUEEN_URL, handleSignals: false })

try {
  console.log(`broker ${QUEEN_URL}`)

  await queen.queue(ORDERS).config({ leaseTime: 30, retryLimit: 5 }).create()

  console.log('\nqueuing orders')
  for (const orderId of ORDER_IDS) {
    await queen.queue(ORDERS).push({
      transactionId: `order-${orderId}`,
      data: { orderId, cents: 1000 + ORDER_IDS.indexOf(orderId) },
    })
  }
  console.log(`  ${ORDER_IDS.length} orders queued`)

  // ---------------------------------------------------------------- charging
  //
  // handle() is the whole pattern: four steps, in this order. It returns
  // { ran: true } when this delivery charged the card and { ran: false } when it
  // found the order already charged. Every redelivery, whatever caused it, has
  // to come back { ran: false }.
  const observed = []

  const handle = async (msg) => {
    const order = msg.data
    attempts.set(order.orderId, (attempts.get(order.orderId) ?? 0) + 1)

    // 1. Was this order charged already? `found` is its own field because null
    //    is a value you can store.
    const marker = await queen.kv.get(NS, markerKey(order.orderId))
    if (marker.found) {
      // Nothing to do, but the message still has to leave this group's
      // cursor, or it comes back forever.
      await queen.ack(msg, 'completed', { group: msg.consumerGroup })
      return { ran: false }
    }

    // 2. Everything that can fail without touching the card goes here, before
    //    the charge: validation, lookups. This is where ORD-3 fails, once.
    if (order.orderId === CRASHING_ORDER && attempts.get(order.orderId) === 1) {
      throw new Error('card network timed out')
    }

    // 3. The external effect.
    const chargeId = chargeCard(order)

    // 4. The marker and the ack, in ONE transaction.
    //
    //    required: true turns the putIfAbsent into a gate. If the marker already
    //    exists (another delivery of this order got here first), the whole
    //    transaction rolls back, ack included, and only the winner's ack lands.
    //
    //    The ack carries this delivery's lease. If the lease ran out while the
    //    card was being charged, the ack is refused and the marker with it, and
    //    commit() throws with reason 'rejected_ack'. A compare-and-set on the
    //    marker could not do that: a version that still matches succeeds for a
    //    worker that no longer owns the message.
    const res = await queen
      .transaction()
      .kv.putIfAbsent(NS, markerKey(order.orderId), { chargeId, cents: order.cents }, { ttl: '1h', required: true })
      .ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
      .commit()

    // A lost gate comes back as a value (HTTP 200, success: false, reason
    // 'kv_precondition'). It is the normal outcome of a duplicate delivery, so
    // it is handled here and kept out of the error path, where a retry would
    // be the reflex.
    if (res.success === false && res.reason === 'kv_precondition') {
      await queen.ack(msg, 'completed', { group: msg.consumerGroup })
      return { ran: false }
    }
    return { ran: true }
  }

  // One order per pop (batch(1)). A failed ack gives up the partition's lease,
  // so in a batch of several orders every order after a failed one would be
  // charged, refused at its commit, and charged again when it came back.
  const charger = (group, limit) => queen
    .queue(ORDERS)
    .group(group)
    .subscriptionMode('all')
    .autoAck(false) // the ack rides the transaction
    .batch(1)
    .each()
    .limit(limit)
    .idleMillis(PHASE_MS)

  console.log('\ncharging')
  await charger('charger', CHARGE_DELIVERIES).consume(async (msg) => {
    try {
      const { ran } = await handle(msg)
      observed.push({ group: 'charger', orderId: msg.data.orderId, ran })
      console.log(`  ${msg.data.orderId}: ${ran ? 'charged' : 'already charged, skipped'}`)
    } catch (err) {
      // autoAck is off, so the failure is reported explicitly. It spends one
      // retry and brings the order back.
      console.log(`  ${msg.data.orderId}: ${err.message} (will be redelivered)`)
      await queen.ack(msg, 'failed', { group: msg.consumerGroup, error: err.message })

      // The moment to check the claim: right after a failure before the commit.
      const marker = await queen.kv.get(NS, markerKey(msg.data.orderId))
      assert(!marker.found, `${msg.data.orderId} failed before its commit and left no marker behind`)
    }
  })
  assert(
    observed.length === ORDER_IDS.length,
    `the charger decided every order (${ORDER_IDS.length}, got ${observed.length})`
  )

  // ------------------------------------------------------------------ replay
  //
  // A second group reads the same orders from the beginning: the same
  // messages, the same handler, and only the markers stand between them and a
  // second charge.
  console.log('\nreplaying')
  await charger('replay', ORDER_IDS.length).consume(async (msg) => {
    const { ran } = await handle(msg)
    observed.push({ group: 'replay', orderId: msg.data.orderId, ran })
    console.log(`  ${msg.data.orderId}: ran === ${ran}`)
  })

  // ---------------------------------------------------------------- checking
  console.log('\nchecking')
  const perOrder = new Map()
  for (const row of ledger) perOrder.set(row.orderId, (perOrder.get(row.orderId) ?? 0) + 1)
  assert(
    ledger.length === ORDER_IDS.length && ORDER_IDS.every(id => perOrder.get(id) === 1),
    `${ORDER_IDS.length} orders, ${ledger.length} charges, one per order`
  )
  assert(attempts.get(CRASHING_ORDER) >= 2, `${CRASHING_ORDER} came back after its failure and was charged then`)

  const replayed = observed.filter(o => o.group === 'replay')
  assert(
    replayed.length === ORDER_IDS.length && replayed.every(o => !o.ran),
    `the replay received all ${ORDER_IDS.length} orders and charged none of them`
  )

  // The markers are readable state. Each one names the charge it stands for, so
  // "was this order billed, and by which charge" has an answer in the broker.
  const markers = await queen.kv.getMany(NS, ORDER_IDS.map(markerKey))
  assert(
    markers.rows.length === ORDER_IDS.length && markers.missing.length === 0,
    `every order has its marker (${markers.rows.length})`
  )
  assert(
    markers.rows.every(r => ledger.some(l => l.chargeId === r.value.chargeId)),
    'each marker names the charge that was made'
  )
  console.log(`\n  ledger: ${ledger.map(l => `${l.orderId}=${l.chargeId}`).join(', ')}`)

  console.log(`\nPASS: ${checks} checks`)
} catch (err) {
  console.error(`\nFAIL: ${err.message}`)
  process.exitCode = 1
} finally {
  // Clean up in every case, a failed run included. The markers are KV entries
  // and deleting the queue does not remove them; they would expire after their
  // hour, but a program should not leave state behind on a shared broker.
  try {
    for (const orderId of ORDER_IDS) await queen.kv.delete(NS, markerKey(orderId))
    await queen.queue(ORDERS).delete()
  } catch (err) {
    console.error(`  (cleanup incomplete: ${err.message})`)
  }
  await queen.close()
}
examples/apps/py/exactly_once.pypython
#
# Charging each order exactly once, when orders can be delivered twice.
#
# Any consumer that charges cards sees some orders twice: a lease runs out
# while the card network is slow, a node restarts, someone replays the queue.
# The usual guard is an "already charged" flag in a database next to the
# broker, written after the charge and before the ack. The flag and the ack
# then commit separately, and in the gap between them the work can happen
# twice.
#
# Here the flag is a KV entry in the broker itself, written in the same
# transaction as the ack. The order is marked and acknowledged together, or
# neither happens, and a redelivery finds the marker and charges nothing.
#
#   orders
#     |-- group "charger"  marker + ack in ONE transaction
#     `-- group "replay"   reads every order again and must charge nothing
#
# Run it:
#   QUEEN_URL=http://localhost:6632 python3 exactly_once.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")

# A fresh queue and a fresh KV namespace per run. The namespace needs it more
# than the queue does: markers outlive the queue that produced them (they
# expire with their TTL), so a second run in the same namespace would find
# every order already charged and pass without charging anything.
RUN = f"{int(time.time() * 1000):x}"
ORDERS = f"app-py-exactly-once-{RUN}"
NS = f"app-py-exactly-once-{RUN}"
GROUP = "charger"
REPLAY_GROUP = "replay"

# Five orders. The first attempt at ORD-3 fails before it charges anything, the
# failure that must leave no trace at all.
ORDER_IDS = ["ORD-1", "ORD-2", "ORD-3", "ORD-4", "ORD-5"]
CRASHING_ORDER = "ORD-3"

# The charging phase ends after six deliveries, five orders plus the retry of
# ORD-3, and the deadline is there so that a stall fails the run instead of
# hanging it.
CHARGE_DELIVERIES = len(ORDER_IDS) + 1
PHASE_MS = 30000

CHECKS = 0

# The external effect. Each entry is a charge on a real card, and the ledger is
# what the customers' statements would show.
LEDGER: list = []
ATTEMPTS: dict = {}


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 charge_card(order: dict) -> str:
    charge_id = f"ch_{order['orderId']}_{len(LEDGER)}"
    LEDGER.append({"orderId": order["orderId"], "chargeId": charge_id, "cents": order["cents"]})
    return charge_id


def marker_key(order_id: str) -> str:
    return f"charge:{order_id}"


async def main() -> int:
    queen = Queen(url=QUEEN_URL)
    verdict, failed = "", False

    # One order per pop (batch(1)). A failed ack gives up the partition's lease,
    # so in a batch of several orders every order after a failed one would be
    # charged, refused at its commit, and charged again when it came back.
    def charger(group: str, limit: int):
        return (
            queen.queue(ORDERS)
            .group(group)
            .subscription_mode("all")
            # The ack rides the transaction.
            .auto_ack(False)
            .batch(1)
            .each()
            .limit(limit)
            .idle_millis(PHASE_MS)
        )

    try:
        print(f"broker {QUEEN_URL}")

        await queen.queue(ORDERS).config({"lease_time": 30, "retry_limit": 5}).create()

        print("\nqueuing orders")
        for index, order_id in enumerate(ORDER_IDS):
            await queen.queue(ORDERS).push(
                {"transactionId": f"order-{order_id}", "data": {"orderId": order_id, "cents": 1000 + index}}
            )
        print(f"  {len(ORDER_IDS)} orders queued")

        # ----------------------------------------------------------- charging
        #
        # handle() is the whole pattern: four steps, in this order. It returns
        # True when this delivery charged the card and False when it found the
        # order already charged. Every redelivery, whatever caused it, has to
        # come back False.
        observed: list = []

        async def handle(msg) -> bool:
            order = msg["data"]
            order_id = order["orderId"]
            ATTEMPTS[order_id] = ATTEMPTS.get(order_id, 0) + 1
            group = msg.get("consumerGroup")

            # 1. Was this order charged already? "found" is its own key
            #    because None is a value you can store.
            marker = await queen.kv.get(NS, marker_key(order_id))
            if marker["found"]:
                # Nothing to do, but the message still has to leave this
                # group's cursor, or it comes back forever.
                await queen.ack(msg, "completed", {"group": group})
                return False

            # 2. Everything that can fail without touching the card goes here,
            #    before the charge: validation, lookups. This is where ORD-3
            #    fails, once.
            if order_id == CRASHING_ORDER and ATTEMPTS[order_id] == 1:
                raise RuntimeError("card network timed out")

            # 3. The external effect.
            charge_id = charge_card(order)

            # 4. The marker and the ack, in ONE transaction.
            #
            #    required=True turns the put_if_absent into a gate. If the
            #    marker already exists (another delivery of this order got here
            #    first), the whole transaction rolls back, ack included, and
            #    only the winner's ack lands. tx.once(...) is the same gate
            #    under a shorter name; it is spelled out here so that required
            #    is visible.
            #
            #    The ack carries this delivery's lease. If the lease ran out
            #    while the card was being charged, the ack is refused and the
            #    marker with it: the broker answers with reason rejected_ack
            #    and commit() raises. A compare-and-set on the marker could not
            #    do that: a version that still matches succeeds for a worker
            #    that no longer owns the message.
            res = await (
                queen.transaction()
                .kv.put_if_absent(
                    NS,
                    marker_key(order_id),
                    {"chargeId": charge_id, "cents": order["cents"]},
                    # The Python client takes a timedelta where the JavaScript
                    # client takes "1h"; both become ttlSeconds on the wire.
                    ttl=timedelta(hours=1),
                    required=True,
                )
                .ack(msg, "completed", {"consumer_group": group})
                .commit()
            )

            # A lost gate comes back as a value (HTTP 200, success False,
            # reason "kv_precondition"). It is the normal outcome of a
            # duplicate delivery, so it is handled here and kept out of the
            # error path, where a retry would be the reflex.
            if res.get("success") is False and res.get("reason") == "kv_precondition":
                await queen.ack(msg, "completed", {"group": group})
                return False

            return True

        print("\ncharging")

        async def charge(msg) -> None:
            group = msg.get("consumerGroup")
            try:
                ran = await handle(msg)
                observed.append({"group": GROUP, "orderId": msg["data"]["orderId"], "ran": ran})
                print(f"  {msg['data']['orderId']}: {'charged' if ran else 'already charged, skipped'}")
            except RuntimeError as err:
                # auto_ack is off, so the failure is reported explicitly. It
                # spends one retry and brings the order back.
                print(f"  {msg['data']['orderId']}: {err} (will be redelivered)")
                await queen.ack(msg, "failed", {"group": group, "error": str(err)})

                # The moment to check the claim: right after a failure before
                # the commit.
                marker = await queen.kv.get(NS, marker_key(msg["data"]["orderId"]))
                check(
                    not marker["found"],
                    f"{msg['data']['orderId']} failed before its commit and left no marker behind",
                )

        await charger(GROUP, CHARGE_DELIVERIES).consume(charge)

        charged = [o for o in observed if o["group"] == GROUP]
        check(
            len(charged) == len(ORDER_IDS),
            f"the charger decided every order ({len(ORDER_IDS)}, got {len(charged)})",
        )

        # ------------------------------------------------------------ replay
        #
        # A second group reads the same orders from the beginning: the same
        # messages, the same handler, and only the markers stand between them
        # and a second charge.
        print("\nreplaying")

        async def replay(msg) -> None:
            ran = await handle(msg)
            observed.append({"group": REPLAY_GROUP, "orderId": msg["data"]["orderId"], "ran": ran})
            print(f"  {msg['data']['orderId']}: ran is {ran}")

        await charger(REPLAY_GROUP, len(ORDER_IDS)).consume(replay)

        # ----------------------------------------------------------- checking
        print("\nchecking")

        check(
            len(LEDGER) == len(ORDER_IDS),
            f"{len(ORDER_IDS)} orders, {len(LEDGER)} charges",
        )
        per_order = {order_id: sum(1 for row in LEDGER if row["orderId"] == order_id) for order_id in ORDER_IDS}
        check(
            all(count == 1 for count in per_order.values()),
            "one charge per order, none twice and none missing",
        )
        check(
            ATTEMPTS.get(CRASHING_ORDER, 0) >= 2,
            f"{CRASHING_ORDER} came back after its failure and was charged then",
        )

        replayed = [o for o in observed if o["group"] == REPLAY_GROUP]
        check(
            len(replayed) == len(ORDER_IDS),
            f"the replay received all {len(ORDER_IDS)} orders (got {len(replayed)})",
        )
        check(
            all(o["ran"] is False for o in replayed),
            "the replay charged none of them",
        )
        check(
            len(LEDGER) == len(ORDER_IDS),
            f"the ledger still has {len(ORDER_IDS)} charges after the replay",
        )

        # The markers are readable state. Each one names the charge it stands
        # for, so "was this order billed, and by which charge" has an answer in
        # the broker.
        markers = await queen.kv.get_many(NS, [marker_key(o) for o in ORDER_IDS])
        check(
            len(markers["rows"]) == len(ORDER_IDS),
            f"every order has its marker ({len(markers['rows'])})",
        )
        check(len(markers["missing"]) == 0, "no order is missing its marker")
        charge_ids = {row["chargeId"] for row in LEDGER}
        check(
            all(row["value"]["chargeId"] in charge_ids for row in markers["rows"]),
            "each marker names the charge that was made",
        )

        print("\n  ledger: " + ", ".join(f"{row['orderId']}={row['chargeId']}" for row in LEDGER))

        verdict = f"\nPASS: {CHECKS} checks"
    except Exception as err:
        verdict, failed = f"\nFAIL: {err}", True
    finally:
        # Clean up in every case, a failed run included. The markers are KV
        # entries and deleting the queue does not remove them; they would
        # expire after their hour, but a program should not leave state behind
        # on a shared broker. Best effort: a cleanup that raised would replace
        # the real verdict with its own.
        try:
            for order_id in ORDER_IDS:
                await queen.kv.delete(NS, marker_key(order_id))
            await queen.queue(ORDERS).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()))
examples/apps/http/exactly-once.shbash
#!/usr/bin/env bash
#
# Charging each order exactly once, when orders can be delivered twice, with
# nothing but curl.
#
# Any consumer that charges cards sees some orders twice: a lease runs out
# while the card network is slow, a node restarts, someone replays the queue.
# The usual guard is an "already charged" flag in a database next to the
# broker, written after the charge and before the ack. The flag and the ack
# then commit separately, and in the gap between them the work can happen
# twice.
#
# Here the flag is a KV entry in the broker itself, written in the same
# transaction as the ack. The order is marked and acknowledged together, or
# neither happens, and a redelivery finds the marker and charges nothing.
#
#   orders
#     ├── group "charger"  marker + ack in ONE transaction
#     └── group "replay"   reads every order again and must charge nothing
#
# There is no client library here and none is needed, and this file is worth
# reading even if you use one: an SDK's `kv.putIfAbsent(...)` inside a
# transaction is the `kv` array below, and the array is a key of the ROOT of the
# request beside `operations`, never an element of it. Everything an SDK hides
# is written out.
#
# Run it:
#   QUEEN_URL=http://localhost:6632 bash exactly-once.sh

set -euo pipefail

QUEEN_URL="${QUEEN_URL:-http://localhost:6632}"

# A fresh queue and a fresh KV namespace per run, so runs never share state. The
# namespace needs it more than the queue does: markers outlive the queue that
# produced them (they expire with their TTL), so a second run in the same
# namespace would find every order already charged and pass without charging
# anything. $$ is the process id, which keeps two runs in the same second apart.
RUN="$(date +%s)-$$"
ORDERS="app-http-exactly-once-$RUN"
NS="app-http-exactly-once-$RUN"

# The consumer groups. A group's cursor lives on the queue, and the queue name is
# already unique per run, so these need no suffix.
CHARGER=app-http-charger
REPLAY=app-http-replay

# Five orders. The first attempt at ORD-3 fails before it charges anything and
# before it commits: the failure that must leave no trace at all.
ORDER_IDS="ORD-1 ORD-2 ORD-3 ORD-4 ORD-5"
ORDER_COUNT=5
CRASHING_ORDER=ORD-3

# Every pop long-polls for this many milliseconds and no longer.
POLL_MS=1000

# The bound that keeps a stall from becoming a hang. A phase that has not made
# every decision by then stops, and the count check that follows reports what was
# missing. Never wait for silence; wait for a total, with a deadline.
PHASE_MS=30000

command -v jq >/dev/null 2>&1 || { echo "FAIL: jq is not installed"; exit 1; }

CHECKS=0
TMP="$(mktemp -d)"

# The external effect. Each line is a charge on a real card, and the ledger is
# what the customers' statements would show.
LEDGER="$TMP/ledger"
: > "$LEDGER"

# One line per delivery handled: "<group> <order> <ran>". `ran` is true when
# this delivery charged the card and false when it found the order already
# charged. Every redelivery, whatever caused it, has to come back false.
OBSERVED="$TMP/observed"
: > "$OBSERVED"

# One line per delivery, so the redelivery of the failed order is counted.
ATTEMPTS="$TMP/attempts"
: > "$ATTEMPTS"

# 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, and there are two
# things to remove: the queue, and the markers. The markers are KV entries, and
# deleting the queue does not remove them. They would expire after their hour
# (the ttlSeconds below), but a program should not leave state behind on a
# shared broker. 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 keys body
  keys="$(printf '%s\n' $ORDER_IDS | jq -R 'sub("^"; "charge:")' | 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/$ORDERS" || true
}
trap cleanup EXIT

fail() { FAILURE="$*"; exit 1; }

# check <actual> <expected> <description>
check() {
  [ "$1" = "$2" ] || fail "$3 (expected [$2], got [$1])"
  CHECKS=$((CHECKS + 1))
  echo "  ok: $3"
}

# ok <description>: records a check whose condition was already tested, for the
# assertions that are not an equality.
ok() {
  CHECKS=$((CHECKS + 1))
  echo "  ok: $1"
}

# 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
}

# ---------------------------------------------------------------------------
# The KV surface, in the two shapes this program needs.
#
# Everything goes through POST /api/v1/kv, including the single-key read. The
# path routes (GET|PUT|DELETE /api/v1/kv/:ns/*key) exist and are correct, but the
# batch route keeps the key out of the access log, the proxy's samples and any
# tracing span, which is the same reason a prefix may not travel in a query
# string. A marker names a customer's order; it is not URL material.
# ---------------------------------------------------------------------------

# 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.
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"
}

# marker_key <order>: the key of the entry that says this order has been
# charged.
marker_key() { printf 'charge:%s' "$1"; }

echo "broker $QUEEN_URL"

# Every broker serves /api/v1/kv: there is no flag that turns it 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)"

# ---------------------------------------------------------------------------
# A crashed worker's messages come back when its lease expires, and retryLimit
# bounds how often a failing message (one a worker acks as `failed`) is retried
# before it goes to the dead-letter queue. A lease that runs out spends none of
# that budget.
#
# /configure merges: an option not named here keeps whatever the queue already
# has. This queue name is unique per run, so what is not named lands on its
# default.
# ---------------------------------------------------------------------------
configure_body="$(jq -n --arg queue "$ORDERS" \
  '{queue: $queue, options: {leaseTime: 30, retryLimit: 5}}')"
request POST /api/v1/configure "$configure_body"
[ "$STATUS" = 200 ] || fail "configure returned HTTP $STATUS"
check "$(jq -r .configured "$OUT")" true 'the queue was created with a 30 second lease'

# ---------------------------------------------------------------------- queuing
echo
echo "queuing orders"
cents=1000
for order in $ORDER_IDS; do
  body="$(jq -n --arg queue "$ORDERS" --arg order "$order" --argjson cents "$cents" \
    '{items: [{queue: $queue, transactionId: ("order-" + $order),
               payload: {orderId: $order, cents: $cents}}]}')"
  request POST /api/v1/push "$body"
  [ "$STATUS" = 201 ] || fail "push of $order 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 $order came back $(jq -r '.[0].status' "$OUT")"
  cents=$((cents + 1))
done
echo "  $ORDER_COUNT orders queued"

# ---------------------------------------------------------------------------
# handle <group>: one delivery, from the pop response in $TMP/pop.
#
# The whole pattern: four steps, in this order.
# ---------------------------------------------------------------------------
handle() {
  local group="$1" order cents txn partition lease marker charge_id body ack_body

  order="$(jq -r '.messages[0].data.orderId' "$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")"

  printf '%s\n' "$order" >> "$ATTEMPTS"

  # 1. Was this order charged already?
  marker="$(kv_get "$(marker_key "$order")")"
  if [ "$(printf '%s' "$marker" | jq -r '.found')" = true ]; then
    # Nothing to do, but the message still has to leave this group's cursor,
    # or it comes back forever.
    ack_body="$(jq -cn --arg txn "$txn" --arg pid "$partition" --arg grp "$group" --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 false\n' "$group" "$order" >> "$OBSERVED"
    echo "  $order: already charged, skipped"
    return 0
  fi

  # 2. Everything that can fail without touching the card goes here, before
  #    the charge and before the commit: validation, lookups. This is where
  #    ORD-3 fails, once.
  if [ "$order" = "$CRASHING_ORDER" ] \
     && [ "$(grep -c "^$CRASHING_ORDER$" "$ATTEMPTS")" -eq 1 ]; then
    # The failure is reported explicitly, with a `failed` ack. It spends one
    # retry and leaves the cursor below this message, which is what brings the
    # order back.
    ack_body="$(jq -cn --arg txn "$txn" --arg pid "$partition" --arg grp "$group" --arg lease "$lease" \
      '{transactionId: $txn, partitionId: $pid, consumerGroup: $grp, leaseId: $lease,
        status: "failed", error: "card network timed out"}')"
    request POST /api/v1/ack "$ack_body"
    [ "$STATUS" = 200 ] || fail "the failing ack returned HTTP $STATUS"
    echo "  $order: card network timed out (will be redelivered)"

    # The moment to check the claim: right after a failure before the commit.
    marker="$(kv_get "$(marker_key "$order")")"
    check "$(printf '%s' "$marker" | jq -r '.found')" false \
      "$order failed before its commit and left no marker behind"
    return 0
  fi

  # 3. The external effect.
  charge_id="ch_${order}_$(wc -l < "$LEDGER" | tr -d ' ')"
  printf '%s %s %s\n' "$order" "$charge_id" "$cents" >> "$LEDGER"

  # 4. The marker and the ack, in ONE transaction.
  #
  #    `kv` is a key of the ROOT of this body, beside `operations`. The two
  #    arrays are separate top-level fields so that no client can send them
  #    under one key by accident.
  #
  #    required: true turns the putIfAbsent into a gate. If the marker already
  #    exists (another delivery of this order got here first), the whole
  #    transaction rolls back, ack included, and only the winner's ack lands.
  #    Without it a lost race would come back applied:false and the ack would
  #    still commit.
  #
  #    ttlSeconds is mandatory on every KV write, so a marker nobody deletes
  #    still goes away.
  #
  #    The ack carries this delivery's lease. If the lease ran out while the
  #    card was being charged, the ack is refused and the marker with it: the
  #    answer is HTTP 200 with success:false and reason "rejected_ack", and
  #    nothing was written. A compare-and-set on the marker could not do that:
  #    a version that still matches succeeds for a worker that no longer owns
  #    the message.
  body="$(jq -cn --arg ns "$NS" --arg key "$(marker_key "$order")" \
    --arg charge "$charge_id" --argjson cents "$cents" \
    --arg txn "$txn" --arg pid "$partition" --arg grp "$group" --arg lease "$lease" '
    {operations: [{type: "ack", transactionId: $txn, partitionId: $pid,
                   consumerGroup: $grp, leaseId: $lease, status: "completed"}],
     kv: [{op: "putIfAbsent", ns: $ns, key: $key,
           value: {chargeId: $charge, cents: $cents},
           ttlSeconds: 3600, required: true}]}')"
  request POST /api/v1/transaction "$body"

  # A lost gate comes back as a value: HTTP 200 with success:false and reason
  # "kv_precondition". It is the normal outcome of a duplicate delivery, so it is
  # handled here and kept out of the error path, where a retry would be the
  # reflex.
  [ "$STATUS" = 200 ] || fail "the commit for $order returned HTTP $STATUS: $(cat "$OUT")"
  if [ "$(jq -r '.success' "$OUT")" != true ]; then
    [ "$(jq -r '.reason' "$OUT")" = kv_precondition ] \
      || fail "the commit for $order failed: $(jq -r '.error' "$OUT")"
    # Nothing was written, ack included, and the message still has to leave
    # this group's cursor, so it is acknowledged on its own.
    ack_body="$(jq -cn --arg txn "$txn" --arg pid "$partition" --arg grp "$group" --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 false\n' "$group" "$order" >> "$OBSERVED"
    echo "  $order: charged by another delivery first, skipped"
    return 0
  fi

  printf '%s %s true\n' "$group" "$order" >> "$OBSERVED"
  echo "  $order: charged"
}

# ---------------------------------------------------------------------------
# drain <group> <decisions>: pop and handle until this group has made that many
# decisions, or the phase deadline passes. The count is the bound and the
# deadline is the net; neither is a wait for silence.
# ---------------------------------------------------------------------------
drain() {
  local group="$1" wanted="$2" deadline
  deadline=$(( $(now_ms) + PHASE_MS ))

  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 takes one order per pop. A `failed` ack gives up the
    # partition's lease, so in a batch of several orders every order after a
    # failed one would be charged, refused at its commit, and charged again when
    # it came back.
    request GET "/api/v1/pop/queue/$ORDERS?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"
    handle "$group"
  done
}

# --------------------------------------------------------------------- charging
echo
echo "charging"
drain "$CHARGER" "$ORDER_COUNT"

check "$(grep -c "^$CHARGER " "$OBSERVED" || true)" "$ORDER_COUNT" \
  'the charger decided every order'

# ---------------------------------------------------------------------- replay
#
# A second consumer group reads the same orders from the beginning: the same
# messages, the same handler, and only the markers stand between them and a
# second charge.
echo
echo "replaying"
drain "$REPLAY" "$ORDER_COUNT"

# --------------------------------------------------------------------- checking
echo
echo "checking"

check "$(wc -l < "$LEDGER" | tr -d ' ')" "$ORDER_COUNT" \
  "$ORDER_COUNT orders, $ORDER_COUNT charges"
check "$(awk '{print $1}' "$LEDGER" | sort -u | wc -l | tr -d ' ')" "$ORDER_COUNT" \
  'one charge per order: none charged twice and none skipped'

[ "$(grep -c "^$CRASHING_ORDER$" "$ATTEMPTS")" -ge 2 ] \
  || fail "$CRASHING_ORDER was never redelivered after it failed"
ok "$CRASHING_ORDER came back after its failure and was charged then"

check "$(grep -c "^$REPLAY " "$OBSERVED" || true)" "$ORDER_COUNT" \
  "the replay received all $ORDER_COUNT orders"
check "$(awk -v g="$REPLAY" '$1 == g && $3 == "false"' "$OBSERVED" | wc -l | tr -d ' ')" \
  "$ORDER_COUNT" 'every order on the second pass reported that it did not run'
check "$(wc -l < "$LEDGER" | tr -d ' ')" "$ORDER_COUNT" \
  "the replay charged nothing: the ledger still has $ORDER_COUNT charges"

# The markers are readable state. Each one names the charge it stands for, so
# "was this order billed, and by which charge" 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' $ORDER_IDS | jq -R 'sub("^"; "charge:")' | jq -sc .)"
many_body="$(jq -cn --arg ns "$NS" --argjson keys "$keys" \
  '{operations: [{op: "getMany", ns: $ns, keys: $keys}]}')"
request POST /api/v1/kv "$many_body"
[ "$STATUS" = 200 ] || fail "kv getMany returned HTTP $STATUS"
check "$(jq -r '.results[0].rows | length' "$OUT")" "$ORDER_COUNT" \
  'every order has its marker'
check "$(jq -r '.results[0].missing | length' "$OUT")" 0 \
  'no order is missing its marker'
check "$(jq -r '[.results[0].rows[].value.chargeId] | sort | join(",")' "$OUT")" \
  "$(awk '{print $2}' "$LEDGER" | sort | paste -sd, -)" \
  'each marker names the charge that was made'

echo
echo "  ledger: $(awk '{printf "%s=%s ", $1, $2}' "$LEDGER")"

echo
echo "PASS: $CHECKS checks"

Limits

The charge itself can still happen twice. Three windows remain, all between the provider’s answer and the commit: the worker dies, its lease expires during the call (the commit is then refused and the next worker charges again), or two deliveries of one order are in flight at once. The marker narrows a double charge to those windows and cannot close them, because the provider is outside the commit. Keep leaseTime well above the provider’s timeout, and pass the order id as the provider’s idempotency key so that a second charge for the same order is refused on their side.

Do the work that can fail before the charge. Validation, lookups and anything else that can throw should run first, so a failure leaves no external effect behind.

A marker lives as long as its TTL, and an expired marker reads as absent, so the TTL must outlast any redelivery or replay you will run. Markers are KV rows, and deleting the queue does not delete them.

The fence covers one pop’s lease. The SDKs send the lease beside the acks, and the broker attaches it only when the transaction names a single lease, so a transaction that acks messages from two different pops is not fenced. Commit one transaction per pop, as the code above does.

Delivery itself stays at-least-once: exactly-once holds for what the transaction writes. The transaction and the dedup are part of the Jepsen runs; see guarantees for the conditions.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close