---
title: "Rate limiters"
description: "Two ways to hold a tenant to a quota with the count kept in the broker: a streaming window that throttles, and an admission controller that moves work to the next window instead of refusing it."
---

> 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

# Rate limiters

A rate limiter needs a count that is exact and survives restarts, and the usual place for it is a
Redis next to the broker. These two programs keep it in the broker instead, in the same log as the
requests, and they answer the question every limiter has to answer differently. The first counts
requests per API key in a two-second window and turns every window over the quota into a throttle
decision for a gateway to enforce. The second never says no: it admits a tenant's requests up to
the quota and puts the rest on a timer that brings them back when the window opens again.

## Run them

```bash
docker run --platform linux/amd64 -d --name queen -p 6632:6632 ghcr.io/queen-mq/queen:latest
npm install queen-mq
node rate-limiter.mjs     # the counter, saved from the first set of tabs below
node deferred-work.mjs    # the admission controller, saved from the second
```

## The counter

What the counter printed against a 2.0.0-beta.6 node:

```text
starting the counter

taking traffic
  key-quiet: 3 requests
  key-noisy: 20 requests

enforcing
  key-noisy: 9 in a window, over by 4
  key-quiet: 3 in a window, within quota
  key-noisy: 11 in a window, over by 6

checking
  ok: every request reached a closed window before the deadline
  ok: the quiet key was counted exactly
  ok: the noisy key was counted exactly
  ok: the noisy key was throttled
  ok: the quiet key was never throttled, so the limiter is not just firing at everything
  ok: every decision is on the queue the gateway reads
  ok: each decision carries the count and the quota that produced it

PASS: 7 checks
```

This run caught a window boundary in the middle of the noisy burst, so its twenty requests were
counted as 9 and 11, both over the quota of 5. That is why the program sends twenty: split any way
at all, one side holds at least ten, so the noisy key is always throttled, while the quiet key's
three requests can never reach five.

### JavaScript

```js title="examples/apps/js/rate-limiter.mjs"
//
// A rate limiter built from a streaming query.
//
// The textbook rate limiter counts requests per API key in a fixed window,
// usually with a counter in Redis, which is one more system to run and one
// more place where the count can drift from the requests.
//
// Here the counter is a windowed aggregation over the request queue itself.
// Each cycle of the stream commits the window state, the closed windows it
// emits and the ack of the requests it counted as one entry in the broker's
// log, so the count cannot drift from the requests, and it survives a restart
// of this process because the state is in the broker.
//
//   api-requests (one partition per API key)
//     └── streaming query: tumbling window, count per key
//           └── api-usage  -> the gate: over quota becomes a throttle decision
//                 └── api-throttled
//
// Run it:
//   QUEEN_URL=http://localhost:6632 node rate-limiter.mjs

import { Queen, Stream } from 'queen-mq'

const QUEEN_URL = process.env.QUEEN_URL || 'http://localhost:6632'
const RUN = Date.now().toString(36)
const REQUESTS = `app-js-api-requests-${RUN}`
const USAGE = `app-js-api-usage-${RUN}`
const THROTTLED = `app-js-api-throttled-${RUN}`
const QUERY_ID = `app-js-rate-limiter-${RUN}`

const WINDOW_SECONDS = 2
const QUOTA_PER_WINDOW = 5

// Two tenants. One is a well behaved integration, the other is a runaway script
// someone left in a loop.
const QUIET_KEY = 'key-quiet'
const NOISY_KEY = 'key-noisy'
const QUIET_REQUESTS = 3
const NOISY_REQUESTS = 20

// Why those numbers make the check deterministic: a window is a slice of time,
// so a burst can land on either side of a boundary. Twenty requests split in
// any way at all leave at least ten on one side, which is over a quota of five,
// so the noisy key is always caught. Three requests cannot reach five however
// they are split, so the quiet key is never caught by accident.

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

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

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

  for (const q of [REQUESTS, USAGE, THROTTLED]) {
    await queen.queue(q).config({ leaseTime: 30, retryLimit: 3 }).create()
  }

  // ------------------------------------------------------------- the counter
  //
  // The stream runs in this process, as a consumer group of its own. A new
  // group starts at the tail of the queue unless it asks otherwise, and a
  // request pushed while the stream was still starting would be missed, so it
  // asks for subscriptionMode 'all': every request in the queue is counted.
  //
  // The partition is the aggregation key, so the window state is per API key
  // without a word about keys here: the producer decides, by partition.
  console.log('\nstarting the counter')
  stream = await Stream
    .from(queen.queue(REQUESTS))
    .windowTumbling({ seconds: WINDOW_SECONDS, idleFlushMs: 800 })
    .aggregate({
      // The extractors receive the payload itself, not the envelope.
      requests: () => 1,
      cost: (r) => r.cost ?? 1,
    })
    .to(queen.queue(USAGE))
    .run({
      queryId: QUERY_ID,
      url: QUEEN_URL,
      subscriptionMode: 'all',
      batchSize: 200,
      maxPartitions: 8,
      maxWaitMillis: 200,
    })

  // ------------------------------------------------------------- the traffic
  console.log('\ntaking traffic')
  const send = async (key, n) => {
    for (let i = 1; i <= n; i++) {
      await queen.queue(REQUESTS).partition(key).push({
        data: { key, path: '/v1/things', cost: 1, at: Date.now() },
      })
    }
    console.log(`  ${key}: ${n} requests`)
  }
  await send(QUIET_KEY, QUIET_REQUESTS)
  await send(NOISY_KEY, NOISY_REQUESTS)

  // ---------------------------------------------------------------- the gate
  //
  // The enforcement point. It reads each closed window and turns the ones over
  // quota into throttle decisions. It is separate from the counter on purpose:
  // the counting is exact and stays the same, while the policy is yours and
  // changes on its own schedule.
  console.log('\nenforcing')
  const counted = {}
  const decisions = []
  const complete = () =>
    (counted[QUIET_KEY] ?? 0) === QUIET_REQUESTS && (counted[NOISY_KEY] ?? 0) === NOISY_REQUESTS
  const deadline = Date.now() + 30000

  while (!complete() && Date.now() < deadline) {
    const windows = await queen
      .queue(USAGE)
      .group('rate-limiter-gate')
      .subscriptionMode('all')
      .batch(50)
      .partitions(10)
      .wait(true)
      .timeoutMillis(2000)
      .pop()

    for (const w of windows) {
      const key = w.partition
      counted[key] = (counted[key] ?? 0) + w.data.requests
      const overBy = w.data.requests - QUOTA_PER_WINDOW

      if (overBy > 0) {
        // The decision is a message: whatever enforces it (an edge worker, a
        // gateway, the API itself) reads this queue and gets the decisions in
        // order, per key. It commits with the ack of the window it came from,
        // so a crash in between cannot lose a decision or make two.
        await queen
          .transaction()
          .queue(THROTTLED).partition(key)
          .push({ data: { key, window: w.data.requests, quota: QUOTA_PER_WINDOW, overBy } })
          .ack(w, 'completed', { consumerGroup: 'rate-limiter-gate' })
          .commit()
        decisions.push({ key, overBy })
        console.log(`  ${key}: ${w.data.requests} in a window, over by ${overBy}`)
      } else {
        await queen.ack(w, true, { group: 'rate-limiter-gate' })
        console.log(`  ${key}: ${w.data.requests} in a window, within quota`)
      }
    }
  }

  // --------------------------------------------------------------- checking
  console.log('\nchecking')
  assert(complete(), 'every request reached a closed window before the deadline')
  assert(counted[QUIET_KEY] === QUIET_REQUESTS, 'the quiet key was counted exactly')
  assert(counted[NOISY_KEY] === NOISY_REQUESTS, 'the noisy key was counted exactly')

  assert(decisions.length > 0, 'the noisy key was throttled')
  assert(
    decisions.every(d => d.key === NOISY_KEY),
    'the quiet key was never throttled, so the limiter is not just firing at everything'
  )

  // The decisions are readable by whatever enforces them, in order, per key.
  const throttled = await queen
    .queue(THROTTLED)
    .batch(50)
    .partitions(10)
    .wait(true)
    .pop()
  assert(throttled.length === decisions.length, 'every decision is on the queue the gateway reads')
  assert(
    throttled.every(m => m.data.window > m.data.quota),
    'each decision carries the count and the quota that produced it'
  )

  await stream.stop()
  stream = null
  for (const q of [REQUESTS, USAGE, THROTTLED]) await queen.queue(q).delete()

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

### Python

```python title="examples/apps/py/rate_limiter.py"
#
# A rate limiter built from a streaming query.
#
# The textbook rate limiter counts requests per API key in a fixed window,
# usually with a counter in Redis, which is one more system to run and one
# more place where the count can drift from the requests.
#
# Here the counter is a windowed aggregation over the request queue itself.
# Each cycle of the stream commits the window state, the closed windows it
# emits and the ack of the requests it counted as one entry in the broker's
# log, so the count cannot drift from the requests, and it survives a restart
# of this process because the state is in the broker.
#
#   api-requests (one partition per API key)
#     `-- streaming query: tumbling window, count per key
#           `-- api-usage  -> the gate: over quota becomes a throttle decision
#                 `-- api-throttled
#
# Run it:
#   QUEEN_URL=http://localhost:6632 python3 rate_limiter.py

import asyncio
import os
import sys
import time

from queen import Queen, Stream

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

# The names are 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}"
REQUESTS = f"app-py-api-requests-{RUN}"
USAGE = f"app-py-api-usage-{RUN}"
THROTTLED = f"app-py-api-throttled-{RUN}"

# The query id is this streaming query's identity in the broker. Its window
# state is keyed by it, so restarting the program with the same id resumes the
# same windows instead of starting new ones.
QUERY_ID = f"app-py-rate-limiter-{RUN}"

WINDOW_SECONDS = 2
QUOTA_PER_WINDOW = 5

# Two tenants. One is a well behaved integration, the other is a runaway script
# someone left in a loop.
QUIET_KEY = "key-quiet"
NOISY_KEY = "key-noisy"
QUIET_REQUESTS = 3
NOISY_REQUESTS = 20

# Why those numbers make the check deterministic: a window is a slice of time,
# so a burst can land on either side of a boundary. Twenty requests split in
# any way at all leave at least ten on one side, which is over a quota of five,
# so the noisy key is always caught. Three requests cannot reach five however
# they are split, so the quiet key is never caught by accident.

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 main() -> int:
    # The whole client is async: every call below is awaited, and this is the
    # one event loop they all run on, including the stream's polling task.
    queen = Queen(url=QUEEN_URL)
    stream = None
    verdict, failed = "", False

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

        for name in (REQUESTS, USAGE, THROTTLED):
            # The config keys are snake_case in Python and the client converts
            # them to the camelCase the broker expects.
            await queen.queue(name).config({"lease_time": 30, "retry_limit": 3}).create()

        # ------------------------------------------------------------ counter
        #
        # The stream runs in this process, as a consumer group of its own. A
        # new group starts at the tail of the queue unless it asks otherwise,
        # and a request pushed while the stream was still starting would be
        # missed, so it asks for subscription_mode "all": every request in the
        # queue is counted.
        #
        # The partition is the aggregation key, so the window state is per API
        # key without a word about keys here: the producer decides, by
        # partition.
        print("\nstarting the counter")
        stream = await (
            # from_ carries a trailing underscore because `from` is a Python
            # keyword; it is the same entry point as the other clients'.
            Stream.from_(queen.queue(REQUESTS))
            # Tumbling: fixed, non-overlapping windows, one set per partition.
            # A window closes when its time is up; idle_flush_ms also closes one
            # whose partition has gone quiet, which is what lets a short program
            # finish.
            .window_tumbling(seconds=WINDOW_SECONDS, idle_flush_ms=800)
            .aggregate(
                {
                    # The extractors receive the payload itself, without the
                    # envelope around it: the cost is r["cost"].
                    "requests": lambda r: 1,
                    "cost": lambda r: r.get("cost", 1),
                }
            )
            .to(queen.queue(USAGE))
            # run() registers the query and then leaves the polling loop running
            # as an asyncio task. Its options are keyword arguments here, in
            # snake_case.
            .run(
                query_id=QUERY_ID,
                url=QUEEN_URL,
                subscription_mode="all",
                batch_size=200,
                max_partitions=8,
                max_wait_millis=200,
            )
        )

        # ------------------------------------------------------------ traffic
        print("\ntaking traffic")

        async def send(key: str, n: int) -> None:
            for _ in range(n):
                await queen.queue(REQUESTS).partition(key).push(
                    {"data": {"key": key, "path": "/v1/things", "cost": 1, "at": int(time.time() * 1000)}}
                )
            print(f"  {key}: {n} requests")

        await send(QUIET_KEY, QUIET_REQUESTS)
        await send(NOISY_KEY, NOISY_REQUESTS)

        # --------------------------------------------------------------- gate
        #
        # The enforcement point. It reads each closed window and turns the ones
        # over quota into throttle decisions. It is separate from the counter on
        # purpose: the counting is exact and stays the same, while the policy is
        # yours and changes on its own schedule.
        print("\nenforcing")
        counted: dict = {}
        decisions = []

        def complete() -> bool:
            return (
                counted.get(QUIET_KEY, 0) == QUIET_REQUESTS
                and counted.get(NOISY_KEY, 0) == NOISY_REQUESTS
            )

        # One key's burst can fall on either side of a window boundary and
        # arrive as two windows. This adds the windows up per key and waits for
        # the totals it expects, with a deadline. Waiting for a quiet period
        # instead would race the timer that closes the last window.
        deadline = time.monotonic() + 30

        while not complete() and time.monotonic() < deadline:
            windows = await (
                queen.queue(USAGE)
                .group("rate-limiter-gate")
                .subscription_mode("all")
                .batch(50)
                # partitions(10) lets this one call claim up to ten partitions,
                # with batch as the total budget across all of them. Left unset,
                # the broker chooses.
                .partitions(10)
                .wait(True)
                .timeout_millis(2000)
                .pop()
            )

            # This loop pops instead of consuming, so it acks every window
            # itself, and names the group each time: the lease belongs to the
            # group, and an ack without it is refused.
            for w in windows:
                # The window's key is the partition it was computed for.
                key = w["partition"]
                counted[key] = counted.get(key, 0) + w["data"]["requests"]
                over_by = w["data"]["requests"] - QUOTA_PER_WINDOW

                if over_by > 0:
                    # The decision is a message: whatever enforces it (an edge
                    # worker, a gateway, the API itself) reads this queue and
                    # gets the decisions in order, per key. It commits with the
                    # ack of the window it came from, so a crash in between
                    # cannot lose a decision or make two.
                    await (
                        queen.transaction()
                        .queue(THROTTLED)
                        .partition(key)
                        .push(
                            {
                                "data": {
                                    "key": key,
                                    "window": w["data"]["requests"],
                                    "quota": QUOTA_PER_WINDOW,
                                    "overBy": over_by,
                                }
                            }
                        )
                        .ack(w, "completed", {"consumer_group": "rate-limiter-gate"})
                        .commit()
                    )
                    decisions.append({"key": key, "overBy": over_by})
                    print(f"  {key}: {w['data']['requests']} in a window, over by {over_by}")
                else:
                    await queen.ack(w, True, {"group": "rate-limiter-gate"})
                    print(f"  {key}: {w['data']['requests']} in a window, within quota")

        # ----------------------------------------------------------- checking
        print("\nchecking")
        check(complete(), "every request reached a closed window before the deadline")
        check(counted[QUIET_KEY] == QUIET_REQUESTS, "the quiet key was counted exactly")
        check(counted[NOISY_KEY] == NOISY_REQUESTS, "the noisy key was counted exactly")

        check(len(decisions) > 0, "the noisy key was throttled")
        check(
            all(d["key"] == NOISY_KEY for d in decisions),
            "the quiet key was never throttled, so the limiter is not just firing at everything",
        )

        # The decisions are readable by whatever enforces them, in order, per
        # key. No group is named, so this read goes through the queue's own
        # cursor, which starts at the beginning.
        throttled = await queen.queue(THROTTLED).batch(50).partitions(10).wait(True).pop()
        check(
            len(throttled) == len(decisions),
            "every decision is on the queue the gateway reads",
        )
        check(
            all(m["data"]["window"] > m["data"]["quota"] for m in throttled),
            "each decision carries the count and the quota that produced it",
        )

        # stop() cancels the idle-flush timer, waits for the polling loop to
        # finish the cycle it is in, and drains a flush already in flight, so
        # nothing is still writing when the queues go away.
        await stream.stop()
        stream = None

        # Clean up on success only: a failed run leaves the queues on the broker
        # to be looked at.
        for name in (REQUESTS, USAGE, THROTTLED):
            await queen.queue(name).delete()

        verdict = f"\nPASS: {CHECKS} checks"
    except Exception as err:
        verdict, failed = f"\nFAIL: {err}", True
    finally:
        if stream:
            await stream.stop()
        # 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/rate-limiter/main.go"
//
// A rate limiter built from a streaming query.
//
// The textbook rate limiter counts requests per API key in a fixed window,
// usually with a counter in Redis, which is one more system to run and one
// more place where the count can drift from the requests.
//
// Here the counter is a windowed aggregation over the request queue itself.
// Each cycle of the stream commits the window state, the closed windows it
// emits and the ack of the requests it counted as one entry in the broker's
// log, so the count cannot drift from the requests, and it survives a restart
// of this process because the state is in the broker.
//
//	api-requests (one partition per API key)
//	  |-- streaming query: tumbling window, count per key
//	        |-- api-usage  -> the gate: over quota becomes a throttle decision
//	              |-- api-throttled
//
// Run it:
//
//	QUEEN_URL=http://localhost:6632 GOWORK=off go run ./rate-limiter
package main

import (
	"context"
	"fmt"
	"os"
	"strconv"
	"sync/atomic"
	"time"

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

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

var (
	requestsQueue  = "app-go-api-requests-" + runID
	usageQueue     = "app-go-api-usage-" + runID
	throttledQueue = "app-go-api-throttled-" + runID

	// The query id is this streaming query's identity in the broker: its
	// window state is keyed by it, so a restart with the same id resumes the
	// same windows.
	queryID = "app-go-rate-limiter-" + runID
)

const (
	windowSeconds  = 2
	quotaPerWindow = 5
	gateGroup      = "rate-limiter-gate"
	quietKey       = "key-quiet"
	noisyKey       = "key-noisy"
	quietRequests  = 3
	noisyRequests  = 20
)

// Why those numbers make the check deterministic: a window is a slice of time,
// so a burst can land on either side of a boundary. Twenty requests split in
// any way at all leave at least ten on one side, which is over a quota of five,
// so the noisy key is always caught. Three requests cannot reach five however
// they are split, so the quiet key is never caught by accident.

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
}

// stopping is set just before the stream is shut down, and read by the logger
// below.
var stopping atomic.Bool

// streamLogger is what the streaming runner reports through. Stopping the
// runner cancels whatever poll it had in flight, and the pop loop reports that
// cancellation on its way out. That report is the shutdown itself, so it is
// dropped. Everything else is printed, because a query failing to commit its
// windows would otherwise fail this run with no explanation.
type streamLogger struct{}

func (streamLogger) Info(msg string, ctx map[string]interface{}) {}

func (streamLogger) Warn(msg string, ctx map[string]interface{}) {
	fmt.Fprintf(os.Stderr, "  stream warning: %s %v\n", msg, ctx)
}

func (streamLogger) Error(msg string, ctx map[string]interface{}) {
	if stopping.Load() {
		return
	}
	fmt.Fprintf(os.Stderr, "  stream error: %s %v\n", msg, ctx)
}

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, including the streaming runner it
	// starts, 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)

	for _, q := range []string{requestsQueue, usageQueue, throttledQueue} {
		if _, err := client.Queue(q).
			Config(queen.QueueConfig{LeaseTime: 30, RetryLimit: 3}).
			Create().Execute(ctx); err != nil {
			return fmt.Errorf("create %s: %w", q, err)
		}
	}

	// ------------------------------------------------------------- the counter
	//
	// The stream runs in this process, as a consumer group of its own. A new
	// group starts at the tail of the queue unless it asks otherwise, and a
	// request pushed while the stream was still starting would be missed, so it
	// asks for SubscriptionMode "all": every request in the queue is counted.
	//
	// The partition is the aggregation key, so the window state is per API key
	// without a word about keys here: the producer decides, by partition.
	fmt.Println("\nstarting the counter")
	runner, err := streams.
		// AsStreamSource adapts a queue builder to what the streaming engine
		// reads from; To takes the queue builder itself, since a sink is only
		// a name.
		From(client.Queue(requestsQueue).AsStreamSource()).
		WindowTumbling(windowSeconds, streams.WithIdleFlushMs(800)).
		// The extractors receive the payload itself, not the envelope, and as
		// an interface{}: nothing about its shape is checked by the compiler,
		// so the fallback lives in the extractor (see cost below, which counts
		// a request with no cost of its own as one). The field order is passed
		// explicitly after the map because a Go map has no order of its own and
		// that order goes into the query's identity hash: left out, the client
		// falls back to sorting the names, which hashes to a different query
		// than the JavaScript object literal's insertion order.
		Aggregate(map[string]streams.ExtractorFn{
			"requests": func(m interface{}) (float64, error) { return 1, nil },
			"cost":     func(m interface{}) (float64, error) { return cost(m), nil },
		}, "requests", "cost").
		To(client.Queue(usageQueue)).
		Run(ctx, streams.RunOptions{
			QueryID:          queryID,
			URL:              brokerURL,
			BatchSize:        200,
			MaxPartitions:    8,
			MaxWaitMillis:    200,
			SubscriptionMode: queen.SubscriptionModeAll,
			Logger:           streamLogger{},
		})
	if err != nil {
		return fmt.Errorf("start the counter: %w", err)
	}
	// Stop waits for the pop loop and the idle-flush loop to leave, and is
	// idempotent, so it is safe both here as a guard and explicitly below.
	stop := func() {
		stopping.Store(true)
		runner.Stop()
	}
	defer stop()

	// ------------------------------------------------------------- the traffic
	fmt.Println("\ntaking traffic")
	send := func(key string, n int) error {
		for i := 1; i <= n; i++ {
			if _, err := client.Queue(requestsQueue).
				Partition(key).
				Push(map[string]interface{}{
					"key":  key,
					"path": "/v1/things",
					"cost": 1,
					"at":   time.Now().UnixMilli(),
				}).
				Execute(ctx); err != nil {
				return fmt.Errorf("push request for %s: %w", key, err)
			}
		}
		fmt.Printf("  %s: %d requests\n", key, n)
		return nil
	}
	if err := send(quietKey, quietRequests); err != nil {
		return err
	}
	if err := send(noisyKey, noisyRequests); err != nil {
		return err
	}

	// ---------------------------------------------------------------- the gate
	//
	// The enforcement point. It reads each closed window and turns the ones over
	// quota into throttle decisions. It is separate from the counter on purpose:
	// the counting is exact and stays the same, while the policy is yours and
	// changes on its own schedule.
	fmt.Println("\nenforcing")
	type decision struct {
		key    string
		overBy int
	}
	counted := map[string]int{}
	var decisions []decision

	// The loop waits for the totals it expects, with a deadline. Stopping on a
	// quiet period would be a race: a window closes when its timer fires,
	// whatever the reader is doing, and a burst that straddles a boundary
	// arrives as two windows.
	complete := func() bool {
		return counted[quietKey] == quietRequests && counted[noisyKey] == noisyRequests
	}
	deadline := time.Now().Add(30 * time.Second)

	for !complete() && time.Now().Before(deadline) {
		windows, err := client.Queue(usageQueue).
			Group(gateGroup).
			SubscriptionMode(queen.SubscriptionModeAll).
			Batch(50).
			// Each key's windows land in that key's partition. Partitions(10)
			// lets one pop take the windows of both keys, with Batch as the
			// budget they share.
			Partitions(10).
			Wait(true).
			TimeoutMillis(2000).
			Pop(ctx)
		if err != nil {
			return fmt.Errorf("read closed windows: %w", err)
		}

		for _, w := range windows {
			// The window's key is the partition it was computed for.
			key := w.Partition
			requests, ok := w.Data["requests"].(float64)
			if !ok {
				return fmt.Errorf("window on %s has no numeric count", key)
			}
			counted[key] += int(requests)
			overBy := int(requests) - quotaPerWindow

			if overBy > 0 {
				// The decision is a message: whatever enforces it (an edge
				// worker, a gateway, the API itself) reads this queue and gets
				// the decisions in order, per key. It commits with the ack of
				// the window it came from, so a crash in between cannot lose a
				// decision or make two.
				if _, err := client.Transaction().
					Queue(throttledQueue).
					Partition(key).
					Push(map[string]interface{}{
						"key":    key,
						"window": int(requests),
						"quota":  quotaPerWindow,
						"overBy": overBy,
					}).
					Ack(w, "completed", queen.AckOptions{ConsumerGroup: gateGroup}).
					Commit(ctx); err != nil {
					return fmt.Errorf("commit throttle decision: %w", err)
				}
				decisions = append(decisions, decision{key: key, overBy: overBy})
				fmt.Printf("  %s: %d in a window, over by %d\n", key, int(requests), overBy)
			} else {
				// A Pop leaves the ack to the caller, and the ack has to name
				// the consumer group: without it the same windows come back on
				// the next turn and every count is added twice.
				if _, err := client.Ack(ctx, w, true, queen.AckOptions{ConsumerGroup: gateGroup}); err != nil {
					return fmt.Errorf("ack window: %w", err)
				}
				fmt.Printf("  %s: %d in a window, within quota\n", key, int(requests))
			}
		}
	}

	// --------------------------------------------------------------- checking
	fmt.Println("\nchecking")
	if err := assert(complete(), "every request reached a closed window before the deadline"); err != nil {
		return err
	}
	if err := assert(counted[quietKey] == quietRequests, "the quiet key was counted exactly"); err != nil {
		return err
	}
	if err := assert(counted[noisyKey] == noisyRequests, "the noisy key was counted exactly"); err != nil {
		return err
	}

	if err := assert(len(decisions) > 0, "the noisy key was throttled"); err != nil {
		return err
	}
	onlyNoisy := true
	for _, d := range decisions {
		if d.key != noisyKey {
			onlyNoisy = false
		}
	}
	if err := assert(
		onlyNoisy,
		"the quiet key was never throttled, so the limiter is not just firing at everything",
	); err != nil {
		return err
	}

	// The decisions are readable by whatever enforces them, in order, per key.
	// No consumer group is named, so this read goes through the queue's own
	// cursor, and TimeoutMillis bounds the call at two seconds (the default
	// long poll is 30 s).
	throttled, err := client.Queue(throttledQueue).
		Batch(50).
		Partitions(10).
		Wait(true).
		TimeoutMillis(2000).
		Pop(ctx)
	if err != nil {
		return fmt.Errorf("read throttle decisions: %w", err)
	}

	if err := assert(
		len(throttled) == len(decisions),
		"every decision is on the queue the gateway reads",
	); err != nil {
		return err
	}
	carriesCounts := true
	for _, m := range throttled {
		window, okW := m.Data["window"].(float64)
		quota, okQ := m.Data["quota"].(float64)
		if !okW || !okQ || window <= quota {
			carriesCounts = false
		}
	}
	if err := assert(
		carriesCounts,
		"each decision carries the count and the quota that produced it",
	); err != nil {
		return err
	}

	stop()

	// Clean up on success only: a failed run returns before this and leaves the
	// three queues, and the query's window state, on the broker.
	for _, q := range []string{requestsQueue, usageQueue, throttledQueue} {
		if _, err := client.Queue(q).Delete().Execute(ctx); err != nil {
			return fmt.Errorf("delete %s: %w", q, err)
		}
	}

	return nil
}

// cost reads the request's cost out of a payload. Extractors are handed the
// decoded payload as an interface{}, so the shape is checked here at run time,
// and a request that carries no cost counts as one.
func cost(m interface{}) float64 {
	payload, ok := m.(map[string]interface{})
	if !ok {
		return 1
	}
	v, ok := payload["cost"].(float64)
	if !ok {
		return 1
	}
	return v
}
```

### Rust

```rust title="examples/apps/rust/src/bin/rate_limiter.rs"
//
// A rate limiter built from a streaming query.
//
// The textbook rate limiter counts requests per API key in a fixed window,
// usually with a counter in Redis, which is one more system to run and one
// more place where the count can drift from the requests.
//
// Here the counter is a windowed aggregation over the request queue itself.
// Each cycle of the stream commits the window state, the closed windows it
// emits and the ack of the requests it counted as one entry in the broker's
// log, so the count cannot drift from the requests, and it survives a restart
// of this process because the state is in the broker.
//
//   api-requests (one partition per API key)
//     └── streaming query: tumbling window, count per key
//           └── api-usage  -> the gate: over quota becomes a throttle decision
//                 └── api-throttled
//
// Run it:
//   QUEEN_URL=http://localhost:6632 cargo run --bin rate_limiter

use std::collections::HashMap;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};

use queen_mq::streams::{RunOptions, Stream};
use queen_mq::{Config, Queen, QueueOptions, SubscriptionMode};
use serde_json::json;

const WINDOW_SECONDS: i64 = 2;
const QUOTA_PER_WINDOW: i64 = 5;

// Two tenants. One is a well behaved integration, the other is a runaway script
// someone left in a loop.
const QUIET_KEY: &str = "key-quiet";
const NOISY_KEY: &str = "key-noisy";
const QUIET_REQUESTS: i64 = 3;
const NOISY_REQUESTS: i64 = 20;

// Why those numbers make the check deterministic: a window is a slice of time,
// so a burst can land on either side of a boundary. Twenty requests split in
// any way at all leave at least ten on one side, which is over a quota of five,
// so the noisy key is always caught. Three requests cannot reach five however
// they are split, so the quiet key is never caught by accident.

const GATE_GROUP: &str = "rate-limiter-gate";

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 requests = format!("app-rust-api-requests-{run_id}");
    let usage = format!("app-rust-api-usage-{run_id}");
    let throttled = format!("app-rust-api-throttled-{run_id}");

    // The query id is this streaming query's identity in the broker. Its
    // window state is keyed by it, and the runner derives its consumer group
    // from it as `streams.{query_id}`.
    let query_id = format!("app-rust-rate-limiter-{run_id}");

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

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

    for q in [&requests, &usage, &throttled] {
        queen
            .queue(q)
            .configure(QueueOptions {
                lease_time: Some(30),
                retry_limit: Some(3),
                ..Default::default()
            })
            .await
            .map_err(|e| e.to_string())?;
    }

    // ------------------------------------------------------------- the counter
    //
    // The stream runs in this process, as a consumer group of its own. A new
    // group starts at the tail of the queue unless it asks otherwise, and a
    // request pushed while the stream was still starting would be missed, so it
    // asks for SubscriptionMode::All: every request in the queue is counted.
    //
    // The partition is the aggregation key, so the window state is per API key
    // without a word about keys here: the producer decides, by partition.
    //
    // Where the JavaScript client takes one options object for the window and
    // one for the aggregates, this client spells each of them as its own step in
    // the chain: window_tumbling, idle_flush_ms, then one aggregate_* per output
    // field.
    println!("\nstarting the counter");
    let counter = Stream::from(queen.queue(&requests))
        .window_tumbling(WINDOW_SECONDS)
        .idle_flush_ms(800)
        // aggregate_count is the count of records in the window. The extractors
        // receive a Record over the payload itself, not the envelope, so it is
        // r.number("cost") and not the message's `data` field. A missing or
        // non-numeric field yields None, so a request that carries no cost is
        // billed as one.
        .aggregate_count("requests")
        .aggregate_sum("cost", |r| Some(r.number("cost").unwrap_or(1.0)))
        .to(queen.queue(&usage))
        .run(
            &queen,
            RunOptions::new(&query_id)
                .batch_size(200)
                .max_partitions(8)
                .max_wait(Duration::from_millis(200))
                .subscription_mode(SubscriptionMode::All),
        )
        .await
        .map_err(|e| e.to_string())?;

    // ------------------------------------------------------------- the traffic
    println!("\ntaking traffic");
    for (key, n) in [(QUIET_KEY, QUIET_REQUESTS), (NOISY_KEY, NOISY_REQUESTS)] {
        for _ in 0..n {
            let at = SystemTime::now()
                .duration_since(UNIX_EPOCH)
                .unwrap()
                .as_millis() as i64;
            queen
                .queue(&requests)
                .partition(key)
                .push(json!({ "key": key, "path": "/v1/things", "cost": 1, "at": at }))
                .await
                .map_err(|e| e.to_string())?;
        }
        println!("  {key}: {n} requests");
    }

    // ---------------------------------------------------------------- the gate
    //
    // The enforcement point. It reads each closed window and turns the ones over
    // quota into throttle decisions. It is separate from the counter on purpose:
    // the counting is exact and stays the same, while the policy is yours and
    // changes on its own schedule.
    //
    // A window is a slice of time, so a burst can arrive as two windows. That is
    // why this accumulates per key and waits for the totals it expects, with a
    // deadline. Waiting for a quiet period would be a race: the last window
    // closes when its timer fires, whatever the reader is doing.
    println!("\nenforcing");
    let mut counted: HashMap<String, i64> = HashMap::new();
    let mut decisions: Vec<(String, i64)> = Vec::new();
    let complete = |counted: &HashMap<String, i64>| {
        counted.get(QUIET_KEY).copied().unwrap_or(0) == QUIET_REQUESTS
            && counted.get(NOISY_KEY).copied().unwrap_or(0) == NOISY_REQUESTS
    };
    let deadline = Instant::now() + Duration::from_secs(30);

    while !complete(&counted) && Instant::now() < deadline {
        // Each key's windows land in that key's partition. partitions(10) lets
        // one pop take the windows of both keys, with batch as the budget they
        // share.
        let windows = queen
            .queue(&usage)
            .group(GATE_GROUP)
            .subscription_mode(SubscriptionMode::All)
            .batch(50)
            .partitions(10)
            .wait(true)
            .poll_timeout(Duration::from_secs(2))
            .pop()
            .await
            .map_err(|e| e.to_string())?;

        for w in &windows {
            // The window's key is the partition it was computed for.
            let key = w.partition.clone();
            // The aggregates come back as JSON floating-point numbers (the
            // accumulator is an f64 whatever it counted), so `20` arrives as
            // `20.0` and as_i64() on it would be None. Read it as f64 and round.
            let in_window = w.data["requests"].as_f64().unwrap_or(0.0).round() as i64;
            *counted.entry(key.clone()).or_insert(0) += in_window;
            let over_by = in_window - QUOTA_PER_WINDOW;

            if over_by > 0 {
                // The decision is a message: whatever enforces it (an edge
                // worker, a gateway, the API itself) reads this queue and gets
                // the decisions in order, per key. It commits with the ack of
                // the window it came from, so a crash in between cannot lose a
                // decision or make two.
                queen
                    .transaction()
                    .push_to(
                        &throttled,
                        &key,
                        json!({
                            "key": key,
                            "window": in_window,
                            "quota": QUOTA_PER_WINDOW,
                            "overBy": over_by,
                        }),
                    )
                    .map_err(|e| e.to_string())?
                    .ack(w)
                    .commit()
                    .await
                    .map_err(|e| e.to_string())?;
                decisions.push((key.clone(), over_by));
                println!("  {key}: {in_window} in a window, over by {over_by}");
            } else {
                // pop() takes a lease and leaves the ack to the caller. This
                // client reads the consumer group and the lease id off the
                // message, so the ack cannot be pointed at the wrong cursor by
                // forgetting one.
                queen.ack(w).await.map_err(|e| e.to_string())?;
                println!("  {key}: {in_window} in a window, within quota");
            }
        }
    }

    // Stop the runner before checking, so nothing is still writing to the queues
    // the assertions read. stop() waits for the in-flight cycle and its flush,
    // and it consumes the handle: a stopped stream cannot be restarted by
    // mistake.
    counter.stop().await.map_err(|e| e.to_string())?;

    // --------------------------------------------------------------- checking
    println!("\nchecking");
    checks.assert(
        complete(&counted),
        "every request reached a closed window before the deadline",
    )?;
    checks.assert(
        counted.get(QUIET_KEY).copied().unwrap_or(0) == QUIET_REQUESTS,
        "the quiet key was counted exactly",
    )?;
    checks.assert(
        counted.get(NOISY_KEY).copied().unwrap_or(0) == NOISY_REQUESTS,
        "the noisy key was counted exactly",
    )?;

    checks.assert(!decisions.is_empty(), "the noisy key was throttled")?;
    checks.assert(
        decisions.iter().all(|(key, _)| key == NOISY_KEY),
        "the quiet key was never throttled, so the limiter is not just firing at everything",
    )?;

    // The decisions are readable by whatever enforces them, in order, per key.
    let gateway = queen
        .queue(&throttled)
        .batch(50)
        .partitions(10)
        .wait(true)
        .poll_timeout(Duration::from_secs(5))
        .pop()
        .await
        .map_err(|e| e.to_string())?;
    checks.assert(
        gateway.len() == decisions.len(),
        "every decision is on the queue the gateway reads",
    )?;
    checks.assert(
        gateway.iter().all(|m| {
            m.data["window"].as_i64().unwrap_or(0) > m.data["quota"].as_i64().unwrap_or(i64::MAX)
        }),
        "each decision carries the count and the quota that produced it",
    )?;

    // Clean up on success only: a failed run leaves the queues, and the query's
    // window state, on the broker to be looked at.
    for q in [&requests, &usage, &throttled] {
        queen.queue(q).delete().await.map_err(|e| e.to_string())?;
    }

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

    Ok(checks.0)
}
```

Each API key's requests land in a partition named after the key, and a stream counts them: a
tumbling window of two seconds and an aggregate of one per request. The stream runs inside the
program as a consumer group of its own, and each of its cycles commits the window state, the
windows it closed and the ack of the requests it counted as one entry in the log, so the count
can neither lose a request nor count one twice when the process restarts. It starts with
`subscriptionMode: 'all'`, because a new group otherwise begins at the tail and can miss a
request pushed while it is starting. Before the program asked for it, that happened in 3 of 14
runs; after, in none of 20.

**Figure.** The rate limiter's counter. Requests land in the requests queue, in a partition per API key. A stream counts them in tumbling two-second windows, and each of its cycles commits the window state, the windows it closed and the ack of the requests it counted as one entry. A gate consumes the closed windows, and for each window over the quota it pushes a throttle decision to the queue a gateway reads, in the same transaction as the ack of the window.

Counting and policy are separate consumers, each ending in one transaction: the count stays exact through restarts, and the gate's rules change on their own schedule.

- requests: a partition per API key
- counting stream: 2 s windows, count
- closed windows: per key and window
- gate: over the quota?
- throttle decisions: read by a gateway
- requests → counting stream: pop
- counting stream → closed windows: cycle
- closed windows → gate: pop
- gate → throttle decisions: push + ack

The gate is a plain consumer of the closed windows. For each window over the quota it pushes a
throttle decision to the queue a gateway would read, in the same transaction as the ack of the
window, so a crash between the two can neither lose a decision nor make two. Keeping the gate out
of the stream is deliberate: the counting stays exact and fixed, and the policy, which is yours,
changes on its own schedule. [Streams](/guides/streams/) covers the operators, the windows and
what one cycle commits.

## The admission controller

What the admission controller printed on the same node, with a quota of three exports per tenant
per eight-second window:

```text
a duplicate request, and the budget it must not spend
  G-1: admitted
  globex counter: 1
  G-1: already admitted, nothing written
  globex counter: 1
  ok: the duplicate lost the gate
  ok: the duplicate produced no second export
  ok: the rolled-back transaction gave the budget back (counter still 1)

six requests against a quota of 3
  R-1: admitted
  R-2: admitted
  R-3: admitted
  R-4: over quota, moved to the next window in 7975 ms
  R-5: over quota, moved to the next window in 7967 ms
  R-6: over quota, moved to the next window in 7959 ms
  ok: 3 of acme's requests ran now and none was refused
  ok: the other 3 were moved to the next window

withdrawing one deferred request
  ok: R-6 was cancelled while it waited

the next window
  R-5: admitted
  R-4: admitted

the work that ran
  exported R-1
  exported R-2
  exported R-3
  exported R-5
  exported R-4
  exported G-1

checking
  ok: exactly the requests not withdrawn ran, once each (G-1, R-1, R-2, R-3, R-4, R-5)
  ok: every export that ran was admitted by the counter
  ok: the gate refused only the one duplicate, never a returning request
  ok: no timer is left pending

  ran: R-1, R-2, R-3, R-5, R-4, G-1; withdrawn: R-6

PASS: 10 checks
```

Six requests met a quota of three, and all five that were still wanted ran: three at once, two
when the window rolled over eight seconds later, with nobody retrying anything.

### JavaScript

```js title="examples/apps/js/deferred-work.mjs"
//
// A quota that moves work to the next window instead of refusing it.
//
// A limiter that answers 429 above its quota pushes the problem back to the
// caller, and most callers retry, so refusals come back fastest when the
// service is busiest. For work somebody asked for, a report or an export, the
// caller wants the result and does not much mind when it is ready. This
// admission controller never refuses: it asks the counter for room, and when
// there is none it puts the same request on a timer that fires when the window
// opens again.
//
// Two primitives do it. incr with a max is the admission decision: the
// increment applies, or nothing is written and the answer says
// applied: false, reason: 'limit', so a refused request spends no budget. A
// timer carries the request forward: one message, scheduled now, delivered
// when the window opens, and cancellable by name until then.
//
//   requests (one partition per tenant)
//     └── group "admission"   ONE transaction: incr + gate + push + ack
//           ├── work           admitted now
//           └── a timer back onto `requests`, for the next window
//
// Run it:
//   QUEEN_URL=http://localhost:6632 node deferred-work.mjs

import { Queen } from 'queen-mq'

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

const REQUESTS = `app-js-deferred-requests-${RUN}`
const WORK = `app-js-deferred-work-${RUN}`
const NS = `app-js-deferred-${RUN}`

// Three exports per tenant per window. The window is eight seconds here and an
// hour in a real service; nothing else in the program changes with it.
const QUOTA = 3
const WINDOW_S = 8

// Timers fire on the leader's next tick after they are due, so waiting for a
// deferred request needs the window plus a margin.
const TIMER_DEADLINE_MS = (WINDOW_S + 20) * 1000

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

const quotaKey = (tenant) => `quota:${tenant}`
const admittedKey = (requestId) => `admitted:${requestId}`
const timerKey = (requestId) => `req:${requestId}`

const admitted = []
const deferred = []
const alreadyAdmitted = []

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

// One partition per tenant, so a tenant's requests are decided in the order
// they were made.
let submitted = 0
const submit = (tenant, requestId) =>
  queen.queue(REQUESTS).partition(tenant).push({
    transactionId: `submit-${submitted++}-${requestId}`,
    data: { tenant, requestId },
  })

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

  for (const q of [REQUESTS, WORK]) {
    await queen.queue(q).config({ leaseTime: 30, retryLimit: 3 }).create()
  }

  // --------------------------------------------------------------- admission
  //
  // One transaction decides, records and dispatches a request. The counter is
  // incremented BEFORE the duplicate check, so a redelivered request that was
  // already admitted spends budget on its way to being refused. It gets the
  // budget back because both are in one transaction: the lost gate rolls the
  // increment back with everything else.
  const decide = async (msg) => {
    const { tenant, requestId } = msg.data

    const res = await queen
      .transaction()
      .kv.incr(NS, quotaKey(tenant), 1, { max: QUOTA, ttl: `${WINDOW_S}s`, required: true })
      .kv.putIfAbsent(NS, admittedKey(requestId), { tenant }, { ttl: '1h', required: true })
      // No transactionId: the push commits only with the admitted:<id> entry,
      // so that entry is its idempotency key.
      .queue(WORK).partition(tenant).push({ data: { tenant, requestId } })
      .ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
      .commit()

    if (res.success !== false) {
      admitted.push({ tenant, requestId })
      console.log(`  ${requestId}: admitted`)
      return
    }

    // A refused gate comes back as a value (HTTP 200, success: false, reason
    // 'kv_precondition'), and kvReason says which gate refused.
    if (res.kvReason === 'exists') {
      // A request that was admitted before. Nothing was written, the increment
      // included.
      alreadyAdmitted.push(requestId)
      await queen.ack(msg, 'completed', { group: msg.consumerGroup })
      console.log(`  ${requestId}: already admitted, nothing written`)
      return
    }
    if (res.kvReason !== 'limit') throw new Error(`${requestId}: unexpected verdict ${res.kvReason}`)

    // Over the quota. The counter's expiry is the window boundary: incr sets the
    // TTL only when it creates the key, so the counter lives exactly one window,
    // and the wait is read off the entry that refused the request.
    const counter = await queen.kv.get(NS, quotaKey(tenant))
    const delayMs = counter.found ? Math.max(0, new Date(counter.expiresAt).getTime() - Date.now()) : 0

    // The deferral is a transaction too: no timer, no ack, so the request comes
    // back and is decided again.
    const out = await queen
      .transaction()
      .timer(REQUESTS).key(timerKey(requestId)).partition(tenant).delayMs(delayMs)
      .payload({ tenant, requestId }).schedule()
      .ack(msg, 'completed', { consumerGroup: msg.consumerGroup })
      .commit()
    if (out.success === false) throw new Error(`${requestId}: deferral failed (${out.reason})`)

    deferred.push(requestId)
    console.log(`  ${requestId}: over quota, moved to the next window in ${delayMs} ms`)
  }

  // Each phase below ends after `limit` decisions, with the idle bound as the
  // deadline behind it.
  const admission = (limit, idleMillis) => queen
    .queue(REQUESTS)
    .group('admission')
    .subscriptionMode('all')
    .autoAck(false)
    .each()
    .limit(limit)
    .timeoutMillis(1000)
    .idleMillis(idleMillis)
    .consume(decide)

  // --------------------------------------------------- a duplicate request
  console.log('\na duplicate request, and the budget it must not spend')
  await submit('globex', 'G-1')
  await admission(1, 20_000)
  const afterFirst = await queen.kv.get(NS, quotaKey('globex'))
  console.log(`  globex counter: ${afterFirst.value}`)

  // The same request again, as a new message: what a redelivery looks like
  // from the admission controller's side.
  await submit('globex', 'G-1')
  await admission(1, 20_000)
  const afterDuplicate = await queen.kv.get(NS, quotaKey('globex'))
  console.log(`  globex counter: ${afterDuplicate.value}`)

  assert(alreadyAdmitted.join(',') === 'G-1', 'the duplicate lost the gate')
  assert(admitted.length === 1, 'the duplicate produced no second export')
  assert(
    afterDuplicate.value === afterFirst.value,
    `the rolled-back transaction gave the budget back (counter still ${afterDuplicate.value})`
  )

  // ------------------------------------------------------ over the quota
  console.log(`\nsix requests against a quota of ${QUOTA}`)
  const ACME = ['R-1', 'R-2', 'R-3', 'R-4', 'R-5', 'R-6']
  for (const requestId of ACME) await submit('acme', requestId)
  await admission(ACME.length, 20_000)

  const acmeAdmitted = () => admitted.filter(a => a.tenant === 'acme')
  assert(acmeAdmitted().length === QUOTA, `${QUOTA} of acme's requests ran now and none was refused`)
  assert(deferred.length === ACME.length - QUOTA, `the other ${ACME.length - QUOTA} were moved to the next window`)

  // ---------------------------------------------------- withdrawing one
  //
  // A deferred request can be inspected and called off by name while it waits.
  console.log('\nwithdrawing one deferred request')
  const withdrawn = deferred[deferred.length - 1]
  const cancelled = await queen.timer(REQUESTS).key(timerKey(withdrawn)).cancel()
  assert(cancelled.status === 'cancelled', `${withdrawn} was cancelled while it waited`)

  // ------------------------------------------------------ the next window
  //
  // The timers deliver the requests back onto the same queue, and the same
  // admission controller decides them again. A request coming back is nothing
  // special, which is what keeps the loop safe.
  console.log('\nthe next window')
  await admission(deferred.length - 1, TIMER_DEADLINE_MS)
  // A short pass with room for the withdrawn request, to show it never returns.
  await admission(1, 4000)

  // --------------------------------------------------------------- the work
  console.log('\nthe work that ran')
  const ran = []
  await queen
    .queue(WORK)
    .group('exporter')
    .subscriptionMode('all')
    .each()
    .timeoutMillis(1000)
    .idleMillis(3000)
    .consume(async (msg) => {
      ran.push(msg.data.requestId)
      console.log(`  exported ${msg.data.requestId}`)
    })

  // ----------------------------------------------------------------- checking
  console.log('\nchecking')
  const expected = ['G-1', ...ACME.filter(id => id !== withdrawn)].sort()
  assert(
    JSON.stringify([...ran].sort()) === JSON.stringify(expected),
    `exactly the requests not withdrawn ran, once each (${[...ran].sort().join(', ')})`
  )
  assert(
    JSON.stringify(admitted.map(a => a.requestId).sort()) === JSON.stringify(expected),
    'every export that ran was admitted by the counter'
  )
  assert(alreadyAdmitted.length === 1, 'the gate refused only the one duplicate, never a returning request')
  const left = await queen.timer(REQUESTS).list({ limit: 50 })
  assert(left.rows.length === 0, 'no timer is left pending')
  console.log(`\n  ran: ${ran.join(', ')}; withdrawn: ${withdrawn}`)

  for (const q of [REQUESTS, WORK]) await queen.queue(q).delete()
  console.log(`\nPASS: ${checks} checks`)
} catch (err) {
  console.error(`\nFAIL: ${err.message}`)
  process.exitCode = 1
} finally {
  await queen.close()
}
```

One transaction decides a request: an `incr` of the tenant's counter with `max` set to the quota,
a `putIfAbsent` of an `admitted:` entry for the request, the push of the work, and the ack. Both KV
operations carry `required: true`, so either of them can stop the whole transaction, and the
answer's `kvReason` says which one did. `limit` means the tenant is over the quota for this window;
`exists` means this request was admitted before, which is what a redelivery looks like.

**Figure.** The admission controller. Each request is decided by one transaction: an incr of the tenant's counter with max set to the quota, a putIfAbsent of an admitted entry for the request, the push of the work and the ack, both KV operations with required true. If it commits, the request is admitted and its work pushed. If the counter is at the quota (kvReason limit), nothing is written, and a second transaction schedules a timer that pushes the request back at the end of the window, together with the ack. If the admitted entry already exists (kvReason exists), it is a redelivery and is only acked.

A refused request spends no budget, because the incr rolled back with everything else. Deferred work comes back on its own when the window rolls over, with nobody retrying.

- request R-4: tenant acme
- one transaction: incr counter, max 3 putIfAbsent admitted:R-4 push the work, ack
- admitted: the work is pushed
- deferred: kvReason limit: a timer pushes it back later
- already admitted: kvReason exists: ack
- request R-4 → one transaction
- one transaction → admitted
- one transaction → deferred
- one transaction → already admitted

When the quota is spent nothing is written, so the refused request has not used any budget. The
program reads the counter's expiry, which is the end of the window, because `incr` sets a TTL
only when it creates the key, and schedules a timer that pushes the same request back onto the
same partition at that moment. The timer and the ack go in one transaction too, so a request is
either deferred or still on the queue. Until it fires, the timer can be read and cancelled by
name, which is how R-6 was withdrawn.

The duplicate shows the other half. G-1's second submission increments the counter on its way to
losing the `admitted:` gate, and the counter still reads 1 afterwards, because the lost gate
rolled the increment back together with everything else in the transaction.

## Why they are short

A Redis counter and a broker commit separately, so a crash between counting and acknowledging
either counts a request that was never processed or processes one that was never counted, and the
code that tries to repair that is longer than the limiter. A refusing limiter also creates work
somewhere else: every 429 comes back as a retry, fastest when the service is busiest. Here the
count commits with the work it counts, and the deferral is a timer instead of a refusal, so both
programs are mostly the policy.

## Limits

- Windows are fixed. A burst that straddles a boundary is split between two windows, as in the
  run above. A sliding window (`windowSliding`), or one of the ready-made gates described in
  [streams](/guides/streams/), counts across the boundary.
- Stream cycles are not part of the Jepsen pass, and the idle flush that closes a window on a
  quiet partition commits without a lease, so in rare cases a window can be emitted twice. The
  [streams guide](/guides/streams/) has the details.
- Deferred requests compete again. When the window opens, they are decided like new requests,
  with no priority over requests made in the meantime, and requests whose timers come due at the
  same moment come back in no particular order (R-5 before R-4 above).
- Every admitted request leaves a KV entry (`admitted:<id>`) that lives as long as its TTL,
  an hour here. KV inside a transaction is not rate-limited; the counter read on the deferral path
  is a KV call of its own and counts against the tenant's read rate.

## Next

The [booking saga](/examples/saga/) uses the same timer-in-a-transaction for a different kind of
deadline, and [KV](/concepts/kv/) covers `incr`, bounds and expiry.

Source: https://queenmq.com/examples/rate-limiter/index.mdx
