Skip to content

fix(test): remove the broker store after late consumer writes land - #667

Closed
EricAndrechek wants to merge 2 commits into
mainfrom
fix/ingest-teardown-race
Closed

EricAndrechek wants to merge 2 commits into
mainfrom
fix/ingest-teardown-race

Conversation

@EricAndrechek

Copy link
Copy Markdown
Member

Summary

Tests that start the embedded broker could fail in t.TempDir cleanup with directory not empty under obs/<consumer>/, after their assertions had passed. The late writer is nats-server's own consumer-state flusher: consumerFileStore.flushLoop runs on a goroutine that neither Shutdown nor WaitForShutdown joins, and consumerFileStore.Stop waits for it only while state is still dirty. So a write already under way (o.dat.tmp, then a rename onto o.dat) can land after EmbeddedNATS.Close returns, and the one-shot RemoveAll then finds the new entry. The worker's own goroutines were already joined; the issue's "double ack" theory is out of date.

  • New internal/testutil/storedir: storedir.New(t) makes the store directory. Its cleanup runs after the broker's Close and retries RemoveAll only on ENOTEMPTY, 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 because internal/mq's tests cannot import testutil.
  • Every test that puts a broker store on disk now uses it: testutil.NewEmbeddedMQ, the mq and mqtest suites (replacing two 50×20 ms sleep-retry copies), the app tests, cmd/wavehouse's boot test and three integration tests.
  • TestStartIngestWorker_StopFunc_RespectsShutdownDeadline now releases its blocked insert and joins the worker before the broker closes.
  • The comment on EmbeddedNATS.Close records 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 after Close returns, or never if the process exits first. This PR does not change production behaviour.

Test plan

  • Measured under six parallel race-enabled test processes with CPU contention (-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 ci green locally on 3235598 (the first run failed one e2e SDK test on a fetch failed settings reload, a path this branch does not touch; the re-run passed).
  • go test -race for internal/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

EricAndrechek and others added 2 commits September 26, 2026 05:54
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
@coderabbitai

coderabbitai Bot commented Sep 26, 2026

Copy link
Copy Markdown

Important

Review skipped

Auto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the .coderabbit.yaml file in this repository. To trigger a single review, invoke the @coderabbitai review command.

⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Advanced

Run ID: c8858d74-fe93-482f-9883-1ec05c1bd2e5

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added documentation Improvements or additions to documentation go Pull requests that update go code area/ingest Ingest pipeline (Bento, batching, DLQ) area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/app Process wiring (internal/app): component build, run, release labels Sep 26, 2026
@github-actions

Copy link
Copy Markdown

📚 Docs preview is live → https://dae3b883-wavehouse-docs.wave-rf.workers.dev

  • Commit — 3235598: test(cmd): keep the boot test's broker store in storedir
  • Author — @EricAndrechek, Claude Opus 5.5 (1M context)
  • Committed — 2026-09-26 08:28 (UTC-04:00)
  • Deployed — 2026-09-26 08:41 EDT

@github-code-quality

Copy link
Copy Markdown
Contributor

Code Coverage Overview

Languages: Go

Go

The overall line coverage in commit 3235598 in the fix/ingest-teardown-... branch remains at 94%, unchanged from commit 5004cd2 in the main branch.

EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
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
EricAndrechek added a commit that referenced this pull request Sep 26, 2026
#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>
@EricAndrechek

Copy link
Copy Markdown
Member Author

Closing: this fix landed on main through #625 (squash 4ff5074), which also closed #442.

@github-project-automation github-project-automation Bot moved this from Backlog to Done in WaveHouse Task Board Sep 26, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area/app Process wiring (internal/app): component build, run, release area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/ingest Ingest pipeline (Bento, batching, DLQ) documentation Improvements or additions to documentation go Pull requests that update go code

Projects

Status: Done

Development

Successfully merging this pull request may close these issues.

fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main

1 participant