fix(ingest): windowed reserve/publish/commit; 503 when dedupe down - #629
Merged
Merged
Conversation
…ailable Each window of up to 256 records is prepared, reserved in one dedupe call, published in order and committed in one call. Deduped records are published under an idempotency key (Nats-Msg-Id), and each tenant's ingest stream keeps an explicit two-minute duplicate window, so a publish whose outcome is unknown leaves its claim to lapse and the retry's copy is dropped by the queue. A dedupe store that cannot answer is a 503 with Retry-After: 5. Fixes #384. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
The idempotency key drops a retry only within two minutes of the first publish, and a 200 with a failed commit is counted, not committed: the docs now say so. The uncertain-publish test uses a 2s lease so a stall cannot lapse the claim early, and the mq test pins the explicit duplicate window rather than the server's matching default. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
…name Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
|
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 |
dedupe.retention (required, "0" = forever) and its per-table override are read per record and passed to Commit. Validation refuses a finite retention below the queue's two-minute duplicate window. The embedded Pebble store deletes expired keys and the version-0 keys in an hourly background sweep, counted by wavehouse_dedupe_swept_keys_total. Part of #613. Refs #220. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Adds a sweep hook so a Commit can race into the gap between a chunk's read and delete; the test fails (3/3) with commitMu removed. Docs: the sweep follows the shared instance, reads (not deletes) 1,024 keys per chunk, and "0" is the one unitless retention. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
This was referenced Sep 25, 2026
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
dedupe.retention was required in every config.json. It is now optional: missing means "0" (forever) at the tenant level, and a table override without it inherits the tenant's, as for the other override fields. A present value is validated as before — unparseable, negative, or finite below the queue's duplicate window is refused — so an existing directory needs no change to upgrade. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
# Conflicts: # CHANGELOG.md # docs/src/content/docs/architecture.md # internal/dedupe/key.go
# Conflicts: # CHANGELOG.md # docs/src/content/docs/architecture.md # internal/settings/validate_test.go
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
# Conflicts: # AGENTS.md # CHANGELOG.md # docs/src/content/docs/architecture.md # docs/src/content/docs/settings-directory.mdx # internal/dedupe/key.go
# Conflicts: # CHANGELOG.md # docs/src/content/docs/architecture.md
# Conflicts: # AGENTS.md # CHANGELOG.md
# Conflicts: # CHANGELOG.md
The version-0 read test planted only a tenant-NUL-id key, which the text layout never looks up, so it passed whatever Reserve made of a legacy value. It now also plants a bare legacy id that spells a current key and checks it reads as absent and is overwritten by the commit. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
A sweep chunk held commitMu while it iterated to its 1,024th live key. Pebble skips point tombstones inside Next, so a chunk that started over a run of them (a tenant's version-0 block deleted by the first pass, or a table whose ids had all expired) held every Commit, and with it every deduped ingest response, for the whole run, on every pass until the run was compacted. The chunk now reads without the lock and collects the keys it would delete, then takes commitMu only to re-read each one and delete those still expired or version-0, without fsync as before. A Commit waits for at most 1,024 point reads and one unsynced batch, however many tombstones lie between the keys. The mid-chunk test splits in two, one per gap a Commit can land in: after the unlocked read, where the re-read keeps the key, and after the re-read, where the lock holds the Commit off until the delete is done. Each fails when its guard is removed. A new test races Commits against a chunk that starts over 300,000 tombstones: under the race detector the old chunk kept one waiting 241-278 ms, the new one at most 11 ms. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
dedupe.retention is the one config.json key the binary defaults (missing means "0", forever). The ingest handler's dedupe comment, the seed's doc comment, a registry test comment and the settings-directory CHANGELOG entry still said there were none. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Brings in main's mq.ErrUnavailable -> 503 mapping (previously on the per-record publish path this branch replaced), the mq.Broker conformance suite (internal/mq/mqtest), and the commit-failed metric rename to wavehouse_ingest_dedupe_commit_failed_total. Resolves the merge's semantic fallout in the same commit, since the merge does not build/pass tests without them: - publishFailed treats mq.ErrUnavailable as an uncertain failure (the same claim handling as any non-ErrQueueFull error: commit before k, release after k, leave k to lapse) but answers 503 instead of the generic 500. Retry-After is the dedupe lease, rounded up to whole seconds, when k held a claim left to lapse; otherwise the flat 5 seconds main's per-record path used, since there is no lapsing claim to wait out. - TestIngest_Dedup_FailedReleaseKeepsThePublishError now fails its publish with mq.ErrQueueFull (a definite failure) instead of a plain error, since an uncertain failure no longer releases its own record's claim in a single-record window (issue #384's fix) -- there would be nothing to release. Expects 503, not 500. - Renamed lingering wavehouse_dedupe_commit_failed_total references in durability.md, CHANGELOG.md and settings-directory.mdx to the renamed metric. docs/api.md, architecture.md and sdk/reference.md updated to match. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The in-flight 503 for an uncertain publish sends the FULL dedupe lease as Retry-After, so an obedient client's retry can land up to ~2*lease after the original Reserve -- not just one lease later -- and a claim's expiry can itself round up by up to a second on some backends (DynamoDB, for one). "lease <= window" understates what the embedded queue's duplicate window actually has to cover. Pin the real invariant in TestIngest_DedupeLeaseFitsTheDuplicateWindow (2*dedupe.DefaultLease + time.Second <= mq.EmbeddedDuplicateWindow), correct the embedded.go comment that said "must not exceed", and reword the same claim in durability.md, api.md, architecture.md and CHANGELOG.md. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The takeStock fix landed with the previous commit (embedded.go), since both are about the same staleness question; this adds its coverage. TestNewEmbedded_TakeStockRefreshesAStaleDuplicateWindow reopens a store whose ingest stream was left with a Duplicates window other than EmbeddedDuplicateWindow, then calls SetMaxBytes with the SAME budget as before and checks the window is brought forward -- the path TestEmbeddedNATS_Publish_IdempotencyKeyDropsARepeat did not cover (its stream is created directly, never recorded by takeStock, so SetMaxBytes's same-budget early return never applies to it). Corrected that test's comment to say so, pointing at the new one. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
mqtest.Run had no case for mq.WithIdempotencyKey, yet the uncertain- publish rule in the ingest handler (a claim left to lapse after a publish whose outcome is unknown, relying on the queue to drop the retry's second copy) depends on every Broker honouring it -- not just the embedded one, which already had its own duplicate-key test. Adds IdempotencyKeyDropsARepeat: a publish repeated under one key is stored once and both calls return nil; a different key is stored separately, checked via ReplaySince. Wired into the suite's case list unconditionally (no Caps flag), since every Broker must honour it. Kept internal/mq's own TestEmbeddedNATS_Publish_IdempotencyKeyDropsARepeat: it additionally pins the embedded-specific mechanics of a stream opened with one duplicate window picking up EmbeddedDuplicateWindow on its next SetMaxBytes, which is outside the generic Broker contract. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
…Reserve internal/app/wire.go's wirePebbleDedupe comment still described a store that fails to open on reload as a 500 "dedupe failed" -- that mapping moved to 503 "dedupe store unavailable" with Retry-After: 5 earlier in this branch's history (internal/api/reserve's dedupe.ErrUnavailable case). Update the comment to match. reserve()'s generic dd.Reserve error branch logged every failure at ERROR, including one caused by the request's own context ending (the client went away, or its deadline passed) while Reserve was in flight -- not a backend problem, and not worth paging an operator over. Check ctx.Err() and log at Debug instead when it is set; a real backend failure still logs ERROR. The response status is unchanged (moot: nothing is listening for it). TestIngest_Dedup_ReserveError_ContextEnded_NotLoggedAsError pins the log level via logtest; TestIngest_Dedup_ReserveError continues to pin the real-backend-failure ERROR case. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Absorb feat/dedupe-reserve's newer commits: a no-op fast path in Managed.Apply, the Reserve-on-error contract reworded for network backends, a start barrier in the commit-race conformance case, and development.md's tree line — all taken as-is. Resolved every docs/test conflict to this PR's semantics, which already distinguish a definite publish failure (mq.ErrQueueFull: releases the record's claim and the rest of the window) from an uncertain one (anything else, including mq.ErrUnavailable: commits records before it, leaves its own claim to lapse with the lease, and releases the rest) — not the base's per-record model, where every publish failure released the id. CHANGELOG.md and api.md keep this PR's rows describing the fixed behavior (they already superseded the base's forward-references to #629); the base's own dedupe-reserve CHANGELOG bullet picks up its "#629 closes that" wording. In ingest_test.go, TestIngest_Dedup_FailedPublishReleasesTheID stays scoped to the one failure that actually releases (ErrQueueFull) so its name keeps telling the truth; the base's new mq.ErrUnavailable coverage moves into TestIngest_Dedup_UncertainPublishLeavesTheClaim (now table-driven) with the expectations this PR's contract actually produces: the claim lapses rather than releasing, the response is 503, and Retry-After is the lease rounded up to whole seconds (30s default) rather than the base's flat 5s. go build, go vet -tags integration, and go test -race across internal/api, internal/mq, internal/dedupe and internal/app all pass. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Absorb feat/dedupe-windowed-ingest's newer history, including its own merge of feat/dedupe-reserve (step 1 of this pass) and origin/main (#615 leases, #618 backend selection, #622 process roles, #627 ClickHouse error classing, and the rest through #623). Resolved six conflicts, combining both sides' facts rather than picking one: - internal/settings/settings.go: kept MinDedupeRetention (this PR's 2-minute floor, tied to the embedded queue's duplicate window) next to windowed-ingest's reworded DLQConfig comment ("a row ClickHouse still rejects", reflecting the outage-retry split from #613). - AGENTS.md: kept this PR's dedupe/ package-inventory line (the sweep detail) alongside windowed-ingest's updated config/ and new coord/ lines pulled in from main. - CHANGELOG.md, docs/architecture.md, docs/durability.md, docs/settings-directory.mdx: superseded this branch's now-stale copies (old `evt%2D123` key spelling from before refactor/keyenc kept '-'; the metric name `wavehouse_dedupe_commit_failed_total` before it gained the `ingest_` prefix; the simpler "fits inside" duplicate-window wording before the `2×lease+1s` invariant was pinned) with windowed-ingest's current, code-matching text, then folded this PR's retention-specific additions back in: the Upgrade note now says the pre-#222 keys are "deleted by the retention sweep" instead of "nothing removing them yet", and durability.md's duplicate-window paragraph keeps its closing sentence tying `dedupe.retention`'s 2-minute floor to that same window. internal/api/ingest.go, ingest_test.go, ingest_window_test.go, app/wire.go, app/app_test.go, dedupe/embedded_test.go and settings/settings.go (the rest of it) auto-merged with no textual conflict; verified by reading the result rather than trusting that: every Commit path windowed-ingest added (commitClaims before the failing record in publishFailed, and after a clean window in ingestWindow) already passes commitClaims the full pendingRecord slice, and commitClaims groups by each record's resolved retention and issues one Commit per distinct value — so both PRs' Commit-path changes compose correctly. The sweep's commitMu (embedded.go/sweep.go) and Managed.Apply's fast path (managed.go) touch disjoint locks and did not need reconciling. go build, go vet -tags integration (whole repo), and go test -race across internal/dedupe, internal/settings, internal/api and internal/mq all pass (re-run with -count=1 after one flaky timing assertion in TestEmbedded_SweepChunkOverTombstonesDoesNotHoldCommits — a 100ms budget the sweep raced past once under parallel-package load — passed clean on every subsequent run, including three solo runs and a full fresh run; unrelated to this merge). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… ones TestEmbedded_SweepChunkOverTombstonesDoesNotHoldCommits (e52b336, this PR) raced a Commit against a sweep chunk over 300k tombstones and asserted the slowest attempt stayed under a 100ms wall-clock budget. That budget itself flaked under contention: 313ms measured when the whole internal/dedupe/... race suite ran, since test-unit runs packages in parallel under -race in make ci and a wall-clock bound moves with scheduler load. Replace it with two assertions load cannot move: a new sweepTouchHook (embedded.go, wired into deleteSweepable in sweep.go) tallies how many candidates the locked phase re-reads — asserted <= 1, the actual candidate count in this scenario, proving the locked phase's work is bounded by the candidates the unlocked read found sweepable, not by however many tombstones it silently stepped over inside Pebble to get there. A direct, non-blocking commitMu.TryLock() from within the existing sweepScanHook (fired after the unlocked read, before deleteSweepable's lock) proves the lock was free at that point without racing a goroutine or a clock at all: on the same goroutine that just did the read, TryLock fails instead of blocking if a regression left the lock held, so a pass is a direct proof, not an inference from a race won in time. Verified the new test actually catches the regression it guards against: temporarily wrapped sweepChunk's whole body (including the unlocked read) in e.commitMu.Lock()/Unlock() and dropped deleteSweepable's own lock to avoid a self-deadlock, simulating the pre-fix "sweep under the lock" design — the test failed exactly on the new TryLock assertion ("commitMu was free right after the unlocked read of 300k tombstones"), then reverted; `git diff` on sweep.go before adding the touch hook showed only the hook call, confirming a clean revert. GOTOOLCHAIN=go1.26.6 go test -race -count=20 -run Sweep ./internal/dedupe/ passed 20/20 while a concurrent `go test -race ./internal/...` ran in the background for contention (both processes exited 0). go build ./... and go vet -tags integration ./... are clean. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
TestIngest_Dedup_ReserveError_ContextEnded_NotLoggedAsError's comment claimed TestIngest_Dedup_ReserveError pins the real-backend-failure case staying ERROR — it doesn't; that test only asserts status, body and publish count, never the log. Two mutations passed the whole internal/api suite as a result: `if ctx.Err() != nil` -> `if true` (collapses every Reserve failure onto the DEBUG branch), and the client-gone line's DebugContext -> WarnContext (both leave a bare `"level":"DEBUG"` check satisfied by an unrelated "debug: span started for ingest" line the package logs on every request). Replace the single-case test with one non-parallel, two-case TestIngest_Dedup_ReserveError_LogLevel: a live context asserts the line `"level":"ERROR","msg":"dedupe reserve failed"`; a cancelled one asserts `"level":"DEBUG","msg":"dedupe reserve failed: request context ended"`. Matching level and msg as one adjacent substring (the exact order slog's JSON handler emits them in) ties the level to the specific line rather than to the buffer as a whole, so neither mutation above can hide behind the unrelated DEBUG line. Verified both mutations now fail the new test (live-context case fails under `if true`; context-ended case fails under WarnContext), then reverted each — `git diff` on ingest.go showed no residue before the two doc fixes below were applied. Also, [MAY]: requestAbort.RetryAfter's inline comment repeated a narrower, now-stale cause list (missing the dedupe-store-unavailable 503) that duplicates and drifts from the requestAbort doc comment above it, which already owns that list — trimmed to the field's own job. mq.ErrUnavailable's doc said the API answers it with "a short Retry-After"; publishFailed sends the full dedupe lease (rounded up to whole seconds) when the failing record held a claim, and only falls back to a flat few seconds otherwise — reworded to say so. GOTOOLCHAIN=go1.26.6 go test -race ./internal/api/... ./internal/mq/... passes. 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
The caveat about a released id letting a retry publish a genuine second copy sat only on the mq.ErrUnavailable 503 rows, which the same text says the embedded broker never returns — so as written it described a case the default deployment can't hit. With the embedded broker the uncertain publish is the 500 (a client disconnect after the broker had already stored the message), so state the caveat there too, cross- referenced from the 503 rows for the external-broker case. Link the follow-up as #629 instead of the unlinked "windowed-ingest follow-up". 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
TestEmbedded_SweepChunkOverTombstonesDoesNotHoldCommits claimed more
than it checked. sweepTouchHook fired once per element of `candidates`,
and the fixture has exactly one visible key, so `touched <= 1` could
never fail: a full keyspace walk added inside the locked phase would
still pass (R4). The TryLock ran in sweepScanHook, which fires AFTER
sweepCandidates returns, so a regression that read under commitMu and
released the lock just before that hook still passed (R3) — only
"the whole chunk under one lock" was actually caught (R5). The
300k-tombstone fixture changed no verdict either way (a handful of
tombstones behaves identically) while costing 1.65s alone / 3.1s under
package load in a 15s-timeout package, and the comment recorded this
test's own history (citing a commit that a squash would erase) instead
of what it asserts.
Drop sweepTouchHook (field, wiring, assertion). To pin "the read runs
unlocked" for real, sweepCandidates now takes an onKey hook called
once per key from inside its own loop, before evaluating it — wired
through sweepChunk as e.sweepReadHook. Renamed
TestEmbedded_SweepChunkOverTombstonesDoesNotHoldCommits to
TestEmbedded_SweepReadRunsUnlocked: TryLock/Unlock from inside that
hook, while the read is still running, so a lock held anywhere during
the read is caught in the act rather than inferred from whether it was
released before some later checkpoint. The fixture shrinks to a
handful of tombstones ahead of the one live key, since the verdict
never depended on the count. Comment cut to the one thing the test
asserts; the old wall-clock/hook history belongs here instead.
Verified by mutation, each applied then reverted (`git diff
internal/dedupe/sweep.go` clean before the real change was made):
- R3 (only the read under the lock, released right after): wrapped
just the sweepCandidates call in sweepChunk with
commitMu.Lock()/Unlock() — TestEmbedded_SweepReadRunsUnlocked failed
("commitMu must be free while sweepCandidates' read is running").
- R5 (the whole chunk under one lock): wrapped sweepChunk's body in
commitMu.Lock()/defer Unlock() and dropped deleteSweepable's own
lock to avoid a self-deadlock — ran ONLY the target test (not the
package: TestEmbedded_SweepNeverDeletesACommitLandingMidChunk
self-deadlocks under this mutation, racing a Commit against a sweep
that never releases the lock) — failed with the same assertion.
GOTOOLCHAIN=go1.26.6 go test -race -count=5 -run Sweep
./internal/dedupe/ passes (5/5, including
TestEmbedded_SweepReadRunsUnlocked and every other Sweep-prefixed
test); the full package (go test -race -count=1 ./internal/dedupe/...)
passes too, now in ~5s rather than the prior fixture's ~57s at
-count=20.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
TestEmbeddedNATS_Publish_IdempotencyKeyDropsARepeat and the api package's realPipeline opened the embedded broker on a bare t.TempDir(). Both replay through a disk-backed consumer, so they were exposed to the late consumer-state write (#442) that storedir absorbs. 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
TestEmbedded_SweepReadRunsUnlocked asserted only that the hook never saw commitMu held, which also holds if the hook never fires: passing nil for the hook left it green. It now counts the visits and requires one; with the hook unwired it fails. 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
A throttle wraps it too, and logged as "dedupe store is not open". The classify comment now says what ingest answers today: 500 for both, until #629 maps ErrUnavailable to a 503. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Fixes #384. Part of #613.
Part of the remote-dedupe design. It is stacked on #625 (
feat/dedupe-reserve).What changes
Ingest runs in windows (
internal/api/ingest.go).processRecordis split.prepareRecorddoes validate → permissions → canonicalize → resolve the dedupe id → encode.ingestWindowthen drives each window of up to 256 records through three phases:Reservefor the window's keyed records.Claimedrecord is published undermq.WithIdempotencyKey(dedupe.IdempotencyKey(tenant, key)).Commitfor the claims that were published.The single-object path is a window of one. On Pebble, each Commit is one fsync, so a batch now costs one fsync per window instead of one per record.
How a failure settles:
dedupe.ErrUnavailable(store not open; later a remote backend throttled or unreachable)503 {"error":"dedupe store unavailable"},Retry-After: 5. Nothing in the window is published. Earlier windows stay published and committed.500 dedupe failed, unchangedInFlightkey503with the lease asRetry-After, unchanged.ErrDisabledmq.ErrQueueFull(definite: the broker refused it)< kare committed and records≥ kreleased.503withRetry-After: 30.mq.ErrUnavailable(uncertain: it may have been stored)< kare committed. k's claim is left to lapse with the lease, and records> kare released.503withRetry-Afterequal to the lease in whole seconds when k held a claim (a sooner retry only gets the in-flight503), else5.< kare committed. k's claim is left to lapse with the lease, and records> kare released.500 publish failed, noRetry-After.The uncertain-publish rule. The earlier design released the claim on every publish failure. Now only a definite failure releases it, so a retry of an event that JetStream actually stored cannot publish a second copy:
503.Nats-Msg-Id. JetStream drops that copy and reports success.This relies on 2×lease + 1s ≤ duplicate window (the in-flight
503sends the full lease asRetry-After, so an obeying client can republish up to about twice the lease after the reserve; a network backend may round expiry up by 1 s). The embedded ingest stream now setsDuplicates: mq.EmbeddedDuplicateWindow(2 min) explicitly. Existing on-disk streams get it through theCreateOrUpdateStreamthatSetMaxBytesalready does, and a test pins that.TestIngest_DedupeLeaseFitsTheDuplicateWindowpinsdedupe.DefaultLease(30 s) against the invariant.mq (
internal/mq/mq.go+14,embedded.go+11):WithIdempotencyKey(key) PublishOptsetsNats-Msg-Idinsideinternal/mq, so nothing outside it names a NATS header. EveryBrokerimplementation is conformance-tested to honour it (IdempotencyKeyDropsARepeat). dedupe (key.go):IdempotencyKey(tenant, key)ishex(sha256(stored key))[:32], which is the design's formula over the layout #625 owns.503 sentinel. #625 already defines
dedupe.ErrUnavailableinmanaged.go("wrapped by a backend's error when a retry later can succeed"). This PR maps it; no new sentinel. The DynamoDB backend must wrap its throttling and timeout errors in it.Deviations from / additions to the design
mq.ErrQueueFullis the only definite publish failure. The design also counts "a failure before anything was sent (marshal, invalid topic)" as definite. Marshal now happens in prepare, before Reserve. An invalid topic cannot be told apart from other errors without an mq change, and treating it as uncertain only delays a retry by one lease.mq.ErrUnavailable("cannot be reached or does not answer in time") is its own uncertain case: a timeout may have stored the event, so it gets its own503row above rather than the generic 500.TestIngest_JSONArray_SyntaxError_Fatalnow asserts nothing was published. The api.md at-least-once caution says so.MockDeduplicatornow answersDuplicatefor a key repeated in oneReservecall, asManageddoes. It also gainedErrAfter, theReserves/Commitscounters andCommitted(k).TestIngest_Dedup_ReserveError_LogLevelpins both by exact level and message.Deliberately left to later PRs
dedupe.ErrUnavailableto get the503.ingestWindowpasses0(forever), as fix(dedupe)!: reserve/commit ids, windowed ingest, retention, dynamodb #625 did.dedupe.leaseboot key (feat(app): choose the DynamoDB dedupe backend at boot #635): boot must refuse or warn when the lease exceeds the embedded window (mq.EmbeddedDuplicateWindow), including a test with a non-default lease once it's wired. For an external broker, it must warn when the lease exceeds the operator'sduplicate_window. Today the lease is always the 30 s default.Test evidence
New in
internal/api/ingest_window_test.go:dedupe.ErrUnavailable→503+Retry-After: 5in these cases: store not open, backend throttled, second window throttled (the first window stays published), and single object.internal/mq: a repeated key is dropped as a success, and a pre-existing stream gains the window on its next budget apply. The stream config assertsDuplicates.internal/dedupe/key_test.go:IdempotencyKeyis stable, 32 hex characters, and distinct across tenant, table, id, and a hashed long id.Fsyncs on Pebble (measured under heavy load,
go test -bench BenchmarkIngest_DedupBatchOnPebble -benchtime 20x ./internal/api/), for a 1,000-record NDJSON batch:TestIngest_Windows_OneSyncPerWindowOnPebblepins 4 Commits for 1,000 records.make cipasses; every coverage gate passed; Go total coverage 94.5%;internal/apitook 5.2 s under-race.🤖 Generated with Claude Code
https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds