feat(dedupe): add a DynamoDB backend, tested but not yet wired - #628
Merged
Merged
Conversation
mq.backend, cache.backend, dedupe.backend and coord.backend select each layer's implementation; only today's in-process one exists per layer and it is the default. Validate refuses an unknown value, internal/app picks the implementation in one switch per layer, data_dir is probed only when a selected backend keeps state there, and boot logs Config.Warnings. Part of #613. 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
…ENTS.md Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
One shared table for every tenant, pk (binary) = the dedupe key, no sort key. Reserve is a conditional PutItem per key (ALL_OLD on failure answers Duplicate or InFlight without a read), Commit a BatchWriteItem with unprocessed-item retries, Release a DeleteItem conditional on the token. An item whose ex has passed is absent to Reserve whether or not TTL has deleted it. Throttles, server faults, timeouts and connection errors wrap ErrUnavailable; a breaker short-circuits Reserve after five in a second. Constructible and tested, not yet selectable at boot (F5). Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
… honest Review round 1. Commit and Release no longer cancel their siblings on the first failure: the records are already published, and a claim left behind holds its id for a lease. A failed Reserve sends no put after the first failure and releases only puts it sent; a sibling cancelled by that failure no longer resets the breaker (and is counted as outcome "canceled"). Docs: TTL reclaims only lapsed claims until retention lands (#220), no future boot-key names, the breaker counts consecutive unavailable claims. The e2e coverage suite excludes dynamodb.go, which the e2e binary never runs; unit and integration cover it and the merged total counts it. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Review round 2. The fakes' BatchWriteItem and DeleteItem honour their context, and the Commit test runs one chunk at a time, so reverting Commit/Release to cancel-on-first-error fails both tests (checked against 108499f). The failed-Reserve test asserts the sent puts are released rather than every key, which the unsent-put cutoff made flaky. forEach no longer sits between expiresAt and its doc comment. 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 |
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
dedupe.backend: dynamodb selects the shared table (F3's dedupe.Dynamo), configured by a dedupe.dynamodb block; dedupe.lease and dedupe.reserve_concurrency join the dedupe boot block. Boot checks the table, creating it first only with create_table on dynamodb-local; a failed check refuses a flat boot and fails switched-on tenants closed over a nested directory until a reload's check passes (Factory.Gated). The lease is capped at the embedded queue's 2m duplicate window. Part of #613 (PR F5). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
A nested directory has no watcher, so a table check that failed at boot (a throttle, credentials not yet issued) left dedupe'd ingest failing closed until someone reloaded. The check now retries in the background, backing off 1s to 30s, and reconciles once it passes. NewDynamo refuses a config that resolves no region, a certain error caught at boot. Docs: reserve_concurrency has no effect while ingest sends one id per call; 0 = default; Pebble-only sentences scoped to dedupe.backend: pebble. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…lause Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
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
# Conflicts: # docs/src/content/docs/architecture.md
# Conflicts: # docs/src/content/docs/architecture.md
The dedupe key is now ASCII text, so the table's partition key is a String: it reads as-is in the console and in get-item output, with the same 2,048-byte limit. Check() now requires pk to be S, CreateTable declares it S, and the Deployment table and Terraform example follow. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Brings in origin/main (process roles, leases, backend selection, the broker conformance suite, ClickHouse error classes), the keyenc change that keeps '-', and the conformance cases that race Reserve against a Commit and widen the concurrent-Reserve race to 20 keys. Conflicts, all prose: - AGENTS.md: took the base's config/ and new coord/ entries and re-added Dynamo to the dedupe/ entry. - docs/src/content/docs/architecture.md: took the base's tree with coord/, keeping "Pebble, DynamoDB" on the dedupe/ line. - CHANGELOG.md: kept the DynamoDB entry above the base's new entries. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
keyenc keeps '-' since the merge from the base, so the example id evt-123 is stored as acme/clicks/evt-123, not acme/clicks/evt%2D123. Fixed in the Deployment page's table of attributes, which now also says what the escaping keeps, and in the CHANGELOG entry. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
While the table check failed, every settings reload ran it from the AfterAdopt hook, which the registry runs under the lock that serializes reloads: 2.5s per reload against a hung endpoint at the default timeout, the hooks registered after it waiting, and a tenant the reload switched on answering ErrDisabled meanwhile, so its records published un-deduped. The hook now applies every store against the last check's result, which fails a switched-on store closed (ErrUnavailable) while the check has not passed, and wakes the background retry through a one-slot channel, so a reload still retries at once. Only boot and that loop run the check, and once passed it stays passed. A flat directory's boot is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The log said ingest with dedupe on fails closed until a reload passes the check, which is wrong in both shapes: a flat directory refuses boot, and a nested one retries the check in the background. The boot line and the retry's line now say that it is retried and ingest fails closed meanwhile. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The base brings five DynamoDB fixes (a retried put keeps its own claim, jittered retries within half the call timeout, an idle pool sized to reserve_concurrency, a throttled commit retried as a round) and its own merge of main. keyenc takes the base's comment; the docs combine both sides line by line, keeping this branch's facts that dynamodb is selectable at boot and the lease is configurable. The timeout, max_attempts and reserve_concurrency docs now state the new backoff ceiling and the idle-pool sizing. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
MaxIdleConnsPerHost was set to ReserveConcurrency outright, so a small ReserveConcurrency (a boot key soon) shrank the pool below the SDK's default of 10 per host. It is now the larger of the two, as MaxIdleConns already was. TestNewDynamo_SizesTheIdlePool adds 4 and asserts both limits stay at or above the SDK's defaults; it failed at 4 and 8 before. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The Metrics bullet on the Deployment page named three of the backend's four metrics; wavehouse_dedupe_dynamodb_short_circuits_total appeared only in passing one bullet earlier. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Combines the Reserve-on-error contract reword, the Managed.Apply no-op fast path, and the commit-race conformance start barrier with this branch's DynamoDB backend. Both merged cleanly with no textual conflicts; the development.md tree line was updated by hand afterward to list DynamoDB alongside Pebble, since this branch builds it but the line had not been touched yet. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Conflicts in docs/src/content/docs/api.md and architecture.md: this branch had reordered the ingest-error-table rows and added a dedupe.lease config link, while feat/dedupe-dynamodb had reworded the uncertain-publish caveat and the idle-connection-pool sizing. Kept this branch's row order and link, combined with the other branch's updated wording. The architecture.md dynamodb.go entry kept this branch's "selected by dedupe.backend: dynamodb" framing (accurate now that wiring makes it selectable) and its no-region-resolves refusal, combined with the other branch's "at least as many idle connections" fix. All other files merged with no textual conflict. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The old check allowed a lease up to 59.5s (2*lease+1s <= 2m), but the DynamoDB claim it bounds rounds its expiry up to the next whole second, and the in-flight 503 it drives can go out a second late. The true worst case is lease + ceil(lease) + 1s, which only stays under the embedded queue's 2-minute duplicate window through 59s exactly: 59.5s (and anything else over 59s) already crosses it once ceil(lease) steps to the next second. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…at a reload waits on Every remaining 59.5s (config.yaml, CHANGELOG.md, configuration.mdx, architecture.md) now reads 59s with the lease+ceil(lease)+1s rule spelled out, matching the previous commit's config change. Adds the boot/background table check's own 10x-timeout deadline (2.5s by default) to the dedupe.dynamodb.timeout row, since that call is not one of the per-request calls the row otherwise describes. Marks the development.md tree line "Pebble or DynamoDB" now that this PR makes the backend selectable, replacing the prior "not yet selectable at boot" framing that architecture.md's dynamodb.go entry already dropped. Qualifies "a reload never waits on the table" everywhere it appears (configuration.mdx, deployment.md, CHANGELOG.md, architecture.md's wire.go paragraph): true for a tenant whose dedupe.enabled did not change, thanks to Managed.Apply's no-op fast path settling under a read lock alone — but switching a tenant's dedupe off is a genuine transition, which takes the write lock and so waits for that tenant's in-flight Reserve/Commit/Release calls to finish first. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…B call Adds a fake per-op hang (fakeDynamo.setHangOn) alongside the existing whole-endpoint one, then exercises Managed.Apply's no-op fast path end to end: a tenant with dedupe on has a Commit in flight against a BatchWriteItem that never answers, and asserts that a reload naming no change for that tenant still returns in well under a second, and that a concurrent Reserve for the same tenant is not blocked either — which it would be if the reload's Apply took the write lock unconditionally, since a pending writer blocks new readers too. Verified by hand: temporarily skipping the settled check in internal/dedupe/managed.go's Apply made the reload assertion fail at 2.475s (Commit's own retry/backoff giving up on the hung call, not a deadlock) instead of passing in 0.02–0.03s; reverted, `git diff` on that file confirmed clean before this commit. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The Reserve only started after Reload had already returned, so no writer was ever pending when it ran, and it could not fail for the reason its comment gave — with the fast path removed only the reload-time assertion actually failed. No clean sync point exists into "Reload is inside the dedupe hook" without instrumenting production code for the test alone, so drop the block and the sentence claiming it; the reload-time assertion already pins the regression on its own. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
configuration.mdx's dedupe.reserve_concurrency row now says the DynamoDB idle-connection pool never goes below the SDK's own default (10), matching newHTTPClient's max(...) floor. configuration.mdx, deployment.md and CHANGELOG.md said switching a tenant's dedupe off was the one case a reload waits on. It is not the only one: a reload that removes or rejects a tenant whose store was open closes it the same way (Stores.Retain -> Managed.Close -> Apply(false)), taking the same write lock and waiting on the same in-flight calls. Reworded all three to match architecture.md's existing "on a genuine flip" framing, which already covered both cases. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Takes the base's AGENTS.md tree line (naming DynamoDB beside Pebble here) and #667's teardown-flake fix. No conflicts. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Takes the base's AGENTS.md tree line and #667's teardown-flake fix. No conflicts. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
make ci's e2e coverage measured 59.7% against a 60% floor: the e2e binary always boots with dedupe.backend: pebble, so wireDynamoDedupe, errDynamoUnchecked and the table-check retry component never ran there, same shape as #628's internal/dedupe/dynamodb.go exclusion. Pure move, no behavior change: wireDynamoDedupe and errDynamoUnchecked move verbatim into the new internal/app/wire_dynamodb.go; wireDedupe's switch stays in wire.go untouched. Adds a matching e2e-only exclusion for the new file in .testcoverage.yml, next to dynamodb.go's, so unit and integration keep covering it and the merged total still counts it -- wire.go itself stays out of the exclude list. Updates architecture.md's wire.go bullet (the dynamodb case now points at the new file's own bullet) and the CHANGELOG's file-provenance list for the dedupe.backend: dynamodb entry. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
Puts already sent, and the release after them, run on context.WithoutCancel(ctx), each call still bounded by its Timeout. A client that disconnects mid-Reserve no longer leaves an abandoned put to land after its release and hold the id InFlight for the lease; a Reserve whose caller cancelled after every put answered also releases them. The breaker's exemption for cancelled puts is gone: a sent put can no longer be cancelled by its caller, so its answer is always the table's. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
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
Over a flat settings directory, boot was refused on any failed table check. Now it is refused only when the failure is a misconfiguration (not ErrUnavailable: a missing table, the wrong key schema, access denied) and a tenant has dedupe on. A transient failure, or a misconfigured table no tenant uses yet, boots with the switched-on stores closed and the check retried in the background, as a nested directory already did; a tenant a reload switches on fails closed until it passes. The misconfiguration is logged at ERROR. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
No put applies before the caller cancels, so neither subtest depends on scheduling. The CHANGELOG line now names what a cancel can still leave held: a put cut off by its own timeout, or a failed release. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
Brings in the caller-cancel fix for DynamoDB Reserve and the renamed ErrUnavailable message. The CHANGELOG and architecture conflicts keep this branch's entries with the cancel wording carried onto them. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
CreateTable returned its errors unclassified, so with create_table on a dynamodb-local not listening yet read as a misconfiguration and refused boot over a flat directory. Its errors now go through classify, and a test pins that a transient create failure boots and is retried. settings-directory.mdx still said any failed open refuses boot; that is now scoped to Pebble, with the DynamoDB rule stated beside it. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
…n rule The table section still said boot refuses a table whose key schema does not match, in every case. Drop a sentence configuration.mdx said twice. 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.
Part of #613. Part of the remote-dedupe design, stacked on #625 (
feat/dedupe-reserve).Addresses #648 (a Reserve-on-error flake; closes when this stack lands on main).
What changes
dedupe.Dynamo(internal/dedupe/dynamodb.go) implements theDeduplicatorcontract from #625 on DynamoDB. Every tenant's keys live in one shared table.pkis a String holding the escaped keytenant/table/id(AppendKey-encoded), and is the only key; there is no sort key.Tenant(id) *Managedis theFactory, mirroringEmbedded.Tenant. The client is shared, so opening a tenant's store is free.ReservePutItem{pk, st=1, ex, tk}per key, parallel up toReserveConcurrency(64)attribute_not_exists(pk) OR ex <= :now. On a failed condition,ReturnValuesOnConditionCheckFailure: ALL_OLDreturns the live item, sostsaysDuplicateorInFlightwith no extra read. dynamodb-local 3.3.1 supports this (measured: the Duplicate cases pass).CommitBatchWriteItemof{pk, st=2, tk, ex?}, 25 per call, in parallel, every chunk attempted even if another fails (the records are already published)UnprocessedItemswith jittered backoff (doubling from 25 ms up to 200 ms, 8 rounds), then fails withErrUnavailable. A key handed over twice is written once, because BatchWriteItem refuses duplicates. Retention 0 writes noex.ReleaseDeleteItemconditional ontk = :tk AND st = :pending, every claim attempted even if another failsex, in epoch seconds, rounded up, so a lease or retention never ends early. The item is live whilenow < ex, which is why the condition isex <= :nowrather than the design'sex < now. The two differ only in the boundary second. Rounding up plus the half-open interval lets the suite's whole-second sleeps work without an injected clock. Correctness never depends on TTL having deleted anything. This is pinned by planting past-exitems directly and checking they are claimable.Claimed. The token is generated client-side, and a conditional delete of an item the token does not hold is a no-op. A sent put runs under the caller's own context rather than the group's, so a sibling's failure cancels only puts not yet sent, and a sent put's release waits for it to answer first — though a put cut off by the caller's own cancellation or deadline (not a sibling's failure) can still land after its release and hold the key until the lease ends. This makes "all-or-nothing on error" hold even for a put whose outcome is unknown.PutItemwhose first attempt was applied (a500or a connection reset after the write), the returned old item is pending and carries this put's own token, soReserveanswersClaimedinstead ofInFlight— which would otherwise lock the id for the whole lease and give the producer a spurious503.ProvisionedThroughputExceeded,Throttling,RequestLimitExceeded(the SDK's throttle codes),InternalServerError, any server fault,ReplicatedWriteConflictException, deadline exceeded, and connection errors all wrapdedupe.ErrUnavailable. fix(ingest): windowed reserve/publish/commit; 503 when dedupe down #629 turns that into503+Retry-After.ResourceNotFound,AccessDeniedand validation errors stay plain errors (500), because they are configuration bugs. The sentinel is fix(dedupe)!: reserve/commit ids, windowed ingest, retention, dynamodb #625's existingErrUnavailable; no new sentinel was needed.PutItems within a second short-circuitReservefor a second (ErrUnavailable, no request sent). Only Reserve's puts feed it. A Release or Commit answer says nothing about whether a new claim would get through, and the conditional deletes of a failed Reserve's cleanup would otherwise reset the count on every failure. A put cancelled because a sibling failed does not touch the breaker either, and is counted as outcomecanceled; a put cancelled by the caller itself (not a sibling) leaves the breaker's count unchanged too.NewDynamo(ctx, cfg, extra...)usesconfig.LoadDefaultConfig, so EKS Pod Identity and IRSA need no WaveHouse code. TheRegionoverride is optional.EndpointsetsBaseEndpointfor dynamodb-local. The retryer isstandardoradaptivewithMaxAttempts(default 3) and full-jitter backoff, uniform over[0, min(25ms·2^attempt, ceiling)]withceiling = Timeout/(2·(MaxAttempts−1))— a call's retries wait at most half itsTimeout(62.5 ms per retry, 125 ms total at the defaults). A throttled call ends as a max-attempts error mapped toErrUnavailable, notDeadlineExceeded. Each call gets aTimeoutdeadline (default 250 ms, SDK retries included).extrapassesconfig.LoadOptionsthrough, which is how tests inject static credentials and a fault-injecting HTTP client.max(10, ReserveConcurrency)idle connections per host, so a warm Reserve ofReserveConcurrencykeys does not re-dial for each one (measured: a 64-key Reserve opened 46–54 new connections before this fix, 0 after).Check(ctx)is the boot check, for feat(app): choose the DynamoDB dedupe backend at boot #635 to call. It refuses a table whose key schema is notpk(HASH,S) alone, and warns (without refusing) when TTL is not enabled onex.CreateTable(ctx)is the dev-only helper. It returnsErrCreateTableNeedsEndpointunlessEndpointis set, createspkS on-demand, waits for the table, and enables TTL onex. An existing table is left alone.wavehouse_dedupe_dynamodb_requests_total{op,outcome}(outcome ok / condition_failed / unavailable / canceled / error),wavehouse_dedupe_dynamodb_request_duration_seconds{op},wavehouse_dedupe_dynamodb_unprocessed_items_total,wavehouse_dedupe_dynamodb_short_circuits_total.Dependencies.
github.com/aws/aws-sdk-go-v2core,service/dynamodbandconfig, plus what those require: credentials, sts/sso/ssooidc/signin, imds, smithy-go. The design allowsconfigbecause the default credential chain is the requirement. Nothing else was added, and no existing module version moved (go.moddiff is additions only).Docs.
deployment.mdgains "A shared dedupe table on DynamoDB", which covers the required attributes, an example Terraform table (on-demand,pkS, TTL onex, deletion protection, SSE) and the least-privilege IAM policy (PutItem,DeleteItem,BatchWriteItem,DescribeTable,DescribeTimeToLive; no Scan, no CreateTable). The example tags the table with neutral placeholder values.architecture.md,AGENTS.mdandCHANGELOG.mdalso describe the backend as built but not selectable.Not wired: deliberately left to later PRs
dedupe.backend: dynamodb, thededupe.dynamodb.*boot keys (table,region,endpoint,timeout,max_attempts,retry_mode,create_table),reserve_concurrency, thewireDedupecase, callingCheckat boot, and refusingcreate_tablewithoutendpointat config load. Until then no deployment can select this backend; every deployment keeps Pebble.503+Retry-AfterforErrUnavailable, and the windowed ingest that makes a network backend viable for batches. Today's per-record sequential ingest would call DynamoDB serially.0today, so a committed item carries noexand is kept forever. TTL reclaims only lapsed claims until then, and the deployment page says so.internal/dedupe/dynamodb_bench_test.go(build tagdynamobench, never in CI) measures Reserve+Commit for 1 key and for a 256-key window against any table named byDEDUPE_BENCH_TABLE. It has not been run against a real AWS table, because this PR creates no AWS resources. Run against dynamodb-local only (1 key ≈ 4 ms, 256 keys ≈ 140 ms). Those figures say nothing about the service, so design §5.3's latency figures remain inferred until someone runs it against a real table.Test evidence
tests/integration/dedupe_dynamodb_test.go, againstamazon/dynamodb-local:3.3.1, pinned, in-memory, no persistent volume, started beside ClickHouse insetup_test.go):dedupetest.Runcovers all 15 conformance cases.FactoryandPeerare two separate clients over one table, with a fresh table per case and noAdvance, so the suite sleeps.FailNextReserveis implemented as an HTTP-level fault on the (n+1)thPutItem.Claimed, 31InFlight, then 32Duplicateafter the commit.ThrottlingException,ProvisionedThroughputExceededExceptionandRequestLimitExceededare injected as real DynamoDB JSON error responses. Each maps toErrUnavailableafter the SDK's retry.ErrUnavailable, and the breaker opens after five failures.ErrUnavailable.Checkrefuses a table keyed on a binarypk.CreateTableis idempotent and enables TTL onex.exis the lease end, a commit'sexis the retention end, retention 0 writes noex, and items planted with a pastex(both pending and committed) are claimable.CreateTableis refused without an endpoint.amazon/dynamodb-local:3.3.1, 10/10 runs.internal/dedupe/dynamodb_test.go, fake API, for the paths dynamodb-local cannot produce):ALL_OLDstate read into Duplicate / InFlight;ErrUnavailable;Checkschema and TTL paths;CreateTablerefusal;expiresAtrounding.internal/dedupeunit coverage is 93.1%..testcoverage.yml:internal/dedupe/dynamodb.gois excluded from the e2e suite gate only. The e2e binary runs Pebble dedupe, and with the file included e2e measured 58.6%, under its 60% floor. The precedent is theinternal/settingse2e exclusion. Unit and integration still cover the file, and the merged total counts it.make cipassed locally: unit 93.3%, integration 50.9%, e2e 61.3%, Go total 94.6%. One run hit an unrelated flaky test ininternal/mq(TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen, amkdirrace); it passed on re-run.🤖 Generated with Claude Code
https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds