Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -446,7 +446,7 @@ internal/query/ → Structured query AST + SQL builder
internal/settings/ → Settings directory (validate, adopted snapshot + reload, watcher, embedded seed)
internal/stream/ → SSE fan-out (event Hub: project once per role, Subscriber outbound queue, Bucket fan-out, keepalive Heartbeater wheel)
internal/tenant/ → Tenant id (type, grammar, reserved default, request header name)
internal/testutil/ → Shared test helpers (mocks, JWT + schema helpers; logtest/ captures or silences the default logger)
internal/testutil/ → Shared test helpers (mocks, JWT + schema helpers; logtest/ captures or silences the default logger; storedir/ is the embedded broker's store directory in tests, removed once late consumer-state writes land)
tests/ → Integration & E2E tests
tests/integration/ → Go integration tests (//go:build integration; ClickHouse testcontainer)
tests/e2e/ → E2E test stack (scripts/orchestrator boots a ClickHouse testcontainer + the wavehouse-cov binary)
Expand Down
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),

### Fixed

- **Tests that start the embedded broker no longer fail removing its store after passing** (`internal/testutil/storedir` (new, + tests), `internal/testutil/testutil.go`, `internal/mq/embedded.go` (comment), `internal/mq/{embedded_test,mqtest/embedded_test}.go`, `internal/ingest/worker_test.go`, `internal/app/{app,roles}_test.go`, `cmd/wavehouse/main_test.go`, `tests/integration/{ingest_outage,query_errors,tenants}_test.go`, `AGENTS.md`, `docs/src/content/docs/development.md`): [#442](https://github.com/Wave-RF/WaveHouse/issues/442). The NATS server writes each durable consumer's state (`obs/<consumer>/o.dat`, through a temporary file renamed into place) from a goroutine that neither `Shutdown` nor `WaitForShutdown` joins, and its consumer store waits for that goroutine at close only when state is still unwritten, for at most 100ms — so a write already under way lands after `EmbeddedNATS.Close` returns, and `t.TempDir`'s one-shot `RemoveAll` met the late entry as `directory not empty`. Under parallel test processes it failed about 4% of the ingest worker tests (78 of 1,800 runs). Every store a test puts on disk now comes from `storedir.New(t)`, whose cleanup — after the broker's `Close` — removes it again whenever a directory was refilled between being read and being removed: each late write adds at most two entries and none once its directory is gone, so the removal ends without a timer (0 of 1,800 under the same load). It replaces two sleep-and-retry copies in the `internal/mq` tests. `TestStartIngestWorker_StopFunc_RespectsShutdownDeadline` also joins the worker its deadline abandons before the broker closes, rather than leaving it to ack on a closed connection.
- **A failed ClickHouse query answers by what went wrong, not a flat `500`/`502`** (`internal/api/ch_errors.go` (new, + tests), `internal/api/{errors,query,structured_query,pipes,schema,ch_settings}.go`, `internal/chconn/errclass.go` (`HTTPStatus` exported), `clients/ts/src/errors.ts` (+ tests), `tests/integration/query_errors_test.go` (new), `tests/integration/query_limits_test.go`, `internal/app/app_test.go`, `tests/e2e/sdk/{admin,query}.test.ts`, `AGENTS.md`, `docs/src/content/docs/{api,architecture}.md`, `docs/src/content/docs/{access-control,configuration}.mdx`, `docs/src/content/docs/sdk/{reference.md,index.mdx}`): fixes [#403](https://github.com/Wave-RF/WaveHouse/issues/403) and [#271](https://github.com/Wave-RF/WaveHouse/issues/271), part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). ClickHouse answers a syntax error, a missing grant and an overloaded server alike with HTTP `500`, so `/v1/ops/query` turned a bad statement into a `502` and `/v1/query` and pipes into a `500` the SDK retried. All three now class the failure with `chconn.Classify` through one helper, `writeCHError`: a statement ClickHouse refused is `400 clickhouse.rejected`; a query over a rows/bytes limit, the role's own memory cap, or its time cap where that is no longer than `query_timeout` is `400 clickhouse.limit_exceeded`; `ACCESS_DENIED` is `403 clickhouse.access_denied`; credentials, user or database refused, or a redirect or `4xx` with no exception code from whatever fronts ClickHouse, is `502 clickhouse.misconfigured`; ClickHouse down, unreachable or overloaded is `503 clickhouse.unavailable` with `Retry-After: 5`; a failure with no verdict stays `500` (`502` on the proxy) as `clickhouse.unknown`. The error envelope gains `code` and `retryable` next to `error` on these responses — additive. A role with `max_execution_time` now queries with no context deadline and a cancel two seconds past the cap instead: clickhouse-go overwrote the cap's `max_execution_time` with deadline+5s for any deadline over 1s, so an overrun came back as a bare deadline, indistinguishable from waiting for a pooled connection; ClickHouse now enforces the cap itself and reports `TIMEOUT_EXCEEDED`. `POST /v1/ops/schema/refresh` against an unreachable ClickHouse is a `503` with `Retry-After` instead of a `500`. **SDK:** `WaveHouseError.code` and `retryable` now take the server's `code`/`retryable` when the body has them (`HTTP_<status>` and "5xx retries" otherwise), so a rejected query is `clickhouse.rejected` rather than `HTTP_500`, and is not retried.
- **An unavailable ClickHouse is retried with backoff instead of dead-lettering every row** (`internal/chconn/errclass.go` (new, + tests), `internal/ingest/{worker,backoff}.go` (`backoff.go` new, + tests), `internal/mq/{mq,embedded}.go`, `internal/testutil/mocks.go`, `tests/integration/ingest_outage_test.go` (new), `AGENTS.md`, `README.md`, `docs/src/content/docs/{ingest-pipeline,architecture,api,deployment,why-wavehouse}.md`, `docs/src/content/docs/{settings-directory,index,access-control}.mdx`): workstream A of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). 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 twice and parked the whole batch on the DLQ. `chconn.Classify` now classes the failure first — `Rejected` (any ClickHouse exception code outside the availability and credential lists: the server read the row and refused it), `Unavailable` (connection refused/reset, timeouts, `TOO_MANY_SIMULTANEOUS_QUERIES`, `SERVER_OVERLOADED`, `MEMORY_LIMIT_EXCEEDED`, `TOO_MANY_PARTS`, `READONLY`, `TABLE_IS_READ_ONLY`, `KEEPER_EXCEPTION`, …), `Denied` (`AUTHENTICATION_FAILED`, `ACCESS_DENIED`, …) or `Unknown` (no code, no recognizable transport failure). Only `Rejected` is isolated and dead-lettered as before, and a multi-row batch refused with `TOO_MANY_PARTS` or `MEMORY_LIMIT_EXCEEDED` is split row by row first (`chconn.Splittable`), because a batch spanning too many partitions or too much memory can fail where each of its rows inserts; every other class hands the batch back to the queue with a delayed nak (`mq.Message.NakWithDelay`, new) under a jittered 1 s → 30 s backoff shared by every table on the same ClickHouse pool (a failure of one table — read-only, too many parts or mutations, a grant missing on it, `chconn.TableScoped` — backs off that table alone), which turns rows away without a request while it runs and probes once per window, and ClickHouse going away mid-isolation stops isolation and retries the rows it had not settled. Counted by the new `wavehouse_ingest_retries_total{table, reason}`; logged at `WARN` when an outage starts and at most every 30 s during it. A long outage now shows as a growing ingest stream and, at `mq.max_bytes_gb`, ingest `503`s — not as a full DLQ; a lasting failure of one table holds back its tenant's other tables once its waiting rows reach `maxAckPending`. Retried rows come back out of arrival order, which matters only to a `ReplacingMergeTree` without a version column or a `CollapsingMergeTree`.
- **Schema discovery's retry loop jitters its backoff** (`internal/discovery/discovery.go` (+ tests), `internal/app/wire.go`, `internal/api/errors.go`, `AGENTS.md`, `docs/src/content/docs/{architecture,api,deployment}.md`): `RetryRefresh` slept exactly `2s * 2^n` capped at 60s, so instances retrying against one recovering ClickHouse fired in lockstep, every 60s on the same second. Each sleep is now drawn uniformly from below the backoff (full jitter), spreading the retries over the whole window and halving the mean wait — so a failing tenant's retries, their log lines and `wavehouse_schema_refresh_failures_total` come about twice as often ([#141](https://github.com/Wave-RF/WaveHouse/issues/141)).
Expand Down
3 changes: 2 additions & 1 deletion cmd/wavehouse/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (

"github.com/Wave-RF/WaveHouse/internal/config"
"github.com/Wave-RF/WaveHouse/internal/settings"
"github.com/Wave-RF/WaveHouse/internal/testutil/storedir"
)

// run reads the whole boot config from the environment here (no config
Expand Down Expand Up @@ -78,7 +79,7 @@ func seedSettings(t *testing.T) string {
func TestRun_BootsAndStopsOnCancel(t *testing.T) {
hermeticEnv(t)
t.Setenv(config.EnvSettingsDir, seedSettings(t))
t.Setenv("WH_DATA_DIR", t.TempDir())
t.Setenv("WH_DATA_DIR", storedir.New(t))
_, port, err := net.SplitHostPort(closedAddr(t))
require.NoError(t, err)
t.Setenv("WH_SERVER_PORT", port)
Expand Down
2 changes: 1 addition & 1 deletion docs/src/content/docs/development.md
Original file line number Diff line number Diff line change
Expand Up @@ -347,7 +347,7 @@ Each test target writes `covdata` to `tmp/coverage/<suite>/data/`, renders a tex
- **Unit tests** live beside the code they test (e.g., `internal/discovery/discovery_test.go`). They use mocks or embedded NATS (in-process, no Docker needed).
- **Integration tests** use the `//go:build integration` build tag. `TestMain` starts one ClickHouse testcontainer and boots the production wiring against it through `app.New` (embedded NATS, ingest worker, sweeper, hub, the API server on a random loopback port); tests reach it via `env(t)` and create their own tables. DLQ tests use `assert.Eventually` with a 30-second timeout for the 5-second ingest worker batch window.

Shared test utilities live in `internal/testutil/`. The packages log through `slog.Default()`, so tests reach log output through `internal/testutil/logtest`: `logtest.Silence()` in a package's `TestMain` discards it, and `logtest.Capture(t, level)` routes it to a buffer for a test that asserts on log lines — such a test must not call `t.Parallel()`, because the default logger is process-wide.
Shared test utilities live in `internal/testutil/`. The packages log through `slog.Default()`, so tests reach log output through `internal/testutil/logtest`: `logtest.Silence()` in a package's `TestMain` discards it, and `logtest.Capture(t, level)` routes it to a buffer for a test that asserts on log lines — such a test must not call `t.Parallel()`, because the default logger is process-wide. A test that starts the embedded broker keeps its store in `internal/testutil/storedir`'s `storedir.New(t)` rather than a bare `t.TempDir()` (`testutil.NewEmbeddedMQ` does): the NATS server can finish writing a consumer's state after `Close` returns, which fails `t.TempDir`'s one-shot removal, and `storedir` removes the store again until those writes have landed ([#442](https://github.com/Wave-RF/WaveHouse/issues/442)).

### Adding New Tests

Expand Down
3 changes: 2 additions & 1 deletion internal/app/app_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ import (
"github.com/Wave-RF/WaveHouse/internal/tenant"
"github.com/Wave-RF/WaveHouse/internal/testutil"
"github.com/Wave-RF/WaveHouse/internal/testutil/logtest"
"github.com/Wave-RF/WaveHouse/internal/testutil/storedir"
)

// None of these tests run in parallel: New installs a process-wide default
Expand Down Expand Up @@ -94,7 +95,7 @@ func writeSettings(t *testing.T, patch map[string]any) string {
func testConfig(t *testing.T, settingsDir string) *config.Config {
t.Helper()
return &config.Config{
DataDir: t.TempDir(),
DataDir: storedir.New(t),
Server: config.Server{Port: closedPort(t), ShutdownTimeout: 2},
MQ: config.MQ{Backend: config.MQEmbedded},
Cache: config.Cache{Backend: config.CacheLocal, L1MaxCost: 1 << 20},
Expand Down
3 changes: 2 additions & 1 deletion internal/app/roles_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"github.com/Wave-RF/WaveHouse/internal/config"
"github.com/Wave-RF/WaveHouse/internal/coord"
"github.com/Wave-RF/WaveHouse/internal/settings"
"github.com/Wave-RF/WaveHouse/internal/testutil/storedir"
)

// Each role wires its own components and nothing else; the settings registry,
Expand Down Expand Up @@ -109,7 +110,7 @@ func TestNew_OpsOnlyRouter(t *testing.T) {
sweeperCfg := *cfg
sweeperCfg.Roles = []config.Role{config.RoleSweeper}
// Its own store: full's embedded JetStream is still open on cfg.DataDir.
sweeperCfg.DataDir = t.TempDir()
sweeperCfg.DataDir = storedir.New(t)
a := newApp(t, &sweeperCfg, Options{})

for _, path := range []string{"/livez", "/readyz", "/healthz", "/version"} {
Expand Down
11 changes: 7 additions & 4 deletions internal/ingest/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -254,10 +254,7 @@ func TestStartIngestWorker_StopFunc_RespectsShutdownDeadline(t *testing.T) {
<-release
w.WriteHeader(http.StatusOK)
}))
t.Cleanup(func() {
close(release)
chSrv.Close()
})
t.Cleanup(chSrv.Close)
u, _ := url.Parse(chSrv.URL)
host, port, _ := net.SplitHostPort(u.Host)

Expand All @@ -269,6 +266,12 @@ func TestStartIngestWorker_StopFunc_RespectsShutdownDeadline(t *testing.T) {
return chconn.Target{URL: fmt.Sprintf("http://%s:%s", host, port), Username: "u", Password: "p", Database: "db"}
}, nil)
require.NoError(t, err)
// The deadline below abandons the worker, not its insert: finish the insert
// and join the worker before the broker closes under its ack.
t.Cleanup(func() {
close(release)
assert.NoError(t, stopFn(context.Background()))
})

// Publish so there's an in-flight insert blocking on `release`.
err = emb.Publish(ctx, mq.Topic{Tenant: tenant.Default, Table: "events"}, makeEnvelope(t, "events", "", map[string]any{"id": 1}))
Expand Down
5 changes: 4 additions & 1 deletion internal/mq/embedded.go
Original file line number Diff line number Diff line change
Expand Up @@ -1086,7 +1086,10 @@ func (e *EmbeddedNATS) Close() error {
// Owning the lifecycle (NoSigs, #287) means waiting it out: without this,
// run()'s remaining defers unwind while JetStream is still tearing down
// and the process can exit mid-shutdown (as-if-crashed stream state).
// Milliseconds for an in-process server.
// Milliseconds for an in-process server. It does not join a durable's
// state flusher: a write under way can land after Close returns, or never
// if the process exits first, leaving the durable's previous ack state on
// disk (#665).
e.server.WaitForShutdown()
return nil
}
Loading
Loading