Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
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
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
a8e43bd
fix(app): retry a failed DynamoDB table check; refuse no region
EricAndrechek Sep 25, 2026
0b60451
docs(dedupe): name every region source; untangle the dynamodb check c…
EricAndrechek Sep 25, 2026
841a2f6
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
c11133d
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
48a6a77
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
1526c2b
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
484c72e
docs(changelog): the boot-backends entry predates dedupe.dynamodb
EricAndrechek Sep 25, 2026
59b0f0c
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 25, 2026
e61307b
Merge remote-tracking branch 'origin/main' into feat/dedupe-dynamodb-…
EricAndrechek Sep 26, 2026
f40c76a
fix(config): dedupe defaults in defaults(); refuse an explicit zero
EricAndrechek Sep 26, 2026
8b616b4
fix(config): cap dedupe.lease at 59.5s over the embedded queue
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
56727fb
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 26, 2026
97197f2
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
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
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
4424338
Merge branch 'feat/dedupe-dynamodb' into feat/dedupe-dynamodb-wiring
EricAndrechek Sep 26, 2026
7f38b9b
fix(app): move DynamoDB dedupe wiring out of wire.go for the e2e gate
EricAndrechek Sep 26, 2026
138b166
feat(app): refuse boot only for a misconfigured DynamoDB table in use
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
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: 6 additions & 0 deletions .testcoverage.yml
Original file line number Diff line number Diff line change
Expand Up @@ -84,3 +84,9 @@ exclude:
# (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$
4 changes: 2 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,9 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared
- **`cache/`** — `Cache` interface → `LocalCache` (Ristretto: one pool for every tenant) + `VersionManager` (the invalidation index). Every key leads with the tenant ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 8) — `<tenant>:query:<sha>` for a result and its singleflight, `<tenant>.<tenant version>.<table>.<table version>.<scope>` for a namespace — so no cached read or coalesced flight crosses tenants, a bump through `Invalidate` names one tenant's namespaces and no other's, and `InvalidateTenant` advances the tenant version that leads every namespace key of one tenant, orphaning every cached query keyed by its tables in one step (a pipe result names no table and keeps its TTL, [#343](https://github.com/Wave-RF/WaveHouse/pull/343)); the one crossing is the wiring's, above the package: `internal/app` hands the ingest worker the cache through `sharedTables`, which repeats each of the worker's bumps under every tenant on the same ClickHouse address and database (`chconn.Pools.SharingTables`, whatever their user or tls block — they read the same tables), and orphans the table-keyed cache (the structured-query results) of a tenant back on a pool after an absence, since it was out of that fan-out while away, or moved to another address or database, since it now reads other tables (story 6)
- **`chconn/`** — `Pools`, one `Manager` (a `driver.Conn`) per distinct `Identity{Addr, Database, Username, Password, TLS}` tuple among the served tenants, reconciled from the settings registry's `AfterAdopt` after every reload ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6): tenants naming one tuple share its pool, sized to their largest `max_open_conns`/`max_idle_conns`; a tenant whose tuple changed is repointed; a tuple no tenant names is released after the longest `query_timeout` among the tenants it had (never dials; a resize swaps the connection with the same grace). The boot config's `clickhouse.max_total_conns` bounds the open pools' `max_open_conns` together: boot refuses naming sum and ceiling; at a reload a resize above it keeps the pool's size, and a tuple that cannot be opened (the ceiling, an unreadable certificate, or options the driver refuses) leaves its tenants on the pool they had or on none — logged, retried by the next reload. Every consumer resolves its tenant's pool per call: `For` (nil for a tenant on no pool, a `503`), `Target` (the tenant's own HTTP wiring over its pool's TLS config), `SharingTables`, `Ping` (every pool at once, ready at the first answer). `HTTPClients` keeps one `http.Client` per TLS config. `Classify` (`errclass.go`) says what a failed ClickHouse request means for the request — `Unavailable`, `Denied`, `Rejected` (any unlisted exception code: the server read it and refused it), or `Unknown` (no code, no recognizable transport failure) — over the driver's error types and the HTTP interface's `HTTPError`; the ingest worker and the query handlers (`api/ch_errors.go` `writeCHError`, [#403](https://github.com/Wave-RF/WaveHouse/issues/403), [#271](https://github.com/Wave-RF/WaveHouse/issues/271)) both use it
- **`chsql/`** — dependency-free ClickHouse SQL helpers shared by `query`/`policy` (avoids an import cycle): `QuoteIdent` (backtick-quote every identifier) + `BindUnsafe` (reject names with a literal `?`)
- **`config/`** — YAML + env var config loading (cleanenv); strict on both sides (undeclared YAML key, unbound `WH_*` variable) and probes `data_dir` writability when a selected backend keeps state there (`NeedsDataDir`); `backends.go` holds each layer's `<layer>.backend` (only the in-process value today); `config.go` holds `roles` (`Has(Role)`) and `instance_id`, and `Validate` refuses a role split the backends cannot serve (any split over the embedded MQ; `api` without `ingest`, or the reverse, over a local cache) — boot is the validator, there is no dry run
- **`config/`** — YAML + env var config loading (cleanenv); strict on both sides (undeclared YAML key, unbound `WH_*` variable) and probes `data_dir` writability when a selected backend keeps state there (`NeedsDataDir`); `backends.go` holds each layer's `<layer>.backend` (the in-process value by default; `dedupe.backend` also takes `dynamodb`, with its `dedupe.dynamodb` sub-block); `config.go` holds `roles` (`Has(Role)`) and `instance_id`, and `Validate` refuses a role split the backends cannot serve (any split over the embedded MQ; `api` without `ingest`, or the reverse, over a local cache) — boot is the validator, there is no dry run
- **`coord/`** — leases for work that must run in one process at a time: `Coordinator.TryAcquire(ctx, name)` → a `Term` (fencing `Token`, strictly increasing per name; `Done`/`Err`, `ErrLost` on loss; `Resign`), `ErrHeld` while another holder's — or this coordinator's own — term is live; `RunElected` runs a loop only while holding its lease, resigning when the loop returns and campaigning again every `RetryPeriod`. `Local` is the in-process implementation (first taker wins, never expires; `Peer` is a second handle over the same table for tests); every implementation runs `coordtest.Conformance`. Imports only the standard library, so a distributed backend lives beside its connection (NATS KV in `internal/mq`). `internal/app`'s `wireCoord` opens the one `coord.backend` selects and the sweeper runs through `RunElected` under the `sweeper` lease
- **`dedupe/`** — `Deduplicator` interface (two-phase `Reserve`/`Commit`/`Release` over `Key{Table, ID}`; every backend passes the `dedupetest` conformance suite) → `Embedded` (Pebble: every tenant's seen ids in one instance at `data_dir/pebble`, each key led by its tenant and table, pending claims in memory, open while any tenant's store is — the layout is the implementation's call, and the wiring hands it `data_dir` once; its `Stats` feed the system gauges) or `Dynamo` (one shared DynamoDB table, conditional `PutItem` claims; built and conformance-tested against dynamodb-local but not yet selectable at boot), wrapped by `Managed` whose open/closed state follows the hot-reloadable `dedupe.enabled` in the settings directory's `config.json`; `Stores` holds one `Managed` per tenant, built through a `Factory` (`func(tenant.ID) *Managed`, `Embedded.Tenant` in production; `Managed` opens its store through a function, so every backend gets the same switch), and reconciled from the registry's `AfterAdopt` hook — open exactly when the tenant is served with its switch on, closed with its seen ids kept otherwise ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) stories 7 and 3)
- **`dedupe/`** — `Deduplicator` interface (two-phase `Reserve`/`Commit`/`Release` over `Key{Table, ID}`; every backend passes the `dedupetest` conformance suite) → `Embedded` (Pebble: every tenant's seen ids in one instance at `data_dir/pebble`, each key led by its tenant and table, pending claims in memory, open while any tenant's store is — the layout is the implementation's call, and the wiring hands it `data_dir` once; its `Stats` feed the system gauges) or `Dynamo` (one shared DynamoDB table, conditional `PutItem` claims; conformance-tested against dynamodb-local, selected by `dedupe.backend: dynamodb`; boot checks the table and never creates it outside dynamodb-local), wrapped by `Managed` whose open/closed state follows the hot-reloadable `dedupe.enabled` in the settings directory's `config.json`; `Stores` holds one `Managed` per tenant, built through a `Factory` (`func(tenant.ID) *Managed`, `Embedded.Tenant` or, gated on the table check (`Factory.Gated`), `Dynamo.Tenant` in production; `Managed` opens its store through a function, so every backend gets the same switch), and reconciled from the registry's `AfterAdopt` hook — open exactly when the tenant is served with its switch on, closed with its seen ids kept otherwise ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) stories 7 and 3)
- **`discovery/`** — `SchemaRegistry`, one per served tenant over a `Source` read once per refresh — the tenant's pool's connection and the database that pool was opened for, one snapshot, so a refused move keeps discovering the database the tenant's queries still use (`internal/app`'s `discoveries` builds, runs and stops them from `AfterAdopt` and `App.Close`: `RetryRefresh` until the first success, then `StartAutoRefresh` with a random first tick; `Lookup` answers `ErrNotLoaded` before the first success — the handlers' `503` with `Retry-After` — and `ErrUnknownTable` after; a failed loop attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`), that introspects ClickHouse `system.columns` (name/type/nullability plus `default_expression` and 1-based `position`) and `system.tables` (each table's `create_table_query`, kept in-process and never serialized — an external-engine table renders its wiring there unconditionally — endpoint, bucket/host, database, username, S3 access key id; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so the exposure is the topology, not the secret), records the server version, + `Validate()` for ingest payloads + `CanonicalizeTimestamps()` rewriting top-level `DateTime`/`DateTime64` column values to the canonical RFC 3339 UTC wire form pre-publish (Key Design Decision #19)
- **`ingest/`** — Ingest worker pipeline (`worker.go`: JetStream input → per-table batch INSERT with DLQ output). The pipeline is **insert-only**. The wire format `EventMessage` (`types.go`) carries `{table_name, scope, received_timestamp, format, columns, row}` and nothing else — `row` is one positional `JSONCompactEachRow` line and `columns` names its slots, the table's **insertable** columns (a `MATERIALIZED`/`ALIAS` column cannot be named in an `INSERT`); the worker batches per (tenant, table, column list), the tenant read off each message's `mq.Topic`, and inserts each batch into its tenant's own ClickHouse (`chconn.Pools.Target`); the worker accepts whatever table name the envelope carries (table existence was already checked by the HTTP ingest handler, which `404`s an unknown table before publish; the worker doesn't re-validate), then bulk-INSERTs. In the embedded-NATS deployment (the default), the server runs with `DontListen: true` (`internal/mq/embedded.go`), so the only Publishers reachable on the `ingest.>` subjects are in-process Go code — today, only the HTTP `/v1/ingest?table={table}` handler. Non-insert mutations (`DELETE`/`UPDATE`/`TRUNCATE`/…) must go through `POST /v1/ops/query` under the admin role (the same `RequireAdmin` gate as the rest of `/v1/ops/*`), so non-admin callers never reach the proxy. A request with no token (or an invalid one) resolves to the `default_role`, which in a production config is not the admin role (setting them equal is a loudly-warned dev-only setting), so it can't reach this endpoint. Plus `Sweeper` (Active Sweeper for NATS message lifecycle) + `EventMessage`/`BufferConsumerName` types (`types.go`)
- **`keyenc/`** — the one escaping composite keys are built from: `Escape` keeps `[A-Za-z0-9_-]` (exactly the tenant-id grammar, so a tenant id is its own escaped form) and writes every other byte as `%XX`, `Unescape` is `url.PathUnescape` (lenient: either hex case, and a byte left unescaped reads as itself, so a `%2D` an earlier build wrote still reads), `Join`/`AppendJoin` escape each field and put a separator between them (they panic on no fields, and on a separator the escaping could write or one outside ASCII) and `Split` reverses them. NATS subjects (`Join`/`Split` after the verbatim tenant), the cache's namespace tokens and the dedupe keys (`<tenant>/<table>/<id>`) use it; changing what it keeps orphans every stored key
Expand Down
Loading
Loading