Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
142 commits
Select commit Hold shift + click to select a range
ef93154
feat(mq): give every tenant a queue of its own
taitelee Sep 24, 2026
8ad9865
fix(mq): reopen detached from the request, and the upgrade runbook swept
taitelee Sep 24, 2026
650a28e
fix(app): boot opens queues under New's context; docs review fixes
taitelee Sep 24, 2026
999c1db
fix(mq): boot re-applies a budget to a split queue pair; review fixes
taitelee Sep 24, 2026
db2d20a
fix(mq): a failed resize restores the ingest stream's own cap; review…
taitelee Sep 25, 2026
07c6a91
fix(app): a rejected tenant keeps its replay history; review fixes
taitelee Sep 25, 2026
04aad9f
fix(mq): share the hub bridge's fetch-ahead across tenants; review fixes
taitelee Sep 25, 2026
57870c4
fix(mq): publish only into a queue the broker recorded open; review f…
taitelee Sep 25, 2026
7cdb794
docs(mq): a consumer that cannot join a queue opened at runtime; revi…
taitelee Sep 25, 2026
e199a03
fix(mq): pace publish-side retries of a queue that cannot open; revie…
taitelee Sep 25, 2026
55137d8
test(mq): one tenant's failed purge stops no other; review fixes
taitelee Sep 25, 2026
060ca9e
docs(mq): size a dead-letter stream the shrink guard kept; review fixes
taitelee Sep 25, 2026
3022b92
feat(config): choose each layer's implementation at boot
EricAndrechek Sep 25, 2026
651956a
fix(discovery): jitter the refresh retry backoff
EricAndrechek Sep 25, 2026
77ef65d
test(mq): pin interest partitions with a sourced history (S1)
EricAndrechek Sep 25, 2026
fde17ba
docs(config): say coord.backend is reserved; sync the boot-config lists
EricAndrechek Sep 25, 2026
431084e
fix(cache): snapshot versions at lookup; tenant token on every key
EricAndrechek Sep 25, 2026
13422da
feat(coord): leases, in-process implementation
EricAndrechek Sep 25, 2026
a82f74b
docs(coord): state what an unfenced sweeper overlap can cost
EricAndrechek Sep 25, 2026
f129d57
docs(config): no backend has a sub-block yet; index backends.go in AG…
EricAndrechek Sep 25, 2026
bb027da
test(mq): one conformance suite for every Broker
EricAndrechek Sep 25, 2026
d913a18
docs(coord): name ErrClosed as RunElected's other exit; changelog files
EricAndrechek Sep 25, 2026
a131fb9
fix(ingest): retry ClickHouse outages instead of dead-lettering
EricAndrechek Sep 25, 2026
1adc286
test(mq): run the embedded conformance in a test binary of its own
EricAndrechek Sep 25, 2026
80d6c22
docs(mq): keep the per-tenant no-queue case in ErrQueueFull's contract
EricAndrechek Sep 25, 2026
58d9e77
fix(dedupe)!: reserve, commit or release ids keyed by tenant and table
EricAndrechek Sep 25, 2026
735fd78
feat(mq): topology spec and verifier for operator-owned NATS
EricAndrechek Sep 25, 2026
8ae9614
fix(ingest): back off a failing table alone; floor the probe-out delay
EricAndrechek Sep 25, 2026
bf11ecc
test(mq): pin the exactly-once failed report through durable deletion
EricAndrechek Sep 25, 2026
0ec037e
fix(dedupe): state what Managed guarantees a backend; default a zero …
EricAndrechek Sep 25, 2026
2a55fb1
docs(discovery): sweep the remaining unjittered backoff claims
EricAndrechek Sep 25, 2026
356125f
fix(ingest): back off a table denied a grant alone; fix password docs
EricAndrechek Sep 25, 2026
ef4153c
docs(api): list the unavailable broker among the request aborts
EricAndrechek Sep 25, 2026
9344f1b
fix(mq): pace a park's reopen, and warn only on a missing consumer; r…
taitelee Sep 25, 2026
e97edc8
docs: say the DLQ takes rejected rows, not every failed insert
EricAndrechek Sep 25, 2026
f484511
fix(dedupe): release only claimed claims on a wrong-count answer
EricAndrechek Sep 25, 2026
68b74ff
fix(mq): refuse per-subject eviction, unattached sources, shared part…
EricAndrechek Sep 25, 2026
67df317
perf(cache): flat version index, pruned per tenant
EricAndrechek Sep 25, 2026
8792d89
test(mq): end delivery only once the pulls are live
EricAndrechek Sep 25, 2026
5a9d1f1
Merge origin/feat/boot-backends into feat/coord-leases
EricAndrechek Sep 25, 2026
e0705f2
feat(coord): coord.backend selects the coordinator
EricAndrechek Sep 25, 2026
20c28d1
Merge remote-tracking branch 'origin/mq-tenant-streams' into feat/ded…
EricAndrechek Sep 25, 2026
386ce74
docs(dedupe): link the old-key sweep to #220 rather than promise it
EricAndrechek Sep 25, 2026
ed8022b
docs(cache): name the cache prune hook; pin LocalCache as a pruner
EricAndrechek Sep 25, 2026
50b7e0c
feat(cache): Redis-compatible shared cache backend
EricAndrechek Sep 25, 2026
3425243
Merge remote-tracking branch 'origin/mq-tenant-streams' into feat/mq-…
EricAndrechek Sep 25, 2026
beab0fd
feat(app): process roles
EricAndrechek Sep 25, 2026
01ea690
fix(api): map ClickHouse query failures by class, not HTTP status
EricAndrechek Sep 25, 2026
e346426
fix(cache): count only transport failures against the breaker
EricAndrechek Sep 25, 2026
e2434d6
fix(mq): a replay whose connection closed is not caught up
EricAndrechek Sep 25, 2026
e78ea2a
test(mq): keep the unit suite inside its per-package timeout
EricAndrechek Sep 25, 2026
a6fdcfd
docs(changelog): leave the other entries' spelling as it was
EricAndrechek Sep 25, 2026
5b5f988
Merge remote-tracking branch 'origin/mq-tenant-streams' into feat/mq-…
EricAndrechek Sep 25, 2026
d2a87ab
fix(cache): replace token and value keys of the wrong type
EricAndrechek Sep 25, 2026
5de4fd0
docs(app): scope instance_id and sweeper claims to what ships today
EricAndrechek Sep 25, 2026
0fb7637
test(mq): give the S1 tests a package of their own
EricAndrechek Sep 25, 2026
c63ec9a
ci(cov): keep the external-NATS topology out of the e2e suite's gate
EricAndrechek Sep 25, 2026
4ff30f3
fix(api): let ClickHouse report a role time cap overrun itself
EricAndrechek Sep 25, 2026
aa9a655
Merge remote-tracking branch 'origin/feat/mq-conformance' into feat/m…
EricAndrechek Sep 25, 2026
d67a46a
test(cov): keep the unreachable Redis backend out of the e2e gate
EricAndrechek Sep 25, 2026
148a160
test(app): tell a verified token by the ClickHouse error, not the status
EricAndrechek Sep 25, 2026
3a314af
test(app): tell a verified token by the ClickHouse error, not the status
EricAndrechek Sep 25, 2026
44cec4d
Merge remote-tracking branch 'origin/feat/boot-backends' into feat/ca…
EricAndrechek Sep 25, 2026
f485505
test: a role cap overrun is 400 limit_exceeded in the limits test
EricAndrechek Sep 25, 2026
108499f
feat(dedupe): DynamoDB backend, conformance-tested on dynamodb-local
EricAndrechek Sep 25, 2026
37ef757
fix(ingest): windowed reserve/publish/commit; 503 when dedupe is unav…
EricAndrechek Sep 25, 2026
85803f0
Merge remote-tracking branch 'origin/feat/cache-flat-versions' into f…
EricAndrechek Sep 25, 2026
71ab824
fix(ingest): qualify the duplicate-window claims; steadier tests
EricAndrechek Sep 25, 2026
6f7944a
fix(dedupe): attempt every Commit and Release chunk; keep the breaker…
EricAndrechek Sep 25, 2026
ff047d2
test(dedupe): make the Dynamo failure tests fail without their fix
EricAndrechek Sep 25, 2026
0aaeddb
feat(mq): a Broker over an external NATS cluster
EricAndrechek Sep 25, 2026
9824663
docs(ingest): Release gives back definite failures only; tighten wording
EricAndrechek Sep 25, 2026
99a5eee
docs(ingest): say the Pebble benchmark stubs the queue; finish the re…
EricAndrechek Sep 25, 2026
c068084
feat(config): cache.backend=redis selects the shared cache
EricAndrechek Sep 25, 2026
83a6949
Merge origin/feat/boot-backends into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
c452cb6
fix(mq): count a replay down, and size the duplicate window to retries
EricAndrechek Sep 25, 2026
c8ad3f0
fix(mq): keep a replay consumer through a slow client's batch
EricAndrechek Sep 25, 2026
ced118c
feat(app): choose the DynamoDB dedupe backend at boot
EricAndrechek Sep 25, 2026
13c6cb3
fix(cache): review round — e2e gate, trust boundary, #386, addr spaces
EricAndrechek Sep 25, 2026
b83e670
feat(dedupe): retention per tenant and table, and an expiry sweep
EricAndrechek Sep 25, 2026
a8e43bd
fix(app): retry a failed DynamoDB table check; refuse no region
EricAndrechek Sep 25, 2026
6c74b65
test(dedupe): pin that a commit mid-sweep-chunk survives; docs wording
EricAndrechek Sep 25, 2026
1134d5c
fix(pipes): run write pipes every call, uncached and uncoalesced
EricAndrechek Sep 25, 2026
0b60451
docs(dedupe): name every region source; untangle the dynamodb check c…
EricAndrechek Sep 25, 2026
252ef7b
fix(config): keep an explicit false/0/"" from config.yaml
EricAndrechek Sep 25, 2026
74a7bf4
fix(pipes): no-store on write pipes; reconcile the mutation-path docs
EricAndrechek Sep 25, 2026
7bb8c27
docs(pipes): name the operator as a write pipe's author; no-store in …
EricAndrechek Sep 25, 2026
2757e64
fix(test): retry removing an embedded broker store in testutil
EricAndrechek Sep 25, 2026
dbd4661
Merge remote-tracking branch 'origin/feat/process-roles' into feat/mq…
EricAndrechek Sep 25, 2026
f5d8f48
test(mq): keep the streams directory occupied through a failed open
taitelee Sep 25, 2026
83ef8d0
feat(app): mq.backend selects embedded or external NATS
EricAndrechek Sep 25, 2026
2bfc9ee
fix(config): refuse credentials in mq.nats.urls; review fixes to docs
EricAndrechek Sep 25, 2026
18d47de
fix(mq): drain partitions a lower N leaves; quiet close
EricAndrechek Sep 25, 2026
6d42f5a
fix(mq): cover a publishing old-N process; name the stream delete
EricAndrechek Sep 25, 2026
d4f7450
feat(mq): coord leases on a NATS KV bucket
EricAndrechek Sep 25, 2026
a1aedc8
Merge remote-tracking branch 'origin/fix/config-yaml-zero' into integ…
EricAndrechek Sep 25, 2026
fd9b341
Merge remote-tracking branch 'origin/fix/discovery-retry-jitter' into…
EricAndrechek Sep 25, 2026
71fdc07
fix(mq): time lease step-down from the renewal sent; review fixes
EricAndrechek Sep 25, 2026
6339f47
Merge remote-tracking branch 'origin/fix/query-error-classes' into in…
EricAndrechek Sep 25, 2026
ddf0101
test(integration): open the outage test's queue on #612's broker API
EricAndrechek Sep 25, 2026
5bc481f
Merge remote-tracking branch 'origin/fix/mq-external-close-and-shrink…
EricAndrechek Sep 25, 2026
a08a6d6
fix(config): give the backend, roles and mq.nats defaults to defaults()
EricAndrechek Sep 25, 2026
b29929e
fix(mq): a per-term id in the lease value; lease tests off the unit b…
EricAndrechek Sep 25, 2026
0a64b9e
Merge remote-tracking branch 'origin/feat/cache-redis-wiring' into in…
EricAndrechek Sep 25, 2026
eaf69ef
fix(config): cache.redis defaults in defaults(); compress_min_bytes 0…
EricAndrechek Sep 25, 2026
aea078e
Merge remote-tracking branch 'origin/fix/mutation-pipes-uncached' int…
EricAndrechek Sep 25, 2026
00e165c
docs(deployment): the shared cache no longer caches write pipes
EricAndrechek Sep 25, 2026
8a9af4b
test(mq): the unit verifier cases skip the lease bucket they do not c…
EricAndrechek Sep 25, 2026
605a6c6
Merge remote-tracking branch 'origin/feat/dedupe-retention' into inte…
EricAndrechek Sep 25, 2026
d2ce776
test(api): an unreachable broker mid-window is uncertain, a 503
EricAndrechek Sep 25, 2026
6a69c03
fix(mq): ExternalNATS honours the idempotency key; conformance pins it
EricAndrechek Sep 25, 2026
aa8a83a
Merge remote-tracking branch 'origin/feat/dedupe-dynamodb-wiring' int…
EricAndrechek Sep 25, 2026
cc53a6a
fix(config): dedupe defaults in defaults(); pin the lease cap to mq's…
EricAndrechek Sep 25, 2026
09f2d3d
feat(mq): nats duplicate_window covers dedupe.lease; warn on short re…
EricAndrechek Sep 25, 2026
c696b2c
docs(deployment): say what nats and dynamodb share; retention reaches…
EricAndrechek Sep 25, 2026
5666252
docs(changelog): one list for #613's Unreleased entries
EricAndrechek Sep 25, 2026
b89ee1f
Merge branch 'feat/coord-nats-kv' into integration/distributed
EricAndrechek Sep 25, 2026
713b782
docs(deployment): a split by role has four shared backends to choose …
EricAndrechek Sep 25, 2026
9c1d252
test(e2e): api, ingest and sweeper in separate processes
EricAndrechek Sep 25, 2026
2a19d7e
fix(mq): refuse an embedded store it cannot create at once
EricAndrechek Sep 25, 2026
96e2112
test(mq): run the embedded broker's tests in parallel, without fsync
EricAndrechek Sep 25, 2026
bd059f2
test(app): boot without the embedded broker's fsync per write
EricAndrechek Sep 25, 2026
8776b4d
test(app): retry the DynamoDB table check sooner under test
EricAndrechek Sep 25, 2026
09d5c14
test(app): a 1ms topology_wait where the outcome cannot change
EricAndrechek Sep 25, 2026
e830f34
test(integration): name roles and backends in hand-built configs
EricAndrechek Sep 25, 2026
3d959b9
fix(mq): create the embedded store at 0700; changelog the fail-fast
EricAndrechek Sep 25, 2026
3c8a93b
docs(changelog): scope the fail-fast store entry to what mkdir catches
EricAndrechek Sep 25, 2026
e3d83c7
Merge branch 'perf/unit-test-budget' into integration/distributed
EricAndrechek Sep 25, 2026
d002146
refactor(app): move the shared-backend wiring out of wire.go
EricAndrechek Sep 25, 2026
f10b9de
refactor(app): wireCoord's nats case becomes wireNATSCoord in wire_na…
EricAndrechek Sep 25, 2026
468b4fe
test(integration): remove the binary C2 builds when the suite ends
EricAndrechek Sep 25, 2026
8c6057e
docs: claims the combined backends made false, round two
EricAndrechek Sep 25, 2026
d6b4d4b
docs: retention's refusal names 2m; a split needs a shared cache only…
EricAndrechek Sep 25, 2026
49a3b01
Merge remote-tracking branch 'origin/main' into integration/distributed
EricAndrechek Sep 25, 2026
6789325
test(integration): run own-stack tests in parallel; build C2's binary…
EricAndrechek Sep 25, 2026
e954536
test(integration): parallelize the last own-stack tests; name C2's bi…
EricAndrechek Sep 25, 2026
48dc144
docs(changelog): the integration timing measured on the tree it ships
EricAndrechek Sep 25, 2026
0e776c6
test(integration): wait for NATS's and Redis's host ports, not only t…
EricAndrechek Sep 25, 2026
079fa9a
test(integration): cite #613 in the roles tests' comments
EricAndrechek Sep 25, 2026
e08367c
test(cache): seed the breaker tests' fill outside their 100ms budget
EricAndrechek Sep 25, 2026
5f392e7
test(cache): name the setup timeout, not the breaker, as the failure
EricAndrechek Sep 25, 2026
1914eac
docs: neutral tags in the DynamoDB example; neutral benchmark wording
EricAndrechek Sep 25, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .github/labeler.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,11 @@
- any-glob-to-any-file:
- "internal/cache/**"

"area/coord":
- changed-files:
- any-glob-to-any-file:
- "internal/coord/**"

"area/dedupe":
- changed-files:
- any-glob-to-any-file:
Expand Down
43 changes: 43 additions & 0 deletions .testcoverage.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,15 @@ exclude:
# HTTP assertions). It is imported only from *_test.go files, never from
# production code, so there's nothing meaningful to cover.
- ^internal/testutil/
# internal/mq/mqtest/ is the Broker conformance suite: test code that
# lives outside *_test.go only so each backend's tests can import it.
- ^internal/mq/mqtest/
# internal/mq/natstest/ stands up NATS as an operator deploys it, for
# tests outside internal/mq; test code, like mqtest.
- ^internal/mq/natstest/
# The coord conformance suite: test helpers every Coordinator's tests
# run, imported only from *_test.go like testutil.
- ^internal/coord/coordtest/
- ^tests/
# scripts/ holds Go helpers (cov, orchestrator) that drive the build but
# aren't part of the shipped binary; they show up in `-coverpkg=./...`
Expand All @@ -73,3 +82,37 @@ exclude:
- ^internal/settings/
- ^cmd/wavehouse/validate\.go$
- ^cmd/wavehouse/bootstrap\.go$
# The external-NATS topology spec, verifier and manifest generator
# (and the `mq manifests` CLI) run against an operator's NATS, which
# the e2e stack (embedded broker) never has: unit territory, covered
# there by a fixture server. Excluding them keeps e2e at ~61%. The
# external broker is the same, covered by the integration suite.
- ^internal/mq/nats_topology\.go$
- ^internal/mq/nats_manifests\.go$
- ^internal/mq/subject_nats\.go$
- ^internal/mq/external\.go$
- ^internal/mq/lease\.go$
- ^cmd/wavehouse/mq\.go$
# The in-process cache backend: the e2e stack runs cache.backend=redis
# (#613), so the binary carries LocalCache and its version index but e2e
# never reaches them. The unit suite and the integration suite's main
# app (cache.backend=local) cover them; the merged total still counts them.
- ^internal/cache/(local|version_manager)\.go$
# The DynamoDB dedupe backend: the e2e binary runs Pebble dedupe, so
# this file measured 0% there and pulled e2e to 58.6%. The unit
# (fake API) and integration (dynamodb-local) suites cover it, and the
# merged total still counts it.
- ^internal/dedupe/dynamodb\.go$
# Their wiring, moved out of wire.go for this: the NATS queue and lease
# wiring and the retention warning under it, the DynamoDB dedupe
# wiring, and the ops-only listener of a process without the api role
# (the e2e binary runs every role). Integration and unit suites cover
# them; the merged total still counts them.
- ^internal/app/wire_(nats|dynamodb|ops)\.go$
unit:
# The external NATS broker's tests start a server per case, which the
# unit suite's 15s per package cannot hold: they are integration-tagged
# (make test-integration), and the merged total counts them.
- ^internal/mq/external\.go$
# The NATS KV leases, the same way.
- ^internal/mq/lease\.go$
48 changes: 26 additions & 22 deletions AGENTS.md

