Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions AGENTS.md

Large diffs are not rendered by default.

8 changes: 5 additions & 3 deletions CHANGELOG.md

Large diffs are not rendered by default.

18 changes: 11 additions & 7 deletions clients/ts/src/dlq.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { err, ok } from "./errors.js";
import { request } from "./http.js";
import { request, tenantParam } from "./http.js";
import type { StreamController } from "./stream/controller.js";
import type { DLQStats, HttpContext, Result, StreamOptions } from "./types.js";
import type { DLQStats, HttpContext, OpsRequestOptions, Result, StreamOptions } from "./types.js";

type CreateStreamFn = (table: string, opts?: StreamOptions) => StreamController;

Expand All @@ -15,23 +15,27 @@ export class DLQNamespace {
this._createStream = createStream;
}

/** Get DLQ statistics (message counts per table). */
async list(opts?: { signal?: AbortSignal }): Promise<Result<DLQStats>> {
/**
* Get DLQ statistics (message counts per table) — of `opts.tenant`, the
* default tenant without it. A tenant with no dead-letter queue is a `404`.
*/
async list(opts?: OpsRequestOptions): Promise<Result<DLQStats>> {
const { data, error } = await request<DLQStats>(this._ctx, {
method: "GET",
path: "/v1/ops/dlq/stats",
params: tenantParam(opts),
signal: opts?.signal,
});
if (error) return err(error);
return ok(data!);
}

/** Get DLQ stats filtered by table name. */
async table(name: string, opts?: { signal?: AbortSignal }): Promise<Result<DLQStats>> {
/** Get DLQ stats filtered by table name — of `opts.tenant`, the default tenant without it. */
async table(name: string, opts?: OpsRequestOptions): Promise<Result<DLQStats>> {
const { data, error } = await request<DLQStats>(this._ctx, {
method: "GET",
path: "/v1/ops/dlq/stats",
params: { table: name },
params: { table: name, ...tenantParam(opts) },
signal: opts?.signal,
});
if (error) return err(error);
Expand Down
19 changes: 19 additions & 0 deletions clients/ts/src/namespaces.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,25 @@ describe("DLQNamespace", () => {
expect(fetchSpy.mock.calls[0][0]).toContain("table=clicks");
});

it("list() and table() send opts.tenant as ?tenant=, and nothing without it", async () => {
fetchSpy.mockImplementation(
async () => new Response(JSON.stringify({ tables: {}, total: 0 }), { status: 200 }),
);
const ns = new DLQNamespace(makeCtx(), mockStream);

await ns.list({ tenant: "acme" });
await ns.table("clicks", { tenant: "acme" });
await ns.list();
await ns.table("clicks");

const urls = fetchSpy.mock.calls.map((call) => new URL(call[0]));
expect(urls[0].pathname + urls[0].search).toBe("/v1/ops/dlq/stats?tenant=acme");
expect(urls[1].searchParams.get("table")).toBe("clicks");
expect(urls[1].searchParams.get("tenant")).toBe("acme");
expect(urls[2].search).toBe("");
expect(urls[3].search).toBe("?table=clicks");
});

it("stream() delegates to createStream", () => {
const ctrl = {} as any;
mockStream.mockReturnValue(ctrl);
Expand Down
4 changes: 2 additions & 2 deletions clients/ts/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -471,8 +471,8 @@ export interface PipeRequestOptions {
/**
* Options for a call to one of the admin routes that address a tenant:
* `wh.pipes.list()`, `wh.pipes.get()`, `wh.settings.reload()`,
* `wh.schema.list()`, `wh.schema.refresh()`, `wh.from(t).schema()` and
* `wh.sql()`.
* `wh.schema.list()`, `wh.schema.refresh()`, `wh.from(t).schema()`,
* `wh.dlq.list()`, `wh.dlq.table()` and `wh.sql()`.
*/
export interface OpsRequestOptions {
signal?: AbortSignal;
Expand Down
15 changes: 9 additions & 6 deletions docs/src/content/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -274,7 +274,7 @@ The body is a **flat JSON object** whose keys must match column names in the tar
| 500 | `{"error":"dedupe failed"}` | Deduplication backend error |
| 503 | `{"error":"schema not loaded yet"}` | The tenant's first schema discovery has not succeeded yet (its ClickHouse unreachable, or [no pool for it](/settings-directory#clickhouse)), so whether the table exists is not known; `Retry-After: 5`. Decided before the body is read |
| 500 | `{"error":"publish failed"}` | Message queue error |
| 503 | `{"error":"service unavailable"}` | NATS JetStream stream full (backpressure). Response includes `Retry-After: 30` header. |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure, for that tenant alone) or not open (see [Message Queue](/settings-directory#message-queue)). Response includes `Retry-After: 30` header. |
| 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) |

**curl example:**
Expand Down Expand Up @@ -385,7 +385,7 @@ A `200` is returned whenever the body was read and the records were processed
| 413 | `{"error":"request body exceeded 16777216 bytes"}` | Request body over the 16 MiB cap |
| 415 | `{"error":"no Content-Type: ingest requires one of application/json, application/x-ndjson, …"}` (declared variant: `Content-Type "text/plain": ingest requires one of …` — see the note above on how declarations are echoed; conflicting variant: `conflicting Content-Type declarations "application/json", "application/x-ndjson": ingest reads one format per request, and requires one of …`) | The request declared no `Content-Type`, one whose media type is unsupported or does not parse, a comma-bearing value that does not parse as a single media type, or repeated header lines that disagree — different formats, or one supported and one not. Checked before the body is parsed |
| 500 | `{"error":"publish failed"}` / `{"error":"dedupe failed"}` | Message-queue or dedup-backend failure mid-batch |
| 503 | `{"error":"service unavailable"}` | NATS JetStream full (backpressure) mid-batch; includes `Retry-After: 30` |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure) or not open, mid-batch; includes `Retry-After: 30` |
| 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) |

:::caution[At-least-once on retry]
Expand Down Expand Up @@ -745,21 +745,24 @@ Triggers an immediate re-discovery of the `?tenant=`'s ClickHouse table schemas

#### `GET /v1/ops/dlq/stats` — DLQ Statistics

Returns per-table message counts in the Dead Letter Queue — a table's count summed across tenants, since one queue serves every tenant until each has its own. Admin-only, like the rest of this section. Whether a poison row lands here is the settings directory's [`dlq.enabled`](/settings-directory#dead-letter-queue) switch (global or per table); the stream and this endpoint always exist. Before any failure has ever occurred, the endpoint returns `200` with `{"tables":{},"total":0}`.
Returns per-table message counts in one tenant's Dead Letter Queue: the [tenant](/deployment#the-nested-settings-directory) an optional `?tenant=<id>` names, the default tenant `0` without it, which is the whole settings directory unless it is nested. The tenant is looked up in the message queue, not the settings, so a tenant whose folder was rejected or removed is read like one being served, since its queue is kept (nothing deletes it). The query string is parsed strictly, as on the other admin reads. Admin-only, like the rest of this section. Whether a poison row lands here is the settings directory's [`dlq.enabled`](/settings-directory#dead-letter-queue) switch (global or per table); a tenant's dead-letter stream is opened when the tenant is first served, and this endpoint always exists. Before any failure has ever occurred, the endpoint returns `200` with `{"tables":{},"total":0}`.

**Error responses:**

| Status | Body | Cause |
| ------ | ---- | ----- |
| 401 | `{"error":"invalid token"}` / `{"error":"token expired"}` | A present-but-invalid/expired token was supplied and denied (the gate surfaces the token reason) |
| 400 | `{"error":"invalid query string: …"}` / `{"error":"invalid ?tenant: …"}` | The query string does not parse (`?tenant=acme;x=1`, a bad `%` escape), or `tenant` is empty, repeated, or not a tenant id |
| 403 | `{"error":"forbidden"}` | Caller's role is not the policy `admin_role` (`"admin"` by default) |
| 404 | `{"error":"no dead-letter queue for tenant: <id>"}` | The tenant has no dead-letter queue: it has never been served on this data directory, its queue could not be opened (see [Message Queue](/settings-directory#message-queue)), or the id names no tenant |
| 500 | `{"error":"stream info failed"}` | NATS JetStream stream-info lookup failed |
| 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while tenant `0`'s JWKS has not been fetched yet (the ops tree verifies as tenant `0`); refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) |

**Query Parameters:**

| Param | Type | Default | Description |
| ----- | ---- | ------- | ----------- |
| `tenant` | string | `0` | The tenant whose dead-letter queue is read. |
| `table` | string | — | Filter stats to a specific table name (e.g., `?table=clicks` returns only the `clicks` count). |

**Response:**
Expand Down Expand Up @@ -855,7 +858,7 @@ The message format used on NATS JetStream between ingest and the batch consumer:
| `columns` | string[] | The table's **insertable** column names, in declaration order — what each position in `row` means. A `MATERIALIZED` or `ALIAS` column is computed by ClickHouse and cannot be named in an `INSERT`, so it never appears here. |
| `row` | array | One `JSONCompactEachRow` line: one value per entry in `columns`, in that order. A column the request body omitted is `null` here; for a **non-nullable** column the insert turns that back into the column's default (`input_format_null_as_default`), but a `Nullable(T) DEFAULT …` column stores `NULL` — only an *absent* key ever took the default, and a positional row has one slot per column and no way to express absence. Parseable `DateTime`/`DateTime64` values are rewritten to canonical RFC 3339 UTC (see [timestamp canonicalization](#timestamp-canonicalization)); other values as originally sent. |

`columns` and `row` are only meaningful together: a reader that cannot pair them — a length mismatch, an undecodable row, a `columns` list naming one column twice — has no way to map a value to a column. Both readers also refuse an envelope whose `format` they do not recognize, which is what a pre-v2 message looks like. Either way the SSE fan-out withholds such an envelope rather than guess, and the batch consumer parks it on the DLQ with `X-DLQ-*` headers — acking and dropping it only where the DLQ is switched off for that table, since it can never insert on retry. Both outcomes increment `wavehouse_ingest_poison_total`, separated by its `disposition` label (`parked` / `dropped`).
`columns` and `row` are only meaningful together: a reader that cannot pair them — a length mismatch, an undecodable row, a `columns` list naming one column twice — has no way to map a value to a column. Both readers also refuse an envelope whose `format` they do not recognize. Either way the SSE fan-out withholds such an envelope rather than guess, and the batch consumer parks it on the DLQ with `X-DLQ-*` headers — acking and dropping it only where the DLQ is switched off for that table, since it can never insert on retry. Both outcomes increment `wavehouse_ingest_poison_total`, separated by its `disposition` label (`parked` / `dropped`).

### Client-Facing Format (SSE)

Expand All @@ -873,9 +876,9 @@ Three values, where the envelope above has four: this is the frame a role restri

## Dead Letter Queue (DLQ)

When a batch insert to ClickHouse fails (e.g., type errors, connection issues), the worker re-inserts the batch row by row: rows that succeed are acked, and only the rows that fail again are published to the DLQ NATS stream (`WAVEHOUSE_DLQ`) under subjects `dlq.{tenant}.{table}` (the tenant the row was ingested under; `0` for a settings directory that holds the four files). This prevents infinite retry loops — those messages are ACKed from the main stream and moved to the DLQ for inspection. A batch whose tenant has no ClickHouse connection — one no longer served, or one no pool could be opened for (such as by the connection ceiling) — skips that retry, which no row of it could pass, and is parked whole; only a served tenant whose DLQ is off for the table leaves it for redelivery, since a tenant no longer served has no switch to read. A second class lands here too: an envelope the worker cannot *read* at all — malformed JSON, an unknown **or absent** `format` (a pre-v2 message has no `format` field at all, which is how it presents here), or `columns` and `row` that do not pair — is parked without ever reaching a table batch, which is what an operator sees after upgrading across the wire change without draining first. **Two different body shapes land here, and a consumer must not assume one decoder.** A row that failed its INSERT is parked as the `EventMessage` envelope above. An envelope the worker could not *read* is parked as **its original bytes, verbatim** — `parkOnDLQ` republishes what arrived — so it is whatever the producer sent: a pre-v2 `data` object, malformed JSON, or a v2 envelope whose `columns` and `row` do not pair. Being undecodable as an `EventMessage` is precisely why it was parked, so decode defensively and fall back on the `X-DLQ-Error` header, which names the reason. For the first shape the body is the published `EventMessage` envelope (`{"table_name":…,"scope":"","received_timestamp":…,"format":…,"columns":[…],"row":[…]}` — the failed row is the `row` array, read against `columns`, its `DateTime`/`DateTime64` values as published: canonicalized where WaveHouse could parse them, otherwise the producer's original spelling — see [timestamp canonicalization](#timestamp-canonicalization)); the failure reason, table, and time travel in the `X-DLQ-Table` / `X-DLQ-Error` / `X-DLQ-Timestamp` message headers.
When a batch insert to ClickHouse fails (e.g., type errors, connection issues), the worker re-inserts the batch row by row: rows that succeed are acked, and only the rows that fail again are published to the tenant's own DLQ NATS stream (`DLQ_{tenant}`) under subjects `dlq.{tenant}.{table}` (the tenant the row was ingested under; `0` for a settings directory that holds the four files). This prevents infinite retry loops — those messages are ACKed from the main stream and moved to the DLQ for inspection. A batch whose tenant has no ClickHouse connection — one no longer served, or one no pool could be opened for (such as by the connection ceiling) — skips that retry, which no row of it could pass, and is parked whole; only a served tenant whose DLQ is off for the table leaves it for redelivery, since a tenant no longer served has no switch to read. A second class lands here too: an envelope the worker cannot *read* at all — malformed JSON, an unknown **or absent** `format`, or `columns` and `row` that do not pair — is parked without ever reaching a table batch. **Two different body shapes land here, and a consumer must not assume one decoder.** A row that failed its INSERT is parked as the `EventMessage` envelope above. An envelope the worker could not *read* is parked as **its original bytes, verbatim** — `parkOnDLQ` republishes what arrived — so it is whatever the producer sent: malformed JSON, an envelope of an unknown `format`, or a v2 envelope whose `columns` and `row` do not pair. Being undecodable as an `EventMessage` is precisely why it was parked, so decode defensively and fall back on the `X-DLQ-Error` header, which names the reason. For the first shape the body is the published `EventMessage` envelope (`{"table_name":…,"scope":"","received_timestamp":…,"format":…,"columns":[…],"row":[…]}` — the failed row is the `row` array, read against `columns`, its `DateTime`/`DateTime64` values as published: canonicalized where WaveHouse could parse them, otherwise the producer's original spelling — see [timestamp canonicalization](#timestamp-canonicalization)); the failure reason, table, and time travel in the `X-DLQ-Table` / `X-DLQ-Error` / `X-DLQ-Timestamp` message headers.

Use `GET /v1/ops/dlq/stats` to monitor DLQ depth.
Use `GET /v1/ops/dlq/stats` to monitor DLQ depth, per tenant (`?tenant=`).

## Generating a JWT for Testing

Expand Down
Loading
Loading