fix(test): remove the broker store after late consumer writes land - #667
EricAndrechek wants to merge 2 commits into
Conversation
The embedded NATS server writes each durable consumer's state (obs/<consumer>/o.dat, via a temp file renamed into place) from a goroutine neither Shutdown nor WaitForShutdown joins; its consumer store waits for it at close only while state is unwritten, and for at most 100ms. A write already under way therefore lands after EmbeddedNATS.Close returns, and t.TempDir's one-shot RemoveAll met the late entry as "directory not empty". The worker's own teardown was already joined; the late writer is the server's. storedir.New(t) is now the store directory for every test that puts a broker on disk (testutil.NewEmbeddedMQ, the mq and mqtest suites, the app tests' data_dir, three integration tests). Its cleanup runs after the broker's Close and removes the store 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. It replaces two sleep-and-retry copies in the mq tests. A rename-aside fence was tried first and does not work: the flusher's rename resolves its paths before the fence and lands after it. TestStartIngestWorker_StopFunc_RespectsShutdownDeadline now joins the worker its deadline abandons before the broker closes. Under six parallel race-enabled test processes on 14 CPU hogs, the dispatch-loop and StartIngestWorker tests failed 78 of 1800 runs before and 0 of 1800 after. With the helper instrumented to log its retries, 102 of 1800 teardowns needed a second removal and every one succeeded on it. Closes #442 Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
TestRun_BootsAndStopsOnCancel boots the full process through run(), embedded broker included, but still gave it a bare t.TempDir() as data_dir, so its durables' late state writes could fail the removal the same way #442's did. It now uses storedir.New(t); run() returns inside the test body, so the store is removed after the broker has closed. The TestRun_RefusesToBoot cases fail before the broker opens and keep t.TempDir(). The comment on EmbeddedNATS.Close now names the production side of the same gap, #665: the flusher's write can land after Close returns, or never if the process exits first, leaving the durable's previous ack state on disk. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
|
Important Review skippedAuto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Advanced Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
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://dae3b883-wavehouse-docs.wave-rf.workers.dev
|
Takes #667's fix for the #442 teardown flake so this branch's make ci can pass; it drops out of the diff once #667 lands on main. CHANGELOG: both new entries kept. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Takes the base's AGENTS.md tree line (naming DynamoDB beside Pebble here) and #667's teardown-flake fix. No conflicts. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Takes the base's AGENTS.md tree line and #667's teardown-flake fix. CHANGELOG: both new entries kept. The stale-duplicate-window mq test now takes its store from storedir.New, the helper #667 replaced storeDir with. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Takes the base's log-level test, its storedir moves and #667's teardown-flake fix (via #625). No conflicts. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Takes the base's AGENTS.md tree line and #667's teardown-flake fix. No conflicts. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
#625) Part of #613. This PR carries the whole remote-dedupe stack: the reserve/commit contract, windowed ingest, retention, and the DynamoDB backend with its boot wiring. #629, #633, #628 and #635 were reviewed on their own and merged into this branch, and #667's test fix came in with them. ## Summary - **Reserve, Commit, Release** (fixes #390, #222, #370). `Deduplicator.CheckAndMark` is replaced by a two-phase contract that every backend implements: - `Reserve(ctx, keys, lease)` claims each key atomically and answers `Claimed`, `Duplicate` or `InFlight` for it. It is all-or-nothing on error. - `Commit(ctx, claims, retention)` marks the claims as seen. `Release(ctx, claims)` gives them back, matched by token. - A claim that is neither committed nor released lapses after its lease, so a request that dies mid-publish never strands an id. - Keys are scoped per tenant and table in the readable `keyenc` format `<tenant>/<escaped table>/<escaped id>`, for example `acme/clicks/evt%2D123`. An id whose escaped form is over 1,024 bytes is stored as `#<sha256 hex>`. The same id in two tables is now two ids (#222). - Two concurrent requests with one id now publish once (#390). An explicit `null` id counts as missing (#370). - Pebble keeps pending claims in memory in 64 locked shards beside the instance, and writes each Commit in one batch with one fsync. - The conformance suite `internal/dedupe/dedupetest` runs every backend through the same cases. - **Windowed ingest** (fixes #384). Ingest runs in windows of up to 256 records. Each window makes one `Reserve`, publishes in record order, then makes one `Commit`. - A deduped record is published under a `Nats-Msg-Id` idempotency key derived from its tenant, table and id. The embedded ingest stream sets its duplicate window to 2 minutes explicitly. - Only a definite publish failure (the queue refused it) releases the claim. After an uncertain failure, the claim lapses with the lease, and a retry is dropped by the stream's duplicate window. The event is stored once and never lost. - A dedupe store that cannot answer (`dedupe.ErrUnavailable`, now "dedupe store unavailable") answers `503 {"error":"dedupe store unavailable"}` with `Retry-After: 5`. It used to answer `500 dedupe failed`. - On Pebble, a 1,000-record batch now costs 4 fsyncs instead of 1,000. - **Per-table retention** (fixes #220). `dedupe.retention` in `config.json`, with a per-table override in `dedupe.tables.<table>.retention`, sets how long a committed id stays a duplicate. The default is `"0"`, which keeps ids forever, so an existing `config.json` needs no change. - A finite retention below 2 minutes (the queue's duplicate window) is refused, not clamped. - A background sweep on the Pebble instance deletes expired ids and the old-format keys. It runs a minute after open and then hourly, and it never deletes a key that was committed again after the sweep read it. `wavehouse_dedupe_swept_keys_total{reason}` counts what it deletes. - **DynamoDB backend**, selected by `dedupe.backend: dynamodb`. Every tenant and every process share one table, so an id ingested through one pod is a duplicate through every other. - `Reserve` is a conditional `PutItem` per key. `Commit` is `BatchWriteItem` with retries. `Release` is a conditional `DeleteItem`. Expiry is the native TTL attribute `ex`, and correctness never waits on TTL. - Throttling, timeouts and connection failures wrap `ErrUnavailable` and answer `503`. A circuit breaker short-circuits `Reserve` for a second after five unavailable claims in a row. - New boot keys: `dedupe.lease` (the lease is now configurable, 30 s by default, at most 59 s with the embedded queue), `dedupe.reserve_concurrency`, and the `dedupe.dynamodb.*` block. `create_table` is refused unless `endpoint` is set, so WaveHouse never creates a table in AWS. - **Boot rule:** boot checks the table whether or not any tenant has dedupe on. A misconfigured table (missing, the wrong key schema, access denied) refuses boot only with a flat settings directory whose tenant has dedupe on. In every other case, including transient failures, nested directories, and no tenant with dedupe on yet, the process boots, and every tenant with dedupe on fails closed with the `503`. The check is retried in the background and again right after every reload. - **Caller-cancel fix** (addresses #648). When a caller cancels mid-`Reserve`, the puts not yet sent are skipped. A put already sent runs to its answer before it is released. Only its own call deadline can cut it off, and then it holds its key at most until the lease ends, as a crashed request's claim does. - **Test teardown** (absorbs #667). Tests that start the embedded broker no longer fail in `t.TempDir` cleanup when the broker's consumer-state flusher writes after `Close`. The new `internal/testutil/storedir` retries the removal. ## Behaviour and compatibility notes - **Old-format dedupe keys are swept, not migrated.** An id seen before the upgrade is accepted once more after it. The retention sweep deletes the old keys on its first pass. Nothing released depends on them. - A dedupe backend that cannot answer returns `503` + `Retry-After: 5` where it used to return `500`. The SDK already retries a `503`. - A mid-body read error or a prepare failure now drops the open window unpublished. Before, the records ahead of it were published. - The in-flight `503` sends the lease as `Retry-After`. That is 30 s by default, as before. ## Known follow-ups - #660: row-by-row isolation silently drops identical rows on a deduplicating table. - #665: a durable's last ack can land after `Close`, or never if the process exits. - #668: this PR does its three items: `config.embeddedDuplicateWindow` is pinned to `mq.EmbeddedDuplicateWindow` by a test, the `reserve_concurrency` wording is updated, and the lease rule is stated once. Close it by hand after this lands. - #651: cross-region dedupe on DynamoDB MRSC needs a sweeper for lapsed claims. - #652: accept events durably while the dedupe backend is down. ## Tests - **Conformance:** `dedupetest.Run` runs against Pebble twice (on an injected clock and on the real clock) and against `amazon/dynamodb-local:3.3.1`. It covers claim, duplicate and in-flight, release then re-claim, lease lapse, retention expiry, a 64-way concurrent Reserve, a Reserve racing a Commit, key isolation per tenant and table, hashed long ids, and all-or-nothing on a mid-call failure. - **Ingest** (`internal/api/ingest_window_test.go`, `ingest_retention_test.go`): - window boundaries and publish failures at chosen records; - the `503` for an unavailable store; - the #384 scenario end to end over the real broker and Pebble; - the uncertain-publish retry, mutation-checked against a missing idempotency key; - retention reaching `Commit` and changing on reload, including mid-window. - **Pebble sweep** (`internal/dedupe/sweep_test.go`): chunk boundaries, expired and old-format keys, and a Commit racing a chunk. Each case is mutation-checked. - **DynamoDB:** unit tests against a fake API cover error classification, Reserve cleanup, a Reserve cancelled by its caller leaving nothing claimed, a retried put keeping its own claim, Commit retries, the breaker, and `Check`. Integration tests against dynamodb-local cover the conformance suite, 32 clients racing one id, throttling, an unreachable endpoint, TTL and expiry, and two `app.New` instances sharing seen ids through one table. - **Boot wiring** (`internal/app/dedupe_dynamodb_test.go`, `internal/config/backends_test.go`): the table check in both directory shapes, the background retry, reloads that make no table call, and every new config key and refusal. - **Pinning tests:** the retention floor is at least the queue's duplicate window, and `config.embeddedDuplicateWindow` equals `mq.EmbeddedDuplicateWindow`. - `make ci` passes on the merged stack: every coverage gate passed. Unit 94.0%, integration 52.4%, e2e 60.2% (60% floor), Go total 95.0%. Fixes #390. Fixes #222. Fixes #370. Fixes #384. Fixes #220. Closes #442. Closes #648. Part of #613. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq --------- Co-authored-by: taitelee <taitelee@umich.edu> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Summary
Tests that start the embedded broker could fail in
t.TempDircleanup withdirectory not emptyunderobs/<consumer>/, after their assertions had passed. The late writer is nats-server's own consumer-state flusher:consumerFileStore.flushLoopruns on a goroutine that neitherShutdownnorWaitForShutdownjoins, andconsumerFileStore.Stopwaits for it only while state is still dirty. So a write already under way (o.dat.tmp, then a rename ontoo.dat) can land afterEmbeddedNATS.Closereturns, and the one-shotRemoveAllthen finds the new entry. The worker's own goroutines were already joined; the issue's "double ack" theory is out of date.internal/testutil/storedir:storedir.New(t)makes the store directory. Its cleanup runs after the broker'sCloseand retriesRemoveAllonly onENOTEMPTY, at most 32 times, with no sleep. It terminates because each late write adds at most two entries and nothing can be added once its directory is gone. It is its own package becauseinternal/mq's tests cannot importtestutil.testutil.NewEmbeddedMQ, themqandmqtestsuites (replacing two 50×20 ms sleep-retry copies), the app tests,cmd/wavehouse's boot test and three integration tests.TestStartIngestWorker_StopFunc_RespectsShutdownDeadlinenow releases its blocked insert and joins the worker before the broker closes.EmbeddedNATS.Closerecords the production side, tracked as bug(mq): a durable's last ack can land after Close, or never if the process exits #665: the same write can land afterClosereturns, or never if the process exits first. This PR does not change production behaviour.Test plan
-run 'TestDispatchLoop|TestStartIngestWorker'): 78 of 1,800 runs failed before, 0 of 1,800 after. An instrumented run needed a second removal in 102 of 1,800 teardowns, and each succeeded on it.make cigreen locally on 3235598 (the first run failed one e2e SDK test on afetch failedsettings reload, a path this branch does not touch; the re-run passed).go test -raceforinternal/ingest,internal/mq,internal/app,internal/api,internal/testutil,cmd/wavehouse.Related Issues
Closes #442. Related: #665 (the production late write), #383.
🤖 Generated with Claude Code
https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds