feat(mq): a Broker over an external NATS cluster - #636
Merged
Merged
Conversation
mq.backend, cache.backend, dedupe.backend and coord.backend select each layer's implementation; only today's in-process one exists per layer and it is the default. Validate refuses an unknown value, internal/app picks the implementation in one switch per layer, data_dir is probed only when a selected backend keeps state there, and boot logs Config.Warnings. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
New internal/coord: Coordinator/TryAcquire/Term with a fencing Token, Done/Err and Resign; RunElected for leader loops; Local, the in-process implementation; and coordtest.Conformance, the suite every backend runs. The sweeper now runs through RunElected under the "sweeper" lease, over a Local coordinator that wireCoord opens until coord.backend lands. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
A handoff overlap cannot lose ClickHouse data (every sweep stops at the ack floor) but can trim SSE replay history when the holders' settings views differ. Also lists coord/ in development.md's package tree. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
…ENTS.md Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
mqtest.Run states the mq.Broker contract as behavior, through the interfaces alone, so the external-NATS backend (#613) runs the same cases as the embedded one; mqtest.Caps covers the places where their semantics legitimately differ. The embedded broker passes it. The suite found that a durable deleted on several tenants' queues could report on failed more than once; fixed. The interface comments now allow a partition as the delivery unit, a CreateConsumer that finds rather than creates, an operator-owned retention, and zero dead-letter counts without a per-tenant queue. mq.ErrUnavailable is new, and the ingest handler answers it with 503 and Retry-After: 5. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
internal/mq's unit tests already take ~10s of their 15s budget under load, and the suite pushed them over. The embedded run moves to mqtest/embedded_test.go and ends delivery by closing the broker, so it needs no hook into mq's internals; the exactly-once failed report gets a deterministic test in internal/mq. The api.md rows for ErrUnavailable say that no backend returns it yet. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Deletes the durable on one tenant's queue, drains the report, then on the next: the real path, rather than calling fail by hand. The replay polls in mqtest pause between attempts. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
A durable deleted before its pull reaches the server ends nothing the client sees, so the exactly-once tests waited out their 5s and failed. Both wait for a delivery on each tenant first. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
wireCoord becomes a switch on coord.backend like the other layers, and New refuses a Config that names no coordinator. Docs stop calling the key reserved. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
…conformance # Conflicts: # docs/src/content/docs/api.md # docs/src/content/docs/architecture.md # internal/api/ingest.go # internal/mq/mq.go
roles (WH_ROLES, default api,ingest,sweeper) picks which components a process wires, and instance_id (WH_INSTANCE_ID, default <hostname>-<8 hex>) names it. Discovery, dedupe, the token verifiers, the hub bridge and keepalive stay per API process; the ingest worker is the ingest role; the sweeper is the sweeper role and stays lease-elected through a.elected. A process without api serves an ops-only router: probes, /version, metrics, and the settings reload behind the operator key alone. Boot refuses any split over the embedded MQ, and api without ingest (or the reverse) over a local cache. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Stressing the conformance suite showed two more flakes. A pull that raced the broker closing could end in a timeout, which ReplaySince read as caught up; it is an error now unless the connection is still open. The AckWait case tolerated no slow DoubleAck under load; its wait is 500ms and it counts deliveries instead of expecting an exact order. The embedded harness's store dir retries its removal, since a consumer's state file can land after Close under parallel load. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Review round 1: instance_id is only logged until a shared coord.backend records it; sweeper exclusivity across processes needs a shared coord.backend; the ops listener serves the probe aliases and answers 403 before 404 under /v1/ops; architecture.md's config and router sections cover roles and NewOpsRouter. The YAML roles test uses a non-default order so it can tell the file from the env default. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
…q-external-broker # Conflicts: # CHANGELOG.md
ExternalNATS implements every mq.Broker method over the operator-owned topology of #624 and creates nothing but auto-expiring consumers on the history stream. It passes the mqtest conformance suite connected as the shipped restricted wavehouse user. Its tests are integration-tagged, so the internal/mq unit binary keeps its 15s budget (#617). Not yet selectable: config and wiring come in D4. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Review round 1 on D3: - ReplaySince sends the events counted when it began, pulled in batches, rather than chasing the live tail one round trip at a time. - The duplicate-window floor covers every publish attempt, not two timeouts, so a publish retried twice is still stored once. - The source gauge reads -1 for a source that never attached (the server's -1ns). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
A replay fetches up to 256 events, then sends them to the SSE client with no pull waiting; a 5s inactive threshold let the server delete the consumer under a client slower than ~20ms/event, losing the rest of the gap-fill. The threshold is now a minute (a finished replay deletes its consumer anyway), pinned by a slow multi-batch replay test that fails on the old value. AGENTS.md: "purges" in the ExternalNATS invariant. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
NewEmbeddedMQ over t.TempDir() fails its one-shot RemoveAll when a consumer state file lands after Close, which failed internal/ingest in 4 of 5 make ci runs here under load (#442). StoreDir retries the removal, the same pattern the mqtest embedded run uses. Refs #442. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
|
Important Review skippedAuto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Advanced Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
…614) ## Summary This PR makes the query cache correct under concurrent writes and shareable across instances. It folds in #621, #634, #626 and #630, which were reviewed separately against this branch. - **Version snapshot at lookup (fixes #382).** The `cache.Cache` interface is now `Lookup(ctx, tenant, sha, deps) (Entry, Snapshot, error)` and `Set(ctx, Snapshot, value, ttl)`. A result's versions are read once, before its query runs and before the handler takes the tenant's ClickHouse pool, and the fill is filed under what was read. Previously `POST /v1/query` and pipe execution rebuilt the version-folded key after the query, so an insert that landed mid-query filed pre-insert rows under post-insert versions and they were served as fresh until their TTL. A reload that moves a tenant to another address or database now orphans a fill taken from the old pool the same way. The singleflight key and coalescing are unchanged. - **A tenant token on every key.** Every entry key folds the tenant's version, a pipe's dependency-free key included, so `InvalidateTenant` now drops a returning or moved tenant's cached pipe results too. A `Lookup` whose deps name another tenant is refused (`ErrForeignDependency`). `cache.Namespace` carries raw table and scope names and the cache escapes them itself with `internal/keyenc`, so `query.SafeEncodeToken` is gone and names that would run together under an unescaped join can no longer share a key. - **Flat local version index (part of #262).** `VersionManager` holds one version per tenant, per (tenant, table) and per (tenant, table, scope), bumped in place, so the index no longer grows with every bump. A tenant's version is a process-unique generation: `InvalidateTenant` drops the tenant's index and the next key gets a fresh generation. After every settings reload, `LocalCache.Prune` drops the index of each tenant no longer served. - **Write pipes run uncached (fixes #386).** A pipe whose bound SQL `IsMutation` classifies as a write skips the cache lookup, the fill and singleflight, and runs on every call. Before, a repeat within the TTL answered `200` without writing, and N concurrent identical calls became one write. The classifier now reads the leading keyword the way ClickHouse's lexer does (comments, quoted text, heredocs, the whitespace ClickHouse accepts), classifies a `WITH`-led statement by `INSERT INTO` alone, and looks through `EXECUTE AS`. An integration test checks every case, and every keyword in `system.keywords` in 12 `WITH` shapes, against ClickHouse's own parser. - **A Redis-compatible shared cache backend.** `cache.RedisCache` runs against Redis, Valkey, Dragonfly, ElastiCache and MemoryDB, standalone or cluster, using only `GET`, `SET` and `MGET`. Versions are random 8-byte tokens under the tenant's hash tag, and a lookup is one round trip. A lost token can only cause a miss. Values of 1 KiB or more are zstd-compressed, and stored values are capped at 1 MiB. Every operation has a 100 ms timeout, and a failure is a miss, a skipped fill or a deferred invalidation, never a failed query. A circuit breaker opens after 5 consecutive failures, or at once on a reply that refuses writes (`READONLY`, `OOM`, …) or the credentials, and only a successful probe write closes it. Deferred invalidations are retried until they land, and while a process owes one it bypasses the lookups that invalidation would orphan. Eight `wavehouse_cache_*` metrics, all labeled `backend="redis"`, report hits, round-trip time, breaker state, owed invalidations, value size and failed fills. - **`cache.backend: redis`.** A new `cache.redis` boot-config block (`WH_CACHE_REDIS_*`): `addrs`, `mode` (`standalone` or `cluster`), credentials, `db`, TLS files, `key_prefix`, `timeout` and `dial_timeout` (each capped at `1s`), `max_value_bytes`, `compress_min_bytes` and `version_ttl`. `wireCache` builds the backend from it. It is the shared cache that splitting the `api` and `ingest` roles into separate processes needs, and the boot error for such a split now names it; every split is still refused while the queue is embedded. The e2e suite now runs on Redis. ## Behaviour and compatibility notes - **Write pipes answer `X-Cache: BYPASS` with `Cache-Control: no-store` and are never coalesced.** Each call executes, so identical concurrent calls are that many writes. A failed write keeps the read path's status and `code` but always answers `retryable: false` with no `Retry-After`, since the statement may have run. Read pipes are unchanged. - **A write pipe does not invalidate cached reads of the table it writes** (#394, and #343 for read pipes), and its rows do not reach `/v1/stream` subscribers (#362). - **`InvalidateTenant` drops more than before**: a tenant's cached pipe results as well as its query results. Inserts still do not reach pipe results (#343). - **In-process cache keys changed** (the caller's query key is now escaped inside the entry key). They are in-process only, so a restart is the whole migration. - **Shared-cache token keys are a protocol between builds.** Every process sharing a server reads and bumps them for itself, so a later change to that layout needs a rolling-upgrade plan. A change to value keys only orphans entries and is safe to roll. - **Boot with Redis down or refusing the password succeeds, degraded.** The cache starts bypassed and keeps reconnecting. When a closed breaker opens, it logs one `WARN`, or one `ERROR` for rejected credentials. A failed probe reopening it logs at `DEBUG`, unless it failed for another cause than the one last logged (rejected credentials after a restart, say), which is logged at its own level. A long outage is one line. - **`mode: sentinel` refuses boot** until #656. A URL-style address is refused without echoing it, and a standalone server takes exactly one address. - **The ingest worker logs an invalidation that did not land at `WARN`**, not `ERROR`: the shared backend defers and retries it. `wavehouse_cache_invalidations_pending` is the signal to alert on. - **Run the server with an evicting `maxmemory-policy` and without persistence.** Under `noeviction` a full server refuses the token writes. Restoring a snapshot, or a crash-restart that reloads the last save, is a rollback that serves previously invalidated entries until their TTL. The deployment guide covers both. - **Known follow-ups:** - #662: a quoted placeholder lets a bound value break out of its literal. - #663: a write pipe answers `GET`, which proxies and clients may replay. - #666: `BACKUP`, `RESTORE`, `UNDROP` and `MOVE` pipes are not classified as writes. - #671: `SET`, `USE` and `EXECUTE AS` in a pipe leak into the pooled session. - Also still open: per-table scope cardinality (#262, until #235 populates `scope`), the rest of #664 (the probe reconnects one connection of several), Sentinel (#656), and the near-cache. ## Tests - **Conformance suite** `internal/testutil/cachetest.Run`: miss, hit and TTL, dependency order, tenant isolation, foreign deps, the scope lattice, raw names that would run together, `Invalidate`/`InvalidateTenant`, a bump during the query (#382), oversize values, zero snapshots, and concurrent use under `-race`. Shared backends also get two-instances-over-one-store cases. `LocalCache` runs it, and so does `RedisCache` against pinned Redis, Valkey, Dragonfly and a Redis Cluster node. - **#382**: on both cached routes, a bump from inside the ClickHouse call, and one from inside the pool lookup, each give MISS, MISS, HIT. The read's namespace and the ingest worker's bump are pinned to meet for raw table names. - **Flat index**: 10,000 rounds of interleaved bumps and key reads leave the index at its settled size. Generations never repeat, a table bump drops its scopes, and a bump against a tenant with no index records nothing. A reload prunes the index to the tenants still served. - **Write pipes**: `INSERT`, `WITH … INSERT` and `ALTER … DELETE` pipes each run on every call. Three identical concurrent calls are three writes in flight. A read pipe over a table named like a write verb stays cached. Every row of the ClickHouse error table on a write pipe answers `retryable: false`. Over the Redis-backed e2e stack, a write pipe called twice leaves both rows, and a failed one answers `400 clickhouse.rejected` with `retryable: false`. There are 152 table-driven classifier cases, each also checked against `EXPLAIN AST` on the pinned ClickHouse, and 17,928 keyword-named `WITH` statements where the classifier must agree with the parser. - **Redis backend (integration)**: lost tokens miss, and so does a flushed server. On a paused server, lookups fail within the bound, the breaker opens, invalidations are deferred, and all of it recovers. An owed bump holds its lookups. Refused writes (`READONLY`, `OOM`) open the breaker at once. A failover behind a stable address delivers the owed bump. A cluster topology read is bounded. A slow reconnect still closes the breaker, and a slow server stays bypassed. Rotated credentials open the breaker. A restored snapshot behaves as a rollback. Unit tests cover the key schema, the codec and its zip-bomb refusal, the breaker state machine, pending coalescing and collapse, and the breaker logging each opening once and a changed cause at its own level. - **Config and wiring**: defaults, env, YAML, the validation table, URL-style addresses (neither the address nor the secret is echoed), boot against a closed port (bypassed, not failed), and an unreadable TLS file refusing boot. Two `app.New` instances over one Redis: an ingest on one invalidates the other, and with Redis paused, queries bypass and still succeed. - Most behavioural tests are mutation-checked: each fails with the fix removed. - `make ci` passes: static checks, unit, integration, e2e on Redis, and coverage. Fixes #382. Fixes #386. Part of #262. Part of #613. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq --------- Co-authored-by: taitelee <taitelee@umich.edu> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
#625) Part of #613. This PR carries the whole remote-dedupe stack: the reserve/commit contract, windowed ingest, retention, and the DynamoDB backend with its boot wiring. #629, #633, #628 and #635 were reviewed on their own and merged into this branch, and #667's test fix came in with them. ## Summary - **Reserve, Commit, Release** (fixes #390, #222, #370). `Deduplicator.CheckAndMark` is replaced by a two-phase contract that every backend implements: - `Reserve(ctx, keys, lease)` claims each key atomically and answers `Claimed`, `Duplicate` or `InFlight` for it. It is all-or-nothing on error. - `Commit(ctx, claims, retention)` marks the claims as seen. `Release(ctx, claims)` gives them back, matched by token. - A claim that is neither committed nor released lapses after its lease, so a request that dies mid-publish never strands an id. - Keys are scoped per tenant and table in the readable `keyenc` format `<tenant>/<escaped table>/<escaped id>`, for example `acme/clicks/evt%2D123`. An id whose escaped form is over 1,024 bytes is stored as `#<sha256 hex>`. The same id in two tables is now two ids (#222). - Two concurrent requests with one id now publish once (#390). An explicit `null` id counts as missing (#370). - Pebble keeps pending claims in memory in 64 locked shards beside the instance, and writes each Commit in one batch with one fsync. - The conformance suite `internal/dedupe/dedupetest` runs every backend through the same cases. - **Windowed ingest** (fixes #384). Ingest runs in windows of up to 256 records. Each window makes one `Reserve`, publishes in record order, then makes one `Commit`. - A deduped record is published under a `Nats-Msg-Id` idempotency key derived from its tenant, table and id. The embedded ingest stream sets its duplicate window to 2 minutes explicitly. - Only a definite publish failure (the queue refused it) releases the claim. After an uncertain failure, the claim lapses with the lease, and a retry is dropped by the stream's duplicate window. The event is stored once and never lost. - A dedupe store that cannot answer (`dedupe.ErrUnavailable`, now "dedupe store unavailable") answers `503 {"error":"dedupe store unavailable"}` with `Retry-After: 5`. It used to answer `500 dedupe failed`. - On Pebble, a 1,000-record batch now costs 4 fsyncs instead of 1,000. - **Per-table retention** (fixes #220). `dedupe.retention` in `config.json`, with a per-table override in `dedupe.tables.<table>.retention`, sets how long a committed id stays a duplicate. The default is `"0"`, which keeps ids forever, so an existing `config.json` needs no change. - A finite retention below 2 minutes (the queue's duplicate window) is refused, not clamped. - A background sweep on the Pebble instance deletes expired ids and the old-format keys. It runs a minute after open and then hourly, and it never deletes a key that was committed again after the sweep read it. `wavehouse_dedupe_swept_keys_total{reason}` counts what it deletes. - **DynamoDB backend**, selected by `dedupe.backend: dynamodb`. Every tenant and every process share one table, so an id ingested through one pod is a duplicate through every other. - `Reserve` is a conditional `PutItem` per key. `Commit` is `BatchWriteItem` with retries. `Release` is a conditional `DeleteItem`. Expiry is the native TTL attribute `ex`, and correctness never waits on TTL. - Throttling, timeouts and connection failures wrap `ErrUnavailable` and answer `503`. A circuit breaker short-circuits `Reserve` for a second after five unavailable claims in a row. - New boot keys: `dedupe.lease` (the lease is now configurable, 30 s by default, at most 59 s with the embedded queue), `dedupe.reserve_concurrency`, and the `dedupe.dynamodb.*` block. `create_table` is refused unless `endpoint` is set, so WaveHouse never creates a table in AWS. - **Boot rule:** boot checks the table whether or not any tenant has dedupe on. A misconfigured table (missing, the wrong key schema, access denied) refuses boot only with a flat settings directory whose tenant has dedupe on. In every other case, including transient failures, nested directories, and no tenant with dedupe on yet, the process boots, and every tenant with dedupe on fails closed with the `503`. The check is retried in the background and again right after every reload. - **Caller-cancel fix** (addresses #648). When a caller cancels mid-`Reserve`, the puts not yet sent are skipped. A put already sent runs to its answer before it is released. Only its own call deadline can cut it off, and then it holds its key at most until the lease ends, as a crashed request's claim does. - **Test teardown** (absorbs #667). Tests that start the embedded broker no longer fail in `t.TempDir` cleanup when the broker's consumer-state flusher writes after `Close`. The new `internal/testutil/storedir` retries the removal. ## Behaviour and compatibility notes - **Old-format dedupe keys are swept, not migrated.** An id seen before the upgrade is accepted once more after it. The retention sweep deletes the old keys on its first pass. Nothing released depends on them. - A dedupe backend that cannot answer returns `503` + `Retry-After: 5` where it used to return `500`. The SDK already retries a `503`. - A mid-body read error or a prepare failure now drops the open window unpublished. Before, the records ahead of it were published. - The in-flight `503` sends the lease as `Retry-After`. That is 30 s by default, as before. ## Known follow-ups - #660: row-by-row isolation silently drops identical rows on a deduplicating table. - #665: a durable's last ack can land after `Close`, or never if the process exits. - #668: this PR does its three items: `config.embeddedDuplicateWindow` is pinned to `mq.EmbeddedDuplicateWindow` by a test, the `reserve_concurrency` wording is updated, and the lease rule is stated once. Close it by hand after this lands. - #651: cross-region dedupe on DynamoDB MRSC needs a sweeper for lapsed claims. - #652: accept events durably while the dedupe backend is down. ## Tests - **Conformance:** `dedupetest.Run` runs against Pebble twice (on an injected clock and on the real clock) and against `amazon/dynamodb-local:3.3.1`. It covers claim, duplicate and in-flight, release then re-claim, lease lapse, retention expiry, a 64-way concurrent Reserve, a Reserve racing a Commit, key isolation per tenant and table, hashed long ids, and all-or-nothing on a mid-call failure. - **Ingest** (`internal/api/ingest_window_test.go`, `ingest_retention_test.go`): - window boundaries and publish failures at chosen records; - the `503` for an unavailable store; - the #384 scenario end to end over the real broker and Pebble; - the uncertain-publish retry, mutation-checked against a missing idempotency key; - retention reaching `Commit` and changing on reload, including mid-window. - **Pebble sweep** (`internal/dedupe/sweep_test.go`): chunk boundaries, expired and old-format keys, and a Commit racing a chunk. Each case is mutation-checked. - **DynamoDB:** unit tests against a fake API cover error classification, Reserve cleanup, a Reserve cancelled by its caller leaving nothing claimed, a retried put keeping its own claim, Commit retries, the breaker, and `Check`. Integration tests against dynamodb-local cover the conformance suite, 32 clients racing one id, throttling, an unreachable endpoint, TTL and expiry, and two `app.New` instances sharing seen ids through one table. - **Boot wiring** (`internal/app/dedupe_dynamodb_test.go`, `internal/config/backends_test.go`): the table check in both directory shapes, the background retry, reloads that make no table call, and every new config key and refusal. - **Pinning tests:** the retention floor is at least the queue's duplicate window, and `config.embeddedDuplicateWindow` equals `mq.EmbeddedDuplicateWindow`. - `make ci` passes on the merged stack: every coverage gate passed. Unit 94.0%, integration 52.4%, e2e 60.2% (60% floor), Go total 95.0%. Fixes #390. Fixes #222. Fixes #370. Fixes #384. Fixes #220. Closes #442. Closes #648. Part of #613. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq --------- Co-authored-by: taitelee <taitelee@umich.edu> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ll (#680) ## Summary When a tenant's queue opened at runtime and the ingest worker's consumer could not join it, the broker reported that as the consumer's delivery ending. The worker failed and the process restarted, so one tenant's failure took ingest down for every tenant. The stream hub's consumer failing the same way was only logged, and that tenant's streams got no live rows until a restart. A queue a consumer cannot join is now left unrecorded as open, whichever consumer it is: - `SetMaxBytes` returns the error and `MaxBytes` does not report the budget, so the next reload retries the join. - The tenant's publishes answer `503` and retry the join too, through the existing pacing for a queue that cannot open (one shared attempt, then refusals for five seconds). - No other tenant is affected and nothing restarts. Unchanged: a consumer that had already joined and whose delivery then ends on its own still fails the worker, and a flat directory whose queue cannot open still refuses boot. ## Test plan - [x] One tenant's failed join refuses that tenant's publishes, reports nothing on `failed`, and leaves a second tenant publishing and consuming - [x] A later publish or reload joins the queue and the row reach - [x] Both cases run for the worker's consumer and for the hub's - [x] Existing tests still cover a terminal failure of a joined consumer and the flat boot refusal ## Related Issues Closes #675 Part of #583 <!-- Checklist for the author (not kept in the squash commit message): - `make ci` passes locally - Docs updated per AGENTS.md "Documentation & Consistency Sync" rules - CHANGELOG.md [Unreleased] entry added - Tests cover new / changed behavior (70 % minimum, 80 %+ preferred) The PR title is the squash commit subject — use Conventional Commits (`feat:`, `fix:`, `docs:`, `refactor:`, `test:`, `chore:`, `ci:`, `deps:`, `build:`, `perf:`, `revert:`, `style:`). The PR body below is the squash commit message, so keep it tight. -->
…at/mq-nats-wiring Collapses #654, #646 and #644 into #639. Conflicts in CHANGELOG.md and architecture.md kept both sides. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Takes #636's per-table dead-letter counts (deadletter.go) into #639 and clears its conflict with its base. CHANGELOG.md: #639's entries kept, with #636's deadletter.go wording applied to the external-broker line. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
…opology The stack absorbed #615, #618 and #622 before their review rounds ended, and main carries them as squashes. Merging the PR head main squashed (its tree is 5b6efe0's) brings their final content in with real ancestry, so the merge of main that follows only has to reconcile this stack's own changes. Conflicts: CHANGELOG entries and docs prose kept both sides; internal/config/backends.go keeps mq.nats/coord.nats and drops the env-default on the two backend keys, as main does. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
…topology Same reason as the #622 head: main squashed it (its tree is 5004cd2's), and the stack held an earlier round. Conflicts kept this stack's changes over the final content: docs and CHANGELOG prose, the roles tests' comments, Warnings' nats checks, validateTopology's nats rules. The dead-letter count cases take main's (every scope under its table, a dotted table counted apart), which ExternalNATS shares through deadLetterTables. configuration.mdx kept one "Process roles" section. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Brings in #614 (shared redis cache), #625 (dedupe reserve/commit, dynamodb, dedupe.lease, WithIdempotencyKey), #680 (a failed queue join is the tenant's alone) and #623's squash, whose content the stack already had from its PR head. Resolutions: - config: Warnings keeps the nats warnings for every role and main's api-only cache/redis/dedupe ones; the unknown-backend test lists nats beside redis and dynamodb. - mq: main's WithIdempotencyKey, ErrUnavailable and Consume contracts; the conformance suite is main's, idempotency case included. - testutil: main's storedir replaces this stack's StoreDir. - ingest: main's publishFailed already answers ErrUnavailable with 503. - integration setup keeps startNATS beside startDynamoDBLocal; the Makefile runs internal/cache with the tagged suites as main does. - docs and CHANGELOG keep both sides; the backend lists name nats, redis and dynamodb together. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Since #632 an env-default tag is refused: cleanenv re-applies it to a YAML zero, so `mq.nats.partitions: 0` came back 1. The block's defaults now come from defaultMQNATS() inside defaults(), each non-zero one has a zeroCases entry (the block is validated only under backend=nats, so its zeros load as written), the docs test reads a derived default such as `<prefix>_coord` and *(none)* as the zero, and the TLS pair gets a row per key. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
publish set jetstream.WithMsgID(nuid.Next()) on every publish, which overwrites the Nats-Msg-Id header WithIdempotencyKey sets, so ingest's retry of an uncertain publish was stored twice under mq.backend=nats. The caller's key is now the message id; the conformance suite's IdempotencyKeyDropsARepeat, from main, pins it for this broker too. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
…tention Under mq.backend=nats the operator owns the duplicate window, so main's caps for the embedded queue (dedupe.lease within 2m, dedupe.retention at least 2m) do not reach it. The verifier now requires every partition's duplicate_window to cover the lease's republish span (the lease, the lease rounded up to a second, and a second: config's rule for the embedded window), and boot and every reload warn about a served tenant with dedupe on whose finite retention, default or per table, is under the partitions' shortest window (ExternalNATS.DuplicateWindow, Store.DedupeRetentions). The NATS wiring moves out of wire.go into wire_nats.go (wireNATSMQ, the new wireNATSCoord, coordBucket and the retention warning), which the e2e coverage gate excludes as it does wire_dynamodb.go: the e2e binary never runs mq.backend=nats. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Main's docs still said the message queue is embedded wherever it described what instances share, that a split needs a shared cache whatever the roles, and listed only cache.redis and dedupe.dynamodb as a shared backend's connection. They now name mq.nats and coord.nats too, and dedupe.retention says to cover a nats partition's duplicate_window as well. startNATS in the integration suite waits for the host port as well as the log line, which can come first when several containers start at once. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
`wavehouse mq manifests` gains --dedupe-lease (default 30s), and the generated partitions' duplicate_window covers it, so a lease longer than 59s no longer yields manifests the boot check refuses. Docs: the external dependencies name every shared backend; the DLQ section says what differs under nats; the per-tenant-queue claims in the multi-tenant guide and the ingest pipeline are scoped to the embedded broker; the ops listener's 403/401 matches the router; the ignored cache.redis block's warning is the api role's. CHANGELOG states the net behaviour instead of a warning no release logs. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
…ease Main's #625 made dedupe.lease and dedupe.reserve_concurrency required (> 0); a Config built without Load carries neither, so Validate refused the NATS integration tests' configs before they booted. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
make ci measured the e2e coverage gate at 59.7%, under its 60% floor: the e2e stack never selects mq.backend: nats, so the mq.nats and coord.nats checks, their Warnings lines and the two nats rules of validateTopology were statements it could not reach. They move unchanged into internal/config/mq_nats.go (natsWarnings, validateNATSTopology, CoordNATSConfig.validate, trimURLs), excluded from the e2e gate only, as cache_redis.go is; the unit and merged totals still count them. No behaviour changes. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Same reason as wire_nats.go, and the same move #645 made: the e2e binary runs every role, so wireOpsAuth and wireOpsHTTP were statements the e2e gate counts but can never reach. Pure move, excluded from the e2e gate only; the unit and merged totals still count them. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Part of #613. This is PR D3 of the external-NATS workstream.
Stacking: the base is
feat/mq-nats-topology(#624, D2). This branch also mergesfeat/mq-conformance(#623, D1), so the diff contains #623's commits too. Review it after #623 and #624. This PR's own change is everything after the merge commitaa9a6553.What
mq.ExternalNATS(internal/mq/external.go) is anmq.Brokerover an operator-owned NATS cluster: the topology in #624 (N interest-retention ingest partitions, a history stream that sources them, and one dead-letter stream). It never creates, changes, purges or deletes a stream or a durable. It creates only auto-expiring consumers on the history stream: one perSubscribe(the hub) and one per replay.Nothing selects it yet. The config (
mq.backend: nats,mq.nats.*), thewireMQcase, the integration test throughapp, and the user-facing docs are D4, on G1.NewNATScan be constructed and is tested. Onlyinternal/mqimports NATS, as before.NewNATSconnects and then runs feat(mq): external nats broker with sharded, pinned ingest workers #624'sawaitNATSTopologyuntilTopologyWaitruns out. The wait also bounds each check's requests. Recommended findings are logged. If a required finding is still there when the wait ends, boot fails with a*TopologyErrorthat lists every finding. A cluster never reached isErrUnavailable(with passwords in the URLs redacted). Partition and dead-letter streams are then found by subject.wh.ingest.<fnv32a(tenant)%N>.<tenant>.<table>[.<scope>]withWithExpectStream(partition), and each attempt is bounded byPublishTimeout.Nats-Msg-Id, so the partition stores it once.maximum bytes exceededandmaximum messages per subject exceeded→ErrQueueFull.ErrUnavailable. A deleted stream is also logged atErrorwith its name and dropswavehouse_mq_topology_okat once.CreateConsumermapsbuffer-consumer(or the operator's own name) towh-ingest. It finds that durable on every partition, checksack_wait ≥ cfg.AckWaitandmax_ack_pending > 0, and never creates it. Any other name isErrConsumerNotFound.Consumeruns one pull per partition and splits the prefetch between them. A delivery the client ends on its own (durable deleted, connection closed for good) is reported onfailedexactly once, and never afterstop(mq: propagate terminal JetStream consumer errors to the ingest worker #587).Subscribeuses a per-pod ordered, ack-less consumer on the history stream withDeliverNew.consumerNameis ignored. A handler error is logged, because there is no redelivery to ask for.ReplaySinceuses an ack-less consumer on the history stream, filtered to the topic, starting atsince. It pulls in batches of 256 and counts down from the events the server counted at the start, so it never chases the live tail: the SSE handler has already subscribed live before gap-fill, so tail events would be duplicates. A pull or connection failure before it finishes is an error. The consumer is deleted afterwards, and its 1-minute inactive threshold outlasts sending one batch to a slow client.DeadLetterpublishes towh.dlq.<topic key>with a fresh id.DeadLetterCountsis one subject-filtered info call on the shared stream. A tenant that has parked nothing gets zero counts.PurgeAcked: returns(false, nil)with no I/O. It warns once per tenant whose gap window is longer than the history'smax_age. A consumer name that does not map isErrConsumerNotFound.SetMaxBytes: records the budget and logs once that budgets are not enforced per tenant. Enforcing them is D6.Close: stops every consumer (so nothing reportsfailed) and the watch loop, then drains the connection.wavehouse_mq_connected,wavehouse_mq_topology_ok, and per sourcewavehouse_mq_history_source_lagandwavehouse_mq_history_source_last_active_seconds. The source lag matters beyond SSE, because a stalled source holds rows on every partition (S1).Additions and deviations from the design
active: -1).Activekeeps counting up from the last contact until the source re-attaches. In one probe, one partition's source re-attached at ~11s and the other three were still unattached at 14s. An "attached" gauge would have read 1 throughout. So a source re-attaching is not a topology fault, as S1 asked. The verifier's one transient finding ("not attached yet",active < 0, which is only seen before a source's first attach) is markedtransientin feat(mq): external nats broker with sharded, pinned ingest workers #624'sFindingand kept offtopology_ok."", which gave a false required finding, "cannot read server version". The connection has its own gauge.Closecancels an in-flight re-check. Without this, a re-check started during an outage could holdClosefor up to 2 × 30s.NATSConfigembeds feat(mq): external nats broker with sharded, pinned ingest workers #624'sNATSTopologyinstead of repeating its fields flat. D4 fills both halves frommq.nats.*.//go:build integrationininternal/mqitself. They need feat(mq): external nats broker with sharded, pinned ingest workers #624's fixture, which lives in packagemq's test files, and the conformance run has to be inmq_testbecausemqtestimportsmq. Soexternal_export_test.goexposes the fixture to it.make test-integrationruns these tests by name (^Test(ExternalNATS|NewNATS|NATSPermissions_Refuse)) in a secondgotestsumcall, so the package's untagged tests stay in the unit suite only..testcoverage.ymlexcludesexternal.gofrom the unit and e2e per-suite gates. The merged total still counts it.≥ 2 × publish_timeout, but the last retry arrives3·T + 2·250msafter the first attempt.minDuplicateWindow()is now shared by the verifier and the manifest generator. The shippedjetstream.yamlis unchanged (2m is more than 15.5s).testutil.StoreDir(refs fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main #442):NewEmbeddedMQnow removes its store with retries, the same patternmqtestuses. The fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main #442 cleanup race failedinternal/ingestin 4 of 5make ciruns on this branch under heavy load (measured; 0/8 in isolation on both base and HEAD). This is a mitigation, not the drain-before-return fix fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main #442 asks for, so fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main #442 stays open.Finding.transient(unexported) and the fixture keeping itsOptions, so a test can restart the server.Tests (all integration-tagged; ~6s together)
mqtest.RunagainstExternalNATS, connected as the restrictedwavehouseuser fromdeployments/nats/values.yaml, with each case on its own server and the shipped N=4 manifests. Caps:PerTenantBudget,PurgesAcked,UnbudgetedNotFoundandConfiguresDurablesare all false.EndDeliverydeleteswh-ingeston every partition as the operator.Fillshrinks the tenant's partition to 4 KiB. 18/18 applicable cases pass.AckWaitRedeliversis skipped by its cap. This proves the published permission set for publish, consume, ack, DLQ, replay and the hub, beyond feat(mq): external nats broker with sharded, pinned ingest workers #624's read-only checks (measured). The run was stable over-count=8under heavy load.ErrUnavailableand notErrQueueFull. After a restart, publishing resumes, the consumer delivers again, andfailedstays quiet.Nats-Msg-Idand the partition holds one message. A new publish gets a new id. Losing every answer isErrUnavailable.max_msgs_per_subjectisErrQueueFull, and only that topic.ErrUnavailablenaming the stream, and the gauge drops. The stream is not recreated.ErrTopologynaming it. An unreachable server isErrUnavailablewithin the wait. Conflicting options are refused.ErrUnavailable.topology_ok, and recreating it restores it. After a server restart the sources read as silent andtopology_okstays 1.PurgeAckedwarns once per tenant.CreateConsumerchecks the durable and never creates one.wavehouseuser is refused stream create, update, purge and delete, a durable on a partition, and deletingwh-ingest.WH_MQ_MEASURE=1; in-process single-node server, file storage, under heavy load; inferred to be optimistic for an R3 cluster):Evidence
make ci: green on 2757e64 (unit 93.3%, integration 51.6%, e2e 61.2%, Go total 94.0%). Earlier runs on c8ad3f0 failed only ininternal/mq's 15s unit budget (test(app): internal/app unit tests use 8–17 s of the 15 s budget, time out under load #617, under heavy load) and in the fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main #442 ingest cleanup race, which the last commit mitigates.pre-push-reviewer(opus), four rounds. Round 1 raised three [SHOULD]s: the source gauge read -1e-9 not -1; the duplicate-window floor was below the retry schedule; ReplaySince made a round trip per event and chased the tail. Round 2 raised one [SHOULD]: batching under a 5s inactive threshold lost a slow client's gap-fill (measured). Rounds 3 and 4 (theStoreDircommit) were ship_it.docs-reviewer(opus): ship_it in round 1; one [MAY] in round 2 ("purges" in AGENTS.md); ship_it in round 3 (c8ad3f0). It was skipped on record for 2757e64, a test-helper-only commit.Deliberately left to later PRs
mq.backend/mq.nats.*config and validation, thewireMQcase,tests/integration/mq_nats_test.gothroughapp, and the docs pages (configuration, deployment "External NATS (Kubernetes)", architecture, ingest-pipeline, durability, api).Sharded/ConsumerConfig.Shards, so a worker consumes only the partitions it claims.SetMaxBytesonly records.Follow-up: dead letters counted per table (7847c14, 50a1170)
DeadLetterCountson bothExternalNATSandEmbeddedNATSfolded a scoped topic into the nametable + "." + scope, so a table literally nameda.bcould share one count with tableascopedb, and the?table=filter matched only a table's unscoped subject. Both brokers now key their per-subject counts through onedeadLetterTables(internal/mq/deadletter.go): every scope of a table counts under the table itself, and the?table=filter keeps all of that table's scopes.deadletter.goand its test are byte-identical to the same files on #655, so they merge without conflict (measured withgit merge-tree; the branches as a whole still conflict, since this one predates main's #612).Totalis unchanged on each broker (external sums every subject of the tenant; embedded uses the stream's message count). Scope is always empty today, soGET /v1/ops/dlq/statsdoes not change visibly; themqtestconformance suite (deadLetterCounts,deadLetterKeepsTheTopicAndDoesNotAck) pins the new behavior on both brokers, and #623's copy of the suite has to flip the same way (noted on #623).Evidence for this follow-up:
make cigreen locally at 50a1170 withJOBS=1. At higher parallelisminternal/mq's 15 s unit budget trips on this branch, which predates main's test(mq,app): fit the unit budget; fail fast on an uncreatable store #647 speed-up; the mq unit tests took 12.6 s before this change and 12.3 s after it, unloaded, so the change does not cause it.pre-push-revieweranddocs-reviewership_it at 50a1170 (round 1 at 7847c14 found the code correct; its findings — this addition's wording and the changelog's file list — are fixed).🤖 Generated with Claude Code
https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL