---
title: "Charge each order once"
description: "A card charger that survives redelivery: the charge marker and the ack commit in one transaction, so five orders, a failure and a full replay still make five charges."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# Charge each order once

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. This program makes the second delivery
harmless. The "already charged" marker lives in the broker, next to the queue, and it is written
in the same transaction as the ack, so an order is either marked and acknowledged or neither. It
charges five orders, fails one of them before it reaches the card, then replays all five through
a second consumer group, and checks that the ledger holds exactly five charges.

## Run it

```bash
docker run --platform linux/amd64 -d --name queen -p 6632:6632 ghcr.io/queen-mq/queen:latest
npm install queen-mq
node exactly-once.mjs     # the JavaScript tab below, saved as exactly-once.mjs
```

What it printed against a 2.0.0-beta.6 node:

```text
queuing orders
  5 orders queued

charging
  ORD-1: charged
  ORD-2: charged
  ORD-3: card network timed out (will be redelivered)
  ok: ORD-3 failed before its commit and left no marker behind
  ORD-3: charged
  ORD-4: charged
  ORD-5: charged
  ok: the charger decided every order (5, got 5)

replaying
  ORD-1: ran === false
  ORD-2: ran === false
  ORD-3: ran === false
  ORD-4: ran === false
  ORD-5: ran === false

checking
  ok: 5 orders, 5 charges, one per order
  ok: ORD-3 came back after its failure and was charged then
  ok: the replay received all 5 orders and charged none of them
  ok: every order has its marker (5)
  ok: each marker names the charge that was made

  ledger: ORD-1=ch_ORD-1_0, ORD-2=ch_ORD-2_1, ORD-3=ch_ORD-3_2, ORD-4=ch_ORD-4_3, ORD-5=ch_ORD-5_4

PASS: 7 checks
```

The replay is the interesting half. It is the same messages through the same handler, as if an
operator had rewound the queue, and the only thing between those five orders and a second charge
is the marker each one left.

## The program

### JavaScript

```js title="examples/apps/js/exactly-once.mjs"
//
// 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()
}
```

### Python

```python title="examples/apps/py/exactly_once.py"
#
# 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()))
```

### curl

```bash title="examples/apps/http/exactly-once.sh"
#!/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"
```

## How it works

The handler is four steps, and their order is the design. It reads the marker first, and if the
order was charged already it only acknowledges the message. Then it does everything that can fail
without touching the card; ORD-3's scripted failure sits there, which is why it leaves no marker
and comes back to be charged normally. Then it charges. Then it writes the marker and the ack in
one `POST /api/v1/transaction`, with the marker as a `putIfAbsent` that carries
`required: true`.

**Figure.** 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.

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.

1. worker → Queen: pop orders, batch 1
2. Queen → worker: order 1042 + lease
3. worker → Queen: kv.get charge:1042
4. Queen → worker: not found
5. worker → card provider: charge the card (order id as idempotency key)
6. card provider → worker: chargeId
7. worker → Queen: transaction (putIfAbsent marker + ack)
   (Queen: lease held, no marker yet (required: true): one log entry)
8. Queen → worker: success

*the same commit, after the lease ran out*

   (Queen: rejected_ack: no marker, no ack; the order goes to the next worker)

Source: `examples/apps/js/exactly-once.mjs, clients/client-js/client-v2/builders/TransactionBuilder.js`.

That last call is where the guarantee comes from. `required` turns the `putIfAbsent` into a gate:
if two deliveries of one order race to the commit, the loser's whole transaction rolls back, its
ack included, and comes back as a value (`success: false`, `reason: 'kv_precondition'`) because a
duplicate is a normal outcome and not an error to retry. The ack carries the delivery's lease, so
a worker whose lease ran out during the charge has its ack refused and its marker with it.
[Charge a card once](/guides/exactly-once/) walks through why each of these holds, and
[transactions](/concepts/transactions/) covers the call itself.

The consumer pops one order at a time with `.batch(1)`. A nack (an ack with status `failed`)
releases the partition's lease, so in a batch of several orders, every order after a failed one
would still be handled, charged, refused at its commit because the lease is gone, and charged again
when it came back. We checked: with broker-sized batches, ORD-4 and ORD-5 were charged and then
refused, and the run failed.

## Why it is short

The usual guard against a double charge is an idempotency table in the application's database,
written after the charge and before the ack. The table and the broker then commit separately, and
the program has to reason about every crash between the two. Here the marker is a KV entry in the
same replicated log as the queue, so "charged" and "acknowledged" are one commit, and the
program needs no second store and no job to reconcile the two.

## Limits

- The card can still be charged twice in one window. If the worker dies, or loses its lease,
  between the provider's answer and the commit, the charge happened and no marker exists, so the
  redelivery charges again. Pass the order id to the payment provider as its idempotency key and
  the provider refuses the second charge. Effects that are themselves writes to Queen (a push, an
  ack, a KV write) go inside the transaction and have no such window.
- The marker lives as long as its TTL. The programs use an hour. Choose one longer than any
  redelivery or replay you will run; an expired marker reads as absent.
- The read before the charge is a KV call of its own. KV calls outside a transaction count
  against the tenant's rate on each node, 200 reads a second by default (burst 400,
  `QUEEN_KV_READ_RATE`), answered with 429 past it. KV inside the transaction is not counted.
- Delivery stays at least once. Exactly once holds for what the transaction writes. See
  [guarantees](/concepts/guarantees/) for what the Jepsen tests covered.

## Next

The [booking saga](/examples/saga/) puts a state machine and a timer in the same kind of
transaction, and [dedup](/concepts/dedup/) shows `once`, the short form of this gate.

Source: https://queenmq.com/examples/exactly-once/index.mdx
