Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
153 commits
Select commit Hold shift + click to select a range
ef93154
feat(mq): give every tenant a queue of its own
taitelee Sep 24, 2026
8ad9865
fix(mq): reopen detached from the request, and the upgrade runbook swept
taitelee Sep 24, 2026
650a28e
fix(app): boot opens queues under New's context; docs review fixes
taitelee Sep 24, 2026
999c1db
fix(mq): boot re-applies a budget to a split queue pair; review fixes
taitelee Sep 24, 2026
db2d20a
fix(mq): a failed resize restores the ingest stream's own cap; review…
taitelee Sep 25, 2026
07c6a91
fix(app): a rejected tenant keeps its replay history; review fixes
taitelee Sep 25, 2026
04aad9f
fix(mq): share the hub bridge's fetch-ahead across tenants; review fixes
taitelee Sep 25, 2026
57870c4
fix(mq): publish only into a queue the broker recorded open; review f…
taitelee Sep 25, 2026
7cdb794
docs(mq): a consumer that cannot join a queue opened at runtime; revi…
taitelee Sep 25, 2026
e199a03
fix(mq): pace publish-side retries of a queue that cannot open; revie…
taitelee Sep 25, 2026
55137d8
test(mq): one tenant's failed purge stops no other; review fixes
taitelee Sep 25, 2026
060ca9e
docs(mq): size a dead-letter stream the shrink guard kept; review fixes
taitelee Sep 25, 2026
3022b92
feat(config): choose each layer's implementation at boot
EricAndrechek Sep 25, 2026
fde17ba
docs(config): say coord.backend is reserved; sync the boot-config lists
EricAndrechek Sep 25, 2026
f129d57
docs(config): no backend has a sub-block yet; index backends.go in AG…
EricAndrechek Sep 25, 2026
58d9e77
fix(dedupe)!: reserve, commit or release ids keyed by tenant and table
EricAndrechek Sep 25, 2026
0ec037e
fix(dedupe): state what Managed guarantees a backend; default a zero …
EricAndrechek Sep 25, 2026
9344f1b
fix(mq): pace a park's reopen, and warn only on a missing consumer; r…
taitelee Sep 25, 2026
f484511
fix(dedupe): release only claimed claims on a wrong-count answer
EricAndrechek Sep 25, 2026
20c28d1
Merge remote-tracking branch 'origin/mq-tenant-streams' into feat/ded…
EricAndrechek Sep 25, 2026
386ce74
docs(dedupe): link the old-key sweep to #220 rather than promise it
EricAndrechek Sep 25, 2026
108499f
feat(dedupe): DynamoDB backend, conformance-tested on dynamodb-local
EricAndrechek Sep 25, 2026
37ef757
fix(ingest): windowed reserve/publish/commit; 503 when dedupe is unav…
EricAndrechek Sep 25, 2026
71ab824
fix(ingest): qualify the duplicate-window claims; steadier tests
EricAndrechek Sep 25, 2026
6f7944a
fix(dedupe): attempt every Commit and Release chunk; keep the breaker…
EricAndrechek Sep 25, 2026
ff047d2
test(dedupe): make the Dynamo failure tests fail without their fix
EricAndrechek Sep 25, 2026
9824663
docs(ingest): Release gives back definite failures only; tighten wording
EricAndrechek Sep 25, 2026
99a5eee
docs(ingest): say the Pebble benchmark stubs the queue; finish the re…
EricAndrechek Sep 25, 2026
83a6949
Merge origin/feat/boot-backends into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
ced118c
feat(app): choose the DynamoDB dedupe backend at boot
EricAndrechek Sep 25, 2026
b83e670
feat(dedupe): retention per tenant and table, and an expiry sweep
EricAndrechek Sep 25, 2026
a8e43bd
fix(app): retry a failed DynamoDB table check; refuse no region
EricAndrechek Sep 25, 2026
6c74b65
test(dedupe): pin that a commit mid-sweep-chunk survives; docs wording
EricAndrechek Sep 25, 2026
0b60451
docs(dedupe): name every region source; untangle the dynamodb check c…
EricAndrechek Sep 25, 2026
18e0b6c
Merge remote-tracking branch 'origin/main' into feat/dedupe-reserve
EricAndrechek Sep 25, 2026
fc4085b
docs(durability): describe the benchmark machine neutrally
EricAndrechek Sep 25, 2026
02d7056
docs(deployment): use neutral tags in the DynamoDB Terraform example
EricAndrechek Sep 25, 2026
841a2f6
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
9a7f849
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
d3032cc
fix(dedupe): length-prefix the key's fields; refuse no table name
EricAndrechek Sep 25, 2026
7b9595e
feat(settings): a missing dedupe.retention keeps ids forever
EricAndrechek Sep 25, 2026
12741f8
docs(dedupe): Reserve refuses no key; scope the table-name CHANGELOG …
EricAndrechek Sep 25, 2026
a9ed5b8
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
dae1b46
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
e2d97de
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 25, 2026
c11133d
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
76b6a9e
docs(settings): name dedupe.retention as the one compiled default
EricAndrechek Sep 25, 2026
bea9e70
docs(settings): the last every-key-required comment; rewrap Store's
EricAndrechek Sep 25, 2026
ce14795
refactor(keyenc): one escaping for subject and cache key tokens
EricAndrechek Sep 25, 2026
3a02719
Merge branch 'refactor/keyenc' into feat/dedupe-reserve
EricAndrechek Sep 25, 2026
2f11973
refactor(dedupe): key ids as readable escaped text
EricAndrechek Sep 25, 2026
5b723f0
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
fed9e04
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 25, 2026
51db11c
docs(keyenc): name only the keys built from it today
EricAndrechek Sep 25, 2026
27f7608
docs(architecture): keyenc is not yet every composite key's escaping
EricAndrechek Sep 25, 2026
cfdc06c
feat(dedupe): key the DynamoDB table by a String pk
EricAndrechek Sep 25, 2026
935638e
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
48a6a77
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
7b757ad
Merge branch 'refactor/keyenc' into feat/dedupe-reserve
EricAndrechek Sep 25, 2026
78e24aa
docs(keyenc): name the dedupe keys among its users
EricAndrechek Sep 25, 2026
35989e4
test(dedupe): pair / with its escape in the keyspace case; review fixes
EricAndrechek Sep 25, 2026
9f42d99
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 25, 2026
fc0d509
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
b299d93
test(dedupe): keep the table/id boundary pair across '/'
EricAndrechek Sep 25, 2026
ca48065
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
c4ab354
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 25, 2026
c5e4f3b
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
1526c2b
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
5d77211
test(dedupe): a legacy value under a current key reads as absent
EricAndrechek Sep 25, 2026
7082907
fix(dedupe): let a sent DynamoDB put answer before its release
EricAndrechek Sep 25, 2026
484c72e
docs(changelog): the boot-backends entry predates dedupe.dynamodb
EricAndrechek Sep 25, 2026
4c2c2a4
test(dedupe): a caller-cancelled put does not reset the breaker
EricAndrechek Sep 25, 2026
59b0f0c
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
673f461
Merge remote-tracking branch 'origin/main' into refactor/keyenc
EricAndrechek Sep 25, 2026
2c31eff
refactor(keyenc): keep '-', join keys, count dead letters per table
EricAndrechek Sep 25, 2026
52a8145
Merge branch 'refactor/keyenc' into feat/dedupe-reserve
EricAndrechek Sep 25, 2026
999e92c
docs(keyenc): pin the %2D compatibility to builds since #612
EricAndrechek Sep 25, 2026
28f3886
refactor(dedupe): build keys with keyenc.AppendJoin; keep '-'
EricAndrechek Sep 25, 2026
deccbf9
Merge branch 'refactor/keyenc' into feat/dedupe-reserve
EricAndrechek Sep 25, 2026
79ab364
test(mq): cover a byte left unescaped in the lenient-token case
EricAndrechek Sep 25, 2026
8794bdf
Merge branch 'refactor/keyenc' into feat/dedupe-reserve
EricAndrechek Sep 25, 2026
7335888
docs(dedupe): '-' is kept in the hashing bound and the keyenc listings
EricAndrechek Sep 26, 2026
bfebcee
test(dedupe): pin a pre-#222 collision and the hashed marker's escaping
EricAndrechek Sep 26, 2026
f507a2e
test(dedupe): pin the length half of committedLive; document the lease
EricAndrechek Sep 26, 2026
155138d
test(dedupe): a failed release keeps the publish's error; drop "valid…
EricAndrechek Sep 26, 2026
44e0a69
test(dedupe): cover Status.String
EricAndrechek Sep 26, 2026
a103ab7
Merge remote-tracking branch 'origin/main' into feat/dedupe-reserve
EricAndrechek Sep 26, 2026
6ea5dbe
Merge remote-tracking branch 'origin/main' into feat/dedupe-reserve
EricAndrechek Sep 26, 2026
3248e7e
fix(dedupe): a retried DynamoDB put keeps its own claim
EricAndrechek Sep 26, 2026
8895863
fix(dedupe): jittered DynamoDB retries that fit the call timeout
EricAndrechek Sep 26, 2026
e61307b
Merge remote-tracking branch 'origin/main' into feat/dedupe-dynamodb-…
EricAndrechek Sep 26, 2026
4b377cb
fix(dedupe): size the DynamoDB client's idle pool to its fan-out
EricAndrechek Sep 26, 2026
da1079e
test(dedupe): race Reserve against a Commit in progress
EricAndrechek Sep 26, 2026
e52b336
fix(dedupe): sweep holds the commit lock only to re-check and delete
EricAndrechek Sep 26, 2026
ffa4dd6
docs(settings): name dedupe.retention in the last no-defaults claims
EricAndrechek Sep 26, 2026
eb5dea9
fix(dedupe): retry a throttled DynamoDB commit like unprocessed items
EricAndrechek Sep 26, 2026
f40c76a
fix(config): dedupe defaults in defaults(); refuse an explicit zero
EricAndrechek Sep 26, 2026
c56103e
docs(dedupe): scope the DynamoDB Reserve undo to what it guarantees
EricAndrechek Sep 26, 2026
8b616b4
fix(config): cap dedupe.lease at 59.5s over the embedded queue
EricAndrechek Sep 26, 2026
4ac7ec8
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 26, 2026
96b9968
docs(dedupe): the DynamoDB key example keeps '-' now
EricAndrechek Sep 26, 2026
a694418
fix(app): the dedupe reload hook makes no DynamoDB call
EricAndrechek Sep 26, 2026
9b8ea3c
fix(app): say what a failed dedupe table check leads to
EricAndrechek Sep 26, 2026
29ddf36
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 26, 2026
56727fb
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 26, 2026
0425a55
fix(dedupe): never size the DynamoDB idle pool below the SDK default
EricAndrechek Sep 26, 2026
5c64ace
docs(dedupe): list the short-circuit counter with the DynamoDB metrics
EricAndrechek Sep 26, 2026
860d929
test(mq): state the lease/duplicate-window invariant as 2*lease+1s
EricAndrechek Sep 26, 2026
6552478
test(mq): pin the takeStock boot fix for a stale duplicate window
EricAndrechek Sep 26, 2026
67edd57
test(mq): add an idempotency-key case to the Broker conformance suite
EricAndrechek Sep 26, 2026
a6b2108
fix(app,api): correct a stale comment; don't ERROR-log a client-gone …
EricAndrechek Sep 26, 2026
04da801
test(dedupe): bound the commit-race case under low GOMAXPROCS
EricAndrechek Sep 26, 2026
f0f39eb
docs(api): the uncertain-publish caveat covers every publish failure
EricAndrechek Sep 26, 2026
2c27538
test(api): cover ErrUnavailable in FailedPublishReleasesTheID
EricAndrechek Sep 26, 2026
cf24115
docs(dedupe): reword Reserve's on-error guarantee for network backends
EricAndrechek Sep 26, 2026
719a3f8
perf(dedupe): skip Managed.Apply's write lock when unchanged
EricAndrechek Sep 26, 2026
aea0600
docs(development): note dedupe's Reserve/Commit/Release in the tree
EricAndrechek Sep 26, 2026
f177f4d
fix(test): remove the broker store after late consumer writes land
EricAndrechek Sep 26, 2026
1e7e12f
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 26, 2026
97197f2
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 26, 2026
fe44591
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 26, 2026
5c45de2
docs(agents): name Reserve/Commit/Release in the dedupe tree line
EricAndrechek Sep 26, 2026
af6205d
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 26, 2026
3513823
fix(config): cap the embedded-mq dedupe lease by lease+ceil(lease)+1s
EricAndrechek Sep 26, 2026
2c42f0a
docs(dedupe): sweep the lease cap, the table check's deadline, and wh…
EricAndrechek Sep 26, 2026
d0a61e4
test(app): a reload must not wait behind a tenant's in-flight DynamoD…
EricAndrechek Sep 26, 2026
3235598
test(cmd): keep the boot test's broker store in storedir
EricAndrechek Sep 26, 2026
5e1ff2d
test(dedupe): replace the sweep-lock timing assertion with structural…
EricAndrechek Sep 26, 2026
b16c193
test(app): drop the concurrent-Reserve half of the reload/commit test
EricAndrechek Sep 26, 2026
674746a
docs(dedupe): the idle-pool floor, and what a reload's wait covers
EricAndrechek Sep 26, 2026
b97e4d7
Merge branch 'fix/ingest-teardown-race' into feat/dedupe-reserve
EricAndrechek Sep 26, 2026
167496e
Merge branch 'feat/dedupe-reserve' into feat/dedupe-dynamodb
EricAndrechek Sep 26, 2026
51716a2
test(api): pin the ERROR/DEBUG split on a Reserve failure by exact msg
EricAndrechek Sep 26, 2026
4424338
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 26, 2026
d4d50a3
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 26, 2026
c4147ec
test(dedupe): pin sweepCandidates' read as unlocked from inside its loop
EricAndrechek Sep 26, 2026
2138068
test: keep this branch's two broker stores in storedir
EricAndrechek Sep 26, 2026
b691d26
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 26, 2026
a76bbb9
test(dedupe): assert the sweep's read hook ran
EricAndrechek Sep 26, 2026
7f38b9b
fix(app): move DynamoDB dedupe wiring out of wire.go for the e2e gate
EricAndrechek Sep 26, 2026
dc9a331
fix(dedupe): a caller's cancel no longer strands a DynamoDB claim
EricAndrechek Sep 26, 2026
07e4f77
fix(dedupe): name ErrUnavailable "dedupe store unavailable"
EricAndrechek Sep 26, 2026
138b166
feat(app): refuse boot only for a misconfigured DynamoDB table in use
EricAndrechek Sep 26, 2026
23ac57a
test(dedupe): order the cancel test's puts after the cancel, not a sleep
EricAndrechek Sep 26, 2026
ade6f7b
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 26, 2026
8ac2fca
fix(dedupe): classify create_table errors; scope the boot rule in docs
EricAndrechek Sep 26, 2026
6096ddd
docs(deployment): a mismatched key schema follows the misconfiguratio…
EricAndrechek Sep 26, 2026
5a8ef7e
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-reserve
EricAndrechek Sep 26, 2026
d5d35f2
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-reserve
EricAndrechek Sep 26, 2026
bb35726
fix(dedupe): reconcile the windowed ingest with the configurable lease
EricAndrechek Sep 26, 2026
c361832
docs(dedupe): retention reaches DynamoDB through TTL, not the sweep
EricAndrechek Sep 26, 2026
34c12f0
Merge origin/main (#614 shared cache) into feat/dedupe-reserve
EricAndrechek Sep 26, 2026
6490203
docs: name the shared dedupe backend where the merge left it out
EricAndrechek Sep 26, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions .testcoverage.yml
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,17 @@ exclude:
- ^internal/settings/
- ^cmd/wavehouse/validate\.go$
- ^cmd/wavehouse/bootstrap\.go$
# The DynamoDB dedupe backend: the e2e binary runs Pebble dedupe, so
# this file measured 0% there and pulled e2e to 58.6%. The unit
# (fake API) and integration (dynamodb-local) suites cover it, and the
# merged total still counts it.
- ^internal/dedupe/dynamodb\.go$
# wireDynamoDedupe and its retry component (internal/app/wire_dynamodb.go):
# same reason as dynamodb.go above — the e2e binary never selects
# dedupe.backend: dynamodb, so this file measured 0% there and pulled
# e2e to 59.7%. The unit and integration suites cover it, and the
# merged total still counts it.
- ^internal/app/wire_dynamodb\.go$
# The in-process cache backend: the e2e stack runs cache.backend=redis
# (#613), so the binary carries LocalCache and its version index but e2e
# never reaches them. The unit suite and the integration suite's main
Expand Down
18 changes: 9 additions & 9 deletions AGENTS.md

Large diffs are not rendered by default.

12 changes: 9 additions & 3 deletions CHANGELOG.md

Large diffs are not rendered by default.

3 changes: 2 additions & 1 deletion cmd/wavehouse/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (

"github.com/Wave-RF/WaveHouse/internal/config"
"github.com/Wave-RF/WaveHouse/internal/settings"
"github.com/Wave-RF/WaveHouse/internal/testutil/storedir"
)

// run reads the whole boot config from the environment here (no config
Expand Down Expand Up @@ -78,7 +79,7 @@ func seedSettings(t *testing.T) string {
func TestRun_BootsAndStopsOnCancel(t *testing.T) {
hermeticEnv(t)
t.Setenv(config.EnvSettingsDir, seedSettings(t))
t.Setenv("WH_DATA_DIR", t.TempDir())
t.Setenv("WH_DATA_DIR", storedir.New(t))
_, port, err := net.SplitHostPort(closedAddr(t))
require.NoError(t, err)
t.Setenv("WH_SERVER_PORT", port)
Expand Down
2 changes: 1 addition & 1 deletion cmd/wavehouse/validate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ func writeSettingsDir(t *testing.T, policies string) string {
"roles.json": `{"roles": ["public"]}`,
"policies.json": policies,
"pipes.json": `{}`,
"config.json": `{"clickhouse": {"addr": "localhost:9000", "http_port": 8123, "http_scheme": "http", "database": "default", "username": "default", "query_timeout": 30, "tls": {"enabled": false, "ca_file": "", "cert_file": "", "key_file": "", "insecure_skip_verify": false, "server_name": ""}, "headers": {}, "max_open_conns": 10, "max_idle_conns": 5}, "auth": {"jwks_url": "", "role_claim": "role"}, "dedupe": {"enabled": false, "id_field": "event_id", "require_id": false}, "dlq": {"enabled": true}, "query": {"default_max_rows": 10000, "timestamp_bucket_seconds": 60}, "schema": {"refresh_interval": 60}, "stream": {"keepalive_interval": 30, "keepalive_buckets": 3, "gap_window_minutes": 15}, "mq": {"max_bytes_gb": 1}, "cors": {"allowed_origins": ["*"]}}`,
"config.json": `{"clickhouse": {"addr": "localhost:9000", "http_port": 8123, "http_scheme": "http", "database": "default", "username": "default", "query_timeout": 30, "tls": {"enabled": false, "ca_file": "", "cert_file": "", "key_file": "", "insecure_skip_verify": false, "server_name": ""}, "headers": {}, "max_open_conns": 10, "max_idle_conns": 5}, "auth": {"jwks_url": "", "role_claim": "role"}, "dedupe": {"enabled": false, "id_field": "event_id", "require_id": false, "retention": "0"}, "dlq": {"enabled": true}, "query": {"default_max_rows": 10000, "timestamp_bucket_seconds": 60}, "schema": {"refresh_interval": 60}, "stream": {"keepalive_interval": 30, "keepalive_buckets": 3, "gap_window_minutes": 15}, "mq": {"max_bytes_gb": 1}, "cors": {"allowed_origins": ["*"]}}`,
}
for name, content := range files {
require.NoError(t, os.WriteFile(filepath.Join(dir, name), []byte(content), 0o600))
Expand Down
22 changes: 16 additions & 6 deletions config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -51,12 +51,22 @@ clickhouse:
password: ""
max_total_conns: 0 # ceiling on open native connections across pools; 0 = none

# Each layer's implementation, chosen at boot. The in-process backend is the
# default for each, and the only one for these three.
# Each layer's implementation, chosen at boot. The in-process backend is
# each layer's default.
mq:
backend: embedded # NATS JetStream under <data_dir>/nats
dedupe:
backend: pebble # Pebble under <data_dir>/pebble
backend: pebble # Pebble under <data_dir>/pebble; or dynamodb (below)
lease: 30s # how long a claimed id stays pending; at most 59s with the embedded mq (lease + ceil(lease) + 1s within its 2m duplicate window)
reserve_concurrency: 64 # parallel calls per Reserve/Commit/Release to a remote backend, and DynamoDB's idle connections per host; ingest sends a window of up to 256 ids per call
# dynamodb: # read only when backend is dynamodb; credentials from the AWS SDK chain
# table: wavehouse-dedupe-prod
# region: "" # empty = AWS_REGION
# endpoint: "" # dynamodb-local only
# timeout: 250ms
# max_attempts: 3
# retry_mode: standard # or adaptive
# create_table: false # dynamodb-local only; refused without endpoint
coord:
backend: local # leases (the sweeper's) held in this process

Expand Down Expand Up @@ -88,12 +98,12 @@ auth:
# clickhouse wiring (addr, http_port, http_scheme, database, username,
# query_timeout, tls, headers, max_open_conns, max_idle_conns), auth
# (jwks_url, role_claim), dedupe (enabled/id_field/
# require_id + per-table overrides), dlq.enabled (+ per table),
# require_id/retention + per-table overrides), dlq.enabled (+ per table),
# query.default_max_rows / timestamp_bucket_seconds,
# schema.refresh_interval, stream keepalive_interval / keepalive_buckets /
# gap_window_minutes, mq.max_bytes_gb, cors.allowed_origins — and every key
# is required: the
# binary has no compiled defaults, so what's adopted is exactly what the
# is required except dedupe.retention (missing = "0", forever): the binary
# has no other compiled default, so what's adopted is exactly what the
# files say. The server validates the directory at boot (invalid or missing
# refuses to start) and reloads it on SIGHUP, on file change, or via
# POST /v1/ops/settings/reload; a reload that fails validation keeps the
Expand Down
1 change: 1 addition & 0 deletions deployments/compose/settings/config.json
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
"enabled": false,
"id_field": "event_id",
"require_id": false,
"retention": "0",
"tables": {}
},
"dlq": {
Expand Down
20 changes: 12 additions & 8 deletions docs/src/content/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -285,7 +285,7 @@ The body is a **flat JSON object** whose keys must match column names in the tar
| 400 | `{"error":"invalid json"}` | Malformed request body |
| 400 | `{"error":"unknown column ... for table ..."}` (also: `missing required column ...`, `type mismatch for column ...`, `null value for non-nullable column ...`) | Schema validation failure (unknown fields, type mismatches, missing required columns, null in a non-nullable column with no default). The body is the validator's message verbatim — there is no `validation failed:` prefix. |
| 400 | `{"error":"column \"x\" of table \"t\" is materialized and cannot be inserted"}` (also `… is alias …`) | The record supplies a value for a column ClickHouse computes. Omit it — the server fills it in. Refused rather than dropped: the published row has one slot per insertable column, so the value would otherwise vanish behind a `200` |
| 400 | `{"error":"missing dedupe id field \"event_id\""}` | Only when dedupe is enabled with `dedupe.require_id: true` and the row lacks the configured `id_field`. With `require_id: false` (the default) the row is instead published un-deduped. Either way — reject or publish — the row is logged at `WARN` and counted by `wavehouse_ingest_dedupe_missing_id_total`. In a batch this is a per-record failure, not a whole-request error. |
| 400 | `{"error":"missing dedupe id field \"event_id\""}` | Only when dedupe is enabled with `dedupe.require_id: true` and the row lacks the configured `id_field` or sets it to `null`. With `require_id: false` (the default) the row is instead published un-deduped. Either way — reject or publish — the row is logged at `WARN` and counted by `wavehouse_ingest_dedupe_missing_id_total`. In a batch this is a per-record failure, not a whole-request error. |
| 401 | `{"error":"invalid token"}` / `{"error":"token expired"}` | A present-but-invalid/expired token was supplied and denied (the gate surfaces the token reason rather than silently falling back to `default_role`) |
| 403 | `{"error":"forbidden"}` (empty-role variant: `forbidden: request has no role and no public default_role is configured`) | The resolved role lacks `insert` on the table |
| 403 | `{"error":"column \"x\" not allowed for insert"}` | The record names a column the role's `allow_columns`/`deny_columns` forbids ([Access control → Column permissions](/access-control#column-permissions)) |
Expand All @@ -295,10 +295,12 @@ The body is a **flat JSON object** whose keys must match column names in the tar
| 413 | `{"error":"request body exceeded 16777216 bytes"}` | Request body over the 16 MiB cap |
| 415 | `{"error":"no Content-Type: ingest requires one of application/json, application/x-ndjson, …"}` (declared variant: `Content-Type "text/plain": ingest requires one of …` — see the note above on how declarations are echoed; conflicting variant: `conflicting Content-Type declarations "application/json", "application/x-ndjson": ingest reads one format per request, and requires one of …`) | The request declared no `Content-Type`, one whose media type is unsupported or does not parse, a comma-bearing value that does not parse as a single media type, or repeated header lines that disagree — different formats, or one supported and one not. Checked before the body is parsed |
| 500 | `{"error":"dedupe failed"}` | Deduplication backend error |
| 503 | `{"error":"dedupe store unavailable"}` | Dedupe is on and its store cannot answer now: it is not open (for example, it failed to open on a reload), or a DynamoDB table is throttling, timing out or unreachable; `Retry-After: 5`. Nothing was published, so the retry is safe |
| 503 | `{"error":"schema not loaded yet"}` | The tenant's first schema discovery has not succeeded yet (its ClickHouse unreachable, or [no pool for it](/settings-directory#clickhouse)), so whether the table exists is not known; `Retry-After: 5`. Decided before the body is read |
| 500 | `{"error":"publish failed"}` | Message queue error |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure, for that tenant alone) or not open (see [Message Queue](/settings-directory#message-queue)). Response includes `Retry-After: 30` header. |
| 503 | `{"error":"service unavailable"}` | The message queue could not be reached or did not answer in time (a transient broker failure, not a full queue). Response includes `Retry-After: 5` header. Reserved for an external broker ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)): the embedded broker never reports this, and its publish failures are the `500` above. |
| 500 | `{"error":"publish failed"}` | Message queue error whose outcome is unknown, other than a full queue or an unreachable broker (below): the event may have been stored. With dedupe on, the record's id is left to lapse with the dedupe lease ([`dedupe.lease`](/configuration#dedupe), 30 seconds by default) rather than given back: a retry inside the lease answers the in-flight `503`, and one after it is published under the same idempotency key, which the queue drops if the first copy was stored. The queue's duplicate window (two minutes) covers up to ~2×lease plus a margin, not just the lease itself, so a retry timed off `Retry-After` anywhere in this flow stores no second copy; a much later one is stored again. |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure, for that tenant alone) or not open (see [Message Queue](/settings-directory#message-queue)). Response includes `Retry-After: 30` header. With dedupe on, the record's id is given back, so the retry is published rather than reported as a duplicate. |
| 503 | `{"error":"service unavailable"}` | The message queue could not be reached or did not answer in time (`mq.ErrUnavailable`, a transient broker failure, not a full queue) — reserved for an external broker ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)): the embedded broker never reports this, and its publish failures are the `500` above. As for the `500`, the record's id is left to lapse rather than given back, so a retry cannot land as a second copy; `Retry-After` is that lease, rounded up to whole seconds, when dedupe was on for the record, else the flat `Retry-After: 5`. |
| 503 | `{"error":"a request with the same dedupe id is in flight"}` | Dedupe is on and another request carrying the same id is still being published — usually a client's timeout-retry racing its own original. Its outcome decides whether this record is a duplicate, so retry after the `Retry-After` header (the dedupe lease, [`dedupe.lease`](/configuration#dedupe), 30 seconds by default). |
| 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) |

**curl example:**
Expand Down Expand Up @@ -408,13 +410,15 @@ A `200` is returned whenever the body was read and the records were processed
| 403 | `{"error":"forbidden"}` (empty-role variant: `forbidden: request has no role and no public default_role is configured`) | The resolved role lacks `insert` on the table (checked once, before any record) |
| 413 | `{"error":"request body exceeded 16777216 bytes"}` | Request body over the 16 MiB cap |
| 415 | `{"error":"no Content-Type: ingest requires one of application/json, application/x-ndjson, …"}` (declared variant: `Content-Type "text/plain": ingest requires one of …` — see the note above on how declarations are echoed; conflicting variant: `conflicting Content-Type declarations "application/json", "application/x-ndjson": ingest reads one format per request, and requires one of …`) | The request declared no `Content-Type`, one whose media type is unsupported or does not parse, a comma-bearing value that does not parse as a single media type, or repeated header lines that disagree — different formats, or one supported and one not. Checked before the body is parsed |
| 500 | `{"error":"publish failed"}` / `{"error":"dedupe failed"}` | Message-queue or dedup-backend failure mid-batch |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure) or not open, mid-batch; includes `Retry-After: 30` |
| 503 | `{"error":"service unavailable"}` | The message queue could not be reached or did not answer in time, mid-batch; includes `Retry-After: 5`. Reserved for an external broker ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)): the embedded broker never reports this, and its publish failures are the `500` above |
| 500 | `{"error":"publish failed"}` / `{"error":"dedupe failed"}` | Message-queue or dedup-backend failure mid-batch, other than a full queue or an unreachable broker (below). After a publish failure the records before it keep their ids, so a whole-batch retry reports those as duplicates; the failing record's id is left to lapse as on the single-object path, and the rest of its window's ids are given back |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure) or not open, mid-batch; includes `Retry-After: 30`. The records before the refused one keep their ids, and its id and the rest of its window's are given back |
| 503 | `{"error":"service unavailable"}` | The message queue could not be reached or did not answer in time (`mq.ErrUnavailable`), mid-batch — reserved for an external broker ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)): the embedded broker never reports this, and its publish failures are the `500` above. As for the `500`, the failing record's id is left to lapse rather than given back — so `Retry-After` is that record's dedupe lease, rounded up to whole seconds, when it was deduped; a record published un-deduped has no lapsing claim to wait out, so `Retry-After: 5` |
| 503 | `{"error":"dedupe store unavailable"}` | Dedupe is on and its store cannot answer now; `Retry-After: 5`. Nothing in the window being reserved was published; the windows before it were, and keep their ids |
| 503 | `{"error":"a request with the same dedupe id is in flight"}` | A record's dedupe id is held by another request still being published; includes `Retry-After` (the dedupe lease, [`dedupe.lease`](/configuration#dedupe), 30 seconds by default). Nothing in that record's window was published; the windows before it were |
| 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) |

:::caution[At-least-once on retry]
A batch aborted partway (a `503`/`500`, a JSON-array syntax error, or an NDJSON line over the 10 MiB line bound, after some leading records were already published) re-publishes those leading records when the whole batch is retried. A whole-body read failure is **not** one of these: a `413`, or the `400 invalid request body` of an upload cut off in transit, is decided before any record is processed, so nothing is published — safe to retry, once split for a `413`. Enable deduplication if duplicate suppression matters — this is the same at-least-once property the single-object path already has (the SDK retries both on `503`).
A batch aborted partway (a `503`/`500`, a JSON-array syntax error, or an NDJSON line over the 10 MiB line bound, after some leading records were already published) re-publishes those leading records when the whole batch is retried. Records are published in windows of 256, in order: a read error or a dedupe failure drops the open window unpublished, so what an aborted batch published is the windows before it, plus, after a publish failure, the records of its window before the failing one. A whole-body read failure is **not** one of these: a `413`, or the `400 invalid request body` of an upload cut off in transit, is decided before any record is processed, so nothing is published — safe to retry, once split for a `413`. Enable deduplication if duplicate suppression matters — this is the same at-least-once property the single-object path already has (the SDK retries both on `503`).
:::

**curl example (JSON array):**
Expand Down
Loading
Loading