test(mq): one conformance suite for every Broker - #623
Conversation
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
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
…conformance # Conflicts: # docs/src/content/docs/api.md # docs/src/content/docs/architecture.md # internal/api/ingest.go # internal/mq/mq.go
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
|
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 |
# Conflicts: # AGENTS.md # CHANGELOG.md # docs/src/content/docs/api.md # docs/src/content/docs/architecture.md # internal/api/ingest.go # internal/app/app_test.go # internal/mq/embedded.go # internal/mq/embedded_test.go # internal/mq/mq.go
|
📚 Docs preview is live → https://212d460e-wavehouse-docs.wave-rf.workers.dev
|
Code Coverage OverviewLanguages: Go GoThe overall line coverage in commit ff67ab5 in the Show a line coverage summary of the most impacted files.
Updated |
|
Heads-up from #655: it changes how dead letters are counted, and two cases in #655 counts every scope of a table under the table itself ( What has to flip here once both are in:
Please don't restore the |
) Part of #613. One shared escaping for the composite keys WaveHouse builds, in the new `internal/keyenc` package. NATS subjects use it now, the cache's namespace tokens use it through `query.SafeEncodeToken`, and the dedupe keys adopt it in #625. The house rule it sets: a package that builds a composite key takes raw names and builds the key with `keyenc.Join`/`AppendJoin`, so no field can reach a key unescaped. ## What changes - **`internal/keyenc`** (new): - `Escape`/`AppendEscape` keep `[A-Za-z0-9_-]` and write every other byte as `%XX` in uppercase hex, so `.` `*` `>`, whitespace, `%`, `/`, `:`, `|`, `#`, `{` `}`, NUL and every non-ASCII byte are escaped. The kept bytes are exactly the tenant-id grammar, so a tenant id is its own escaped form. A `%` in a name is itself escaped, so names that look escaped (`b%2Dc`) never share a key with the name they resemble (`b-c`). - `Unescape` is `url.PathUnescape`: `%XX` in either case decodes, and any other byte reads as itself. - `Join`/`AppendJoin` escape each field and put a separator between them; `Split` reverses them. They panic on no fields (a key of no fields could not be told from one empty field) and on a separator the escaping could write, `%`, or a byte outside ASCII. - **NATS subjects** (`internal/mq`): the private `encodeToken`/`decodeToken` are gone. `Topic.key()` writes the tenant verbatim, then `AppendJoin`s the table and scope; `parseTopicKey` reads the tenant token verbatim (as `keyTenant` does) and `Split`s the rest. - **Dead-letter counts** (`internal/mq/deadletter.go`): `GET /v1/ops/dlq/stats` counts every scope of a table under the table itself, and `?table=` keeps all of its scopes. A scoped message used to count under `table.scope`, a name a dotted table could share. Scope is always empty today, so the response is unchanged. - **Cache namespace tokens**: `query.SafeEncodeToken` is a one-line delegate to `keyenc.Escape`; #614 moves the escaping into the cache itself and deletes it. ## What is deliberately not byte-identical `-` is kept rather than escaped as `%2D`, so dashed table names, scopes and (in #625) ids read as themselves. Every other byte escapes exactly as v0.1.0's encoder did (`TestEscape_MatchesV010ButDash`, all 256 byte values, and `FuzzEscapeRoundTrip` against a verbatim copy of it). Why this is safe to change now: no released queue survives into this build. v0.1.0 queued under the shared `WAVEHOUSE`/`WAVEHOUSE_DLQ` streams with subjects that carry no tenant, and #612 deletes both at boot. What remains is a queue an unreleased build since #612 wrote. It still reads, because decoding is unchanged: `%2D` decodes to `-`, and a dead-letter count merges both forms. One path notices: a `/v1/stream` client resuming across such an upgrade (`Last-Event-ID` or `since`) on a table whose name holds `-` misses that table's events queued before it, since the replay filters on the table's exact subject. The in-process cache starts empty on restart, so its keys changing costs nothing. ## Work that follows in other PRs Each of these owes a change once it takes this branch (a note is on each PR): - **#623** (broker conformance suite): `deadLetterKeepsTheTopicAndDoesNotAck` expects a scoped topic under `{"t.s": 1}` and `deadLetterCounts` counts `t1` and `t1.s` apart; both must expect every scope under its table (`{"t": 1}`, `t1` summed). The fold is not the `Broker` contract any more — restoring it would undo this PR. - **#625** (dedupe keys): builds its keys with `AppendJoin` and pins `evt-123` rather than `evt%2D123`. - **#614** (cache): `cache.Namespace` carries raw names and the cache escapes its own keys with `Join`; `query.SafeEncodeToken` is deleted. - **#626** (Redis cache): builds its token keys with `AppendJoin` and hashes escaped fields. - **#636** (external NATS broker): counts dead letters through the same `deadLetterTables`, and its copy of the conformance suite flips as #623's does. ## Tests - `TestSubject_Golden`: full subjects on the ingest and DLQ prefixes for dotted, wildcard, whitespace, `-`, `%`, `/`, brace, NUL, `\xff`, 2- and 3-byte UTF-8, empty-table and scopeless topics. - `TestParseTopicKey_LenientTokens`: a lowercase escape, a byte left unescaped (`~`), and an earlier build's `%2D` form read as the same topic. - `TestParseTopicKey_ForeignTailKeepsItself`: tails this package could not have written fall back to a topic of no tenant; an escaped tenant token (`a%2Db`) is refused, since the tenant is read verbatim. - `TestEscape_KeepsExactlyTheTenantGrammar`, `TestEscape_LookalikesStayDistinct`, `TestSeparators` (every byte value as a separator: refused, or round-trips), `TestJoin_RefusesZeroFields`. - `TestDeadLetterTables`: a dotted table and a table + scope pair count apart, scopes sum under their table, the filter keeps every scope, and `%2D` and `-` subjects merge. ## Checks - `make ci` passes locally at 79ab364. - Pre-push reviewers: `pre-push-reviewer` and `docs-reviewer` ship_it at 79ab364, after two iterate rounds with no code defects: the changelog tied the `%2D` compatibility to v0.1.0 (only queues since #612 are affected), two package listings missed `keyenc`, a stale cache line, the follow-up list missed #623, and the lenient-token test no longer exercised an unescaped byte. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01B1tJWUp6oaDoH1usLwGtLF --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Takes #655's per-table dead-letter counts: the conformance cases now expect every scope of a table counted under the table itself, a table filter keeping all of its scopes, and a dotted table name counted apart from a table + scope pair. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The ops API encodes Tables as-is and documents {"tables":{}}; the
conformance suite accepted a nil map, so a backend could pass it while
answering "tables": null.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Conflicts, both sides kept: .testcoverage.yml (main's coordtest exclude, then mqtest's), CHANGELOG.md (this PR's entry above main's new ones), and architecture.md (main's new ch_errors.go bullet, and both sides' edits to the ingest.go and mq.go bullets, merged word by word). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01E9RQYWAojiLjxBb2LvxsaW
main now carries #612 and #618 as squashes, plus #632, #622, #627, #615, #619, #647, #616, #655 and #623. The merge was resolved against the pre-squash #618 head (f129d57) as its base, so main's version wins for everything this stack does not own and only the cache stack's changes (#614, #621, #626 as merged here, and this PR) are re-applied on top. Warnings keeps main's api-role gate for the cache.redis warnings too: a split's Deployments differ only in roles, so the API's cover the others'. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017aS7rLrH1RKkUMem7X4ckd
Absorb feat/dedupe-windowed-ingest's newer history, including its own merge of feat/dedupe-reserve (step 1 of this pass) and origin/main (#615 leases, #618 backend selection, #622 process roles, #627 ClickHouse error classing, and the rest through #623). Resolved six conflicts, combining both sides' facts rather than picking one: - internal/settings/settings.go: kept MinDedupeRetention (this PR's 2-minute floor, tied to the embedded queue's duplicate window) next to windowed-ingest's reworded DLQConfig comment ("a row ClickHouse still rejects", reflecting the outage-retry split from #613). - AGENTS.md: kept this PR's dedupe/ package-inventory line (the sweep detail) alongside windowed-ingest's updated config/ and new coord/ lines pulled in from main. - CHANGELOG.md, docs/architecture.md, docs/durability.md, docs/settings-directory.mdx: superseded this branch's now-stale copies (old `evt%2D123` key spelling from before refactor/keyenc kept '-'; the metric name `wavehouse_dedupe_commit_failed_total` before it gained the `ingest_` prefix; the simpler "fits inside" duplicate-window wording before the `2×lease+1s` invariant was pinned) with windowed-ingest's current, code-matching text, then folded this PR's retention-specific additions back in: the Upgrade note now says the pre-#222 keys are "deleted by the retention sweep" instead of "nothing removing them yet", and durability.md's duplicate-window paragraph keeps its closing sentence tying `dedupe.retention`'s 2-minute floor to that same window. internal/api/ingest.go, ingest_test.go, ingest_window_test.go, app/wire.go, app/app_test.go, dedupe/embedded_test.go and settings/settings.go (the rest of it) auto-merged with no textual conflict; verified by reading the result rather than trusting that: every Commit path windowed-ingest added (commitClaims before the failing record in publishFailed, and after a clean window in ingestWindow) already passes commitClaims the full pendingRecord slice, and commitClaims groups by each record's resolved retention and issues one Commit per distinct value — so both PRs' Commit-path changes compose correctly. The sweep's commitMu (embedded.go/sweep.go) and Managed.Apply's fast path (managed.go) touch disjoint locks and did not need reconciling. go build, go vet -tags integration (whole repo), and go test -race across internal/dedupe, internal/settings, internal/api and internal/mq all pass (re-run with -count=1 after one flaky timing assertion in TestEmbedded_SweepChunkOverTombstonesDoesNotHoldCommits — a 100ms budget the sweep raced past once under parallel-package load — passed clean on every subsequent run, including three solo runs and a full fresh run; unrelated to this merge). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
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 formq.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:Subscribe, andSubscribeseeing every tenant;NakandAckWaitredelivery;DeadLetterkeeps the topic and does not ack;Tableswhen nothing is parked;failedreport, and none afterstop;MaxBytes,Stats,ErrQueueFull, andPurgeAckedsemantics.Each case runs as a parallel subtest on a fresh broker.
mqtest.Capsflags the four places where the external backend legitimately differs:PerTenantBudgetPurgesAckedUnbudgetedNotFound(DLQ counts of a tenant never given a budget)ConfiguresDurables(whetherCreateConsumerappliesAckWaitor only finds an operator-made durable)The embedded broker passes the suite. Its run is
internal/mq/mqtest/embedded_test.go.mq.gocontract wording is updated as the design specifies:ErrQueueFullis a byte limit, and the per-tenant no-queue case applies only to an implementation that opens queues per tenant.DeadLetterCountsmay return zero counts in place ofErrNoDeadLetterQueue.CreateConsumermay find rather than create.PurgeAckedmay remove nothing.Subscribeguarantees delivery only for events published after it returns.mq.ErrUnavailable(new) is mapped by the ingest handler to503+Retry-After: 5. Before, it would have been the500"publish failed". No backend returns it yet; D3's will.api.mdsays so.Two embedded bugs found by the suite are fixed:
failedmore 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.ReplaySinceread 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
internal/mq/mqtest/embedded_test.go, notinternal/mq/embedded_conformance_test.go.internal/mq's unit binary already takes about 10s of its 15s-racebudget 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 anexport_test.gohook into mq's internals.Harness.DeleteIngestDurablebecameHarness.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 ininternal/mq's own tests.Caps.NeverParkedNotFoundbecameUnbudgetedNotFound. Embedded returns zero counts for a budgeted tenant that has parked nothing. Only a tenant with no budget getsErrNoDeadLetterQueue.ConfiguresDurables. The AckWait-redelivery case cannot pass against an operator-made durable with a 60sack_wait.mqtest.Durable = "buffer-consumer". It is the worker's name, so a backend that maps durable names has one to find.Runsets the global W3C propagator for its duration. The trace case needs it, soRunmust not be called from a parallel test.Left to later PRs
ExternalNATSplusexternal_conformance_test.go, which runsmqtest.Runwith its ownCaps.Shardedcap and the shard-subset case.ConsumerConfig.Shardsdoes not exist yet, so neither is here.Evidence
make ci: green on 6ebfc6f, after the merge ofmainwith refactor(keyenc): one key escaping, keep '-', DLQ counts per table #655 (all coverage gates passed; Go total 94.4%).-race -count=20 -cpu 1,4(before the refactor(keyenc): one key escaping, keep '-', DLQ counts per table #655 merge).pre-push-revieweranddocs-reviewerboth returnedship_itbefore the refactor(keyenc): one key escaping, keep '-', DLQ counts per table #655 merge, after 6 and 4 rounds. After it,pre-push-reviewerreturnedship_iton 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.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL