feat(app): choose the DynamoDB dedupe backend at boot - #635
Merged
EricAndrechek merged 31 commits intoSep 26, 2026
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
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>
|
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 |
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
# Conflicts: # docs/src/content/docs/architecture.md
# Conflicts: # CHANGELOG.md # docs/src/content/docs/architecture.md # docs/src/content/docs/settings-directory.mdx
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
EricAndrechek
added a commit
that referenced
this pull request
Sep 25, 2026
…647) Fixes #617. Part of #613. `internal/app` and `internal/mq` each took 11–12 s of the 15 s unit-test budget on main under a full parallel `make test-unit`, and timed out when the machine was loaded. After this PR they take under 5 s each. This PR also makes `NewEmbedded` fail at once on a store directory it cannot create, instead of waiting out the 5 s readiness check. ## What changes | Commit | Change | |---|---| | 2a494bd | `NewEmbedded` runs `MkdirAll` on the store first and returns `nats store: <mkdir error>`. Before this, a regular file at `<data_dir>/nats` failed JetStream in the background, and boot reported only `nats server not ready` after 5 s. New test: `TestNewEmbedded_AStoreItCannotCreateFailsAtOnce`. | | d46e273, f897394 | The directory is created at `0700`, the same mode nats-server uses (`defaultDirPerms`). CHANGELOG entry under Fixed, limited to what mkdir catches: an existing directory that cannot be written to still fails the old way. | | 98acc11 | The `mq` tests call `t.Parallel`. A new `mq.EmbeddedSyncAlways` is `true` in production and is set to false only in `TestMain`: on macOS the fsync on every JetStream write was more than half the run. Each store lives in `storeDir`, which retries its removal (#442). | | e441f57 | `internal/app`'s new `TestMain` also turns the fsync off. | | 668c2c4 | Removes the "occupied"-directory hunk that #612's squash added to `TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen` (see below). | These are 2a19d7e, 96e2112, bd059f2, 3d959b9 and 3c8a93b from `perf/unit-test-budget`, cherry-picked onto main. Two hunks from the perf branch are left out because they belong to stacks that are not on main: the `t.Parallel` lines for `embedded_failed_test.go` and for `TestEmbeddedNATS_Publish_IdempotencyKeyDropsARepeat`. Neither test exists on main; they come from the integration tree (#645). The CHANGELOG entry was re-applied under main's `### Fixed`. ## The `AQueueThatCannotOpen` conflict: measured #612 (f5d8f48) and the perf change fix the same race in different ways. The race: after a failed open, the server removes the empty `streams/` and `$G` directories on a goroutine of its own, and that removal overlaps the next open. f5d8f48 has three hunks: - **AQ-occ**: an ignored `occupied` directory keeps `streams/` non-empty during acme's failed open in `SetMaxBytes_AQueueThatCannotOpen`. This was added to main's squash. - **P-globex**: `PacesTheRetriesOfAQueueThatCannotOpen` opens globex first, so its streams keep the directory occupied. - **App-globex**: `internal/app`'s `TestNew_QueueOpenFailure` "nested" subtest blocks globex rather than acme. The perf change's version of the fix is **AQ-wait**: `require.Eventually` until the server has removed `$G`, and only then open globex. These runs were on this branch, with `-race` and `GOTOOLCHAIN=go1.26.6`, on 2026-09-25 on an otherwise idle machine. "Isolated" means `-run '^Name$' -count=20`. "pkg" means the whole parallel `internal/mq` package with `-count=5`. "Mutation" means `JetStreamMaxStore: math.MaxInt64` instead of `/ 2`, which is the refusal the test exists to catch, at `-count=5`. | P-globex | AQ variant | AQueue isolated | Paces isolated | AQueue in pkg | Paces in pkg | other pkg fails | Mutation caught | |---|---|---|---|---|---|---|---| | kept | **wait only (this PR)** | **20/20 pass** | **20/20** | **5/5** | **5/5** | **0** | **5/5 fail (caught)** | | kept | occupied only | 20/20 | 20/20 | 5/5 | 5/5 | 0 | 5/5 fail (caught) | | kept | both | **0/20** | 20/20 | 0/5 | 5/5 | — | — | | dropped | wait only | 20/20 | 20/20 | 5/5 | **4/5** | 1 | 5/5 caught | | dropped | occupied only | 20/20 | 20/20 | 5/5 | 5/5 | 0 | 5/5 caught | | dropped | both | 0/20 | 20/20 | 0/5 | 3/5 | 2 | — | | App-globex | `TestNew_QueueOpenFailure` isolated ×20 | in the `internal/app` package ×3 | |---|---|---| | kept (main) | 20/20 | 3/3 | | reverted to acme | 20/20 | 3/3 | What the runs show: - **Both AQ fixes together always fail**, 20/20. The `occupied` directory keeps `$G` from ever being removed, so the wait for its removal times out. The finding from #645 reproduces. - **P-globex is what `PacesTheRetries` needs.** Without it, the pacing test raced in the full parallel package (1/5 and 2/5 failures), even though it passed 20/20 when run alone. This is the finding from #646. It is a different test from `AQueueThatCannotOpen`, so the two findings do not conflict. P-globex is on main and this PR keeps it. - **For `AQueueThatCannotOpen`, AQ-wait alone and AQ-occ alone both passed every run, and both caught the mutation in this matrix.** The earlier finding, that the occupied variant still passes with `MaxInt64`, did **not** reproduce here. On main as merged (occupied only), the mutation also fails the test 5/5 (measured). I kept AQ-wait because it is the version the perf change was written and measured against. It is also what the test's comment describes: no queue is kept open, and no directory is kept around to hold the reservation count up. Dropping AQ-occ instead of AQ-wait would work equally well by these numbers. - **App-globex made no difference in either direction** in 20 isolated runs and 3 package runs. It is left as main has it. ## Timings: full parallel `make test-unit` Four runs, alternating between main at 2a2b886 and this branch, with `-count=1 -race -cover`, on an otherwise idle machine for every run. | | `internal/app` | `internal/mq` | whole run (`DONE … in`) | |---|---|---|---| | main (2a2b886) | 11.30 / 11.37 / 11.41 / 11.45 s | 12.05 / 12.18 / 12.30 / 12.38 s | 12.07–12.39 s | | this PR | 4.71 / 4.72 / 4.73 / 4.81 s | 3.42 / 3.64 / 3.64 / 3.94 s | 5.07–5.27 s | An earlier baseline of main alone, under heavy load, gave `internal/app` 12.2–13.5 s and `internal/mq` 12.8–15.1 s. The 15.1 s run was already over the budget. ## Verification - `make ci` passed through the shared queue (`GOTOOLCHAIN=go1.26.6`, the known golangci-lint toolchain workaround), including all coverage gates. - Pre-push reviewers were both run on HEAD 668c2c4 (opus, fresh context): - `pre-push-reviewer`: **ship_it**, with 0 MUST, 0 SHOULD and 0 MAY findings. - `docs-reviewer`: **ship_it**, with 0 findings. The CHANGELOG entry was checked against nats-server's `defaultDirPerms` and `wireMQ`'s error path. No docs page needed a change. - Because of #454, the hook wrote the markers against the main checkout's HEAD, not this branch's HEAD. No marker was written by hand. ## Left for later - `internal/testutil.NewEmbeddedMQ`, used by the `internal/ingest` and `internal/api` tests, still fsyncs on every write. Those packages are not near the budget today. The reviewer noted it as the next place to get the same speed-up. - 8776b4d (#635) and 09d5c14 (#639) from the perf branch belong to their own stacks and are not included here. 🤖 Generated with [Claude Code](https://claude.com/claude-code) https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…wiring Main carries the final boot-backends change, process roles, coord leases, zero-keeping config defaults, the mq conformance suite and the keyenc '-' change. Resolved by taking main's version of what those own and re-applying this branch's dynamodb wiring on top: the dynamodb case sits inside main's api-role gate, and the coord wording is main's. Dedupe key examples follow keyenc keeping '-' (evt-123, not evt%2D123). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The dedupe keys carried cleanenv env-default tags, which main no longer allows: cleanenv applies them after the YAML decode to any field still zero, so an explicit zero in config.yaml silently became the default. dedupe.lease, reserve_concurrency and the dynamodb block's timeout, max_attempts and retry_mode now default in defaults(), and a zero lease, concurrency, timeout or attempt count, or an empty retry_mode, refuses boot (the dynamodb block's only while dynamodb is selected). Each is a refusedZeros case; the docs-defaults test gains a duration parser. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds
The cap sat at the duplicate window itself (2m), which is the wrong boundary. The in-flight 503 answers with the whole lease as Retry-After, so a client that obeys it after a publish whose outcome it never learned republishes up to twice the lease after the claim, and DynamoDB rounds a claim's expiry up to the second. With mq.backend=embedded a lease is now refused unless twice it plus one second fits the 2m window. The 30s default is unchanged; 59s and 59.5s pass, 60s is refused. 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
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 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
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
This was referenced Sep 26, 2026
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 #628 (
feat/dedupe-dynamodb). Commits of its own:feat(app): choose the DynamoDB dedupe backend at boot(ced118c),fix(app): retry a failed DynamoDB table check; refuse no region(a8e43bd),docs(dedupe): …(0b60451).This branch also carries #618's commits (
feat/boot-backends), merged in at 83a6949 because this PR extends itsdedupe.backendenum andwireDedupeswitch. The resolution:wireMQ(ctx)now switches towireEmbeddedMQ(ctx), because main's ctx-takingwireMQmet #618's switch. In AGENTS.md, CHANGELOG.md and architecture.md, both sides were kept.What changes
dedupe.backend: dynamodbselects #628'sdedupe.Dynamo: one table that every tenant and every process shares. An id ingested through one pod is then a duplicate through every other pod.dedupe.backendWH_DEDUPE_BACKENDpebbledynamodbdedupe.leaseWH_DEDUPE_LEASE30sIngestHandler.DedupeLease. It is also the in-flight503'sRetry-After. Refused whenlease + ceil(lease) + 1s > 2mwhilemq.backend=embedded(so at most59s) — that mirrors the embedded duplicate window from #629 as a local constant. An explicit0is refused rather than read as the default.dedupe.reserve_concurrencyWH_DEDUPE_RESERVE_CONCURRENCY64DynamoConfig.ReserveConcurrency; Pebble ignores it. Also sizes the HTTP idle-connection pool, never below the SDK default of 10. An explicit0is refused. Has no effect yet, since ingest still sends one id per call — #629's windowed ingest changes that.dedupe.dynamodb.tableWH_DEDUPE_DYNAMODB_TABLEdedupe.dynamodb.region…_REGIONdedupe.dynamodb.endpoint…_ENDPOINTdedupe.dynamodb.timeout…_TIMEOUT250ms0refused; the table check runs under 10× this (2.5s by default)dedupe.dynamodb.max_attempts…_MAX_ATTEMPTS30refuseddedupe.dynamodb.retry_mode…_RETRY_MODEstandardadaptive; an empty value is refuseddedupe.dynamodb.create_table…_CREATE_TABLEfalseendpointis setConfig (
internal/config/backends.go):DedupeDynamoDBis added todedupeBackends. TheDedupeDynamoDBConfigsub-block is validated only when the backend is selected. TheDedupeDynamoDBname is used for both the enum constant and the struct. Negative durations and counts are refused, and the zero-value refusals above matchDynamoConfig.NeedsDataDiralready returned false for dedupe once the backend is notpebble. There are no credential keys: the SDK default chain supplies credentials.Wiring (
internal/app/wire.go,wireDynamoDedupe):NewDynamois built from the block. Boot checks the table (Dynamo.Check, preceded byCreateTablewhencreate_tableis on) whether or not any tenant has dedupe on. Boot is the validator, and a DynamoDB check is one cheap call. A failed check follows the registry's rule for the directory shape, the same rule as a Pebble instance that cannot open:dedupe open: …);ErrUnavailable(fails closed). A reload applies stores against the last table-check result — it makes no DynamoDB call itself — and wakes the background retry below.A nested directory has no settings watcher, so a check that fails at boot is also retried in the background: the component "dedupe table check" starts at 1s and backs off to 30s, then reconciles the stores once the table passes. A mutex keeps its reconcile and the reload hook's from running at once. Without the retry, a throttle or a credential delay during a rolling deploy would leave dedupe'd ingest failing closed until someone sent a reload. A settings reload also no longer waits on any in-flight DynamoDB call: it waits only for a tenant whose store it actually closes (dedupe switched off, or the tenant removed or rejected) — measured at 0.02 s versus 2.5 s against a hung endpoint before this fix.
No region refuses boot in both shapes (
NewDynamo:awsCfg.Region == "", a small change to feat(dedupe): add a DynamoDB backend, tested but not yet wired #628'sdynamodb.go). That config error is certain, so it is caught at boot rather than surfacing as a check failure later.The gate is a new
dedupe.Factory.Gated(ready)instores.go: a store opens only oncereadyreturns nil.Dynamo.Tenantstays the Factory, as feat(dedupe): add a DynamoDB backend, tested but not yet wired #628 intended. Without the gate, a tenant's open would succeed for free against a missing table.Stores.Closeis the component's release. There are no Pebble gauges (dedupeStatsnil).Per-tenant
dedupe.enabled/id_field/require_idstay in each tenant'sconfig.json. Nothing moved.Tests
internal/config/backends_test.go: env (everyWH_DEDUPE_*variable), YAML (durations parse), defaults, unknown sub-block keys (dedupe.dynamodb.access_key_id,dedupe.redis) refused, unbound-env knows every new variable, and a validation table. The table covers create_table refused without an endpoint, a missing table, retry mode, negatives, the lease at 2m (ok before this round, refused after) and the current cap boundary, and the block ignored under pebble.internal/app/dedupe_dynamodb_test.go(a fake DynamoDB JSON endpoint; AWS env pinned):create_table;create_tablecreates the table and enables TTL;ErrUnavailable; a dedupe-off tenant still getsErrDisabled) and a reload opens the store once the table exists;TestRun_DynamoDBDedupeRetriesTheTableCheck);TestReload_DynamoDBDedupeMakesNoTableCall), wakes the background retry (TestRun_DynamoDBDedupeReloadWakesTheRetry), and does not wait on an in-flight commit (TestReload_DynamoDBDedupeDoesNotWaitOnInFlightCommit).Mutation checks: with the gate removed, the nested case fails; with the retry loop disabled, the retry test fails.
internal/dedupe/stores_test.go:Gatedfails closed, then opens.Integration
tests/integration/dedupe_dynamodb_app_test.go: twoapp.Newinstances, each with its own data_dir, embedded queue and ingest worker, anddedupe.backend: dynamodbover one dynamodb-local table (both created it withcreate_table, and the second found it already there). An id accepted by pod A is{"duplicate":true}on pod B, and the reverse. A claim held by a third client makes both pods answer503withRetry-After: 7, which showsdedupe.leasereached ingest. After release, the id goes through once. ClickHouse ends with exactly one row per id, each from the pod that accepted it.The integration suite's 9 DynamoDB tests (Conformance, Check, Expiry, ThirtyTwoClientsOneID, Throttled, Unreachable, CreateTableNeedsEndpoint, ConfigErrorsAreNotUnavailable, TwoInstancesShareSeenIDs) pass against
amazon/dynamodb-local:3.3.1.internal/app/wire_dynamodb.gois excluded from the e2e coverage gate only (the e2e binary runs Pebble), matching feat(dedupe): add a DynamoDB backend, tested but not yet wired #628's exclusion fordynamodb.go; unit and integration still cover it.make ciis green: unit 93.4%, integration 51.9%, e2e coverage 60.3–60.4% against a 60% floor, Go total 94.7%.TempDircleanup race ininternal/ingest'sTestStartIngestWorker_EndToEnd, and the pre-existinginternal/mqflake inTestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen(measured failing 7/15 in isolation on the unmodified feat(dedupe): add a DynamoDB backend, tested but not yet wired #628 base, 8/15 on this branch — it comes fromobstruct()making directories under the JetStream streams directory while the server is also writing there). Both cleared on re-run.Deliberately left to later PRs
503onErrUnavailable,Nats-Msg-Idandmq.EmbeddedDuplicateWindow. This PR's lease cap is a local constant (config.embeddedDuplicateWindow) until fix(ingest): windowed reserve/publish/commit; 503 when dedupe down #629 lands, at which point it should be replaced bymq.EmbeddedDuplicateWindowdirectly.mq.backend=embedded.dedupe.retention, per tenant). Ingest still commits with retention 0.dedupe.Factory.Gatedworks for any remote backend that checks a shared resource once, and the background table-check retry is a plain run component that a future backend (a Redis-backed dedupe store, say) could reuse.🤖 Generated with Claude Code
https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds