Skip to content

feat(dedupe): add a DynamoDB backend, tested but not yet wired - #628

Merged
EricAndrechek merged 56 commits into
feat/dedupe-reservefrom
feat/dedupe-dynamodb
Sep 26, 2026
Merged

EricAndrechek merged 56 commits into
feat/dedupe-reservefrom
feat/dedupe-dynamodb

Conversation

@EricAndrechek

@EricAndrechek EricAndrechek commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

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 the Deduplicator contract from #625 on DynamoDB. Every tenant's keys live in one shared table. pk is a String holding the escaped key tenant/table/id (AppendKey-encoded), and is the only key; there is no sort key. Tenant(id) *Managed is the Factory, mirroring Embedded.Tenant. The client is shared, so opening a tenant's store is free.

Call DynamoDB request Notes
Reserve conditional PutItem{pk, st=1, ex, tk} per key, parallel up to ReserveConcurrency (64) condition attribute_not_exists(pk) OR ex <= :now. On a failed condition, ReturnValuesOnConditionCheckFailure: ALL_OLD returns the live item, so st says Duplicate or InFlight with no extra read. dynamodb-local 3.3.1 supports this (measured: the Duplicate cases pass).
Commit unconditional BatchWriteItem of {pk, st=2, tk, ex?}, 25 per call, in parallel, every chunk attempted even if another fails (the records are already published) retries UnprocessedItems with jittered backoff (doubling from 25 ms up to 200 ms, 8 rounds), then fails with ErrUnavailable. A key handed over twice is written once, because BatchWriteItem refuses duplicates. Retention 0 writes no ex.
Release DeleteItem conditional on tk = :tk AND st = :pending, every claim attempted even if another fails a failed condition counts as success: the claim lapsed, was re-claimed, or was committed.
  • Expiry is the native TTL attribute ex, in epoch seconds, rounded up, so a lease or retention never ends early. The item is live while now < ex, which is why the condition is ex <= :now rather than the design's ex < 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-ex items directly and checking they are claimable.
  • Reserve on error stops sending puts after the first failure and releases, by token, every put it sent that may have landed. That covers a put that errored or timed out, not only those that answered 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.
  • A retried put keeps its own claim. When the SDK retries a conditional PutItem whose first attempt was applied (a 500 or a connection reset after the write), the returned old item is pending and carries this put's own token, so Reserve answers Claimed instead of InFlight — which would otherwise lock the id for the whole lease and give the producer a spurious 503.
  • Error mapping. ProvisionedThroughputExceeded, Throttling, RequestLimitExceeded (the SDK's throttle codes), InternalServerError, any server fault, ReplicatedWriteConflictException, deadline exceeded, and connection errors all wrap dedupe.ErrUnavailable. fix(ingest): windowed reserve/publish/commit; 503 when dedupe down #629 turns that into 503 + Retry-After. ResourceNotFound, AccessDenied and 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 existing ErrUnavailable; no new sentinel was needed.
  • Circuit breaker (design §7): five unavailable PutItems within a second short-circuit Reserve for 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 outcome canceled; a put cancelled by the caller itself (not a sibling) leaves the breaker's count unchanged too.
  • Client. NewDynamo(ctx, cfg, extra...) uses config.LoadDefaultConfig, so EKS Pod Identity and IRSA need no WaveHouse code. The Region override is optional. Endpoint sets BaseEndpoint for dynamodb-local. The retryer is standard or adaptive with MaxAttempts (default 3) and full-jitter backoff, uniform over [0, min(25ms·2^attempt, ceiling)] with ceiling = Timeout/(2·(MaxAttempts−1)) — a call's retries wait at most half its Timeout (62.5 ms per retry, 125 ms total at the defaults). A throttled call ends as a max-attempts error mapped to ErrUnavailable, not DeadlineExceeded. Each call gets a Timeout deadline (default 250 ms, SDK retries included). extra passes config.LoadOptions through, which is how tests inject static credentials and a fault-injecting HTTP client.
  • Connection pool: the HTTP client keeps at least max(10, ReserveConcurrency) idle connections per host, so a warm Reserve of ReserveConcurrency keys 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 not pk (HASH, S) alone, and warns (without refusing) when TTL is not enabled on ex. CreateTable(ctx) is the dev-only helper. It returns ErrCreateTableNeedsEndpoint unless Endpoint is set, creates pk S on-demand, waits for the table, and enables TTL on ex. An existing table is left alone.
  • Metrics: 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-v2 core, service/dynamodb and config, plus what those require: credentials, sts/sso/ssooidc/signin, imds, smithy-go. The design allows config because the default credential chain is the requirement. Nothing else was added, and no existing module version moved (go.mod diff is additions only).

Docs. deployment.md gains "A shared dedupe table on DynamoDB", which covers the required attributes, an example Terraform table (on-demand, pk S, TTL on ex, 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.md and CHANGELOG.md also describe the backend as built but not selectable.

Not wired: deliberately left to later PRs

  • feat(app): choose the DynamoDB dedupe backend at boot #635: dedupe.backend: dynamodb, the dedupe.dynamodb.* boot keys (table, region, endpoint, timeout, max_attempts, retry_mode, create_table), reserve_concurrency, the wireDedupe case, calling Check at boot, and refusing create_table without endpoint at config load. Until then no deployment can select this backend; every deployment keeps Pebble.
  • fix(ingest): windowed reserve/publish/commit; 503 when dedupe down #629: 503 + Retry-After for ErrUnavailable, and the windowed ingest that makes a network backend viable for batches. Today's per-record sequential ingest would call DynamoDB serially.
  • feat(dedupe): retention per tenant and table, and an expiry sweep #633: per-tenant/table retention. Ingest passes 0 today, so a committed item carries no ex and is kept forever. TTL reclaims only lapsed claims until then, and the deployment page says so.
  • Real-table benchmark. internal/dedupe/dynamodb_bench_test.go (build tag dynamobench, never in CI) measures Reserve+Commit for 1 key and for a 256-key window against any table named by DEDUPE_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.
  • Creating a production table is out of scope for this repo.

Test evidence

  • Integration (tests/integration/dedupe_dynamodb_test.go, against amazon/dynamodb-local:3.3.1, pinned, in-memory, no persistent volume, started beside ClickHouse in setup_test.go):
    • dedupetest.Run covers all 15 conformance cases. Factory and Peer are two separate clients over one table, with a fresh table per case and no Advance, so the suite sleeps. FailNextReserve is implemented as an HTTP-level fault on the (n+1)th PutItem.
    • 32 separate clients race one id: exactly one Claimed, 31 InFlight, then 32 Duplicate after the commit.
    • Throttling: ThrottlingException, ProvisionedThroughputExceededException and RequestLimitExceeded are injected as real DynamoDB JSON error responses. Each maps to ErrUnavailable after the SDK's retry.
    • An unreachable endpoint gives ErrUnavailable, and the breaker opens after five failures.
    • A missing table is not ErrUnavailable.
    • Check refuses a table keyed on a binary pk. CreateTable is idempotent and enables TTL on ex.
    • TTL/expiry: a pending item's ex is the lease end, a commit's ex is the retention end, retention 0 writes no ex, and items planted with a past ex (both pending and committed) are claimable.
    • CreateTable is refused without an endpoint.
    • The conformance suite (28 cases, including the Reserve-race and retried-put cases) passes against amazon/dynamodb-local:3.3.1, 10/10 runs.
  • Mutation check: weakening the Reserve condition failed 11 of the conformance cases plus the 32-client test.
  • Unit (internal/dedupe/dynamodb_test.go, fake API, for the paths dynamodb-local cannot produce):
    • the error classification table;
    • ALL_OLD state read into Duplicate / InFlight;
    • a failed Reserve releasing every possibly-landed put by its exact token;
    • Commit's unprocessed-item retries (3 chunks, each retried once) and dedupe of a repeated key;
    • Commit giving up with ErrUnavailable;
    • Commit attempting every chunk, and Release every claim, when one fails. The fakes honour cancellation, and both tests fail against the first commit's cancel-on-first-error code (measured, 5/5);
    • a multi-key Reserve with one throttled put: only sent puts are released, and the breaker still trips after five such Reserves;
    • Release treating a failed condition as done;
    • breaker trip, cool-down and reset;
    • Check schema and TTL paths;
    • config defaults, retry modes and CreateTable refusal;
    • expiresAt rounding.
      internal/dedupe unit coverage is 93.1%.
  • .testcoverage.yml: internal/dedupe/dynamodb.go is 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 the internal/settings e2e exclusion. Unit and integration still cover the file, and the merged total counts it.
  • make ci passed locally: unit 93.3%, integration 50.9%, e2e 61.3%, Go total 94.6%. One run hit an unrelated flaky test in internal/mq (TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen, a mkdir race); it passed on re-run.

🤖 Generated with Claude Code

https://claude.ai/code/session_01G6Cz4H5k1spJAPk5CZeGds

EricAndrechek and others added 6 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
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
@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: e5cedf13-db52-40db-adce-74a679dd1c13

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 dependencies Pull requests that update a dependency file go Pull requests that update go code area/dedupe Deduplication (Pebble, ScyllaDB) area/docs Documentation, site/, README area/infra CI, build, deploy, Docker, release labels Sep 25, 2026
EricAndrechek and others added 4 commits September 25, 2026 02:56
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>
EricAndrechek and others added 6 commits September 25, 2026 11:25
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
# 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
EricAndrechek and others added 21 commits September 26, 2026 05:00
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
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
@github-actions github-actions Bot added the area/app Process wiring (internal/app): component build, run, release label Sep 26, 2026
@EricAndrechek
EricAndrechek merged commit d5d35f2 into feat/dedupe-reserve Sep 26, 2026
2 checks passed
@EricAndrechek
EricAndrechek deleted the feat/dedupe-dynamodb branch September 26, 2026 19:38
@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 area/infra CI, build, deploy, Docker, release dependencies Pull requests that update a dependency file 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