Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared
- **`app/`** — the process wiring: `New` builds every component from the boot config and the settings directory (each one wired in one place — what it opens, what it loops, what it releases — with the settings registry handed to its wiring function whole, the injection point of the per-tenant registry of #583: store-keyed getters for the handlers, `perTenant` for the async paths (with the tenant each message's `mq.Topic` names for the stream hub and the ingest worker), the `chconn.Pools`, the per-tenant `discoveries` and the per-tenant keepalive wheels (`keepalives`: one `stream.Heartbeater` per served tenant, at its own `stream.keepalive_*`, turned by the `keepalive` component, [#597](https://github.com/Wave-RF/WaveHouse/issues/597)) reconciled from `AfterAdopt`, `gapWindows` handing the sweeper each tenant's own gap window (a rejected tenant's as its folder last had it, unbounded for one rejected since boot) and the `mq.max_bytes_gb` reconcile each served tenant's byte budget, and `defaultPolicy` for the one setting that still follows tenant `0`, a flat directory's ops-gate admin role; the auth verifiers are per tenant, reconfigured (rebuilt only on changed wiring) and pruned from `AfterAdopt`, the same hook's `Hub.Prune` ends the open streams of a tenant no longer served, and `wireCache`'s hook drops, through `LocalCache.Prune`, the cache version index of a tenant no longer served ([#262](https://github.com/Wave-RF/WaveHouse/issues/262))), `Run` drives the long-lived ones under one `errgroup` until the context is cancelled or one fails, `Close` releases them in reverse order. `New` wires only what the process's `roles` need (discovery, dedupe, auth verifiers, the hub bridge and keepalive per API process; the ingest worker per ingest process; the sweeper under its lease through `elected`, on the embedded MQ only); a process without `api` serves `api.NewOpsRouter` — probes, `/version`, metrics, and the settings reload behind the operator key alone. `cmd/wavehouse` and `tests/integration` both boot through it
- **`auth/`** — JWT auth middleware: HMAC **or** JWKS verification with `alg` pinned to the active verifier, role extraction from a configurable claim path; always runs, never rejects (bad token → empty role + stashed reason). One verifier per tenant ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 9): `Authenticator` keys them by `tenant.ID` — the request store's `settings.Store.Tenant()`, through an injected `TenantSource`; `tenant.Default` on the tenant-exempt routes — built from each tenant's `auth` block by `Reconfigure`, dropped by `Prune` once the tenant stops being served (rejected or removed), released by `Close`; the secrets (`Config`) are boot-level and shared. A JWKS key set is fetched off the boot and reload paths: until one has been stored the verifier is pending and a token-bearing request gets `503` + `Retry-After` from `api.refuseUnverifiable` (`auth.ErrVerifierPending`), never a `default_role` evaluation; refresh is library-managed (Eric, 2026-09-22), response capped at 1 MiB; the operator key's admin role is the request tenant's
- **`cache/`** — `Cache` interface → `LocalCache` (Ristretto: one pool for every tenant) + `VersionManager` (the invalidation index), and `RedisCache`, the Redis-compatible shared backend (random version tokens under the tenant's hash tag, one-round-trip lookups, bypass on failure behind a circuit breaker, deferred invalidations retried; selected by `cache.backend: redis`, configured by the boot config's `cache.redis` block — [#613](https://github.com/Wave-RF/WaveHouse/issues/613)). Every key carries the tenant (in `RedisCache`, after the key prefix: `<prefix>:{<tenant>}:…` for a version token, `<prefix>:q:<tenant>:…` for a value); in `LocalCache` and the version index it leads ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 8) — `<tenant>:query:<sha>` for the caller's query key and its singleflight, escaped whole as the lead field of the stored key `<query key>|<tenant>.<gen>|<dependency>|…`, where each raw table and scope name is escaped by `keyenc` (a `Namespace` carries them raw, so no caller escapes); the index holds a version per tenant, per (tenant, table) and per (tenant, table, scope), keyed by raw name and bumped in place (one entry per live namespace however often it is bumped, [#262](https://github.com/Wave-RF/WaveHouse/issues/262)) — so no cached read or coalesced flight crosses tenants, a bump through `Invalidate` names one tenant's namespaces and no other's, and `InvalidateTenant` drops the tenant's index so its next key gets a process-unique generation, orphaning its every cached result in one step, pipe results included (no insert reaches a pipe result until [#343](https://github.com/Wave-RF/WaveHouse/pull/343)); `Lookup` returns a `Snapshot` of the versions it read, taken before the handler chooses any input a bump invalidates — the tenant's connection included — and `Set` files the fill under it, so a write landing mid-query, or a reload moving the tenant to another address or database after the request took its connection, orphans the fill ([#382](https://github.com/Wave-RF/WaveHouse/issues/382)), and every backend runs the conformance suite `internal/testutil/cachetest`; 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 whole cache 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
- **`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, and drops a released tuple's, its idle connections closed, after the grace the tuple's pool gets (`Release`, handed the tuples `Reconcile` returns as released by the wiring's pools hook; the tuples with a zero `tls` block share one client, which stays; [#713](https://github.com/Wave-RF/WaveHouse/issues/713)). `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` (the in-process value by default; `mq.backend` also takes `nats`, with its `mq.nats` sub-block of file-path-only credentials; `coord.backend` takes `nats`, whose `coord.nats` block names only the lease bucket and rides `mq.nats`'s connection (both blocks, their rules and warnings are `mq_nats.go`); `cache.backend` takes `redis`, whose sub-block is `cache_redis.go`; and `dedupe.backend` takes `dynamodb`, with its `dedupe.dynamodb` sub-block) and `Warnings`, the valid combinations boot logs at `WARN`; `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; `coord.backend=nats` without `mq.backend=nats`; `mq.backend=nats` with `coord.backend=local` in a process running `ingest`; a process running only `sweeper` under `mq.backend=nats`) — boot is the validator, there is no dry run
- **`coord/`** — leases for work that must run in one process at a time (`Observer.Held` reads whether one is held without campaigning): `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: `coord.backend: nats` is `internal/mq/lease.go` (`ExternalNATS.Leases`), a key per lease in the operator's KV bucket, the KV revision as the fencing token, and expiry judged on the candidate's own clock (the same revision seen unchanged for 15s), never by a server TTL. `internal/app`'s `wireCoord` opens the one `coord.backend` selects and the sweeper runs through `RunElected` under the `sweeper` lease
Expand Down
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),

### Fixed

- **A ClickHouse tuple's HTTP clients are released with its pool** (`internal/chconn/chconn.go` (+ tests), `internal/app/{app,wire}.go` (+ tests), `internal/ingest/worker.go` (+ tests), `internal/api/query.go` (+ tests), `tests/integration/{ingest_outage,shard_order}_test.go`, `tests/e2e/sdk/{admin.test,settings}.ts`, `docs/src/content/docs/architecture.md`, `AGENTS.md`): closes [#713](https://github.com/Wave-RF/WaveHouse/issues/713), part of [#583](https://github.com/Wave-RF/WaveHouse/issues/583). The ingest worker and the `/v1/ops/query` proxy each keep one HTTP client per TLS config (`chconn.HTTPClients`), and a connection tuple with a `tls` block gets a config of its own each time its pool opens, so the client each had made for such a tuple stayed for the life of the process once the tuple was gone, its idle connections open until the transport's own 90-second idle timeout rather than closing with the pool: a test that opens and releases one tuple five times counted five clients, each still holding its connection after the native pool had closed. `Pools.Reconcile` now also returns the tuples it released, whether their tenants were removed or moved to another address, database, user or `tls` block, and the pools hook has both caches drop each one's client and close its idle connections after the grace the native pool gets, the longest `query_timeout` among the tenants the tuple had, so a request in flight finishes. A tuple another tenant still names keeps its client, and one that opens again gets a fresh one. The tuples with a zero `tls` block are unchanged: they share one client per cache, which stays.
- **`make ci` no longer kills healthy test runs at the `go test` time limits** (`Makefile`, `docs/src/content/docs/development.md`): its parallel phase runs the unit suite beside the lint and build jobs on every core, where `internal/app` and `internal/mq`, about 10s each alone, run past the suite's 15s per-package budget, so that phase now gives the suite 60s (`CI_UNIT_TIMEOUT`). `make test-unit` and the unit job in CI, which run the suite on its own, keep the 15s (`UNIT_TIMEOUT`). The integration suite's limit goes from 480s to 900s (`INTEGRATION_TIMEOUT`): `tests/integration` takes about 6m on a CI runner and up to 8m under Docker Desktop, which 480s left half a minute. Each limit is a Makefile variable, so one can be raised for a run (`make ci CI_UNIT_TIMEOUT=90s`).
- **Four `internal/mq` tests no longer fail when the embedded NATS server removes its streams directory under them** (`internal/mq/embedded_test.go`): after a stream is deleted, nats-server 2.14.6 removes its streams directory and the account's, once they are empty, on a goroutine of its own. A test that deleted every stream and then opened one raced that removal ([nats-io/nats-server#8725](https://github.com/nats-io/nats-server/issues/8725)), and the open was refused as `error creating store for stream`: rare (one failure in 6,400 local runs), but enough to fail the unit job on a pull request that changed no Go. Each of the four now keeps another stream on the server, so the directory is never empty, as `TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen` already did.
- **`make` pins the Go toolchain to the one `go.mod` names (its `toolchain` line, else its `go` line), so local runs and CI use the same Go** (`Makefile`, `scripts/ci/go-toolchain{,.test}.sh` (new), `.github/actions/setup-env/action.yml`, `.github/workflows/README.md`, `docs/src/content/docs/development.md`, `AGENTS.md`), closing [#580](https://github.com/Wave-RF/WaveHouse/issues/580): `GOTOOLCHAIN=auto` means the newer of the local Go and `go.mod`'s, so a machine on a newer Go built, tested and linted on a different toolchain than CI, and `make lint-go` panicked there (golangci-lint is built with an older Go and cannot type-check a newer standard library). The Makefile now exports `GOTOOLCHAIN` derived from `go.mod`, and the setup-env action exports the same value to `$GITHUB_ENV` for CI steps that run Go outside `make`, which also keeps a future runner image with a newer Go from taking over. Raising the `go` directive to a new minor version needs a golangci-lint built with that minor or newer (otherwise `make lint-go` refuses to load its config); patch bumps don't. One script, `scripts/ci/go-toolchain.sh`, derives the pin for both the Makefile and CI (with a fixture test in `make verify`): it reads `go.mod`'s `toolchain` line, else its `go` line, and fails with a message unless that is `x.y.z`, so bump with `go get go@1.N.P` rather than `go mod edit -go=1.27`. `GOTOOLCHAIN=local` in the environment is ignored; use `make GOTOOLCHAIN=local ...`.
Expand Down
Loading
Loading