Skip to content

Streams

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.

Updated View as Markdown

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

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.

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.salespartition cust-17runnerfilter, window, aggregateone cycleone transactionsales-per-minutepartition cust-17window stateper source partitionpopPOST /streams/v1/cycleack the batchpushupsert, delete
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. 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:

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

examples/apps/js/rate-limiter.mjsjs
//
// 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()
}
examples/apps/py/rate_limiter.pypython
#
# 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()))
examples/apps/go/rate-limiter/main.gogo
//
// 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
}
examples/apps/rust/src/bin/rate_limiter.rsrust
//
// 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 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 a cycle is built on, and the KV store that holds the same kind of state for your own code.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close