Skip to content

feat(app): choose the DynamoDB dedupe backend at boot - #635

Merged
EricAndrechek merged 31 commits into
feat/dedupe-dynamodbfrom
feat/dedupe-dynamodb-wiring
Sep 26, 2026
Merged

EricAndrechek merged 31 commits into
feat/dedupe-dynamodbfrom
feat/dedupe-dynamodb-wiring

Conversation

@EricAndrechek

@EricAndrechek EricAndrechek commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

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 its dedupe.backend enum and wireDedupe switch. The resolution: wireMQ(ctx) now switches to wireEmbeddedMQ(ctx), because main's ctx-taking wireMQ met #618's switch. In AGENTS.md, CHANGELOG.md and architecture.md, both sides were kept.

What changes

dedupe.backend: dynamodb selects #628's dedupe.Dynamo: one table that every tenant and every process shares. An id ingested through one pod is then a duplicate through every other pod.

key env default notes
dedupe.backend WH_DEDUPE_BACKEND pebble now also dynamodb
dedupe.lease WH_DEDUPE_LEASE 30s → IngestHandler.DedupeLease. It is also the in-flight 503's Retry-After. Refused when lease + ceil(lease) + 1s > 2m while mq.backend=embedded (so at most 59s) — that mirrors the embedded duplicate window from #629 as a local constant. An explicit 0 is refused rather than read as the default.
dedupe.reserve_concurrency WH_DEDUPE_RESERVE_CONCURRENCY 64 → DynamoConfig.ReserveConcurrency; Pebble ignores it. Also sizes the HTTP idle-connection pool, never below the SDK default of 10. An explicit 0 is refused. Has no effect yet, since ingest still sends one id per call — #629's windowed ingest changes that.
dedupe.dynamodb.table WH_DEDUPE_DYNAMODB_TABLE — required when selected
dedupe.dynamodb.region …_REGION empty = SDK chain
dedupe.dynamodb.endpoint …_ENDPOINT empty dynamodb-local
dedupe.dynamodb.timeout …_TIMEOUT 250ms explicit 0 refused; the table check runs under 10× this (2.5s by default)
dedupe.dynamodb.max_attempts …_MAX_ATTEMPTS 3 explicit 0 refused
dedupe.dynamodb.retry_mode …_RETRY_MODE standard or adaptive; an empty value is refused
dedupe.dynamodb.create_table …_CREATE_TABLE false refused at config load unless endpoint is set
  • Config (internal/config/backends.go): DedupeDynamoDB is added to dedupeBackends. The DedupeDynamoDBConfig sub-block is validated only when the backend is selected. The DedupeDynamoDB name is used for both the enum constant and the struct. Negative durations and counts are refused, and the zero-value refusals above match DynamoConfig. NeedsDataDir already returned false for dedupe once the backend is not pebble. There are no credential keys: the SDK default chain supplies credentials.

  • Wiring (internal/app/wire.go, wireDynamoDedupe): NewDynamo is built from the block. Boot checks the table (Dynamo.Check, preceded by CreateTable when create_table is 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:

    • a flat directory refuses boot (dedupe open: …);
    • a nested directory boots, and every tenant with dedupe on answers 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's dynamodb.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) in stores.go: a store opens only once ready returns nil. Dynamo.Tenant stays 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.Close is the component's release. There are no Pebble gauges (dedupeStats nil).

  • Per-tenant dedupe.enabled / id_field / require_id stay in each tenant's config.json. Nothing moved.

Tests

  • internal/config/backends_test.go: env (every WH_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):

    • boot checks the table and never creates it without create_table;
    • a claim and its commit reach the table;
    • create_table creates the table and enables TTL;
    • a missing table refuses flat boot whether dedupe is on or off;
    • nested fails closed (ErrUnavailable; a dedupe-off tenant still gets ErrDisabled) and a reload opens the store once the table exists;
    • after the table appears, the background retry opens the store with no reload (TestRun_DynamoDBDedupeRetriesTheTableCheck);
    • an unresolved region refuses boot, flat and nested;
    • a settings reload makes no DynamoDB call and applies against the last check result (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: Gated fails closed, then opens.

  • Integration tests/integration/dedupe_dynamodb_app_test.go: two app.New instances, each with its own data_dir, embedded queue and ingest worker, and dedupe.backend: dynamodb over one dynamodb-local table (both created it with create_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 answer 503 with Retry-After: 7, which shows dedupe.lease reached 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.go is 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 for dynamodb.go; unit and integration still cover it.

  • make ci is green: unit 93.4%, integration 51.9%, e2e coverage 60.3–60.4% against a 60% floor, Go total 94.7%.

    • A couple of runs hit unrelated unit-test flakes in packages this PR does not change: a TempDir cleanup race in internal/ingest's TestStartIngestWorker_EndToEnd, and the pre-existing internal/mq flake in TestEmbeddedNATS_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 from obstruct() making directories under the JetStream streams directory while the server is also writing there). Both cleared on re-run.

Deliberately left to later PRs

  • fix(ingest): windowed reserve/publish/commit; 503 when dedupe down #629: windowed ingest, 503 on ErrUnavailable, Nats-Msg-Id and mq.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 by mq.EmbeddedDuplicateWindow directly.
  • For an external NATS stream, the operator owns the duplicate window; a corresponding boot check for that broker is separate work. This PR applies the cap only while mq.backend=embedded.
  • feat(dedupe): retention per tenant and table, and an expiry sweep #633: retention (dedupe.retention, per tenant). Ingest still commits with retention 0.
  • Provisioning a production table is out of scope, along with the real-table benchmark. No AWS resource was created; everything ran against dynamodb-local or a fake endpoint.
  • dedupe.Factory.Gated works 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

EricAndrechek and others added 7 commits September 24, 2026 23:27
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>
@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: fe49fe32-e9e3-4663-969f-e57e15cc7ead

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added documentation Improvements or additions to documentation go Pull requests that update go code area/dedupe Deduplication (Pebble, ScyllaDB) area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release area/app Process wiring (internal/app): component build, run, release labels Sep 25, 2026
EricAndrechek 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>
EricAndrechek and others added 4 commits September 26, 2026 04:52
…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
EricAndrechek and others added 10 commits September 26, 2026 05:08
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
EricAndrechek and others added 3 commits September 26, 2026 14:53
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
EricAndrechek merged commit 6096ddd into feat/dedupe-dynamodb Sep 26, 2026
1 check passed
@EricAndrechek
EricAndrechek deleted the feat/dedupe-dynamodb-wiring branch September 26, 2026 19:24
@github-project-automation github-project-automation Bot moved this from Backlog to Done in WaveHouse Task Board Sep 26, 2026
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/app Process wiring (internal/app): component build, run, release area/dedupe Deduplication (Pebble, ScyllaDB) area/docs Documentation, site/, README 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