---
title: "Deliver webhooks"
description: "Deliver webhooks in order per endpoint, with the broker's retry budget and a dead-letter queue, so an endpoint that is down slows down nobody else."
---

> 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

# Deliver webhooks

Put every delivery in the partition of the endpoint it is for, let the broker's retry budget do the
retrying, and let what never succeeds land in the dead-letter queue with its error. Deliveries to
one endpoint then arrive in the order the events happened, and an endpoint that is down backs up
its own partition while every other endpoint keeps receiving.

## The shape

```js
await queen.queue('webhooks')
  .config({ leaseTime: 30, retryLimit: 5, dlqAfterMaxRetries: true })
  .create()

// Producer: one partition per endpoint, the event id as transactionId.
await queen.queue('webhooks').partition(endpoint).push([
  { transactionId: event.id, data: { url, event } },
])

// Sender: return to ack, throw to spend one retry.
await queen.queue('webhooks')
  .group('sender')
  .subscriptionMode('all')        // a new group starts at the tail unless told otherwise
  .concurrency(20)
  .partitions(1)                  // one endpoint per pop
  .each()
  .consume(async (msg) => {
    const res = await fetch(msg.data.url, { method: 'POST', body: JSON.stringify(msg.data.event) })
    if (!res.ok) throw new Error(`${msg.data.url} answered ${res.status}`)
  })
```

The first push that names an endpoint creates its partition, so there is nothing to set up per
customer, and order holds inside it. The 20 senders spread over endpoints and never share one,
because a group leases a partition to one consumer at a time. `.partitions(1)` keeps each pop to a
single endpoint. A pop can otherwise lease several partitions at once, and after a failed delivery
the client skips the rest of what it popped. Deliveries to other endpoints that came in the same
pop would then come back only when their lease runs out, so a healthy endpoint would wait on a dead
one.

**Figure.** Deliveries in one queue, one partition per endpoint. Senders of the group sender each lease one endpoint's partition at a time. The endpoint hooks.acme.io is down and answers 503: its sender keeps failing that delivery, each failure spends one retry, and after retryLimit failures the delivery is filed in the dead-letter queue with its last error, and the endpoint's next delivery goes out. Meanwhile the senders of hooks.globex.io and hooks.initech.io keep delivering, in order.

A dead endpoint backs up its own partition and nothing else. The retry count lives in the broker, so it survives a sender dying halfway.

- hooks.acme.io: partition
- hooks.globex.io: partition
- hooks.initech.io: partition
- sender: retrying
- sender
- sender
- acme: 503
- globex: 200
- initech: 200
- dead-letter queue: after retryLimit failures
- hooks.acme.io → sender
- hooks.globex.io → sender
- hooks.initech.io → sender
- sender → acme: POST
- sender → globex: POST
- sender → initech: POST
- sender → dead-letter queue: budget spent

The retry budget belongs to the broker. With `autoAck` on, the default, a handler that returns
acknowledges the delivery and a handler that throws gives it back with the error, which spends one
unit of `retryLimit`. After `retryLimit` failures the next one files the delivery in the
dead-letter queue with its last error (so it was tried `retryLimit` + 1 times), and the endpoint's
next delivery goes out. Because the count lives in the broker, the retries survive the sender
dying halfway, which a retry loop inside the handler would not. The `transactionId` makes the
enqueue idempotent too: within the queue's dedup window (an hour by default), an application that
retries its own emit does not create a second delivery.

## The dead letters

Read them, and re-queue them once the endpoint is fixed:

```js
const dlq = await queen.queue('webhooks').dlq().limit(50).get()
for (const row of dlq.messages.reverse()) {      // newest first from the broker: re-queue oldest first
  console.log(row.partition, row.errorMessage, row.retryCount)
  await queen.admin.retryMessage(row.partitionId, row.transactionId)
}
```

A re-queued delivery is a new message at the tail of its partition, with transactionId
`dlq:<row id>`, and the dead-letter row is removed in the same transaction.

Some failures are permanent, a `410 Gone` or a deleted subscription, and retrying them only spends
time. To skip the rest of the budget, consume with `.autoAck(false)` and `.batch(1)`, and ack each
delivery yourself: `completed` on success, `failed` with the error on a transient failure, and
`queen.ack(msg, 'dlq', { group: msg.consumerGroup, error: reason })` on a permanent one, which files
the row at once whatever the queue's flags say. One delivery per pop keeps a dead-lettered delivery
from having anything behind it in the same batch. See [consuming](/concepts/consuming/) for the
statuses.

## The verified programs

Each program queues three events for three endpoints, one of which answers 500 to everything. It
checks that the two healthy endpoints received their events in order, that the dead one was tried
the full budget and never skipped, and that its three deliveries sit in the dead-letter queue with
the error attached.

### JavaScript

```js title="examples/apps/js/webhooks.mjs"
//
// A webhook sender: ordered per endpoint, retried by the broker, and
// dead-lettered with its error when it never succeeds.
//
// Deliveries to one endpoint have to arrive in the order the events happened,
// an endpoint that is down must not slow anybody else down, a failure is
// retried a bounded number of times, and what never succeeds has to end up
// somewhere a person can read it. Each endpoint gets a partition of its own,
// created by its first delivery, so a dead endpoint backs up its own partition
// and nothing else. Retries are the broker's retry budget, and an exhausted
// delivery lands in the dead-letter queue with the error attached.
//
//   webhook-deliveries (one partition per endpoint)
//     └── group "sender"  POSTs each delivery; a throw spends one retry
//           └── retryLimit spent -> dead-letter queue, with the error
//
// Run it:
//   QUEEN_URL=http://localhost:6632 node webhooks.mjs

import { Queen } from 'queen-mq'

const QUEEN_URL = process.env.QUEEN_URL || 'http://localhost:6632'
const RUN = Date.now().toString(36)
const DELIVERIES = `app-js-webhooks-${RUN}`

// Three subscribers. One of them answers 500 to everything. It is listed first,
// so its deliveries are the oldest in the queue and its partition is usually
// handed out first: a sender that let a failing endpoint hold up the others
// would fail the checks below.
const ENDPOINTS = {
  'initech.example': { healthy: false },
  'acme.example': { healthy: true },
  'globex.example': { healthy: true },
}
const EVENTS_PER_ENDPOINT = 3
const RETRY_LIMIT = 2

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

// Stands in for the HTTP POST to the subscriber. A real sender calls fetch and
// throws on any status that is not 2xx, which is what this does.
const postToEndpoint = async (endpoint, event) => {
  if (!ENDPOINTS[endpoint].healthy) throw new Error(`${endpoint} answered 500`)
  return { status: 200 }
}

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

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

  // retryLimit is the delivery budget: a delivery that fails RETRY_LIMIT + 1
  // times is filed in the dead-letter queue with its last error, because
  // dlqAfterMaxRetries is on. leaseTime is how long the broker waits for a
  // sender that took a delivery and never came back before it hands the
  // delivery to another sender.
  await queen.queue(DELIVERIES).config({
    leaseTime: 30,
    retryLimit: RETRY_LIMIT,
    dlqAfterMaxRetries: true,
  }).create()

  // ---------------------------------------------------------------- queuing
  //
  // The application emits events. Each delivery goes into the partition of the
  // endpoint it is for, which is what makes "in order per subscriber" a
  // property of the storage instead of something the sender has to arrange.
  console.log('\nqueuing deliveries')
  for (let seq = 1; seq <= EVENTS_PER_ENDPOINT; seq++) {
    for (const endpoint of Object.keys(ENDPOINTS)) {
      await queen.queue(DELIVERIES).partition(endpoint).push({
        // The event id. An application that retries its own emit does not
        // create a second delivery.
        transactionId: `${endpoint}-evt-${seq}`,
        data: { endpoint, seq, type: 'invoice.paid', invoiceId: `INV-${seq}` },
      })
    }
  }
  console.log(`  ${EVENTS_PER_ENDPOINT * Object.keys(ENDPOINTS).length} deliveries queued`)

  // ---------------------------------------------------------------- sending
  //
  // The sender pool. A handler that returns acknowledges the delivery and a
  // handler that throws gives it back with the error: the broker redelivers it
  // and counts one retry, and msg.deliveryAttempt says which attempt this is.
  // The retries live in the broker, so they survive the sender dying halfway,
  // which a retry loop inside the handler would not.
  //
  // partitions(1): every pop takes ONE endpoint. After a failed delivery the
  // client skips the rest of that pop, so deliveries to other endpoints that
  // came in the same pop would wait for their lease to run out.
  console.log('\nsending')
  const deliveredTo = new Map()
  const attempts = new Map()

  await queen
    .queue(DELIVERIES)
    .group('sender')
    .subscriptionMode('all')
    .concurrency(3)
    .partitions(1)
    .each()
    // Long polls end after a second and a sender stops after three quiet
    // seconds, so the program finishes. A service runs without these two.
    .timeoutMillis(1000)
    .idleMillis(3000)
    .consume(async (msg) => {
      const { endpoint, seq } = msg.data
      attempts.set(endpoint, (attempts.get(endpoint) ?? 0) + 1)
      try {
        await postToEndpoint(endpoint, msg.data)
      } catch (err) {
        console.log(`  ${endpoint} <- event ${seq} failed on attempt ${msg.deliveryAttempt}: ${err.message}`)
        throw err
      }
      deliveredTo.set(endpoint, [...(deliveredTo.get(endpoint) ?? []), seq])
      console.log(`  ${endpoint} <- event ${seq}`)
    })

  // ---------------------------------------------------------------- checking
  console.log('\nchecking')
  for (const [endpoint, { healthy }] of Object.entries(ENDPOINTS)) {
    if (!healthy) continue
    const seqs = deliveredTo.get(endpoint) ?? []
    assert(
      seqs.join(',') === '1,2,3',
      `${endpoint} received its ${EVENTS_PER_ENDPOINT} events in the order they happened (got ${seqs.join(',') || 'none'})`
    )
  }
  assert(!deliveredTo.has('initech.example'), 'the dead endpoint received nothing')
  assert(
    attempts.get('initech.example') === EVENTS_PER_ENDPOINT * (RETRY_LIMIT + 1),
    `each dead delivery was tried ${RETRY_LIMIT + 1} times before it was given up (${attempts.get('initech.example')} attempts)`
  )

  // The dead letters are records you can query. Each one keeps the payload, so
  // it names the endpoint and the invoice, and the last error, which is what
  // answers "why did this customer not get the webhook".
  const dlq = await queen.queue(DELIVERIES).dlq().limit(50).get()
  const dead = dlq.messages.filter(m => m.data.endpoint === 'initech.example')
  assert(dead.length === EVENTS_PER_ENDPOINT, `all ${EVENTS_PER_ENDPOINT} dead deliveries are in the dead-letter queue`)
  assert(
    dead.every(m => (m.errorMessage ?? '').includes('answered 500')),
    'each dead letter carries the error that killed it'
  )
  assert(dlq.messages.length === dead.length, 'no healthy delivery ended up in the dead-letter queue')
  console.log(`\n  dead letters: ${dead.map(m => `${m.data.endpoint}/${m.data.invoiceId}: ${m.errorMessage}`).join('; ')}`)

  await queen.queue(DELIVERIES).delete()
  console.log(`\nPASS: ${checks} checks`)
} catch (err) {
  console.error(`\nFAIL: ${err.message}`)
  process.exitCode = 1
} finally {
  await queen.close()
}
```

