Skip to content

feat(mq): give every tenant a queue of its own - #612

Merged
taitelee merged 15 commits into
mainfrom
mq-tenant-streams
Sep 25, 2026
Merged

taitelee merged 15 commits into
mainfrom
mq-tenant-streams

Conversation

@taitelee

@taitelee taitelee commented Sep 24, 2026 •

Copy link
Copy Markdown
Member

Summary

Story 5b of the multi-tenant epic: every tenant gets a message queue of its own inside the embedded NATS server, and the internal/mq interfaces speak per tenant, never per stream, so nothing outside the package assumes that layout (an external implementation, where the Kubernetes operator owns stream config, can keep one shared stream). A settings directory that holds the four files sees no change beyond the stream names: tenant 0's own window and budget, as before.

  • A stream pair per tenant. INGEST_<tenant> (ingest.<tenant>.>, DiscardNew) at the tenant's own mq.max_bytes_gb, and DLQ_<tenant> (dlq.<tenant>.>, DiscardOld) at a tenth of it. The prefixes differ in their first letter, so no tenant id makes one kind's name the other's (pinned by TestStreamNames_NeverCollide). Subjects are unchanged.
  • Opened with the first budget, kept on removal. SetMaxBytes(ctx, tenant, bytes) opens a tenant's pair the first time (dead-letter stream first, so no row is queued that could not be parked) and resizes it after; the wiring hands every served tenant's budget at boot and after every reload, replacing the tenant-0-only onDefaultAdopt hook and its boot warning — and with it the tracking of tenant 0's last adopted store (defaultSetting), whose one reader left, a flat directory's ops gate, now reads tenant 0 through the registry (defaultPolicy). A publish or park that finds a stream missing reopens it at the budget last asked for, and so does a publish to a queue the broker has not recorded open: an open that timed out can leave a stream JetStream creates after all, which no consumer would hold. The publishes and parks that find a queue not open share one attempt, and after one fails the tenant's publishes and parks are refused at once for five seconds — so a queue that cannot open never holds the broker's lock for every retrying client, which every other tenant's reload would wait on — while a reload retries regardless; with none asked yet (the instant between a reload adopting a new tenant and its budget arriving) the publish is refused as ErrQueueFull, a 503. A removed or rejected tenant keeps its pair at the last budget, and boot takes stock of the pairs on disk, so its queued rows are still delivered and parked on its own dead-letter queue.
  • Isolation. The ingest worker's and the hub bridge's durables are held on every tenant's stream, those opened later included; each tenant has its own ack floor, its own MaxAckPending, and a delivery goroutine of its own, so one tenant's backlog or stuck handler holds back no other. A publish that reopens a missing queue does it detached from the request's cancellation, and the consumers join on a budget of their own, so a client that goes away mid-open cannot leave a queue no consumer holds — which would fail the worker and stop the process. The worker's prefetch, and the hub bridge's fetch-ahead (the client default of 500), are each shared across the tenants' streams (at least one each) rather than multiplied by them. The durables are looked up before anything is written, so a boot over many queues writes nothing it need not (measured: 0.6 s at 1,000 tenants, against 15 s rewriting every stream and durable).
  • Sweeper. PurgeAcked takes each tenant's cutoff; the sweeper hands it every tenant's own stream.gap_window_minutes (gapWindows replaces longestGapWindow). A rejected tenant keeps the window its folder last had — a rejection is the common reload failure, and its clients resume from Last-Event-ID once the folder is fixed (settings.Registry.Known now yields each tenant's last adopted store for this), and one rejected since boot, whose window this process never read, keeps all of its history until its folder validates — while a removed tenant's stream keeps no acknowledged history. Per-tenant purge lines are Debug, with one Info summary per sweep.
  • GET /v1/ops/dlq/stats?tenant=, parsed strictly with opsTenant. Absent is tenant 0, the ops-read convention, and for a flat directory exactly the previous answer; no parameter no longer sums every tenant. The tenant is looked up in the MQ, not the settings, so a rejected or removed tenant's parked rows are read by name; a tenant with no dead-letter queue is a 404. The SDK's wh.dlq.list() and .table() take a tenant option, as the pipe reads did in feat(settings): nested settings directory with a per-tenant registry #598.
  • Dead-letter shrink guard (mq(dlq): shrinking mq.max_bytes_gb silently deletes the oldest dead letters — investigate how the reload should treat a non-empty DLQ #532, interim). A reload that shrinks a budget never caps a dead-letter stream below the bytes it holds: it keeps what it has, and a warning names the tenant and both figures. mq(dlq): shrinking mq.max_bytes_gb silently deletes the oldest dead letters — investigate how the reload should treat a non-empty DLQ #532's ClickHouse-backed dead-letter table stays there.
  • Disk accounting (decided: option A). JetStream counts every stream's cap as reserved disk and refuses a stream once they pass 75% of the free disk at boot; per tenant that refused the 569th tenant at the 1 GB minimum on an 834 GB disk, and would refuse the 12th at the seed's 50 GB. The embedded server's limit is set out of reach, so a budget is a cap and never a reservation; what the budgets add up to against the disk is mq: validate WH_MQ_MAX_BYTES_GB upper bound to prevent disk over-reservation #138's. The one flat difference beyond the stream names: a budget above three quarters of the free disk now boots where it was refused.
  • Old data directories. WAVEHOUSE and WAVEHOUSE_DLQ claim ingest.> and dlq.>, which JetStream will not let overlap a new stream, so boot deletes them, logging what they held. What that made dead is gone: parseTopicKey's one-token subjects and deployment.md's "Upgrading across the tenant subject token". The v2-envelope upgrade runbook now says the old queue is deleted rather than carried over, so draining and replaying the parked rows has to happen before the upgrade, and the worker's comments and dead-letter hint no longer assume an undrained pre-v2 row can reach it.
  • Sweep. Every "story 5b" and "until 5b", and every sentence this makes false, is updated: the mq interface comments, the sweeper, wire.go and app.go, app_test.go, the hub subscriber's one-goroutine comment, the settings comments, and deployment.md, settings-directory.mdx, api.md, architecture.md, ingest-pipeline.md (diagrams included), durability.md, configuration.mdx, why-wavehouse.md, the SDK admin and reference pages, AGENTS.md and the CHANGELOG.

Cost check

Measured in a scratch program against the same nats-server 2.14.6 and server options (fsync on every write), with both durables consuming on each tenant's stream, at 10 / 100 / 1,000 tenants, against the one shared pair:

  • Resident memory, idle: 25 / 53 / 280–315 MB, against 21–35 MB shared — about 0.3 MB per tenant. Goroutines: 323 / 2,123 / 20,123, against 43.
  • Under load, a tenant written in the last ~10 s holds a block buffer: at 1,000 tenants one write to each took the heap from 172 to 439 MB, back to 172 MB within 15 s. One open file per stream written, closed within about two minutes of quiet.
  • Boot with data: 12 / 59 / 611 ms. A tenant's first creation costs about 20 ms (the pair, two durables, delivery), so a first boot of 1,000 new tenants pays about 20 s once.
  • Publish throughput: 32 concurrent publishers 1,565 / 7,155 / 2,390 msg/s against ~1,170 shared; one publisher round-robin across tenants 999 / 560 / 503 against ~1,100.
  • One sweep: 31 ms / 296 ms / 2.75 s, once a minute.
  • Not measured, but by construction: while ClickHouse stalls, the worker can hold up to maxAckPending (10,000) unacknowledged rows per tenant with traffic — about 10 million at 1,000 tenants — where the shared stream held 10,000 for the whole process (see Follow-ups).

An upstream quirk

When a stream's store fails to open, nats-server 2.14.6 releases a disk reservation it never made, so its reserved count goes negative. Harmless under the default limit, but with the limit at math.MaxInt64 the subtraction overflows and every later stream is refused until restart. The limit is half the int64 range instead, and TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen pins it: one tenant's failed open, then the next tenant's queue opens.

Test plan

  • make ci passes locally
  • Pre-push reviewers (round 1: both iterate — fixed: the reopen detach, the v2 upgrade runbook, the remaining sweeper and create stream lines, the 503 rows, the disk-full sizing rule, the per-tenant consume callbacks in the goroutine topology; round 2: both iterate — fixed: a test for the boot rule of a queue that cannot open, the boot pass honoring New's context, the last pre-v2 mentions, the resume hole after a rejected folder, the runbook's replay step the old build cannot carry out, the boot warning's count; round 3: code iterate on one SHOULD — boot now re-applies a tenant's budget to a pair it finds split, or missing its dead-letter stream, as every boot rewrote both streams before; docs iterate — the boot/reload contexts in architecture.md, the runbook's opening, the per-tenant consumer line, and a full disk needing a restart; round 4: code iterate — a failed resize's undo now restores the ingest stream's actual cap (after a boot that found a pair split it would have restored 0, which JetStream reads as no cap), and a test pins the report of a queue a consumer cannot join; docs iterate — both open-timeout errors and when a later boot writes, three garden-path sentences; round 5: code iterate — a rejected tenant keeps its replay history at its folder's last window instead of losing it at the next sweep, and all of it while rejected since boot; docs iterate — disk sizing counts every tenant ever served, since removed tenants' queues are kept; round 6: code iterate — the hub bridge's fetch-ahead is shared across the tenants like the worker's prefetch, instead of 500 per tenant; docs iterate — a stale "stop the sweeper" sentence from the shared stream, and a dead-letter stream is opened, not guaranteed to exist, when its tenant is first served; round 7: code iterate — a publish goes by the broker's record of an open queue rather than a stream answering, so a stream an open gave up on but JetStream created anyway is joined before rows land in it, and the queue-full and publish-failure log lines name the tenant; docs iterate — the disk-sizing guidance gets its own paragraph, and the leftover tenant-0 tracking clause goes, with the tracking itself; round 8: code ship it; docs iterate — what a consumer that cannot join a queue opened at runtime does (the worker exits the process, the hub logs and the tenant's streams get no live rows until a restart), and a failed purge lookup ends that tenant's purge, not the sweep; round 9: code iterate — a tenant whose queue cannot open no longer holds the broker's lock for every retrying publish (publishes share one attempt, and after one fails are refused at once for five seconds; a reload retries regardless), and the worker's per-tenant memory ceiling is recorded here, for epic(settings): multi-tenant settings directory and tenant-scoped runtime #583's deferred worker rework; docs iterate — two wording fixes; round 10: code iterate — tests for one tenant's failed purge stopping no other and for a durable reused across a restart; docs iterate — the SDK replay note names the upgrade that deletes the old queue, and each tenant's own gap window; round 11: code ship it; docs iterate — disk sizing counts a dead-letter stream the shrink guard kept above a tenth, and the stats route says where the tenant is looked up; round 12: docs ship it — the round-12 commit is docs prose only, so the code reviewer's round-11 ship it stands for it, recorded as a logged skip)
  • CodeRabbit round 1 (changes requested, two findings, fixed): the sweeper no longer logs other tenants' purge failures at WARN beside one tenant's missing consumer, and a park that finds the dead-letter stream missing is paced like a publish
  • CodeRabbit round 2: approved
  • Flake fix after Eric's heads-up from the stacked epic(distributed): shared backends and standalone workers for multi-node deployments #613 builds (f5d8f484): TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen re-created its obstacle while nats-server's post-failure goroutine was removing the emptied streams and account directories (EINVAL on APFS, ENOENT on Linux; 3 of 25 runs here); the test and the one app test with the pattern now keep another tenant's streams in the directory. CodeRabbit round 3: approved. eb9f7409 closes the same goroutine's much narrower window in TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen with a non-stream directory under streams/ (reserves nothing, so the test still pins the reservation overflow); Eric's second measurement (9 of 10 on feat/mq-nats-wiring) was on a branch forked before f5d8f484 — verified by applying the fix there (5 of 20 → 0 of 20). CodeRabbit round 4: approved
  • mq: a reopen outlives the caller's cancellation and the consumer still joins (fails without the detach); a tenant's pair opens with its first budget, names and caps as above, and no other tenant's; a publish reopens a missing pair at the last budget and is refused with none asked; a publish opens a queue whose open gave up on a stream JetStream created anyway, and its row reaches a running consumer (fails without the gate); after a failed attempt a publish or park is refused without waiting on the broker's lock, and a publish tries again once the window has passed (fails without the pacing); one tenant at its budget is refused while another publishes; a tenant at MaxAckPending and one with a stuck handler hold back no other, and a tenant's messages arrive in order; both consumer paths join a queue opened after they started; one tenant's deleted durable reports on failed, and a stop is not a failure; the prefetch share, and the hub bridge's shared fetch-ahead; per-tenant resize, rollback and cancellation; the shrink guard keeps every parked row; per-tenant purge at each cutoff, and no history for a tenant not named; one tenant's failed purge — its durable or its stream gone — stops no later tenant's, and an ended sweep touches none (fails with either continue turned into a break); a durable on disk is reused across a restart, or updated in place when its settings differ, and delivery resumes past what it acknowledged; per-tenant dead-letter counts, ErrNoDeadLetterQueue, and a missing dead-letter stream reopened; a queue that cannot open costs that tenant alone; boot deletes the old shared pair; boot over existing pairs reads their budgets back and delivers a queued row of a tenant given no budget, and re-applies the budget to a split pair or one missing its dead-letter stream while leaving a guarded one alone; a failed resize's undo restores the ingest stream's own cap; a consumer that cannot join a queue opened later reports it on failed
  • ingest: the sweeper hands each tenant's own cutoff, re-read every sweep, and logs a failed sweep at WARN only when every tenant's failure is a missing buffer consumer
  • settings: Known yields a rejected tenant's last adopted store, and none for a folder that has not validated since boot
  • api: dlq/stats without the parameter reads tenant 0, names a tenant's queue alone, 404s without a queue, and refuses a malformed, empty, repeated or misparsed tenant
  • app: a queue that cannot open refuses a flat boot and costs a nested directory that tenant alone, whose next publish opens it; a cancelled boot opens no queue; each tenant's budget follows its own folder, and a rejected or removed folder keeps its queue at the last budget; gapWindows names each served tenant, a rejected one at its folder's last window (unbounded when rejected since boot), and no removed one
  • SDK: tenant on wh.dlq.list() and .table()
  • Manual: a nested directory with two tenants; fill one tenant's queue and see only its ingest answer 503; remove it and reload, then read its parked rows with GET /v1/ops/dlq/stats?tenant=

Follow-ups

  • The nats-server reservation quirk above is worth reporting upstream.
  • A consumer that cannot join a queue opened at runtime fails the ingest worker, so the process restarts. Now that a publish goes only into a queue recorded open, apply could leave such a queue unrecorded instead, so the tenant answers 503 and publishes and reloads retry the join, with no restart.
  • The worker's in-memory bound is per tenant now: while ClickHouse stalls it can hold up to maxAckPending (10,000) rows for every tenant with traffic, documented at the constant and in the ingest pipeline's backpressure list. A process-wide bound would bring back one tenant holding the others back; it belongs with epic(settings): multi-tenant settings directory and tenant-scoped runtime #583's deferred worker rework.
  • nats-server logs three Info lines per stream at every boot (restore and consumer recovery), about 6,000 at 1,000 tenants, through slogNATSLogger.
  • sdk/admin.md's older DLQ example declares const { data } twice in one block, and shows a "users": 0 count the server never reports (both predate this branch).
  • api.md's ?table= row says it returns only that table's count; total stays the tenant's whole count (predates this branch).
  • settings-directory.mdx's DLQ switch still calls the unreadable-envelope exception "new in this release", which goes stale with the next one (predates this branch).

Related Issues

Part of #583 (story 5b).

@coderabbitai

coderabbitai Bot commented Sep 24, 2026 •

Copy link
Copy Markdown

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Advanced

Run ID: cbcce041-f03e-456e-9a0b-1eb65e4d6c09

📥 Commits

Reviewing files that changed from the base of the PR and between f5d8f48 and eb9f740.

📒 Files selected for processing (1)
  • internal/mq/embedded_test.go

Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review.

📜 Recent review details
⏰ Context from checks skipped due to timeout. (11)
  • GitHub Check: Docs build
  • GitHub Check: Unit tests
  • GitHub Check: Integration tests
  • GitHub Check: Coverage
  • GitHub Check: E2E tests
  • GitHub Check: Lint
  • GitHub Check: Analyze (actions)
  • GitHub Check: Analyze (go)
  • GitHub Check: Analyze (javascript-typescript)
  • GitHub Check: Analyze (go)
  • GitHub Check: Analyze (javascript-typescript)
🔇 Additional comments (1)
internal/mq/embedded_test.go (1)

441-449: LGTM!


📝 Summary

Summary by CodeRabbit

  • New Features
    • Message queues, dead-letter queues, storage budgets, and history retention are managed separately for each tenant.
    • Dead-letter statistics can be requested for a specific tenant; requests for tenants without a dead-letter queue return a not-found response.
    • The TypeScript admin client supports tenant selection when listing dead-letter statistics.
  • Bug Fixes
    • Queue-capacity errors and delivery issues are isolated to the affected tenant.
    • Removing or rejecting a tenant preserves its last adopted queue budget and applicable history while other tenants continue operating.

Walkthrough

The pull request replaces shared JetStream streams with tenant-specific ingest and dead-letter streams. It adds per-tenant queue budgets and purge windows, updates message delivery and replay, and lets DLQ statistics requests select a tenant.

Changes

Tenant-scoped queue flow

Layer / File(s) Summary
Tenant queue contracts and routing
internal/mq/mq.go, internal/mq/subject.go, internal/mq/subject_test.go, CHANGELOG.md, AGENTS.md
MQ contracts and subject helpers now identify queues and message subjects by tenant. One-token topic tails are no longer assigned to the default tenant.
Embedded tenant queue lifecycle
internal/mq/embedded.go, internal/mq/embedded_test.go, internal/ingest/worker.go, internal/ingest/worker_test.go, internal/stream/*, internal/testutil/*, internal/api/ingest.go, internal/api/router_test.go, docs/src/content/docs/{api,architecture,durability,ingest-pipeline,settings-directory,deployment}.md*, CHANGELOG.md, internal/settings/settings.go
The broker manages tenant-specific streams, budgets, consumers, queue reopening, purge cutoffs, and replay. Tests cover per-tenant queue state, recovery, delivery, and purge behavior.
Per-tenant settings and purge windows
internal/app/*, internal/ingest/sweeper*, internal/settings/registry*, docs/src/content/docs/{architecture,configuration,deployment,settings-directory,sdk/streaming}.md*, CHANGELOG.md, AGENTS.md, internal/testutil/mocks.go
App wiring reconciles budgets for served tenants. The sweeper uses a cutoff for each known tenant, and the settings registry retains adopted stores for rejected tenants.
Tenant-specific DLQ statistics
internal/api/dlq*, clients/ts/src/{dlq.ts,namespaces.test.ts,types.ts}, docs/src/content/docs/{api,deployment,sdk/admin,sdk/reference,why-wavehouse}.md*, AGENTS.md, internal/settings/settings.go
DLQ statistics requests accept a tenant query. The API returns 404 when the selected tenant has no queue, and the TypeScript methods accept tenant request options.

Sequence Diagram(s)

sequenceDiagram
  participant AppWire as app.wireMQ
  participant Broker as EmbeddedNATS
  participant Queue as Tenant JetStream streams
  participant FanIn as Tenant consumer fan-in
  participant Worker as Ingest worker
  AppWire->>Broker: Reconcile tenant queue budget
  Broker->>Queue: Open or resize tenant streams
  Worker->>FanIn: Register durable consumer
  FanIn->>Queue: Join tenant ingest stream
  Queue->>FanIn: Deliver tenant messages
  FanIn->>Worker: Pass messages to worker
Loading

Priority: ➖ Normal

Estimated code review effort: 4 (Complex) | ~60 minutes

Change: Feature

Suggested reviewers: ericandrechek

Merge Risk: ⚪ Minimal · up to eb9f7

No actionable merge-blocking risk remains from the reviewed changes.

Security Architecture Review

Security architecture risk: 🟡 Moderate · up to eb9f7

Tenant-specific limits improve isolation, but the combined queues can now consume more disk than the previous server-wide guard allowed. Upgrading also deletes the old queues, including parked rows and replay history. The administrative endpoint remains protected.

Retained concerns

  • Medium · security · inferred: Per-tenant queue caps no longer have an effective server-wide disk reservation limit. Traffic that fills allowed tenant queues can exhaust their shared volume and affect other tenants, despite each tenant staying within its own configured budget.
  • Medium · reliability · observed: Startup deletes the legacy ingest and dead-letter streams rather than migrating them. Draining can protect uninserted events, but parked rows and pre-upgrade SSE replay history still have no preservation path described for this upgrade.
Security review details

Security Blast Radius

  • inferred — Authorized ingestion is limited by its tenant's stream cap, but tenants share a physical JetStream store without an aggregate cap tied to free disk. Exhausting that store could disrupt more than the tenant producing traffic.

Security Findings and Attack Paths

  • inferred — A caller permitted to ingest for a tenant can generate queued data up to that tenant's configured limit. If combined configured limits exceed available disk, the removed server-wide guard no longer prevents aggregate consumption from affecting the shared store; this is an availability risk, not evidence of cross-tenant data access.

Trust Boundaries and Controls

  • observed — DLQ tenant selection parses a single valid tenant ID inside the existing admin-gated operations route. A failed queue open is surfaced as queue-full backpressure, which the ingest handler maps to a retryable 503 rather than acknowledging publication.

Resilience and Maintainability Implications

  • observed — Consumer-join failure is reported but does not undo an opened stream pair. The available failed-open test verifies stream-opening refusal and recovery, not a guarantee that publication waits for every consumer to attach.

Hardening Proposals

  • proposed — Consider an aggregate disk or reservation policy in addition to per-tenant caps, so a tenant's permitted queue growth cannot consume the shared volume's operational headroom.
  • proposed — Consider an explicit backup or migration gate before legacy-stream deletion, particularly where parked rows or replay continuity must survive deployment or rollback.
🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 68.47% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 111 functions across 27 files. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
Title check ✅ Passed The title clearly and concisely states the main change: each tenant receives its own message queue. It matches the pull request objectives and changes.
Description check ✅ Passed The description directly explains the per-tenant queue architecture, tenant-specific budgets and consumers, API and SDK changes, migration behavior, tests, and known manual verification gap.
  • Fix all pre-merge checks with AI
✨ Finishing Touches 💡 1
📝 Generate docstrings 💡
  • Commit to this branch
  • Create a new PR
🧪 Generate unit tests (beta)
  • Commit to this branch
  • Create a new PR
✨ Simplify code
  • Commit to this branch
  • Create a new PR

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 go Pull requests that update go code area/api HTTP handlers, routing, middleware area/ingest Ingest pipeline (Bento, batching, DLQ) area/sdk TypeScript SDK (clients/ts/) area/docs Documentation, site/, README area/app Process wiring (internal/app): component build, run, release labels Sep 24, 2026
@github-actions

github-actions Bot commented Sep 24, 2026 •

Copy link
Copy Markdown

📚 Docs preview is live → https://e94fefb2-wavehouse-docs.wave-rf.workers.dev

  • Commit — eb9f740: test(mq): keep the streams directory occupied through the first failed open too
  • Author — @taitelee
  • Committed — 2026-09-25 08:35 (UTC-04:00)
  • Deployed — 2026-09-25 09:12 EDT

@github-code-quality

github-code-quality Bot commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Code Coverage Overview

Languages: Go

Go

The overall line coverage in commit eb9f740 in the mq-tenant-streams branch remains at 93%, unchanged from commit 93d8019 in the main branch.

Show a line coverage summary of the most impacted files.
File main 93d8019 mq-tenant-streams eb9f740 +/-
internal/app/wire.go 92% 92% 0%
internal/stream/hub.go 97% 97% 0%
internal/ingest/worker.go 97% 97% 0%
internal/api/dlq.go 100% 100% 0%
internal/ingest/sweeper.go 100% 100% 0%
internal/mq/subject.go 100% 100% 0%
internal/settings/registry.go 100% 100% 0%
internal/mq/embedded.go 85% 88% +3%

Updated September 25, 2026 12:38 UTC

@taitelee

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Sep 25, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 2


ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Advanced

Run ID: eac01c7c-fbee-4177-8af5-bf2260fa31ce

📥 Commits

Reviewing files that changed from the base of the PR and between 93d8019 and 060ca9e.

📒 Files selected for processing (40)
  • AGENTS.md
  • CHANGELOG.md
  • clients/ts/src/dlq.ts
  • clients/ts/src/namespaces.test.ts
  • clients/ts/src/types.ts
  • docs/src/content/docs/api.md
  • docs/src/content/docs/architecture.md
  • docs/src/content/docs/configuration.mdx
  • docs/src/content/docs/deployment.md
  • docs/src/content/docs/durability.md
  • docs/src/content/docs/ingest-pipeline.md
  • docs/src/content/docs/sdk/admin.md
  • docs/src/content/docs/sdk/reference.md
  • docs/src/content/docs/sdk/streaming.md
  • docs/src/content/docs/settings-directory.mdx
  • docs/src/content/docs/why-wavehouse.md
  • internal/api/dlq.go
  • internal/api/dlq_test.go
  • internal/api/ingest.go
  • internal/api/router_test.go
  • internal/app/app.go
  • internal/app/app_test.go
  • internal/app/wire.go
  • internal/ingest/sweeper.go
  • internal/ingest/sweeper_test.go
  • internal/ingest/worker.go
  • internal/ingest/worker_test.go
  • internal/mq/embedded.go
  • internal/mq/embedded_test.go
  • internal/mq/mq.go
  • internal/mq/subject.go
  • internal/mq/subject_test.go
  • internal/settings/registry.go
  • internal/settings/registry_test.go
  • internal/settings/settings.go
  • internal/settings/store.go
  • internal/stream/hub.go
  • internal/stream/subscriber.go
  • internal/testutil/mocks.go
  • internal/testutil/testutil.go

Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review.

📜 Review details
⏰ Context from checks skipped due to timeout. (9)
  • GitHub Check: Integration tests
  • GitHub Check: E2E tests
  • GitHub Check: Unit tests
  • GitHub Check: Coverage
  • GitHub Check: Docs build
  • GitHub Check: Lint
  • GitHub Check: Analyze (go)
  • GitHub Check: Analyze (go)
  • GitHub Check: Analyze (javascript-typescript)
🧰 Additional context used
📓 Path-based instructions (5)
Fail-closed — preserve it when touching `internal/api`.

📄 CodeRabbit inference engine (AGENTS.md)

Files:

  • internal/api/ingest.go
  • internal/api/router_test.go
  • internal/api/dlq.go
  • internal/api/dlq_test.go
See [AGENTS.md](AGENTS.md) for project conventions, architecture notes, and AI agent instructions.

📄 CodeRabbit inference engine (CLAUDE.md)

Files:

  • AGENTS.md
See [AGENTS.md](../AGENTS.md) for project conventions, architecture notes, and AI agent instructions.

📄 CodeRabbit inference engine (.github/copilot-instructions.md)

Files:

  • AGENTS.md
In MDX, leave a blank line between a JSX tag and a code fence.

📄 CodeRabbit inference engine (AGENTS.md)

Files:

  • docs/src/content/docs/configuration.mdx
  • docs/src/content/docs/settings-directory.mdx
WH001 applies to every tracked Markdown file, with no carve-out Editors see WH001 in `.md` only.

📄 CodeRabbit inference engine (AGENTS.md)

Files:

  • docs/src/content/docs/sdk/reference.md
  • docs/src/content/docs/why-wavehouse.md
  • docs/src/content/docs/sdk/streaming.md
  • docs/src/content/docs/sdk/admin.md
  • docs/src/content/docs/durability.md
  • AGENTS.md
  • docs/src/content/docs/deployment.md
  • docs/src/content/docs/architecture.md
  • docs/src/content/docs/ingest-pipeline.md
  • docs/src/content/docs/api.md
  • CHANGELOG.md
🧠 Learnings (2)
📚 Learning: 2026-05-23T01:23:59.268Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 174
File: internal/api/ingest_test.go:111-111
Timestamp: 2026-05-23T01:23:59.268Z
Learning: In WaveHouse Go tests in internal/api/**/*_test.go, use internal/testutil.AssertJSONErrorResponse(t, w) for HTTP error-path JSON assertions. Do not use (or reintroduce) package-local assertJSONErrorResponse helpers. AssertJSONErrorResponse verifies the response Content-Type is application/json, includes the X-Content-Type-Options: nosniff header, and that the JSON body contains an "error" field.

Applied to files:

  • internal/api/dlq_test.go
📚 Learning: 2026-06-26T12:23:22.696Z
Learnt from: EricAndrechek
Repo: Wave-RF/WaveHouse PR: 346
File: internal/stream/subscriber_test.go:9-28
Timestamp: 2026-06-26T12:23:22.696Z
Learning: In this Go repository, prefer table-driven tests (e.g., `[]struct{...}` with `t.Run(...)`) only for tests that cover multiple scenarios/inputs and can be cleanly enumerated. Do not artificially rewrite a clear single-scenario sequential behavioral-flow test into a table-driven form just to fit the pattern; if there’s only one meaningful scenario, keep the test as a straightforward linear flow (as in `TestSubscriber_SendDeliversThenDropsWhenFull`).

Applied to files:

  • internal/mq/embedded_test.go
🪛 LanguageTool
docs/src/content/docs/configuration.mdx

[style] ~100-~100: Since ownership is already implied, this phrasing may be redundant.
Context: Each tenant's queue has its own disk budget, mq.max_bytes_gb, a hot-r...

(PRP_OWN)

docs/src/content/docs/durability.md

[style] ~36-~36: Consider using the typographical ellipsis character here instead.
Context: ...eload that first serves the tenant — is open dlq stream: ... context deadline exceeded, or `open in...

(ELLIPSIS)


[style] ~36-~36: Consider using the typographical ellipsis character here instead.
Context: ...eam: ... context deadline exceeded, or open ingest stream: ...` (the two share the budget); a boot tha...

(ELLIPSIS)


[style] ~96-~96: Consider using the typographical ellipsis character here instead.
Context: ...eam: ... context deadline exceeded, or open ingest stream: ...`, when a tenant's queue first opens, at...

(ELLIPSIS)


[style] ~97-~97: Consider using the typographical ellipsis character here instead.
Context: ...very ended; ingestion has stoppedwithjoin its queue: ... context deadline exceeded`, and the pro...

(ELLIPSIS)

AGENTS.md

[style] ~43-~43: Consider using the typographical ellipsis character here instead.
Context: ...equest; settings.Store in production, Static(q...) in tests) - policy/ — Hasura-st...

(ELLIPSIS)


[style] ~47-~47: Since ownership is already implied, this phrasing may be redundant.
Context: ... table — and evaluates each event under its own tenant's policy and schema registry; `P...

(PRP_OWN)

docs/src/content/docs/deployment.md

[style] ~388-~388: This word has been used in one of the immediately preceding sentences. Using a synonym could make your text more interesting to read, unless the repetition is intentional.
Context: ...t the parameter — and on SIGHUP — the whole directory is reloaded and mirrors its f...

(EN_REPEATEDWORDS_WHOLE)


[style] ~388-~388: This word has been used in one of the immediately preceding sentences. Using a synonym could make your text more interesting to read, unless the repetition is intentional.
Context: ...ved: delete its folder, then reload the whole directory. Its open streams end, its ro...

(EN_REPEATEDWORDS_WHOLE)


[style] ~388-~388: Since ownership is already implied, this phrasing may be redundant.
Context: ...queued rows are parked on the DLQ under its own subject; nothing it stored is deleted —...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: ...cides.** A request is evaluated against its own tenant's policies.json and `pipes.jso...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: ...s-directory#clickhouse) — and discovers its own tables from its own database on its own...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: ...se) — and discovers its own tables from its own database on its own `schema.refresh_int...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: .../v1/streamconnection is authorized by its own tenant'spolicies.json` and receives i...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: ...n tenant's policies.json and receives its own tenant's rows alone, the ingest worker ...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: ...e, the ingest worker inserts a row into its own tenant's ClickHouse, a failed row is pa...

(PRP_OWN)


[style] ~392-~392: Since ownership is already implied, this phrasing may be redundant.
Context: ...lickHouse, a failed row is parked under its own tenant's dlq.enabled and subject (`dl...

(PRP_OWN)


[style] ~394-~394: Since ownership is already implied, this phrasing may be redundant.
Context: ... while every other tenant's routes keep their own list. Tenant 0's own dedupe store clo...

(PRP_OWN)

docs/src/content/docs/architecture.md

[grammar] ~85-~85: Ensure spelling is correct
Context: ...ave-RF/WaveHouse/issues/319)). Gap-fill replay (mq.Replayer.ReplaySince on the conne...

(QB_NEW_EN_ORTHOGRAPHY_ERROR_IDS_1)


[style] ~92-~92: This phrase is redundant. Consider writing “last”.
Context: ...nt. The SIGHUP registration is released last of all. Handler, Registry, and MQ expose...

(LAST_OF_ALL)


[style] ~92-~92: Since ownership is already implied, this phrasing may be redundant.
Context: ...ons.Listenerlets one serve the API on its own listener instead ofserver.port`. - **...

(PRP_OWN)


[style] ~148-~148: Since ownership is already implied, this phrasing may be redundant.
Context: ...o a blocking handler is backpressure on its own tenant alone, spreads the prefetch acro...

(PRP_OWN)


[style] ~148-~148: A comma is missing here.
Context: ...ending on its own — ErrDeliveryEnded, e.g. a deleted consumer, a closed connection...

(EG_NO_COMMA)


[style] ~148-~148: Since ownership is already implied, this phrasing may be redundant.
Context: ...terer.DeadLetter` (park a message under its own topic, in its tenant's dead-letter queu...

(PRP_OWN)


[style] ~151-~151: ‘out of reach’ might be wordy. Consider a shorter alternative.
Context: ...75% of the free disk by default) is set out of reach, so a budget is a cap and never a reser...

(EN_WORDINESS_PREMIUM_OUT_OF_REACH)


[style] ~151-~151: For a more expressive style, consider rephrasing the sentence in the active voice.
Context: ...rts the budget last applied in full, so a failed resize is retried by the next reload. The consumers CreateConsumer and `Su...

(PASSIVE_VOICE_SIMPLE)

docs/src/content/docs/ingest-pipeline.md

[grammar] ~64-~64: Ensure spelling is correct
Context: ... as ClickHouse's calendar/epoch shapes, where basic read a plain DateTime column'...

(QB_NEW_EN_ORTHOGRAPHY_ERROR_IDS_1)


[grammar] ~64-~64: Ensure spelling is correct
Context: ...nds (shorter runs it rejected outright, where best_effort reads "2026" as a year)...

(QB_NEW_EN_ORTHOGRAPHY_ERROR_IDS_1)


[grammar] ~64-~64: Ensure spelling is correct
Context: ...effort "20260711"stores 2026-07-11, wherebasicstored 1970-08-23.DateTime64`...

(QB_NEW_EN_ORTHOGRAPHY_ERROR_IDS_1)


[style] ~205-~205: Consider using “who” when you are referring to people instead of objects.
Context: ...rker's own failed channel. A consumer that cannot start at all takes the same path...

(THAT_WHO)


[style] ~207-~207: To elevate your writing, try using a synonym here.
Context: ...at could delete a durable) this path is hard to reach; the likeliest way in is a ten...

(HARD_TO)


[style] ~227-~227: Consider an alternative for the overused word “exactly”.
Context: ...an fsync and therefore slow, which is exactly why acks run in the background (ackWg...

(EXACTLY_PRECISELY)

CHANGELOG.md

[typographical] ~13-~13: Consider using an em dash in dialogues and enumerations.
Context: - **A tenant removed or rejected at runti...

(DASH_RULE)


[style] ~13-~13: This sentence is over 40 words long. Consider splitting it up, as shorter sentences make the text easier to read.
Context: - A tenant removed or rejected at runtime has its open streams ended (internal/stream/{hub,subscriber,bucket}.go (+ tests), internal/api/stream.go (+ tests), internal/api/router.go, internal/settings/{registry,tree}.go (+ tests), internal/ingest/worker.go (+ tests), internal/app/wire.go (+ tests), clients/ts/src/stream/sse.ts, docs/src/content/docs/{api,deployment,architecture,ingest-pipeline}.md, docs/src/content/docs/settings-directory.mdx, AGENTS.md): story 3 of the multi-tenant epic (#583). A GET /v1/stream used to outlive i...

(TOO_LONG_SENTENCE)


[style] ~13-~13: This word has been used in one of the immediately preceding sentences. Using a synonym could make your text more interesting to read, unless the repetition is intentional.
Context: ...ld pass, and meets its DLQ switch once, whole, logged once per batch rather than twic...

(EN_REPEATEDWORDS_WHOLE)


[style] ~13-~13: Since ownership is already implied, this phrasing may be redundant.
Context: ...on, so its queued rows are parked under its own subject rather than dropped, inserted i...

(PRP_OWN)


[style] ~27-~27: Since ownership is already implied, this phrasing may be redundant.
Context: ...t-reloadable half of configuration gets its own page; configuration.mdx is boot confi...

(PRP_OWN)


[style] ~27-~27: For conciseness, consider replacing this expression with an adverb.
Context: ....tables` (resolved by the ingest worker at the moment a poison row is isolated: on → park it ...

(AT_THE_MOMENT)


[style] ~27-~27: This sentence is over 40 words long. Consider splitting it up, as shorter sentences make the text easier to read.
Context: ...out materializing settings directories. The settings directory is also the runtime authority for access control and named pipes (internal/settings/store.go, internal/policy/source.go (new), internal/pipes/pipes.go, internal/api/{policy,pipes,router}.go, internal/stream/hub.go, internal/auth/auth.go, cmd/wavehouse/main.go, Makefile, deployments/compose/settings/{policies,roles}.json, clients/ts/src/settings.ts (new); closes #229, #33, #461, #514, #460, #363; advances #48 and #214): roles.json, policies.json, and pipes.json are adopted with config.json as one snapshot and re-adopted on the same three triggers, and files are the only write path — standalone, the operator edits them on the host; on WaveHouse Cloud the control plane writes them — so there is no stored copy that can skip validation: every adoption runs the current rules (strict decode rejecting unknown and duplicate keys, the full policy validation including the claim-template grammar, pipe name/SQL/parameter-type rules, and the cross-file check that every role a grant or allowed_roles names is declared in roles.json), and a rejected edit keeps the previous good policy and pipes in effect. policies.json is one policy document ...

(TOO_LONG_SENTENCE)


[style] ~31-~31: Since ownership is already implied, this phrasing may be redundant.
Context: ...ter,LiveDemo}.astro`): the site tracked its own CTAs but nothing a reader did on the wa...

(PRP_OWN)


[style] ~31-~31: Since ownership is already implied, this phrasing may be redundant.
Context: ...docs area without each tracker carrying its own copy; it's stamped at capture time by a...

(PRP_OWN)


[style] ~31-~31: Since ownership is already implied, this phrasing may be redundant.
Context: ...ressive Code, and the 404 route all own their own markup — some of it created after page ...

(PRP_OWN)


[typographical] ~35-~35: Consider using an em dash in dialogues and enumerations.
Context: - **Each tenant has a message queue of it...

(DASH_RULE)


[style] ~35-~35: This sentence is over 40 words long. Consider splitting it up, as shorter sentences make the text easier to read.
Context: - Each tenant has a message queue of its own (internal/mq/{mq,subject,embedded}.go (+ tests), internal/ingest/{sweeper,worker}.go (+ tests), internal/api/{dlq,ingest}.go (+ tests), internal/app/{app,wire}.go (+ tests), internal/stream/{subscriber,hub}.go, internal/settings/{settings,store,registry}.go (+ tests), internal/testutil/{mocks,testutil}.go, clients/ts/src/{dlq,types}.ts (+ tests), docs/src/content/docs/{deployment,api,architecture,ingest-pipeline,durability,why-wavehouse}.md, docs/src/content/docs/sdk/{admin,reference,streaming}.md, docs/src/content/docs/{settings-directory,configuration}.mdx, AGENTS.md): story 5b of the multi-tenant epic (#583). The embedded NATS server keeps each ...

(TOO_LONG_SENTENCE)


[style] ~35-~35: This sentence is over 40 words long. Consider splitting it up, as shorter sentences make the text easier to read.
Context: ...g outside internal/mq names a stream. A tenant at its budget gets 503 while the others keep publishing, and the ingest worker's and the hub bridge's durables are held on every tenant's stream, each with its own ack floor and MaxAckPending, so one tenant's backlog holds back neither another's delivery nor its purge; the worker's prefetch and the hub bridge's fetch-ahead are each shared across the tenants' streams. A removed or rejected tenant's stream i...

(TOO_LONG_SENTENCE)


[style] ~37-~37: Since ownership is already implied, this phrasing may be redundant.
Context: ...ll topic and evaluates each event under its own tenant's policy, so a subscriber on one...

(PRP_OWN)

🔇 Additional comments (40)
clients/ts/src/dlq.ts (1)

2-2: LGTM!

Also applies to: 4-4, 18-22, 26-26, 33-34, 38-38

clients/ts/src/namespaces.test.ts (1)

149-166: LGTM!

clients/ts/src/types.ts (1)

474-475: LGTM!

docs/src/content/docs/sdk/admin.md (1)

68-75: LGTM!

docs/src/content/docs/sdk/reference.md (1)

107-108: LGTM!

docs/src/content/docs/why-wavehouse.md (1)

56-56: LGTM!

Also applies to: 159-159, 193-193, 224-224

internal/api/dlq.go (2)

10-10: LGTM!

Also applies to: 22-28, 34-35, 39-44


30-30: 🔒 Security & Privacy | 🛡️ Detected with Advanced Tier

Nested DLQ reads already require the operator key. The /v1/ops gate disables policy-based admin access for nested settings, so a non-operator admin token receives 403 before DLQHandler.Stats runs.

Likely an incorrect or invalid review comment.

internal/api/dlq_test.go (1)

18-22: LGTM!

Also applies to: 27-34, 36-49, 52-52, 58-58, 61-61, 64-64, 78-82, 85-85, 89-89, 91-91, 104-104, 109-109, 112-115, 122-161

AGENTS.md (1)

32-32: LGTM!

Also applies to: 41-41, 48-48, 59-59

CHANGELOG.md (1)

13-13: LGTM!

Also applies to: 27-27, 35-37

internal/mq/mq.go (1)

144-169: LGTM!

Also applies to: 192-228, 239-285, 299-299, 309-319

internal/mq/subject.go (1)

15-61: LGTM!

Also applies to: 97-115

internal/mq/subject_test.go (1)

4-7: LGTM!

Also applies to: 140-186

docs/src/content/docs/api.md (1)

277-277: LGTM!

Also applies to: 388-388, 748-765, 861-861, 879-881

docs/src/content/docs/architecture.md (1)

80-93: LGTM!

Also applies to: 142-151, 193-193, 231-235, 255-258

docs/src/content/docs/durability.md (1)

36-37: LGTM!

Also applies to: 84-84, 96-97, 104-104

docs/src/content/docs/ingest-pipeline.md (1)

25-33: LGTM!

Also applies to: 49-49, 64-64, 73-73, 92-96, 205-207, 215-217, 224-225, 231-239, 251-251

docs/src/content/docs/deployment.md (1)

382-394: LGTM!

Also applies to: 422-426, 438-445

docs/src/content/docs/settings-directory.mdx (1)

34-34: LGTM!

Also applies to: 127-135, 215-226, 234-234

internal/ingest/worker.go (1)

119-128: LGTM!

Also applies to: 227-245, 476-481, 497-497, 798-800

internal/ingest/worker_test.go (1)

125-125: LGTM!

Also applies to: 155-155, 246-246, 295-295, 1089-1089, 1180-1180, 1268-1271, 1782-1782

internal/api/ingest.go (1)

710-713: LGTM!

internal/api/router_test.go (1)

334-334: LGTM!

internal/mq/embedded_test.go (1)

5-7: LGTM!

Also applies to: 19-74, 135-135, 144-167, 178-178, 254-278, 319-319, 358-420, 422-643

internal/testutil/mocks.go (1)

224-227: LGTM!

Also applies to: 236-250

internal/testutil/testutil.go (1)

17-17: LGTM!

Also applies to: 43-60

internal/settings/settings.go (1)

166-168: LGTM!

Also applies to: 222-231

internal/stream/hub.go (1)

46-48: LGTM!

Also applies to: 397-398

internal/stream/subscriber.go (1)

58-68: LGTM!

docs/src/content/docs/configuration.mdx (1)

100-100: LGTM!

docs/src/content/docs/sdk/streaming.md (1)

119-119: LGTM!

internal/app/app.go (1)

21-21: LGTM!

Also applies to: 91-92, 143-145, 176-176

internal/app/app_test.go (1)

240-240: LGTM!

Also applies to: 249-250, 400-468, 574-617, 733-776

internal/app/wire.go (1)

9-9: LGTM!

Also applies to: 69-86, 108-132, 165-169, 522-541, 559-585, 602-606

internal/ingest/sweeper_test.go (1)

10-10: LGTM!

Also applies to: 19-53, 60-60, 72-72

internal/settings/registry.go (1)

157-167: LGTM!

internal/settings/registry_test.go (1)

309-347: LGTM!

internal/settings/store.go (1)

196-196: LGTM!

internal/mq/embedded.go (1)

637-642: 🩺 Stability & Availability

The concurrent Hub.Broadcast calls do not expose an unsynchronized Hub state race. Broadcast snapshots routes under h.mu.RLock and uses local projection state. deliver protects schema state with schemaMu, and Subscriber.Send uses a channel. Each subscriber receives events from one tenant's delivery goroutine, so different tenant broadcasts do not concurrently update the same subscriber state.

Comment thread internal/ingest/sweeper.go
Comment thread internal/mq/embedded.go
@EricAndrechek

EricAndrechek commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

@/tmp/pr_edits/comment5831496372.md

@taitelee

Copy link
Copy Markdown
Member Author

Those runs were on the pre-fix test: feat/mq-nats-wiring and feat/mq-external-broker both fork from 9344f1b8, one commit before f5d8f484, and their embedded_test.go is byte-identical to it. On a checkout of feat/mq-nats-wiring here the unfixed test fails 5 of 20 runs; with f5d8f484 applied it's 0 of 20, the app test with the same pattern 5 of 5, and the patch applies cleanly (nothing in #639's embedded.go changes touches the open path). eb9f7409 closes a much narrower window of the same kind in TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen. So rather than a second test-only fix on the integration branch, merge #612's head into feat/mq-external-broker — two fixes to the same test would conflict when #612 lands.

@taitelee

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Sep 25, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@taitelee
taitelee marked this pull request as ready for review September 25, 2026 13:09
@taitelee
taitelee requested review from a team and EricAndrechek September 25, 2026 13:09
@taitelee
taitelee merged commit 2a2b886 into main Sep 25, 2026
30 checks passed
@taitelee
taitelee deleted the mq-tenant-streams branch September 25, 2026 13:12
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
#612 was squash-merged into main as 2a2b886 while this branch was open,
and origin/mq-tenant-streams was deleted, so this branch now bases on
main. The squash's content is f5d8f48 (this branch's base) plus one
hunk in internal/mq/embedded_test.go: an "occupied" directory keeping
the streams directory alive through TestEmbeddedNATS_SetMaxBytes_
AQueueThatCannotOpen's failed open. This side already holds all of #612,
so the tree is kept as is (-s ours). The extra hunk is left out on
purpose: #617's perf branch (merged here in e3d83c7) fixes the same race
by waiting for the failed open's cleanup, and the two together cannot
pass (measured: with the hunk the test fails, since the occupied
directory stops the removal it waits for; without it, 5/5 pass).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
Absorbs #612's squash (2a2b886), superseding the older mq-tenant-streams
commit this branch carried. Every file outside this PR's own diff takes
main's version; this PR's hunks were reapplied on top (wire.go comment,
architecture.md and deployment.md cache prose).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
Absorbs #612 (2a2b886), which gave every tenant a queue of its own.
Doc conflicts resolved onto #612's per-tenant wording; the outage test
opens its tenant's queue with SetMaxBytes and names the tenant in
DeadLetterCounts, and the outage docs scope the ack floor, maxAckPending
and the byte budget to the tenant.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
…query-error-classes

Brings in the parent's merge of main (#612) and its review fixes.
Conflicts: AGENTS.md (this PR's api/ line beside main's app/ line) and
CHANGELOG.md (this PR's entry beside the parent's updated one).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
Sync #615 with its parent, which now includes main's #612 squash.
Resolved a conflict in docs/architecture.md (app/wiring section):
kept boot-backends' rewritten prose and combined it with coord-leases'
own additions (the lease coordinator in app.go's boot order, `wireCoord`
in wire.go's backend switch, and the sweeper-lease clause).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
Sync #621 with its parent, which now includes main's #612 squash and
subsequent main history. Resolved conflicts in AGENTS.md and
docs/architecture.md (app/wiring sections): kept cache-snapshot's
rewritten prose and combined it with cache-flat-versions' own
additions (the wireCache/LocalCache.Prune clauses naming cache in the
per-tenant prune set).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
Sync #622 with its updated parent (#615, now synced with main through
#612). Resolved conflicts in AGENTS.md and docs/architecture.md
(app/wiring sections): kept coord-leases' rewritten prose (which
already folds in main's own edits) and combined it with process-roles'
own additions (the roles-aware app.go boot order and the `elected`
wrapper around RunElected in wire.go's sweeper-lease clause).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
…647)

Fixes #617. Part of #613.

`internal/app` and `internal/mq` each took 11–12 s of the 15 s unit-test
budget on main under a full parallel `make test-unit`, and timed out
when the machine was loaded. After this PR they take under 5 s each.
This PR also makes `NewEmbedded` fail at once on a store directory it
cannot create, instead of waiting out the 5 s readiness check.

## What changes

| Commit | Change |
|---|---|
| 2a494bd | `NewEmbedded` runs `MkdirAll` on the store first and
returns `nats store: <mkdir error>`. Before this, a regular file at
`<data_dir>/nats` failed JetStream in the background, and boot reported
only `nats server not ready` after 5 s. New test:
`TestNewEmbedded_AStoreItCannotCreateFailsAtOnce`. |
| d46e273, f897394 | The directory is created at `0700`, the same mode
nats-server uses (`defaultDirPerms`). CHANGELOG entry under Fixed,
limited to what mkdir catches: an existing directory that cannot be
written to still fails the old way. |
| 98acc11 | The `mq` tests call `t.Parallel`. A new
`mq.EmbeddedSyncAlways` is `true` in production and is set to false only
in `TestMain`: on macOS the fsync on every JetStream write was more than
half the run. Each store lives in `storeDir`, which retries its removal
(#442). |
| e441f57 | `internal/app`'s new `TestMain` also turns the fsync off. |
| 668c2c4 | Removes the "occupied"-directory hunk that #612's squash
added to `TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen` (see
below). |

These are 2a19d7e, 96e2112, bd059f2, 3d959b9 and 3c8a93b from
`perf/unit-test-budget`, cherry-picked onto main. Two hunks from the
perf branch are left out because they belong to stacks that are not on
main: the `t.Parallel` lines for `embedded_failed_test.go` and for
`TestEmbeddedNATS_Publish_IdempotencyKeyDropsARepeat`. Neither test
exists on main; they come from the integration tree (#645). The
CHANGELOG entry was re-applied under main's `### Fixed`.

## The `AQueueThatCannotOpen` conflict: measured

#612 (f5d8f48) and the perf change fix the same race in different ways.
The race: after a failed open, the server removes the empty `streams/`
and `$G` directories on a goroutine of its own, and that removal
overlaps the next open. f5d8f48 has three hunks:

- **AQ-occ**: an ignored `occupied` directory keeps `streams/` non-empty
during acme's failed open in `SetMaxBytes_AQueueThatCannotOpen`. This
was added to main's squash.
- **P-globex**: `PacesTheRetriesOfAQueueThatCannotOpen` opens globex
first, so its streams keep the directory occupied.
- **App-globex**: `internal/app`'s `TestNew_QueueOpenFailure` "nested"
subtest blocks globex rather than acme.

The perf change's version of the fix is **AQ-wait**:
`require.Eventually` until the server has removed `$G`, and only then
open globex.

These runs were on this branch, with `-race` and `GOTOOLCHAIN=go1.26.6`,
on 2026-09-25 on an otherwise idle machine. "Isolated" means `-run
'^Name$' -count=20`. "pkg" means the whole parallel `internal/mq`
package with `-count=5`. "Mutation" means `JetStreamMaxStore:
math.MaxInt64` instead of `/ 2`, which is the refusal the test exists to
catch, at `-count=5`.

| P-globex | AQ variant | AQueue isolated | Paces isolated | AQueue in
pkg | Paces in pkg | other pkg fails | Mutation caught |
|---|---|---|---|---|---|---|---|
| kept | **wait only (this PR)** | **20/20 pass** | **20/20** | **5/5**
| **5/5** | **0** | **5/5 fail (caught)** |
| kept | occupied only | 20/20 | 20/20 | 5/5 | 5/5 | 0 | 5/5 fail
(caught) |
| kept | both | **0/20** | 20/20 | 0/5 | 5/5 | — | — |
| dropped | wait only | 20/20 | 20/20 | 5/5 | **4/5** | 1 | 5/5 caught |
| dropped | occupied only | 20/20 | 20/20 | 5/5 | 5/5 | 0 | 5/5 caught |
| dropped | both | 0/20 | 20/20 | 0/5 | 3/5 | 2 | — |

| App-globex | `TestNew_QueueOpenFailure` isolated ×20 | in the
`internal/app` package ×3 |
|---|---|---|
| kept (main) | 20/20 | 3/3 |
| reverted to acme | 20/20 | 3/3 |

What the runs show:

- **Both AQ fixes together always fail**, 20/20. The `occupied`
directory keeps `$G` from ever being removed, so the wait for its
removal times out. The finding from #645 reproduces.
- **P-globex is what `PacesTheRetries` needs.** Without it, the pacing
test raced in the full parallel package (1/5 and 2/5 failures), even
though it passed 20/20 when run alone. This is the finding from #646. It
is a different test from `AQueueThatCannotOpen`, so the two findings do
not conflict. P-globex is on main and this PR keeps it.
- **For `AQueueThatCannotOpen`, AQ-wait alone and AQ-occ alone both
passed every run, and both caught the mutation in this matrix.** The
earlier finding, that the occupied variant still passes with `MaxInt64`,
did **not** reproduce here. On main as merged (occupied only), the
mutation also fails the test 5/5 (measured). I kept AQ-wait because it
is the version the perf change was written and measured against. It is
also what the test's comment describes: no queue is kept open, and no
directory is kept around to hold the reservation count up. Dropping
AQ-occ instead of AQ-wait would work equally well by these numbers.
- **App-globex made no difference in either direction** in 20 isolated
runs and 3 package runs. It is left as main has it.

## Timings: full parallel `make test-unit`

Four runs, alternating between main at 2a2b886 and this branch, with
`-count=1 -race -cover`, on an otherwise idle machine for every run.

| | `internal/app` | `internal/mq` | whole run (`DONE … in`) |
|---|---|---|---|
| main (2a2b886) | 11.30 / 11.37 / 11.41 / 11.45 s | 12.05 / 12.18 /
12.30 / 12.38 s | 12.07–12.39 s |
| this PR | 4.71 / 4.72 / 4.73 / 4.81 s | 3.42 / 3.64 / 3.64 / 3.94 s |
5.07–5.27 s |

An earlier baseline of main alone, under heavy load, gave `internal/app`
12.2–13.5 s and `internal/mq` 12.8–15.1 s. The 15.1 s run was already
over the budget.

## Verification

- `make ci` passed through the shared queue (`GOTOOLCHAIN=go1.26.6`, the
known golangci-lint toolchain workaround), including all coverage gates.
- Pre-push reviewers were both run on HEAD 668c2c4 (opus, fresh
context):
- `pre-push-reviewer`: **ship_it**, with 0 MUST, 0 SHOULD and 0 MAY
findings.
- `docs-reviewer`: **ship_it**, with 0 findings. The CHANGELOG entry was
checked against nats-server's `defaultDirPerms` and `wireMQ`'s error
path. No docs page needed a change.
- Because of #454, the hook wrote the markers against the main
checkout's HEAD, not this branch's HEAD. No marker was written by hand.

## Left for later

- `internal/testutil.NewEmbeddedMQ`, used by the `internal/ingest` and
`internal/api` tests, still fsyncs on every write. Those packages are
not near the budget today. The reviewer noted it as the next place to
get the same speed-up.
- 8776b4d (#635) and 09d5c14 (#639) from the perf branch belong to
their own stacks and are not included here.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
EricAndrechek added a commit that referenced this pull request Sep 25, 2026
A v0.1.0 queue is deleted at boot, so the lenient %2D decoding and the
stream-resume caveat only concern queues an unreleased build since #612
wrote; say so in the changelog and the comments. List keyenc in the
development.md tree and AGENTS.md's file structure, and correct the
stale "L2" cache line beside it.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01B1tJWUp6oaDoH1usLwGtLF
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
)

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>
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
Part of #613 (workstream A).

## What changes

A failed batch insert used to go through row-by-row isolation whatever
the failure, so a ClickHouse that was down, overloaded or read-only
failed every row a second time and parked the whole batch on the DLQ.
Now the worker asks what the failure was first:

- **ClickHouse rejected the row** (`chconn.Rejected`): unchanged.
Row-by-row isolation, the rows rejected again go to the DLQ, or stay
unacked with `dlq.enabled` off.
- **ClickHouse refused a multi-row batch for its size**
(`TOO_MANY_PARTS`, `MEMORY_LIMIT_EXCEEDED`: `chconn.Splittable`): the
same row-by-row isolation first. `TOO_MANY_PARTS` is also what
ClickHouse answers for one INSERT spanning more than
`max_partitions_per_insert_block` partitions: measured on 26.6.3.62, 150
rows over 150 day-partitions fail with 252, and each row inserts on its
own. Retried whole, such a batch re-forms the same way on redelivery and
never lands. If a row fails any way but `Rejected` (typically the first
row failing the same way, because the server is the problem after all),
isolation stops there and that row and the rest are retried under the
backoff below; in the typical case that costs one extra request.
ClickHouse's own Distributed async inserts split a batch on the same
codes (`isSplittableErrorCode`).
- **ClickHouse could not take it** (`Unavailable`, `Denied`, `Unknown`):
no isolation, no DLQ. The batch goes back to the queue with a delayed
nak (`mq.Message.NakWithDelay`, new), under a backoff per ClickHouse
pool (URL + user + database). The backoff starts at 1 s, doubles to a 30
s cap, and each delay is jittered down to as little as half. While the
window runs, every table on that pool is turned away with no request,
and a row that arrives is handed straight back rather than buffered.
When the window ends, one flush probes. Any answer that is not an outage
closes the backoff, and that includes a rejected row. While a probe is
out, arriving rows are handed back with a delay of at least 0.5 s, so a
slow probe does not cycle the backlog. A failure of **one table**
(`TABLE_IS_READ_ONLY`, `TABLE_IS_PERMANENTLY_READ_ONLY`,
`TOO_MANY_PARTS`, `TOO_MANY_MUTATIONS`, `ACCESS_DENIED`:
`chconn.TableScoped`) backs off that table alone: its neighbours keep
inserting until its waiting rows fill the tenant's `maxAckPending` (see
below), and their successes don't close its backoff. The retries are
counted by `wavehouse_ingest_retries_total{table, reason}`. An outage is
logged at `WARN` when it starts, at most every 30 s while it lasts, and
at `INFO` when ClickHouse takes inserts again.
- **ClickHouse goes away mid-isolation**: isolation stops at that row.
The rows already inserted stay acked and the rows already rejected stay
parked. The row that hit the outage, and every row after it, go back to
the queue.

The classifier lives in `internal/chconn/errclass.go`: `Classify(err)
Class`, `ClassOfCode`, `ExceptionCode`, and `HTTPError`/`NewHTTPError`
for the HTTP interface. It reads `*clickhouse.Exception` and
`*clickhouse.HTTPError` from the driver, the HTTP interface's
`X-ClickHouse-Exception-Code` / `Code: NNN.` body, the clickhouse-go
sentinels (`ErrAcquireConnTimeout`, `ErrConnectionClosed`,
`driver.ErrBadConn`) and the network errors under them. So the query
path can reuse it.

## How the ambiguous cases are classed, and why

The line is drawn at **the exception code**.

- **Unlisted exception code → `Rejected` (DLQ).** A code means
ClickHouse was up and read the request. Nearly all of its several
hundred codes are verdicts on what it read, so the availability codes
are the enumerated exception. The cost of each wrong guess settles it.
If an availability code is missing from the list, its rows get parked:
nothing is lost, and they wait on the DLQ (reading them back is #237).
If a genuine bad row were retried forever, it would pin the ingest
stream's ack floor and a share of `maxAckPending` for every table behind
it, and nothing would clear it. A code added by a future ClickHouse
falls into the same bucket.
- **No code, and not a recognisable transport failure → `Unknown`
(retried).** Examples: a bare `500` from something that is not
ClickHouse, or a TLS setup error. Nothing says a row was judged, so
isolating the batch would only multiply the requests, and dead-lettering
it would park good rows.
- **Schema drift: `UNKNOWN_TABLE`, `NO_SUCH_COLUMN_IN_TABLE`,
`UNKNOWN_DATABASE` → `Rejected`.** The row cannot insert into the table
as it now is, and the DLQ keeps it rather than retrying it forever.
`TestDLQ_PopulatedOnIngestWorkerFailure` still pins this.
- **Auth: `AUTHENTICATION_FAILED`, `ACCESS_DENIED`, `UNKNOWN_USER`,
`WRONG_PASSWORD`, `REQUIRED_PASSWORD`, `IP_ADDRESS_NOT_ALLOWED`,
`DATABASE_ACCESS_DENIED`, `USER_EXPIRED` → `Denied` (retried).** These
are the connection's configuration, not the rows. A grant, a corrected
`clickhouse.username`, or `WH_CH_PASSWORD` and a restart makes the whole
batch insert as it is. `ACCESS_DENIED` is usually a grant missing on one
table, so it backs off that table alone (see below). Parking would move
every row of every table on that pool to the DLQ.
- **Codes classed `Unavailable`** (names checked with `errorCodeToName`
on 26.8, and against `system.errors` on 26.6.3.62, the pinned test
image): 95, 96, 159, 160, 164, 173, 201, 202, 203, 209, 210, 225, 236,
241, 242, 243, 244, 252, 254, 265, 279, 285, 286, 289, 297, 319, 364,
369, 394, 415, 425, 439, 499, 574, 692, 700, 735, 745, 749, 762, 774,
776, 777, 778, 904, 999, 1000. `TABLE_IS_PERMANENTLY_READ_ONLY` (774) is
here even though it does not clear by itself: its rows are not at fault,
and an operator fixes it the way one fixes a grant.
`UNKNOWN_STATUS_OF_INSERT` (319) is retried: a duplicate is the lesser
harm. ClickHouse's insert deduplication will not catch it, because the
retried rows come back regrouped with newer ones; the docs now say
delivery after an ambiguous failure is at-least-once.
`SERVER_OVERLOADED` (745) was missing from the first cut: the
CPU-overload check at the top of every query throws it, and so does the
workload scheduler's `max_waiting_queries`.
- **HTTP status with no code** (a proxy answering for ClickHouse):
`408`/`429`/`502`/`503`/`504` → `Unavailable`. `401`/`403`/`407` →
`Denied`. `413` → `Rejected`, because isolation shrinks the body.

Unchanged on purpose: a batch whose tenant has **no ClickHouse
connection** (`parkBatch`) still meets its DLQ switch whole. That rule
exists so that a tenant no longer served does not pin the shared ingest
queue.

## What an operator sees

Rows waiting out an outage stay unacked in the ingest stream. They count
toward `maxAckPending`, and the sweeper cannot purge past them. A long
outage therefore fills the stream to `mq.max_bytes_gb`, and ingest then
answers `503`: backpressure, with nothing lost and nothing parked. The
DLQ no longer fills up. `NakWithDelay` keeps the message on the
consumer's pending list (verified in nats-server `processNak`), so the
queue holds the backlog, not the worker.

A failure of one table counts toward the same budget. One that lasts (a
missing grant, a permanently read-only table) suspends delivery for all
of that tenant's tables once its waiting rows reach `maxAckPending`, and
the tenant's ingest backs up to `503` as in an outage. A retry also does
not keep arrival order; that matters only to a `ReplacingMergeTree`
without a version column or a `CollapsingMergeTree`, and the docs say
so.

## Tests

- `internal/chconn/errclass_test.go`: a table of classifier cases
covering a real refused dial and a real client timeout through
`net/http`, `*clickhouse.Exception`, `*clickhouse.HTTPError`, the driver
sentinels, header- and body-coded HTTP answers, and codeless proxy
statuses. Also pins that the two code lists are disjoint, that every
splittable code is a retried one, and `NewHTTPError` parsing.
- `internal/ingest/backoff_test.go`: escalation to the cap, jitter
bounds, no escalation from late reports, one probe at a time,
rate-limited outage logging, and one backoff per pool.
- `internal/ingest/worker_test.go`:
- (a) an availability failure, in eleven shapes, never dead-letters and
naks with a delay: one request, or two for a splittable code (the
split's first row fails the same way);
- a batch refused for its size (`TOO_MANY_PARTS` as too many partitions,
`MEMORY_LIMIT_EXCEEDED`) is split, every row inserts, and no backoff
opens; a lone row is not split again;
  - (b) a `CANNOT_PARSE_NUMBER` row is still isolated and parked;
- (c) ClickHouse going down mid-isolation stops isolation and retries
the unsettled rows;
  - an outage stops later column groups;
- tables on a down pool back off together, and one probe closes the
backoff;
  - another pool is unaffected;
- the batcher buffers nothing while its pool backs off, or while a probe
is out;
- a read-only table backs off alone, and a healthy neighbour's success
does not close its backoff.
- `tests/integration/ingest_outage_test.go`: stops a real ClickHouse
container under a running worker and publishes three rows. It checks
that none is parked for 12 s, then restarts ClickHouse and waits until
all three land, still with none parked. Run against the old `worker.go`,
it fails on both assertions.

## Review

The `pre-push-reviewer` and `docs-reviewer` subagents (opus) ran in
fresh context against this worktree, four rounds:

- **Round 1** (`a131fb95`): both `iterate`. Code review: table-scoped
codes flapped the shared pool breaker, and the delay while a probe was
out was near zero. Docs review: an overclaim that the query handlers
already use `Classify`, `NakWithDelay` missing from the `mq` surface
list, an incomplete `Denied` list. Fixed in `8ae9614b`.
- **Round 2** (`8ae9614b`): both `iterate`. Code review: `ACCESS_DENIED`
is usually one table's grant, so it should be table-scoped; a stale
bound in the backoff map comment. Docs review: the docs pointed
operators at a ClickHouse password in `config.json` (it is
`WH_CH_PASSWORD`, boot config), and overstated the pool key. Fixed in
`356125f3`.
- **Round 3** (`356125f3`): code `ship_it`; docs `iterate`, because the
README, landing page, why page and architecture diagram still said
failed inserts go to the DLQ. Fixed in `e97edc8b`.
- **Round 4** (`e97edc8b`): **both `ship_it`**.
- **Round 5** (`c70d3324`..`d1ef1302`): the size split for
`TOO_MANY_PARTS`/`MEMORY_LIMIT_EXCEEDED`, `SERVER_OVERLOADED`, and doc
corrections. Code review of `c70d3324`: **`ship_it`**, no findings. Docs
review took nine passes to reach **`ship_it`** on `d1ef1302`. Along the
way it fixed:
- a drain runbook that described this build's retry signals for a drain
run on the earlier build (which dead-letters an outage, so ClickHouse
must stay healthy for the whole drain);
  - the split exception missing from two summaries;
  - a lock-free claim the shared backoffs broke;
- duplicates after an ambiguous insert failure, which neither
ClickHouse's insert deduplication nor `dedupe.enabled` catches;
  - version-column and `FINAL` advice where a user chooses an engine;
  - a set of missed copies of "a failed insert goes to the DLQ".

**Markers for round 5:** the SubagentStop hook wrote none. The reports
came back through the subagent hand-back, so the payload it reads
evidently carried no `VERDICT:` line; fed the report text, the hook
writes a marker (tested in a scratch repo). Both verdicts are recorded
with `scripts/skip-pre-push-review.sh`, and each reason names the run.
`docs-reviewer` ran on HEAD. `pre-push-reviewer` ran on `c70d3324`;
everything after it is docs prose and two doc-comment rewordings.

**Review markers (#454):** the `review-marker.sh` SubagentStop hook
writes to `$CLAUDE_PROJECT_DIR/tmp`, keyed to that checkout's HEAD (the
main checkout, on `main`), so no `tmp/<reviewer>-passed-e97edc8b…`
marker exists for this branch. The verdicts above are the record; no
marker was hand-written or skipped.

## Follow-ups (not in this PR)

- #403 / #271: map `/v1/query` and `/v1/ops/query` failures through
`chconn.Classify` and `ExceptionCode`. `Rejected` → 4xx (`ACCESS_DENIED`
→ 403), `Unavailable` → 503/502 with `retryable: true`. The classifier
is exported for that.
- #613 C (distributed workers): see the seams below.

## Seams left for workstream C

- The backoff state is in-process: `IngestWorker.backoffs`, keyed by
`chconn.Target` (URL, user, database), plus the table for table-scoped
failures. With several worker processes, each one backs off on its own.
That is harmless, because each probes at most once per window, but a
shared backoff could live behind the coordination primitive (B).
- `flushTable` hands back exactly the unsettled rows through
`retryLater`. Any batching/claiming scheme that owns acks can keep that
contract.
- `tableBatcher.add` asks the backoff before buffering. A per-tenant
consumer (#612) could pause that tenant's `Consume` here instead of
nak'ing arrivals.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
https://claude.ai/code/session_01LYbFK7Na4RTNZxzL3h9Hs9

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
Part of #613. Stacks on #612 (`mq-tenant-streams`). This is PR **G1** of
the #613 core design.

## What

Each layer's implementation is now chosen once, at boot:

| key | env | values today | default |
|---|---|---|---|
| `mq.backend` | `WH_MQ_BACKEND` | `embedded` | `embedded` |
| `cache.backend` | `WH_CACHE_BACKEND` | `local` | `local` |
| `dedupe.backend` | `WH_DEDUPE_BACKEND` | `pebble` | `pebble` |
| `coord.backend` | `WH_COORD_BACKEND` | `local` (reserved: nothing
reads it until B1) | `local` |

- `internal/config/backends.go`: one enum type per layer, its list of
valid values, and a `validate()` per layer block. `Validate` refuses a
value this build has no backend for and names the valid ones. A
`<layer>.<backend>` sub-block written before its backend lands (for
example `mq.nats`) is refused by the strict loader as an unknown key.
- `Config.Distributed()` (true when the mq is not `embedded`),
`Config.NeedsDataDir()` (gates `CheckDataDir` in `main`), and
`Config.Warnings()`, which `app.New` logs at WARN after observability is
wired. The two warnings keyed on a shared queue (local cache, pebble
dedupe) ship now. Nothing can trigger them until D1, so they are
unit-tested with a literal non-embedded value.
- `internal/app/wire.go`: `wireMQ`, `wireCache` and `wireDedupe` are
each a `switch` on the layer's backend. The old bodies are now
`wireEmbeddedMQ` and `wirePebbleDedupe`, unchanged, and cache's single
case is inline. The `default:` case (`unreachableBackend`) refuses boot.
Only a `config.Config` built without `config.Load` can reach it, because
the zero value is not the default. The three hand-built configs
(app_test, two integration tests) now name every backend.
- Docs: a new Backends section in `configuration.mdx` (table, sub-block
convention, the shared `mq`/`dedupe` block names with config.json),
updates to the example file and env, the data_dir probe step,
`architecture.md` (the wire.go switches), `settings-directory.mdx`, the
`TenantConfig` comment in `internal/settings/settings.go`, the `config/`
entry in AGENTS.md, and root `config.yaml`. CHANGELOG entry.

## Verification

- `make ci` is green at f129d57.
- pre-push-reviewer returned ship_it at fde17ba. The only commit after
it is prose, so its skip is logged.
- docs-reviewer returned ship_it at f129d57. Its SubagentStop hook
wrote no marker in this worktree, so the result is recorded with the
skip script, and the skip reason says so.

## Left to later PRs (by design)

- `roles`, `instance_id`, `Has(Role)` and cross-layer rules 2 and 5 go
to **C1**.
- `coord.backend=nats` and rules 3 and 4 go to **B2**. Rule 4 cannot
land before B2, or D1's `mq.backend=nats` would be unbootable until B2
exists.
- `wireCoord` goes to **B1**, which should switch on
`a.cfg.Coord.Backend` the same way.
- The `mq.nats` sub-block, the `MQNATS` value and the "`mq.max_bytes_gb`
is not applied" warning go to **D1**. `cache.redis` goes to **E1** and
`dedupe.dynamodb` to **F1**.

Adding a backend takes one constant appended to the layer's list in
`backends.go`, a sub-block field plus a case in that layer's
`validate()`, and a case in the layer's wire switch.

🤖 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>
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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>
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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
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/docs Documentation, site/, README area/ingest Ingest pipeline (Bento, batching, DLQ) area/sdk TypeScript SDK (clients/ts/) documentation Improvements or additions to documentation go Pull requests that update go code

Projects

Status: Done

Development

Successfully merging this pull request may close these issues.

2 participants