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
5 changes: 5 additions & 0 deletions .github/labeler.yml
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,11 @@
- any-glob-to-any-file:
- "internal/api/**"

"area/tenant":
- changed-files:
- any-glob-to-any-file:
- "internal/tenant/**"

"area/ingest":
- changed-files:
- any-glob-to-any-file:
Expand Down
14 changes: 8 additions & 6 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,10 +26,10 @@ One binary:

- **`cmd/wavehouse/`** — Standalone mode (all-in-one with embedded NATS, optional Pebble dedup): argv dispatch, the logger, `config.Load`, and the signal context; everything else is `internal/app`

Seventeen internal packages under `internal/` (plus `internal/testutil/` for shared test helpers):
Eighteen internal packages under `internal/` (plus `internal/testutil/` for shared test helpers):

- **`api/`** — Chi HTTP router, JWT/JWKS middleware (from `auth/`), ingest/query/structured-query/SSE/schema/DLQ/pipes handlers
- **`app/`** — the process wiring: `New` builds every component from the boot config and the settings directory (each one wired in one place — what it opens, what it loops, what it releases — with the settings store handed to its wiring function whole, the injection point the per-tenant registry of #583 lands on), `Run` drives the long-lived ones under one `errgroup` until the context is cancelled or one fails, `Close` releases them in reverse order. `cmd/wavehouse` and `tests/integration` both boot through it
- **`app/`** — the process wiring: `New` builds every component from the boot config and the settings directory (each one wired in one place — what it opens, what it loops, what it releases — with the settings store handed to its wiring function whole, the injection point of the per-tenant registry of #583: store-keyed getters for the handlers, `perTenant` for the async paths), `Run` drives the long-lived ones under one `errgroup` until the context is cancelled or one fails, `Close` releases them in reverse order. `cmd/wavehouse` and `tests/integration` both boot through it
- **`auth/`** — JWT auth middleware: HMAC **or** JWKS verification with `alg` pinned to the active verifier, role extraction from a configurable claim path; always runs, never rejects (bad token → empty role + stashed reason)
- **`cache/`** — `Cache` interface → `LocalCache` (Ristretto) + `SharedCache` (TBD) + `TieredCache` (singleflight)
- **`chconn/`** — `Manager`, the one ClickHouse `driver.Conn` every consumer holds; `Reconfigure` swaps the connection behind it after a settings reload changes the wiring (never dials; the old connection closes after a `query_timeout` grace)
Expand All @@ -43,8 +43,9 @@ Seventeen internal packages under `internal/` (plus `internal/testutil/` for sha
- **`pipes/`** — Named query pipes: `NamedQuery` type + `BindParams` + `Source` (read per request; `settings.Store` in production, `Static(q...)` in tests)
- **`policy/`** — Hasura-style access control, **role-first**: `TablePolicy` is `map[string]RolePermissions`, and a role's grant splits by operation into `SelectPermissions` (columns, row `filter`, aggregations, the `max_*` limits) and `InsertPermissions` (columns, `check`) — so a field only one side honors does not exist on the other. `Evaluate()` resolves ONE operation and leaves the other side **nil** (`Select *ResolvedSelect` / `Insert *ResolvedInsert`), which every accessor fails closed on — nil is "not resolved", distinct from an empty side, which is "unrestricted" (what the admin return builds). Claim templating (`{{ jwt.claim.path }}`) resolves during that call. Policies come from `Source`, a `func() *Policy` read per call (`settings.Store.Policy` in production, `Static(p)` in tests)
- **`query/`** — Structured query AST types + SQL builder with schema validation, structural policy predicate/limit emission, timestamp bucketing
- **`settings/`** — the settings directory: `Validate` (strict JSON, per-file rules, cross-file role references), `Store` (the adopted snapshot + serialized `Reload`, typed accessors read per call, `AfterAdopt` hooks), the fsnotify `Watch`, and the `go:embed`ded seed `wavehouse bootstrap` writes
- **`settings/`** — the settings directory: `Validate` (strict JSON, per-file rules, cross-file role references), `Store` (the adopted snapshot + serialized `Reload`, typed accessors read per call, `AfterAdopt` hooks), `Registry` (tenant id → `Store`; holds the one store under `tenant.Default`), the fsnotify `Watch`, and the embedded (`go:embed`) seed `wavehouse bootstrap` writes
- **`stream/`** — SSE fan-out: rows travel POSITIONALLY, so each connection is told its projected column list in an `event: schema` frame before its first row and again on drift — **not** guaranteed after a gap-fill across a column change, which can leave a connection reading live rows against a stale list until it reconnects ([#543](https://github.com/Wave-RF/WaveHouse/issues/543)) — (tracked per connection; replay tracks its own). The event `Hub` (registers subscribers by `(topic, role)`; `Broadcast` projects + serializes each event once per role, the #294 delivery hot path — a role carrying a row-level `filter` keeps the shared projection but delivers per subscriber, each subscriber's claims evaluated against the row, #319), `Subscriber` (per-connection outbound `Frame` queue, `Send`/`Frames`; claims fixed at construction, immutable), the `Bucket` fan-out set (`subscriberSet`, one per `(topic, role)`), the `Heartbeater` keepalive wheel, and `Metrics` (the `wavehouse_sse_*` stream instruments)
- **`tenant/`** — the tenant identifier ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)): `ID` (a validated string), `Parse` (letters, digits, `_`, `-`; ≤ 64 bytes — safe as a folder name and as an MQ subject token), `Default` (`"0"`), and `Header` (`X-Tenant-ID`). Imports nothing from the rest of the repo. `api.TenantMW` resolves the header against `settings.Registry` before auth on every `/v1` route outside `/v1/ops/*` (`400` malformed, `404` unknown) and puts the resolved `*settings.Store` in the request context; handlers read it once (`api.StoreFromContext`) and pass it down as an argument, and nothing below a handler reads context. The async paths (ingest worker, sweeper, stream hub, schema registry) are constructed with a `tenant.ID` and their getters take it

## Key Design Decisions

Expand Down Expand Up @@ -74,10 +75,10 @@ The invariant index — what must stay true. Full narrative and rationale live i
## Code Conventions

- **Go 1.26**, strict formatting (`gofumpt`, enforced by CI)
- **Structured logging** with `log/slog` (JSON handler)
- **Structured logging** with `log/slog` (JSON handler), through the default logger: call `slog.InfoContext(ctx, …)` and its siblings (the context carries the trace ids the handler stamps) rather than taking a `*slog.Logger` parameter or field. `cmd/wavehouse` and `internal/app` install the default; tests silence or capture it with `internal/testutil/logtest` (a capturing test must not call `t.Parallel()`)
- **Chi v5** for HTTP routing
- **Error handling**: Return errors, don't panic. Wrap with `fmt.Errorf("context: %w", err)`.
- **No global state**: Dependencies are passed explicitly (constructor injection).
- **No global state**: Dependencies are passed explicitly (constructor injection). The `slog` default logger is the one sanctioned exception.
- **Package naming**: Lowercase, single word (or abbreviated). `internal/` enforces module privacy.

## Craftsmanship
Expand Down Expand Up @@ -440,7 +441,8 @@ internal/policy/ → Access control policies (types, evaluation, Source)
internal/query/ → Structured query AST + SQL builder
internal/settings/ → Settings directory (validate, adopted snapshot + reload, watcher, embedded seed)
internal/stream/ → SSE fan-out (event Hub: project once per role, Subscriber outbound queue, Bucket fan-out, keepalive Heartbeater wheel)
internal/testutil/ → Shared test helpers (NopLogger, etc.)
internal/tenant/ → Tenant id (type, grammar, reserved default, request header name)
internal/testutil/ → Shared test helpers (mocks, JWT + schema helpers; logtest/ captures or silences the default logger)
tests/ → Integration & E2E tests
tests/integration/ → Go integration tests (//go:build integration; ClickHouse testcontainer)
tests/e2e/ → E2E test stack (scripts/orchestrator boots a ClickHouse testcontainer + the wavehouse-cov binary)
Expand Down
6 changes: 6 additions & 0 deletions CHANGELOG.md

Large diffs are not rendered by default.

6 changes: 5 additions & 1 deletion docs/src/content/docs/access-control.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -212,6 +212,10 @@ On a structured query (`POST /v1/query?table={table}`) the allowlist is a **hard

This produces `WHERE (tenant_id = ?)` with the caller's `app_metadata.tenant_id` claim bound as the parameter — so a `viewer` only ever sees rows for their own tenant, and the value comes from the signed token, not from anything the client sends.

:::note[Two meanings of "tenant"]
On this page a *tenant* is a row-scoping value carried in the signed token. The [`X-Tenant-ID` header](/deployment#multi-tenant-deployments) is a different axis: it selects which settings directory — and so which `policies.json` — serves the request. It is client-supplied, resolved before authentication, and isolates no rows. A settings directory is one such tenant (`0`), so most deployments never send it.
:::

Supported comparison operators:

| Operator | SQL | Meaning |
Expand Down Expand Up @@ -360,7 +364,7 @@ A few more edges worth knowing when you write a policy — the stream evaluates
- **A non-scalar event value** (array/object/null) under a filtered column withholds the row.
- **Insert-time numeric narrowing is simulated, not skipped.** The insert narrows a payload carrying more precision than the column's declared type — a `Decimal` **truncates** at its scale (`1.005`, `1.006` and `1.009` all store as `1.00` in a `Decimal(10, 2)`), a `Float32` **rounds** to its nearest representable value (`16777217` stores as `16777216`) — and ClickHouse applies the same narrowing to a bound filter constant at compare time. The stream narrows **both operands** identically before comparing, so its verdict matches the query path's on narrowing columns under every operator: a `_gt: "1.004"` filter on a `Decimal(10, 2)` column withholds a `1.005` payload exactly as the query path hides the stored `1.00`. (An earlier revision of this feature compared the raw payload and could deliver such an event; that fail-open is closed, and an integration test holds every in-range numeric stream verdict equal to a live ClickHouse's — for the out-of-range operands the range gate refuses, it asserts the half that matters: the stream never admits a row ClickHouse hides.) What remains payload-vs-stored: an event whose insert later **fails outright** (an out-of-range value, a batch error, the DLQ) was already streamed to whichever subscribers the filter admitted, and its row never becomes queryable.

One more boundary is temporal: a subscriber's claims (and role) are captured when the SSE connection is established and are never re-read, while the policy itself is re-read on every live event. A gap-fill replay is the one exception: it runs under the single policy snapshot taken when the connection opened, so a policy change landing mid-replay applies from the first live event after it. Tightening a policy therefore applies from the next live event, but a token that expires — or claims revoked at the identity provider — keeps its open stream until the client disconnects, so treat connection lifetime as the revocation window for stream row-scoping.
One more boundary is temporal: a subscriber's claims (and role) are captured when the SSE connection is established and are never re-read, while the policy itself is re-read on every event, live and replayed alike — so a policy adopted mid-gap-fill applies to the next replayed row. Tightening a policy therefore applies from the next event either way, but a token that expires — or claims revoked at the identity provider — keeps its open stream until the client disconnects, so treat connection lifetime as the revocation window for stream row-scoping.
Comment thread
taitelee marked this conversation as resolved.

Each row withheld **from a subscriber** by row-level security is counted in `wavehouse_sse_rows_withheld_total` (labeled by table and role; a row withheld from three subscribers counts three times) — check it before concluding a quiet stream simply has no matching rows.

Expand Down
4 changes: 2 additions & 2 deletions docs/src/content/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,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. 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 bad `X-Tenant-ID` answers `400`/`404` before the probe runs. 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.

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 @@ -616,7 +616,7 @@ Each SSE connection is bound to a single `?table=`; to consume multiple tables,

Values of top-level `DateTime`/`DateTime64` columns inside `row` arrive in the canonical RFC 3339 UTC form (ingest rewrites them before publishing — see [timestamp canonicalization](#timestamp-canonicalization)), so a live event and a `/v1/query` read of the same row agree on the instant in zone-explicit form — a zone-less spelling no longer parses as local time in a browser ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)). The two renderings are byte-identical regardless of the declared time zone or a `Nullable` wrapper — a column declared with a non-UTC zone also streams as `Z`, and `/v1/query` normalizes it (nullable or not) to UTC before rendering. Canonicalization is fail-open at ingest, so a value outside the accepted input forms streams in whatever spelling the producer sent — and for exactly those events the byte-identity above does not hold: a spelling ClickHouse accepts anyway is stored and still queries back canonical, while one it too rejects lands in the DLQ and never becomes queryable at all.

**Note:** When access control policies are active, streamed events are filtered per the caller's role: tables without `select` permission are skipped, denied columns are removed from each event, and the role's [row-level `filter`](/access-control#row-level-security) is evaluated per subscriber against the caller's JWT claims — supplied by the connection's token (the `Authorization` header, or the `?token=` fallback above), with replayed gap-fill events filtered the same way. For a filter constant the query path's SQL also accepts ([the enforcement caution](/access-control#where-each-rule-is-enforced) gives per-type guidance), a connection is never delivered a row the query path would hide for that role — every comparison the stream can't prove fails closed and withholds the row instead. Numeric comparisons run in the column's storage domain — both operands narrowed the way ClickHouse narrows the stored value and the bound constant — so columns that narrow on insert (`Float32`/`Float64` width, a `Decimal`'s scale) agree with the query path too; the residual payload-vs-stored case is an event whose insert later fails into the DLQ, which the caution documents. The connection's claims are captured once, when the stream is established — a policy change applies from the next live event (an in-flight gap-fill finishes under the policy snapshot taken when the stream opened), but an expired token or changed claims take effect only when the client reconnects.
**Note:** When access control policies are active, streamed events are filtered per the caller's role: tables without `select` permission are skipped, denied columns are removed from each event, and the role's [row-level `filter`](/access-control#row-level-security) is evaluated per subscriber against the caller's JWT claims — supplied by the connection's token (the `Authorization` header, or the `?token=` fallback above), with replayed gap-fill events filtered the same way. For a filter constant the query path's SQL also accepts ([the enforcement caution](/access-control#where-each-rule-is-enforced) gives per-type guidance), a connection is never delivered a row the query path would hide for that role — every comparison the stream can't prove fails closed and withholds the row instead. Numeric comparisons run in the column's storage domain — both operands narrowed the way ClickHouse narrows the stored value and the bound constant — so columns that narrow on insert (`Float32`/`Float64` width, a `Decimal`'s scale) agree with the query path too; the residual payload-vs-stored case is an event whose insert later fails into the DLQ, which the caution documents. The connection's claims are captured once, when the stream is established — a policy change applies from the next event, replayed or live (a gap-fill re-reads the policy per event too), but an expired token or changed claims take effect only when the client reconnects.

**CORS:** `/v1/stream` honors the `cors.allowed_origins` allowlist (settings directory) like every endpoint. Note that a **header-authenticated stream preflights before it connects** — `Authorization` is not CORS-safelisted — where a bare `EventSource` never preflighted at all: its request is not a `fetch()`, so Fetch's unsafe-request flag is never set and `Last-Event-ID` rides on the plain `GET`. Both headers are allow-listed, so an allowed origin connects *and* resumes cross-origin.

Expand Down
Loading
Loading