diff --git a/.testcoverage.yml b/.testcoverage.yml index 5af6e5e5..99d3b1cd 100644 --- a/.testcoverage.yml +++ b/.testcoverage.yml @@ -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=./...` diff --git a/AGENTS.md b/AGENTS.md index 7cc4e23c..f6580a22 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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..`, 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` +- **`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, `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_`/`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`. 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) @@ -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) diff --git a/CHANGELOG.md b/CHANGELOG.md index 38c1c4db..ccdd58bb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Added +- **One conformance suite for every `mq.Broker`, and a transient broker failure is a `503`** (`internal/mq/mqtest/` (new: the suite and the embedded broker's run of it), `internal/mq/{mq,embedded}.go` (+ tests), `internal/api/ingest.go` (+ tests), `.testcoverage.yml`, `docs/src/content/docs/{api,architecture}.md`, `AGENTS.md`): the first piece of the external-NATS workstream of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). `mqtest.Run` states the `Broker` contract as behavior — round trips with names that need encoding, per-tenant order, `Nak` and `AckWait` redelivery, the trace context reaching `Subscribe`, dead-lettering that keeps the topic and leaves the original unacked, per-tenant dead-letter counts, replay bounds and isolation, and exactly one `failed` report when delivery ends underneath a consumer — through the interfaces alone, so the external backend runs the same cases, with `mqtest.Caps` for the four places where its semantics legitimately differ. The embedded broker passes it; writing it turned up that a durable deleted on several tenants' queues could report on `failed` more than once, which is fixed and pinned by a test that deletes it on one queue after another. It also turned up a replay that lost its connection mid-pull passing for a caught-up one when the pull ended in a timeout; that is an error now, as the `Replayer` contract says. The interface comments now allow a delivery unit that is a partition holding several tenants, a `CreateConsumer` that finds a durable rather than creating one, a `PurgeAcked` that leaves retention to the operator, and zero dead-letter counts where there is no per-tenant queue. A new sentinel, `mq.ErrUnavailable`, is a broker that cannot be reached or does not answer in time: the ingest handler answers it with `503` + `Retry-After: 5` rather than the `500` "publish failed" it would have been. Nothing returns it yet; the external backend of #613 will. - **Process roles: the API and the background workers can run in separate processes** (`internal/config/config.go` (+ `roles_test.go`, `defaults_test.go`), `internal/config/backends.go`, `internal/app/{app,wire}.go` (+ `roles_test.go`), `internal/api/router.go` (+ tests), `tests/integration/{setup,tenants}_test.go`, `config.yaml`, `docs/src/content/docs/{configuration.mdx,deployment.md,architecture.md}`, `AGENTS.md`), part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). The boot config gains `roles` (`WH_ROLES`, default `api,ingest,sweeper`, set in `defaults()` like every boot default, so an explicit `roles: []` refuses boot) and `instance_id` (`WH_INSTANCE_ID`, default `-<8 hex>`, fresh at every boot; logged at boot, and recorded as a lease holder once a shared `coord.backend` exists). A process wires only what its roles need: `api` runs the HTTP API with schema discovery, the token verifiers, the dedupe stores and the SSE hub (all per API process); `ingest` runs the ingest worker; `sweeper` runs the sweeper under its lease. A process without `api` serves an ops-only listener on `server.port` (the probes and their aliases, `/version`, the same-port metrics path, and `POST /v1/ops/settings/reload`, which takes the operator key alone); every other route answers 404, under `/v1/ops` once the operator key has passed. Boot refuses any split over the embedded MQ, which no other process can reach, and `api` without `ingest` (or the reverse) over a local cache, which the ingest worker's invalidations would never reach. Until a shared `mq.backend` exists, every process therefore runs every role, which is the default, so nothing changes for an existing deployment. `data_dir` is probed for Pebble only in a process running `api`. A `config.Config` built without `config.Load` must now name its roles (`config.AllRoles()` for all of them): `app.New` refuses an empty set. - **Leases for work that must run in one process at a time, and the sweeper runs under one** (`internal/coord/` (new: `coord.go`, `local.go`, `elect.go`, `coordtest/`, + tests), `internal/app/{app,wire}.go` (+ tests), `internal/config/backends.go`, `internal/ingest/sweeper.go`, `.github/labeler.yml`, `.testcoverage.yml`, `docs/src/content/docs/{architecture,development,ingest-pipeline}.md`, `docs/src/content/docs/configuration.mdx`, `config.yaml`, `AGENTS.md`): part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). `coord.Coordinator` hands out named leases (`TryAcquire` → a `Term` with a strictly increasing fencing `Token`, a `Done` channel and `Resign`; `ErrHeld` while another holder's term is live), and `coord.RunElected` runs a loop only while its process holds the lease, resigning when the loop returns and campaigning again every 2s. `coord.Local` is the in-process implementation, and `coordtest.Conformance` is the suite every implementation runs — the NATS KV backend that lets several replicas share one queue comes next. The sweeper now runs through `RunElected` under the `sweeper` lease; with the in-process coordinator the one process always holds it, so nothing changes for a single-process deployment beyond one `coord: elected` log line at startup. `coord.backend` now selects the coordinator (`local`, the only value), so a `config.Config` built without `config.Load` must name it as well as the other three layers' backends. @@ -331,7 +332,7 @@ The first public release. Everything below shipped in it — the sections are gr - **BREAKING (SDK): `PipeRef.fetch` no longer accepts a `limit` it silently ignored** (`clients/ts/src/pipes.ts`, `clients/ts/src/client.test.ts`, `docs/src/content/docs/sdk/pipes.md`, `docs/src/content/docs/sdk/reference.md`): closes #464, raised by CodeRabbit on #456. It took the same per-call options type as the query builder — which carries `limit` — but forwarded only `signal`, so `wh.pipe('top_pages').fetch({ limit: 10 })` type-checked, ran, and quietly returned whatever the pipe's SQL returned. `QueryBuilder.fetch` and `TableRef.fetch` both honour `limit`, so the inconsistency sat inside one shared type. There is nothing to forward: the endpoint binds the request body as the pipe's *parameters* (`internal/api/pipes.go` → `pipes.BindParams`), and a key the SQL doesn't declare is ignored, so a client-side row cap is not something the pipes surface offers. The parameter is now a dedicated `PipeRequestOptions` (exported) declaring `signal?: AbortSignal` and `limit?: never`, making the dead option a compile error rather than a silent no-op. `never` rather than simply omitting `limit`, because omitting it only rejects fresh object literals — TypeScript's excess-property check doesn't apply to a *variable*, so a shared `const opts: RequestOptions` carrying a limit would still have passed and still been dropped, which is the defect rather than a narrower version of it. Both cases are pinned by `@ts-expect-error` tests. **Note the collateral effect**, which is the half most consumers will actually meet: a value *declared* `RequestOptions` no longer assigns to a pipe `.fetch()` at all, even when it carries no limit at runtime, because the declared type permits one and assignability is decided on the type. Type a shared options object as `PipeRequestOptions` — the table and query-builder `.fetch()` accept it too, so it works everywhere — or inline `{ signal }` at the pipe call. Structural wrappers are unaffected: method parameters compare bivariantly, so an `interface Fetchable { fetch(opts?: RequestOptions): … }` is still satisfied by `PipeRef`. **Migration:** declare a `{{limit}}` parameter in the pipe's SQL and pass it as a pipe parameter — `wh.pipe(name, { limit })` — which is what the docs already showed. Pre-existing rather than introduced by #456, folded in there because that PR renames the type in question. -- **BREAKING (SDK): `FetchOptions` is renamed `RequestOptions`** (`clients/ts/src/types.ts`, `clients/ts/src/index.ts`, `clients/ts/src/query-builder.ts`, `clients/ts/src/table.ts`, `clients/ts/src/pipes.ts`): the per-call options type accepted by `.fetch()`. The old name collided conceptually with the new `options.fetchOptions` — which, following OpenAI, Anthropic, and the wider ecosystem, means "extra `RequestInit` fields", not "options for our `.fetch()` method". Shipping both would have left `FetchOptions` and `fetchOptions` in the same SDK one capital letter apart, meaning unrelated things. `RequestOptions` is what Anthropic's SDK calls the identical concept. No deprecated alias: the type is unreferenced by anything consuming the pre-1.0 package, and keeping it would preserve exactly the ambiguity the rename removes. Renaming the import is the whole migration for this entry — note the separate `PipeRef.fetch` narrowing above, which is a behavioural break in the same file. The module-private `RequestOptions` in `http.ts` — the internal request descriptor — becomes `RequestSpec` to free the name. +- **BREAKING (SDK): `FetchOptions` is renamed `RequestOptions`** (`clients/ts/src/types.ts`, `clients/ts/src/index.ts`, `clients/ts/src/query-builder.ts`, `clients/ts/src/table.ts`, `clients/ts/src/pipes.ts`): the per-call options type accepted by `.fetch()`. The old name collided conceptually with the new `options.fetchOptions` — which, following OpenAI, Anthropic, and the wider ecosystem, means "extra `RequestInit` fields", not "options for our `.fetch()` method". Shipping both would have left `FetchOptions` and `fetchOptions` in the same SDK one capital letter apart, meaning unrelated things. `RequestOptions` is what Anthropic's SDK calls the identical concept. No deprecated alias: the type is unreferenced by anything consuming the pre-1.0 package, and keeping it would preserve exactly the ambiguity the rename removes. Renaming the import is the whole migration for this entry — note the separate `PipeRef.fetch` narrowing above, which is a behavioral break in the same file. The module-private `RequestOptions` in `http.ts` — the internal request descriptor — becomes `RequestSpec` to free the name. - **`@wavehouse/sdk` `engines.node` floor back to `>=22`, matching the only line we test** (`clients/ts/package.json`, `clients/ts/README.md`, `docs/src/content/docs/sdk/index.mdx`, `docs/src/content/docs/sdk/queries.md`, `pnpm-workspace.yaml`): the floor was relaxed to `>=18` when the browser-first distribution landed (see the entry below), on the reasoning that the runtime needs only `fetch`. Nothing ever tested 18, though — `.nvmrc` pins 22 and `.github/actions/setup-env` consumes it via `node-version-file`, so 22 is the single version CI exercises — and Node 18 and 20 have both since reached upstream end-of-life. Declaring a floor we neither test nor is supported upstream promises more than it can back, so it returns to `>=22`. **Consumer impact:** installing on Node < 22 now warns with `EBADENGINE` under npm, and fails outright under pnpm with `engine-strict` enabled. The SDK README and the docs' Runtime support section state the requirement, which they previously either omitted or quoted as 18. @@ -633,7 +634,7 @@ The first public release. Everything below shipped in it — the sections are gr - **Hub wildcard subscriptions** (`internal/api/hub.go`, `internal/api/hub_test.go`): dropped the NATS-style `*` / `>` pattern matching from `Hub.Broadcast`, the wildcard pattern loop, the `sent` dedup map, the `matchTopic` helper, and the eight wildcard tests (plus `TestMatchTopic`). After the #89 MVP cuts every producer publishes a concrete `ingest.
` subject and the SDK only ever subscribes to one concrete subject, so the wildcard fan-out was unused machinery. Closes #100 (part of #87). Net −210 lines (mostly tests). -- **`project-orchestrator.yml` workflow + its three composite-action artifacts** (`.github/workflows/project-orchestrator.yml`, `.github/actions/board-upsert-status/`, `.github/actions/set-linked-issues-status/`, `.github/scripts/board-fetch-item.sh`, `AGENTS.md`, `CHANGELOG.md`): −887 lines net. The orchestrator was the largest single source of cross-trigger complexity on this repo (3-4 workflow_run-chained runs per PR push, `statusCheckRollup` GraphQL perms quirks, integration-token `NONE` for private-org members) for behaviour that is mostly either provided natively by GitHub or a one-click manual operation on a 4-person team. Replaced by: reviewer-assign step in `housekeeping.yml` that fires once on `pull_request_target: opened` / `ready_for_review` (not per-synchronize, so it doesn't re-spam after `dismiss_stale_reviews_on_push`), plus GitHub's native Projects v2 workflows (`Auto-add to project`, `Item added`, `Pull request merged`) configured in the project UI. Trade-offs explicit in the PR body: drafts no longer auto-flip on bot-clean, `CHANGES_REQUESTED` doesn't auto-move the board card, linked-issue card mirroring is dropped. AGENTS.md §"Governance Files" + §"Task Board state machine" + §"Review tooling reference" all rewritten to match. `dependabot-automerge.yml` trimmed in parallel: no more board-upsert step (native handles placement), `PROJECT_BOARD_TOKEN` guard removed (no longer used in this workflow), reviewer list sourced from `board-config.env`'s `ADMINS` via `replace()`, major-bump comment uses the marker-comment upsert pattern from `housekeeping.yml`. +- **`project-orchestrator.yml` workflow + its three composite-action artifacts** (`.github/workflows/project-orchestrator.yml`, `.github/actions/board-upsert-status/`, `.github/actions/set-linked-issues-status/`, `.github/scripts/board-fetch-item.sh`, `AGENTS.md`, `CHANGELOG.md`): −887 lines net. The orchestrator was the largest single source of cross-trigger complexity on this repo (3-4 workflow_run-chained runs per PR push, `statusCheckRollup` GraphQL perms quirks, integration-token `NONE` for private-org members) for behavior that is mostly either provided natively by GitHub or a one-click manual operation on a 4-person team. Replaced by: reviewer-assign step in `housekeeping.yml` that fires once on `pull_request_target: opened` / `ready_for_review` (not per-synchronize, so it doesn't re-spam after `dismiss_stale_reviews_on_push`), plus GitHub's native Projects v2 workflows (`Auto-add to project`, `Item added`, `Pull request merged`) configured in the project UI. Trade-offs explicit in the PR body: drafts no longer auto-flip on bot-clean, `CHANGES_REQUESTED` doesn't auto-move the board card, linked-issue card mirroring is dropped. AGENTS.md §"Governance Files" + §"Task Board state machine" + §"Review tooling reference" all rewritten to match. `dependabot-automerge.yml` trimmed in parallel: no more board-upsert step (native handles placement), `PROJECT_BOARD_TOKEN` guard removed (no longer used in this workflow), reviewer list sourced from `board-config.env`'s `ADMINS` via `replace()`, major-bump comment uses the marker-comment upsert pattern from `housekeeping.yml`. - **`STATUS_*` and old `ADMINS` consumers in `board-config.env`** — STATUS option IDs had only orchestrator-side consumers and are now unreferenced. `ADMINS` was restored to `board-config.env` after the initial orchestrator-removal commit dropped it (Gemini and Claude both flagged the resulting drift across three inlined copies); both `housekeeping.yml` and `dependabot-automerge.yml` now load `ADMINS` from `board-config.env`. `admin-approval.yml` keeps its own inline copy with the documented latency-avoidance reasoning. diff --git a/docs/src/content/docs/api.md b/docs/src/content/docs/api.md index 13992778..a32e2198 100644 --- a/docs/src/content/docs/api.md +++ b/docs/src/content/docs/api.md @@ -296,6 +296,7 @@ The body is a **flat JSON object** whose keys must match column names in the tar | 503 | `{"error":"schema not loaded yet"}` | The tenant's first schema discovery has not succeeded yet (its ClickHouse unreachable, or [no pool for it](/settings-directory#clickhouse)), so whether the table exists is not known; `Retry-After: 5`. Decided before the body is read | | 500 | `{"error":"publish failed"}` | Message queue error | | 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure, for that tenant alone) or not open (see [Message Queue](/settings-directory#message-queue)). Response includes `Retry-After: 30` header. | +| 503 | `{"error":"service unavailable"}` | The message queue could not be reached or did not answer in time (a transient broker failure, not a full queue). Response includes `Retry-After: 5` header. Reserved for an external broker ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)): the embedded broker never reports this, and its publish failures are the `500` above. | | 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) | **curl example:** @@ -407,6 +408,7 @@ A `200` is returned whenever the body was read and the records were processed | 415 | `{"error":"no Content-Type: ingest requires one of application/json, application/x-ndjson, …"}` (declared variant: `Content-Type "text/plain": ingest requires one of …` — see the note above on how declarations are echoed; conflicting variant: `conflicting Content-Type declarations "application/json", "application/x-ndjson": ingest reads one format per request, and requires one of …`) | The request declared no `Content-Type`, one whose media type is unsupported or does not parse, a comma-bearing value that does not parse as a single media type, or repeated header lines that disagree — different formats, or one supported and one not. Checked before the body is parsed | | 500 | `{"error":"publish failed"}` / `{"error":"dedupe failed"}` | Message-queue or dedup-backend failure mid-batch | | 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure) or not open, mid-batch; includes `Retry-After: 30` | +| 503 | `{"error":"service unavailable"}` | The message queue could not be reached or did not answer in time, mid-batch; includes `Retry-After: 5`. Reserved for an external broker ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)): the embedded broker never reports this, and its publish failures are the `500` above | | 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) | :::caution[At-least-once on retry] diff --git a/docs/src/content/docs/architecture.md b/docs/src/content/docs/architecture.md index d5ac03e3..954379b1 100644 --- a/docs/src/content/docs/architecture.md +++ b/docs/src/content/docs/architecture.md @@ -83,7 +83,7 @@ The API layer uses [Chi](https://github.com/go-chi/chi) for routing with Request - **pipes.go** — Named query pipe handlers: admin listing (`GET /v1/ops/pipes[/{name}]`, read per request from its `pipes.Source`) and execution with parameter binding. `pipes.json` is the only write path. - **structured_query.go** — Handler for `POST /v1/query?table={table}`: validates query AST, enforces permissions, builds and executes SQL. - **ch_errors.go** — `writeCHError`, the one mapping from a failed ClickHouse query to a response, shared by `/v1/query`, pipes and `/v1/ops/query` so they cannot drift apart: `chconn.Classify` decides the class, and the class the status, `code` and `retryable` ([ClickHouse errors on the query paths](/api#clickhouse-errors-on-the-query-paths)). -- **ingest.go** — Accepts `POST /v1/ingest?table={table}` in three body shapes: one flat JSON object, a JSON array of them, or NDJSON. The **required** `Content-Type` chooses the format *family* — `application/json` versus the four NDJSON spellings — and within the JSON family the body's first non-whitespace byte picks array versus single object; the bytes never choose the family. Anything that is not exactly one readable media type is a `415`, decided before the body is read: the header is parsed per RFC 9110 §8.3, and because `Content-Type` is a singleton field, repeated header lines must all resolve to the same format and a value carrying a comma is refused unless the value as a whole parses as one media type — a comma inside a *quoted* parameter value is data, so `application/json; a=", application/x-ndjson; b="` is accepted. It then reads the whole (`MaxBytesReader`-capped) body into a pooled buffer and runs the per-format record readers over those bytes, so the `413` lands before any record is processed and peak memory per request is O(body) rather than O(record). Then it validates each record against the discovered schema, optional dedup, and publishes each row through `mq.Publisher` on `mq.Topic{Tenant, Table, Scope}` (the request's tenant, read off its resolved store — `store.Tenant()` — and raw names; the subject it becomes is `internal/mq`'s; a full queue comes back as `mq.ErrQueueFull`, which is the `503` + `Retry-After`). When dedup is on, a row missing the configured `id_field` can't be deduped: it is logged at `WARN` and counted by `wavehouse_ingest_dedupe_missing_id_total` (labeled by `table`), then published un-deduped — or rejected when `dedupe.require_id` is set ([#219](https://github.com/Wave-RF/WaveHouse/issues/219)). +- **ingest.go** — Accepts `POST /v1/ingest?table={table}` in three body shapes: one flat JSON object, a JSON array of them, or NDJSON. The **required** `Content-Type` chooses the format *family* — `application/json` versus the four NDJSON spellings — and within the JSON family the body's first non-whitespace byte picks array versus single object; the bytes never choose the family. Anything that is not exactly one readable media type is a `415`, decided before the body is read: the header is parsed per RFC 9110 §8.3, and because `Content-Type` is a singleton field, repeated header lines must all resolve to the same format and a value carrying a comma is refused unless the value as a whole parses as one media type — a comma inside a *quoted* parameter value is data, so `application/json; a=", application/x-ndjson; b="` is accepted. It then reads the whole (`MaxBytesReader`-capped) body into a pooled buffer and runs the per-format record readers over those bytes, so the `413` lands before any record is processed and peak memory per request is O(body) rather than O(record). Then it validates each record against the discovered schema, optional dedup, and publishes each row through `mq.Publisher` on `mq.Topic{Tenant, Table, Scope}` (the request's tenant, read off its resolved store — `store.Tenant()` — and raw names; the subject it becomes is `internal/mq`'s; a full queue comes back as `mq.ErrQueueFull`, which is the `503` + `Retry-After: 30`, and a broker that cannot be reached or does not answer in time as `mq.ErrUnavailable`, the `503` + `Retry-After: 5`). When dedup is on, a row missing the configured `id_field` can't be deduped: it is logged at `WARN` and counted by `wavehouse_ingest_dedupe_missing_id_total` (labeled by `table`), then published un-deduped — or rejected when `dedupe.require_id` is set ([#219](https://github.com/Wave-RF/WaveHouse/issues/219)). - **query.go** — Proxies raw SQL for `POST /v1/ops/query` straight to the `?tenant=`'s ClickHouse HTTP interface (`chconn.Pools.Target` by the resolved store's tenant; the zero target — no pool — is a `503` with `Retry-After`). **Not cached** — sets `Cache-Control: no-store` so every request hits ClickHouse; DateTime is rendered ISO-8601 via `date_time_output_format=iso` (the Go-side type conversion lives in the structured-query / pipes path, not here). - **stream.go** — Real-time streaming via SSE. Callers select a table with the `?table=` query parameter. Each connection registers one `Subscriber` (the `stream/` package) with both the event `Hub` (under its `(topic, role)`) and the shared keepalive wheel, then drains both from a single byte-pump — so idle streams keep emitting `:` keepalive comments (surviving reverse-proxy idle timeouts) while live events arrive already projected and serialized. Per-event projection/serialization happens **once per role** in the `Hub`, not once per subscriber ([#294](https://github.com/Wave-RF/WaveHouse/issues/294)); the handler also snapshots the connection's JWT claims onto the `Subscriber`, which the `Hub` evaluates per subscriber when the role carries a row-level `filter` ([#319](https://github.com/Wave-RF/WaveHouse/issues/319)). Gap-fill replay (`mq.Replayer.ReplaySince` on the connection's `mq.Topic` — a `DeliverByStartTime` consumer inside `internal/mq`) stays per-connection (low-volume, one-time on connect). A stream ends, a gap-fill in progress included, when the server begins shutting down (`Closing`) or its `Subscriber` is evicted because its tenant is no longer served (`Hub.Prune`); one admitted just before the reload that stopped serving its tenant, and registered just after the prune, is ended right after it registers (`Served`). - **schema.go** — Schema discovery API of one tenant, the `?tenant=` (`opsStore`): list all schemas, get one table, trigger refresh. `lookupSchema`, shared with the ingest and structured-query handlers, is the one reading of a `SchemaRegistry.Lookup` miss: `503` with `Retry-After` before the tenant's first discovery (`ErrNotLoaded`, or no registry built yet), `404` for a table the discovered schema lacks; the list answers the same `503` rather than `[]`. A refresh of a tenant on no pool (`discovery.ErrNoConnection`) is a `503` with `Retry-After` too. The handlers hold `RegistrySource`, `func(*settings.Store) *discovery.SchemaRegistry`, and the query paths a `func(*settings.Store) driver.Conn` beside it — each resolves the request's tenant per call, and a nil connection (a tenant no pool could be opened for, such as by the connection ceiling) is a `503` ahead of the cache, so nothing cached before is served. @@ -158,11 +158,12 @@ 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()`, `NakWithDelay(d)`, which falls back to `Nak` for a message built without `WithNakDelay`); `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. +- **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()`, `NakWithDelay(d)`, which falls back to `Nak` for a message built without `WithNakDelay`); `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: 30`; `ErrUnavailable` when the broker cannot be reached or does not answer in time — the 503 + `Retry-After: 5`, which no backend returns yet: the embedded broker's publish failures are the `500`), `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 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. +- **mqtest/** — The conformance suite for `Broker` (`mqtest.Run`): the behavior the rest of the process relies on — publish and consume round trips with names that need encoding, per-tenant order, redelivery, dead-lettering and its counts, replay bounds and isolation, the one `failed` report of a consumer whose delivery ends underneath it — checked through the interfaces alone, with no stream or subject name in sight. Each implementation runs it from a test of its own — the embedded one from `mqtest/embedded_test.go`, a test binary apart from `internal/mq`'s so the two share no 15s budget — handing it a fresh broker per case and flags (`mqtest.Caps`) for the few places where backends legitimately differ: whether a full queue refuses its own tenant alone, whether `PurgeAcked` removes anything, whether a tenant never given a budget has a dead-letter queue to report on, and whether `CreateConsumer` configures the durable or only finds one. ### `observability/` — OpenTelemetry Pipeline diff --git a/internal/api/ingest.go b/internal/api/ingest.go index 10daddea..ca89bb83 100644 --- a/internal/api/ingest.go +++ b/internal/api/ingest.go @@ -124,8 +124,9 @@ type recordReject struct { // abandons the remaining records rather than silently losing the tail. // // Most causes are TRANSIENT system conditions, where abandoning the tail is what -// makes the batch safe to retry: publish backpressure (503), a publish/marshal -// failure (500), a dedup backend error (500). +// makes the batch safe to retry: publish backpressure (503), an unreachable +// broker (503, mq.ErrUnavailable), a publish/marshal failure (500), a dedup +// backend error (500). // // One is not. An insert grant that resolved for the other operation is a 403 and // a caller/config bug — retrying cannot help. It aborts rather than rejecting @@ -135,7 +136,7 @@ type recordReject struct { type requestAbort struct { Status int Message string - RetryAfter string // non-empty → emit a Retry-After header (503 backpressure) + RetryAfter string // non-empty → emit a Retry-After header (503: backpressure or an unavailable broker) } func (h *IngestHandler) Handle(w http.ResponseWriter, r *http.Request) { @@ -710,6 +711,11 @@ func (h *IngestHandler) processRecord( slog.WarnContext(ctx, "ingest queue is full", "tenant", store.Tenant(), "error", err, "table", table, "scope", scope) return false, nil, &requestAbort{Status: http.StatusServiceUnavailable, Message: "service unavailable", RetryAfter: "30"} } + if errors.Is(err, mq.ErrUnavailable) { + // A broker blip, not a full queue: a sooner retry is likely to land. + slog.WarnContext(ctx, "ingest queue unavailable", "tenant", store.Tenant(), "error", err, "table", table, "scope", scope) + return false, nil, &requestAbort{Status: http.StatusServiceUnavailable, Message: "service unavailable", RetryAfter: "5"} + } slog.ErrorContext(ctx, "failed to publish to the ingest queue", "tenant", store.Tenant(), "error", err, "table", table, "scope", scope) return false, nil, &requestAbort{Status: http.StatusInternalServerError, Message: "publish failed"} } diff --git a/internal/api/ingest_test.go b/internal/api/ingest_test.go index 2ae205e3..a87a4df9 100644 --- a/internal/api/ingest_test.go +++ b/internal/api/ingest_test.go @@ -251,6 +251,20 @@ func TestIngest_PublishError_503(t *testing.T) { testutil.AssertJSONErrorResponse(t, w) } +func TestIngest_PublishUnavailable_503(t *testing.T) { + t.Parallel() + pub := &testutil.MockPublisher{Err: fmt.Errorf("%w: nats: timeout", mq.ErrUnavailable)} + h := NewIngestHandler(fixedRegistry(testRegistry(t)), pub) + + req := ingestRequest(t, "clicks", map[string]any{"page": "/home"}) + w := httptest.NewRecorder() + h.Handle(w, withTenant(req)) + + assert.Equal(t, http.StatusServiceUnavailable, w.Code) + assert.Equal(t, "5", w.Header().Get("Retry-After")) + testutil.AssertJSONErrorResponse(t, w) +} + func TestIngest_PublishError_500(t *testing.T) { t.Parallel() pub := &testutil.MockPublisher{Err: errors.New("some other error")} @@ -1055,6 +1069,20 @@ func TestIngest_NDJSON_Backpressure_503(t *testing.T) { testutil.AssertJSONErrorResponse(t, w) } +func TestIngest_NDJSON_Unavailable_503(t *testing.T) { + t.Parallel() + pub := &testutil.MockPublisher{Err: fmt.Errorf("%w: nats: no responders", mq.ErrUnavailable)} + h := NewIngestHandler(fixedRegistry(testRegistry(t)), pub) + + req := ndjsonRequest(t, "clicks", jsonLine(t, map[string]any{"page": "/a"})) + w := httptest.NewRecorder() + h.Handle(w, withTenant(req)) + + assert.Equal(t, http.StatusServiceUnavailable, w.Code) + assert.Equal(t, "5", w.Header().Get("Retry-After")) + testutil.AssertJSONErrorResponse(t, w) +} + func TestIngest_NDJSON_PublishError_500(t *testing.T) { t.Parallel() pub := &testutil.MockPublisher{Err: errors.New("some other error")} diff --git a/internal/mq/embedded.go b/internal/mq/embedded.go index efd9f05a..15ca485f 100644 --- a/internal/mq/embedded.go +++ b/internal/mq/embedded.go @@ -680,14 +680,13 @@ func (e *EmbeddedNATS) CreateConsumer(ctx context.Context, cfg ConsumerConfig) ( failed: make(chan error, 1), } c.fail = func(err error) { - // Exactly one error, and nothing once stop has been called. - if c.stopped.Load() { + // Exactly one error, and nothing once stop has been called: a durable + // deleted on several tenants' queues ends each delivery, and a caller + // that already drained the first must not see the next. + if c.stopped.Load() || !c.reported.CompareAndSwap(false, true) { return } - select { - case c.failed <- err: - default: - } + c.failed <- err } if err := e.register(ctx, c.fanIn); err != nil { return nil, fmt.Errorf("create consumer: %w", err) @@ -873,7 +872,8 @@ func (f *fanIn) start(deliver func(jetstream.Msg), prefetch int, watch bool) (st // failed channel its contract promises. type workerConsumer struct { *fanIn - failed chan error + failed chan error + reported atomic.Bool } func (c *workerConsumer) Consume(handler func(msg *Message), prefetch int) (func(), <-chan error, error) { @@ -1057,7 +1057,9 @@ func (e *EmbeddedNATS) ReplaySince(ctx context.Context, topic Topic, since time. } msg, err := cons.Next(jetstream.FetchMaxWait(500 * time.Millisecond)) if err != nil { - if errors.Is(err, jetstream.ErrNoMessages) || errors.Is(err, nats.ErrTimeout) { + // A pull that raced the connection closing can end in either + // answer too, and that is not caught up. + if (errors.Is(err, jetstream.ErrNoMessages) || errors.Is(err, nats.ErrTimeout)) && !e.conn.IsClosed() { return nil // caught up } return fmt.Errorf("replay next: %w", err) diff --git a/internal/mq/embedded_failed_test.go b/internal/mq/embedded_failed_test.go new file mode 100644 index 00000000..9ab73091 --- /dev/null +++ b/internal/mq/embedded_failed_test.go @@ -0,0 +1,48 @@ +package mq + +import ( + "testing" + "time" + + "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/stretchr/testify/require" +) + +// A durable deleted on several tenants' queues ends each delivery; a caller +// that drained the first report must not see the next. +func TestEmbeddedNATS_Consume_ReportsOnceHoweverManyDeliveriesEnd(t *testing.T) { + e := newTestEmbedded(t, "acme", "globex") + ctx := t.Context() + cons, err := e.CreateConsumer(ctx, ConsumerConfig{Durable: "doomed", MaxAckPending: 10}) + require.NoError(t, err) + delivered := make(chan struct{}, 2) + stop, failed, err := cons.Consume(func(*Message) { delivered <- struct{}{} }, 4) + require.NoError(t, err) + t.Cleanup(stop) + // A delivery on each tenant proves both pulls are live: a durable deleted + // before its pull reaches the server ends nothing the client sees. + for _, id := range []tenant.ID{"acme", "globex"} { + require.NoError(t, e.Publish(ctx, Topic{Tenant: id, Table: "t"}, []byte("x"))) + } + for range 2 { + select { + case <-delivered: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for a delivery on each tenant") + } + } + + require.NoError(t, e.js.DeleteConsumer(ctx, "INGEST_globex", "doomed")) + select { + case err := <-failed: + require.ErrorIs(t, err, ErrDeliveryEnded) + case <-time.After(5 * time.Second): + t.Fatal("delivery ended underneath the consumer and nothing was reported") + } + require.NoError(t, e.js.DeleteConsumer(ctx, "INGEST_acme", "doomed")) + select { + case err := <-failed: + t.Fatalf("a second failure was reported: %v", err) + case <-time.After(300 * time.Millisecond): + } +} diff --git a/internal/mq/mq.go b/internal/mq/mq.go index 11ae7bea..089d2b44 100644 --- a/internal/mq/mq.go +++ b/internal/mq/mq.go @@ -5,7 +5,8 @@ // ingest queue, park a message on the dead-letter queue, replay since a time, // drop what is both written and expired — in the types below. How that maps to // subjects, streams, sequences, and consumers is the implementation's -// (EmbeddedNATS), so a broker change lands here once. +// (EmbeddedNATS), so a broker change lands here once. The behavior below is +// what mqtest checks: every implementation passes its suite. package mq import ( @@ -167,32 +168,44 @@ func WithHeader(key, value string) PublishOpt { } } -// ErrQueueFull is returned by Publisher.Publish when the topic's tenant's -// ingest queue refuses new events — it is at its byte budget, or the tenant -// has no queue open yet — the backpressure signal the API turns into a 503 -// with Retry-After. +// ErrQueueFull is returned by Publisher.Publish when the queue that holds the +// topic's tenant refuses new events because it is at a byte limit — the +// backpressure signal the API turns into a 503 with Retry-After. Which limits +// there are, and which tenants share one, is the implementation's (see +// Broker.SetMaxBytes). An implementation that opens a queue per tenant also +// returns it for a tenant whose queue it cannot open yet. var ErrQueueFull = errors.New("ingest queue is full") +// ErrUnavailable is returned when the broker cannot be reached or does not +// answer in time — a transient failure, not a refusal, that the API turns +// into a 503 with a short Retry-After. Only a backend whose broker is out of +// process returns it; the embedded one's publish failures are plain errors. +var ErrUnavailable = errors.New("message queue unavailable") + // Publisher appends events to the ingest queue. type Publisher interface { - // Publish stores data as one event on topic, in the ingest queue of the - // topic's tenant. ErrQueueFull when that queue is at its byte budget, or - // the tenant has no queue open yet (see Broker.SetMaxBytes). + // Publish stores data as one event on topic, in the ingest queue that + // holds the topic's tenant. A topic without a valid tenant is refused + // before anything is sent. ErrQueueFull when that queue refuses the event + // at a byte limit (or, per tenant, cannot be opened yet), ErrUnavailable + // when the broker cannot take it now. Publish(ctx context.Context, topic Topic, data []byte, opts ...PublishOpt) error Close() error } // Subscriber delivers every event on the ingest queue, across all tenants -// and topics: each tenant's in the order it was published, and different -// tenants' concurrently. +// and topics: each tenant's in the order it was published. type Subscriber interface { - // Subscribe registers a handler for incoming events under a durable - // consumer named consumerName, held on every tenant's queue — those - // opened after Subscribe included. The handler runs on one delivery - // goroutine per tenant, one message at a time, so it must be safe to - // call concurrently for different tenants. The messages fetched ahead of - // it are a fixed number split across the tenants, as Consumer.Consume's - // prefetch is, so they do not grow with the number of tenants. + // Subscribe registers a handler for incoming events, across every + // tenant — those whose queues open after Subscribe included. Every event + // published after Subscribe returns is delivered; whether earlier ones + // are is the implementation's, and so is whether consumerName names a + // durable consumer. The handler runs one message at a time on each + // delivery unit — a tenant's queue, or the partition that holds it — so + // it must be safe to call concurrently for different units. The messages + // fetched ahead of it are a fixed number split across the units, as + // Consumer.Consume's prefetch is, so they do not grow with the number of + // tenants. // // CONTRACT: If the handler intends to return an error to trigger automatic // redelivery, it MUST NOT manually call msg.Ack() or msg.Nak() beforehand. @@ -202,7 +215,7 @@ type Subscriber interface { // error return. // // CONTRACT: Calling msg.DoubleAck(ctx) and then returning a non-nil error is - // undefined behaviour — the consume loop will Nak() after a successful + // undefined behavior — the consume loop will Nak() after a successful // broker-confirmed Ack. Call DoubleAck, then return nil on success. Subscribe(ctx context.Context, consumerName string, handler func(msg *Message) error) error Close() error @@ -215,31 +228,33 @@ type ConsumerConfig struct { // AckWait is the redelivery timeout: a message not acked within it is // delivered again. AckWait time.Duration - // MaxAckPending caps unacked messages broker-side, per tenant: delivery - // of a tenant's events pauses when that tenant's unacked ones hit it - // (backpressure), and no other tenant's does. + // MaxAckPending caps unacked messages broker-side, per delivery unit (a + // tenant's queue, or the partition that holds it): delivery from a unit + // pauses when its unacked messages hit it (backpressure), and no other + // unit's does. MaxAckPending int } // Consumer is a live durable consumer created by ConsumerManager. type Consumer interface { - // Consume delivers each message to handler on a delivery goroutine of its - // tenant's: one per tenant, so a tenant's messages arrive in order, one at - // a time, while different tenants' arrive concurrently — handler must be - // safe for that. A handler that blocks holds back its tenant's delivery — + // Consume delivers each message to handler on the delivery goroutine of + // its delivery unit — the tenant's queue, or the partition that holds + // it: one per unit, so a tenant's messages arrive in order, one at a + // time, while different units' arrive concurrently — handler must be + // safe for that. A handler that blocks holds back its unit's delivery — // that is the backpressure the ingest worker relies on. About prefetch - // messages are fetched ahead across the tenants together: the tenants' - // queues when delivery starts split it, and a queue joined later fetches - // ahead its share of it at that point, at least one message each (0 = the - // client default, per tenant). The returned stop asks delivery to end and - // returns without waiting: a handler invocation already in flight, or one - // for a message already queued client-side, may still run after stop - // returns, so a handler must not write to anything the caller tears down - // right after stopping. + // messages are fetched ahead across the units together: the units when + // delivery starts split it, and a queue joined later fetches ahead its + // share of it at that point, at least one message each (0 = the client + // default, per unit). The returned stop asks delivery to end and returns + // without waiting: a handler invocation already in flight, or one for a + // message already queued client-side, may still run after stop returns, + // so a handler must not write to anything the caller tears down right + // after stopping. // // Delivery can also end on its own after Consume has returned: the broker // or the client gives up on the consumer (it was deleted, the connection - // closed), or a tenant's queue opened later could not be joined. That is + // closed), or a queue opened later could not be joined. That is // reported on failed — exactly one error, and nothing once stop has been // called — because no message will ever arrive to say so. A caller that // ignores failed waits forever on a dead consumer. @@ -250,8 +265,10 @@ type Consumer interface { // broker's reason when it gave one. var ErrDeliveryEnded = errors.New("consumer delivery ended") -// ConsumerManager creates durable consumers on the ingest queue, held on -// every tenant's queue — those opened later included. A delivered +// ConsumerManager gives access to durable consumers on the ingest queue, held +// on every tenant's queue — those opened later included. Whether +// CreateConsumer creates the durable, or only finds one someone else made and +// checks it against the config, is the implementation's. A delivered // Message.Ctx is the ctx given to CreateConsumer: unlike Subscriber, the // consumer path does not extract the trace context carried in the message // headers, because its one consumer (the ingest worker) batches across @@ -271,7 +288,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. + // Tables maps table name → parked messages, for the tables asked about; + // empty, never nil, when none has any. // Every scope of a table counts under the table; scope is not broken out // yet (it is inert until #235). Tables map[string]uint64 @@ -280,8 +298,9 @@ type DeadLetterCounts struct { } // ErrNoDeadLetterQueue is returned by DeadLetterStats.DeadLetterCounts when -// the tenant has no dead-letter queue (nothing can have been parked for it). -// Any other failure to read it is a plain error. +// the tenant has no dead-letter queue of its own (nothing can have been +// parked for it). An implementation whose tenants share one queue returns +// zero counts instead. Any other failure to read it is a plain error. var ErrNoDeadLetterQueue = errors.New("dead-letter queue not found") // DeadLetterStats reports on the dead-letter queues. @@ -289,7 +308,8 @@ 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 (all of its - // scopes). + // scopes). A tenant with nothing parked has zero counts, or + // ErrNoDeadLetterQueue when it has no queue at all. DeadLetterCounts(ctx context.Context, id tenant.ID, table string) (DeadLetterCounts, error) } @@ -308,7 +328,9 @@ type Purger interface { // acknowledged goes. Reports whether anything was removed, and joins // each failed tenant's error — ErrConsumerNotFound for one whose queue the // consumer has not been created on; the other tenants' are purged all the - // same. + // same. An implementation whose retention the broker's operator owns + // removes nothing and reports false: either way, no unacked event is + // removed. PurgeAcked(ctx context.Context, consumer string, olderThan map[tenant.ID]time.Time) (purged bool, err error) } @@ -318,7 +340,8 @@ type Replayer interface { // since, in order, until send returns false or the queue is caught up. // Running out of events is the normal end; failing to start the replay, or // a delivery failure before it catches up, is an error. A done ctx stops - // the replay and returns ctx's error. + // the replay and returns ctx's error. A topic without a valid tenant is + // refused as Publish refuses it. ReplaySince(ctx context.Context, topic Topic, since time.Time, send func(data []byte) bool) error } @@ -336,10 +359,12 @@ type Broker interface { // SetMaxBytes applies tenant id's byte budget (its hot-reloadable // mq.max_bytes_gb) to that tenant's queues — how it is split between them // is the implementation's — opening them if the tenant has none yet. No - // other tenant's queues are touched. On an error the implementation - // restores the previous budget where it can (best effort: the error says - // when it could not, and a canceled ctx abandons the restore too), and - // MaxBytes keeps reporting the previous budget so the next call retries. + // other tenant's queues are touched. An implementation whose tenants + // share queues may only record the budget, and say so where it does. On + // an error the implementation restores the previous budget where it can + // (best effort: the error says when it could not, and a canceled ctx + // abandons the restore too), and MaxBytes keeps reporting the previous + // budget so the next call retries. // MaxBytes reports the budget last applied in full for id, 0 when none // has been. SetMaxBytes(ctx context.Context, id tenant.ID, maxBytes int64) error diff --git a/internal/mq/mqtest/cases.go b/internal/mq/mqtest/cases.go new file mode 100644 index 00000000..699b5194 --- /dev/null +++ b/internal/mq/mqtest/cases.go @@ -0,0 +1,527 @@ +package mqtest + +import ( + "context" + "errors" + "fmt" + "slices" + "testing" + "time" + + "github.com/Wave-RF/WaveHouse/internal/mq" + "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/trace" +) + +// delivery is one message as a handler saw it. +type delivery struct { + topic mq.Topic + data string + msg *mq.Message +} + +func publish(t *testing.T, b mq.Broker, topic mq.Topic, data string, opts ...mq.PublishOpt) { + t.Helper() + require.NoError(t, b.Publish(ctx(t), topic, []byte(data), opts...), "publish %q on %+v", data, topic) +} + +// consume runs the suite's durable on b with handle called before each +// delivery is reported on the returned channel. Stopped at cleanup. +func consume(c context.Context, t *testing.T, b mq.Broker, cfg mq.ConsumerConfig, handle func(*mq.Message)) (<-chan delivery, func(), <-chan error) { + t.Helper() + if cfg.Durable == "" { + cfg.Durable = Durable + } + cons, err := b.CreateConsumer(c, cfg) + require.NoError(t, err) + got := make(chan delivery, 256) + stop, failed, err := cons.Consume(func(m *mq.Message) { + if handle != nil { + handle(m) + } + got <- delivery{topic: m.Topic(), data: string(m.Data), msg: m} + }, 16) + require.NoError(t, err) + t.Cleanup(stop) + return got, stop, failed +} + +// ackEach DoubleAcks every message, reporting a failed ack on t. +func ackEach(t *testing.T) func(*mq.Message) { + return func(m *mq.Message) { + assert.NoError(t, m.DoubleAck(m.Ctx)) + } +} + +// next waits for n deliveries. +func next(t *testing.T, got <-chan delivery, n int) []delivery { + t.Helper() + out := make([]delivery, 0, n) + timeout := time.After(wait) + for len(out) < n { + select { + case d := <-got: + out = append(out, d) + case <-timeout: + t.Fatalf("timed out after %d of %d deliveries: %+v", len(out), n, out) + } + } + return out +} + +// none asserts nothing arrives on ch for a while. +func none[T any](t *testing.T, ch <-chan T, what string) { + t.Helper() + select { + case v := <-ch: + t.Fatalf("%s: %+v", what, v) + case <-time.After(quiet): + } +} + +func replay(t *testing.T, b mq.Broker, topic mq.Topic, since time.Time) []string { + t.Helper() + got := []string{} + require.NoError(t, b.ReplaySince(ctx(t), topic, since, func(data []byte) bool { + got = append(got, string(data)) + return true + })) + return got +} + +// replayEventually waits for a replay of topic since to be want: a backend +// may serve replays from a store that trails the ingest queue. +func replayEventually(t *testing.T, b mq.Broker, topic mq.Topic, since time.Time, want []string) { + t.Helper() + deadline := time.Now().Add(wait) + for { + got := replay(t, b, topic, since) + if slices.Equal(got, want) { + return + } + if time.Now().After(deadline) { + assert.Equal(t, want, got, "replay of %+v since %v", topic, since) + return + } + time.Sleep(retryPause) + } +} + +// replayReaches waits until a replay of topic from the start holds at least +// n events: that they are stored where the backend replays from. It stops the +// replay at n, so it never waits out a caught-up. +func replayReaches(t *testing.T, b mq.Broker, topic mq.Topic, n int) { + t.Helper() + deadline := time.Now().Add(wait) + for { + got := 0 + require.NoError(t, b.ReplaySince(ctx(t), topic, time.Time{}, func([]byte) bool { + got++ + return got < n + })) + if got >= n { + return + } + require.False(t, time.Now().After(deadline), "a replay of %+v never reached %d events", topic, n) + time.Sleep(retryPause) + } +} + +type ctxKey struct{} + +// A topic whose names need encoding comes back as it went in, with its data, +// under its tenant; the consumer path delivers with CreateConsumer's ctx. +func roundTrip(t *testing.T, h Harness) { + b := h.New(t) + topics := []mq.Topic{ + {Tenant: Acme, Table: "events"}, + {Tenant: Acme, Table: "a.b*c> d%e", Scope: "s.1 *>%"}, + {Tenant: Globex, Table: "events", Scope: "x"}, + } + for i, topic := range topics { + publish(t, b, topic, fmt.Sprint(i)) + } + c := context.WithValue(ctx(t), ctxKey{}, "worker") + got, _, _ := consume(c, t, b, mq.ConsumerConfig{MaxAckPending: 100}, ackEach(t)) + + byTopic := map[mq.Topic]string{} + for _, d := range next(t, got, len(topics)) { + byTopic[d.topic] = d.data + assert.Equal(t, "worker", d.msg.Ctx.Value(ctxKey{}), "a delivered Message.Ctx is CreateConsumer's") + assert.NotEmpty(t, d.msg.TopicKey()) + } + for i, topic := range topics { + assert.Equal(t, fmt.Sprint(i), byTopic[topic], "%+v", topic) + } +} + +// Nothing lands on a tenant by omission (#583), and an invalid tenant is not +// backpressure a retry could clear. +func refusesATopicWithoutATenant(t *testing.T, h Harness) { + b := h.New(t) + for _, topic := range []mq.Topic{{Table: "events"}, {Tenant: "a.b", Table: "events"}, {Tenant: "*", Table: "events"}} { + err := b.Publish(ctx(t), topic, []byte("x")) + require.Error(t, err, "%+v", topic) + assert.NotErrorIs(t, err, mq.ErrQueueFull, "%+v", topic) + require.Error(t, b.ReplaySince(ctx(t), topic, time.Time{}, func([]byte) bool { return true }), "%+v", topic) + } +} + +// The trace context of the publishing request reaches the Subscribe handler +// through the message's headers, alongside any the options set. +func subscribeCarriesTheTraceContext(t *testing.T, h Harness) { + b := h.New(t) + got := make(chan context.Context, 4) + require.NoError(t, b.Subscribe(t.Context(), "hub-bridge", func(m *mq.Message) error { + got <- m.Ctx + return nil + })) + + sc := trace.NewSpanContext(trace.SpanContextConfig{ + TraceID: trace.TraceID{0x4b, 0xf9, 0x2f, 0x35, 0x77, 0xb3, 0x4d, 0xa6, 0xa3, 0xce, 0x92, 0x9d, 0x0e, 0x0e, 0x47, 0x36}, + SpanID: trace.SpanID{0x00, 0xf0, 0x67, 0xaa, 0x0b, 0xa9, 0x02, 0xb7}, + TraceFlags: trace.FlagsSampled, + }) + pubCtx := trace.ContextWithSpanContext(ctx(t), sc) + require.NoError(t, b.Publish(pubCtx, mq.Topic{Tenant: Acme, Table: "traced"}, []byte("x"), mq.WithHeader("X-Test", "1"))) + + select { + case c := <-got: + have := trace.SpanContextFromContext(c) + assert.Equal(t, sc.TraceID(), have.TraceID()) + assert.Equal(t, sc.SpanID(), have.SpanID()) + assert.True(t, have.IsRemote()) + case <-time.After(wait): + t.Fatal("the subscriber was never called") + } +} + +// Every tenant's events reach one Subscribe, whichever tenant published them. +func subscribeSeesEveryTenant(t *testing.T, h Harness) { + b := h.New(t) + got := make(chan mq.Topic, 8) + require.NoError(t, b.Subscribe(t.Context(), "hub-bridge", func(m *mq.Message) error { + got <- m.Topic() + return nil + })) + want := []mq.Topic{{Tenant: Acme, Table: "t"}, {Tenant: Globex, Table: "t"}} + for _, topic := range want { + publish(t, b, topic, "x") + } + var have []mq.Topic + timeout := time.After(wait) + for len(have) < len(want) { + select { + case topic := <-got: + have = append(have, topic) + case <-timeout: + t.Fatalf("timed out; delivered %+v", have) + } + } + assert.ElementsMatch(t, want, have) +} + +// Each tenant's events arrive in the order they were published, however the +// tenants interleave. +func eachTenantInOrder(t *testing.T, h Harness) { + b := h.New(t) + const n = 5 + for i := range n { + publish(t, b, mq.Topic{Tenant: Acme, Table: "a"}, fmt.Sprint(i)) + publish(t, b, mq.Topic{Tenant: Globex, Table: "b"}, fmt.Sprint(i)) + } + got, _, _ := consume(ctx(t), t, b, mq.ConsumerConfig{MaxAckPending: 100}, ackEach(t)) + order := map[tenant.ID][]string{} + for _, d := range next(t, got, 2*n) { + order[d.topic.Tenant] = append(order[d.topic.Tenant], d.data) + } + want := make([]string, n) + for i := range want { + want[i] = fmt.Sprint(i) + } + assert.Equal(t, want, order[Acme]) + assert.Equal(t, want, order[Globex]) +} + +// A Nak'd message comes back; a DoubleAck is confirmed. +func nakRedelivers(t *testing.T, h Harness) { + b := h.New(t) + publish(t, b, mq.Topic{Tenant: Acme, Table: "n"}, "x") + seen := 0 + got, _, _ := consume(ctx(t), t, b, mq.ConsumerConfig{MaxAckPending: 100}, func(m *mq.Message) { + seen++ // one tenant: one delivery goroutine + if seen == 1 { + assert.NoError(t, m.Nak()) + return + } + assert.NoError(t, m.DoubleAck(m.Ctx)) + }) + d := next(t, got, 2) + assert.Equal(t, "x", d[0].data) + assert.Equal(t, "x", d[1].data) +} + +// A message not acked within the consumer's AckWait is delivered again; one +// that was acked is not. +func ackWaitRedelivers(t *testing.T, h Harness) { + b := h.New(t) + publish(t, b, mq.Topic{Tenant: Acme, Table: "w"}, "acked") + publish(t, b, mq.Topic{Tenant: Acme, Table: "w"}, "left") + // Long enough that a DoubleAck under load lands inside it. + got, _, _ := consume(ctx(t), t, b, mq.ConsumerConfig{AckWait: 500 * time.Millisecond, MaxAckPending: 100}, func(m *mq.Message) { + if string(m.Data) == "acked" { + assert.NoError(t, m.DoubleAck(m.Ctx)) + } + }) + seen := map[string]int{} + for seen["left"] < 2 { + seen[next(t, got, 1)[0].data]++ + } + assert.Equal(t, 1, seen["acked"], "an acked message is not redelivered") +} + +// DeadLetter parks a delivered message under its own topic and leaves the +// original unacked: a Nak after parking still brings it back. +func deadLetterKeepsTheTopicAndDoesNotAck(t *testing.T, h Harness) { + b := h.New(t) + publish(t, b, mq.Topic{Tenant: Acme, Table: "t", Scope: "s"}, "x") + seen := 0 + got, _, _ := consume(ctx(t), t, b, mq.ConsumerConfig{MaxAckPending: 100}, func(m *mq.Message) { + seen++ + if seen == 1 { + assert.NoError(t, b.DeadLetter(m.Ctx, m, mq.WithHeader("X-Error", "boom"))) + assert.NoError(t, m.Nak()) + return + } + assert.NoError(t, m.DoubleAck(m.Ctx)) + }) + next(t, got, 2) + + counts, err := b.DeadLetterCounts(ctx(t), Acme, "") + require.NoError(t, err) + assert.Equal(t, map[string]uint64{"t": 1}, counts.Tables, "a scoped topic counts under its table") + assert.Equal(t, uint64(1), counts.Total) +} + +// Counts are per tenant and per table, every scope of a table under the table +// itself (so a dotted table name never shares a count with a table + scope +// pair), a table filter narrows Tables but not Total, and a tenant with +// nothing parked has zero counts. +func deadLetterCounts(t *testing.T, h Harness) { + b := h.New(t) + c := ctx(t) + + empty, err := b.DeadLetterCounts(c, Globex, "") + require.NoError(t, err, "a tenant with a budget and nothing parked") + assert.Equal(t, map[string]uint64{}, empty.Tables, "empty, not nil: the ops API encodes it as {}") + assert.Zero(t, empty.Total) + + park := func(topic mq.Topic, n int) { + for range n { + require.NoError(t, b.DeadLetter(c, mq.NewMessage(c, topic, []byte("x"), time.Now(), nil, nil, nil))) + } + } + park(mq.Topic{Tenant: Acme, Table: "t1"}, 2) + park(mq.Topic{Tenant: Acme, Table: "t2"}, 1) + park(mq.Topic{Tenant: Acme, Table: "t1", Scope: "s"}, 1) + park(mq.Topic{Tenant: Acme, Table: "t1.s"}, 1) + park(mq.Topic{Tenant: Globex, Table: "t1"}, 1) + + tests := []struct { + name string + id tenant.ID + table string + tables map[string]uint64 + total uint64 + }{ + {"every table", Acme, "", map[string]uint64{"t1": 3, "t2": 1, "t1.s": 1}, 5}, + {"one table, all of its scopes", Acme, "t1", map[string]uint64{"t1": 3}, 5}, + {"a dotted table", Acme, "t1.s", map[string]uint64{"t1.s": 1}, 5}, + {"a table with nothing parked", Acme, "none", map[string]uint64{}, 5}, + {"the other tenant", Globex, "", map[string]uint64{"t1": 1}, 1}, + } + for _, tt := range tests { + counts, err := b.DeadLetterCounts(c, tt.id, tt.table) + require.NoError(t, err, tt.name) + assert.Equal(t, tt.tables, counts.Tables, tt.name) + assert.Equal(t, tt.total, counts.Total, tt.name) + } + + unbudgeted, err := b.DeadLetterCounts(c, "initech", "") + if h.Caps.UnbudgetedNotFound { + require.ErrorIs(t, err, mq.ErrNoDeadLetterQueue) + return + } + require.NoError(t, err) + assert.Equal(t, map[string]uint64{}, unbudgeted.Tables) + assert.Zero(t, unbudgeted.Total) +} + +// A replay sends one topic's events in order from since on, and nothing of +// another table, scope or tenant. +func replaySince(t *testing.T, h Harness) { + b := h.New(t) + topic := mq.Topic{Tenant: Acme, Table: "r"} + publish(t, b, topic, "one") + // Waiting until "one" replays puts it before since in whatever store the + // backend replays from. + replayReaches(t, b, topic, 1) + since := time.Now() + publish(t, b, topic, "two") + publish(t, b, topic, "three") + publish(t, b, mq.Topic{Tenant: Acme, Table: "r2"}, "other table") + publish(t, b, mq.Topic{Tenant: Acme, Table: "r", Scope: "s"}, "scoped") + publish(t, b, mq.Topic{Tenant: Globex, Table: "r"}, "other tenant") + + tests := []struct { + name string + topic mq.Topic + since time.Time + want []string + }{ + {"since", topic, since, []string{"two", "three"}}, + {"everything", topic, time.Time{}, []string{"one", "two", "three"}}, + {"future", topic, time.Now().Add(time.Hour), []string{}}, + {"scoped", mq.Topic{Tenant: Acme, Table: "r", Scope: "s"}, time.Time{}, []string{"scoped"}}, + {"other tenant", mq.Topic{Tenant: Globex, Table: "r"}, time.Time{}, []string{"other tenant"}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + replayEventually(t, b, tt.topic, tt.since, tt.want) + }) + } +} + +// A done ctx ends a replay before the next event, with ctx's error. +func replaySinceStopsWhenContextIsDone(t *testing.T, h Harness) { + b := h.New(t) + topic := mq.Topic{Tenant: Acme, Table: "r"} + publish(t, b, topic, "one") + publish(t, b, topic, "two") + replayReaches(t, b, topic, 2) + + c, cancel := context.WithCancel(ctx(t)) + defer cancel() + var got []string + err := b.ReplaySince(c, topic, time.Time{}, func(data []byte) bool { + got = append(got, string(data)) + cancel() + return true + }) + require.ErrorIs(t, err, context.Canceled) + assert.Equal(t, []string{"one"}, got) +} + +// A replay that loses the broker before catching up says so, rather than +// passing for a caught-up one. +func replaySincePullFailureIsAnError(t *testing.T, h Harness) { + b := h.New(t) + topic := mq.Topic{Tenant: Acme, Table: "r"} + publish(t, b, topic, "one") + publish(t, b, topic, "two") + replayReaches(t, b, topic, 2) + + var got []string + err := b.ReplaySince(ctx(t), topic, time.Time{}, func(data []byte) bool { + got = append(got, string(data)) + assert.NoError(t, b.Close()) + return true + }) + require.Error(t, err) + assert.False(t, errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded), "not a ctx error: %v", err) + assert.Equal(t, []string{"one"}, got) +} + +// Delivery ended underneath a running Consume is reported on failed exactly +// once. +func failedOnceWhenDeliveryEnds(t *testing.T, h Harness) { + b := h.New(t) + got, _, failed := consume(ctx(t), t, b, mq.ConsumerConfig{MaxAckPending: 100}, nil) + // A delivery on each tenant proves the pulls are live before delivery is + // ended underneath them. + publish(t, b, mq.Topic{Tenant: Acme, Table: "t"}, "x") + publish(t, b, mq.Topic{Tenant: Globex, Table: "t"}, "x") + next(t, got, 2) + h.EndDelivery(t, b) + select { + case err := <-failed: + require.ErrorIs(t, err, mq.ErrDeliveryEnded) + case <-time.After(wait): + t.Fatal("delivery ended underneath the consumer and nothing was reported") + } + none(t, failed, "a second failure was reported") +} + +// A delivery the caller stopped is not a failure, even if delivery would +// have ended afterwards. +func failedNeverAfterStop(t *testing.T, h Harness) { + b := h.New(t) + _, stop, failed := consume(ctx(t), t, b, mq.ConsumerConfig{MaxAckPending: 100}, nil) + stop() + h.EndDelivery(t, b) + none(t, failed, "a stopped consumer reported a failure") +} + +func maxBytesReportsTheBudget(t *testing.T, h Harness) { + b := h.New(t) + for _, n := range []int64{32 << 20, 48 << 20} { + require.NoError(t, b.SetMaxBytes(ctx(t), Acme, n)) + assert.Equal(t, n, b.MaxBytes(Acme)) + } +} + +func stats(t *testing.T, h Harness) { + b := h.New(t) + s, err := b.Stats() + require.NoError(t, err) + assert.GreaterOrEqual(t, s.Connections, int64(1), "the broker's own connection") +} + +// A full queue refuses with ErrQueueFull; with per-tenant budgets, only its +// own tenant. +func queueFull(t *testing.T, h Harness) { + b := h.New(t) + h.Fill(t, b, Acme) + err := b.Publish(ctx(t), mq.Topic{Tenant: Acme, Table: "full"}, []byte("x")) + require.ErrorIs(t, err, mq.ErrQueueFull) + if h.Caps.PerTenantBudget { + publish(t, b, mq.Topic{Tenant: Globex, Table: "full"}, "x") + } +} + +// PurgeAcked never removes an unacked event. A backend that purges removes +// the acked ones past the cutoff; one that leaves retention to the operator +// reports nothing purged. +func purgeAcked(t *testing.T, h Harness) { + b := h.New(t) + topic := mq.Topic{Tenant: Acme, Table: "p"} + for _, data := range []string{"a", "b", "c", "left"} { + publish(t, b, topic, data) + } + got, _, _ := consume(ctx(t), t, b, mq.ConsumerConfig{MaxAckPending: 100}, func(m *mq.Message) { + if string(m.Data) != "left" { + assert.NoError(t, m.DoubleAck(m.Ctx)) + } + }) + next(t, got, 4) + + future := time.Now().Add(time.Hour) + purged, err := b.PurgeAcked(ctx(t), Durable, map[tenant.ID]time.Time{Acme: future, Globex: future}) + require.NoError(t, err) + if h.Caps.PurgesAcked { + assert.True(t, purged) + replayEventually(t, b, topic, time.Time{}, []string{"left"}) + return + } + assert.False(t, purged) + replayEventually(t, b, topic, time.Time{}, []string{"a", "b", "c", "left"}) +} + +func purgeAckedUnknownConsumer(t *testing.T, h Harness) { + b := h.New(t) + _, err := b.PurgeAcked(ctx(t), "no-such-consumer", nil) + require.ErrorIs(t, err, mq.ErrConsumerNotFound) +} diff --git a/internal/mq/mqtest/embedded_test.go b/internal/mq/mqtest/embedded_test.go new file mode 100644 index 00000000..98b619d9 --- /dev/null +++ b/internal/mq/mqtest/embedded_test.go @@ -0,0 +1,76 @@ +// The embedded broker's run lives here rather than in internal/mq so it is a +// test binary of its own, clear of that package's 15s budget. +package mqtest_test + +import ( + "os" + "path/filepath" + "testing" + "time" + + "github.com/Wave-RF/WaveHouse/internal/mq" + "github.com/Wave-RF/WaveHouse/internal/mq/mqtest" + "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/stretchr/testify/require" +) + +func TestEmbeddedNATS_Conformance(t *testing.T) { + mqtest.Run(t, mqtest.Harness{ + New: func(t *testing.T) mq.Broker { + e, err := mq.NewEmbedded(storeDir(t)) + require.NoError(t, err) + t.Cleanup(func() { _ = e.Close() }) + for _, id := range []tenant.ID{mqtest.Acme, mqtest.Globex} { + require.NoError(t, e.SetMaxBytes(t.Context(), id, 64<<20)) + } + return e + }, + // Closing the broker ends every tenant's delivery at once, the + // connection-closed half of the #587 path; internal/mq's own tests + // delete the durable, one tenant's queue and then another's. + EndDelivery: func(t *testing.T, b mq.Broker) { + require.NoError(t, b.Close()) + }, + // A tiny budget, then publishes until the tenant's own stream refuses + // even the smallest event, so no later one fits. + Fill: func(t *testing.T, b mq.Broker, id tenant.ID) { + require.NoError(t, b.SetMaxBytes(t.Context(), id, 4<<10)) + for _, size := range []int{1 << 10, 1} { + payload := make([]byte, size) + for i := 0; ; i++ { + require.Less(t, i, 1<<10, "the queue never filled") + err := b.Publish(t.Context(), mq.Topic{Tenant: id, Table: "f"}, payload) + if err != nil { + require.ErrorIs(t, err, mq.ErrQueueFull) + break + } + } + } + }, + Caps: mqtest.Caps{ + PerTenantBudget: true, + PurgesAcked: true, + UnbudgetedNotFound: true, + ConfiguresDurables: true, + }, + }) +} + +// storeDir is a temporary store directory whose removal retries briefly: under +// parallel load a consumer's state file can land after Close has returned, +// which fails t.TempDir's one-shot RemoveAll. The retrying cleanup runs first +// (cleanups are LIFO), leaving t.TempDir an empty directory to remove. +func storeDir(t *testing.T) string { + dir := filepath.Join(t.TempDir(), "store") + var err error + t.Cleanup(func() { + for range 50 { + if err = os.RemoveAll(dir); err == nil { + return + } + time.Sleep(20 * time.Millisecond) + } + t.Errorf("remove %s: %v", dir, err) + }) + return dir +} diff --git a/internal/mq/mqtest/mqtest.go b/internal/mq/mqtest/mqtest.go new file mode 100644 index 00000000..5978a5b6 --- /dev/null +++ b/internal/mq/mqtest/mqtest.go @@ -0,0 +1,121 @@ +// Package mqtest is the conformance suite for mq.Broker: the behavior the +// rest of the process relies on, stated once and run by every implementation +// from its own tests. The cases address events by mq.Topic alone and assume +// no layout — no stream, subject or partition names — so a backend passes by +// behaving, not by being built like the embedded one. Where backends +// legitimately differ, a Caps flag says which way; nothing else is optional. +package mqtest + +import ( + "context" + "testing" + "time" + + "github.com/Wave-RF/WaveHouse/internal/mq" + "github.com/Wave-RF/WaveHouse/internal/tenant" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/propagation" +) + +const ( + // Acme and Globex are the tenants Harness.New makes ready to publish. + Acme tenant.ID = "acme" + Globex tenant.ID = "globex" + // Durable is the consumer the suite creates, consumes and purges by: the + // ingest worker's (ingest.BufferConsumerName), so a backend that only + // finds durables an operator made has one to find. + Durable = "buffer-consumer" + // wait bounds every wait for something that should happen. + wait = 5 * time.Second + // quiet is how long a case watches for something that must not happen. + quiet = 300 * time.Millisecond + // retryPause spaces the polls of a backend whose replay store trails. + retryPause = 20 * time.Millisecond +) + +// Harness is what a backend gives the suite. +type Harness struct { + // New returns a fresh broker, isolated from every other New's, in which + // Acme and Globex can publish, budgets already applied as the wiring + // would. Its cleanup is registered on t and must tolerate the broker + // having been closed already. + New func(t *testing.T) mq.Broker + // EndDelivery ends delivery underneath a running consumer of Durable, as + // the broker's operator or the network could (the #587 failure path): + // deleting the durable, or closing the connection for good. + EndDelivery func(t *testing.T, b mq.Broker) + // Fill makes the next Publish for id refuse with mq.ErrQueueFull. nil + // skips the cases that need it. + Fill func(t *testing.T, b mq.Broker, id tenant.ID) + Caps Caps +} + +// Caps records where a backend's semantics legitimately differ. +type Caps struct { + // PerTenantBudget: a full queue refuses its own tenant alone, so Fill on + // one tenant leaves another publishing. + PerTenantBudget bool + // PurgesAcked: PurgeAcked removes acknowledged events past the cutoff, + // rather than leaving retention to the broker's operator. + PurgesAcked bool + // UnbudgetedNotFound: DeadLetterCounts of a tenant never given a budget + // is mq.ErrNoDeadLetterQueue rather than zero counts. + UnbudgetedNotFound bool + // ConfiguresDurables: CreateConsumer applies cfg.AckWait to the durable, + // rather than checking it against one the operator configured. + ConfiguresDurables bool +} + +type testCase struct { + name string + // need, when false, skips the case: the backend lacks what it checks. + need bool + run func(t *testing.T, h Harness) +} + +// Run runs every case against h, each as a parallel subtest on a broker of +// its own. It sets the global W3C trace-context propagator for its duration +// (the trace case needs one), so it must not be called from a parallel test. +func Run(t *testing.T, h Harness) { + prev := otel.GetTextMapPropagator() + otel.SetTextMapPropagator(propagation.TraceContext{}) + t.Cleanup(func() { otel.SetTextMapPropagator(prev) }) + + cases := []testCase{ + {"RoundTrip", true, roundTrip}, + {"RefusesATopicWithoutATenant", true, refusesATopicWithoutATenant}, + {"SubscribeCarriesTheTraceContext", true, subscribeCarriesTheTraceContext}, + {"SubscribeSeesEveryTenant", true, subscribeSeesEveryTenant}, + {"EachTenantInOrder", true, eachTenantInOrder}, + {"NakRedelivers", true, nakRedelivers}, + {"AckWaitRedelivers", h.Caps.ConfiguresDurables, ackWaitRedelivers}, + {"DeadLetterKeepsTheTopicAndDoesNotAck", true, deadLetterKeepsTheTopicAndDoesNotAck}, + {"DeadLetterCounts", true, deadLetterCounts}, + {"ReplaySince", true, replaySince}, + {"ReplaySinceStopsWhenContextIsDone", true, replaySinceStopsWhenContextIsDone}, + {"ReplaySincePullFailureIsAnError", true, replaySincePullFailureIsAnError}, + {"FailedOnceWhenDeliveryEnds", true, failedOnceWhenDeliveryEnds}, + {"FailedNeverAfterStop", true, failedNeverAfterStop}, + {"MaxBytesReportsTheBudget", true, maxBytesReportsTheBudget}, + {"Stats", true, stats}, + {"QueueFull", h.Fill != nil, queueFull}, + {"PurgeAcked", true, purgeAcked}, + {"PurgeAckedUnknownConsumer", true, purgeAckedUnknownConsumer}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if !c.need { + t.Skip("the backend's capabilities exclude this case") + } + t.Parallel() + c.run(t, h) + }) + } +} + +// ctx is a test's context with the suite's overall bound. +func ctx(t *testing.T) context.Context { + c, cancel := context.WithTimeout(t.Context(), 4*wait) + t.Cleanup(cancel) + return c +}