Skip to content

feat(mq): a Broker over an external NATS cluster - #636

Merged
EricAndrechek merged 90 commits into
feat/mq-nats-topologyfrom
feat/mq-external-broker
Sep 29, 2026
Merged

EricAndrechek merged 90 commits into
feat/mq-nats-topologyfrom
feat/mq-external-broker

Conversation

@EricAndrechek

@EricAndrechek EricAndrechek commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

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 merges feat/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 commit aa9a6553.

What

mq.ExternalNATS (internal/mq/external.go) is an mq.Broker over 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 per Subscribe (the hub) and one per replay.

Nothing selects it yet. The config (mq.backend: nats, mq.nats.*), the wireMQ case, the integration test through app, and the user-facing docs are D4, on G1. NewNATS can be constructed and is tested. Only internal/mq imports NATS, as before.

  • Boot: NewNATS connects and then runs feat(mq): external nats broker with sharded, pinned ingest workers #624's awaitNATSTopology until TopologyWait runs 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 *TopologyError that lists every finding. A cluster never reached is ErrUnavailable (with passwords in the URLs redacted). Partition and dead-letter streams are then found by subject.
  • Auth: user plus password file, nkey seed file, creds file, TLS, and mutual TLS. Secrets are file paths only. Conflicting auth, or half a key pair, is refused before any dial.
  • Publish: goes to wh.ingest.<fnv32a(tenant)%N>.<tenant>.<table>[.<scope>] with WithExpectStream(partition), and each attempt is bounded by PublishTimeout.
    • A lost answer is retried at most twice with the same Nats-Msg-Id, so the partition stores it once.
    • maximum bytes exceeded and maximum messages per subject exceeded → ErrQueueFull.
    • No answer, a lost or closed connection, or a partition stream the operator deleted → ErrUnavailable. A deleted stream is also logged at Error with its name and drops wavehouse_mq_topology_ok at once.
  • Worker: CreateConsumer maps buffer-consumer (or the operator's own name) to wh-ingest. It finds that durable on every partition, checks ack_wait ≥ cfg.AckWait and max_ack_pending > 0, and never creates it. Any other name is ErrConsumerNotFound. Consume runs 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 on failed exactly once, and never after stop (mq: propagate terminal JetStream consumer errors to the ingest worker #587).
  • Hub: Subscribe uses a per-pod ordered, ack-less consumer on the history stream with DeliverNew. consumerName is ignored. A handler error is logged, because there is no redelivery to ask for.
  • Replay: ReplaySince uses an ack-less consumer on the history stream, filtered to the topic, starting at since. 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.
  • DLQ: DeadLetter publishes to wh.dlq.<topic key> with a fresh id. DeadLetterCounts is 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's max_age. A consumer name that does not map is ErrConsumerNotFound.
  • SetMaxBytes: records the budget and logs once that budgets are not enforced per tenant. Enforcing them is D6.
  • Close: stops every consumer (so nothing reports failed) and the watch loop, then drains the connection.
  • Re-check and gauges: the topology is verified again every 5 minutes and the history's sources are read every 30s. The gauges are wavehouse_mq_connected, wavehouse_mq_topology_ok, and per source wavehouse_mq_history_source_lag and wavehouse_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

Tests (all integration-tagged; ~6s together)

  • Conformance: mqtest.Run against ExternalNATS, connected as the restricted wavehouse user from deployments/nats/values.yaml, with each case on its own server and the shipped N=4 manifests. Caps: PerTenantBudget, PurgesAcked, UnbudgetedNotFound and ConfiguresDurables are all false. EndDelivery deletes wh-ingest on every partition as the operator. Fill shrinks the tenant's partition to 4 KiB. 18/18 applicable cases pass. AckWaitRedelivers is 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=8 under heavy load.
  • Reconnect: the server is stopped mid-consume. A publish while it is down is ErrUnavailable and not ErrQueueFull. After a restart, publishing resumes, the consumer delivers again, and failed stays quiet.
  • Idempotent retry: the answers to the first two attempts are lost. All three attempts carry the same Nats-Msg-Id and the partition holds one message. A new publish gets a new id. Losing every answer is ErrUnavailable.
  • A topic at max_msgs_per_subject is ErrQueueFull, and only that topic.
  • A deleted partition is ErrUnavailable naming the stream, and the gauge drops. The stream is not recreated.
  • Boot: a missing DLQ is ErrTopology naming it. An unreachable server is ErrUnavailable within the wait. Conflicting options are refused.
  • Auth: nkey seed, creds file (operator mode, in-process), and mutual TLS with a CA generated in the test. Missing credentials are ErrUnavailable.
  • Re-check: a DLQ deleted after boot drops topology_ok, and recreating it restores it. After a server restart the sources read as silent and topology_ok stays 1.
  • Gauges read through an OTel manual reader. PurgeAcked warns once per tenant. CreateConsumer checks the durable and never creates one.
  • Replay: a slow client (25ms/event) over more than one batch gets the whole replay (this fails on the old 5s threshold with "no responders"). Events published during a replay are not sent by it.
  • Duplicate window: a window of exactly 2 × PublishTimeout is refused.
  • Permissions: the wavehouse user is refused stream create, update, purge and delete, a durable on a partition, and deleting wh-ingest.
  • Measurements (WH_MQ_MEASURE=1; in-process single-node server, file storage, under heavy load; inferred to be optimistic for an R3 cluster):
    • publish-to-hub latency over 2000 events: p50 62µs, p99 180µs, max 525µs. That is far under design risk 3's 50ms threshold, so the hub stays on the history.
    • one tenant's publish throughput: 16,000 × 256 B from 32 callers in 0.20s, ≈80k/s (design risk 4).

Evidence

Deliberately left to later PRs

  • D4: mq.backend/mq.nats.* config and validation, the wireMQ case, tests/integration/mq_nats_test.go through app, and the docs pages (configuration, deployment "External NATS (Kubernetes)", architecture, ingest-pipeline, durability, api).
  • D5: Sharded / ConsumerConfig.Shards, so a worker consumes only the partitions it claims.
  • D6: per-tenant budgets via the backlog KV. Until then SetMaxBytes only records.

Follow-up: dead letters counted per table (7847c14, 50a1170)

DeadLetterCounts on both ExternalNATS and EmbeddedNATS folded a scoped topic into the name table + "." + scope, so a table literally named a.b could share one count with table a scoped b, and the ?table= filter matched only a table's unscoped subject. Both brokers now key their per-subject counts through one deadLetterTables (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.go and its test are byte-identical to the same files on #655, so they merge without conflict (measured with git merge-tree; the branches as a whole still conflict, since this one predates main's #612). Total is unchanged on each broker (external sums every subject of the tenant; embedded uses the stream's message count). Scope is always empty today, so GET /v1/ops/dlq/stats does not change visibly; the mqtest conformance 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 ci green locally at 50a1170 with JOBS=1. At higher parallelism internal/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 reviewers: pre-push-reviewer and docs-reviewer ship_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

EricAndrechek and others added 23 commits September 24, 2026 23:27
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
@coderabbitai

coderabbitai Bot commented Sep 25, 2026 •

Copy link
Copy Markdown

Important

Review skipped

Auto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the .coderabbit.yaml file in this repository. To trigger a single review, invoke the @coderabbitai review command.

⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Advanced

Run ID: f4d95978-4299-4445-a951-e0b2dfce5ba3

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added documentation Improvements or additions to documentation dependencies Pull requests that update a dependency file go Pull requests that update go code area/api HTTP handlers, routing, middleware labels Sep 25, 2026
EricAndrechek and others added 18 commits September 26, 2026 17:19
…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
Collapses the external-NATS stack (#636, #639, #644, #646, #654) into
#624. No conflicts.

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
@EricAndrechek
EricAndrechek merged commit 7fa5e46 into feat/mq-nats-topology Sep 29, 2026
@EricAndrechek
EricAndrechek deleted the feat/mq-external-broker branch September 29, 2026 18:21
@github-actions github-actions Bot added github_actions Pull requests that update GitHub Actions code area/observability Metrics, logs, traces, health, profiling area/ingest Ingest pipeline (Bento, batching, DLQ) area/query Structured query AST, SQL builder area/cache Local / shared / tiered caching area/dedupe Deduplication (Pebble, ScyllaDB) area/pipes Named query pipes area/sdk TypeScript SDK (clients/ts/) area/app Process wiring (internal/app): component build, run, release area/coord Leases and leader election (internal/coord) labels Sep 29, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area/api HTTP handlers, routing, middleware area/app Process wiring (internal/app): component build, run, release area/cache Local / shared / tiered caching area/coord Leases and leader election (internal/coord) area/dedupe Deduplication (Pebble, ScyllaDB) area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/ingest Ingest pipeline (Bento, batching, DLQ) area/observability Metrics, logs, traces, health, profiling area/pipes Named query pipes area/query Structured query AST, SQL builder area/sdk TypeScript SDK (clients/ts/) dependencies Pull requests that update a dependency file documentation Improvements or additions to documentation github_actions Pull requests that update GitHub Actions code go Pull requests that update go code

Projects

Status: Done

Development

Successfully merging this pull request may close these issues.

2 participants