Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
ef93154
feat(mq): give every tenant a queue of its own
taitelee Sep 24, 2026
8ad9865
fix(mq): reopen detached from the request, and the upgrade runbook swept
taitelee Sep 24, 2026
650a28e
fix(app): boot opens queues under New's context; docs review fixes
taitelee Sep 24, 2026
999c1db
fix(mq): boot re-applies a budget to a split queue pair; review fixes
taitelee Sep 24, 2026
db2d20a
fix(mq): a failed resize restores the ingest stream's own cap; review…
taitelee Sep 25, 2026
07c6a91
fix(app): a rejected tenant keeps its replay history; review fixes
taitelee Sep 25, 2026
04aad9f
fix(mq): share the hub bridge's fetch-ahead across tenants; review fixes
taitelee Sep 25, 2026
57870c4
fix(mq): publish only into a queue the broker recorded open; review f…
taitelee Sep 25, 2026
7cdb794
docs(mq): a consumer that cannot join a queue opened at runtime; revi…
taitelee Sep 25, 2026
e199a03
fix(mq): pace publish-side retries of a queue that cannot open; revie…
taitelee Sep 25, 2026
55137d8
test(mq): one tenant's failed purge stops no other; review fixes
taitelee Sep 25, 2026
060ca9e
docs(mq): size a dead-letter stream the shrink guard kept; review fixes
taitelee Sep 25, 2026
bb027da
test(mq): one conformance suite for every Broker
EricAndrechek Sep 25, 2026
1adc286
test(mq): run the embedded conformance in a test binary of its own
EricAndrechek Sep 25, 2026
80d6c22
docs(mq): keep the per-tenant no-queue case in ErrQueueFull's contract
EricAndrechek Sep 25, 2026
bf11ecc
test(mq): pin the exactly-once failed report through durable deletion
EricAndrechek Sep 25, 2026
ef4153c
docs(api): list the unavailable broker among the request aborts
EricAndrechek Sep 25, 2026
9344f1b
fix(mq): pace a park's reopen, and warn only on a missing consumer; r…
taitelee Sep 25, 2026
8792d89
test(mq): end delivery only once the pulls are live
EricAndrechek Sep 25, 2026
3425243
Merge remote-tracking branch 'origin/mq-tenant-streams' into feat/mq-…
EricAndrechek Sep 25, 2026
e2434d6
fix(mq): a replay whose connection closed is not caught up
EricAndrechek Sep 25, 2026
5570296
Merge remote-tracking branch 'origin/main' into feat/mq-conformance
EricAndrechek Sep 25, 2026
22d30b0
Merge remote-tracking branch 'origin/main' into feat/mq-conformance
EricAndrechek Sep 26, 2026
6ebfc6f
test(mq): nothing parked is an empty Tables map, never nil
EricAndrechek Sep 26, 2026
ff67ab5
Merge remote-tracking branch 'origin/main' into feat/mq-conformance
EricAndrechek Sep 26, 2026
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
3 changes: 3 additions & 0 deletions .testcoverage.yml
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,9 @@ exclude:
# The coord conformance suite: test helpers every Coordinator's tests
# run, imported only from *_test.go like testutil.
- ^internal/coord/coordtest/
# internal/mq/mqtest/ is the Broker conformance suite: test code that
# lives outside *_test.go only so each backend's tests can import it.
- ^internal/mq/mqtest/
- ^tests/
# scripts/ holds Go helpers (cov, orchestrator) that drive the build but
# aren't part of the shipped binary; they show up in `-coverpkg=./...`
Expand Down
4 changes: 2 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared
- **`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`)
- **`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`
- **`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, `ErrUnavailable` a broker that cannot be reached β€” both a `503`, with `Retry-After` `30` and `5`), `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`. Every implementation passes the conformance suite in `internal/mq/mqtest` (`mqtest.Run`), which states the `Broker` contract as behavior; a new backend runs it from its own test, with `mqtest.Caps` only where its semantics legitimately differ
- **`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 @@ -438,7 +438,7 @@ internal/dedupe/ β†’ Optional deduplication (interface + embedded/distrib
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/mq/ β†’ MQ boundary (the only NATS/JetStream importer: owned message/consumer/stream types + embedded server; mqtest/ is the Broker conformance suite)
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)
internal/policy/ β†’ Access control policies (types, evaluation, Source)
Expand Down
Loading
Loading