Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
12 changes: 6 additions & 6 deletions AGENTS.md

Large diffs are not rendered by default.

5 changes: 3 additions & 2 deletions CHANGELOG.md

Large diffs are not rendered by default.

13 changes: 8 additions & 5 deletions clients/ts/src/stream/sse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -344,11 +344,14 @@ export class SSETransport<T = Record<string, unknown>> implements StreamTranspor
// is transient even though the stream still ends (#469). Note WaveHouse
// doesn't reject here *for authentication* — the endpoint is ungated and
// answers a bad token with a reduced view, which is why `auth` is re-read
// per attempt (#239 tracks enforcing expiry). The one 4xx it raises itself
// is a 400 on this route for a missing or empty `table`
// (internal/api/stream.go); a 404/405 means the request missed the route
// entirely (router.go's chi handlers), and anything else comes from
// something in front — a gateway, or a proxy.
// per attempt (#239 tracks enforcing expiry). The 4xx it raises itself are
// a 400 for a missing or empty `table` (internal/api/stream.go) or a
// malformed `X-Tenant-ID`, and a 404 `unknown tenant` for a tenant it does
// not serve (internal/api/tenant.go) — never held, or removed, which is
// what a stream the server ended with its tenant reconnects into; a
// rejected tenant's 503 is retried like any 5xx. Any other 404/405 means
// the request missed the route entirely (router.go's chi handlers), and
// anything else comes from something in front — a gateway, or a proxy.
return { liveMs: 0, terminal: !error.retryable };
}

Expand Down
12 changes: 4 additions & 8 deletions config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,10 @@
# declare (a typo, or a tunable that moved to the settings directory) refuses
# to boot and is named — whether it arrives as a YAML key here or as a WH_*
# variable no key binds.
# Root for embedded state. NATS lives at <data_dir>/nats; each tenant's Pebble
# dedupe store (when its dedupe is enabled) at <data_dir>/<tenant>/dedupe —
# tenant 0's for a settings directory of the four files. An earlier layout's
# <data_dir>/pebble is moved there at boot when <data_dir>/0/dedupe is absent;
# with both present, boot uses <data_dir>/0/dedupe, leaves the old directory
# alone, and warns; a move that fails refuses boot. In a container, this MUST
# be on a host-backed volume — the relative default is for local binary use
# only.
# Root for embedded state. NATS lives at <data_dir>/nats; Pebble, holding every
# tenant's dedupe store while any tenant has dedupe enabled, at
# <data_dir>/pebble. In a container, this MUST be on a host-backed volume —
# the relative default is for local binary use only.
data_dir: ./data

server:
Expand Down
2 changes: 1 addition & 1 deletion deployments/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ RUN --mount=type=cache,target=/go/pkg/mod \
CGO_ENABLED=0 go build -tags="${BUILD_TAGS}" -ldflags="-s -w" -o /bin/wavehouse ./cmd/wavehouse

# Create the parent state and settings directories owned by the nonroot
# user (UID 65532 in distroless). The binary creates `nats/` and `<tenant>/dedupe/`
# user (UID 65532 in distroless). The binary creates `nats/` and `pebble/`
# subdirectories itself when NATS/Pebble open their stores — pre-creating
# them here would buy nothing and obscures intent. Named-volume copy-up runs
# against `/app/data` regardless; bind mounts mask the image dir entirely
Expand Down
2 changes: 1 addition & 1 deletion deployments/Dockerfile.goreleaser
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
# variant in lockstep with deployments/Dockerfile (source-build) so the two
# images expose the same /app filesystem layout to operators.
#
# Only the parents are pre-created — the binary mkdirs `nats/` and `<tenant>/dedupe/`
# Only the parents are pre-created — the binary mkdirs `nats/` and `pebble/`
# subdirs itself when NATS/Pebble open their stores. Named-volume copy-up
# runs against `/app/data`; bind mounts mask the image dir entirely and
# require the host directory to be writable by UID 65532.
Expand Down
5 changes: 2 additions & 3 deletions deployments/compose/standalone.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,8 @@ services:
ports:
- "8080:8080"
environment:
# NATS lives at /app/data/nats, each tenant's Pebble dedupe store at
# /app/data/<tenant>/dedupe (tenant 0's for the four-file settings
# directory) — subdirs are convention, not config. The Dockerfile
# NATS lives at /app/data/nats, Pebble (every tenant's dedupe store) at
# /app/data/pebble — subdirs are convention, not config. The Dockerfile
# pre-creates /app/data owned by the nonroot user; bind-mount any
# persistent volume here.
WH_DATA_DIR: /app/data
Expand Down
8 changes: 4 additions & 4 deletions docs/src/content/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ Status code: `503 Service Unavailable`

The boot-degraded response lets an operator `curl /livez` to learn why the gateway isn't ready to serve traffic yet, instead of grepping a restart-loop log. The binary is bound on `:8080` and serves diagnostics, but is not yet accepting ingest/query traffic. Schema discovery retries with exponential backoff (2s → 60s); once a Refresh succeeds, `/livez` flips to `200` and stays there for the rest of the process lifetime — transient ClickHouse blips after that point are reflected in `/readyz`, not `/livez`.

Over a [nested settings directory](/deployment#the-nested-settings-directory) the probe reads every tenant together: `/livez` is `503` while **no** tenant has completed a first discovery — the diagnostic names the tenant whose attempt it reports (`schema discovery: tenant acme: …`), and reads `no tenant has completed a first discovery yet` before any attempt or when the directory serves no tenant — and `200` from the first tenant's success on, for the rest of the process lifetime. A tenant whose ClickHouse is unreachable after that is a log line and the `wavehouse_schema_refresh_failures_total{tenant}` counter, never a probe failure. A tenant that has not completed its own first discovery answers `503` (`schema not loaded yet`) on its schema-aware routes until it does; one whose ClickHouse goes down after that answers query errors, as a single-tenant server does.
Over a [nested settings directory](/deployment#the-nested-settings-directory) the probe reads every tenant together: `/livez` is `503` while **no** tenant has completed a first discovery — the diagnostic names the tenant whose attempt it reports (`schema discovery: tenant acme: …`), and reads `no tenant has completed a first discovery yet` before any attempt, when the directory serves no tenant, and once the tenant it named stops being served — and `200` from the first tenant's success on, for the rest of the process lifetime. A tenant whose ClickHouse is unreachable after that is a log line and the `wavehouse_schema_refresh_failures_total{tenant}` counter, never a probe failure. A tenant that has not completed its own first discovery answers `503` (`schema not loaded yet`) on its schema-aware routes until it does; one whose ClickHouse goes down after that answers query errors, as a single-tenant server does.

---

Expand Down Expand Up @@ -152,7 +152,7 @@ Status code: `503 Service Unavailable`

### `GET /v1/health` — Liveness ping (public, content-free)

Returns **`200 OK` with an empty body** once the gateway is past boot, or **`503 Service Unavailable`** (also empty) while boot-time schema discovery is still failing. Like every `/v1` route outside `/v1/ops/*` it [resolves a tenant](/deployment#multi-tenant-deployments) first, so a malformed or unknown `X-Tenant-ID` answers `400`/`404` before the probe runs, and — over a [nested settings directory](/deployment#the-nested-settings-directory) — a tenant whose settings folder was rejected answers a `503` that carries the usual JSON error body rather than this route's empty one, as does a request carrying a token, with no valid operator key, while that tenant's [JWKS has not been fetched yet](#authentication) (`Retry-After: 30`; the SDK sends its token on this ping too). No authentication required and no response body — the caller only branches on the status code, so there's nothing to JSON-encode or cache per request.
Returns **`200 OK` with an empty body** once the gateway is past boot, or **`503 Service Unavailable`** (also empty) while boot-time schema discovery is still failing. Like every `/v1` route outside `/v1/ops/*` it [resolves a tenant](/deployment#multi-tenant-deployments) first, so a malformed or unknown `X-Tenant-ID` answers `400`/`404` before the probe runs, and — over a [nested settings directory](/deployment#the-nested-settings-directory) — a tenant whose settings folder was rejected answers a `503` that carries the usual JSON error body rather than this route's empty one, as does a request carrying a token, with no valid operator key, while that tenant's [JWKS has not been fetched yet](#authentication) (`Retry-After: 30`; the SDK sends its token on this ping too). No authentication required and no response body — the caller only branches on the status code, so there's nothing to JSON-encode or cache per request. Resolving the tenant makes the ping an unauthenticated answer to whether a tenant is served, deliberately: every tenant route gives an unknown tenant the same `404` before authenticating, since authenticating takes that tenant's own verifier, so the ping reveals nothing the others don't. It also answers a served tenant's ping from that tenant's own CORS list, where a tenant-exempt route answers from tenant `0`'s, which a nested directory need not have.

This is what the SDK's `wh.sys.health()` calls, and the endpoint to use when choosing among multiple servers in a distributed setup. It mirrors `/livez` under the hood but is intentionally a `/v1` API route rather than a Kubernetes probe path: an operator may filter the bare probe paths (`/livez`, `/readyz`, `/healthz`) out at the reverse proxy since they're internal probes, so the SDK relies on `/v1/health`, which is documented public API surface meant to stay reachable. It does **not** ping ClickHouse — readiness-based load balancing is the proxy/LB's job (via `/readyz`), not the client's.

Expand Down Expand Up @@ -612,7 +612,7 @@ Opens a persistent SSE connection for real-time event streaming. Supports histor
| ------ | ----------- |
| `Last-Event-ID` | RFC 3339 timestamp of the last received event. If present, overrides the `since` query parameter for automatic reconnection (standard `EventSource` behavior). |

**Response:** SSE stream (`text/event-stream`). Data events include an `id:` field set to the event's `received_timestamp`. The stream opens with a `: connected` comment and emits a minimal `:` keepalive comment periodically (every 30 seconds by default), which keeps a quiet connection from being closed by a proxy; both are standard SSE comments that `EventSource` ignores (raw consumers should skip `:`-prefixed lines). When the server stops (see [Stopping](/deployment#stopping)) it ends every open stream immediately rather than holding it for the drain; `EventSource` reconnects on its own and resumes from `Last-Event-ID`.
**Response:** SSE stream (`text/event-stream`). Data events include an `id:` field set to the event's `received_timestamp`. The stream opens with a `: connected` comment and emits a minimal `:` keepalive comment periodically (every 30 seconds by default), which keeps a quiet connection from being closed by a proxy; both are standard SSE comments that `EventSource` ignores (raw consumers should skip `:`-prefixed lines). When the server stops (see [Stopping](/deployment#stopping)) it ends every open stream immediately rather than holding it for the drain; `EventSource` reconnects on its own and resumes from `Last-Event-ID`. A reload that stops serving the stream's tenant — its folder removed or rejected, over a [nested settings directory](/deployment#the-nested-settings-directory) — ends that tenant's open streams the same way, and the reconnect then gets its `404` (removed) or `503` (rejected): the SDK stops on the `404` and retries the `503`, resuming from `Last-Event-ID` once the folder is back, while a browser `EventSource` treats either as fatal. A browser going cross-origin reads either refusal only when it passes CORS: it is decorated from tenant `0`'s list ([multi-tenant deployments](/deployment#multi-tenant-deployments)), so where tenant `0` is not served or its list does not admit the page's origin, the SDK sees a network error instead and keeps re-dialing.

**Row values arrive positionally, and the column names are announced separately.** Before the first row, and again whenever the column list changes, the stream sends an `event: schema` frame naming the columns of the rows that follow — in order, already reduced to what the caller's role may read. That re-announcement is **not** guaranteed after a gap-fill across a column change; see the arity note below. Every data frame's `row` array then has exactly one value per announced column, in that order. `schema` is a **named** SSE event, so a browser `EventSource` must `addEventListener('schema', …)` — it never reaches `onmessage`. A schema frame carries **no** `id:` line, so it never moves the client's `Last-Event-ID`. In the example below the table has its own `received_timestamp` **column**, which collides by name with the frame's top-level `received_timestamp` **field** — they are different values: the field is when WaveHouse received the event, the row slot is that column as published (`null` where the record omitted it, which ClickHouse replaces with the column's default on insert).

Expand Down Expand Up @@ -873,7 +873,7 @@ 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 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 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.

Use `GET /v1/ops/dlq/stats` to monitor DLQ depth.

Expand Down
Loading
Loading