Skip to content

feat(dedupe): retention per tenant and table, and an expiry sweep - #633

Merged
EricAndrechek merged 17 commits into
feat/dedupe-windowed-ingestfrom
feat/dedupe-retention
Sep 26, 2026
Merged

EricAndrechek merged 17 commits into
feat/dedupe-windowed-ingestfrom
feat/dedupe-retention

Conversation

@EricAndrechek

@EricAndrechek EricAndrechek commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

Part of #613. Part of the remote-dedupe design, stacked on #629 (feat/dedupe-windowed-ingest).

Fixes #220

What changes

dedupe.retention (internal/settings/{settings,validate,store}.go). This is a new optional config.json key, and dedupe.tables.<table>.retention overrides it per table. It says how long a committed id stays a duplicate, as a Go duration string ("720h"). "0" keeps ids forever: that is the seed value and today's behaviour. A missing key means "0" (forever), and a table override without one inherits the tenant's, so an existing config.json needs no change.

  • Validation. It refuses a value that is not a duration ("30d", a bare number), a negative one, and a finite one below settings.MinDedupeRetention = 2 min, which is the embedded queue's duplicate window. It refuses rather than clamps, so the file never means something other than what it says. The reason for the floor (as noted in fix(ingest): windowed reserve/publish/commit; 503 when dedupe down #629): a deduped record is published under a Nats-Msg-Id derived from its id. If an id is re-sent after a shorter retention but inside the window, it is claimed again, and then JetStream drops it as a copy while the client gets 200. TestIngest_MinDedupeRetentionCoversTheDuplicateWindow pins the constant ≥ mq.EmbeddedDuplicateWindow. settings does not import mq, which would pull NATS into the control plane's import of settings.
  • Resolution. Store.DedupeFor(table) now returns a settings.Dedupe{Enabled, IDField, RequireID, Retention} struct, not three values, and resolves every field from one snapshot, as before. Ingest reads it per record, so a reload changes retention for the next record. Each pendingRecord carries its retention. commitClaims issues one Commit per distinct retention in a window: one in practice, two only if a reload lands mid-window.

Pebble sweep (internal/dedupe/sweep.go, embedded.go). The read side already honoured the committed 0x02 ‖ expiry value. Now:

  • A background loop starts when the instance opens and is stopped (and waited for) before it closes. Its first pass runs one minute after open, so a legacy value goes promptly, and then it runs hourly.
  • A pass walks the whole keyspace in chunks of 1,024 keys, with a 10 ms pause between chunks, reading each chunk without holding Embedded.commitMu. It collects candidates whose committed expiry has passed, plus any key whose value is not a commit at all (isCommit: a 9-byte committedMark value — only commits are ever stored in Pebble, so an older, 8-byte value from a prior key layout is never a commit). It then takes commitMu only to re-Get each candidate and delete those still expired or still not a commit (NoSync), so a key the sweep read as expired cannot be re-committed between that read and the delete. A Commit waits at most for about 1,024 point reads and one batch, however many stale values the read stepped over (measured under -race: a Commit racing a chunk over 300k stale values waited 241–278 ms with a single lock over the whole chunk, 5–11 ms with this two-phase read).
  • Deletes are NoSync. A delete lost to a crash is redone on the next pass, so the sweep adds no fsync to ingest.
  • A new metric, wavehouse_dedupe_swept_keys_total{reason="expired"|"version_0"}, counts what it deletes, and a pass that deletes anything logs one line at Info.
  • expiry() saturates at MaxInt64 rather than wrapping into the past for an absurd retention.

Cadence is a constant (1 min first, then 1 h), not a config key. The sweep only reclaims space: correctness never waits on it, which makes it a candidate for a scheduled background worker later, per the design.

Deviations from / additions to the design

  • A legacy (non-commit) value is deleted on the first pass, not "once older than the largest retention". Nothing reads a value that isn't a commit, so keeping it longer protects nothing.
  • Refuse, not clamp, for a retention under the window (above).
  • The design's Pebble value encodes expiry in unix seconds; fix(dedupe)!: reserve/commit ids, windowed ingest, retention, dynamodb #625 shipped UnixNano. This PR keeps fix(dedupe)!: reserve/commit ids, windowed ingest, retention, dynamodb #625's encoding, which is the on-disk format already.
  • A reload mid-window splits the window's commit by retention. The design says only "changing retention affects ids committed after the change".
  • Disk space of deleted keys comes back as Pebble compacts. This PR does not force a compaction, because an hourly manual compaction of the whole keyspace would rewrite it every hour, and a legacy value can be anywhere in the keyspace so the cost could not be bounded. When space is reclaimed is inferred, not measured.

Deliberately left to later PRs

Test evidence

  • internal/dedupe/sweep_test.go:
    • A sweep over 2,058 expired keys interleaved with 2,058 live keys crosses chunk boundaries. It deletes exactly the expired ones and 3 legacy (version-0) keys from two tenants. It keeps a forever key and a key that expired and was committed again before the sweep, which is still a duplicate afterwards. A second pass finds nothing.
    • Retention is honoured on read to the nanosecond before any sweep.
    • The loop runs on its own while the instance is open, and closing waits for it.
    • A cancelled sweep stops after one chunk.
    • expiry saturates.
    • TestEmbedded_SweepNeverDeletesACommitLandingMidChunk: a test-only hook races a Commit into the gap between a chunk reading its keys and deleting them; mutation-checked (fails 3/3 with the lock removed).
    • TestEmbedded_SweepKeepsAKeyCommittedAfterItsRead, TestEmbedded_SweepReadRunsUnlocked (asserts the chunk read genuinely runs without holding commitMu): each fails on the regression it names, checked by mutation (skipping the re-read, releasing the lock before the delete, reading under the lock, or running the whole chunk under one lock).
  • internal/settings:
    • A missing key loads and resolves to forever (tenant and inherited table level); a non-duration, a bare number, a negative value, and a value under the window are each refused, at the top level and per table.
    • "0", "0s", "2m" and table overrides (shorter, forever) are accepted.
    • DedupeFor cascades all four fields.
    • The seed resolves to retention 0.
  • internal/api/ingest_retention_test.go:
    • Real settings directory. Each table's retention reaches Commit, including a per-table "0".
    • A reload changes it: new ids get the new retention, ids committed earlier keep theirs, and the removed override is gone.
    • A reload mid-window gives two Commits with the right retention each.
    • The floor ≥ the queue window.
  • A mid-chunk Commit survives the sweep (above).
  • The conformance suite's existing "a commit expires after its retention" case still covers read-side expiry on both clocks.
  • make ci is green: unit 93.4%, integration 45.1%, e2e 61.6%, Go total 94.5%, and every coverage gate passed.

🤖 Generated with Claude Code

https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds

EricAndrechek and others added 2 commits September 25, 2026 03:49
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
@coderabbitai

coderabbitai Bot commented Sep 25, 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: f5c41fb5-cd67-4f4e-9705-3ed284ea6f07

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/api HTTP handlers, routing, middleware area/dedupe Deduplication (Pebble, ScyllaDB) area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/app Process wiring (internal/app): component build, run, release labels Sep 25, 2026
EricAndrechek and others added 7 commits September 25, 2026 11:26
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/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:
#	CHANGELOG.md
#	docs/src/content/docs/architecture.md
EricAndrechek and others added 8 commits September 25, 2026 13:37
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
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
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
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
EricAndrechek merged commit a76bbb9 into feat/dedupe-windowed-ingest Sep 26, 2026
3 checks passed
@EricAndrechek
EricAndrechek deleted the feat/dedupe-retention branch September 26, 2026 19:24
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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area/api HTTP handlers, routing, middleware area/app Process wiring (internal/app): component build, run, release area/dedupe Deduplication (Pebble, ScyllaDB) area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release 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.

1 participant