### Python

```python title="examples/apps/py/webhooks.py"
#
# A webhook sender: ordered per endpoint, retried by the broker, and
# dead-lettered with its error when it never succeeds.
#
# Deliveries to one endpoint have to arrive in the order the events happened,
# an endpoint that is down must not slow anybody else down, a failure is
# retried a bounded number of times, and what never succeeds has to end up
# somewhere a person can read it. Each endpoint gets a partition of its own,
# created by its first delivery, so a dead endpoint backs up its own partition
# and nothing else. Retries are the broker's retry budget, and an exhausted
# delivery lands in the dead-letter queue with the error attached.
#
#   webhook-deliveries (one partition per endpoint)
#     `-- group "sender"  POSTs each delivery; a raise spends one retry
#           `-- retry_limit spent -> dead-letter queue, with the error
#
# Run it:
#   QUEEN_URL=http://localhost:6632 python3 webhooks.py

import asyncio
import os
import sys
import time

from queen import Queen

QUEEN_URL = os.environ.get("QUEEN_URL", "http://localhost:6632")

# The name is prefixed per language and suffixed per run, so every application
# in every language can share one broker and two runs never read each other's
# messages.
RUN = f"{int(time.time() * 1000):x}"
DELIVERIES = f"app-py-webhooks-{RUN}"
GROUP = "sender"

# Three subscribers. One of them answers 500 to everything. It is listed first,
# so its deliveries are the oldest in the queue and its partition is usually
# handed out first: a sender that let a failing endpoint hold up the others
# would fail the checks below.
ENDPOINTS = {
    "initech.example": {"healthy": False},
    "acme.example": {"healthy": True},
    "globex.example": {"healthy": True},
}
EVENTS_PER_ENDPOINT = 3
RETRY_LIMIT = 2

CHECKS = 0


def check(condition: bool, description: str) -> None:
    """Record one verified fact, or abort the run.

    This raises instead of using the `assert` statement, because `python3 -O`
    removes `assert` and the checks are the whole point of the program.
    """
    global CHECKS
    if not condition:
        raise AssertionError(description)
    CHECKS += 1
    print(f"  ok: {description}")


async def post_to_endpoint(endpoint: str, event: dict) -> dict:
    """Stands in for the HTTP POST to the subscriber.

    A real sender calls httpx and raises on any status that is not 2xx, which
    is what this does.
    """
    if not ENDPOINTS[endpoint]["healthy"]:
        raise RuntimeError(f"{endpoint} answered 500")
    return {"status": 200}


async def main() -> int:
    # The whole client is async: every call below is awaited, and this is the
    # one event loop they all run on. Unlike the JavaScript client there is no
    # handleSignals switch, so SIGINT and SIGTERM are always handled for you.
    queen = Queen(url=QUEEN_URL)
    verdict, failed = "", False

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

        # retry_limit is the delivery budget: a delivery that fails
        # RETRY_LIMIT + 1 times is filed in the dead-letter queue with its last
        # error, because dlq_after_max_retries is on. lease_time is how long
        # the broker waits for a sender that took a delivery and never came
        # back before it hands the delivery to another sender. The config keys
        # are snake_case in Python and the client converts them to the
        # camelCase the broker expects.
        await queen.queue(DELIVERIES).config(
            {"lease_time": 30, "retry_limit": RETRY_LIMIT, "dlq_after_max_retries": True}
        ).create()

        # ------------------------------------------------------------ queuing
        #
        # The application emits events. Each delivery goes into the partition
        # of the endpoint it is for, which is what makes "in order per
        # subscriber" a property of the storage instead of something the
        # sender has to arrange.
        print("\nqueuing deliveries")
        for seq in range(1, EVENTS_PER_ENDPOINT + 1):
            for endpoint in ENDPOINTS:
                await queen.queue(DELIVERIES).partition(endpoint).push(
                    {
                        # The event id. An application that retries its own
                        # emit does not create a second delivery. The key is
                        # camelCase because it is the broker's wire name.
                        "transactionId": f"{endpoint}-evt-{seq}",
                        "data": {
                            "endpoint": endpoint,
                            "seq": seq,
                            "type": "invoice.paid",
                            "invoiceId": f"INV-{seq}",
                        },
                    }
                )
        print(f"  {EVENTS_PER_ENDPOINT * len(ENDPOINTS)} deliveries queued")

        # ------------------------------------------------------------ sending
        #
        # The sender pool. auto_ack is on by default, so a handler that returns
        # acknowledges the delivery and a handler that raises gives it back
        # with the error: the broker redelivers it and counts one retry, and
        # msg["deliveryAttempt"] says which attempt this is. The retries live
        # in the broker, so they survive the sender dying halfway, which a
        # retry loop inside the handler would not.
        #
        # partitions(1): every pop takes ONE endpoint. After a failed delivery
        # the client skips the rest of that pop, so deliveries to other
        # endpoints that came in the same pop would wait for their lease to
        # run out.
        print("\nsending")
        delivered_to: dict = {}
        attempts: dict = {}

        async def send(msg) -> None:
            endpoint = msg["data"]["endpoint"]
            seq = msg["data"]["seq"]
            attempts[endpoint] = attempts.get(endpoint, 0) + 1
            try:
                await post_to_endpoint(endpoint, msg["data"])
            except Exception as err:
                print(f"  {endpoint} <- event {seq} failed on attempt {msg['deliveryAttempt']}: {err}")
                raise
            delivered_to.setdefault(endpoint, []).append(seq)
            print(f"  {endpoint} <- event {seq}")

        await (
            queen.queue(DELIVERIES)
            .group(GROUP)
            # A group created after the messages were pushed starts at the tail,
            # so without this it would see nothing.
            .subscription_mode("all")
            .concurrency(3)
            .partitions(1)
            .each()
            # Long polls end after a second and a sender stops after three
            # quiet seconds, so the program finishes. A service runs without
            # these two.
            .timeout_millis(1000)
            .idle_millis(3000)
            .consume(send)
        )

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

        for endpoint, meta in ENDPOINTS.items():
            if not meta["healthy"]:
                continue
            seqs = delivered_to.get(endpoint, [])
            got = ",".join(map(str, seqs)) or "none"
            check(
                len(seqs) == EVENTS_PER_ENDPOINT,
                f"{endpoint} received its {EVENTS_PER_ENDPOINT} events (got {len(seqs)})",
            )
            check(seqs == [1, 2, 3], f"{endpoint} received them in the order they happened (got {got})")

        check("initech.example" not in delivered_to, "the dead endpoint received nothing")
        check(
            attempts.get("initech.example", 0) == EVENTS_PER_ENDPOINT * (RETRY_LIMIT + 1),
            f"each dead delivery was tried {RETRY_LIMIT + 1} times before it was given up "
            f"({attempts.get('initech.example', 0)} attempts)",
        )

        # The dead letters are records you can query. Each one keeps the
        # payload, so it names the endpoint and the invoice, and the last
        # error, which is what answers "why did this customer not get the
        # webhook". They come back as plain dicts, with the payload under
        # "data" and the error under the broker's wire name, "errorMessage".
        dlq = await queen.queue(DELIVERIES).dlq().limit(50).get()
        dead = [m for m in dlq["messages"] if m["data"]["endpoint"] == "initech.example"]

        check(
            len(dead) == EVENTS_PER_ENDPOINT,
            f"all {EVENTS_PER_ENDPOINT} dead deliveries are in the dead-letter queue",
        )
        check(
            all("answered 500" in (m.get("errorMessage") or "") for m in dead),
            "each dead letter carries the error that killed it",
        )
        check(
            len(dlq["messages"]) == len(dead),
            "no healthy delivery ended up in the dead-letter queue",
        )

        listing = "; ".join(
            f"{m['data']['endpoint']}/{m['data']['invoiceId']}: {m.get('errorMessage')}" for m in dead
        )
        print(f"\n  dead letters: {listing}")

        # Clean up on success only: a failed run leaves the queue on the broker
        # to be looked at.
        await queen.queue(DELIVERIES).delete()

        verdict = f"\nPASS: {CHECKS} checks"
    except Exception as err:
        verdict, failed = f"\nFAIL: {err}", True
    finally:
        # close() flushes the client-side buffers and closes the HTTP pool. It
        # narrates its own shutdown on stdout, which is why the verdict is
        # printed after it: PASS or FAIL stays the last line of a run.
        await queen.close()

    # A failure goes to stderr, like the rest of the set. Flush stdout first so
    # the verdict still lands last when the two are piped into one file.
    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()))
```

### Go

```go title="examples/apps/go/webhooks/main.go"
//
// A webhook sender: ordered per endpoint, retried by the broker, and
// dead-lettered with its error when it never succeeds.
//
// Deliveries to one endpoint have to arrive in the order the events happened,
// an endpoint that is down must not slow anybody else down, a failure is
// retried a bounded number of times, and what never succeeds has to end up
// somewhere a person can read it. Each endpoint gets a partition of its own,
// created by its first delivery, so a dead endpoint backs up its own partition
// and nothing else. Retries are the broker's retry budget, and an exhausted
// delivery lands in the dead-letter queue with the error attached.
//
//	webhook-deliveries (one partition per endpoint)
//	  `-- group "sender"  POSTs each delivery; an error spends one retry
//	        `-- RetryLimit spent -> dead-letter queue, with the error
//
// Run it:
//
//	QUEEN_URL=http://localhost:6632 GOWORK=off go run ./webhooks
package main

