Skip to content
Merged
6 changes: 4 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ 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`

Eighteen internal packages under `internal/` (plus `internal/testutil/` for shared test helpers):
Nineteen 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 registry 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 (with the tenant each message's `mq.Topic` names for the stream hub and the ingest worker), the `chconn.Pools` and the per-tenant `discoveries` reconciled from `AfterAdopt`, `shortestKeepalive` for the one setting folded over every tenant served, `gapWindows` handing the sweeper each tenant's own gap window (a rejected tenant's as its folder last had it, unbounded for one rejected since boot) and the `mq.max_bytes_gb` reconcile each served tenant's byte budget, and `defaultPolicy` for the one setting that still follows tenant `0`, a flat directory's ops-gate admin role; the auth verifiers are per tenant, reconfigured (rebuilt only on changed wiring) and pruned from `AfterAdopt`, and the same hook's `Hub.Prune` ends the open streams of a tenant no longer served), `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
Expand All @@ -38,7 +38,8 @@ Eighteen internal packages under `internal/` (plus `internal/testutil/` for shar
- **`dedupe/`** β€” `Deduplicator` interface β†’ `Embedded` (Pebble: every tenant's seen ids in one instance at `data_dir/pebble`, each key led by its tenant, open while any tenant's store is β€” the layout is the implementation's call, and the wiring hands it `data_dir` once; its `Stats` feed the system gauges), wrapped by `Managed` whose open/closed state follows the hot-reloadable `dedupe.enabled` in the settings directory's `config.json`; `Stores` holds one `Managed` per tenant, built through a `Factory` (`func(tenant.ID) *Managed`, `Embedded.Tenant` in production; `Managed` opens its store through a function, so every backend gets the same switch), and reconciled from the registry's `AfterAdopt` hook β€” open exactly when the tenant is served with its switch on, closed with its seen ids kept otherwise ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) stories 7 and 3)
- **`discovery/`** β€” `SchemaRegistry`, one per served tenant over a `Source` read once per refresh β€” the tenant's pool's connection and the database that pool was opened for, one snapshot, so a refused move keeps discovering the database the tenant's queries still use (`internal/app`'s `discoveries` builds, runs and stops them from `AfterAdopt` and `App.Close`: `RetryRefresh` until the first success, then `StartAutoRefresh` with a random first tick; `Lookup` answers `ErrNotLoaded` before the first success β€” the handlers' `503` with `Retry-After` β€” and `ErrUnknownTable` after; a failed loop attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`), that introspects ClickHouse `system.columns` (name/type/nullability plus `default_expression` and 1-based `position`) and `system.tables` (each table's `create_table_query`, kept in-process and never serialized β€” an external-engine table renders its wiring there unconditionally β€” endpoint, bucket/host, database, username, S3 access key id; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so the exposure is the topology, not the secret), records the server version, + `Validate()` for ingest payloads + `CanonicalizeTimestamps()` rewriting top-level `DateTime`/`DateTime64` column values to the canonical RFC 3339 UTC wire form pre-publish (Key Design Decision #19)
- **`ingest/`** β€” Ingest worker pipeline (`worker.go`: JetStream input β†’ per-table batch INSERT with DLQ output). The pipeline is **insert-only**. The wire format `EventMessage` (`types.go`) carries `{table_name, scope, received_timestamp, format, columns, row}` and nothing else β€” `row` is one positional `JSONCompactEachRow` line and `columns` names its slots, the table's **insertable** columns (a `MATERIALIZED`/`ALIAS` column cannot be named in an `INSERT`); the worker batches per (tenant, table, column list), the tenant read off each message's `mq.Topic`, and inserts each batch into its tenant's own ClickHouse (`chconn.Pools.Target`); the worker accepts whatever table name the envelope carries (table existence was already checked by the HTTP ingest handler, which `404`s an unknown table before publish; the worker doesn't re-validate), then bulk-INSERTs. In the embedded-NATS deployment (the default), the server runs with `DontListen: true` (`internal/mq/embedded.go`), so the only Publishers reachable on the `ingest.>` subjects are in-process Go code β€” today, only the HTTP `/v1/ingest?table={table}` handler. Non-insert mutations (`DELETE`/`UPDATE`/`TRUNCATE`/…) must go through `POST /v1/ops/query` under the admin role (the same `RequireAdmin` gate as the rest of `/v1/ops/*`), so non-admin callers never reach the proxy. A request with no token (or an invalid one) resolves to the `default_role`, which in a production config is not the admin role (setting them equal is a loudly-warned dev-only setting), so it can't reach this endpoint. Plus `Sweeper` (Active Sweeper for NATS message lifecycle) + `EventMessage`/`BufferConsumerName` types (`types.go`)
- **`mq/`** β€” the message-queue boundary: the **only** package that imports NATS/JetStream (Key Design Decision #20). and the only one that knows how the broker works. Everything else addresses events by `Topic{Tenant, Table, Scope}` (a validated tenant id and raw names β€” the tenant leads every subject, `ingest.<tenant>.<table>`, so one wildcard selects a tenant's traffic, and a topic without one is refused) and states intent through the interfaces β€” `Publisher` (`ErrQueueFull` is the backpressure signal), `Subscriber`, `ConsumerManager`/`Consumer`/`ConsumerConfig` (the ingest worker's durable consumer), `DeadLetterer` and `DeadLetterStats` (park a message, count what is parked), `Purger` (drop what is both acked and older than a cutoff β€” the sweeper), `Replayer` (SSE gap-fill) β€” composed into `Broker`, which adds each tenant's byte budget (`SetMaxBytes`/`MaxBytes`: the `mq.max_bytes_gb` reload, which opens a tenant's queue the first time) and `Stats` (the system gauges' source). Every interface speaks per tenant, never per stream: the embedded implementation gives each tenant a queue of its own (a stream pair, `INGEST_<tenant>`/`DLQ_<tenant>`), and nothing outside the package may assume that layout β€” an external implementation may keep one shared stream. Subjects, prefixes, wildcards, token encoding, stream names, sequences, and ack floors are private to the one implementation, `EmbeddedNATS` (`embedded.go`, `subject.go`, `purge.go`); `internal/app` constructs it and hands everything else a `mq.Broker`
- **`keyenc/`** β€” the one escaping composite keys are built from: `Escape` keeps `[A-Za-z0-9_-]` (exactly the tenant-id grammar, so a tenant id is its own escaped form) and writes every other byte as `%XX`, `Unescape` is `url.PathUnescape` (lenient: either hex case, and a byte left unescaped reads as itself, so a `%2D` an earlier build wrote still reads), `Join`/`AppendJoin` escape each field and put a separator between them (they panic on no fields, and on a separator the escaping could write or one outside ASCII) and `Split` reverses them. NATS subjects (`Join`/`Split` after the verbatim tenant) and the cache's namespace tokens use it; changing what it keeps orphans every stored key
- **`mq/`** β€” the message-queue boundary: the **only** package that imports NATS/JetStream (Key Design Decision #20). and the only one that knows how the broker works. Everything else addresses events by `Topic{Tenant, Table, Scope}` (a validated tenant id and raw names β€” the tenant leads every subject, `ingest.<tenant>.<table>`, so one wildcard selects a tenant's traffic, and a topic without one is refused) and states intent through the interfaces β€” `Publisher` (`ErrQueueFull` is the backpressure signal), `Subscriber`, `ConsumerManager`/`Consumer`/`ConsumerConfig` (the ingest worker's durable consumer), `DeadLetterer` and `DeadLetterStats` (park a message, count what is parked), `Purger` (drop what is both acked and older than a cutoff β€” the sweeper), `Replayer` (SSE gap-fill) β€” composed into `Broker`, which adds each tenant's byte budget (`SetMaxBytes`/`MaxBytes`: the `mq.max_bytes_gb` reload, which opens a tenant's queue the first time) and `Stats` (the system gauges' source). Every interface speaks per tenant, never per stream: the embedded implementation gives each tenant a queue of its own (a stream pair, `INGEST_<tenant>`/`DLQ_<tenant>`), and nothing outside the package may assume that layout β€” an external implementation may keep one shared stream. Subjects, prefixes, wildcards, stream names, sequences, and ack floors are private to the one implementation, `EmbeddedNATS` (`embedded.go`, `subject.go`, `purge.go`, `deadletter.go`), whose subject tokens are escaped by the shared `internal/keyenc`; `internal/app` constructs it and hands everything else a `mq.Broker`
- **`observability/`** β€” OpenTelemetry pipeline: `InitProvider` wires trace/metric/log providers via OTLP gRPC (each signal independently gated). A top-level `Prometheus` config block drives an optional `/metrics` scrape endpoint that runs independently of OTLP push β€” standalone (Alloy/Mimir scrape, no collector), alongside OTLP, or off. `NewLogger` produces a slog handler that fans out to stdout AND OTLP (stdout always 100%, OTLP sample-rate-aware). `TraceHandler` injects trace_id/span_id from active spans. `tracer.go` provides W3C trace context propagation over message headers (`InjectHeaders`/`ExtractHeaders` on a plain header map; `internal/mq` injects on every publish and extracts on the `Subscribe` path, so this package never sees a NATS type).
- **`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)
Expand Down Expand Up @@ -434,6 +435,7 @@ internal/config/ β†’ Configuration structs + loader
internal/dedupe/ β†’ Optional deduplication (interface + embedded/distributed)
internal/discovery/ β†’ ClickHouse schema introspection + ingest validation
internal/ingest/ β†’ Batch buffer with DLQ + Active Sweeper (NATS message lifecycle)
internal/keyenc/ β†’ One escaping for composite keys (NATS subject tokens, cache namespace tokens)
internal/mq/ β†’ MQ boundary (the only NATS/JetStream importer: owned message/consumer/stream types + embedded server)
internal/observability/ β†’ OpenTelemetry pipeline (traces/metrics/logs providers, Prometheus exporter, slog fan-out, message-header trace propagation)
internal/pipes/ β†’ Named query pipes (types, parameter binding, Source)
Expand Down
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),

### Changed

- **One escaping for composite keys, `-` kept; dead-letter counts per table** (`internal/keyenc` (new, + tests), `internal/mq/{mq,subject,embedded,deadletter}.go` (+ tests), `internal/query/ident.go` (+ tests), `AGENTS.md`, `docs/src/content/docs/{architecture,development}.md`): part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). NATS subject tokens and the cache's namespace tokens each carried a copy of the same encoder; both now call `internal/keyenc`, which the dedupe keys will use too, and subjects are built with its `Join` (each field escaped, with a separator the escaping never writes between them). The escaping now keeps `-` as well as ASCII letters, digits and `_` β€” exactly the tenant-id grammar β€” so a table or scope such as `my-table` is `my-table` in a subject rather than `my%2Dtable`; every other byte is escaped as before (pinned by golden tests, and against v0.1.0's encoder for every other byte value). Upgrading from v0.1.0 notices nothing further, since its queue is deleted at boot (below). A queue an unreleased build since [#612](https://github.com/Wave-RF/WaveHouse/pull/612) wrote still reads, because decoding is `url.PathUnescape` as it was: `%2D` decodes to `-`, and a dead-letter count merges both forms. Only a `/v1/stream` client resuming across such an upgrade (`Last-Event-ID` or `since`) on a table whose name holds `-` misses that table's events queued before it, since the replay filters on the table's exact subject. `GET /v1/ops/dlq/stats` now counts every scope of a table under the table itself, and `?table=` keeps all of its scopes; a scoped message used to count under `table.scope`, a name a dotted table could share. Scope is always empty today, so the response is unchanged.

- **Each tenant has a message queue of its own** (`internal/mq/{mq,subject,embedded}.go` (+ tests), `internal/ingest/{sweeper,worker}.go` (+ tests), `internal/api/{dlq,ingest}.go` (+ tests), `internal/app/{app,wire}.go` (+ tests), `internal/stream/{subscriber,hub}.go`, `internal/settings/{settings,store,registry}.go` (+ tests), `internal/testutil/{mocks,testutil}.go`, `clients/ts/src/{dlq,types}.ts` (+ tests), `docs/src/content/docs/{deployment,api,architecture,ingest-pipeline,durability,why-wavehouse}.md`, `docs/src/content/docs/sdk/{admin,reference,streaming}.md`, `docs/src/content/docs/{settings-directory,configuration}.mdx`, `AGENTS.md`): story 5b of the multi-tenant epic ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)). The embedded NATS server keeps each tenant's events on a pair of JetStream streams of its own β€” `INGEST_<tenant>` (`ingest.<tenant>.>`, `DiscardNew`) at the tenant's own `mq.max_bytes_gb`, and `DLQ_<tenant>` (`dlq.<tenant>.>`, `DiscardOld`) at a tenth of it β€” opened when the tenant is first served and kept, at the budget it last had, when its folder is rejected or removed; subjects are unchanged, and nothing outside `internal/mq` names a stream. A tenant at its budget gets `503` while the others keep publishing, and the ingest worker's and the hub bridge's durables are held on every tenant's stream, each with its own ack floor and `MaxAckPending`, so one tenant's backlog holds back neither another's delivery nor its purge; the worker's prefetch and the hub bridge's fetch-ahead are each shared across the tenants' streams. A removed or rejected tenant's stream is still consumed, so its queued rows reach the worker and are parked on its own dead-letter queue. The sweeper purges each tenant's stream at that tenant's own `stream.gap_window_minutes` β€” a rejected tenant's as its folder last had it (`settings.Registry.Known` now yields each tenant's last adopted store), and all of its history if the folder has been rejected since boot, so its clients resume once the folder is fixed β€” keeping no acknowledged history for a removed tenant, and every served tenant's `mq.max_bytes_gb` is applied after each reload rather than tenant `0`'s alone (the boot warning about a nested directory with no tenant `0` is gone with it, and so is the tracking of tenant `0`'s last adopted store: a flat directory's ops gate, its one reader left, reads tenant `0` through the registry). One tenant's failed purge holds up no other tenant's, and the sweep logs it at `ERROR` unless every failure in it is a buffer consumer not created yet. A reload that shrinks a budget no longer deletes dead letters: a dead-letter stream holding more than a tenth of the new budget keeps what it holds, and that is logged β€” the interim guard of [#532](https://github.com/Wave-RF/WaveHouse/issues/532). JetStream's own check of the streams' caps against the disk, which made them fit 75% of the free disk at boot together, is lifted: a budget is a cap and never a reservation, what the budgets add up to against the disk is [#138](https://github.com/Wave-RF/WaveHouse/issues/138)'s, and a flat directory whose budget exceeds three quarters of the free disk now boots where it used to be refused. A queue that cannot be opened refuses a flat boot like any other store and, over a nested directory, costs its tenant alone β€” its ingest answers `503`, each reload trying again, and a publish at most once every five seconds, so its clients retrying never hold up another tenant's queue. `GET /v1/ops/dlq/stats` reads one tenant's dead-letter queue β€” the one `?tenant=` names, parsed strictly like the other admin reads, and tenant `0`'s without it, no longer the sum across tenants β€” answering for a rejected or removed tenant too and `404` for a tenant with no queue; the SDK's `wh.dlq.list()` and `.table()` take a `tenant` option. The streams an earlier build kept for every tenant together (`WAVEHOUSE`, `WAVEHOUSE_DLQ`) overlap every tenant's subjects and are deleted at boot, what they held with them, and a subject with no tenant token no longer reads as tenant `0`'s.

- **Message-queue subjects lead with the tenant, and the async paths read it off each message** (`internal/mq/{mq,subject,embedded}.go`, `internal/api/{ingest,stream}.go`, `internal/stream/hub.go`, `internal/ingest/{worker,sweeper}.go`, `internal/app/wire.go`, `docs/src/content/docs/{architecture,ingest-pipeline,deployment,api}.md`, `docs/src/content/docs/{settings-directory,access-control}.mdx`, `AGENTS.md`): story 5a of the multi-tenant epic ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)). `mq.Topic` gains a `Tenant`, the leading token of every subject β€” `ingest.<tenant>.<table>[.<scope>]`, `dlq.<tenant>.<table>[.<scope>]` β€” placed verbatim, since the tenant-id grammar makes it one token, so one wildcard selects a tenant's traffic (`ingest.acme.>`); a settings directory that holds the four files produces the same subjects with `0` as the token, and nothing else about them changes. The ingest and stream handlers address the request's tenant, which the resolved store carries (`settings.Store.Tenant`, from story 8). The stream hub indexes subscribers by the full topic and evaluates each event under its own tenant's policy, so a subscriber on one tenant's table never receives another tenant's rows for a table of the same name; a gap-fill and the opening schema frame read the connection's tenant. The ingest worker reads each message's tenant off its topic, batches per tenant table, and resolves the dead-letter switch under the row's own tenant β€” an envelope it cannot read is parked or dropped under the topic's tenant too β€” and bumps that tenant's cache namespaces (story 8's `invalidate` now receives the message's tenant rather than tenant `0`; the wiring's `sharedTables` still repeats each bump under every tenant the directory holds, since every tenant reads the same ClickHouse table until story 6). `Publish` and gap-fill refuse a topic whose tenant is empty or outside the grammar, so nothing lands on tenant `0` by omission.
Expand Down
Loading
Loading