diff --git a/AGENTS.md b/AGENTS.md index e776f162..ceae4d5d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 @@ -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..`, 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_`/`DLQ_`), 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..
`, 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_`/`DLQ_`), 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) @@ -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) diff --git a/CHANGELOG.md b/CHANGELOG.md index 476f3b69..df19a657 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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_` (`ingest..>`, `DiscardNew`) at the tenant's own `mq.max_bytes_gb`, and `DLQ_` (`dlq..>`, `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..
[.]`, `dlq..
[.]` — 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. diff --git a/docs/src/content/docs/architecture.md b/docs/src/content/docs/architecture.md index c9161b2a..d4ca4c0b 100644 --- a/docs/src/content/docs/architecture.md +++ b/docs/src/content/docs/architecture.md @@ -61,6 +61,7 @@ internal/ ├── dedupe/ Optional deduplication (Pebble) ├── discovery/ ClickHouse schema introspection and validation ├── ingest/ Batch buffering, DLQ, and Active Sweeper +├── keyenc/ The one escaping composite keys are built from (NATS subject tokens, cache namespace tokens) ├── mq/ MQ boundary: the only NATS/JetStream importer (owned message/consumer/stream types + embedded server) ├── observability/ OpenTelemetry pipeline (traces/metrics/logs + Prometheus exposition) ├── pipes/ Named query pipes (NamedQuery type, parameter binding, Source) @@ -146,7 +147,8 @@ The SSE fan-out, factored out of `api/` so the delivery hot path ([#294](https:/ The **only** package that imports NATS/JetStream — a `depguard` rule in `.golangci.yml` fails `make lint` on any `github.com/nats-io` import in every package golangci-lint builds; the `integration`-tagged files under `tests/` sit outside its default build context, so the boundary there rests on convention (AGENTS.md Key Design Decision #20). Every other package talks to the broker through the types below, so a subject, stream, or broker change lands here once. - **mq.go** — The owned surface, stated as intent rather than broker mechanics. `Topic{Tenant, Table, Scope}` is the only address the rest of the process handles (a validated tenant id and raw names; comparable, so the SSE hub keys its index by the value). `Message` carries `Data`, its topic (`Topic()` decodes the delivered key on demand — the tenant included, which is how the hub bridge and the worker learn whose event it is; `TopicKey()` is the delivered form, for log lines), and the ack family (`DoubleAck(ctx)`, `Ack()`, `Nak()`); `Headers` is the message header map (`Add`/`Set`/`Get`, exact-key) that `PublishOpt`s such as `WithHeader` shape. Interfaces, each speaking per tenant and never per stream: `Publisher` (`ErrQueueFull` when the tenant's ingest queue is at its byte budget or not open yet — the API's 503 + `Retry-After`), `Subscriber` (every ingest event of every tenant, under a named durable consumer, its fetch-ahead split across the tenants — the hub bridge), `ConsumerManager` → `Consumer` (a durable explicit-ack consumer from a `ConsumerConfig`, whose `MaxAckPending` holds per tenant; `Consume` delivers each tenant's messages on a goroutine of that tenant's, in order, so a blocking handler is backpressure on its own tenant alone, spreads the prefetch across the tenants, and returns a `stop` plus a `failed` channel that reports delivery ending on its own — `ErrDeliveryEnded`, e.g. a deleted consumer, a closed connection, or a tenant's queue that could not be joined — since no message would ever say so) for the ingest worker, `DeadLetterer.DeadLetter` (park a message under its own topic, in its tenant's dead-letter queue; the caller acks), `DeadLetterStats.DeadLetterCounts` (one tenant's; `ErrNoDeadLetterQueue` when it has none), `Purger.PurgeAcked` (drop what is both acked by a consumer and stored before its tenant's cutoff, everything acked for a tenant given none; one error per failed tenant, joined — `ErrConsumerNotFound` for a queue the consumer has not been created on yet, the one failure the sweeper logs as a warning rather than an error) for the sweeper, and `Replayer.ReplaySince` for SSE gap-fill. `Broker` composes them with each tenant's byte budget (`SetMaxBytes`/`MaxBytes`), `Stats`, and `Close`; it is what `internal/app` holds. -- **subject.go** — The embedded broker's naming, private to the package: the stream names (`INGEST_` and `DLQ_` — prefixes that differ in their first letter, so no tenant id makes one kind's name the other's — and the one pair an earlier build kept for every tenant together, `WAVEHOUSE`/`WAVEHOUSE_DLQ`, which boot deletes), the `ingest.`/`dlq.` prefixes and `>` wildcards, the subject-token encoder (alphanumerics and `_` pass, everything else is percent-encoded, so a name can never split or wildcard a subject), and `Topic` ↔ subject conversion. A subject is `.
[.]`: the tenant verbatim — its grammar (`tenant.Parse`) makes it one token, and it is checked on the way to the wire, so a topic without one has no subject — then the table and scope as encoded tokens; tenant first so one wildcard selects a tenant's traffic (`ingest.acme.>`). A topic has the same tail on both streams, so parking on the DLQ is a prefix swap on the delivered subject — nothing is decoded or re-encoded — and the tail's first token picks the tenant's stream. +- **subject.go** — The embedded broker's naming, private to the package: the stream names (`INGEST_` and `DLQ_` — prefixes that differ in their first letter, so no tenant id makes one kind's name the other's — and the one pair an earlier build kept for every tenant together, `WAVEHOUSE`/`WAVEHOUSE_DLQ`, which boot deletes), the `ingest.`/`dlq.` prefixes and `>` wildcards, the subject tokens (`internal/keyenc`: ASCII letters, digits, `_` and `-` pass, everything else is percent-encoded, so a name can never split or wildcard a subject), and `Topic` ↔ subject conversion. A subject is `.
[.]`: the tenant verbatim — its grammar (`tenant.Parse`) makes it one token, and it is checked on the way to the wire, so a topic without one has no subject — then the table and scope as encoded tokens; tenant first so one wildcard selects a tenant's traffic (`ingest.acme.>`). A topic has the same tail on both streams, so parking on the DLQ is a prefix swap on the delivered subject — nothing is decoded or re-encoded — and the tail's first token picks the tenant's stream. +- **deadletter.go** — `deadLetterTables`, the per-table count `DeadLetterCounts` reports: a dead-letter stream's per-subject counts, each subject parsed back to its topic and counted under its table — every scope of a table under the table itself, so a dotted table name never shares a count with a table + scope pair — and a table filter keeps that table with all of its scopes. - **purge.go** — The Active Sweeper's arithmetic over JetStream sequences: purge target = `MIN(consumer ack floor + 1, first sequence stored at or after the cutoff)`, the latter found by binary search over message timestamps (~15 lookups). Every uncertainty resolves toward purging less: a sequence that holds no message is kept as a candidate bound rather than discarding the half below it, and a lookup that fails outright aborts that tenant's purge. It runs on each tenant's stream at that tenant's cutoff, and a failure on one tenant's stream is reported without stopping the sweep of the others. Healthy state keeps exactly the gap window; ClickHouse down freezes purging; a catastrophic outage fills the stream to `MaxBytes` and `DiscardNew` pushes back. - **embedded.go** — `EmbeddedNATS`, the one `Broker`: an in-process NATS server with JetStream, giving each tenant a queue of its own — stream `INGEST_` with subjects `ingest..>`, capped at the tenant's `mq.max_bytes_gb` (`DiscardNew`), and stream `DLQ_` (`dlq..>`, `DiscardOld`) at a tenth of it — with the durable consumers on the ingest one; nothing outside the package sees that layout. JetStream's own check of the streams' caps against the disk (75% of the free disk by default) is set out of reach, so a budget is a cap and never a reservation. Boot deletes the pair an earlier build kept for every tenant together — its subjects overlap every tenant's — and takes stock of the tenants' streams on disk with their budgets, so a consumer created later is held on every one, a tenant no longer served included. `SetMaxBytes` opens a tenant's queue the first time — its dead-letter stream first, so no row is queued that could not be parked — and every registered consumer joins it; a publish or park that finds a stream missing reopens it at the budget last asked for the tenant, and so does a publish to a queue the broker has not recorded open — an open that timed out can leave a stream JetStream creates after all, which no consumer holds — either refused as a full queue with none asked yet. Publishes and parks that find the same queue not open share one attempt (`singleflight`), and after one fails the tenant's publishes and parks are refused at once for five seconds rather than each trying again under the broker's lock, which every tenant's open, resize and reload takes; a reload retries regardless. After that `SetMaxBytes` applies a reloaded budget to the tenant's two live streams as a pair: if the DLQ update fails after the ingest one succeeded, the ingest resize is undone so both stay on the previous budget — best effort, since if that undo also fails the ingest stream stays at the new limit and the DLQ at the previous, and the error says so. A dead-letter stream is never capped below the bytes it holds, which `DiscardOld` would delete to fit ([#532](https://github.com/Wave-RF/WaveHouse/issues/532)): it keeps what it holds, and that is logged. Its JetStream calls are bounded to ten seconds — plus ten more for the consumers joining a queue it has just opened, and five for the rollback of a failed resize, each a budget of its own rather than the one that just expired — since a reload holds the settings store's lock while its hooks run; `MaxBytes` reports the budget last applied in full, so a failed resize is retried by the next reload. The consumers `CreateConsumer` and `Subscribe` build hold one durable on each tenant's stream, looked up before anything is written so a boot over many queues writes nothing it need not, each delivering on a goroutine of its own into the one handler. `PurgeAcked` and `DeadLetterCounts` run per tenant stream. `Stats` reports connection and inbound-message counters for `observability.RegisterSystemMetrics`. Trace context rides in the message headers: `Publish` applies `observability.InjectHeaders`, and a message delivered through `Subscribe` (the hub bridge) carries `observability.ExtractHeaders` on its `Ctx`; the worker's `Consumer` path skips the extraction, since it batches across messages and reads no per-message context. @@ -200,6 +202,10 @@ The hot-reloadable half of configuration: a directory of four JSON files (`confi - **chsql.go** — Dependency-free ClickHouse SQL helpers shared by `query/` and `policy/`, kept in their own package to break an import cycle. `QuoteIdent` is the single place every identifier — column, table, alias — becomes SQL text: always backtick-quoted and escaped, so any ClickHouse-legal name (dots, spaces, unicode, keywords) is safe. `BindUnsafe` reports whether a name contains a literal `?`, which would desync clickhouse-go's positional binder; such names are rejected fail-closed rather than silently mis-bound. +### `keyenc/` — Key Escaping + +- **keyenc.go** — The one escaping composite keys are built from, so a name can never be mistaken for a separator: `Escape` keeps ASCII letters, digits, `_` and `-` — exactly the tenant-id grammar, so a tenant id is its own escaped form — and writes every other byte as `%XX` (uppercase hex); `Unescape` is `url.PathUnescape`, which decodes `%XX` in either case and takes any other byte as itself, so a `%2D` for `-` that an earlier build wrote still reads. `Join`/`AppendJoin` escape each field and put a separator between them, panicking on no fields and on a separator the escaping could write or one outside ASCII, and `Split` reverses them. NATS subject tokens (`internal/mq`) and the cache's namespace tokens (`query.SafeEncodeToken`) both use it. Keys built from it are stored, so changing what it keeps orphans them. + ## Data Flows ### Ingest Path diff --git a/docs/src/content/docs/development.md b/docs/src/content/docs/development.md index 01f82e73..16b65a74 100644 --- a/docs/src/content/docs/development.md +++ b/docs/src/content/docs/development.md @@ -454,13 +454,14 @@ WaveHouse/ │ ├── api/ # HTTP handlers, router, middleware │ ├── app/ # Process wiring (build every component, run under one errgroup, release in reverse) │ ├── auth/ # JWT/JWKS authentication middleware -│ ├── cache/ # L1 (Ristretto) + L2 caching +│ ├── cache/ # Query cache: Ristretto L1 + the tenant-led version index │ ├── chconn/ # ClickHouse pools, one per connection tuple (reconciled on settings reload) │ ├── chsql/ # Shared ClickHouse SQL helpers (quoting + bind-safety) │ ├── config/ # YAML + env var configuration │ ├── dedupe/ # Optional deduplication (Pebble) │ ├── discovery/ # ClickHouse schema introspection + validation │ ├── ingest/ # Batch buffering + DLQ + Active Sweeper +│ ├── keyenc/ # One escaping for composite keys (NATS subject tokens, cache namespace tokens) │ ├── mq/ # MQ boundary: the only NATS/JetStream importer │ ├── observability/ # OpenTelemetry pipeline (traces/metrics/logs + Prometheus) │ ├── pipes/ # Named query pipes (types + parameter binding) diff --git a/internal/keyenc/keyenc.go b/internal/keyenc/keyenc.go new file mode 100644 index 00000000..0e3434bc --- /dev/null +++ b/internal/keyenc/keyenc.go @@ -0,0 +1,100 @@ +// Package keyenc is the one escaping composite WaveHouse keys are built +// from: NATS subject tokens and cache namespace tokens. A field keeps ASCII +// letters, digits, '_' and '-' as they are and writes every other byte as %XX +// (uppercase hex), so no separator, wildcard, whitespace, brace or non-ASCII +// byte ever appears in it unescaped, and any table name ClickHouse accepts +// encodes. The bytes it keeps are exactly a tenant id's (tenant.Parse), so a +// tenant id is its own escaped form. +// +// Keys built from it are stored — queued under NATS subjects, held in caches +// — so a change to what it keeps orphans them. Earlier builds escaped '-' as +// %2D; Unescape still reads that form. +package keyenc + +import ( + "fmt" + "net/url" + "strings" +) + +const upperHex = "0123456789ABCDEF" + +// kept reports whether b is written as itself. +func kept(b byte) bool { + return (b >= 'a' && b <= 'z') || (b >= 'A' && b <= 'Z') || (b >= '0' && b <= '9') || b == '_' || b == '-' +} + +// Escape encodes s as one field. +func Escape(s string) string { + for i := 0; i < len(s); i++ { + if !kept(s[i]) { + return string(AppendEscape(make([]byte, 0, len(s)+2*(len(s)-i)), s)) + } + } + return s +} + +// AppendEscape appends Escape(s) to dst. +func AppendEscape(dst []byte, s string) []byte { + for i := 0; i < len(s); i++ { + b := s[i] + if kept(b) { + dst = append(dst, b) + } else { + dst = append(dst, '%', upperHex[b>>4], upperHex[b&0x0F]) + } + } + return dst +} + +// Unescape reverses Escape. It is url.PathUnescape: %XX in either hex case +// decodes, and any other byte reads as itself, so a field another writer +// left partly unescaped — an earlier build's %2D included — still reads. +func Unescape(s string) (string, error) { + return url.PathUnescape(s) +} + +// checkSep panics unless sep can separate escaped fields: a byte Escape never +// writes, and ASCII, so the key stays valid UTF-8. +func checkSep(sep byte) { + if kept(sep) || sep == '%' || sep >= 0x80 { + panic(fmt.Sprintf("keyenc: %q cannot separate fields", sep)) + } +} + +// Join escapes each field and joins them with sep. It panics on no fields, +// whose key would be one empty field's, and on a separator Escape could +// write. +func Join(sep byte, fields ...string) string { + return string(AppendJoin(nil, sep, fields...)) +} + +// AppendJoin appends Join(sep, fields...) to dst. +func AppendJoin(dst []byte, sep byte, fields ...string) []byte { + checkSep(sep) + if len(fields) == 0 { + panic("keyenc: Join needs at least one field") + } + for i, f := range fields { + if i > 0 { + dst = append(dst, sep) + } + dst = AppendEscape(dst, f) + } + return dst +} + +// Split reverses Join: the fields of key, each unescaped. It panics on a +// separator Join would refuse. +func Split(key string, sep byte) ([]string, error) { + checkSep(sep) + parts := strings.Split(key, string([]byte{sep})) + for i, p := range parts { + f, err := Unescape(p) + if err != nil { + return nil, err + } + parts[i] = f + } + return parts, nil +} diff --git a/internal/keyenc/keyenc_test.go b/internal/keyenc/keyenc_test.go new file mode 100644 index 00000000..877599e2 --- /dev/null +++ b/internal/keyenc/keyenc_test.go @@ -0,0 +1,176 @@ +package keyenc_test + +import ( + "bytes" + "fmt" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/Wave-RF/WaveHouse/internal/keyenc" + "github.com/Wave-RF/WaveHouse/internal/tenant" +) + +// v010Escape is the encoder v0.1.0 shipped as query.SafeEncodeNATS, copied +// verbatim. Escape differs from it only in keeping '-'. +func v010Escape(raw string) string { + var buf bytes.Buffer + for i := 0; i < len(raw); i++ { + b := raw[i] + if (b >= 'a' && b <= 'z') || (b >= 'A' && b <= 'Z') || (b >= '0' && b <= '9') || b == '_' { + buf.WriteByte(b) + } else { + fmt.Fprintf(&buf, "%%%02X", b) + } + } + return buf.String() +} + +func TestEscape_Golden(t *testing.T) { + t.Parallel() + for raw, want := range map[string]string{ + "": "", + "my_table123": "my_table123", + "default.clicks": "default%2Eclicks", + "my table": "my%20table", + "a-b/c": "a-b%2Fc", + "a.*.>": "a%2E%2A%2E%3E", + "100%": "100%25", + "{acme}:x|y": "%7Bacme%7D%3Ax%7Cy", + "a\x00b": "a%00b", + "café": "caf%C3%A9", + "\xff": "%FF", + "tab\tnewline\n": "tab%09newline%0A", + "#hash": "%23hash", + "ABCxyz_0189": "ABCxyz_0189", + "evt-123": "evt-123", + "日本": "%E6%97%A5%E6%9C%AC", + } { + assert.Equal(t, want, keyenc.Escape(raw), "%q", raw) + assert.Equal(t, want, string(keyenc.AppendEscape([]byte("x"), raw))[1:], "%q", raw) + } +} + +// Every byte value, alone and between kept bytes, encodes as v0.1.0 did, +// but for '-'. +func TestEscape_MatchesV010ButDash(t *testing.T) { + t.Parallel() + for b := range 256 { + for _, s := range []string{string([]byte{byte(b)}), "a" + string([]byte{byte(b)}) + "Z"} { + want := v010Escape(s) + if b == '-' { + want = s + } + require.Equal(t, want, keyenc.Escape(s), "byte %#x", b) + } + } +} + +// A tenant id is its own escaped form: the kept bytes are its grammar. +func TestEscape_KeepsExactlyTheTenantGrammar(t *testing.T) { + t.Parallel() + for b := range 256 { + s := string([]byte{byte(b)}) + _, err := tenant.Parse(s) + assert.Equal(t, err == nil, keyenc.Escape(s) == s, "byte %#x", b) + } +} + +func TestEscape_NoAllocWhenNothingToEscape(t *testing.T) { + s := "events_2026-09" + assert.Zero(t, testing.AllocsPerRun(100, func() { _ = keyenc.Escape(s) })) +} + +// Distinct names never share an escaped form, even names that look escaped: +// '%' is itself escaped. +func TestEscape_LookalikesStayDistinct(t *testing.T) { + t.Parallel() + names := []string{"b-c", "b%2Dc", "b%2dc", "b.c", "b%2Ec"} + seen := map[string]string{} + for _, n := range names { + e := keyenc.Escape(n) + require.NotContains(t, seen, e, "%q and %q", seen[e], n) + seen[e] = n + back, err := keyenc.Unescape(e) + require.NoError(t, err) + assert.Equal(t, n, back) + } +} + +func TestUnescape(t *testing.T) { + t.Parallel() + for in, want := range map[string]string{ + "": "", + "plain": "plain", + "default%2Eclicks": "default.clicks", + "lower%2ecase": "lower.case", + "evt%2D123": "evt-123", // an earlier build's form + "%00%FF": "\x00\xff", + "a+b": "a+b", + } { + got, err := keyenc.Unescape(in) + require.NoError(t, err, "%q", in) + assert.Equal(t, want, got, "%q", in) + } + for _, bad := range []string{"%", "%2", "a%2Gb", "%%41", "x%"} { + _, err := keyenc.Unescape(bad) + require.Error(t, err, "%q", bad) + } +} + +func TestJoinSplit(t *testing.T) { + t.Parallel() + assert.Equal(t, "acme/clicks/evt-123", keyenc.Join('/', "acme", "clicks", "evt-123")) + assert.Equal(t, "a%2Fb/c", keyenc.Join('/', "a/b", "c"), "a separator inside a field is escaped") + assert.Equal(t, "a..", keyenc.Join('.', "a", "", "")) + assert.Equal(t, "", keyenc.Join('.', ""), "one empty field") + assert.Equal(t, "p:a", string(keyenc.AppendJoin([]byte("p:"), '/', "a"))) + + for _, fields := range [][]string{{"acme", "a/b", "id"}, {"", "", ""}, {""}, {"%", "/", "%2F"}, {"x"}} { + got, err := keyenc.Split(keyenc.Join('/', fields...), '/') + require.NoError(t, err) + assert.Equal(t, fields, got) + } + _, err := keyenc.Split("a/%zz", '/') + require.Error(t, err) +} + +func TestJoin_RefusesZeroFields(t *testing.T) { + t.Parallel() + assert.Panics(t, func() { keyenc.Join('/') }) + assert.Panics(t, func() { keyenc.AppendJoin(nil, '/') }) +} + +// A separator Escape could write, or one outside ASCII, is refused by Join +// and Split alike; every other byte separates. +func TestSeparators(t *testing.T) { + t.Parallel() + for b := range 256 { + sep := byte(b) + if keyenc.Escape(string([]byte{sep})) == string([]byte{sep}) || sep == '%' || sep >= 0x80 { + assert.Panics(t, func() { keyenc.Join(sep, "x") }, "%#x", b) + assert.Panics(t, func() { _, _ = keyenc.Split("x", sep) }, "%#x", b) + continue + } + fields := []string{"a", string([]byte{sep}), "b" + string([]byte{sep}) + "c"} + got, err := keyenc.Split(keyenc.Join(sep, fields...), sep) + require.NoError(t, err, "%#x", b) + assert.Equal(t, fields, got, "%#x", b) + } +} + +func FuzzEscapeRoundTrip(f *testing.F) { + for _, s := range []string{"", "a.b", "\x00", "café", "%", "a/b c", "b-c", "b%2Dc"} { + f.Add(s) + } + f.Fuzz(func(t *testing.T, s string) { + enc := keyenc.Escape(s) + require.Equal(t, strings.ReplaceAll(v010Escape(s), "%2D", "-"), enc) + require.False(t, strings.ContainsAny(enc, "./:{}|# *>\x00"), enc) + dec, err := keyenc.Unescape(enc) + require.NoError(t, err) + require.Equal(t, s, dec) + }) +} diff --git a/internal/mq/deadletter.go b/internal/mq/deadletter.go new file mode 100644 index 00000000..e6cc98e1 --- /dev/null +++ b/internal/mq/deadletter.go @@ -0,0 +1,18 @@ +package mq + +// deadLetterTables counts a dead-letter stream's parked messages per table, +// from its per-subject counts under prefix. Every scope of a table counts +// under the table itself, so no table + scope pair can share a count with a +// dotted table name. A non-empty table keeps that table alone, all of its +// scopes included. +func deadLetterTables(subjects map[string]uint64, prefix, table string) map[string]uint64 { + tables := make(map[string]uint64, len(subjects)) + for subj, n := range subjects { + t := parseTopicKey(topicKey(prefix, subj)) + if table != "" && t.Table != table { + continue + } + tables[t.Table] += n + } + return tables +} diff --git a/internal/mq/deadletter_test.go b/internal/mq/deadletter_test.go new file mode 100644 index 00000000..aed76094 --- /dev/null +++ b/internal/mq/deadletter_test.go @@ -0,0 +1,24 @@ +package mq + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestDeadLetterTables(t *testing.T) { + t.Parallel() + subjects := map[string]uint64{ + "dlq.0.a%2Eb": 1, // table "a.b" + "dlq.0.a.b": 2, // table "a", scope "b" + "dlq.0.a": 4, // table "a", unscoped + "dlq.0.my-t": 8, + "dlq.0.my%2Dt.org-1": 16, // an earlier build's escaping of '-', scoped + "dlq.0.clicks.org%2E": 32, + } + assert.Equal(t, map[string]uint64{"a.b": 1, "a": 6, "my-t": 24, "clicks": 32}, + deadLetterTables(subjects, dlqPrefix, ""), "a dotted table never shares a count with a table + scope") + assert.Equal(t, map[string]uint64{"a": 6}, deadLetterTables(subjects, dlqPrefix, "a"), "the filter keeps every scope of its table") + assert.Equal(t, map[string]uint64{"a.b": 1}, deadLetterTables(subjects, dlqPrefix, "a.b")) + assert.Empty(t, deadLetterTables(subjects, dlqPrefix, "never_failed")) +} diff --git a/internal/mq/embedded.go b/internal/mq/embedded.go index 98971492..f314840d 100644 --- a/internal/mq/embedded.go +++ b/internal/mq/embedded.go @@ -1003,9 +1003,9 @@ func (e *EmbeddedNATS) PurgeAcked(ctx context.Context, consumer string, olderTha } // DeadLetterCounts reads tenant id's dead-letter stream's per-subject counts -// and keys them by table. The table filter matches that table's unscoped -// subject, so it is applied to the parsed topic rather than as a subject -// filter; a scoped topic counts under "table.scope". +// and keys them by table (deadLetterTables). The table filter matches every +// scope of that table, so it is applied to the parsed topic rather than as a +// subject filter. func (e *EmbeddedNATS) DeadLetterCounts(ctx context.Context, id tenant.ID, table string) (DeadLetterCounts, error) { if _, err := tenant.Parse(string(id)); err != nil { return DeadLetterCounts{}, fmt.Errorf("tenant: %w", err) @@ -1023,20 +1023,7 @@ func (e *EmbeddedNATS) DeadLetterCounts(ctx context.Context, id tenant.ID, table return DeadLetterCounts{}, fmt.Errorf("dlq stream info: %w", err) } - counts := DeadLetterCounts{Tables: make(map[string]uint64, len(state.Subjects)), Total: state.Msgs} - for subj, n := range state.Subjects { - t := parseTopicKey(topicKey(dlqPrefix, subj)) - if table != "" && (t.Table != table || t.Scope != "") { - continue - } - name := t.Table - if t.Scope != "" { - // TODO(#235): break scopes out rather than fold them into the name. - name += "." + t.Scope - } - counts.Tables[name] += n - } - return counts, nil + return DeadLetterCounts{Tables: deadLetterTables(state.Subjects, dlqPrefix, table), Total: state.Msgs}, nil } // ReplaySince creates an ephemeral consumer on topic's ingest subject, in its diff --git a/internal/mq/mq.go b/internal/mq/mq.go index 3f1c45c1..57eae51c 100644 --- a/internal/mq/mq.go +++ b/internal/mq/mq.go @@ -13,6 +13,7 @@ import ( "errors" "time" + "github.com/Wave-RF/WaveHouse/internal/keyenc" "github.com/Wave-RF/WaveHouse/internal/observability" "github.com/Wave-RF/WaveHouse/internal/tenant" ) @@ -33,16 +34,16 @@ type Topic struct { // key is the injective string form of the topic that a subject's tail // carries: the tenant first, verbatim — its grammar makes it one token — then -// the table and scope as encoded tokens. A topic without a tenant has no -// subject, and its key parses back to a topic of no tenant with the whole key -// as its table (parseTopicKey's fallback). Callers key their own maps by the -// Topic value itself. +// the table and scope joined as escaped tokens (keyenc.AppendJoin). A topic +// without a tenant has no subject, and its key parses back to a topic of no +// tenant with the whole key as its table (parseTopicKey's fallback). Callers +// key their own maps by the Topic value itself. func (t Topic) key() string { - key := string(t.Tenant) + "." + encodeToken(t.Table) - if t.Scope != "" { - key += "." + encodeToken(t.Scope) + key := append([]byte(t.Tenant), '.') + if t.Scope == "" { + return string(keyenc.AppendJoin(key, '.', t.Table)) } - return key + return string(keyenc.AppendJoin(key, '.', t.Table, t.Scope)) } // Message represents a message received from the queue. @@ -246,8 +247,8 @@ type DeadLetterer interface { // DeadLetterCounts is what is parked on one tenant's dead-letter queue. type DeadLetterCounts struct { // Tables maps table name → parked messages, for the tables asked about. - // Scope is not broken out yet (it is inert until #235): a message parked - // under a scoped topic counts under "table.scope", not under its table. + // Every scope of a table counts under the table; scope is not broken out + // yet (it is inert until #235). Tables map[string]uint64 // Total is every parked message of the tenant, whatever the filter. Total uint64 @@ -262,8 +263,8 @@ var ErrNoDeadLetterQueue = errors.New("dead-letter queue not found") type DeadLetterStats interface { // DeadLetterCounts counts tenant id's parked messages per table — a // tenant served, rejected, or removed alike, for as long as its queue is - // kept. A non-empty table narrows Tables to that one (its unscoped - // messages). + // kept. A non-empty table narrows Tables to that one (all of its + // scopes). DeadLetterCounts(ctx context.Context, id tenant.ID, table string) (DeadLetterCounts, error) } diff --git a/internal/mq/subject.go b/internal/mq/subject.go index 489ad241..2ee41ea3 100644 --- a/internal/mq/subject.go +++ b/internal/mq/subject.go @@ -1,11 +1,10 @@ package mq import ( - "bytes" "fmt" - "net/url" "strings" + "github.com/Wave-RF/WaveHouse/internal/keyenc" "github.com/Wave-RF/WaveHouse/internal/tenant" ) @@ -59,29 +58,6 @@ func streamTenant(prefix, name string) (tenant.ID, bool) { return id, err == nil } -// encodeToken converts any table or scope name into a safe, single NATS -// subject token. It preserves alphanumerics and underscores, but -// percent-encodes everything else (so '.', ' ', '*' and '>' can never split -// or wildcard a subject). -func encodeToken(raw string) string { - var buf bytes.Buffer - for i := 0; i < len(raw); i++ { - b := raw[i] - if (b >= 'a' && b <= 'z') || (b >= 'A' && b <= 'Z') || (b >= '0' && b <= '9') || b == '_' { - buf.WriteByte(b) - } else { - fmt.Fprintf(&buf, "%%%02X", b) - } - } - return buf.String() -} - -// decodeToken reverses encodeToken. url.PathUnescape handles exactly the %XX -// form encodeToken writes. -func decodeToken(safe string) (string, error) { - return url.PathUnescape(safe) -} - // subject renders a caller's topic under prefix. The tenant is checked // against its grammar here, on the way to the wire: an empty one — a caller // that never set it — must not become a subject of some other tenant's, and @@ -108,24 +84,27 @@ func keyTenant(key string) (tenant.ID, bool) { return id, err == nil } -// parseTopicKey recovers the Topic from a subject tail. Three tokens are -// tenant, table and scope; two are tenant and table. A tail this package -// could not have written — one token, more than three, a token that does not -// decode, a tenant outside the grammar — cannot be split reliably, so the -// whole of it becomes the table of no tenant rather than being dropped. +// parseTopicKey recovers the Topic from a subject tail: the tenant token, +// read verbatim as keyTenant reads it, then one or two escaped tokens, table +// and scope (keyenc.Split). A tail this package could not have written — one +// token, more than three, a token that does not decode, a tenant outside the +// grammar — cannot be split reliably, so the whole of it becomes the table of +// no tenant rather than being dropped. func parseTopicKey(tail string) Topic { - parts := strings.Split(tail, ".") - switch len(parts) { - case 2, 3: - id, idErr := tenant.Parse(parts[0]) - table, tableErr := decodeToken(parts[1]) - scope, scopeErr := "", error(nil) - if len(parts) == 3 { - scope, scopeErr = decodeToken(parts[2]) - } - if idErr == nil && tableErr == nil && scopeErr == nil { - return Topic{Tenant: id, Table: table, Scope: scope} - } + first, rest, ok := strings.Cut(tail, ".") + if !ok { + return Topic{Table: tail} + } + id, idErr := tenant.Parse(first) + fields, fieldsErr := keyenc.Split(rest, '.') + if idErr != nil || fieldsErr != nil { + return Topic{Table: tail} + } + switch len(fields) { + case 1: + return Topic{Tenant: id, Table: fields[0]} + case 2: + return Topic{Tenant: id, Table: fields[0], Scope: fields[1]} } return Topic{Table: tail} } diff --git a/internal/mq/subject_test.go b/internal/mq/subject_test.go index 67536e4e..6f859acd 100644 --- a/internal/mq/subject_test.go +++ b/internal/mq/subject_test.go @@ -9,56 +9,42 @@ import ( "github.com/stretchr/testify/require" ) -func TestEncodeToken(t *testing.T) { +// The subjects are pinned byte for byte: an embedded broker holds messages +// under them across an upgrade. +func TestSubject_Golden(t *testing.T) { t.Parallel() - tests := []struct { - name string - raw string - expected string + for _, tt := range []struct { + topic Topic + want string }{ - {"safe string", "my_table123", "my_table123"}, - {"with dots", "default.clicks", "default%2Eclicks"}, - {"with spaces", "my table", "my%20table"}, - {"with dashes and slashes", "a-b/c", "a%2Db%2Fc"}, - {"wildcards cannot survive", "a.*.>", "a%2E%2A%2E%3E"}, - {"empty string", "", ""}, - {"only safe characters", "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789_", "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789_"}, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - t.Parallel() - assert.Equal(t, tt.expected, encodeToken(tt.raw)) - }) + {Topic{Tenant: "0", Table: "events"}, "0.events"}, + {Topic{Tenant: "acme-co", Table: "default.clicks", Scope: "org_1"}, "acme-co.default%2Eclicks.org_1"}, + {Topic{Tenant: "a", Table: "a.*.>", Scope: "*"}, "a.a%2E%2A%2E%3E.%2A"}, + {Topic{Tenant: "a", Table: "my table", Scope: "tab\there"}, "a.my%20table.tab%09here"}, + {Topic{Tenant: "a", Table: "table-with-dashes", Scope: "org-1"}, "a.table-with-dashes.org-1"}, + {Topic{Tenant: "a", Table: "100%", Scope: "a/b"}, "a.100%25.a%2Fb"}, + {Topic{Tenant: "a", Table: "{acme}:x"}, "a.%7Bacme%7D%3Ax"}, + {Topic{Tenant: "a", Table: "nul\x00", Scope: "\xff"}, "a.nul%00.%FF"}, + {Topic{Tenant: "a", Table: "caf\u00e9", Scope: "\u65e5"}, "a.caf%C3%A9.%E6%97%A5"}, + {Topic{Tenant: "a", Table: ""}, "a."}, + {Topic{Tenant: "a", Table: "t", Scope: ""}, "a.t"}, + } { + assert.Equal(t, tt.want, tt.topic.key(), "%+v", tt.topic) + for _, prefix := range []string{ingestPrefix, dlqPrefix} { + subj, err := subject(prefix, tt.topic) + require.NoError(t, err) + assert.Equal(t, prefix+tt.want, subj) + } } } -func TestDecodeToken(t *testing.T) { +// A token another writer left partly unescaped, or escaped in lowercase, +// still reads as it always did — and so does an earlier build's %2D for '-', +// so a message it queued reads as the same topic. +func TestParseTopicKey_LenientTokens(t *testing.T) { t.Parallel() - tests := []struct { - name string - safe string - expected string - wantErr bool - }{ - {"safe string", "my_table123", "my_table123", false}, - {"encoded dots", "default%2Eclicks", "default.clicks", false}, - {"encoded spaces", "my%20table", "my table", false}, - {"invalid percent encoding", "default%2Gclicks", "", true}, // %2G is not valid hex - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - t.Parallel() - got, err := decodeToken(tt.safe) - if tt.wantErr { - assert.Error(t, err) - } else { - require.NoError(t, err) - assert.Equal(t, tt.expected, got) - } - }) - } + assert.Equal(t, Topic{Tenant: "a", Table: "b~c", Scope: "d.e"}, parseTopicKey("a.b~c.d%2ee")) + assert.Equal(t, parseTopicKey("a.table-with-dashes.org-1"), parseTopicKey("a.table%2Dwith%2Ddashes.org%2D1")) } func TestSubject_RoundTripsEveryTopic(t *testing.T) { @@ -136,6 +122,7 @@ func TestParseTopicKey_ForeignTailKeepsItself(t *testing.T) { "a.b.c.d", // more tokens than any topic renders "0.bad%2Gtoken", // a token that does not decode "a%2Eb.events", // a tenant outside the grammar + "a%2Db.events", // a tenant token is read verbatim, never decoded ".events", // a topic whose tenant was never set "events", // one token: no tenant leads it "bad%2G", // one token that does not decode diff --git a/internal/query/ident.go b/internal/query/ident.go index 320fe097..35f032da 100644 --- a/internal/query/ident.go +++ b/internal/query/ident.go @@ -1,24 +1,8 @@ package query -import ( - "bytes" - "fmt" -) +import "github.com/Wave-RF/WaveHouse/internal/keyenc" -// SafeEncodeToken converts any table or scope name into a single dot-free -// token, for composing the cache's dotted namespace keys. It preserves -// alphanumerics and underscores, but percent-encodes everything else. -func SafeEncodeToken(raw string) string { - var buf bytes.Buffer - for i := 0; i < len(raw); i++ { - b := raw[i] - // Pass through safe characters: a-z, A-Z, 0-9, and _ (underscore) - if (b >= 'a' && b <= 'z') || (b >= 'A' && b <= 'Z') || (b >= '0' && b <= '9') || b == '_' { - buf.WriteByte(b) - } else { - // Hex encode everything else (e.g., '.' becomes '%2E', ' ' becomes '%20') - fmt.Fprintf(&buf, "%%%02X", b) - } - } - return buf.String() -} +// SafeEncodeToken renders a table or scope name as one dot-free token of the +// cache's namespace keys: keyenc's escaping, the same bytes a NATS subject +// carries for the name. +func SafeEncodeToken(raw string) string { return keyenc.Escape(raw) } diff --git a/internal/query/ident_test.go b/internal/query/ident_test.go index 572b005e..ac4fd64a 100644 --- a/internal/query/ident_test.go +++ b/internal/query/ident_test.go @@ -16,7 +16,7 @@ func TestEncodeTable(t *testing.T) { {"safe string", "my_table123", "my_table123"}, {"with dots", "default.clicks", "default%2Eclicks"}, {"with spaces", "my table", "my%20table"}, - {"with dashes and slashes", "a-b/c", "a%2Db%2Fc"}, + {"with dashes and slashes", "a-b/c", "a-b%2Fc"}, {"empty string", "", ""}, {"only safe characters", "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789_", "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789_"}, }