diff --git a/AGENTS.md b/AGENTS.md index f6580a22..dcc13cc9 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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) diff --git a/CHANGELOG.md b/CHANGELOG.md index ccdd58bb..0be38b97 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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//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_` 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)). diff --git a/cmd/wavehouse/main_test.go b/cmd/wavehouse/main_test.go index d0e8404c..ff471474 100644 --- a/cmd/wavehouse/main_test.go +++ b/cmd/wavehouse/main_test.go @@ -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 @@ -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) diff --git a/docs/src/content/docs/development.md b/docs/src/content/docs/development.md index 6ee7ceed..d16835f5 100644 --- a/docs/src/content/docs/development.md +++ b/docs/src/content/docs/development.md @@ -347,7 +347,7 @@ Each test target writes `covdata` to `tmp/coverage//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 diff --git a/internal/app/app_test.go b/internal/app/app_test.go index 6a8bc449..7af8c8f3 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -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 @@ -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}, diff --git a/internal/app/roles_test.go b/internal/app/roles_test.go index 48fb0fa1..c50eb9a3 100644 --- a/internal/app/roles_test.go +++ b/internal/app/roles_test.go @@ -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, @@ -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"} { diff --git a/internal/ingest/worker_test.go b/internal/ingest/worker_test.go index f07500bc..7b884093 100644 --- a/internal/ingest/worker_test.go +++ b/internal/ingest/worker_test.go @@ -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) @@ -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})) diff --git a/internal/mq/embedded.go b/internal/mq/embedded.go index 15ca485f..c4928b3c 100644 --- a/internal/mq/embedded.go +++ b/internal/mq/embedded.go @@ -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 } diff --git a/internal/mq/embedded_test.go b/internal/mq/embedded_test.go index 55e7fec4..e172c903 100644 --- a/internal/mq/embedded_test.go +++ b/internal/mq/embedded_test.go @@ -10,6 +10,7 @@ import ( "time" "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" "github.com/stretchr/testify/assert" @@ -19,26 +20,6 @@ import ( // testBudget is the byte budget newTestEmbedded opens each queue at. const testBudget = 64 << 20 -// storeDir is a temporary directory for a broker's store whose removal -// retries briefly: a consumer's state file can land after Close has returned, -// which fails t.TempDir's one-shot RemoveAll (#442). The retrying cleanup runs -// first (cleanups are LIFO), leaving t.TempDir an empty directory to remove. -func storeDir(t *testing.T) string { - t.Helper() - dir := filepath.Join(t.TempDir(), "store") - t.Cleanup(func() { - var err error - for range 50 { - if err = os.RemoveAll(dir); err == nil { - return - } - time.Sleep(20 * time.Millisecond) - } - t.Errorf("remove %s: %v", dir, err) - }) - return dir -} - // openEmbedded starts an EmbeddedNATS over dir, closed by the test framework. func openEmbedded(t *testing.T, dir string) *EmbeddedNATS { t.Helper() @@ -53,7 +34,7 @@ func openEmbedded(t *testing.T, dir string) *EmbeddedNATS { // at testBudget. func newTestEmbedded(t *testing.T, tenants ...tenant.ID) *EmbeddedNATS { t.Helper() - e := openEmbedded(t, storeDir(t)) + e := openEmbedded(t, storedir.New(t)) if len(tenants) == 0 { tenants = []tenant.ID{tenant.Default} } @@ -167,7 +148,7 @@ func TestEmbeddedNATS_PublishHeaders(t *testing.T) { // at a tenth of it, dropping its oldest when full. No other tenant gets one. func TestEmbeddedNATS_SetMaxBytes_OpensTheTenantsQueue(t *testing.T) { t.Parallel() - e := openEmbedded(t, storeDir(t)) + e := openEmbedded(t, storedir.New(t)) assert.Zero(t, e.MaxBytes("acme"), "no budget applied yet") require.NoError(t, e.SetMaxBytes(t.Context(), "acme", testBudget)) @@ -342,7 +323,7 @@ func TestEmbeddedNATS_DefaultLogger(t *testing.T) { t.Parallel() // NewEmbedded without a logger should not panic — it falls back to the // default slog logger. - e, err := NewEmbedded(storeDir(t)) + e, err := NewEmbedded(storedir.New(t)) require.NoError(t, err) t.Cleanup(func() { _ = e.Close() }) } @@ -458,7 +439,7 @@ func TestEmbeddedNATS_SetMaxBytes_IngestFailureChangesNothing(t *testing.T) { // however recently a publish tried. func TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen(t *testing.T) { t.Parallel() - dir := storeDir(t) + dir := storedir.New(t) // The dead-letter stream is the first of the pair to open. A failed open // removes what was in the way, so the obstacle is put back before each // attempt meant to fail. @@ -509,7 +490,7 @@ func TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen(t *testing.T) { // resize and reload takes. Once the window has passed, a publish tries again. func TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen(t *testing.T) { t.Parallel() - dir := storeDir(t) + dir := storedir.New(t) block := filepath.Join(dir, "jetstream", "$G", "streams", dlqStreamName("acme")) obstruct := func() { t.Helper() @@ -574,7 +555,7 @@ func TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen(t *testing.T) { // joined, so its row reaches them rather than a stream nobody reads. func TestEmbeddedNATS_Publish_OpensAQueueItsOpenGaveUpOn(t *testing.T) { t.Parallel() - dir := storeDir(t) + dir := storedir.New(t) block := filepath.Join(dir, "jetstream", "$G", "streams", ingestStreamName("acme")) require.NoError(t, os.MkdirAll(filepath.Dir(block), 0o750)) require.NoError(t, os.WriteFile(block, nil, 0o600)) @@ -618,7 +599,7 @@ func TestEmbeddedNATS_SetMaxBytes_UndoRestoresTheIngestStreamsCap(t *testing.T) t.Parallel() ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() - dir := storeDir(t) + dir := storedir.New(t) first, err := NewEmbedded(dir) require.NoError(t, err) require.NoError(t, first.SetMaxBytes(ctx, "acme", 8<<20)) @@ -641,7 +622,7 @@ func TestEmbeddedNATS_SetMaxBytes_UndoRestoresTheIngestStreamsCap(t *testing.T) // queue itself is open, so SetMaxBytes succeeds. func TestEmbeddedNATS_Consume_ReportsAQueueItCannotJoin(t *testing.T) { t.Parallel() - e := openEmbedded(t, storeDir(t)) + e := openEmbedded(t, storedir.New(t)) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() // A durable name the client refuses: with no queue yet, nothing checks it. @@ -666,7 +647,7 @@ func TestEmbeddedNATS_Consume_ReportsAQueueItCannotJoin(t *testing.T) { // stream keeps what it holds, capped at that, and every row survives. func TestEmbeddedNATS_SetMaxBytes_NeverShrinksTheDeadLetterQueueBelowWhatItHolds(t *testing.T) { t.Parallel() - e := openEmbedded(t, storeDir(t)) + e := openEmbedded(t, storedir.New(t)) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() require.NoError(t, e.SetMaxBytes(ctx, "acme", 10<<20)) @@ -881,7 +862,7 @@ func TestEmbeddedNATS_DeadLetter_ReopensAMissingQueue(t *testing.T) { func TestEmbeddedNATS_Publish_QueueFull(t *testing.T) { t.Parallel() - e := openEmbedded(t, storeDir(t)) + e := openEmbedded(t, storedir.New(t)) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() require.NoError(t, e.SetMaxBytes(ctx, "acme", 4<<10)) @@ -1364,7 +1345,7 @@ func TestNewEmbedded_AStoreItCannotCreateFailsAtOnce(t *testing.T) { // so no tenant's queue could open beside them. func TestNewEmbedded_DeletesTheStreamsAnEarlierBuildShared(t *testing.T) { t.Parallel() - dir := storeDir(t) + dir := storedir.New(t) old, err := NewEmbedded(dir) require.NoError(t, err) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) @@ -1395,7 +1376,7 @@ func TestNewEmbedded_ASplitPairIsAppliedAgainAtBoot(t *testing.T) { t.Parallel() ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() - dir := storeDir(t) + dir := storedir.New(t) first, err := NewEmbedded(dir) require.NoError(t, err) for _, id := range []tenant.ID{"split", "gone", "guarded"} { @@ -1433,7 +1414,7 @@ func TestNewEmbedded_ASplitPairIsAppliedAgainAtBoot(t *testing.T) { // so what such a tenant had queued still reaches the worker. func TestNewEmbedded_TakesStockOfTheQueuesOnDisk(t *testing.T) { t.Parallel() - dir := storeDir(t) + dir := storedir.New(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() first, err := NewEmbedded(dir) @@ -1484,7 +1465,7 @@ func TestEmbeddedNATS_ADurableOnDiskIsReusedAcrossARestart(t *testing.T) { } { t.Run(tt.name, func(t *testing.T) { t.Parallel() - dir := storeDir(t) + dir := storedir.New(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() topic := Topic{Tenant: "acme", Table: "t"} diff --git a/internal/mq/mqtest/embedded_test.go b/internal/mq/mqtest/embedded_test.go index 98b619d9..f3be090a 100644 --- a/internal/mq/mqtest/embedded_test.go +++ b/internal/mq/mqtest/embedded_test.go @@ -3,21 +3,19 @@ package mqtest_test import ( - "os" - "path/filepath" "testing" - "time" "github.com/Wave-RF/WaveHouse/internal/mq" "github.com/Wave-RF/WaveHouse/internal/mq/mqtest" "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" "github.com/stretchr/testify/require" ) func TestEmbeddedNATS_Conformance(t *testing.T) { mqtest.Run(t, mqtest.Harness{ New: func(t *testing.T) mq.Broker { - e, err := mq.NewEmbedded(storeDir(t)) + e, err := mq.NewEmbedded(storedir.New(t)) require.NoError(t, err) t.Cleanup(func() { _ = e.Close() }) for _, id := range []tenant.ID{mqtest.Acme, mqtest.Globex} { @@ -55,22 +53,3 @@ func TestEmbeddedNATS_Conformance(t *testing.T) { }, }) } - -// storeDir is a temporary store directory whose removal retries briefly: under -// parallel load a consumer's state file can land after Close has returned, -// which fails t.TempDir's one-shot RemoveAll. The retrying cleanup runs first -// (cleanups are LIFO), leaving t.TempDir an empty directory to remove. -func storeDir(t *testing.T) string { - dir := filepath.Join(t.TempDir(), "store") - var err error - t.Cleanup(func() { - for range 50 { - if err = os.RemoveAll(dir); err == nil { - return - } - time.Sleep(20 * time.Millisecond) - } - t.Errorf("remove %s: %v", dir, err) - }) - return dir -} diff --git a/internal/testutil/storedir/storedir.go b/internal/testutil/storedir/storedir.go new file mode 100644 index 00000000..aae61863 --- /dev/null +++ b/internal/testutil/storedir/storedir.go @@ -0,0 +1,47 @@ +// Package storedir gives a test a directory for the embedded message broker's +// store. It imports nothing from the repository, so internal/mq's own tests can +// use it. +package storedir + +import ( + "errors" + "os" + "path/filepath" + "syscall" + "testing" +) + +// maxRemovals bounds the store's removal: a consumer-state write still under way +// when the broker closes adds at most two entries after it — its temporary file, +// then the rename into place — so a removal is refilled at most twice per +// durable. Past this many, something is writing that Close did not stop. +const maxRemovals = 32 + +// New returns an empty directory under t.TempDir for a broker's store, removed +// once the broker is closed: New's cleanup runs after the test's own, Close +// among them (cleanups run last-in, first-out). +// +// A closed broker's store is not yet quiescent: the embedded NATS server writes +// each durable consumer's state from a goroutine its Shutdown does not join, so +// a write under way can land after Close has returned — failing t.TempDir's +// one-shot RemoveAll with "directory not empty" (#442). There is nothing to +// wait on, so the removal is tried again whenever a directory was refilled +// between reading and removing it. That ends without a clock: those writes add a +// bounded number of entries, and none once their directory is gone. +func New(t testing.TB) string { + t.Helper() + dir := filepath.Join(t.TempDir(), "store") + if err := os.Mkdir(dir, 0o700); err != nil { + t.Fatalf("store directory: %v", err) + } + t.Cleanup(func() { + err := os.RemoveAll(dir) + for i := 1; i < maxRemovals && errors.Is(err, syscall.ENOTEMPTY); i++ { + err = os.RemoveAll(dir) + } + if err != nil { + t.Errorf("remove the store: %v", err) + } + }) + return dir +} diff --git a/internal/testutil/storedir/storedir_test.go b/internal/testutil/storedir/storedir_test.go new file mode 100644 index 00000000..6720a534 --- /dev/null +++ b/internal/testutil/storedir/storedir_test.go @@ -0,0 +1,28 @@ +package storedir + +import ( + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestNew_RemovedAfterTheTestsOwnCleanups(t *testing.T) { + var dir string + t.Run("store", func(t *testing.T) { + dir = New(t) + entries, err := os.ReadDir(dir) + require.NoError(t, err) + assert.Empty(t, entries, "the store starts empty") + // Registered after New, as a broker's Close is, so it runs first and + // what it writes is removed with the rest. + t.Cleanup(func() { + obs := filepath.Join(dir, "jetstream", "obs") + require.NoError(t, os.MkdirAll(obs, 0o700)) + require.NoError(t, os.WriteFile(filepath.Join(obs, "o.dat"), nil, 0o600)) + }) + }) + assert.NoDirExists(t, dir) +} diff --git a/internal/testutil/testutil.go b/internal/testutil/testutil.go index db119686..8b3cefb4 100644 --- a/internal/testutil/testutil.go +++ b/internal/testutil/testutil.go @@ -16,6 +16,7 @@ import ( "github.com/Wave-RF/WaveHouse/internal/discovery" "github.com/Wave-RF/WaveHouse/internal/mq" "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" ) // NewTestSchemaRegistry creates a SchemaRegistry pre-loaded with the given @@ -40,13 +41,13 @@ func NewTestSchemaRegistry(t testing.TB, tables []*discovery.TableSchema) *disco // hardcoding the same literal twice. const TestServerVersion = "24.8.1.1" -// NewEmbeddedMQ starts the embedded broker over a temporary directory, closed -// by the test framework, with a queue open for each of tenants — +// NewEmbeddedMQ starts the embedded broker over a storedir.New directory, +// closed by the test framework, with a queue open for each of tenants — // tenant.Default when none is named — at maxBytes: a tenant has a queue once // its budget is applied, as the wiring does for every tenant it serves. func NewEmbeddedMQ(t testing.TB, maxBytes int64, tenants ...tenant.ID) *mq.EmbeddedNATS { t.Helper() - emb, err := mq.NewEmbedded(t.TempDir()) + emb, err := mq.NewEmbedded(storedir.New(t)) require.NoError(t, err) t.Cleanup(func() { _ = emb.Close() }) if len(tenants) == 0 { diff --git a/tests/integration/ingest_outage_test.go b/tests/integration/ingest_outage_test.go index 752d84b2..53faa2d0 100644 --- a/tests/integration/ingest_outage_test.go +++ b/tests/integration/ingest_outage_test.go @@ -18,6 +18,7 @@ import ( "github.com/Wave-RF/WaveHouse/internal/mq" "github.com/Wave-RF/WaveHouse/internal/tenant" "github.com/Wave-RF/WaveHouse/internal/testutil" + "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" ) // TestIngest_ClickHouseOutage_RetriedNotDeadLettered stops a real ClickHouse @@ -39,7 +40,7 @@ func TestIngest_ClickHouseOutage_RetriedNotDeadLettered(t *testing.T) { const table = "outage_events" require.NoError(t, ch.conn.Exec(ctx, "CREATE TABLE "+table+" (id UInt32) ENGINE = MergeTree ORDER BY id")) - broker, err := mq.NewEmbedded(t.TempDir()) + broker, err := mq.NewEmbedded(storedir.New(t)) require.NoError(t, err) t.Cleanup(func() { _ = broker.Close() }) require.NoError(t, broker.SetMaxBytes(ctx, tenant.Default, 64<<20)) diff --git a/tests/integration/query_errors_test.go b/tests/integration/query_errors_test.go index 9eccd851..d29660ef 100644 --- a/tests/integration/query_errors_test.go +++ b/tests/integration/query_errors_test.go @@ -22,6 +22,7 @@ import ( "github.com/Wave-RF/WaveHouse/internal/app" "github.com/Wave-RF/WaveHouse/internal/chconn" "github.com/Wave-RF/WaveHouse/internal/config" + "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" ) // queryError is the error envelope a failed ClickHouse query answers with. @@ -115,7 +116,7 @@ func TestQueryErrors_ClickHouseDown(t *testing.T) { ln, err := lc.Listen(ctx, "tcp", "127.0.0.1:0") require.NoError(t, err) cfg := &config.Config{ - DataDir: t.TempDir(), + DataDir: storedir.New(t), Server: config.Server{ShutdownTimeout: 10}, ClickHouse: config.ClickHouse{Password: testCHPassword}, MQ: config.MQ{Backend: config.MQEmbedded}, diff --git a/tests/integration/tenants_test.go b/tests/integration/tenants_test.go index 16d888ba..e3c756b6 100644 --- a/tests/integration/tenants_test.go +++ b/tests/integration/tenants_test.go @@ -20,6 +20,7 @@ import ( "github.com/Wave-RF/WaveHouse/internal/app" "github.com/Wave-RF/WaveHouse/internal/config" "github.com/Wave-RF/WaveHouse/internal/tenant" + "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" ) // TestNestedDirectory_PerTenantPoolsAndDiscovery boots the real wiring over @@ -58,7 +59,7 @@ func TestNestedDirectory_PerTenantPoolsAndDiscovery(t *testing.T) { ln, err := lc.Listen(ctx, "tcp", "127.0.0.1:0") require.NoError(t, err) cfg := &config.Config{ - DataDir: t.TempDir(), + DataDir: storedir.New(t), Server: config.Server{ShutdownTimeout: 10}, ClickHouse: config.ClickHouse{Password: testCHPassword}, Auth: config.Auth{OperatorKey: operatorKey},