---
title: "Streams"
description: "Window and aggregate a queue in your own process, with the window state kept in Queen and each cycle committing the source ack, the state and the sink push as one entry."
---

> 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

# Streams

A stream is a chain of operators (filter, map, a window, an aggregate, a sink) that runs inside your
own process and reads a queue as a consumer group. The window state is kept in Queen, and each
cycle commits the source ack, the state change and the sink push as one entry, so a crash between
cycles neither loses an event nor counts one twice. There is no stream-processing cluster to run
beside the broker, no changelog topic, and no local state store to rebuild after a restart.

## A windowed total per customer

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

const url = 'http://localhost:6632'
const queen = new Queen({ url })

const handle = await Stream
  .from(queen.queue('sales'))                       // one partition per customer
  .filter((m) => m.data.amount > 0)
  .windowTumbling({ seconds: 60 })
  .aggregate({ count: () => 1, sum: (s) => s.amount })
  .to(queen.queue('sales-per-minute'))
  .run({ queryId: 'sales-per-minute', url })

// later: await handle.stop()
```

The partition is the grouping key. Window state is kept per source partition, so with one
partition per customer each customer gets windows of its own without the code naming a key, and
the result goes to the partition of the same name in the sink unless `.to(q, { partition })` says
otherwise (a name, or a function of the result).

`queryId` is the query's identity, and its state is stored under it, so a process restarted with
the same id resumes the same windows. The runner reads the source as the consumer group
`streams.<queryId>`, which means several processes running the same query share the partitions
between them. Like any new group, a new query starts at the tail of the queue: start it before the
traffic you want counted, or pass `subscriptionMode: 'all'` to `run()` to begin at the oldest
retained message.

Before the window, `filter`, `map` and `eventTime` receive the whole message (`m.data.amount`),
while the extractors of `aggregate` and `reduce` receive the payload, or what the last `map` before
the window returned (`s.amount`). Every `map` before the window sees the original message, so two
of them do not compose; after the window, `map` and `filter` receive the emitted result. In
`aggregate`, `count` adds one per event, `sum` treats a missing field as zero, `min` and `max`
start at `null`, and `avg` skips a value it cannot read.

## What one cycle commits

The runner pops up to `batchSize` messages (200) spread over up to `maxPartitions` partitions (4),
runs the chain for each partition, and posts one cycle per partition to `POST /streams/v1/cycle`.
The broker turns the cycle into one transaction: a positional ack of the source batch, made under
the batch's lease, the state upserts and deletes, and the sink pushes. All of it applies or none
of it does. If the lease expired while the chain ran, the whole cycle is refused and the batch is
processed again by whoever holds it next. If the commit fails for any other reason, the runner does
not nack: the lease runs out, and the batch comes back to the same chain with the same state.

**Figure.** One cycle of a stream. The runner pops a batch of the source queue sales under a lease, for example from partition cust-17, and runs the chain (filter, window, aggregate) on it. It posts one cycle for that partition, which the broker commits as one transaction: the positional ack of the source batch under its lease, the upserts and deletes of the window state, and the pushes of the results to partition cust-17 of the sink queue sales-per-minute.

The window state, the results and the ack of the events they came from are one transaction, under the batch's lease, so the counter cannot drift from the stream it counted.

- sales: partition cust-17
- runner: filter, window, aggregate
- one cycle: one transaction
- sales-per-minute: partition cust-17
- window state: per source partition
- sales → runner: pop
- runner → one cycle: POST /streams/v1/cycle
- one cycle → sales: ack the batch
- one cycle → sales-per-minute: push
- one cycle → window state: upsert, delete

Source: `server/src/rsm/facade/real/phase2/streams.rs, clients/client-js/client-v2/streams/runtime/Runner.js`.

We keep the state in the broker for the same reason we keep KV there. A stream processor that
holds its windows in local storage has to checkpoint them and restore them after a crash, and
needs a changelog to rebuild them elsewhere. Here the window is a row that commits with the ack of
the events it counted, so the counter cannot drift from the stream it was computed from.

## Windows

| Operator | Window | Closes | Idle flush default |
|---|---|---|---|
| `windowTumbling({ seconds })` | Fixed, back to back | At its end plus `gracePeriod` | 5000 ms |
| `windowSliding({ size, slide })` | Overlapping, one every `slide` seconds; `size` is a multiple of `slide` | At each window's end plus grace | 5000 ms |
| `windowSession({ gap })` | Per key, extended while events keep arriving within `gap` seconds | At the last event plus `gap` plus grace | 1000 ms |
| `windowCron({ every })` | Aligned to a `second`, `minute`, `hour`, `day` or `week` (UTC, weeks from Monday) | At the boundary plus grace | 30000 ms |

By default a window runs on processing time, the `createdAt` Queen stamped on each message. A
window on a partition that has gone quiet is closed by the idle flush (`idleFlushMs`, `0` turns it
off), which commits the state delete and the sink push as one entry with no source ack. A session
window needs a reducer after it.

To window on a time carried inside the event, pass `eventTime`:

```js
.windowTumbling({
  seconds: 60,
  eventTime: (m) => m.data.eventTs,   // a Date, epoch ms or a date string
  allowedLateness: 30,                // seconds behind the newest event still accepted
  onLate: 'drop',                     // or 'include'
})
```

Each partition keeps a watermark, the newest event time it has seen minus `allowedLateness`. A
window closes when the watermark passes its end plus grace, and an event older than the watermark
is dropped. With `onLate: 'include'`, the late event opens a fresh accumulator for its window
instead, so the window is emitted a second time with only the late events in it. The watermark is
stored with the window state and survives a restart.

## Other terminals

`.foreach(fn)` runs `fn` for each result and commits the ack after `fn` resolves, so its effect is
at-least-once. When the effect has to happen once, write to a sink queue and act on it from there.

`.gate(fn)` makes an allow or deny decision per message with per-key state, and `tokenBucketGate`
and `slidingWindowGate` are ready-made gates exported next to `Stream`. The first denial stops the
batch: the messages allowed before it are acked, and the denied message comes back, with
everything behind it in order, once the lease expires. A gate cannot share a stream with a window
or a reducer.

## A complete program

A rate limiter built from a stream. Each API key's requests land in its own partition, a 2-second
tumbling window counts them, and a consumer turns every window over the quota into a throttle
decision on another queue. The program checks that the quiet key and the noisy key are counted
exactly, that the noisy one is throttled and that the quiet one never is. It runs in four
languages, because the Python, Go and Rust clients have the same operators and defaults.

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

## Limits

A stream cycle has no retry limit and no dead-letter queue. A message that makes the chain throw
(an extractor that reads `s.a.b` of an event without `a`) comes back after every lease expiry and
holds up its partition until you fix the code or the data, so filter out what the chain cannot
handle before the window. An event whose time `eventTime` cannot parse is reported, skipped and
acked.

The exactly-once property covers the cycle. It does not cover `foreach`, a window emitted again by
`onLate: 'include'`, or the idle flush, which commits without a lease: two processes flushing the
same partition, or a flush retried after an unanswered request, can emit a window twice.

A stream has one window and one reducer. For more stages, chain two streams through a queue.
`.keyBy()` can group by something other than the partition, but state stays per source partition,
so a key spread over several partitions gets a separate window, and a separate partial result, in
each. To aggregate across partitions, push to a queue partitioned by the new key and aggregate
there.

Registering a different operator chain under an existing `queryId` answers `409` unless you pass
`reset: true`, which deletes every state row of the query while the group keeps its position, so
events already counted into the dropped windows are gone from the results. The comparison covers
the operators, their order and their options; the bodies of your functions are outside it, so a
changed `map` body goes unnoticed.

On a partition that stops receiving events, an event-time window stays open until a newer event
moves the watermark, because in event time the idle flush follows the watermark. `windowCron`
takes `every`, not cron syntax.

The cycle route has its own workload in the Jepsen suite (W9 in `test/jepsen`), which checks that
every acknowledged push lands in exactly one sink record while nodes are killed, partitioned and
paused. It has run since 2.0.0-alpha.6; [Jepsen](/benchmarks/jepsen/) lists the campaigns, including
one run under bridge partitions that ended unknown on a liveness problem that is still open. The
runners' windowing, watermarks and idle flush are covered by the client test suites and have no
fault-injection workload, and no published benchmark measures stream throughput.

Next: the [transaction](/concepts/transactions/) a cycle is built on, and the
[KV store](/concepts/kv/) that holds the same kind of state for your own code.

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