Large diffs are not rendered by default.

31 changes: 29 additions & 2 deletions CHANGELOG.md

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ Open a [feature request issue](https://github.com/Wave-RF/WaveHouse/issues/new?t

The pre-push hook (installed by `make tools`) blocks a push until the tree has been validated locally: a code change needs `make ci`, a docs/prose-only change needs only `make verify` (the same split CI makes). `make lint` / `make test` / `make build` are fast inner-loop subsets.

2. Write tests for new functionality. Unit tests go alongside the code in `internal/`. Integration tests go in `tests/` with the `//go:build integration` tag.
2. Write tests for new functionality. Unit tests go alongside the code in `internal/`. Integration tests go in `tests/` with the `//go:build integration` tag. The exception is a test that must import NATS, which only `internal/mq` may do; such tests go in `internal/mq/natsspike`, or in `internal/mq` itself with the `integration` tag when they need its internals (the external NATS broker's tests, which `make test-integration` selects by name).

3. Update documentation if your change affects:
- API endpoints → update `docs/src/content/docs/api.md`
Expand Down
8 changes: 7 additions & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -764,7 +764,13 @@ test-integration: go-mod-download ## Run Go integration tests + render coverage
@rm -rf $(COV_INT)/data && mkdir -p $(COV_INT)/data
@GOCOVERDIR="$(CURDIR)/$(COV_INT)/data" go tool gotestsum --format $(GOTESTSUM_FMT) -- \
-tags="integration $(TAGS)" -timeout 240s -coverpkg=./... -race -count=1 \
./tests/integration/... $(ARGS) \
./tests/integration/... ./internal/mq/natsspike/... ./internal/cache/... $(ARGS) \
-args -test.gocoverdir="$(CURDIR)/$(COV_INT)/data"
@# internal/mq's integration-tagged tests (the external NATS broker) run
@# alone: its untagged tests are the unit suite's.
@GOCOVERDIR="$(CURDIR)/$(COV_INT)/data" go tool gotestsum --format $(GOTESTSUM_FMT) -- \
-tags="integration $(TAGS)" -timeout 240s -coverpkg=./... -race -count=1 \
-run '^Test(ExternalNATS|NewNATS|NATSPermissions_Refuse|Leases)' ./internal/mq $(ARGS) \
-args -test.gocoverdir="$(CURDIR)/$(COV_INT)/data"
@if [ -z "$(COV_DEFER)" ]; then go run ./scripts/cov render integration; fi

Expand Down
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,8 +71,8 @@ ClickHouse is a phenomenal OLAP database, but pointing a frontend right at it le

If you're building user-facing analytics, WaveHouse is like **Supabase for ClickHouse**. Or an **open-source Tinybird** that pushes data to the frontend in real time over SSE, not just pull-based REST.

- **Ingest** — async durable WAL (embedded NATS JetStream), `200 OK` instantly, background batch-flush; schema-validated against `system.columns`; optional ID-based dedup (idempotent ingest); dead-letter queue for failed inserts.
- **Query** — in-process Ristretto cache + `singleflight` coalescing; type-safe structured query AST; Tinybird-style named pipes (parameterized SQL endpoints).
- **Ingest** — async durable WAL (embedded NATS JetStream), `200 OK` instantly, background batch-flush; schema-validated against `system.columns`; optional ID-based dedup (idempotent ingest); dead-letter queue for rows ClickHouse rejects (an unavailable ClickHouse is retried with backoff, not dead-lettered).
- **Query** — result cache (in-process Ristretto, or a Redis shared by every instance) + `singleflight` coalescing; type-safe structured query AST; Tinybird-style named pipes (parameterized SQL endpoints).
- **Real-time** — native SSE push, broadcast *before* the ClickHouse flush, with JetStream gap-fill for late/reconnecting clients.
- **Security** — Hasura-style per-table, per-role column + row policies with JWT claim templating, defined in the hot-reloadable settings directory.
- **Client** — `@wavehouse/sdk`: TypeScript client with query builder, live queries, streaming, and schema codegen; one runtime dependency (an SSE frame parser, ~1.4 KB gzipped).
Expand Down
38 changes: 38 additions & 0 deletions clients/ts/src/errors.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,44 @@ describe("parseErrorResponse", () => {
expect(e.retryable).toBe(true);
});

it("takes the server's code and retryable when it sends them", async () => {
const res = new Response(
JSON.stringify({
error: "Code: 62. Syntax error",
code: "clickhouse.rejected",
retryable: false,
}),
{ status: 400, statusText: "Bad Request" },
);
const e = await parseErrorResponse(res);
expect(e.code).toBe("clickhouse.rejected");
expect(e.retryable).toBe(false);
});

it("lets the server mark a 5xx not retryable", async () => {
const res = new Response(
JSON.stringify({
error: "Authentication failed",
code: "clickhouse.misconfigured",
retryable: false,
}),
{ status: 502, statusText: "Bad Gateway" },
);
const e = await parseErrorResponse(res);
expect(e.code).toBe("clickhouse.misconfigured");
expect(e.retryable).toBe(false);
});

it("ignores a non-string code and a non-boolean retryable", async () => {
const res = new Response(JSON.stringify({ error: "x", code: 123, retryable: "no" }), {
status: 500,
statusText: "Internal Server Error",
});
const e = await parseErrorResponse(res);
expect(e.code).toBe("HTTP_500");
expect(e.retryable).toBe(true);
});

it("marks 4xx as not retryable", async () => {
const res = new Response(JSON.stringify({ error: "forbidden" }), {
status: 403,
Expand Down
8 changes: 6 additions & 2 deletions clients/ts/src/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,14 @@ export async function parseErrorResponse(res: Response): Promise<WaveHouseError>
? body.message
: res.statusText;

const retryable = res.status === 503 || res.status >= 500;
// The server's own `code` and `retryable` win where it sends them (a
// failed ClickHouse query, for one); the status decides otherwise.
const code =
typeof body?.code === "string" && body.code !== "" ? body.code : `HTTP_${res.status}`;
const retryable = typeof body?.retryable === "boolean" ? body.retryable : res.status >= 500;
return {
status: res.status,
code: `HTTP_${res.status}`,
code,
message,
details: body,
retryable,
Expand Down
23 changes: 23 additions & 0 deletions clients/ts/src/http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,29 @@ describe("request", () => {
expect(result.error?.retryable).toBe(false);
});

it("does not retry a 5xx the server marks not retryable", async () => {
fetchSpy.mockResolvedValue(
new Response(
JSON.stringify({
error: "Authentication failed",
code: "clickhouse.misconfigured",
retryable: false,
}),
{ status: 502 },
),
);

const result = await request(makeCtx({ options: { maxRetries: 2 } }), {
method: "POST",
path: "/v1/query?table=clicks",
body: {},
});

expect(fetchSpy).toHaveBeenCalledOnce();
expect(result.error?.code).toBe("clickhouse.misconfigured");
expect(result.error?.retryable).toBe(false);
});

it("returns error for 500 without retry when maxRetries=0", async () => {
fetchSpy.mockResolvedValue(
new Response(JSON.stringify({ error: "internal" }), { status: 500 }),
Expand Down
22 changes: 14 additions & 8 deletions cmd/wavehouse/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,8 @@ func main() {
os.Exit(runValidate(os.Args[2:]))
case "bootstrap":
os.Exit(runBootstrap(os.Args[2:]))
case "mq":
os.Exit(runMQ(os.Args[2:], os.Stdout, os.Stderr))
case "version", "--version", "-v":
fmt.Printf("wavehouse %s (commit %s, built %s)\n", Version, GitCommit, BuildTime)
os.Exit(0)
Expand Down Expand Up @@ -138,6 +140,7 @@ func printUsage(w io.Writer) {
wavehouse start the server
wavehouse validate [dir] validate a settings directory (dir falls back to %[1]s)
wavehouse bootstrap [dir] write a starter settings directory, every key at its default (dir falls back to %[1]s)
wavehouse mq manifests print the nack resources for an external NATS JetStream
wavehouse health liveness self-probe against the local server (container HEALTHCHECK)
wavehouse version print version, commit, and build time
wavehouse help show this help
Expand Down Expand Up @@ -171,14 +174,17 @@ func run(ctx context.Context) int {
return 1
}

// data_dir must be writable before anything dials out, so the refusal
// (and, for the typical cause — a bind mount owned by root rather than
// UID 65532 — the remediation) lands at the top of the log rather than
// after ClickHouse discovery. NATS and Pebble still fail loud on their
// own if the directory changes underneath us.
if err := config.CheckDataDir(cfg.DataDir); err != nil {
logger.Error("check data_dir", "error", err)
return 1
// data_dir, when a selected backend keeps state there, must be writable
// before anything dials out, so the refusal (and, for the typical cause —
// a bind mount owned by root rather than UID 65532 — the remediation)
// lands at the top of the log rather than after ClickHouse discovery.
// NATS and Pebble still fail loud on their own if the directory changes
// underneath us.
if cfg.NeedsDataDir() {
if err := config.CheckDataDir(cfg.DataDir); err != nil {
logger.Error("check data_dir", "error", err)
return 1
}
}

a, err := app.New(ctx, app.Options{
Expand Down
83 changes: 83 additions & 0 deletions cmd/wavehouse/mq.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
package main

import (
"errors"
"flag"
"fmt"
"io"

"github.com/Wave-RF/WaveHouse/internal/mq"
)

// runMQ implements `wavehouse mq <command>`: tooling for the message queue.
// Exit codes: 0 ok, 1 failed, 2 usage.
func runMQ(args []string, stdout, stderr io.Writer) int {
usage := func(w io.Writer) {
_, _ = fmt.Fprint(w, `usage: wavehouse mq <command>

commands:
manifests print the nack resources for an external NATS JetStream
`)
}
if len(args) == 0 {
usage(stderr)
return 2
}
switch args[0] {
case "manifests":
return runMQManifests(args[1:], stdout, stderr)
case "help", "-h", "--help":
usage(stdout)
return 0
default:
_, _ = fmt.Fprintf(stderr, "wavehouse mq: unknown command %q\n\n", args[0])
usage(stderr)
return 2
}
}

// runMQManifests implements `wavehouse mq manifests`: print the nack
// Stream, Consumer and KeyValue resources for the topology WaveHouse checks at boot
// under mq.backend: nats, for the operator to apply.
func runMQManifests(args []string, stdout, stderr io.Writer) int {
fs := flag.NewFlagSet("mq manifests", flag.ContinueOnError)
fs.SetOutput(stderr)
partitions := fs.Int("partitions", mq.DefaultNATSPartitions, "number of ingest partition streams (mq.nats.partitions)")
prefix := fs.String("prefix", mq.DefaultNATSSubjectPrefix, "subject prefix (mq.nats.subject_prefix)")
replicas := fs.Int("replicas", 3, "replicas for every stream and the lease bucket")
bucket := fs.String("coord-bucket", "", "the lease KV bucket (coord.nats.bucket); empty is <prefix>_coord")
fs.Usage = func() {
_, _ = fmt.Fprint(fs.Output(), `usage: wavehouse mq manifests [--partitions N] [--prefix wh] [--replicas 3] [--coord-bucket B]

Print the nack (jetstream.nats.io/v1beta2) Stream, Consumer and KeyValue
resources for the JetStream topology WaveHouse needs under mq.backend: nats
and coord.backend: nats, as YAML for kubectl apply. WaveHouse never creates these itself; it checks them at boot.

`)
fs.PrintDefaults()
}
if err := fs.Parse(args); err != nil {
if errors.Is(err, flag.ErrHelp) {
return 0
}
return 2
}
if fs.NArg() > 0 {
_, _ = fmt.Fprintf(stderr, "wavehouse mq manifests: unexpected argument %q\n", fs.Arg(0))
fs.Usage()
return 2
}
if *replicas < 1 {
_, _ = fmt.Fprintf(stderr, "wavehouse mq manifests: --replicas must be at least 1\n")
return 2
}
err := mq.WriteNATSManifests(stdout, mq.NATSManifestOptions{
Topology: mq.NATSTopology{Prefix: *prefix, Partitions: *partitions, CoordBucket: *bucket},
Replicas: *replicas,
})
if err != nil {
_, _ = fmt.Fprintf(stderr, "wavehouse mq manifests: %v\n", err)
return 1
}
return 0
}
51 changes: 51 additions & 0 deletions cmd/wavehouse/mq_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
package main

import (
"bytes"
"os"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// The shipped manifests are the generator's output for four partitions.
// Regenerate with: go run ./cmd/wavehouse mq manifests --partitions 4 > deployments/nats/jetstream.yaml
func TestRunMQManifests_MatchesShipped(t *testing.T) {
want, err := os.ReadFile("../../deployments/nats/jetstream.yaml")
require.NoError(t, err)
var out, errOut bytes.Buffer
require.Equal(t, 0, runMQ([]string{"manifests", "--partitions", "4"}, &out, &errOut), errOut.String())
assert.Equal(t, string(want), out.String(), "deployments/nats/jetstream.yaml is stale; regenerate it")
}

func TestRunMQ_ExitCodes(t *testing.T) {
cases := map[string]struct {
args []string
code int
}{
"no command": {nil, 2},
"unknown command": {[]string{"frobnicate"}, 2},
"help": {[]string{"help"}, 0},
"manifests help": {[]string{"manifests", "-h"}, 0},
"stray argument": {[]string{"manifests", "extra"}, 2},
"unknown flag": {[]string{"manifests", "--nope"}, 2},
"zero replicas": {[]string{"manifests", "--replicas", "0"}, 2},
"bad prefix": {[]string{"manifests", "--prefix", "a.b"}, 1},
"bad partitions": {[]string{"manifests", "--partitions", "-1"}, 1},
"bad coord bucket": {[]string{"manifests", "--coord-bucket", "a.b"}, 1},
"defaults generate": {[]string{"manifests"}, 0},
}
for name, tc := range cases {
t.Run(name, func(t *testing.T) {
var out, errOut bytes.Buffer
assert.Equal(t, tc.code, runMQ(tc.args, &out, &errOut), errOut.String())
})
}
}

func TestRunMQManifests_NamesTheLeaseBucket(t *testing.T) {
var out, errOut bytes.Buffer
require.Equal(t, 0, runMQ([]string{"manifests", "--prefix", "acme", "--coord-bucket", "acme_leases"}, &out, &errOut), errOut.String())
assert.Contains(t, out.String(), "kind: KeyValue\nmetadata:\n name: acme-coord\nspec:\n bucket: acme_leases\n")
}
2 changes: 1 addition & 1 deletion cmd/wavehouse/validate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ func writeSettingsDir(t *testing.T, policies string) string {
"roles.json": `{"roles": ["public"]}`,
"policies.json": policies,
"pipes.json": `{}`,
"config.json": `{"clickhouse": {"addr": "localhost:9000", "http_port": 8123, "http_scheme": "http", "database": "default", "username": "default", "query_timeout": 30, "tls": {"enabled": false, "ca_file": "", "cert_file": "", "key_file": "", "insecure_skip_verify": false, "server_name": ""}, "headers": {}, "max_open_conns": 10, "max_idle_conns": 5}, "auth": {"jwks_url": "", "role_claim": "role"}, "dedupe": {"enabled": false, "id_field": "event_id", "require_id": false}, "dlq": {"enabled": true}, "query": {"default_max_rows": 10000, "timestamp_bucket_seconds": 60}, "schema": {"refresh_interval": 60}, "stream": {"keepalive_interval": 30, "keepalive_buckets": 3, "gap_window_minutes": 15}, "mq": {"max_bytes_gb": 1}, "cors": {"allowed_origins": ["*"]}}`,
"config.json": `{"clickhouse": {"addr": "localhost:9000", "http_port": 8123, "http_scheme": "http", "database": "default", "username": "default", "query_timeout": 30, "tls": {"enabled": false, "ca_file": "", "cert_file": "", "key_file": "", "insecure_skip_verify": false, "server_name": ""}, "headers": {}, "max_open_conns": 10, "max_idle_conns": 5}, "auth": {"jwks_url": "", "role_claim": "role"}, "dedupe": {"enabled": false, "id_field": "event_id", "require_id": false, "retention": "0"}, "dlq": {"enabled": true}, "query": {"default_max_rows": 10000, "timestamp_bucket_seconds": 60}, "schema": {"refresh_interval": 60}, "stream": {"keepalive_interval": 30, "keepalive_buckets": 3, "gap_window_minutes": 15}, "mq": {"max_bytes_gb": 1}, "cors": {"allowed_origins": ["*"]}}`,
}
for name, content := range files {
require.NoError(t, os.WriteFile(filepath.Join(dir, name), []byte(content), 0o600))
Expand Down
Loading
Loading