---
title: "S3 sink"
description: "Running queen-s3 beside a broker or a proxy, or embedded in the broker: every environment variable it reads, the API key scopes, the bucket policy, the pg_read_all_stats requirement behind an undegraded safeTime, and the lag SLO to alarm on."
---

> 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

# S3 sink

`queen-s3` is a separate binary that reads a queue's log through Queen's own HTTP API and writes
it into an object store in an open format. It is a client of Queen like any SDK is: it holds no
database connection, embeds no broker, and stores nothing durable of its own except three small
documents per queue in Queen's key/value store (`s3:<sink>:<queue>:{intent,committed,lease}`) and
what it puts in the bucket. What it writes, and what a reader gets, is
[the lake reference](/reference/s3).

> **Note**
>
> The connector ships **in the broker image**. `ghcr.io/queen-mq/queen` carries `bin/queen-s3`
> beside the broker binary and the two wire facades from 1.5.0 on, so pulling the broker pulls the
> sink and there is no second artifact to track. `QUEEN_S3_EMBEDDED=true` is the only thing that
> starts it inside the broker: left unset, nothing in the supervisor is read and the image behaves
> exactly as it did before.

It is **not a wire-protocol facade**. Nothing connects to it, and it advertises no address. The
two ports it opens are `/healthz` and `/metrics`, on loopback by default.

## The two shapes

**Beside a broker (self-hosted).** The sink runs anywhere that can reach a broker, points
`QUEEN_URL` at it, and carries one `QUEEN_TOKEN`. It reads through `POST /api/v1/fetch` and
`POST /api/v1/partitions/changed`, and keeps its commit pointers in `POST /api/v1/kv`. Nothing
else on the broker changes: no push path, no pop path, no retention, no table.

**Beside a proxy (Cloud).** The same binary with `QUEEN_URL` pointed at a cell's proxy. Every
call then crosses the proxy's authentication, tenant scoping, quotas and metering exactly as an
SDK's call does, and the sink is an ordinary tenant process reading its own queues. See
[In Queen Cloud](#in-queen-cloud) for the two routes and the key/value carve-out that make that
work.

**Embedded in the broker.** `QUEEN_S3_EMBEDDED=true` makes the broker spawn and supervise the
sink as a child process, wired to the broker's own listener over loopback. One deployment, two
processes, one image. That is [its own section](#embedded-in-the-broker).

The choice between the first and the third is not about capability, it is about where the process
lives. A sink is a client and needs nothing from the broker's container, so running it separately
is what lets it be scaled, restarted and upgraded on its own schedule. Embedded mode exists so a
self-hoster does not have to operate a second thing.

## Build and run

The crate is standalone, so it builds on its own manifest and produces one binary.

```bash
cargo build --release --manifest-path connectors/queen-s3/Cargo.toml
QUEEN_URL=http://localhost:6632 \
QUEEN_S3_QUEUES=orders \
QUEEN_S3_ENDPOINT=https://s3.eu-central-1.amazonaws.com \
QUEEN_S3_REGION=eu-central-1 \
QUEEN_S3_BUCKET=my-lake \
QUEEN_S3_ACCESS_KEY=AKIAEXAMPLE \
QUEEN_S3_SECRET_KEY=a-secret-of-your-own \
  ./connectors/queen-s3/target/release/queen-s3
```

Two rules hold for every variable below. A value that does not parse is a fatal boot and never a
silent fall back to the default, because the default is not there to paper over a typo. And a
variable that is present but empty or whitespace-only counts as unset, because an empty value in a
Compose file or a Helm template is a variable somebody meant to set.

### Boot, and the four exit codes

Boot is short and it is all before the first byte is read: the configuration, the health listener,
one client for Queen and one for S3, then **two probes, both before any queue is claimed**. The
first is an empty `POST /api/v1/partitions/changed`, which answers `safeTime` alone, writes
nothing, and is classed `Consume` at the proxy exactly like the fetch and the commit, so it proves
`QUEEN_URL` and `QUEEN_TOKEN` with the scope the sink needs anyway ([see below](#the-api-key-and-its-scopes)); the
second is a `HEAD` of the prefix, where a `404` is a fine answer (nothing has been written yet) and
a `403` or a connection failure is the credential or the endpoint being wrong. Both failures print
what to check, not just the error.

| Exit | What happened | Worth restarting? |
| --- | --- | --- |
| `0` | A signal arrived and every queue drained | n/a |
| `1` | The health listener could not bind `QUEEN_S3_LISTEN` | maybe: something else has the port |
| `2` | The configuration is wrong; the line names the variable and the accepted values | no, until it is edited |
| `3` | Queen or the bucket could not be reached **at boot** | maybe: URL, credential or network |

Code `3` is deliberately separate from the running failure policy, which is to lag rather than to
die: once the sink is up, an unreachable broker or bucket is a retry behind a backoff and a red
`/healthz`. At boot the same symptom is almost always a wrong URL, a wrong key or a missing
network, and a container that exits says so louder than one that sits there logging.

Then the process is a task per queue, and `main` owns only three things around them. **The queue
set**: a named list is fixed for the life of the process, while `QUEEN_S3_QUEUES=*` is re-listed
every ten discovery intervals (never faster than once a second), so a queue created at ten in the
morning is being mirrored a minute later. Queues are never removed from the set: a task whose queue
disappeared meets `UNKNOWN_TOPIC_OR_PARTITION`, stops, and retries once a minute, which is also
what should happen when the queue comes back. **Ownership**: one lease per queue, retried after its
TTL when another instance holds it. **The signals**: SIGTERM and Ctrl-C stop new reads, let the
window in flight finish, and exit `0`. That is the same drain rule the embedded supervisor's
grace is sized for, and it is [described in full below](#what-a-stop-guarantees-and-what-it-does-not).

## Every variable

The four with no default are the destination: `QUEEN_S3_ENDPOINT`, `QUEEN_S3_REGION`,
`QUEEN_S3_BUCKET` and `QUEEN_S3_QUEUES`, plus the keypair. There is no bucket that is right more
often than it is wrong, so there is no default for one.

| Variable | Default | What it is |
| --- | --- | --- |
| `QUEEN_URL` | `http://localhost:6632` | The broker or proxy the sink calls. Must be an `http`/`https` URL; a trailing slash is normalised away. Checked at boot. |
| `QUEEN_TOKEN` | none | Bearer token for that broker. Never logged: the boot line reports whether there is one, not what it is. [Which scopes it needs](#the-api-key-and-its-scopes). |
| `QUEEN_S3_SINK` | `default` | The sink name. It is a path segment in the bucket **and** a segment of every key/value key, so two sinks on one queue are two independent lakes that cannot collide. 1 to 64 characters of `[A-Za-z0-9._-]`, because anything else would be a sink writing into another sink's prefix. |
| `QUEEN_S3_QUEUES` | **required** | Comma-separated queue names, or `*` for every queue of the tenant, re-listed every **ten** discovery intervals so a queue created after boot gets a task. A sink with no queues has nothing to read, so this is a fatal boot rather than an idle process. [The scopes section](#the-api-key-and-its-scopes) has what `*` costs a token. |
| `QUEEN_S3_PARTITIONS` | unset | A static partition list, `queue:0..1023`. **Refused at boot in this release**, and the message says why: a window may close only at or below `safeTime`, `safeTime` is answered by the discovery call, and a sink naming its own lanes never makes one, so it would have no boundary it is allowed to close at. Leave it unset; discovery finds the lanes at any cardinality. |
| `QUEEN_S3_ENDPOINT` | **required** | The S3 API base URL, `scheme://host[:port]` with **no path**. A path here would silently become part of every object key; the field for that is `QUEEN_S3_PREFIX`. |
| `QUEEN_S3_REGION` | **required** | The region label the SigV4 credential scope is signed with. On AWS it is the real region; on a gateway it is whatever label that gateway accepts, commonly `us-east-1`. |
| `QUEEN_S3_BUCKET` | **required** | The destination bucket: one name, no slash, no space. |
| `QUEEN_S3_PREFIX` | `queen` | The root every object is written under, sidecars included. No leading or trailing slash, no empty or relative segment. One bucket can hold several sinks side by side under different prefixes. |
| `QUEEN_S3_ACCESS_KEY` | **required** | S3 access key id. |
| `QUEEN_S3_SECRET_KEY` | **required** | S3 secret access key. Never logged, never rendered in an error, never in the boot line. |
| `QUEEN_S3_PATH_STYLE` | `false` | `true` addresses the bucket as the first path segment rather than as a DNS label, which is what versitygw and MinIO-shaped hosts want. |
| `QUEEN_S3_SSE` | unset | `AES256` or `aws:kms`, sent as `x-amz-server-side-encryption` on every PUT. Unset sends no header at all, and the bucket's own default still applies. |
| `QUEEN_S3_SSE_KMS_KEY_ID` | unset | The KMS key id, only with `QUEEN_S3_SSE=aws:kms`. Setting it beside `AES256` is a fatal boot rather than a silently ignored policy: the objects would be encrypted with the bucket's own key while the deployment looked correct. The boot line prints the last four characters and no more. |
| `QUEEN_S3_FORMAT` | `jsonl` | `jsonl` (one JSON object per line, every reader takes it) or `parquet`. |
| `QUEEN_S3_COMPRESSION` | `zstd` | `zstd`, `gzip` or `none`, for `jsonl` objects only. This is the one format choice a reader can refuse: Spark's JSONL path goes through Hadoop's `ZStandardCodec` and needs a native library most images do not carry, and pandas needs the `zstandard` package. Both take `gzip`. [The measured matrix](/reference/s3#reader-recipes) is which reader takes what. |
| `QUEEN_S3_PARQUET_CODEC` | `zstd` | `zstd` or `snappy`, for `parquet` objects. A separate variable because a Parquet file's codec lives inside the file, per column chunk, and is not the same decision as compressing a text object. |
| `QUEEN_S3_LAYOUT` | `merged` | `merged` (one object per window, the partition is a column) or `per-partition` (one object per window per partition). `merged` is the only shape that survives a million lanes. |
| `QUEEN_S3_ALIGN` | `hour` | `hour`, `day` or `none`. The Hive bucket a window may not straddle, which is what makes `dt=` and `hour=` exact for every record in an object. |
| `QUEEN_S3_START` | `latest` | Where a queue with **no committed pointer** starts: `latest` (at the current `safeTime`) or `earliest` (backfill everything retention still holds). Read the caution under [scaling out](#scaling-out) before writing `earliest`. |
| `QUEEN_S3_TARGET_MB` | `128` | Close a window at this many **uncompressed** buffered megabytes. In 1 to 5120. Uncompressed on purpose, so the window size does not depend on how well a particular window happens to compress. |
| `QUEEN_S3_MAX_WINDOW_MS` | `300000` | Close a window at this age, whichever comes first. In 100 to 86,400,000. The lag SLO is about this plus the safe lag. |
| `QUEEN_S3_CHECKPOINT_EVERY` | `20` | Windows between position checkpoints. In 1 to 100,000. It bounds the re-read after a restart and nothing else: positions are a cache. |
| `QUEEN_S3_MEMORY_MB` | `1024` | The buffer budget across every queue this process owns. In 1 to 1,048,576. |
| `QUEEN_S3_FETCH_CONCURRENCY` | `4` | In-flight fetch calls per queue. In 1 to 256. Every one of them spends the broker's **pop lane** admission budget, so this is the throttle a backfill is held back with. |
| `QUEEN_S3_DISCOVERY_INTERVAL_MS` | `2000` | How often an idle queue asks which partitions moved. In 10 to 3,600,000. |
| `QUEEN_S3_SAFE_GUARD_MS` | `5000` | Subtracted from the broker's `safeTime` before a window may close. In 0 to 3,600,000. It is **added** to the broker's own five second guard and never subtracted from it. |
| `QUEEN_S3_LEASE_TTL_MS` | `30000` | How long a queue lease survives without a refresh. In 1000 to 3,600,000. |
| `QUEEN_S3_MULTIPART_THRESHOLD_MB` | `64` | Objects at or below this go up as a single PUT with a `Content-MD5`; above it they go up as a multipart upload, 16 MiB per part. In 5 to 5120. |
| `QUEEN_S3_LISTEN` | `127.0.0.1:9333` | Where `/healthz` and `/metrics` are served. The default is loopback, so it is **not scrapable from another pod**: set `0.0.0.0:9333` to scrape it. |
| `QUEEN_S3_INSTANCE` | the hostname | Lease identity. A generated identity is reported as such in the boot line, because a lease held by a name that changes on every restart is never handed back, only expired. |
| `QUEEN_S3_CRASH_AT` | `never` | A fault injection point for the crash matrix: `after_intent`, `mid_upload`, `after_upload`, `before_commit`, `after_commit`. It is configuration and not a test hook because the matrix kills a real process in a real container. Unset it anywhere that is not a test. |
| `QUEEN_S3_LOG_FORMAT` | unset | `json` switches the log format to the structured form a log pipeline wants. Anything else is the human one. |
| `RUST_LOG` | `warn,queen_s3=info` | Tracing filter, `EnvFilter` syntax. The same default the broker and both facades use, so one mental model covers four binaries. |

`QUEEN_S3_BIN`, `QUEEN_S3_EMBEDDED` and `QUEEN_S3_SHUTDOWN_GRACE_MS` are read by the **broker**,
not by this process, and are [in the table below](#embedded-in-the-broker).

## The API key and its scopes

The sink makes three or four kinds of call, and the credential has to cover all of them. Against
a **proxy** the answer is short: **`consume`, and `consume` alone**, unless the queue list is `*`.

| Call | Route | Class at the proxy | Scope |
| --- | --- | --- | --- |
| Read the log | `POST /api/v1/fetch` | `Consume` | `consume` |
| Find the partitions | `POST /api/v1/partitions/changed` | `Consume` | `consume` |
| Commit a window | `POST /api/v1/kv`, keys under `s3:` | `Consume`, by the carve-out below | `consume` |
| Prove the broker is reachable, at boot | `POST /api/v1/partitions/changed`, empty batch | `Consume` | `consume` |
| Resolve `QUEEN_S3_QUEUES=*` | `GET /api/v1/resources/queues` | `Read` | `read` |

**The last row is the one that catches people.** A named queue list never touches the queue
listing, so a consume-only key runs a named sink end to end. `*` resolves through
`GET /api/v1/resources/queues`, which the proxy classifies as a read route, and for an API key
`Read` means `read` or `admin`: a consume-only key is answered `403` there and the process
**exits 3 at boot** with the "cannot reach Queen" line, before it has claimed a queue or touched
the bucket. So `*` costs a token the `read` scope on top of `consume`, and the failure is nowhere
near the bucket. Against a broker rather than a proxy, `read-only` already covers the same call.

Against a **broker with `JWT_ENABLED=true`** the levels are the broker's own rather than the
proxy's, and the binding one is the key/value write: `POST /api/v1/kv` is `read-write`, while both
reads are `read-only`. A `read-write` token therefore covers the sink and a `read-only` token
covers only two thirds of it, which is a sink that reads a whole window and can never commit it.

With authentication off the broker verifies nothing and `QUEEN_TOKEN` is an identity rather than a
credential.

## The S3 side

Any S3 API works: AWS, DigitalOcean Spaces, Cloudflare R2, and versitygw, which is what the tests
run against. The sink deliberately uses **no conditional PUT, no LIST on the correctness path and
no read-after-write assumption**, because a retried upload of a window is byte-identical and
overwriting a key with the same bytes is the whole idempotency story. That is what makes it work
on a gateway whose LIST is eventually consistent and on a service that has no conditional PUT to
offer.

Integrity is checked both ways the protocol offers. A single PUT sends `Content-MD5` and compares
the answer's `ETag` against the MD5 it computed; a completed multipart upload's `ETag` carries a
`-<parts>` suffix, and the part count is compared with the number of parts that were sent. Two
shapes are not the ones integrity is defined over, and both are reported at debug rather than
failed: an `ETag` under SSE-KMS is opaque by an AWS rule, and some gateways answer a multipart
`ETag` with no suffix at all.

**Set an encryption policy on the bucket.** `QUEEN_S3_SSE=AES256` puts the header on every PUT,
and `aws:kms` with `QUEEN_S3_SSE_KMS_KEY_ID` makes the KMS key policy the real access control.
Say plainly what that is protecting: **the bucket holds a decrypted copy of the log.** A fetch
decrypts a payload when the broker has an encryption key configured, so a queue that is encrypted
at rest inside Queen is plaintext in the lake, and at-rest encryption on the bucket is what
restores the property rather than an extra.

### The bucket policy, least privilege

Five verbs and no more, and the fifth is narrower than it looks: the sink **never deletes a data
object**, only its own checkpoints.

```json
{
  "Version": "2012-10-17",
  "Statement": [
    { "Sid": "ListPrefix",
      "Effect": "Allow",
      "Action": ["s3:ListBucket"],
      "Resource": "arn:aws:s3:::my-lake",
      "Condition": { "StringLike": { "s3:prefix": ["queen/*"] } } },
    { "Sid": "WriteAndRead",
      "Effect": "Allow",
      "Action": ["s3:PutObject", "s3:GetObject", "s3:AbortMultipartUpload"],
      "Resource": "arn:aws:s3:::my-lake/queen/*" },
    { "Sid": "GarbageCollectCheckpointsOnly",
      "Effect": "Allow",
      "Action": ["s3:DeleteObject"],
      "Resource": "arn:aws:s3:::my-lake/queen/_queen/*/checkpoint/*" }
  ]
}
```

`queen/` is `QUEEN_S3_PREFIX`. Add an `AbortIncompleteMultipartUpload` lifecycle rule while you
are there: the sink aborts its own multipart uploads, but a process killed between two parts
leaves one behind, and a lifecycle rule is what collects those without anybody watching.

## Grant `pg_read_all_stats` to the broker role

This is the one requirement that lives on the **database** rather than on the sink, and it is
worth doing before the first window closes.

A window may only close at or below `safeTime`, a watermark the broker computes as the oldest
`xact_start` among in-transaction sessions: below it, nothing new can still become visible, which
is what makes a window a deterministic set. PostgreSQL **masks `pg_stat_activity` across roles**.
For a session owned by a role the reader is not a member of, and without `pg_read_all_stats`,
`state` and `xact_start` come back NULL, and a masked row contributes nothing to a minimum while
being exactly the kind of session that could be holding a write transaction open. Taking the
minimum over the visible rows would overestimate the safe instant, which is the one direction that
loses data.

So masked rows are counted, and one is enough: the broker falls back to
`now() - QUEEN_FETCH_SAFE_FLOOR_MS` (default 30 seconds) and answers `safeTimeDegraded: true`.

**Measured, not predicted:** on PostgreSQL 16 `backend_type` is masked too, and PostgreSQL's own
background workers belong to the bootstrap superuser, so for any non-superuser broker role without
`pg_read_all_stats` the degrade is **permanent rather than occasional**.

```sql
GRANT pg_read_all_stats TO queen;
```

Degraded is usable. It costs lag and nothing else: every window closes 30 seconds behind instead
of five, the sink logs the state once and never stalls on it, and no record is lost or duplicated
either way. The grant is what buys the 25 seconds back.

> **Note**
>
> `QUEEN_FETCH_SAFE_FLOOR_MS` is read by the handler directly rather than through the broker's
> `Config`, so it does not appear in [the generated configuration table](/reference/config/). It
> must exceed any write statement's timeout: a floor shorter than that could declare an instant
> settled while a segment is still able to commit below it. Raise it, never lower it.

## The lag SLO, and the retention rule

**`queen_s3_lag_seconds{queue}` is the number to alarm on.** It is `now` minus the end of the last
committed window, so in a healthy sink it oscillates below `QUEEN_S3_MAX_WINDOW_MS` plus the safe
lag: about five and a half minutes at the defaults, or six with a degraded `safeTime`. Alarm at a
multiple of that, not at the number itself.

It is worth alarming on because of what sits behind it. If retention deletes a segment before the
sink read it, that is **the one failure that is data loss**, and the lag gauge is the thing that
was rising for hours beforehand. The loss itself is never silent: the offset gap goes into the
window's manifest as `lost`, counts on `queen_s3_records_lost_total{queue}`, gets one sampled log
line, and the sink resumes from the new log start and **keeps committing**, because a stalled sink
loses more than a lagging one.

The rule that prevents it: **a queue's `retentionSeconds` must comfortably exceed
`QUEEN_S3_MAX_WINDOW_MS` plus the longest outage the sink is expected to survive.** Five minutes
of window against a day of retention is not a decision anybody has to think about; five minutes of
window against ten minutes of retention is one bad deploy away from a hole.

Or make it structural. A queue can name the sink that has to keep up with it, in two options on
`POST /api/v1/configure`:

| Queue option | Default | What it does |
| --- | --- | --- |
| `retentionSinkHold` | `""`, off | The name of a sink, matching `[A-Za-z0-9._-]{0,64}`. Retention will not delete a segment that sink has not committed yet. A name outside the charset is **rejected**, not clamped: the configure answers an error and the queue keeps its old options, because the name is a segment of the sink's own commit key. |
| `retentionSinkHoldMaxSeconds` | `604800` | The ceiling on that hold, seven days. Must be in 60 to 31,536,000, and out of range is rejected the same way. |

The floor retention applies is
`GREATEST(committed window end - 60 s, now() - retentionSinkHoldMaxSeconds)`, and both halves earn
their place. The minute of slack covers the sub-second skew of a watermark-derived bound, paid in
kept bytes and never in lost ones. The cap is what made the option shippable at all: without it a
stopped sink is unbounded retention and the first symptom is a full disk. It also means the hold is
**safe from the moment it is set**, before the sink has ever committed, because a queue with no
pointer yet keeps everything younger than the cap.

Two limits worth knowing before you rely on it. The hold floors the two cutoffs that delete
**segments**, and deliberately not the dedup-hash purge or the max-wait eviction: those delete by
different rules, and a copy in a lake says nothing about when a dedup hash may go or about a
message too old to be worth delivering. And the pointer read degrades rather than fails: if it
cannot be read, every held queue falls back to its cap-only floor, because that value is
tenant-writable and a malformed one must never be able to stop deletion cluster-wide.

## Embedded in the broker

`QUEEN_S3_EMBEDDED=true` makes the broker spawn and supervise the sink as a child process, wired
to the broker's own listener over loopback and derived from the address that listener actually
bound, so `PORT=0` and a wildcard bind both come out right.

It is a child and not a library on purpose. The sink has its own accept loop, its own compression
and Parquet arenas, its own S3 client and its own crash modes, and the broker is built with
`panic = "abort"`, so an allocation failure in a writer sized from `QUEEN_S3_TARGET_MB` would take
the broker down rather than one task. A child keeps the blast radius at the sink: it dies, the
broker keeps serving, and the supervisor brings it back on a backoff that doubles from one second
to a thirty second ceiling and resets after an hour of healthy running. A restart costs lag and
never a gap, because the sink resumes from its committed pointer.

These three are read by the **broker**, not by the sink:

| Variable | Default | What it is |
| --- | --- | --- |
| `QUEEN_S3_EMBEDDED` | `false` | `true` spawns the sink. Off, nothing in the supervisor is read and the broker behaves exactly as it always did, which is why the binary can ship in the default image without changing what the default image does. |
| `QUEEN_S3_BIN` | the `queen-s3` file beside the broker executable | Where the binary is. The default is resolved from the running executable's own directory at boot rather than written down as a path, and it is what makes the image work with no configuration at all. |
| `QUEEN_S3_SHUTDOWN_GRACE_MS` | `30000` | How long a stopping child has between SIGTERM and SIGKILL. Floored at 100 ms, because a grace of zero is a SIGKILL with extra steps. **Thirty seconds and not the facades' five**, and that difference is the point: a stopping facade is closing sockets, a stopping sink has an open window to finish. |

Every other `QUEEN_S3_*` variable means exactly what it means when the sink runs alone. The
environment forwards verbatim, `QUEEN_TOKEN` and the S3 keypair included, because those are the
child's own credentials for the two hops it makes. Four of the broker's secrets are stripped
before exec and the sink reads none of them: `PG_PASSWORD`, `JWT_SECRET`, `QUEEN_ENCRYPTION_KEY`
and `QUEEN_SYNC_SECRET`. A database password has no business being readable in a second process's
environment just because that process happens to be colocated.

Three mistakes are refused at boot rather than discovered as a crash loop: a `QUEEN_S3_BIN` that is
not a file, a binary that exists but has lost its execute bit, and `QUEEN_S3_EMBEDDED=true` with no
`QUEEN_S3_BUCKET`. The third names all four required variables in its message, because nobody
arrives at a bucket name without an endpoint, so the fix is one edit. A fourth is a warning and not
a refusal: `JWT_ENABLED=true` with no `QUEEN_TOKEN`, where every call the child makes is answered
`401` and nothing is ever written. The loopback hop is authenticated like any other client's, and
embedded mode gets no private door into the broker, because a private door is exactly the kind of
thing that is later found open from somewhere else.

The loopback is the **default**, not the only answer. An explicitly set, non-empty `QUEEN_URL` in
the broker's own environment wins, and the child is given that instead, which is what makes
embedded mode usable in a cell where the sink has to reach the broker through the proxy. The boot
line says which branch was taken:

```text
queen-s3 sink started (embedded) pid=18847 bin=/app/bin/queen-s3
  queen_url=http://queen-proxy:6711 queen_url_from="QUEEN_URL (explicit)"
```

`queen_url_from` is either `QUEEN_URL (explicit)` or `loopback (bound listener)`.

The child's stdout and stderr are forwarded into the broker's log, stripped of escape sequences and
control bytes, truncated at 4 KiB a line, and budgeted at 200 lines per ten second window with a
`suppressed` count when the budget is spent. `GET /status` grows an `s3` block when embedded mode
is on, and only then: `mode`, `phase`, `pid`, `restarts`, `lastExit`, `uptimeMs` and `backoffMs`.

### What a stop guarantees, and what it does not

The child is put in its own process group, so a stop signals the whole group and a grandchild dies
with it rather than being re-parented.

- **The broker gets SIGTERM or Ctrl-C**: the supervisor runs after the broker's serve loop drains,
  sends SIGTERM to the group, waits up to `QUEEN_S3_SHUTDOWN_GRACE_MS`, and escalates to SIGKILL.
  A broker stopping with a sink attached therefore takes up to half a minute longer than one
  stopping with a facade. That is not overhead, it is the window being finished.
- **The broker panics or drops the handle**: the child is reaped with it.
- **The broker is itself SIGKILLed**: on **Linux** the child carries `PR_SET_PDEATHSIG`, so the
  kernel kills it whatever the parent died of. On **macOS and BSD** there is no equivalent: a
  SIGKILLed broker leaves the sink running, re-parented to init and still writing to the bucket.
  A development machine caveat rather than a production one, stated rather than papered over.

What the grace buys is bounded and worth being exact about. On SIGTERM the sink stops reading and
decides once, per queue, whether the window it is filling is worth finishing: it is if at least
1 MiB is buffered, or if the window has been open at least ten seconds on **the broker's clock**
and holds anything at all. A window worth finishing is closed, uploaded and committed inside the
grace. Then: **either that commit lands, or the next start redoes the window from its intent and
writes the identical object under the identical key.** Nothing is lost in the second case; the work
is repeated. Cutting the grace to a facade's five seconds does not risk the lake, it just makes the
repeat more likely.

## The image

The repository's `Dockerfile` builds the sink in a stage of its own and copies it next to the
broker binary in `/app/bin`. That adjacency is the contract: with `QUEEN_S3_EMBEDDED=true` and no
`QUEEN_S3_BIN`, the supervisor resolves the child out of its own executable's directory, so
embedded mode needs no extra configuration in that image. The same image runs the sink alone
instead, with `docker run … queen-mq ./bin/queen-s3`, which is the shape that scales out, because
a sink needs nothing from the broker's container.

The connector has a CI job from its first commit, running its unit tests, `rustfmt` and `clippy`
on its own manifest. That is deliberate and it is a correction: the SQS facade shipped in 1.4.0
with 715 tests and no job at all, so for a release those tests only ever ran where somebody
happened to run them. A suite CI never executes is a suite whose green is a memory.

## Scaling out

**Windows are per queue, so one queue's throughput is one instance's.** Scaling out means
splitting the queue list across processes, not adding replicas over the same one.

Instances that do overlap are safe and wasteful rather than dangerous. A queue is owned through a
lease in Queen's key/value store, refreshed on a tick and expiring after `QUEEN_S3_LEASE_TTL_MS`,
and every intent and commit carries that lease as a required precondition. Two instances racing
for one queue means the loser's compare-and-set fails, its batch rolls back whole, and
`queen_s3_commit_precondition_lost_total` counts it. Two instances can never commit two different
versions of the same window.

The practical consequence for Kubernetes: enabling the sink on a StatefulSet with more than one
replica gives every pod the same queue list. Run it on a single-replica release, or split the
queues across releases.

> **Caution**
>
> **`QUEEN_S3_START=earliest` reads everything retention still holds.** On a queue with a week of
> retention that is a week of log through the fetch path, at full speed, and an object-store bill to
> match. Every one of those fetches spends the broker's pop lane admission budget, so a backfill
> competes with real consumers for admission; `QUEEN_S3_FETCH_CONCURRENCY` is the throttle. The
> default, `latest`, is the one that surprises nobody, and a backfill is a decision to take
> deliberately with the concurrency turned down.

## `/healthz` and `/metrics`

Both are served on `QUEEN_S3_LISTEN`, which defaults to loopback.

`/healthz` answers one question, and it is the failure policy rather than a subsystem tree: **the
sink never drops, it only lags.** It is `200` while every queue this process owns has committed a
window inside `3 × QUEEN_S3_MAX_WINDOW_MS` (floored at 30 seconds), and `503` when one has not,
naming the first offender by name so the line is stable across scrapes:

```json
{"ok":false,"queue":"orders","staleMs":1200000,"limitMs":900000}
```

Three windows and not one: a window closes at `MAX_WINDOW_MS`, then has to be uploaded and
committed, and the broker's own `safeTime` lags by its guard, so a one-window threshold would go
red on a healthy sink under load, and a probe that flaps is a probe an operator turns off. A queue
that has **never** committed is green, which is the honest answer for a lane nobody has pushed to
yet and for the seconds after boot.

An unreachable broker, a bucket refusing every PUT, and a lease lost to another instance all arrive
here as the same symptom, because all three have the same consequence.

`/metrics` is Prometheus text exposition, version 0.0.4:

| Metric | Type | Labels | What it is |
| --- | --- | --- | --- |
| `queen_s3_lag_seconds` | gauge | `queue` | `now` minus the last committed window's end. **The SLO.** |
| `queen_s3_safe_lag_seconds` | gauge | none | `now` minus `safeTime`: how far behind the broker's own visibility floor is. A long read-only transaction moves this, and it is latency rather than loss. |
| `queen_s3_windows_committed_total` | counter | `queue` | Committed windows. |
| `queen_s3_records_written_total` | counter | `queue` | Records written to objects. |
| `queen_s3_bytes_written_total` | counter | `queue`, `format` | Object bytes, after compression. |
| `queen_s3_records_lost_total` | counter | `queue` | Records retention deleted before the sink read them. Should be flat forever. |
| `queen_s3_window_records` | histogram | none | Records per window. |
| `queen_s3_window_bytes` | histogram | none | Object bytes per window. |
| `queen_s3_buffer_bytes` | gauge | none | Buffered record bytes across every queue, against `QUEEN_S3_MEMORY_MB`. |
| `queen_s3_fetch_calls_total` | counter | `queue`, `result` | `POST /api/v1/fetch` calls, `result` being `ok` or `error`. |
| `queen_s3_discovery_partitions` | gauge | `queue` | Partitions the last discovery sweep returned. |
| `queen_s3_s3_requests_total` | counter | `op`, `code` | S3 API requests. `op` is `put`, `get`, `head`, `list`, `delete` or a `multipart_*`; `code` is the HTTP status, or `0` when the request never got one. |
| `queen_s3_commit_precondition_lost_total` | counter | none | Key/value batches a lost precondition rolled back. Non-zero means two instances contended for one queue. |
| `queen_s3_checkpoint_age_windows` | gauge | `queue` | Windows committed since the last position checkpoint. It bounds the re-read after a restart. |

The two histograms are deliberately unlabelled: a histogram per queue is buckets times queues in
series, and the per-queue answer worth having is the lag gauge rather than the size distribution.

## On Kubernetes

The repository ships no Helm chart, so a Kubernetes deployment is a manifest an operator writes.
Four things belong in it, and none of them is obvious from the variable table alone.

- **The S3 keypair comes from a Secret**, always. It is the credential for the hop the broker never
  makes, and it does not belong in a values file or a plain `env:` block.
- **`QUEEN_S3_LISTEN` has to be widened to be scraped.** The default is loopback, which is right
  for a sidecar and wrong for a `ServiceMonitor`. Set `0.0.0.0:9333` and add the `containerPort`.
- **One replica**, for the reason under [scaling out](#scaling-out): a second replica gets the same
  queue list and races the first for every window.
- **The readiness probe is `/healthz`**, and its budget is three windows, so a `MAX_WINDOW_MS` of
  five minutes wants a probe timeout comfortably inside fifteen minutes rather than one that
  restarts the pod during a slow window.

In embedded mode there is nothing extra to deploy: the three broker variables above go on the
broker's own workload, and the sink is a process inside the pod that already exists.

## In Queen Cloud

[Queen Cloud](https://queenmq.cloud/try) runs the broker as a managed cell, and has a free tier.

In Cloud the sink does not talk to a broker, it talks to the **cell proxy**, and the proxy talks to
the broker. That is the whole of the difference, and it is one variable: point `QUEEN_URL` at the
proxy's Service. Every call then crosses the proxy's authentication, tenant scoping, quotas and
metering exactly as an SDK's call does.

**Both routes the sink reads through are `Consume`, and neither is ever quota-blocked.**
`POST /api/v1/fetch` and `POST /api/v1/partitions/changed` are classified for the authority they
need rather than for what they write. Neither writes anything, but one hands out message payloads
and the other hands out the partition names, offsets and retention watermarks to read them by, so
both carry the authority of the pop they stand beside instead of the read level every user role
already has. A blocked read would be worse here than for a consumer: a sink refused **discovery**
cannot even learn which partitions it is behind on, so it stops mirroring a queue whose backlog
keeps growing, which is the opposite of what a storage block is for. Both are `POST` on the exact
path only, so every other spelling fails closed.

**The commit pointers are reclassified.** `POST /api/v1/kv` is normally gated on the `kv` feature
and answers the storage block on writes. A batch that addresses **nothing but** the sink's reserved
space, namespace `queen-s3` and key prefix `s3:`, is reclassified `Consume` at the gateway, which
has two consequences worth stating separately. A tenant whose plan has never heard of the `kv`
feature can still run a sink. And **a tenant over its storage quota can still commit a window**,
which is the trap this exists for: refusing a commit pointer of a few hundred bytes does not stop
the tenant growing, it makes the sink re-upload and re-commit the identical window for ever while
its lag grows without bound.

Everything about that carve-out fails closed. An unreadable body, an empty batch, an operation the
proxy has not been told about, one foreign key anywhere in the array: all of them mean "not a sink
batch", which means today's gating. The cost of getting it wrong in that direction is a stalled
sink with a loud lag gauge; the cost in the other direction would be a tenant slipping the feature
gate with an arbitrary key.

**Nothing here is message metered.** A fetch and a discovery call each book one request and zero
messages, exactly as ack and lease extension do. A tenant mirroring a million records into its own
bucket is billed for the requests that carried them, not for the records, and the egress to the
object store is the tenant's own bill with its own provider. That is a pricing decision rather than
a gap, and it is stated here so nobody discovers it from an invoice.

The tenant of every call is the proxy's, from its trusted header; the sink sends only its token.

## What the boot log says

One line, and it is the one to compare against what the bucket actually contains:

```text
queen-s3 1.5.0 sink=default instance=sink-0 queues=orders,clicks partitions=discovery
  queen=http://localhost:6632 token=<set> s3=https://s3.eu-central-1.amazonaws.com
  bucket=my-lake prefix=queen region=eu-central-1 addressing=virtual-host sse=AES256
  format=jsonl compression=zstd layout=merged align=hour target_mb=128 max_window_ms=300000
  start=latest checkpoint_every=20 memory_mb=1024 fetch_concurrency=4
  discovery_interval_ms=2000 safe_guard_ms=5000 lease_ttl_ms=30000
  multipart_threshold_mb=64 listen=127.0.0.1:9333 crash_at=never
```

Every enum is printed the way the **environment** spells it, `per-partition` and not
`PerPartition`, so the line can be pasted back into a manifest rather than translated. `token` is
`<set>` or `<unset>` and never the token; the secret key is never printed at all; and an
`aws:kms` key id is printed as its last four characters, which is enough to confirm the right one
was configured and not enough to name the key to a log aggregator.

Source: https://queenmq.com/deploy/s3/index.mdx
