Repository navigation
feat(app): mq.backend selects embedded or external NATS - #639
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
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
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
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
…-nats-wiring # Conflicts: # .testcoverage.yml # AGENTS.md # CHANGELOG.md # docs/src/content/docs/architecture.md # internal/app/wire.go
mq.backend: nats wires mq.ExternalNATS from a new mq.nats boot-config block (file-path-only credentials, TLS, topology). Role splits boot on it; coord.backend=local and the unapplied mq.max_bytes_gb are warnings. internal/mq/natstest stands NATS up from the shipped deployments/nats files for tests outside internal/mq, and an integration test boots two processes on a NATS container. Docs: External NATS deployment guide, the mq.nats reference, the ops listener, and the nats-mode 503/DLQ. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
A user, password or token in a NATS URL is an inline secret and would sidestep the one-auth-method check. Docs: changing N regenerates the manifests (partition metadata), publish rights on the ingest subjects are trusted as WaveHouse, and embedded-only claims are scoped. 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 |
ExternalNATS's worker consumed only partitions 0..N-1, so rows left in a partition an operator removed by lowering mq.nats.partitions were never written, while the verifier's finding said they were drained. It now also consumes, through wh-ingest, every stream holding ingest subjects outside the N partitions, and the operator deleting one once it is empty ends only that stream's delivery. The finding counts the stream's rows and says when it has no durable to drain it. Close logged "disconnected from nats; reconnecting" at WARN with a nil error; a deliberate close now logs nothing, a lost server still warns. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Review round 1: the shrink test keeps the old-N process publishing during the rollout and waits for the removed partition's delivery to end; the split of prefetch counts only the N partitions; the docs say how to delete a drained stream under preventDelete and qualify the other 'durable deleted ends the worker' claims. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
coord.backend: nats holds leases as keys in an operator-owned KV bucket on the mq.nats connection (ExternalNATS.Leases). The KV revision a term was taken at is its fencing token; a candidate takes another holder's lease only after seeing the same revision unchanged for the lease duration on its own clock, never by server TTL. The bucket joins the topology spec, verifier, manifest generator (nack KeyValue) and the shipped permissions. Boot now refuses coord.backend=nats without mq.backend=nats (rule 3), and mq.backend=nats with coord.backend=local in a sweeper process (rule 4, previously a warning). Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
This was referenced Sep 25, 2026
The renew deadline now runs from when a stored renewal was sent (and from before the acquiring write), on a timer rather than the next tick, so a slow answer cannot let a candidate take over while the holder still runs. Resign cancels a renewal in flight and deletes its own later write, so it honours its ctx. TryAcquire no longer holds the coordinator's lock across requests. Docs: --coord-bucket, allow_direct, takeover window, stale lines. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
…udget Resign's delete of a renewal it cut short, and the lost-reply adoption, matched the coordinator's value, so they could take a later term of the same coordinator for their own. Each term now writes its own id. Tests for a cut-short renewal and for Close during a campaign. The bucket verifier cases and the app missing-bucket boot test move to integration-tagged tests: internal/mq and internal/app sit at their 15s unit budget (#617). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Part of #613. This is PR **D1** of the external-NATS workstream. It is based on `main` (#612, which it was stacked on, has merged). ## What - **`internal/mq/mqtest`** (new) is the conformance suite for `mq.Broker`. `mqtest.Run(t, Harness{New, EndDelivery, Fill, Caps})` states the contract as behavior and uses the interfaces only, with no stream, subject or partition names. It covers: - round trips with names that need encoding; - that a topic without a tenant is refused; - trace context reaching `Subscribe`, and `Subscribe` seeing every tenant; - per-tenant order; - `Nak` and `AckWait` redelivery; - that `DeadLetter` keeps the topic and does not ack; - per-tenant, per-table dead-letter counts, every scope of a table counted under the table itself (#655), and an empty (never nil) `Tables` when nothing is parked; - replay bounds and isolation, ctx cancellation, and that a failed pull is an error; - exactly one `failed` report, and none after `stop`; - `MaxBytes`, `Stats`, `ErrQueueFull`, and `PurgeAcked` semantics. Each case runs as a parallel subtest on a fresh broker. - **`mqtest.Caps`** flags the four places where the external backend legitimately differs: - `PerTenantBudget` - `PurgesAcked` - `UnbudgetedNotFound` (DLQ counts of a tenant never given a budget) - `ConfiguresDurables` (whether `CreateConsumer` applies `AckWait` or only finds an operator-made durable) - **The embedded broker passes the suite.** Its run is `internal/mq/mqtest/embedded_test.go`. - **`mq.go` contract wording** is updated as the design specifies: - The delivery unit is "a tenant's queue, or the partition that holds it". - `ErrQueueFull` is a byte limit, and the per-tenant no-queue case applies only to an implementation that opens queues per tenant. - `DeadLetterCounts` may return zero counts in place of `ErrNoDeadLetterQueue`. - `CreateConsumer` may find rather than create. - `PurgeAcked` may remove nothing. - `Subscribe` guarantees delivery only for events published after it returns. - **`mq.ErrUnavailable`** (new) is mapped by the ingest handler to `503` + `Retry-After: 5`. Before, it would have been the `500` "publish failed". No backend returns it yet; D3's will. `api.md` says so. - **Two embedded bugs found by the suite are fixed:** - A durable deleted on several tenants' queues could report on `failed` more than once. A CAS now allows one report, and it is pinned by a test that deletes the durable on real queues one after another. - `ReplaySince` read a pull that raced the connection closing as "caught up". It is now an error unless the connection is open. ## Deviations from the design doc - **The embedded run lives in `internal/mq/mqtest/embedded_test.go`, not `internal/mq/embedded_conformance_test.go`.** `internal/mq`'s unit binary already takes about 10s of its 15s `-race` budget when the machine is idle, and 21–34s under heavy load on the base branch alone (measured). The suite in that binary pushed it over. In its own binary it takes about 2.5s. It also no longer needs an `export_test.go` hook into mq's internals. - **`Harness.DeleteIngestDurable` became `Harness.EndDelivery`,** which ends delivery under a running consumer. The embedded harness closes the broker; D3 should delete the durable. The durable-deletion path for embedded is covered in `internal/mq`'s own tests. - **`Caps.NeverParkedNotFound` became `UnbudgetedNotFound`.** Embedded returns zero counts for a budgeted tenant that has parked nothing. Only a tenant with no budget gets `ErrNoDeadLetterQueue`. - **New cap `ConfiguresDurables`.** The AckWait-redelivery case cannot pass against an operator-made durable with a 60s `ack_wait`. - **The suite uses a fixed `mqtest.Durable = "buffer-consumer"`.** It is the worker's name, so a backend that maps durable names has one to find. - **`Run` sets the global W3C propagator for its duration.** The trace case needs it, so `Run` must not be called from a parallel test. ## Left to later PRs - **D2:** topology spec and verifier, manifests, S1. - **D3:** `ExternalNATS` plus `external_conformance_test.go`, which runs `mqtest.Run` with its own `Caps`. - **D5:** the `Sharded` cap and the shard-subset case. `ConsumerConfig.Shards` does not exist yet, so neither is here. - **D4:** docs for the nats backend. ## Evidence - `make ci`: green on 6ebfc6f, after the merge of `main` with #655 (all coverage gates passed; Go total 94.4%). - The suite passed 40/40 under `-race -count=20 -cpu 1,4` (before the #655 merge). - Pre-push reviewers: `pre-push-reviewer` and `docs-reviewer` both returned `ship_it` before the #655 merge, after 6 and 4 rounds. After it, `pre-push-reviewer` returned `ship_it` on 6ebfc6f (one round of review fixes: the nil-map assertion); the merge left the PR's own docs unchanged. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --------- Co-authored-by: taitelee <taitelee@umich.edu> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…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 D4 of the external-NATS workstream:
mq.backend: natsselectsmq.ExternalNATS.Stacking: the base is
feat/mq-external-broker(#636, D3). The branch also mergesfeat/process-roles(#622 C1, which carries #615 B1 and #618 G1), so the diff contains those PRs' commits too. Review it after #618, #615, #622 and #636. This PR's own change is everything after the merge commitdbd4661c(the merge kept both sides:wireMQis G1's switch with D3's bootctx).What
internal/config/backends.go):mq.backendtakesnats(MQNATS). The newmq.natsblock (MQNATSConfig,WH_MQ_NATS_*) holdsurls,name,creds_file/nkey_seed_file/user+password_file,tls.{ca_file,cert_file,key_file,server_name,handshake_first},js_domain,subject_prefix(wh),partitions(1),ingest_consumer(wh-ingest),history_stream(empty =<PREFIX>_HISTORY),connect_timeout(5s),publish_timeout(5s) andtopology_wait(60s).mq.nats.passwordorWH_MQ_NATS_PASSWORDis refused as an unknown key or variable.MQ.validatechecks the block only whennatsis selected. It refuses: no URLs, an empty URL, a prefix outside[a-z0-9_-], a URL carrying credentials (@), fewer than one partition, an empty durable name, more than one of creds/nkey/user, a password file without a user, half a cert pair, and a timeout that is not positive.internal/app/wire.go):wireMQgets anatscase,wireNATSMQ. It buildsmq.NATSConfigand callsmq.NewNATSunderNew's ctx. It hands over noSetMaxBytes: the embedded reconcile stays inwireEmbeddedMQ. Both cases shareadoptMQ, which registers the close component and the system gauges.validateTopology): rule 2 refuses a split only on embedded, so a split now boots onnats. Rule 5 still refusesapiwithoutingest(or the reverse) over a local cache. What boots today isapi,ingestreplicas plus asweeper-only process.Config.Warnings()logs these once each at boot:mq.backend=nats+coord.backend=localin a process runningsweeper: every such process holds its own lease. This is harmless becausePurgeAckedremoves nothing under nats. The message says a sharedcoord.backendwill be required once one exists.mq.max_bytes_gbis not applied under nats.mq.natsblock underembeddedis ignored.internal/mq/natstest(new, test code outside_test.go, excluded from coverage likemqtest) stands NATS up from the shippeddeployments/natsfiles:ServerConfigrendersvalues.yaml'sconfig.mergeas a nats.conf,LoadManifestsparsesjetstream.yaml(strictly, so a new generator field fails loudly),Operator.ApplyShipped/DeleteDurable, andStartruns an in-process server. Why an export: tests outsideinternal/mqmay not import NATS (depguard), and D2's fixture lives in packagemq's_test.gofiles, so neitherinternal/appnortests/integrationcould reach it.internal/mq's own fixture now builds onnatstest(one parser, one config renderer), so the shipped manifests stay the single fixture.Additions and deviations
mq.NATSTopologyare not boot config. They are what the ingest worker asks of the durable, so they are left at mq's defaults, which equal the worker's constants. Exposing them would let an operator "tune" a value the worker does not read.history_streamdefault is empty (→<PREFIX>_HISTORY), not the design's literalWH_HISTORY, so a non-default prefix matches whatwavehouse mq manifests --prefixgenerates.wavehouse_nats_connections/wavehouse_nats_in_msgs_totalare registered for the external broker too. They read this client's connection (0/1) and its received count, and the docs say so. Without that, the Pebble gauges registered by the same call would vanish in nats mode.roles=api+roles=ingestsplit is refused by C1's rule 5 until a shared cache (E) exists, so it cannot boot yet. That test is left to C2. The integration test instead runs process A (api,ingest,sweeper) and process B (api,ingest), a split the embedded MQ refuses. B's hub gets live events that A ingested (reconciliation.md ci: bump actions/setup-go from 5 to 6 #3).Tests
internal/config/mq_nats_test.go:defaultMQNATS), every env variable, YAMLunboundEnvknows all 19 variablesinternal/app/mq_nats_test.go(in-processnatstest.Start, unit build):api,ingestprocess wires*mq.ExternalNATS, the expected components and nothing underdata_dir/natswh-ingest, and the operator deleting it endsRunwithingest worker: …ErrDeliveryEndedErrUnavailableErrTopologytests/integration/mq_nats_test.go,TestNATSBackend_EndToEnd: anats:2.14.6-alpinecontainer configured fromvalues.yaml, withjetstream.yamlapplied asnackbefore WaveHouse starts. Two processes connect as the restrictedwavehouseuser over a nested directory with two tenants (two ClickHouse databases). It checks:since=) from the historyglobexcounted under?tenant=globexonly, withacmeat zero on the shared DLQwh-ingestends both processesDocs
deployment.md: a new External NATS section covering:wavehouse mq manifestsand apply order: the history before WaveHouse publishes (S1: rows acked before the source attaches never reach it), and never a partition without its durablediscard: old(S1: history availability couples to ingest admission),duplicate_window,max_deliver: -1wavehouseuser's permissionsconfiguration.mdx: themq.natstable, a Boot warnings section, and the examples.api.md: an ops listener section (C1 owed it), the503+Retry-After: 5as nats-only (no longer "reserved"), a full partition refusing the partition's tenants, and dlq stats200/zeros under nats.architecture.md(external.go, the topology files, natstest, the wireMQ nats case, backends.go),ingest-pipeline.md(scaling section rewritten for what exists, sweeper no-op),durability.md(external contract),settings-directory.mdx(mq.max_bytes_gb,gap_window_minutesunder nats),development.md,AGENTS.md, rootconfig.yaml, CHANGELOG. Stale "not yet selectable" / "nothing returns it yet" lines in the D1–D3 and C1 CHANGELOG entries were updated.Integration items (not in this PR)
mqtesthas no idempotency-key conformance case. The lease must be ≤ the operator'sduplicate_window, and a tenant's finite dedupe retention shorter thanduplicate_windowshould warn. Neither is wired here, and F2 is not merged.coord.backend=local with mq.backend=nats(backends.go, configuration.mdx, the CHANGELOG, AGENTS.md config bullet).api-only /ingest-only processes. C2 owns the trueroles=api+roles=ingesttest.Evidence
make ci(queued,GOTOOLCHAIN=go1.26.6) was green on 83ef8d0 and again on 2bfc9ee. The final run's coverage: integration 56.1%, e2e 60.8%, Go total 94.3%.TestNATSBackend_EndToEndpasses in about 12s.make build-docspasses, with links validated.pre-push-reviewer(opus):mq.nats.urlswas accepted. Now any@in a URL is refused, with test cases and docs.max_bytes_gbWARN. It stays a WARN, because per the design (epic(distributed): shared backends and standalone workers for multi-node deployments #613), it goes inWarnings()and the key is required per tenant. The reviewer accepted this, and the code comment says why.docs-reviewer(opus):topology_ok 0.<prefix>.ingest.>are trusted as WaveHouse itself.Left to later PRs
Sharded, partition claiming), D6 (per-tenant budgets), B2 (NATS KV leases + rule 4), C2 (multi-process test with a shared cache), E (shared cache).🤖 Generated with Claude Code
https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL