Skip to content

fix(ingest): windowed reserve/publish/commit; 503 when dedupe down - #629

Merged
EricAndrechek merged 35 commits into
feat/dedupe-reservefrom
feat/dedupe-windowed-ingest
Sep 26, 2026
Merged

EricAndrechek merged 35 commits into
feat/dedupe-reservefrom
feat/dedupe-windowed-ingest

Conversation

@EricAndrechek

@EricAndrechek EricAndrechek commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

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). processRecord is split. prepareRecord does validate → permissions → canonicalize → resolve the dedupe id → encode. ingestWindow then drives each window of up to 256 records through three phases:

  1. Reserve: one Reserve for the window's keyed records.
  2. Publish: publish in record order. A Claimed record is published under mq.WithIdempotencyKey(dedupe.IdempotencyKey(tenant, key)).
  3. Commit: one Commit for 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:

Where What happens
Reserve → 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.
Reserve → any other error 500 dedupe failed, unchanged
Reserve → an InFlight key The window's claims are released. 503 with the lease as Retry-After, unchanged.
Reserve → ErrDisabled The window is published un-deduped and counted, unchanged
Publish fails at record k with mq.ErrQueueFull (definite: the broker refused it) Records < k are committed and records ≥ k released. 503 with Retry-After: 30.
Publish fails at k with mq.ErrUnavailable (uncertain: it may have been stored) Records < k are committed. k's claim is left to lapse with the lease, and records > k are released. 503 with Retry-After equal to the lease in whole seconds when k held a claim (a sooner retry only gets the in-flight 503), else 5.
Publish fails at k with anything else (uncertain: it may have been stored) Records < k are committed. k's claim is left to lapse with the lease, and records > k are released. 500 publish failed, no Retry-After.
Commit fails after the publishes The records still succeed. The failure is logged and counted per record.
A mid-body read error, or a prepare abort The open window is dropped unpublished

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:

  • A retry inside the lease answers the in-flight 503.
  • A retry after the lease is claimed again and republished under the same Nats-Msg-Id. JetStream drops that copy and reports success.

This relies on 2×lease + 1s ≤ duplicate window (the in-flight 503 sends the full lease as Retry-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 sets Duplicates: mq.EmbeddedDuplicateWindow (2 min) explicitly. Existing on-disk streams get it through the CreateOrUpdateStream that SetMaxBytes already does, and a test pins that. TestIngest_DedupeLeaseFitsTheDuplicateWindow pins dedupe.DefaultLease (30 s) against the invariant.

mq (internal/mq/mq.go +14, embedded.go +11): WithIdempotencyKey(key) PublishOpt sets Nats-Msg-Id inside internal/mq, so nothing outside it names a NATS header. Every Broker implementation is conformance-tested to honour it (IdempotencyKeyDropsARepeat). dedupe (key.go): IdempotencyKey(tenant, key) is hex(sha256(stored key))[:32], which is the design's formula over the layout #625 owns.

503 sentinel. #625 already defines dedupe.ErrUnavailable in managed.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.ErrQueueFull is 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 own 503 row above rather than the generic 500.
  • A mid-body read error or a prepare abort drops the open window unpublished. Before, the records before it were published. The whole request fails either way, and publishing less before a failure means fewer duplicates on retry. TestIngest_JSONArray_SyntaxError_Fatal now asserts nothing was published. The api.md at-least-once caution says so.
  • The windows run sequentially: reserve, then publish, then commit, then the next window. The design text mentions pipelining reserve with publish. That only pays off on a network backend, and it can wait for a benchmark once the DynamoDB backend is wired in.
  • MockDeduplicator now answers Duplicate for a key repeated in one Reserve call, as Managed does. It also gained ErrAfter, the Reserves/Commits counters and Committed(k).
  • A Reserve that fails because the request's own context ended logs at Debug; a real backend failure stays ERROR. TestIngest_Dedup_ReserveError_LogLevel pins both by exact level and message.
  • Measured: each deduped publish holds about 110 B of heap in JetStream's duplicate map for the 2-minute window (200k publishes: 153 vs 44 B/msg) — about 130 MB at 10k deduped events/s.

Deliberately left to later PRs

Test evidence

  • New in internal/api/ingest_window_test.go:

    • Windows of 1/255/256/257/600 records: one Reserve and one Commit per window.
    • A publish failure at k ∈ {1, 100, 256, 257, 590} (refused) and {100, 400} (uncertain). Each case checks what was committed, what was released, and that later windows were never reserved. After a refusal, a whole-batch retry publishes every record exactly once.
    • dedupe.ErrUnavailable → 503 + Retry-After: 5 in these cases: store not open, backend throttled, second window throttled (the first window stays published), and single object.
    • In-flight releases the window.
    • Result order across windows, for NDJSON and JSON-array bodies, covering rejects, repeats inside and across windows, and records with no id.
    • The bug(ingest): dedupe marks the event id before the NATS publish — a failed publish + client retry permanently drops the event #384 scenario end to end over the real embedded broker and real Pebble: refused publish → retry → exactly one event in the queue, and a third try is a duplicate.
    • Uncertain publish + retry, also end to end: the event is stored but the publish reports failure. The retry is in flight until the lease lapses, then republished and dropped by the queue: 3 events for 3 records, then all duplicates. Mutation-checked: with the idempotency key removed, the test counts 4 events.
  • 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 asserts Duplicates.

  • internal/dedupe/key_test.go: IdempotencyKey is 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:

    Window Time per batch Syncs per batch
    1 (per record) 5.71 s 1,000
    256 23.8 ms 4

    TestIngest_Windows_OneSyncPerWindowOnPebble pins 4 Commits for 1,000 records.

  • make ci passes; every coverage gate passed; Go total coverage 94.5%; internal/api took 5.2 s under -race.

🤖 Generated with Claude Code

https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds

EricAndrechek and others added 4 commits September 25, 2026 01:11
…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
@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 labels Sep 25, 2026
@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: f2f3833d-74ba-4619-bf1c-2ae8ae602c26

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.

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
EricAndrechek and others added 12 commits September 25, 2026 11:25
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
EricAndrechek and others added 12 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
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
EricAndrechek and others added 5 commits September 26, 2026 08:49
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
@github-actions github-actions Bot added the area/app Process wiring (internal/app): component build, run, release label Sep 26, 2026
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
@github-actions github-actions Bot added the area/infra CI, build, deploy, Docker, release label Sep 26, 2026
@EricAndrechek
EricAndrechek merged commit 5a8ef7e into feat/dedupe-reserve Sep 26, 2026
4 checks passed
@EricAndrechek
EricAndrechek deleted the feat/dedupe-windowed-ingest branch September 26, 2026 19:38
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