import (
	"context"
	"fmt"
	"os"
	"slices"
	"strconv"
	"strings"
	"sync"
	"time"

	queen "github.com/smartpricing/queen/clients/client-go/v2"
)

var runID = strconv.FormatInt(time.Now().UnixMilli(), 36)

var deliveriesQueue = "app-go-webhooks-" + runID

const group = "sender"

// Three subscribers. One of them answers 500 to everything. It is listed first,
// so its deliveries are the oldest in the queue and its partition is usually
// handed out first: a sender that let a failing endpoint hold up the others
// would fail the checks below. The list is a slice so the queuing order is the
// same on every run, which Go map iteration would not give.
type endpoint struct {
	host    string
	healthy bool
}

var endpoints = []endpoint{
	{host: "initech.example", healthy: false},
	{host: "acme.example", healthy: true},
	{host: "globex.example", healthy: true},
}

const (
	eventsPerEndpoint = 3
	retryLimit        = 2
	deadEndpoint      = "initech.example"
)

func isHealthy(host string) bool {
	for _, e := range endpoints {
		if e.host == host {
			return e.healthy
		}
	}
	return false
}

// postToEndpoint stands in for the HTTP POST to the subscriber. A real sender
// calls net/http and returns an error for any status that is not 2xx, which is
// what this does.
func postToEndpoint(host string, event map[string]interface{}) error {
	if !isHealthy(host) {
		return fmt.Errorf("%s answered 500", host)
	}
	return nil
}

var checks int

func assert(condition bool, description string) error {
	if !condition {
		return fmt.Errorf("%s", description)
	}
	checks++
	fmt.Printf("  ok: %s\n", description)
	return nil
}

func main() {
	if err := run(); err != nil {
		fmt.Fprintf(os.Stderr, "\nFAIL: %v\n", err)
		os.Exit(1)
	}
	fmt.Printf("\nPASS: %d checks\n", checks)
}

func run() error {
	brokerURL := os.Getenv("QUEEN_URL")
	if brokerURL == "" {
		brokerURL = "http://localhost:6632"
	}

	// One context bounds the whole program, so a broker that stops answering
	// fails the run when it expires.
	ctx, cancel := context.WithTimeout(context.Background(), 180*time.Second)
	defer cancel()

	client, err := queen.New(brokerURL)
	if err != nil {
		return fmt.Errorf("create client: %w", err)
	}
	defer client.Close(context.Background())

	fmt.Printf("broker %s\n", brokerURL)

	// RetryLimit is the delivery budget: a delivery that fails RetryLimit + 1
	// times is filed in the dead-letter queue with its last error, because
	// DlqAfterMaxRetries is on. LeaseTime is how long the broker waits for a
	// sender that took a delivery and never came back before it hands the
	// delivery to another sender.
	if _, err := client.Queue(deliveriesQueue).
		Config(queen.QueueConfig{
			LeaseTime:          30,
			RetryLimit:         retryLimit,
			DlqAfterMaxRetries: true,
		}).
		Create().Execute(ctx); err != nil {
		return fmt.Errorf("create %s: %w", deliveriesQueue, err)
	}

	// ------------------------------------------------------------------ queuing
	//
	// The application emits events. Each delivery goes into the partition of
	// the endpoint it is for, which is what makes "in order per subscriber" a
	// property of the storage instead of something the sender has to arrange.
	fmt.Println("\nqueuing deliveries")
	for seq := 1; seq <= eventsPerEndpoint; seq++ {
		for _, e := range endpoints {
			// The event id, which this client takes on the push builder. An
			// application that retries its own emit does not create a second
			// delivery.
			if _, err := client.Queue(deliveriesQueue).
				Partition(e.host).
				Push(map[string]interface{}{
					"endpoint":  e.host,
					"seq":       seq,
					"type":      "invoice.paid",
					"invoiceId": fmt.Sprintf("INV-%d", seq),
				}).
				TransactionID(fmt.Sprintf("%s-evt-%d", e.host, seq)).
				Execute(ctx); err != nil {
				return fmt.Errorf("queue %s/%d: %w", e.host, seq, err)
			}
		}
	}
	fmt.Printf("  %d deliveries queued\n", eventsPerEndpoint*len(endpoints))

	// ------------------------------------------------------------------ sending
	//
	// The sender pool. Auto-ack is the default, so a handler that returns nil
	// acknowledges the delivery and a handler that returns an error gives it
	// back with the error: the broker redelivers it and counts one retry. The
	// retries live in the broker, so they survive the sender dying halfway,
	// which a retry loop inside the handler would not.
	//
	// Partitions(1): every pop takes ONE endpoint. After a failed delivery the
	// client skips the rest of that pop, so deliveries to other endpoints that
	// came in the same pop would wait for their lease to run out.
	fmt.Println("\nsending")
	var mu sync.Mutex
	deliveredTo := map[string][]int{}
	attempts := map[string]int{}
	// The broker numbers every delivery (deliveryAttempt in the pop response),
	// but this client's Message does not expose that field, and RetryCount is
	// filled only on a dead-letter read. So the sender counts its own attempts
	// at each event, by event id, for the log line.
	tried := map[string]int{}

	err = client.Queue(deliveriesQueue).
		Group(group).
		SubscriptionMode(queen.SubscriptionModeAll).
		Concurrency(3).
		Partitions(1).
		Each().
		// Long polls end after a second and a sender stops after three quiet
		// seconds, so the program finishes. A service runs without these two.
		TimeoutMillis(1000).
		IdleMillis(3000).
		Consume(ctx, func(ctx context.Context, msg *queen.Message) error {
			host, _ := msg.Data["endpoint"].(string)
			seq, ok := msg.Data["seq"].(float64)
			if !ok {
				return fmt.Errorf("delivery %s has no numeric seq", msg.TransactionID)
			}

			// Three senders are three goroutines in this handler, so the
			// bookkeeping is behind a mutex.
			mu.Lock()
			attempts[host]++
			tried[msg.TransactionID]++
			attempt := tried[msg.TransactionID]
			mu.Unlock()

			if err := postToEndpoint(host, msg.Data); err != nil {
				fmt.Printf("  %s <- event %d failed on attempt %d: %v\n", host, int(seq), attempt, err)
				return err
			}

			mu.Lock()
			deliveredTo[host] = append(deliveredTo[host], int(seq))
			mu.Unlock()
			fmt.Printf("  %s <- event %d\n", host, int(seq))
			return nil
		}).
		Execute(ctx)
	if err != nil {
		return fmt.Errorf("sending: %w", err)
	}

	// ------------------------------------------------------------------ checking
	fmt.Println("\nchecking")

	for _, e := range endpoints {
		if !e.healthy {
			continue
		}
		seqs := deliveredTo[e.host]
		if err := assert(
			len(seqs) == eventsPerEndpoint,
			fmt.Sprintf("%s received all %d events (got %d)", e.host, eventsPerEndpoint, len(seqs)),
		); err != nil {
			return err
		}
		if err := assert(
			slices.Equal(seqs, []int{1, 2, 3}),
			fmt.Sprintf("%s received its events in the order they happened (got %s)", e.host, orNone(joinInts(seqs))),
		); err != nil {
			return err
		}
	}

	if err := assert(len(deliveredTo[deadEndpoint]) == 0, "the dead endpoint received nothing"); err != nil {
		return err
	}
	if err := assert(
		attempts[deadEndpoint] == eventsPerEndpoint*(retryLimit+1),
		fmt.Sprintf("each dead delivery was tried %d times before it was given up (%d attempts)",
			retryLimit+1, attempts[deadEndpoint]),
	); err != nil {
		return err
	}

	// The dead letters are records you can query. Each one keeps the payload,
	// so it names the endpoint and the invoice, and the last error, which is
	// what answers "why did this customer not get the webhook". DLQ takes a
	// consumer group to filter by; empty means every group on this queue.
	dlq, err := client.Queue(deliveriesQueue).DLQ("").Limit(50).Get(ctx)
	if err != nil {
		return fmt.Errorf("read dead letters: %w", err)
	}

	var dead []queen.Message
	for _, m := range dlq.Messages {
		if host, _ := m.Data["endpoint"].(string); host == deadEndpoint {
			dead = append(dead, m)
		}
	}

	if err := assert(
		len(dead) == eventsPerEndpoint,
		fmt.Sprintf("all %d dead deliveries are in the dead-letter queue", eventsPerEndpoint),
	); err != nil {
		return err
	}

	carriesError := true
	for _, m := range dead {
		if !strings.Contains(m.ErrorMessage, "answered 500") {
			carriesError = false
		}
	}
	if err := assert(carriesError, "each dead letter carries the error that killed it"); err != nil {
		return err
	}
	if err := assert(
		len(dlq.Messages) == len(dead),
		"no healthy delivery ended up in the dead-letter queue",
	); err != nil {
		return err
	}

	names := make([]string, 0, len(dead))
	for _, m := range dead {
		host, _ := m.Data["endpoint"].(string)
		invoice, _ := m.Data["invoiceId"].(string)
		names = append(names, fmt.Sprintf("%s/%s: %s", host, invoice, m.ErrorMessage))
	}
	fmt.Printf("\n  dead letters: %s\n", strings.Join(names, "; "))

	// Clean up on success only: a failed run leaves the queue and its dead
	// letters on the broker to be looked at.
	if _, err := client.Queue(deliveriesQueue).Delete().Execute(ctx); err != nil {
		return fmt.Errorf("delete %s: %w", deliveriesQueue, err)
	}

	return nil
}

func joinInts(values []int) string {
	parts := make([]string, len(values))
	for i, v := range values {
		parts[i] = strconv.Itoa(v)
	}
	return strings.Join(parts, ",")
}

func orNone(s string) string {
	if s == "" {
		return "none"
	}
	return s
}
```

### Rust

```rust title="examples/apps/rust/src/bin/webhooks.rs"
//
// A webhook sender: ordered per endpoint, retried by the broker, and
// dead-lettered with its error when it never succeeds.
//
// Deliveries to one endpoint have to arrive in the order the events happened,
// an endpoint that is down must not slow anybody else down, a failure is
// retried a bounded number of times, and what never succeeds has to end up
// somewhere a person can read it. Each endpoint gets a partition of its own,
// created by its first delivery, so a dead endpoint backs up its own partition
// and nothing else. Retries are the broker's retry budget, and an exhausted
// delivery lands in the dead-letter queue with the error attached.
//
//   webhook-deliveries (one partition per endpoint)
//     └── group "sender"  POSTs each delivery; an Err spends one retry
//           └── retry_limit spent -> dead-letter queue, with the error
//
// Run it:
//   QUEEN_URL=http://localhost:6632 cargo run --bin webhooks

use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime, UNIX_EPOCH};

use queen_mq::{Config, Message, PushItem, Queen, QueueOptions, SubscriptionMode};
use serde_json::json;

const GROUP: &str = "sender";

// Three subscribers. One of them answers 500 to everything. It is listed first,
// so its deliveries are the oldest in the queue and its partition is usually
// handed out first: a sender that let a failing endpoint hold up the others
// would fail the checks below.
//
// (endpoint, healthy)
const ENDPOINTS: [(&str, bool); 3] = [
    ("initech.example", false),
    ("acme.example", true),
    ("globex.example", true),
];
const EVENTS_PER_ENDPOINT: i64 = 3;
const RETRY_LIMIT: i32 = 2;

fn healthy(endpoint: &str) -> bool {
    ENDPOINTS
        .iter()
        .find(|(name, _)| *name == endpoint)
        .map(|(_, ok)| *ok)
        .unwrap_or(false)
}

/// Stands in for the HTTP POST to the subscriber. A real sender calls an HTTP
/// client and returns Err for any status that is not 2xx, which is what this
/// does.
async fn post_to_endpoint(endpoint: &str) -> Result<(), String> {
    if !healthy(endpoint) {
        return Err(format!("{endpoint} answered 500"));
    }
    Ok(())
}

struct Checks(usize);

impl Checks {
    fn assert(&mut self, condition: bool, description: &str) -> Result<(), String> {
        if !condition {
            return Err(description.to_string());
        }
        self.0 += 1;
        println!("  ok: {description}");
        Ok(())
    }
}

// Rust has no exceptions, so the shape the JavaScript gets from try/catch comes
// from `run` returning a Result: every `?` on the way down is a failed check or
// a failed call, and main turns it into FAIL and a non-zero exit.
#[tokio::main]
async fn main() {
    match run().await {
        Ok(checks) => println!("\nPASS: {checks} checks"),
        Err(e) => {
            eprintln!("\nFAIL: {e}");
            std::process::exit(1);
        }
    }
}

async fn run() -> Result<usize, String> {
    let url = std::env::var("QUEEN_URL").unwrap_or_else(|_| "http://localhost:6632".into());
    let run_id = SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .unwrap()
        .as_millis();
    let deliveries = format!("app-rust-webhooks-{run_id}");

    let mut checks = Checks(0);
    println!("broker {url}");

    let queen = Queen::connect(Config::new(&url)).map_err(|e| e.to_string())?;

    // retry_limit is the delivery budget: a delivery that fails RETRY_LIMIT + 1
    // times is filed in the dead-letter queue with its last error, because
    // dlq_after_max_retries is on. lease_time is how long the broker waits for
    // a sender that took a delivery and never came back before it hands the
    // delivery to another sender.
    queen
        .queue(&deliveries)
        .configure(QueueOptions {
            lease_time: Some(30),
            retry_limit: Some(RETRY_LIMIT),
            dlq_after_max_retries: Some(true),
            ..Default::default()
        })
        .await
        .map_err(|e| e.to_string())?;

    // ------------------------------------------------------------------ queuing
    //
    // The application emits events. Each delivery goes into the partition of
    // the endpoint it is for, which is what makes "in order per subscriber" a
    // property of the storage instead of something the sender has to arrange.
    println!("\nqueuing deliveries");
    for seq in 1..=EVENTS_PER_ENDPOINT {
        for (endpoint, _) in ENDPOINTS {
            // The event id. An application that retries its own emit does not
            // create a second delivery. push() would mint a UUIDv7 id of its
            // own, so the item is built by hand and sent with push_items().
            queen
                .queue(&deliveries)
                .partition(endpoint)
                .push_items(vec![PushItem::new(
                    &deliveries,
                    json!({
                        "endpoint": endpoint,
                        "seq": seq,
                        "type": "invoice.paid",
                        "invoiceId": format!("INV-{seq}"),
                    }),
                )
                .partition(endpoint)
                .transaction_id(format!("{endpoint}-evt-{seq}"))])
                .await
                .map_err(|e| e.to_string())?;
        }
    }
    println!(
        "  {} deliveries queued",
        EVENTS_PER_ENDPOINT as usize * ENDPOINTS.len()
    );

    // ------------------------------------------------------------------ sending
    //
    // The sender pool. A handler that returns Ok acknowledges the delivery and
    // a handler that returns Err gives it back with the error: the broker
    // redelivers it and counts one retry. The retries live in the broker, so
    // they survive the sender dying halfway, which a retry loop inside the
    // handler would not.
    //
    // partitions(1): every pop takes ONE endpoint. After a failed delivery the
    // client skips the rest of that pop, so deliveries to other endpoints that
    // came in the same pop would wait for their lease to run out.
    println!("\nsending");
    let delivered_to: Arc<Mutex<HashMap<String, Vec<i64>>>> = Arc::new(Mutex::new(HashMap::new()));
    let attempts: Arc<Mutex<HashMap<String, usize>>> = Arc::new(Mutex::new(HashMap::new()));
    // The broker numbers every delivery (deliveryAttempt in the pop response),
    // but this client's Message does not carry that field, so the sender counts
    // its own attempts at each event, by event id, for the log line.
    let tried: Arc<Mutex<HashMap<String, usize>>> = Arc::new(Mutex::new(HashMap::new()));

    {
        let delivered_to = Arc::clone(&delivered_to);
        let attempts = Arc::clone(&attempts);
        let tried = Arc::clone(&tried);
        queen
            .queue(&deliveries)
            .group(GROUP)
            .subscription_mode(SubscriptionMode::All)
            .concurrency(3)
            .partitions(1)
            // Every healthy delivery once and every dead one RETRY_LIMIT + 1
            // times. This client counts limit() across the three workers, so
            // the pool stops on that count; idle() is the deadline behind it,
            // and long polls end after a second so it is noticed promptly. A
            // service runs without these three.
            .limit(
                (EVENTS_PER_ENDPOINT * 2 + EVENTS_PER_ENDPOINT * (RETRY_LIMIT as i64 + 1)) as u64,
            )
            .poll_timeout(Duration::from_secs(1))
            .idle(Duration::from_secs(3))
            .consume(move |msg: Message| {
                let delivered_to = Arc::clone(&delivered_to);
                let attempts = Arc::clone(&attempts);
                let tried = Arc::clone(&tried);
                async move {
                    let endpoint = msg.data["endpoint"]
                        .as_str()
                        .unwrap_or_default()
                        .to_string();
                    let seq = msg.data["seq"].as_i64().unwrap_or(0);
                    *attempts
                        .lock()
                        .unwrap()
                        .entry(endpoint.clone())
                        .or_insert(0) += 1;
                    let attempt = {
                        let mut tried = tried.lock().unwrap();
                        let n = tried.entry(msg.transaction_id.clone()).or_insert(0);
                        *n += 1;
                        *n
                    };

                    // An Err becomes a nack carrying this string, which is what
                    // the dead letter shows once the budget runs out.
                    if let Err(e) = post_to_endpoint(&endpoint).await {
                        println!("  {endpoint} <- event {seq} failed on attempt {attempt}: {e}");
                        return Err(e);
                    }

                    delivered_to
                        .lock()
                        .unwrap()
                        .entry(endpoint.clone())
                        .or_default()
                        .push(seq);
                    println!("  {endpoint} <- event {seq}");
                    Ok::<_, String>(())
                }
            })
            .await
            .map_err(|e| e.to_string())?;
    }

    // ------------------------------------------------------------------ checking
    println!("\nchecking");

    let delivered_to = delivered_to.lock().unwrap().clone();
    let attempts = attempts.lock().unwrap().clone();

    for (endpoint, ok) in ENDPOINTS {
        if !ok {
            continue;
        }
        let seqs = delivered_to.get(endpoint).cloned().unwrap_or_default();
        let listed: Vec<String> = seqs.iter().map(i64::to_string).collect();
        let listed = if listed.is_empty() {
            "none".to_string()
        } else {
            listed.join(",")
        };
        checks.assert(
            seqs.len() == EVENTS_PER_ENDPOINT as usize,
            &format!(
                "{endpoint} received all {EVENTS_PER_ENDPOINT} events (got {})",
                seqs.len()
            ),
        )?;
        checks.assert(
            seqs == [1, 2, 3],
            &format!("{endpoint} received its events in the order they happened (got {listed})"),
        )?;
    }

    checks.assert(
        !delivered_to.contains_key("initech.example"),
        "the dead endpoint received nothing",
    )?;
    let dead_attempts = attempts.get("initech.example").copied().unwrap_or(0);
    checks.assert(
        dead_attempts == EVENTS_PER_ENDPOINT as usize * (RETRY_LIMIT as usize + 1),
        &format!(
            "each dead delivery was tried {} times before it was given up ({dead_attempts} attempts)",
            RETRY_LIMIT + 1
        ),
    )?;

    // The dead letters are records you can query. Each one keeps the payload,
    // so it names the endpoint and the invoice, and the last error, which is
    // what answers "why did this customer not get the webhook". In this client
    // the page size and the offset are the arguments of dlq(), the only two
    // filters the broker honours besides the queue and the group.
    let dlq = queen
        .queue(&deliveries)
        .dlq(Some(50), None)
        .await
        .map_err(|e| e.to_string())?;
    let dead: Vec<_> = dlq
        .messages
        .iter()
        .filter(|m| m.data["endpoint"] == json!("initech.example"))
        .collect();

    checks.assert(
        dead.len() == EVENTS_PER_ENDPOINT as usize,
        &format!("all {EVENTS_PER_ENDPOINT} dead deliveries are in the dead-letter queue"),
    )?;
    checks.assert(
        dead.iter()
            .all(|m| m.error.as_deref().unwrap_or("").contains("answered 500")),
        "each dead letter carries the error that killed it",
    )?;
    checks.assert(
        dlq.messages.len() == dead.len(),
        "no healthy delivery ended up in the dead-letter queue",
    )?;

    let listed: Vec<String> = dead
        .iter()
        .map(|m| {
            format!(
                "{}/{}: {}",
                m.data["endpoint"].as_str().unwrap_or("?"),
                m.data["invoiceId"].as_str().unwrap_or("?"),
                m.error.as_deref().unwrap_or("")
            )
        })
        .collect();
    println!("\n  dead letters: {}", listed.join("; "));

    // Clean up on success only: a failed run leaves the queue and its dead
    // letters on the broker to be looked at.
    queen
        .queue(&deliveries)
        .delete()
        .await
        .map_err(|e| e.to_string())?;

    queen.close().await.map_err(|e| e.to_string())?;

    Ok(checks.0)
}
```

### PHP

```php title="examples/apps/php/webhooks.php"
//
// A webhook sender: ordered per endpoint, retried by the broker, and
// dead-lettered with its error when it never succeeds.
//
// Deliveries to one endpoint have to arrive in the order the events happened,
// an endpoint that is down must not slow anybody else down, a failure is
// retried a bounded number of times, and what never succeeds has to end up
// somewhere a person can read it. Each endpoint gets a partition of its own,
// created by its first delivery, so a dead endpoint backs up its own partition
// and nothing else. Retries are the broker's retry budget, and an exhausted
// delivery lands in the dead-letter queue with the error attached.
//
//   webhook-deliveries (one partition per endpoint)
//     └── group "sender"  POSTs each delivery; a nack spends one retry
//           └── retryLimit spent -> dead-letter queue, with the error
//
// Run it:
//   QUEEN_URL=http://localhost:6632 php webhooks.php

require __DIR__ . '/vendor/autoload.php';

use Queen\Queen;

$QUEEN_URL = getenv('QUEEN_URL') ?: 'http://localhost:6632';
$RUN = base_convert((string) (int) (microtime(true) * 1000), 10, 36);
$DELIVERIES = "app-php-webhooks-{$RUN}";
$GROUP = 'sender';

// Three subscribers. One of them answers 500 to everything. It is listed first,
// so its deliveries are the oldest in the queue and its partition is usually
// handed out first: a sender that let a failing endpoint hold up the others
// would fail the checks below.
$ENDPOINTS = [
    'initech.example' => ['healthy' => false],
    'acme.example' => ['healthy' => true],
    'globex.example' => ['healthy' => true],
];
$EVENTS_PER_ENDPOINT = 3;
$RETRY_LIMIT = 2;

$checks = 0;
$assert = function (bool $condition, string $description) use (&$checks): void {
    if (!$condition) {
        throw new RuntimeException($description);
    }
    $checks++;
    echo "  ok: {$description}\n";
};

// Stands in for the HTTP POST to the subscriber. A real sender calls Guzzle and
// throws on any status that is not 2xx, which is what this does.
$postToEndpoint = function (string $endpoint, array $event) use ($ENDPOINTS): array {
    if (!$ENDPOINTS[$endpoint]['healthy']) {
        throw new RuntimeException("{$endpoint} answered 500");
    }
    return ['status' => 200];
};

$queen = new Queen($QUEEN_URL);
$exitCode = 0;

try {
    echo "broker {$QUEEN_URL}\n";

    // retryLimit is the delivery budget: a delivery that fails RETRY_LIMIT + 1
    // times is filed in the dead-letter queue with its last error, because
    // dlqAfterMaxRetries is on. leaseTime is how long the broker waits for a
    // sender that took a delivery and never came back before it hands the
    // delivery to another sender. config() fills in the queue defaults around
    // the keys named here.
    $queen->queue($DELIVERIES)->config([
        'leaseTime' => 30,
        'retryLimit' => $RETRY_LIMIT,
        'dlqAfterMaxRetries' => true,
    ])->create()->execute();

    // ------------------------------------------------------------------ queuing
    //
    // The application emits events. Each delivery goes into the partition of the
    // endpoint it is for, which is what makes "in order per subscriber" a
    // property of the storage instead of something the sender has to arrange.
    echo "\nqueuing deliveries\n";
    for ($seq = 1; $seq <= $EVENTS_PER_ENDPOINT; $seq++) {
        foreach (array_keys($ENDPOINTS) as $endpoint) {
            $queen->queue($DELIVERIES)->partition($endpoint)->push([[
                // The event id. An application that retries its own emit does
                // not create a second delivery.
                'transactionId' => "{$endpoint}-evt-{$seq}",
                'data' => [
                    'endpoint' => $endpoint,
                    'seq' => $seq,
                    'type' => 'invoice.paid',
                    'invoiceId' => "INV-{$seq}",
                ],
            ]])->execute();
        }
    }
    echo '  ' . ($EVENTS_PER_ENDPOINT * count($ENDPOINTS)) . " deliveries queued\n";

    // ------------------------------------------------------------------ sending
    //
    // The sender pool. concurrency(3) is three long polls in flight at once on
    // one cURL multi-handle, and the handlers run one after another in this
    // process. A failed delivery goes back to the broker with its error: the
    // broker redelivers it and counts one retry, and deliveryAttempt says which
    // attempt this is. The retries live in the broker, so they survive the
    // sender dying halfway, which a retry loop inside the handler would not.
    //
    // autoAck(false): the handler acknowledges each delivery itself, because
    // this client's automatic nack sends no error and the dead letter would
    // keep none. partitions(1) and batch(1): every pop takes ONE delivery, from
    // one endpoint. A nack releases the endpoint's lease, and the broker then
    // refuses acks for anything else popped under it, so a sender holding more
    // than the failed delivery could not acknowledge the rest.
    echo "\nsending\n";
    $deliveredTo = [];
    $attempts = [];

    $queen
        ->queue($DELIVERIES)
        ->group($GROUP)
        ->subscriptionMode('all')
        ->concurrency(3)
        ->partitions(1)
        ->batch(1)
        ->each()
        ->autoAck(false)
        // Long polls end after a second and a sender stops after three quiet
        // seconds, so the program finishes. A service runs without these two.
        ->timeoutMillis(1000)
        ->idleMillis(3000)
        ->consume(function (array $msg) use ($queen, $postToEndpoint, $GROUP, &$deliveredTo, &$attempts): void {
            $endpoint = $msg['data']['endpoint'];
            $seq = $msg['data']['seq'];
            $attempts[$endpoint] = ($attempts[$endpoint] ?? 0) + 1;

            try {
                $postToEndpoint($endpoint, $msg['data']);
            } catch (Throwable $failure) {
                echo "  {$endpoint} <- event {$seq} failed on attempt {$msg['deliveryAttempt']}: {$failure->getMessage()}\n";
                // The nack spends one retry, and its error is what the dead
                // letter keeps once the budget is spent.
                $nack = $queen->ack($msg, 'failed', ['group' => $GROUP, 'error' => $failure->getMessage()]);
                if (($nack[0]['success'] ?? false) !== true) {
                    throw new RuntimeException("the broker refused the nack for {$endpoint}/{$seq}");
                }
                return;
            }

            $deliveredTo[$endpoint][] = $seq;
            echo "  {$endpoint} <- event {$seq}\n";

            // The ack names the consumer group: without it the broker looks for
            // the lease on the queue's own cursor and refuses the ack. The
            // client's outer success flag only says the request went through,
            // and the broker's verdict is the first item of the reply.
            $ack = $queen->ack($msg, 'completed', ['group' => $GROUP]);
            if (($ack[0]['success'] ?? false) !== true) {
                throw new RuntimeException("the broker refused the ack for {$endpoint}/{$seq}");
            }
        })
        ->execute();

    // ----------------------------------------------------------------- checking
    echo "\nchecking\n";

    foreach ($ENDPOINTS as $endpoint => $meta) {
        if (!$meta['healthy']) {
            continue;
        }
        $seqs = $deliveredTo[$endpoint] ?? [];
        $assert(count($seqs) === $EVENTS_PER_ENDPOINT, "{$endpoint} received all {$EVENTS_PER_ENDPOINT} events");
        $assert(
            $seqs === [1, 2, 3],
            "{$endpoint} received them in the order they happened (got " . (implode(',', $seqs) ?: 'none') . ')'
        );
    }

    $assert(
        count($deliveredTo['initech.example'] ?? []) === 0,
        'the dead endpoint received nothing'
    );
    $deadAttempts = $attempts['initech.example'] ?? 0;
    $assert(
        $deadAttempts === $EVENTS_PER_ENDPOINT * ($RETRY_LIMIT + 1),
        'each dead delivery was tried ' . ($RETRY_LIMIT + 1) . " times before it was given up ({$deadAttempts} attempts)"
    );

    // The dead letters are records you can query. Each one keeps the payload, so
    // it names the endpoint and the invoice, and the last error, which is what
    // answers "why did this customer not get the webhook".
    $dlq = $queen->queue($DELIVERIES)->dlq()->limit(50)->get();
    $messages = $dlq['messages'] ?? [];
    $dead = array_values(array_filter($messages, fn(array $m): bool => $m['data']['endpoint'] === 'initech.example'));

    $assert(count($dead) === $EVENTS_PER_ENDPOINT, "all {$EVENTS_PER_ENDPOINT} dead deliveries are in the dead-letter queue");
    $assert(
        count(array_filter($dead, fn(array $m): bool => str_contains($m['errorMessage'] ?? '', 'answered 500'))) === count($dead),
        'each dead letter carries the error that killed it'
    );
    $assert(
        count(array_filter($messages, fn(array $m): bool => $m['data']['endpoint'] === 'initech.example')) === count($messages),
        'no healthy delivery ended up in the dead-letter queue'
    );

    echo "\n  dead letters: " . implode('; ', array_map(
        fn(array $m): string => "{$m['data']['endpoint']}/{$m['data']['invoiceId']}: {$m['errorMessage']}",
        $dead
    )) . "\n";

    $queen->queue($DELIVERIES)->delete()->execute();

    echo "\nPASS: {$checks} checks\n";
} catch (Throwable $error) {
    fwrite(STDERR, "\nFAIL: " . $error->getMessage() . "\n");
    $exitCode = 1;
} finally {
    $queen->close();
}

exit($exitCode);
```

### C++

```cpp title="examples/apps/cpp/webhooks.cpp"
//
// A webhook sender: ordered per endpoint, retried by the broker, and
// dead-lettered with its error when it never succeeds.
//
// Deliveries to one endpoint have to arrive in the order the events happened,
// an endpoint that is down must not slow anybody else down, a failure is
// retried a bounded number of times, and what never succeeds has to end up
// somewhere a person can read it. Each endpoint gets a partition of its own,
// created by its first delivery, so a dead endpoint backs up its own partition
// and nothing else. Retries are the broker's retry budget, and an exhausted
// delivery lands in the dead-letter queue with the error attached.
//
//   webhook-deliveries (one partition per endpoint)
//     `-- group "sender"  POSTs each delivery; a nack spends one retry
//           `-- retryLimit spent -> dead-letter queue, with the error
//
// Build it. queen_client.hpp needs json.hpp and threadpool.hpp, which this
// repository carries under clients/server, and cpp-httplib, which it does not
// (Homebrew puts httplib.h under /opt/homebrew/include). The client switches on
// cpp-httplib's OpenSSL support, so -lssl -lcrypto are needed even over plain
// http.
//   mkdir -p build
//   c++ -std=c++17 -O1 -pthread \
//       -I../../../clients/client-cpp -I../../../clients/server/vendor \
//       -I/opt/homebrew/include -I"$(brew --prefix openssl)/include" \
//       webhooks.cpp -o build/webhooks \
//       -L"$(brew --prefix openssl)/lib" -lssl -lcrypto -lpthread
//
// Run it:
//   QUEEN_URL=http://localhost:6632 ./build/webhooks

#include "queen_client.hpp"

#include <atomic>
#include <chrono>
#include <cstdlib>
#include <exception>
#include <iostream>
#include <map>
#include <mutex>
#include <sstream>
#include <string>
#include <vector>

using queen::QueenClient;
using json = nlohmann::json;

static std::string run_id() {
    auto millis = std::chrono::duration_cast<std::chrono::milliseconds>(
                      std::chrono::system_clock::now().time_since_epoch())
                      .count();
    std::string out;
    const char* digits = "0123456789abcdefghijklmnopqrstuvwxyz";
    while (millis > 0) {
        out.insert(out.begin(), digits[millis % 36]);
        millis /= 36;
    }
    return out;
}

struct Endpoint {
    std::string host;
    bool healthy;
};

// Three subscribers. One of them answers 500 to everything. It is listed first,
// so its deliveries are the oldest in the queue and its partition is usually
// handed out first: a sender that let a failing endpoint hold up the others
// would fail the checks below.
static const std::vector<Endpoint> ENDPOINTS = {
    {"initech.example", false},
    {"acme.example", true},
    {"globex.example", true},
};
static const int EVENTS_PER_ENDPOINT = 3;
static const int RETRY_LIMIT = 2;
static const std::string GROUP = "sender";

// What a green run looks like: every healthy delivery succeeds once, and each
// dead one is attempted RETRY_LIMIT + 1 times before the budget is gone.
static const int HEALTHY_DELIVERIES = EVENTS_PER_ENDPOINT * 2;
static const int DEAD_DELIVERIES = EVENTS_PER_ENDPOINT;
static const int EXPECTED_ATTEMPTS =
    HEALTHY_DELIVERIES + DEAD_DELIVERIES * (RETRY_LIMIT + 1);

static int checks = 0;

// A throwing check: the failure travels to main() as an exception, which is
// what turns it into "FAIL: <reason>" and a non-zero exit.
static void check(bool condition, const std::string& description) {
    if (!condition) throw std::runtime_error(description);
    ++checks;
    std::cout << "  ok: " << description << std::endl;
}

static bool is_healthy(const std::string& host) {
    for (const Endpoint& endpoint : ENDPOINTS) {
        if (endpoint.host == host) return endpoint.healthy;
    }
    throw std::runtime_error("unknown endpoint " + host);
}

// Stands in for the HTTP POST to the subscriber. A real sender POSTs with an
// HTTP client and throws on any status that is not 2xx, which is what this
// does.
static void post_to_endpoint(const std::string& host) {
    if (!is_healthy(host)) {
        throw std::runtime_error(host + " answered 500");
    }
}

static std::string join(const std::vector<int>& values) {
    std::ostringstream out;
    for (size_t i = 0; i < values.size(); ++i) {
        if (i) out << ",";
        out << values[i];
    }
    return out.str();
}

int main() {
    const char* env_url = std::getenv("QUEEN_URL");
    const std::string QUEEN_URL = env_url ? env_url : "http://localhost:6632";
    const std::string DELIVERIES = "app-cpp-webhooks-" + run_id();

    QueenClient client(QUEEN_URL);

    std::string verdict;
    bool failed = false;

    try {
        std::cout << "broker " << QUEEN_URL << std::endl;

        // retryLimit is the delivery budget: a delivery that fails
        // RETRY_LIMIT + 1 times is filed in the dead-letter queue with its last
        // error, because dlqAfterMaxRetries is on. The broker turns that flag
        // on by default, and it is written out here because the dead-letter
        // checks at the end depend on it. leaseTime is how long the broker
        // waits for a sender that took a delivery and never came back before
        // it hands the delivery to another sender.
        //
        // The C++ QueueConfig struct has no dlqAfterMaxRetries field, so the
        // options are assembled by hand and posted through the client's own
        // HTTP transport, which keeps the base URL, the retries and the 429
        // backoff that every other call in this file uses.
        json configured = client.get_http_client()->post(
            "/api/v1/configure",
            json{{"queue", DELIVERIES},
                 {"options", {{"leaseTime", 30},
                              {"retryLimit", RETRY_LIMIT},
                              {"dlqAfterMaxRetries", true}}}});
        if (!configured.value("configured", false)) {
            throw std::runtime_error("the broker did not configure the queue");
        }

        // -------------------------------------------------------------- queuing
        //
        // The application emits events. Each delivery goes into the partition
        // of the endpoint it is for, which is what makes "in order per
        // subscriber" a property of the storage instead of something the
        // sender has to arrange.
        std::cout << "\nqueuing deliveries" << std::endl;
        for (int seq = 1; seq <= EVENTS_PER_ENDPOINT; ++seq) {
            for (const Endpoint& endpoint : ENDPOINTS) {
                // The event id. An application that retries its own emit does
                // not create a second delivery.
                client.queue(DELIVERIES).partition(endpoint.host).push({
                    json{{"transactionId",
                          endpoint.host + "-evt-" + std::to_string(seq)},
                         {"data", {{"endpoint", endpoint.host},
                                   {"seq", seq},
                                   {"type", "invoice.paid"},
                                   {"invoiceId", "INV-" + std::to_string(seq)}}}}
                });
            }
        }
        std::cout << "  " << EVENTS_PER_ENDPOINT * ENDPOINTS.size()
                  << " deliveries queued" << std::endl;

        // -------------------------------------------------------------- sending
        //
        // The sender pool: three workers, each a thread with a poll loop of its
        // own. A failed delivery goes back to the broker with its error: the
        // broker redelivers it and counts one retry, and deliveryAttempt says
        // which attempt this is. The retries live in the broker, so they
        // survive the sender dying halfway, which a retry loop inside the
        // handler would not.
        //
        // auto_ack(false): the handler acknowledges each delivery itself,
        // because this client's automatic nack sends no error and the dead
        // letter would keep none. partitions(1) and batch(1): every pop takes
        // ONE delivery, from one endpoint. A nack releases the endpoint's
        // lease, and the broker then refuses acks for anything else popped
        // under it, so a worker holding more than the failed delivery could not
        // acknowledge the rest.
        //
        //   wait(false)     every pop answers at once. This client has no
        //                   setter for the long-poll timeout, which is 30
        //                   seconds, and the stop flag and the idle deadline
        //                   are only consulted between pops.
        //   idle_millis     the deadline. A delivery that never comes back
        //                   fails this run instead of hanging it.
        //   limit()         counts per worker, so it is only a backstop. The
        //                   stop flag ends the loop, and it is raised on the
        //                   outcome: every healthy delivery sent and every dead
        //                   one dead-lettered.
        std::cout << "\nsending" << std::endl;
        std::mutex lock;
        std::map<std::string, std::vector<int>> delivered_to;
        std::map<std::string, int> attempts;
        int dead_lettered = 0;
        int delivered_total = 0;
        std::atomic<bool> stop{false};
        std::exception_ptr handler_error;

        client.queue(DELIVERIES)
            .group(GROUP)
            .subscription_mode("all")
            .concurrency(3)
            .partitions(1)
            .batch(1)
            .each()
            .auto_ack(false)
            .limit(EXPECTED_ATTEMPTS)
            .wait(false)
            .idle_millis(6000)
            .consume([&](const json& msg) {
                try {
                    const std::string host = msg["data"]["endpoint"].get<std::string>();
                    const int seq = msg["data"]["seq"].get<int>();

                    std::string error;
                    try {
                        post_to_endpoint(host);
                    } catch (const std::exception& e) {
                        error = e.what();
                    }

                    {
                        // Recorded before the ack: the ack releases the
                        // endpoint's lease, and another worker could take and
                        // record the next delivery before this one otherwise.
                        std::lock_guard<std::mutex> guard(lock);
                        attempts[host] += 1;
                        if (error.empty()) {
                            delivered_to[host].push_back(seq);
                            ++delivered_total;
                            std::cout << "  " << host << " <- event " << seq << std::endl;
                        } else {
                            std::cout << "  " << host << " <- event " << seq
                                      << " failed on attempt "
                                      << msg.value("deliveryAttempt", 0) << ": "
                                      << error << std::endl;
                        }
                    }

                    // The ack names the consumer group: without it the broker
                    // looks for the lease on the queue's own cursor and refuses
                    // the ack. A failure carries its error, which is what the
                    // dead letter keeps once the budget is spent.
                    json context = {{"group", GROUP}};
                    if (!error.empty()) context["error"] = error;
                    json ack = client.ack(msg, error.empty(), context);

                    // A refused acknowledgement still arrives as HTTP 200 with
                    // success: false on the item, so the item is the broker's
                    // verdict. The outer "success" only says the call did not
                    // throw.
                    if (!ack.value("success", false) || !ack["result"].is_array() ||
                        ack["result"].empty() ||
                        !ack["result"][0].value("success", false)) {
                        throw std::runtime_error("the broker refused the acknowledgement for " +
                                                 host + "/" + std::to_string(seq));
                    }

                    std::lock_guard<std::mutex> guard(lock);
                    if (ack["result"][0].value("dlq", false)) {
                        // This failure spent the last of the retry budget, and
                        // the delivery is now a dead letter.
                        ++dead_lettered;
                    }
                    if (delivered_total >= HEALTHY_DELIVERIES &&
                        dead_lettered >= DEAD_DELIVERIES) {
                        stop = true;
                    }
                } catch (...) {
                    // auto_ack is off, so nothing acknowledges behind this
                    // handler's back: an escaped exception leaves the delivery
                    // leased until it expires. Carry the error out and stop.
                    std::lock_guard<std::mutex> guard(lock);
                    if (!handler_error) handler_error = std::current_exception();
                    stop = true;
                }
            }, &stop);
        if (handler_error) std::rethrow_exception(handler_error);

        // ------------------------------------------------------------- checking
        std::cout << "\nchecking" << std::endl;

        std::vector<int> in_order;
        for (int seq = 1; seq <= EVENTS_PER_ENDPOINT; ++seq) in_order.push_back(seq);

        for (const Endpoint& endpoint : ENDPOINTS) {
            if (!endpoint.healthy) continue;
            const std::vector<int>& seqs = delivered_to[endpoint.host];
            check(seqs.size() == static_cast<size_t>(EVENTS_PER_ENDPOINT),
                  endpoint.host + " received all " +
                      std::to_string(EVENTS_PER_ENDPOINT) + " events");
            check(seqs == in_order,
                  endpoint.host + " received them in the order they happened (got " +
                      (seqs.empty() ? std::string("none") : join(seqs)) + ")");
        }

        check(delivered_to["initech.example"].empty(),
              "the dead endpoint received nothing");
        check(attempts["initech.example"] == DEAD_DELIVERIES * (RETRY_LIMIT + 1),
              "each dead delivery was tried " + std::to_string(RETRY_LIMIT + 1) +
                  " times before it was given up (" +
                  std::to_string(attempts["initech.example"]) + " attempts)");

        // The dead letters are records you can query. Each one keeps the
        // payload, so it names the endpoint and the invoice, and the last error,
        // which is what answers "why did this customer not get the webhook".
        //
        // dlq().get() answers a transport failure with an empty page and does
        // not throw, so an unreachable broker shows up here as the next check
        // failing.
        json dlq = client.queue(DELIVERIES).dlq().limit(50).get();
        const json& letters = dlq["messages"];

        std::vector<json> dead;
        for (const json& letter : letters) {
            if (letter["data"]["endpoint"] == "initech.example") dead.push_back(letter);
        }

        check(dead.size() == static_cast<size_t>(EVENTS_PER_ENDPOINT),
              "all " + std::to_string(EVENTS_PER_ENDPOINT) +
                  " dead deliveries are in the dead-letter queue");

        bool every_letter_explains_itself = true;
        for (const json& letter : dead) {
            const std::string error = letter.value("errorMessage", std::string());
            if (error.find("answered 500") == std::string::npos) {
                every_letter_explains_itself = false;
            }
        }
        check(every_letter_explains_itself,
              "each dead letter carries the error that killed it");

        bool only_the_dead_endpoint = true;
        for (const json& letter : letters) {
            if (letter["data"]["endpoint"] != "initech.example") only_the_dead_endpoint = false;
        }
        check(only_the_dead_endpoint,
              "no healthy delivery ended up in the dead-letter queue");

        std::ostringstream summary;
        for (size_t i = 0; i < dead.size(); ++i) {
            if (i) summary << "; ";
            summary << dead[i]["data"]["endpoint"].get<std::string>() << "/"
                    << dead[i]["data"]["invoiceId"].get<std::string>() << ": "
                    << dead[i].value("errorMessage", std::string());
        }
        std::cout << "\n  dead letters: " << summary.str() << std::endl;

        // Clean up on success only: a failed run leaves the queue and its dead
        // letters on the broker to be looked at.
        client.queue(DELIVERIES).del();

        verdict = "\nPASS: " + std::to_string(checks) + " checks";
    } catch (const std::exception& err) {
        verdict = std::string("\nFAIL: ") + err.what();
        failed = true;
    }

    client.close();

    (failed ? std::cerr : std::cout) << verdict << std::endl;
    return failed ? 1 : 0;
}
```

### curl

```bash title="examples/apps/http/webhooks.sh"
#!/usr/bin/env bash
#
# A webhook sender, with nothing but curl: ordered per endpoint, retried by the
# broker, and dead-lettered with its error when it never succeeds.
#
# Deliveries to one endpoint have to arrive in the order the events happened,
# an endpoint that is down must not slow anybody else down, a failure is
# retried a bounded number of times, and what never succeeds has to end up
# somewhere a person can read it. Each endpoint gets a partition of its own,
# created by its first delivery, so a dead endpoint backs up its own partition
# and nothing else. Retries are the broker's retry budget, and an exhausted
# delivery lands in the dead-letter queue with the error attached.
#
#   webhook-deliveries (one partition per endpoint)
#     └── group "sender"  POSTs each delivery; a failure spends one retry
#           └── retryLimit spent -> dead-letter queue, with the error
#
# One sender runs per endpoint, as a background subshell on the
# partition-scoped pop route, so the isolation is real here: the dead endpoint
# spends the whole run failing and retrying while the two healthy ones are
# already finished.
#
# Run it:
#   QUEEN_URL=http://localhost:6632 bash webhooks.sh

set -euo pipefail

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

# The name carries the language and a per-run suffix, so every application in
# every language can share one broker and no run inherits another's state.
RUN="$(date +%s)-$$"
DELIVERIES="app-http-webhooks-$RUN"
GROUP=app-http-sender

# Three subscribers: endpoint and whether it is healthy. One of them answers 500
# to everything. It is listed first, so its deliveries are the oldest in the
# queue and its sender starts first; the checks below show that the healthy
# endpoints get every event anyway, in order.
ENDPOINTS='initech.example no
acme.example yes
globex.example yes'
EVENTS_PER_ENDPOINT=3
RETRY_LIMIT=2

# 1,2,3: what a healthy subscriber must receive, in that order.
EXPECTED_SEQS="$(seq 1 "$EVENTS_PER_ENDPOINT" | paste -sd, -)"

# Every pop long-polls for this many milliseconds and no longer, so a sender
# re-checks its own progress at least once a second.
POLL_MS=1000

# The bound that keeps a stall from becoming a hang. A sender that has not
# finished its endpoint by then stops, and the checks that follow report what is
# missing. Never wait for silence; wait for a total, with a deadline.
SEND_MS=30000

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

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

# 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.
cleanup() {
  local status=$?
  rm -rf "$TMP"
  if [ "$status" -ne 0 ]; then
    echo
    echo "FAIL: ${FAILURE:-a command exited with status $status}"
  fi
  exit "$status"
}
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"
}

# A millisecond clock, for the deadline only. 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.
#
# $OUT is per-process: each sender below runs as a background subshell and points
# it at its own file, so three concurrent pops never overwrite each other's
# response. There is no --fail, because Queen reports outcomes in the body and
# several of the interesting ones arrive as 200.
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
}

# healthy <endpoint>: prints yes or no. The shell this has to run on has no
# associative arrays, so the subscriber table is a few lines of text and awk is
# the lookup.
healthy() {
  printf '%s\n' "$ENDPOINTS" | awk -v e="$1" '$1 == e { print $2 }'
}

# Stands in for the HTTP POST to the subscriber. A real sender calls curl and
# fails on any status that is not 2xx, which is what returning non-zero does
# here.
post_to_endpoint() {
  [ "$(healthy "$1")" = yes ] || return 1
  return 0
}

# lines <file>: how many lines a tally file holds, 0 when it does not exist yet.
lines() {
  [ -f "$1" ] || { echo 0; return; }
  wc -l < "$1" | tr -d ' '
}

echo "broker $QUEEN_URL"

# ---------------------------------------------------------------------------
# retryLimit is the delivery budget: a delivery that fails RETRY_LIMIT + 1
# times is filed in the dead-letter queue with its last error. Dead-lettering
# is on by default on a new queue (deadLetterQueue and dlqAfterMaxRetries); with
# both off, a delivery that spends its budget is skipped and gone. Both flags
# are sent anyway, because /configure merges: an option this body does not name
# keeps whatever the queue already has, and depending on that is how a
# configuration drifts.
#
# leaseTime is how long the broker waits for a sender that took a delivery and
# never came back before it hands the delivery to another sender.
# ---------------------------------------------------------------------------
configure_body="$(jq -n --arg queue "$DELIVERIES" --argjson retry "$RETRY_LIMIT" \
  '{queue: $queue,
    options: {leaseTime: 30, retryLimit: $retry,
              deadLetterQueue: true, dlqAfterMaxRetries: true}}')"
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 delivery budget of $RETRY_LIMIT retries"

# ------------------------------------------------------------------------ queuing
#
# The application emits events. Each delivery goes into the partition of the
# endpoint it is for, which is what makes "in order per subscriber" a property
# of the storage instead of something the sender has to arrange. Nothing was
# declared for a subscriber in advance: the partition comes into existence with
# the first delivery to it.
echo
echo "queuing deliveries"
seq_no=1
while [ "$seq_no" -le "$EVENTS_PER_ENDPOINT" ]; do
  while read -r endpoint is_healthy; do
    # The event id makes the enqueue idempotent: an application that retries its
    # own emit does not create a second delivery. The wire field for the body is
    # "payload"; "data" is what a pop calls it on the way back.
    body="$(jq -n --arg queue "$DELIVERIES" --arg endpoint "$endpoint" \
      --argjson seq "$seq_no" \
      '{items: [{
         queue:     $queue,
         partition: $endpoint,
         transactionId: ($endpoint + "-evt-" + ($seq | tostring)),
         payload: {endpoint: $endpoint, seq: $seq, type: "invoice.paid",
                   invoiceId: ("INV-" + ($seq | tostring))}
       }]}')"
    request POST /api/v1/push "$body"
    [ "$STATUS" = 201 ] || fail "push of $endpoint/$seq_no 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 $endpoint/$seq_no came back $(jq -r '.[0].status' "$OUT")"
  done <<EOF
$ENDPOINTS
EOF
  seq_no=$((seq_no + 1))
done
echo "  $((EVENTS_PER_ENDPOINT * 3)) deliveries queued"

# ------------------------------------------------------------------------ sending
#
# die() is a sender's fail(): a sender is a subshell, so its variables die with
# it and the parent would never see FAILURE. It leaves the reason in a file the
# parent reads after wait().
die() { printf '%s\n' "$*" > "$TMP/sender-error"; exit 1; }

# send_endpoint <endpoint>
#
# One sender, one endpoint. It pops that endpoint's partition by name, posts
# what it gets, and reports the outcome back to the broker; the loop is what an
# SDK's consume() does, and the ack status is what an SDK's autoAck derives from
# a handler that returned or threw.
#
# It stops when the endpoint's deliveries are resolved, either because every
# event was delivered or because every event was dead-lettered, and in any case
# at the deadline.
send_endpoint() {
  local endpoint="$1"
  local deadline popfile delivered_file attempts_file dead_file
  local txn partition_id lease event_seq attempt aborted ack_body

  OUT="$TMP/$endpoint-body"
  popfile="$TMP/$endpoint-pop"
  delivered_file="$TMP/delivered-$endpoint"
  attempts_file="$TMP/attempts-$endpoint"
  dead_file="$TMP/dead-$endpoint"
  : > "$delivered_file"
  : > "$attempts_file"
  : > "$dead_file"
  deadline=$(( $(now_ms) + SEND_MS ))

  while [ "$(( $(lines "$delivered_file") + $(lines "$dead_file") ))" \
          -lt "$EVENTS_PER_ENDPOINT" ]; do
    [ "$(now_ms)" -lt "$deadline" ] || break

    # The partition-scoped pop route claims exactly the partition you name. The
    # queue-scoped one lets the broker pick, and with `partitions` at its
    # default of 1 (and no autopilot=true) it would claim a single partition per
    # call anyway; naming the partition is what makes this sender belong to one
    # endpoint.
    #
    # 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. It seeds a cursor that does not exist yet, and is ignored on
    # every later pop.
    request GET "/api/v1/pop/queue/$DELIVERIES/partition/$endpoint?consumerGroup=$GROUP&subscriptionMode=all&batch=10&wait=true&timeout=$POLL_MS"
    # 204 is an empty pop, with no body at all. Here it means the partition is
    # quiet for the moment; the loop condition decides whether that is the end.
    [ "$STATUS" != 204 ] || continue
    [ "$STATUS" = 200 ] || die "pop on $endpoint returned HTTP $STATUS"
    cp "$OUT" "$popfile"

    # Every popped message carries deliveryAttempt: 1 the first time, one more
    # each time the same delivery comes back.
    lease="$(jq -r .leaseId "$popfile")"
    jq -r '.messages[] | [.transactionId, .partitionId, (.data.seq | tostring),
                          (.deliveryAttempt | tostring)] | @tsv' \
      "$popfile" > "$TMP/$endpoint-batch"

    aborted=0
    while IFS=$'\t' read -r txn partition_id event_seq attempt; do
      echo "$event_seq" >> "$attempts_file"

      if ! post_to_endpoint "$endpoint"; then
        echo "  $endpoint <- event $event_seq failed on attempt $attempt: $endpoint answered 500"
        # -------------------------------------------------------------------
        # The delivery failed. A `failed` ack is the nack: it leaves the cursor
        # just below this message, so this one and everything after it in the
        # batch is delivered again, and it spends one retry of the budget.
        # There is no loop in this sender and no sleep: the redelivery is the
        # retry, and it survives this process dying halfway, which a retry
        # loop inside the sender would not.
        #
        # A lease that merely runs out spends nothing: only a `failed` ack
        # does, so a sender that keeps crashing never uses up a delivery's
        # budget by crashing. deliveryAttempt is how a sender can set a limit
        # of its own for that case.
        #
        # `error` is the reason. A plain nack does not store it; the nack that
        # dead-letters the delivery files it with the dead letter, where a
        # person reads it later.
        #
        # The rest of the batch is left alone. The nack released this
        # partition's lease, so an ack of any later message of this batch would
        # be refused, and those messages are coming back anyway.
        # -------------------------------------------------------------------
        ack_body="$(jq -n --arg txn "$txn" --arg partitionId "$partition_id" \
          --arg group "$GROUP" --arg lease "$lease" \
          --arg error "$endpoint answered 500" \
          '{transactionId: $txn, partitionId: $partitionId, consumerGroup: $group,
            leaseId: $lease, status: "failed", error: $error}')"
        request POST /api/v1/ack "$ack_body"
        [ "$STATUS" = 200 ] || die "ack for $endpoint returned HTTP $STATUS"
        [ "$(jq -r '.[0].success' "$OUT")" = true ] \
          || die "ack refused for $endpoint: $(jq -r '.[0].error' "$OUT")"

        # dlq on the ack result is how a sender on raw HTTP learns the budget ran
        # out: true means this nack filed the dead letter and moved the cursor
        # past the delivery.
        if [ "$(jq -r '.[0].dlq' "$OUT")" = true ]; then
          echo "$event_seq" >> "$dead_file"
          echo "  $endpoint dead-lettered event $event_seq, its retry budget is spent"
        fi
        aborted=1
        break
      fi

      echo "$event_seq" >> "$delivered_file"
      echo "  $endpoint <- event $event_seq"
    done < "$TMP/$endpoint-batch"

    # Everything in the batch was posted, so commit the batch with one ack of its
    # LAST message: an ack is a cursor commit, so that completes every earlier
    # message of this partition for this group too. consumerGroup is mandatory,
    # here as everywhere: omit it and the ack goes to __QUEUE_MODE__, a cursor
    # this sender never read from, so it is refused and the batch comes back
    # every time its lease runs out.
    if [ "$aborted" = 0 ]; then
      ack_body="$(jq -c --arg group "$GROUP" '{
        transactionId: .messages[-1].transactionId,
        partitionId:   .messages[-1].partitionId,
        consumerGroup: $group,
        leaseId:       .leaseId,
        status:        "completed"
      }' "$popfile")"
      request POST /api/v1/ack "$ack_body"
      [ "$STATUS" = 200 ] || die "ack for $endpoint returned HTTP $STATUS"
      [ "$(jq -r '.[0].success' "$OUT")" = true ] \
        || die "ack refused for $endpoint: $(jq -r '.[0].error' "$OUT")"
    fi
  done
}

echo
echo "sending"
rm -f "$TMP/sender-error"
PIDS=""
while read -r endpoint is_healthy; do
  send_endpoint "$endpoint" &
  PIDS="$PIDS $!"
done <<EOF
$ENDPOINTS
EOF
for pid in $PIDS; do
  wait "$pid" \
    || fail "a sender stopped: $(cat "$TMP/sender-error" 2>/dev/null || echo 'no reason recorded')"
done

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

while read -r endpoint is_healthy; do
  [ "$is_healthy" = yes ] || continue
  check "$(lines "$TMP/delivered-$endpoint")" "$EVENTS_PER_ENDPOINT" \
    "$endpoint received all $EVENTS_PER_ENDPOINT events"
  check "$(paste -sd, - < "$TMP/delivered-$endpoint")" "$EXPECTED_SEQS" \
    "$endpoint received them in the order they happened"
done <<EOF
$ENDPOINTS
EOF

check "$(lines "$TMP/delivered-initech.example")" 0 \
  'the dead endpoint received nothing'

check "$(lines "$TMP/attempts-initech.example")" \
  "$((EVENTS_PER_ENDPOINT * (RETRY_LIMIT + 1)))" \
  "each dead delivery was tried $((RETRY_LIMIT + 1)) times before it was given up"

check "$(lines "$TMP/dead-initech.example")" "$EVENTS_PER_ENDPOINT" \
  'the broker reported a dead letter on the ack that spent each budget'

# ---------------------------------------------------------------------------
# The dead letters are records you can query. Each one keeps the payload, so it
# names the endpoint and the invoice, and the last error, which is what answers
# "why did this customer not get the webhook".
#
# Dead letters follow the queue's retention: one goes when retention removes
# the message it holds. With retention off, as here, a dead letter stays until
# you replay it with POST /api/v1/messages/:partitionId/:transactionId/retry,
# delete it, or delete the queue.
# ---------------------------------------------------------------------------
request GET "/api/v1/dlq?queue=$DELIVERIES&limit=50"
[ "$STATUS" = 200 ] || fail "the dead-letter listing returned HTTP $STATUS"
cp "$OUT" "$TMP/dlq"

check "$(jq '[.messages[] | select(.data.endpoint == "initech.example")] | length' "$TMP/dlq")" \
  "$EVENTS_PER_ENDPOINT" \
  "all $EVENTS_PER_ENDPOINT dead deliveries are in the dead-letter queue"
check "$(jq '[.messages[] | select((.errorMessage // "") | contains("answered 500"))] | length' "$TMP/dlq")" \
  "$EVENTS_PER_ENDPOINT" 'each dead letter carries the error that killed it'
check "$(jq '[.messages[] | select(.data.endpoint != "initech.example")] | length' "$TMP/dlq")" 0 \
  'no healthy delivery ended up in the dead-letter queue'

echo
echo "  dead letters: $(jq -r '[.messages[] | .data.endpoint + "/" + .data.invoiceId + ": " + (.errorMessage // "")] | sort | join("; ")' "$TMP/dlq")"

# Clean up on success only: a failed run leaves the queue, and its dead letters,
# on the broker to be looked at. Deleting a queue that does not exist also
# answers 200, with deleted:false, so the field is what to check.
request DELETE "/api/v1/resources/queues/$DELIVERIES"
[ "$(jq -r .deleted "$OUT")" = true ] || fail 'the queue was not deleted'

echo
echo "PASS: $CHECKS checks"
```

## Limits

A failing delivery holds up its own endpoint: later deliveries to it wait until it succeeds or is
dead-lettered. That is the price of order, and other endpoints do not pay it.

Retries are immediate, because a nacked delivery comes back on the next pop. `retryDelay` is
accepted by `/configure` and stored, but the 2.0 consumption engine does not apply it. For a pause
between attempts, ack the delivery and schedule it again as a timer into the same partition, in one
[transaction](/concepts/transactions/); it then goes out behind the newer deliveries to that
endpoint.

Only a nack spends the budget. A sender that dies mid-delivery costs nothing, because the lease
expires and the delivery comes back, which also means a delivery that crashes its sender every
time is retried forever. Messages carry `deliveryAttempt` so a handler can dead-letter such a
delivery past a threshold of its own (the Go and Rust clients do not expose the field yet).

Dead-lettering is on by default (`deadLetterQueue` and `dlqAfterMaxRetries` are both `true` on a
new queue). Turn both off and a delivery that spends its budget is skipped, and it is gone.

The endpoint can receive a delivery twice, for instance when a lease expires after the POST and
before the ack. Send the event id in a header so the receiver can drop repeats.

The dead-letter path has its own workload in the Jepsen suite (W7). It caught a real bug, a dead
letter filed without its payload under kill -9, which was fixed in `7a2a1711`; [Jepsen](/benchmarks/jepsen/)
lists every run.

Source: https://queenmq.com/guides/webhooks/index.mdx
