Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
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
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
b83e670
feat(dedupe): retention per tenant and table, and an expiry sweep
EricAndrechek Sep 25, 2026
6c74b65
test(dedupe): pin that a commit mid-sweep-chunk survives; docs wording
EricAndrechek Sep 25, 2026
fc4085b
docs(durability): describe the benchmark machine neutrally
EricAndrechek Sep 25, 2026
9a7f849
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
7b9595e
feat(settings): a missing dedupe.retention keeps ids forever
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
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
5b723f0
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
935638e
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
fc0d509
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
ca48065
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 25, 2026
c5e4f3b
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 25, 2026
5d77211
test(dedupe): a legacy value under a current key reads as absent
EricAndrechek Sep 25, 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
29ddf36
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
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
fe44591
Merge branch 'feat/dedupe-reserve' into feat/dedupe-windowed-ingest
EricAndrechek Sep 26, 2026
af6205d
Merge branch 'feat/dedupe-windowed-ingest' into feat/dedupe-retention
EricAndrechek Sep 26, 2026
5e1ff2d
test(dedupe): replace the sweep-lock timing assertion with structural…
EricAndrechek Sep 26, 2026
51716a2
test(api): pin the ERROR/DEBUG split on a Reserve failure by exact msg
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
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
6 changes: 3 additions & 3 deletions AGENTS.md

Large diffs are not rendered by default.

6 changes: 4 additions & 2 deletions CHANGELOG.md

Large diffs are not rendered by default.

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
6 changes: 3 additions & 3 deletions config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -81,12 +81,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
16 changes: 9 additions & 7 deletions docs/src/content/docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -293,11 +293,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 is not open (for example, it failed to open on a reload); `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. With dedupe on, the record's id is given back, so a retry is published rather than reported as a duplicate — but if the publish reached the broker before failing (e.g. a client disconnect after the embedded broker had already stored the message), that retry can publish a second copy ([#629](https://github.com/Wave-RF/WaveHouse/pull/629) closes this with an idempotency key). |
| 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 (30 seconds) 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":"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, 30 seconds). |
| 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. As for `publish failed`, the record's id is given back so the retry can publish, and the same uncertain-publish caveat applies — the broker may already have stored the event before the timeout. |
| 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":"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 @@ -407,14 +408,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. After a publish failure the failing record's id is given back and the records before it keep theirs, so a whole-batch retry reports those as duplicates and publishes the rest — but if the failing record's publish reached the broker before failing (e.g. a client disconnect after the embedded broker had already stored it), that retry can publish a second copy of it ([#629](https://github.com/Wave-RF/WaveHouse/pull/629) closes this with an idempotency key) |
| 503 | `{"error":"service unavailable"}` | The tenant's ingest queue is full (backpressure) or not open, mid-batch; includes `Retry-After: 30`. As for `publish failed`, the failing record's id is given back and the records before it keep theirs |
| 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, 30 seconds). The records before it were published |
| 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. As for `publish failed`, the failing record's id is given back, and the same uncertain-publish caveat applies |
| 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, 30 seconds). 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