Repository navigation
refactor(mq): seal the MQ boundary behind an intent-level broker API - #586
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Advanced Run ID: 📒 Files selected for processing (2)
Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review. 📜 Recent review details🧰 Additional context used📓 Path-based instructions (1)Never hard-wrap prose.📄 CodeRabbit inference engine (AGENTS.md) Files:
🪛 LanguageTooldocs/src/content/docs/ingest-pipeline.md[style] ~266-~266: 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. (EN_REPEATEDWORDS_NEED) 🔇 Additional comments (2)
📝 SummarySummary by CodeRabbit
WalkthroughThe pull request seals NATS and JetStream access inside ChangesMQ boundary migration
Priority: ➖ Normal Estimated code review effort: 5 (Critical) | ~90 minutes Change: Refactor Merge Risk: ⚪ Minimal · up to The MQ boundary documentation is corrected, and no current merge-blocking behavior was established in the reviewed paths. 🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
Full details: Docstring CoverageExplanation Docstring coverage is 50.94% which is insufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 159 functions across 39 files. (2 skipped: 2 unsupported.)
✨ Finishing Touches🧪 Generate unit tests (beta)
✨ Simplify code
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 |
|
📚 Docs preview is live → https://be9beaf4-wavehouse-docs.wave-rf.workers.dev |
Code Coverage OverviewLanguages: Go GoThe overall line coverage in commit e7df851 in the Show a line coverage summary of the most impacted files.
Updated |
|
@coderabbitai review |
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 3
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Advanced
Run ID: 82c10037-1ba4-4732-a3e9-a0f6c4440769
📒 Files selected for processing (29)
.golangci.ymlAGENTS.mdCHANGELOG.mdcmd/wavehouse/main.godocs/src/content/docs/architecture.mddocs/src/content/docs/development.mdinternal/api/dlq.gointernal/api/dlq_test.gointernal/api/router.gointernal/api/router_test.gointernal/api/stream.gointernal/cache/local.gointernal/cache/version_manager.gointernal/cache/version_manager_test.gointernal/ingest/sweeper.gointernal/ingest/sweeper_test.gointernal/ingest/worker.gointernal/ingest/worker_test.gointernal/mq/embedded.gointernal/mq/embedded_test.gointernal/mq/mq.gointernal/mq/mq_test.gointernal/observability/metrics.gointernal/observability/metrics_test.gointernal/observability/tracer.gointernal/observability/tracer_test.gointernal/testutil/mocks.gotests/integration/dlq_test.gotests/integration/setup_test.go
💤 Files with no reviewable changes (1)
- internal/api/router.go
Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review.
📜 Review details
🧰 Additional context used
📓 Path-based instructions (1)
**WH001 applies to every tracked Markdown file, with no carve-out** — `AGENTS.md`, `CHANGELOG.md`, `.github/` CI docs and `.claude/` agent prompts included.
📄 CodeRabbit inference engine (AGENTS.md)
Files:
docs/src/content/docs/architecture.mdAGENTS.mdCHANGELOG.mddocs/src/content/docs/development.md
🧠 Learnings (2)
📓 Common learnings
Learnt from: CR
Repo: Wave-RF/WaveHouse
Timestamp: 2026-09-16T00:33:47.592Z
Learning: **Validate locally before every push** — run `make ci` the documented way ([§Running `make ci`](`#running-make-ci-for-agents`)). Don't use CI as your first feedback loop.
📚 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/ingest/sweeper_test.go
🪛 Betterleaks (1.8.1)
internal/ingest/worker_test.go
[high] 50-50: Detected a potential hardcoded password literal, which may expose account credentials.
(generic-password)
🪛 LanguageTool
docs/src/content/docs/architecture.md
[grammar] ~82-~82: Ensure spelling is correct
Context: ...ave-RF/WaveHouse/issues/319)). Gap-fill replay (mq.Replayer.ReplaySince — a `Deliver...
(QB_NEW_EN_ORTHOGRAPHY_ERROR_IDS_1)
[typographical] ~130-~130: Consider using an em dash in dialogues and enumerations.
Context: - worker.go — StartIngestWorker lau...
(DASH_RULE)
CHANGELOG.md
[typographical] ~27-~27: Consider using an em dash in dialogues and enumerations.
Context: - **The MQ boundary is sealed: only `inte...
(DASH_RULE)
docs/src/content/docs/development.md
[grammar] ~430-~430: Please add a punctuation mark at the end of paragraph.
Context: ...verything else goes through an mq-owned type Formatting (gofumpt — strict super...
(PUNCTUATION_PARAGRAPH_END)
🔇 Additional comments (20)
internal/observability/tracer.go (1)
9-13: LGTM!Also applies to: 15-19, 22-24, 26-27, 35-39, 42-46, 49-49
internal/observability/tracer_test.go (1)
16-16: LGTM!Also applies to: 40-40, 43-43, 51-52, 59-59, 70-71, 73-73, 75-75, 81-81, 84-88, 91-91, 95-95
internal/cache/version_manager.go (1)
21-21: LGTM!internal/cache/local.go (1)
27-27: LGTM!internal/cache/version_manager_test.go (1)
11-11: LGTM!Also applies to: 26-26, 40-40, 56-56
.golangci.yml (1)
19-19: LGTM!Also applies to: 66-75
internal/observability/metrics.go (1)
11-17: LGTM!Also applies to: 20-23, 36-40
internal/observability/metrics_test.go (1)
5-5: LGTM!Also applies to: 88-120, 122-147
AGENTS.md (1)
40-41: LGTM!Also applies to: 71-71, 434-435
CHANGELOG.md (1)
27-27: LGTM!docs/src/content/docs/architecture.md (1)
63-63: LGTM!Also applies to: 82-84, 130-130, 137-140, 146-147
docs/src/content/docs/development.md (1)
430-430: LGTM!Also applies to: 463-463
internal/mq/mq_test.go (1)
56-73: LGTM!Also applies to: 75-83
internal/api/dlq.go (1)
15-15: LGTM!Also applies to: 19-20, 26-26, 41-41, 48-48, 58-58
internal/api/dlq_test.go (1)
23-23: LGTM!Also applies to: 50-50, 54-54, 57-57, 60-60, 86-86, 88-88, 90-90
internal/api/router_test.go (1)
337-337: LGTM!internal/api/stream.go (1)
18-18: LGTM!Also applies to: 23-24, 108-109, 136-137, 140-141, 175-180
cmd/wavehouse/main.go (1)
399-399: LGTM!Also applies to: 408-408, 425-425, 454-454, 465-465, 497-497, 507-507
tests/integration/dlq_test.go (1)
65-65: LGTM!Also applies to: 109-109
tests/integration/setup_test.go (1)
151-151: LGTM!Also applies to: 171-171, 326-326, 329-329
# Conflicts: # cmd/wavehouse/main.go # docs/src/content/docs/architecture.md # internal/api/stream.go # tests/integration/setup_test.go
|
@coderabbitai review |
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to GitHub limitations.
🟡 Minor · Document Topic-based publishing instead of raw NATS publishing. · architecture.md:81
docs/src/content/docs/architecture.md:81
📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick winDocument Topic-based publishing instead of raw NATS publishing.
This bullet still says
internal/api/ingest.gopublishes to theingest.{table}NATS subject. After this change, application code should publish throughmq.Publisherwithmq.Topic{Table, Scope}. Keep raw subject naming and token encoding documented only underinternal/mq.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Advanced
Run ID: 1b2d3ec9-64c5-487d-bcb5-d2e75fe9ce5e
📒 Files selected for processing (10)
CHANGELOG.mddocs/src/content/docs/architecture.mddocs/src/content/docs/ingest-pipeline.mdinternal/app/app_test.gointernal/app/wire.gointernal/ingest/worker.gointernal/ingest/worker_test.gointernal/mq/embedded.gointernal/mq/embedded_test.gointernal/mq/mq.go
Included review availability: Your plan provides up to 1 included review per hour; 0 remain after this review.
📜 Review details
🧰 Additional context used
🧠 Learnings (2)
📓 Common learnings
Learnt from: CR
Repo: Wave-RF/WaveHouse
Timestamp: 2026-09-17T19:45:14.078Z
Learning: Run `make lint` and `make test` before considering work complete.
📚 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/app/app_test.go
🪛 LanguageTool
CHANGELOG.md
[style] ~29-~29: Consider using “who” when you are referring to a person instead of an object.
Context: ...d as a missing sequence; and a consumer that dies underneath the ingest worker no lo...
(THAT_WHO)
docs/src/content/docs/ingest-pipeline.md
[style] ~204-~204: 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] ~206-~206: To elevate your writing, try using a synonym here.
Context: ... could delete the durable) this path is hard to reach today; it matters once a remot...
(HARD_TO)
docs/src/content/docs/architecture.md
[style] ~145-~145: Since ownership is already implied, this phrasing may be redundant.
Context: ...terer.DeadLetter(park a message under its own topic; the caller acks),DeadLetterSta...
(PRP_OWN)
🔇 Additional comments (8)
docs/src/content/docs/architecture.md (1)
64-64: LGTM!Also applies to: 83-85, 91-91, 136-139, 143-148, 154-155
docs/src/content/docs/ingest-pipeline.md (1)
18-18: LGTM!Also applies to: 173-173, 202-207, 230-230
internal/mq/mq.go (1)
186-192: LGTM!Also applies to: 195-197
internal/mq/embedded.go (1)
329-344: LGTM!Also applies to: 352-372
internal/mq/embedded_test.go (1)
170-170: LGTM!Also applies to: 549-549, 586-615, 617-632
internal/app/app_test.go (1)
363-396: LGTM!internal/ingest/worker.go (1)
53-57: LGTM!Also applies to: 121-149, 169-170, 193-193, 229-237, 269-278
internal/app/wire.go (1)
338-357: LGTM!
…fy changelog wording
|
@coderabbitai review |
✅ Action performedReview finished.
|
…593) ## Summary Story 1 of the multi-tenant epic: every request resolves to a tenant before authentication, and the tenant is threaded through every settings read. No behavior change for a deployment that sends no tenant header — only tenant `0` exists. - New `internal/tenant`: `ID` (a validated string — letters, digits, `_`, `-`, ≤ 64 bytes, safe as a folder name and an MQ subject token), `Default = "0"`, `Header = "X-Tenant-ID"`. Imports nothing from the repo. - `settings.Registry` (`For(id)`) over the one store `settings.Open` adopts, keyed `0`. Open, reload and the watcher are untouched. - `api.TenantMW` runs ahead of auth on every `/v1` route outside `/v1/ops/*`: absent or empty header is tenant `0`, malformed or repeated is `400`, unknown is `404`; the resolved `*settings.Store` rides the request context. Handlers read it once and pass it down — the ingest, structured-query and pipe getters take the store as a parameter, and a tenant route reached without one answers `500` rather than fall back to a tenant. The async paths (worker, sweeper, hub, schema registry) are constructed with a `tenant.ID` and their getters take it, wired with `tenant.Default` in `internal/app`. - The probes, `/version`, the metrics path and `/v1/ops/*` stay tenant-exempt; the admin pipe reads serve the default tenant. `X-Tenant-ID` joins the CORS allow-headers list so browser clients can send it. - The slog cleanup deferred from #586: no constructor takes a `*slog.Logger` anymore (mq, the handlers, auth, discovery, sweeper, worker, plus `settings.Open`, `chconn.Open` and the `config` data-dir helpers); call sites use the context-aware calls. Tests reach log output through the new `internal/testutil/logtest` (`Silence` in `TestMain`, `Capture` for log-asserting tests, which run serially). `testutil.NopLogger` is removed. The tenant work and the logger cleanup are separate commits; the rest are review rounds. ## Test plan - [x] `make ci` passes locally (unit, integration, e2e, coverage gates; Go total 92.8%) - [x] `internal/tenant`: grammar table test (default, 19-digit id, length cap, dots, slashes, wildcards, non-ASCII) - [x] `api.TenantMW`: no header / empty / `0` / unknown `404` / malformed `400` / repeated `400`; tenant resolves before `AuthMW`; exempt routes ignore the header; tenant-route handlers `500` without a resolved store - [x] `internal/app`: the wired registry answers `404`/`400` end to end; ops ignores the header - [x] async registry miss: the DLQ switch reads as on, the schema refresh survives a zero interval at boot and mid-run (each test fails with its fix reverted) - [x] logger: log-asserting tests in api, auth, discovery, config converted to `logtest.Capture` - [ ] Manual: `curl -H 'X-Tenant-ID: acme' /v1/health` → `404`; no header → as before ## Follow-ups - Which tenant an ops read addresses is story 2's: `?tenant=` on the reload route and the pipe reads, with the `bearerToken` query-rewrite fix, the per-tenant admin gate, and the SDK option to send it. - An unknown tenant on an async path reads as the getter's zero value (nil policy, gap window 0), with two fail-safes: the worker's DLQ switch reads as on (an unreadable message is parked, never dropped) and the schema refresh keeps its cadence rather than take a zero interval. Impossible while only tenant `0` exists; story 3 defines a removed tenant's semantics. - Auth's policy source (story 9) and the `chconn`-backed getters (story 6) are not threaded. ## Related Issues Part of #583 (story 1). Finishes the logger cleanup deferred from #586. <!-- Checklist for the author (not kept in the squash commit message): - `make ci` passes locally - Docs updated per AGENTS.md "Documentation & Consistency Sync" rules - CHANGELOG.md [Unreleased] entry added - Tests cover new / changed behavior (70 % minimum, 80 %+ preferred) The PR title is the squash commit subject — use Conventional Commits (`feat:`, `fix:`, `docs:`, `refactor:`, `test:`, `chore:`, `ci:`, `deps:`, `build:`, `perf:`, `revert:`, `style:`). The PR body below is the squash commit message, so keep it tight. -->
Summary
internal/mqis now the only package that imports NATS/JetStream and the only one that knows how the broker works. The import half is enforced by adepguardrule in.golangci.yml(Key Design Decision #20). The semantic half: nothing outsideinternal/mqbuilds a subject, names a stream, or reasons in sequences, so story 5's tenant token lands inmq.Topicandinternal/mq/subject.gowith no call-site edits.mq.Topic{Table, Scope}with raw names. Theingest./dlq.prefixes,>wildcards, stream names, and the subject-token encoder (formerlyquery.SafeEncodeNATS/SafeDecodeNATS) are private tointernal/mq/subject.go. A deliveredMessageexposesTopicKey()(free; the delivered form) andTopic()(decodes on demand), so the per-message path decodes nothing.Publisher(ErrQueueFullis the backpressure signal; the API no longer matches broker error text),Subscriber,ConsumerManager/Consumer,DeadLettererandDeadLetterStats,Purger.PurgeAcked,Replayer.ReplaySince, composed intomq.Broker.internal/appholds aBroker;mq.NewEmbeddedis the one place the implementation is chosen.internal/ingest/sweeper.go→internal/mq/purge.go; the sweeper keeps the schedule and the gap window), and themq.max_bytes_gbreload frominternal/app(SetMaxBytes: ingest stream and the DLQ at a tenth of it resized as a pair, with the rollback and time bounds from refactor(app): extract the process wiring from main.go into internal/app #585;NewEmbeddedcreates both streams, replacingapi.EnsureDLQStreamandResize).JetStream(),NatsConn(),GetServer(), the unreadapi.Dependencies.JS, and the never-wired cache connection parameter are removed.mq.Headers+WithHeaderreplace*nats.Msgoptions;observability.InjectHeaders/ExtractHeaderswork over a plain header map;RegisterSystemMetricstakes afunc() (MQStats, error).query.SafeEncodeToken(same output, cache keys unchanged).MockMessage,MockPurger,MockDeadLetterStats;MockPublisherrecords the topic and headers of every publish and dead-letter parking.Wire behavior is preserved: subjects, stream names, consumer settings, ack semantics, and the
X-DLQ-*headers are byte-for-byte what they were; parking on the DLQ is a prefix swap on the delivered subject. Merged with main (#585), portinginternal/appto the mq-owned types and keeping its shutdown-cancelled gap-fill (ReplaySincechecks its context between pulls).Deliberate behavior differences, none on the wire:
GET /v1/ops/dlq/statsreturns the documented500when the broker cannot be read; only a genuinely absent queue reads as empty.WARNinstead of ending the replay silently; a cancelled replay is not logged as a failure.Deferred: terminal consumer errors on
mq.Consumer(#587); breaking scope out of DLQ counts (#235). Theintegration-tagged tests sit outside lint's build context, so the boundary there rests on convention (documented).Test plan
1b656c12(lint incl. thedepguardboundary rule, unit, integration, e2e, coverage, docs build)internal/mq: subject round-trip and injectivity, dead-letter parking and counts (incl. a foreign subject parking under the same tail),ErrQueueFull,PurgeAckedend to end,SetMaxBytespair resize / rollback / cancelled context, replay incl. cancellation, consumer config read-back, headers, statsMockPublisher/MockMessageand the real embedded server; app asserts a reload reaches the broker's byte budgetRelated Issues
Part of #583 (story 4, plus the broker abstraction story 5 builds on). Follow-ups: #587, #235.