From 49386a494c3f678bf9caa7ab88be944b8f71832a Mon Sep 17 00:00:00 2001 From: taitelee Date: Tue, 29 Sep 2026 12:59:32 -0400 Subject: [PATCH 1/4] fix(app): rediscover a tenant's schema when a reload moves it (#638) --- AGENTS.md | 2 +- CHANGELOG.md | 1 + docs/src/content/docs/architecture.md | 2 +- docs/src/content/docs/settings-directory.mdx | 2 +- internal/app/app_test.go | 166 +++++++++++++++++++ internal/app/discoveries.go | 26 ++- internal/app/roles_test.go | 17 ++ internal/app/wire.go | 18 +- internal/discovery/discovery.go | 3 +- internal/testutil/testutil.go | 6 + 10 files changed, 232 insertions(+), 11 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 4573713a9..15a814b11 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -37,7 +37,7 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared - **`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 `.backend` (the in-process value by default; `cache.backend` also takes `redis`, whose sub-block is `cache_redis.go`, and `dedupe.backend` 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, committed ids stored with their expiry and deleted by an hourly background sweep along with the version-0 keys from before the table joined the key, 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) +- **`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`, and starts a tenant over on a fresh registry when a reload moves it to another address or database: `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) or through an operator-authored pipe that writes, gated only by its `allowed_roles` (#386). 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. The package that builds a key takes raw names and escapes them itself, so no caller has to and no field reaches a key unescaped: NATS subjects (`Join`/`Split` after the verbatim tenant) and the cache's keys — the version index and the shared backend's Redis keys (`internal/cache`) — and the dedupe keys (`//`) use it; changing what it keeps orphans every stored key (an orphaned dedupe key lets a seen id through again), and on the shared backend, whose keys every process builds for itself, splits them between builds for the length of a rolling upgrade (a bump one build makes misses the entries the other filed, served until their TTL) - **`mq/`** — the message-queue boundary: the **only** package that imports NATS/JetStream (Key Design Decision #20). and the only one that knows how the broker works. Everything else addresses events by `Topic{Tenant, Table, Scope}` (a validated tenant id and raw names — the tenant leads every subject, `ingest..
`, so one wildcard selects a tenant's traffic, and a topic without one is refused) and states intent through the interfaces — `Publisher` (`ErrQueueFull` is the backpressure signal, `ErrUnavailable` a broker that cannot be reached — both a `503`, with `Retry-After` `30` and `5`; `WithIdempotencyKey` makes a republish inside the queue's duplicate window a no-op), `Subscriber`, `ConsumerManager`/`Consumer`/`ConsumerConfig` (the ingest worker's durable consumer), `DeadLetterer` and `DeadLetterStats` (park a message, count what is parked), `Purger` (drop what is both acked and older than a cutoff — the sweeper), `Replayer` (SSE gap-fill) — composed into `Broker`, which adds each tenant's byte budget (`SetMaxBytes`/`MaxBytes`: the `mq.max_bytes_gb` reload, which opens a tenant's queue the first time) and `Stats` (the system gauges' source). Every interface speaks per tenant, never per stream: the embedded implementation gives each tenant a queue of its own (a stream pair, `INGEST_`/`DLQ_`), and nothing outside the package may assume that layout — an external implementation may keep one shared stream. Subjects, prefixes, wildcards, stream names, sequences, and ack floors are private to the one implementation, `EmbeddedNATS` (`embedded.go`, `subject.go`, `purge.go`, `deadletter.go`), whose subject tokens are escaped by the shared `internal/keyenc`; `internal/app` constructs it and hands everything else a `mq.Broker`. Every implementation passes the conformance suite in `internal/mq/mqtest` (`mqtest.Run`), which states the `Broker` contract as behavior; a new backend runs it from its own test, with `mqtest.Caps` only where its semantics legitimately differ diff --git a/CHANGELOG.md b/CHANGELOG.md index b0636db2a..c9d92ac90 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -90,6 +90,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Fixed +- **A tenant moved to another ClickHouse address or database discovers its new schema with the reload** (`internal/app/{wire,discoveries}.go` (+ tests), `internal/discovery/discovery.go` (comment), `internal/testutil/testutil.go`, `docs/src/content/docs/architecture.md`, `settings-directory.mdx`, `AGENTS.md`): closes [#638](https://github.com/Wave-RF/WaveHouse/issues/638). A reload that changed a tenant's `clickhouse.addr` or `clickhouse.database` orphaned the tenant's cache and left its schema registry as it was, so until the tenant's loop fired at `schema.refresh_interval` its queries and inserts were validated against the previous database's schema and run against the new one. The tenants the pools reconcile reports stale now have their registry dropped in the same hook, and the discovery reconcile that follows builds each a fresh one over the pool it is on, as it does for a tenant back after a rejection or removal: the first discovery runs at once in the tenant's own loop, so the reload never waits on ClickHouse, and until it succeeds the tenant's table lookups answer `503` with `Retry-After: 5`, as before any first discovery. A discovery that fails is logged with its tenant (`schema discovery retry failed`), counted in `wavehouse_schema_refresh_failures_total`, and retried with backoff from two seconds to sixty. A tenant whose address and database did not change keeps its registry and its loop, a flat directory's tenant `0` moves the same way, and a process without the api role, which discovers no schema, only repoints its pool. - **One tenant's failed queue join no longer ends ingest for every tenant** (`internal/mq/{mq,embedded}.go` (+ tests), `docs/src/content/docs/{architecture,durability,ingest-pipeline}.md`, `settings-directory.mdx`): closes [#675](https://github.com/Wave-RF/WaveHouse/issues/675). When a tenant's queue opened while the server ran and the ingest worker's consumer could not join it, the broker reported that as the consumer's delivery ending, so the worker failed and the process restarted — one tenant's failure taking ingest down for all of them — while the stream hub's consumer failing the same way was only logged, and the tenant's `GET /v1/stream` connections got no live rows until a restart. A queue a consumer cannot join is now left unrecorded as open, whichever consumer it is: `SetMaxBytes` returns the error and `MaxBytes` keeps not reporting the budget, so the next reload tries the consumers that had not joined again, and the tenant's publishes answer `503` and try them too, paced as a queue that cannot open already is (one shared attempt, then refusals for five seconds). No row is accepted that a consumer would not read, no other tenant is touched, and nothing restarts. A queue reopened after its ingest stream went missing has its consumers' delivery started again on the new stream, where the stopped delivery on the old one used to stand in for it. A consumer that had joined and whose delivery then ends on its own (a deleted durable, a closed connection) still fails the worker, and boot is unchanged: a flat directory whose queue cannot open refuses boot. - **Dedupe claims an id, publishes, then commits it — and keys it by tenant, table and id** (`internal/dedupe/{dedupe,key,embedded,managed}.go` (+ tests), `internal/dedupe/dedupetest/` (new), `internal/api/ingest.go` (+ tests), `internal/settings/validate_test.go`, `internal/keyenc/keyenc.go`, `internal/testutil/mocks.go`, `internal/app/app_test.go`, `AGENTS.md`, `docs/src/content/docs/{architecture,api,deployment,development}.md`, `settings-directory.mdx`, `sdk/reference.md`): `CheckAndMark` is replaced by a two-phase `Reserve` → `Commit` / `Release` contract with a lease on the pending claim, and every backend now runs one conformance suite. Four bugs go with it. Two concurrent requests carrying one id no longer both publish it: Pebble's check and claim happen under one lock, and the loser answers `503` with `Retry-After` while the winner is still publishing ([#390](https://github.com/Wave-RF/WaveHouse/issues/390)). A publish that fails gives its id back, so the retry a `503` asks for is published instead of skipped as a duplicate of a record that never reached the queue ([#384](https://github.com/Wave-RF/WaveHouse/issues/384)) — the residual case is a publish that fails after it already reached the broker (a timeout, a disconnect), where the released id lets the retry through but that retry publishes a genuine second copy; [#629](https://github.com/Wave-RF/WaveHouse/pull/629) closes that with an idempotency key. The same id in two tables is two ids ([#222](https://github.com/Wave-RF/WaveHouse/issues/222)'s keyspace half). An explicit `null` id is a missing id — rejected under `require_id`, published un-deduped otherwise — instead of the one id `""` that made every null record after the first a duplicate ([#370](https://github.com/Wave-RF/WaveHouse/issues/370)). **Upgrade:** the key layout changes, so an id seen before the upgrade is accepted once more after it; nothing is migrated, and the old keys are left in `/pebble`, unread, deleted by the retention sweep (see Added) ([#220](https://github.com/Wave-RF/WaveHouse/issues/220)) ([Deployment → Upgrading across the dedupe key change](https://github.com/Wave-RF/WaveHouse/blob/main/docs/src/content/docs/deployment.md#upgrading-across-the-dedupe-key-change)). The key is readable text, `/
/` (for example `acme/clicks/evt-123`), with the table and id escaped and joined by `internal/keyenc`, the escaping NATS subject tokens already use, so any table name gets a keyspace of its own, including one holding a NUL byte or a `/`. New metrics: `wavehouse_ingest_dedupe_commit_failed_total` (a published record whose id failed to commit; the claim lapses with its lease) and `wavehouse_dedupe_hashed_id_total` (an id over 1,024 bytes once escaped, stored as its SHA-256). - **Ingest runs in windows of 256 records over the dedupe contract, and a dedupe store that cannot answer is a `503`** (`internal/api/ingest.go` (+ tests), `internal/mq/{mq,embedded}.go` (+ tests), `internal/dedupe/key.go` (+ tests), `internal/testutil/mocks.go`, `AGENTS.md`, `docs/src/content/docs/{api,architecture,durability}.md`, `settings-directory.mdx`, `sdk/reference.md`): each window of a request is prepared, then reserved in one dedupe call, published in order, and committed in one call, so a batch costs one dedupe round trip per phase per window rather than per record — on Pebble, one commit `fsync` per window (a 1,000-record batch: four syncs instead of a thousand, 24 ms against 5.7 s of dedupe time measured with the queue stubbed). Every deduped record is published under an idempotency key (`mq.WithIdempotencyKey`, JetStream's message id, derived by `dedupe.IdempotencyKey`), and each tenant's ingest stream now keeps an explicit two-minute duplicate window, sized to `2 × the 30-second lease + 1s`: an uncertain publish's `503` sends the *full* lease as `Retry-After`, so an obedient client's retry can land up to ~2×lease after the original request, and the `+1s` covers a backend whose claim expiry itself rounds up by that much. That closes the last path of [#384](https://github.com/Wave-RF/WaveHouse/issues/384): a publish that fails with an unknown outcome (anything but a full queue) keeps its record's claim until the lease lapses instead of releasing it, and a retry after the lease but within two minutes of the first publish is dropped by the queue if the first copy was stored (a later one is stored again). A dedupe store that is not open or that reports `dedupe.ErrUnavailable` now answers `503 {"error":"dedupe store unavailable"}` with `Retry-After: 5`, which the SDK retries, rather than `500 dedupe failed`. A mid-body read error or dedupe failure now drops the open window unpublished, where records before it used to be published; `wavehouse_ingest_dedupe_commit_failed_total` and `wavehouse_ingest_dedupe_disabled_total` count records, as before, now added a window at a time. diff --git a/docs/src/content/docs/architecture.md b/docs/src/content/docs/architecture.md index 5c3093985..58843fcd7 100644 --- a/docs/src/content/docs/architecture.md +++ b/docs/src/content/docs/architecture.md @@ -153,7 +153,7 @@ The SSE fan-out, factored out of `api/` so the delivery hot path ([#294](https:/ ### `discovery/` — Schema Discovery & Validation -- **discovery.go** — `SchemaRegistry`, one per tenant since [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6, over a `Source` read once per refresh: the tenant's pool's connection and the database that pool was opened for, one snapshot (`App.discoverySource` over `chconn.Pools.For` in production, so a reload that repoints the tenant applies to the next refresh, and a refused move keeps discovering the database the tenant's queries and inserts still use; no pool is `ErrNoConnection`), queries `system.columns` to discover ClickHouse table schemas, keeping each column's `default_kind` so `IsInsertable` / `InsertableColumns` / `InsertableColumnNames` (memoized per table at refresh) can decide the insertable subset the ingest envelope and the SSE announcement are both built from. Each refresh also records the server version (`SELECT version()`), joins `system.tables` for each table's `create_table_query` (kept in-process as `TableSchema.DDL` and marked `json:"-"` — an external-engine table renders its wiring in that statement — endpoint, bucket/host, database, username, access key id — so it must never reach `/v1/ops/schema`; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so what is withheld here is the topology), reads each column's `default_expression` and 1-based `position` alongside its type, discovers the server's default time zone (`SELECT timezone()`) and bakes every `DateTime`/`DateTime64` column's canonicalization spec (precision + resolved zone) into the cached schema, so the per-record ingest path parses no type strings and loads no zones ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)). `Lookup` tells the two misses apart — `ErrNotLoaded` before the first successful refresh, `ErrUnknownTable` after — where `Get` answers nil for both (the stream hub's fail-closed reading). Supports periodic auto-refresh (`StartAutoRefresh`, the first tick at a random point within the interval so tenants adopted together do not refresh together, the cadence re-read after every tick), on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop, each sleep drawn uniformly below the backoff so instances retrying one ClickHouse do not retry in lockstep, used by `internal/app` so a transiently unreachable ClickHouse doesn't crash-loop the binary); a loop's failed attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`. Thread-safe via `sync.RWMutex`. +- **discovery.go** — `SchemaRegistry`, one per tenant since [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6, over a `Source` read once per refresh: the tenant's pool's connection and the database that pool was opened for, one snapshot (`App.discoverySource` over `chconn.Pools.For` in production, so a reload that repoints the tenant to another user or `tls` block, which reads the same tables, applies to the next refresh, one that moves it to another address or database has `internal/app`'s `discoveries` drop its registry and build a fresh one, whose first discovery runs at once and whose lookups answer `ErrNotLoaded` until it succeeds ([#638](https://github.com/Wave-RF/WaveHouse/issues/638)), and a refused move keeps discovering the database the tenant's queries and inserts still use; no pool is `ErrNoConnection`), queries `system.columns` to discover ClickHouse table schemas, keeping each column's `default_kind` so `IsInsertable` / `InsertableColumns` / `InsertableColumnNames` (memoized per table at refresh) can decide the insertable subset the ingest envelope and the SSE announcement are both built from. Each refresh also records the server version (`SELECT version()`), joins `system.tables` for each table's `create_table_query` (kept in-process as `TableSchema.DDL` and marked `json:"-"` — an external-engine table renders its wiring in that statement — endpoint, bucket/host, database, username, access key id — so it must never reach `/v1/ops/schema`; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so what is withheld here is the topology), reads each column's `default_expression` and 1-based `position` alongside its type, discovers the server's default time zone (`SELECT timezone()`) and bakes every `DateTime`/`DateTime64` column's canonicalization spec (precision + resolved zone) into the cached schema, so the per-record ingest path parses no type strings and loads no zones ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)). `Lookup` tells the two misses apart — `ErrNotLoaded` before the first successful refresh, `ErrUnknownTable` after — where `Get` answers nil for both (the stream hub's fail-closed reading). Supports periodic auto-refresh (`StartAutoRefresh`, the first tick at a random point within the interval so tenants adopted together do not refresh together, the cadence re-read after every tick), on-demand refresh, and `RetryRefresh` (boot-time exponential backoff loop, each sleep drawn uniformly below the backoff so instances retrying one ClickHouse do not retry in lockstep, used by `internal/app` so a transiently unreachable ClickHouse doesn't crash-loop the binary); a loop's failed attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`. Thread-safe via `sync.RWMutex`. - **timestamp.go** — `CanonicalizeTimestamps(schema, data)` rewrites every parseable value in a top-level `DateTime`/`DateTime64` column to the canonical RFC 3339 UTC wire form before the event is published ([#372](https://github.com/Wave-RF/WaveHouse/issues/372)): zone-less values are interpreted in the column's declared zone, else the discovered server default — ClickHouse's own rule, so the spelling changes but never the instant. Fail-open: an unparseable value or unresolvable zone passes through verbatim for ClickHouse's own parser to judge; ingest never rejects a record over its timestamp spelling. `Column.TimeParser()` exposes the same grammar as a value→instant parser (nil for a column with no resolved timestamp spec — a non-timestamp column, or one whose declared zone couldn't be loaded), which the stream row-filter uses so filter constants and canonicalized payloads can't disagree on the instant ([#381](https://github.com/Wave-RF/WaveHouse/issues/381)). - **validation.go** — `Validate(schema, data)` checks incoming JSON against the discovered schema: unknown fields, type compatibility, missing required columns, null handling. Also exports the type classifiers `IsNumericType` / `IsStringType` and the storage-model classifier `NumericStorageOf` (all unwrapping `Nullable`/`LowCardinality`; the latter yields a numeric column's float width, `Decimal` scale, or integer exactness), which — together with `Column.TimeParser` from timestamp.go — seed the stream row-filter's `policy.ColumnSpec` comparison. - **discovery_test.go** — Unit tests for validation logic. diff --git a/docs/src/content/docs/settings-directory.mdx b/docs/src/content/docs/settings-directory.mdx index 5d54d7ba0..8d2a6a867 100644 --- a/docs/src/content/docs/settings-directory.mdx +++ b/docs/src/content/docs/settings-directory.mdx @@ -205,7 +205,7 @@ The `clickhouse` block is the connection wiring, minus the password. A reload ap **Per tenant.** Over a [nested settings directory](/deployment#the-nested-settings-directory) each folder's `clickhouse` block is that tenant's own. The process opens one native connection pool per distinct `addr`, `database`, `username`, password and `tls` tuple among the tenants being served — ClickHouse authenticates per connection, so tenants with different credentials never share one — and tenants naming the same tuple share one pool, sized to the largest `max_open_conns` and the largest `max_idle_conns` among them. `http_port`, `http_scheme`, `headers` and `query_timeout` stay each tenant's own whatever it shares: the HTTP hop is addressed per request and the deadline is per query. A reload reconciles the pools: a tenant whose tuple changed moves to the pool that tuple names, opened if it is new; a pool no tenant names any more closes once the longest `query_timeout` among the tenants it had has passed; a pool whose largest ask changed is resized the same way. The boot config's [`clickhouse.max_total_conns`](/configuration#clickhouse) bounds the open pools **together**: at boot, pools that would add up to more refuse to start, naming the sum and the ceiling; on a reload, a resize above it is refused and the pool keeps its size, and a tuple that cannot be opened — the ceiling, a `tls` block whose files cannot be loaded, or options the driver refuses — leaves its tenants on the pool they had, or on none for a tenant that had none. A tenant on no pool answers `503` with `Retry-After` on every route that reaches its ClickHouse until a reload opens one. The refusals are `ERROR` log lines, not findings, and the next reload retries. Schema discovery is per tenant too: each tenant's tables are discovered from its own `database` over its own pool, on its own `schema.refresh_interval`, with the first periodic refresh at a random point within the interval so tenants adopted together do not refresh together; until a tenant's first discovery, its table lookups answer `503` with `Retry-After` rather than `404` — the table may well exist. Until [#529](https://github.com/Wave-RF/WaveHouse/issues/529) every tenant's user authenticates with the one `WH_CH_PASSWORD`, so the tuple is in effect the address, database, user and `tls` block. -Moving to a different ClickHouse leaves ingest traffic unaffected: events land in the embedded queue regardless, and the worker's next flush uses the new target. +Moving to a different ClickHouse address or database starts the tenant's schema discovery over with the reload, not at its next `schema.refresh_interval`: its table lookups answer `503` with `Retry-After` until the new database's tables are discovered, retried with backoff while it cannot be reached, and are never answered from the previous database's schema. A tenant whose address and database did not change keeps its schema. Events already accepted stay in the queue, and the worker's next flush uses the new target. ## Authentication diff --git a/internal/app/app_test.go b/internal/app/app_test.go index 80a2e8c58..9d761fc86 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -22,6 +22,7 @@ import ( "testing" "time" + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" "github.com/golang-jwt/jwt/v5" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -32,6 +33,7 @@ import ( "github.com/Wave-RF/WaveHouse/internal/coord" "github.com/Wave-RF/WaveHouse/internal/dedupe" "github.com/Wave-RF/WaveHouse/internal/dedupe/dedupetest" + "github.com/Wave-RF/WaveHouse/internal/discovery" "github.com/Wave-RF/WaveHouse/internal/mq" "github.com/Wave-RF/WaveHouse/internal/settings" "github.com/Wave-RF/WaveHouse/internal/tenant" @@ -1801,6 +1803,170 @@ func TestNew_SchemaNotLoadedIs503(t *testing.T) { }) } +// fakeClickHouse stands in for the ClickHouse servers: a registry built +// after newFakeClickHouse discovers the tables named for the database its +// tenant's pool was opened for, and from that pool itself — a closed port — +// for a database with none named. +type fakeClickHouse struct { + mu sync.Mutex + conns map[string]driver.Conn +} + +// set names the tables of database. +func (f *fakeClickHouse) set(database string, tables ...string) { + schemas := make([]*discovery.TableSchema, 0, len(tables)) + for _, name := range tables { + schemas = append(schemas, &discovery.TableSchema{Name: name, Columns: []discovery.Column{{Name: "id", Type: "String"}}}) + } + f.mu.Lock() + defer f.mu.Unlock() + f.conns[database] = testutil.NewSchemaConn(schemas) +} + +func (f *fakeClickHouse) conn(database string) driver.Conn { + f.mu.Lock() + defer f.mu.Unlock() + return f.conns[database] +} + +// newFakeClickHouse starts every served tenant's discovery over against +// databases (database → its tables) and waits for each to load. The source +// stays the production one, the fake replacing only the connection. +func newFakeClickHouse(t *testing.T, a *App, databases map[string][]string) *fakeClickHouse { + t.Helper() + f := &fakeClickHouse{conns: map[string]driver.Conn{}} + for database, tables := range databases { + f.set(database, tables...) + } + a.discoveries.build = func(id tenant.ID, _ *settings.Store) *discovery.SchemaRegistry { + source := a.discoverySource(id) + return discovery.NewSchemaRegistry(func() (driver.Conn, string) { + conn, database := source() + if fake := f.conn(database); conn != nil && fake != nil { + return fake, database + } + return conn, database + }, id, perTenant(a.tenants, (*settings.Store).SchemaRefreshInterval)) + } + var served []tenant.ID + for id := range a.tenants.All() { + served = append(served, id) + } + a.discoveries.drop(served) + a.discoveries.reconcile(a.tenants) + for _, id := range served { + require.Eventually(t, a.discoveries.For(id).Loaded, 5*time.Second, 10*time.Millisecond, "tenant %s", id) + } + return f +} + +// databaseSettings is chSettings with the database named. +func databaseSettings(addr, database string) map[string]any { + p := chSettings(addr, "default", 10) + p["clickhouse"].(map[string]any)["database"] = database + return p +} + +// A reload that moves a tenant to another ClickHouse database starts its +// schema discovery over (#638): the schema it held is dropped with the +// reload, so no lookup is answered from the previous database's, and the new +// database's tables are discovered at once rather than at the next +// schema.refresh_interval (60 seconds here). The tenant beside it, and a +// reload that moves nobody, keep the registry and the loop they had. +func TestReload_MovedTenantDiscoversTheNewDatabase(t *testing.T) { + addr := closedAddr(t) + root := writeNestedSettings(t, map[string]map[string]any{ + "acme": databaseSettings(addr, "default"), + "globex": databaseSettings(addr, "default"), + }) + a := newApp(t, testConfig(t, root), Options{}) + newFakeClickHouse(t, a, map[string][]string{"default": {"events"}, "moved_db": {"orders"}}) + acme, globex := a.discoveries.For("acme"), a.discoveries.For("globex") + require.NotNil(t, acme.Get("events")) + loops := *a.discoveries.cur.Load() + + _, adopted := a.tenants.Reload("test") + require.True(t, adopted) + assert.Same(t, acme, a.discoveries.For("acme"), "nobody moved, nobody starts over") + assert.Same(t, globex, a.discoveries.For("globex")) + + rewriteSettings(t, filepath.Join(root, "acme"), databaseSettings(addr, "moved_db")) + _, adopted = a.tenants.Reload("test") + require.True(t, adopted) + moved := a.discoveries.For("acme") + require.NotNil(t, moved) + assert.NotSame(t, acme, moved, "a fresh registry") + assert.Nil(t, moved.Get("events"), "never the previous database's schema") + select { + case <-loops["acme"].done: + case <-time.After(5 * time.Second): + t.Fatal("the loop refreshing the previous registry did not stop") + } + require.Eventually(t, moved.Loaded, 5*time.Second, 10*time.Millisecond, "discovered without waiting for schema.refresh_interval") + _, err := moved.Lookup("orders") + require.NoError(t, err, "the new database's table") + _, err = moved.Lookup("events") + require.ErrorIs(t, err, discovery.ErrUnknownTable, "the previous database's table") + + assert.Same(t, globex, a.discoveries.For("globex"), "the tenant that stayed is untouched") + assert.Same(t, loops["globex"], (*a.discoveries.cur.Load())["globex"], "its loop too") + assert.NotNil(t, globex.Get("events")) +} + +// A flat directory's tenant 0 moves the same way. +func TestReload_MovedTenantDiscoversTheNewDatabase_Flat(t *testing.T) { + addr := closedAddr(t) + dir := writeSettings(t, databaseSettings(addr, "default")) + a := newApp(t, testConfig(t, dir), Options{}) + newFakeClickHouse(t, a, map[string][]string{"default": {"events"}, "moved_db": {"orders"}}) + before := a.Registry() + require.NotNil(t, before.Get("events")) + + rewriteSettings(t, dir, databaseSettings(addr, "moved_db")) + _, adopted := a.tenants.Reload("test") + require.True(t, adopted) + moved := a.Registry() + require.NotNil(t, moved) + assert.NotSame(t, before, moved) + assert.Nil(t, moved.Get("events"), "never the previous database's schema") + require.Eventually(t, moved.Loaded, 5*time.Second, 10*time.Millisecond) + assert.NotNil(t, moved.Get("orders")) +} + +// A moved tenant whose first discovery of the new database fails answers +// like any tenant before its first discovery — 503 with Retry-After, for a +// table the previous database had too — with the failure logged under its +// tenant, and its loop keeps retrying until the database answers. +func TestReload_MovedTenantFailedDiscoveryIsRetried(t *testing.T) { + addr := closedAddr(t) + root := writeNestedSettings(t, map[string]map[string]any{"acme": databaseSettings(addr, "default")}) + a := newApp(t, testConfig(t, root), Options{}) + fake := newFakeClickHouse(t, a, map[string][]string{"default": {"events"}}) + + logs := logtest.Capture(t, slog.LevelWarn) + // No fake names moved_db: its discovery dials the closed port. + rewriteSettings(t, filepath.Join(root, "acme"), databaseSettings(addr, "moved_db")) + _, adopted := a.tenants.Reload("test") + require.True(t, adopted) + require.Eventually(t, func() bool { return strings.Contains(logs.String(), "schema discovery retry failed") }, 5*time.Second, 10*time.Millisecond) + assert.Contains(t, logs.String(), `"tenant":"acme"`) + + req := httptest.NewRequestWithContext(t.Context(), http.MethodPost, "/v1/ingest?table=events", strings.NewReader(`{"id": "1"}`)) + req.Header.Set("Content-Type", "application/json") + req.Header.Set(tenant.Header, "acme") + rec := httptest.NewRecorder() + a.Handler().ServeHTTP(rec, req) + assert.Equal(t, http.StatusServiceUnavailable, rec.Code, "body: %s", rec.Body.String()) + assert.Equal(t, "5", rec.Header().Get("Retry-After")) + assert.Contains(t, rec.Body.String(), "schema not loaded yet") + + fake.set("moved_db", "orders") + moved := a.discoveries.For("acme") + require.Eventually(t, moved.Loaded, 15*time.Second, 10*time.Millisecond, "the loop kept retrying") + assert.NotNil(t, moved.Get("orders")) + assert.Nil(t, moved.Get("events")) +} + // Close stops every tenant's discovery loop within the release budget, and // the pools after them. func TestClose_StopsTheDiscoveryLoops(t *testing.T) { diff --git a/internal/app/discoveries.go b/internal/app/discoveries.go index 60fbee416..d7f3a8e3e 100644 --- a/internal/app/discoveries.go +++ b/internal/app/discoveries.go @@ -21,8 +21,12 @@ import ( // from the settings registry's AfterAdopt hook, after the pools: a newly // served tenant gets a registry over the pool it is on and a loop; a tenant // no longer served — rejected or removed — has its loop stopped and its -// registry dropped, and starts over when it is back. Every loop stops under -// App.Close within the release budget. A lookup is one lock-free load. +// registry dropped, and starts over when it is back. A tenant a reload moved +// to another ClickHouse address or database starts over the same way (drop, +// #638): the schema it held describes tables it no longer reads, so its +// lookups answer not loaded, never the previous database's schema, until the +// first discovery of the new one. Every loop stops under App.Close within the +// release budget. A lookup is one lock-free load. type discoveries struct { // ctx is the loops' parent: the App's stop context. ctx context.Context @@ -82,6 +86,24 @@ func (d *discoveries) reconcile(tenants *settings.Registry) { d.cur.Store(&next) } +// drop stops the loop of each of ids and drops its registry, so the +// reconcile that follows builds each a fresh one, whose loop discovers at +// once and retries until it succeeds. No I/O: the hook calling it holds the +// lock that serializes reloads. A tenant with no registry — new, or back +// after a rejection or removal — is left to that reconcile alone. +func (d *discoveries) drop(ids []tenant.ID) { + d.mu.Lock() + defer d.mu.Unlock() + next := maps.Clone(*d.cur.Load()) + for _, id := range ids { + if td, ok := next[id]; ok { + td.cancel() + delete(next, id) + } + } + d.cur.Store(&next) +} + // adopt registers reg as tenant id's registry and starts its loop: for the // registry boot refreshed synchronously before any loop ran. func (d *discoveries) adopt(id tenant.ID, reg *discovery.SchemaRegistry) { diff --git a/internal/app/roles_test.go b/internal/app/roles_test.go index c50eb9a32..55eceaacf 100644 --- a/internal/app/roles_test.go +++ b/internal/app/roles_test.go @@ -17,6 +17,7 @@ import ( "github.com/Wave-RF/WaveHouse/internal/config" "github.com/Wave-RF/WaveHouse/internal/coord" "github.com/Wave-RF/WaveHouse/internal/settings" + "github.com/Wave-RF/WaveHouse/internal/tenant" "github.com/Wave-RF/WaveHouse/internal/testutil/storedir" ) @@ -155,6 +156,22 @@ func TestNew_OpsOnlyReadiness(t *testing.T) { assert.Nil(t, a.Registry(), "no schema registry without the api role") } +// A reload that moves the tenant repoints the pool of a process without the +// api role, which has no schema registry to start over. +func TestReload_OpsOnlyMovedTenant(t *testing.T) { + addr := closedAddr(t) + dir := writeSettings(t, databaseSettings(addr, "default")) + cfg := testConfig(t, dir) + cfg.Roles = []config.Role{config.RoleIngest} + a := newApp(t, cfg, Options{}) + + rewriteSettings(t, dir, databaseSettings(addr, "moved_db")) + _, adopted := a.tenants.Reload("test") + require.True(t, adopted) + assert.Equal(t, "moved_db", a.pools.For(tenant.Default).Identity().Database) + assert.Nil(t, a.Registry()) +} + func TestNew_OpsOnlyPrometheusInline(t *testing.T) { cfg := testConfig(t, writeSettings(t, nil)) cfg.Roles = []config.Role{config.RoleIngest} diff --git a/internal/app/wire.go b/internal/app/wire.go index a18ab5b45..1893a1adb 100644 --- a/internal/app/wire.go +++ b/internal/app/wire.go @@ -326,6 +326,13 @@ func (a *App) wireClickHouse() error { slog.Warn("cache invalidation of a stale tenant did not land; it may serve stale rows until it does", "tenant", id, "error", err) } } + // The schema a moved tenant discovered is the previous database's + // (#638): its registry is dropped here, and the discovery hook, which + // runs after this one, builds it a fresh one over the pool it is on + // now. Only an API process discovers schemas. + if a.discoveries != nil { + a.discoveries.drop(stale) + } }) return nil } @@ -342,11 +349,12 @@ func (a *App) chConn(id tenant.ID) driver.Conn { } // discoverySource is what tenant id's schema registry discovers from, read -// per refresh so a reload that repoints the tenant or moves its database -// applies to the next one: its pool's connection and the database that pool -// was opened for — never the adopted document's, which a refused move would -// pair with the pool the tenant kept, discovering a database its queries and -// inserts do not use. +// per refresh so a reload that repoints the tenant applies to the next one +// (a move to another address or database starts the tenant over on a fresh +// registry, see wireClickHouse): its pool's connection and the database that +// pool was opened for — never the adopted document's, which a refused move +// would pair with the pool the tenant kept, discovering a database its +// queries and inserts do not use. func (a *App) discoverySource(id tenant.ID) discovery.Source { return func() (driver.Conn, string) { m := a.pools.For(id) diff --git a/internal/discovery/discovery.go b/internal/discovery/discovery.go index 90fe9b0ef..0a7282bfe 100644 --- a/internal/discovery/discovery.go +++ b/internal/discovery/discovery.go @@ -200,7 +200,8 @@ type SchemaRegistry struct { // source supplies the tenant's connection and the database to discover // from, read together once per Refresh, so a settings reload that moves // the tenant to another pool or database is honored by the next refresh - // (the tenant's chconn.Pools entry in production). + // (the tenant's chconn.Pools entry in production, where a tenant moved + // to another address or database gets a fresh registry instead). source Source // tenant is whose tables the registry discovers. tenant tenant.ID diff --git a/internal/testutil/testutil.go b/internal/testutil/testutil.go index 8b3cefb4b..0a8f00219 100644 --- a/internal/testutil/testutil.go +++ b/internal/testutil/testutil.go @@ -36,6 +36,12 @@ func NewTestSchemaRegistry(t testing.TB, tables []*discovery.TableSchema) *disco return reg } +// NewSchemaConn is the mock connection NewTestSchemaRegistry discovers from, +// for a test that builds the registry itself. +func NewSchemaConn(tables []*discovery.TableSchema) driver.Conn { + return &schemaConn{tables: tables} +} + // TestServerVersion is the ClickHouse version NewTestSchemaRegistry's mock // connection reports, so a test can assert against ServerVersion() without // hardcoding the same literal twice. From 8b4eb2a86c1d7ba6bf293c0d918487cc1e8c30cc Mon Sep 17 00:00:00 2001 From: taitelee Date: Tue, 29 Sep 2026 12:59:43 -0400 Subject: [PATCH 2/4] test(dedupe): drop the wall-clock bound from the throttled call test --- internal/dedupe/dynamodb_test.go | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/internal/dedupe/dynamodb_test.go b/internal/dedupe/dynamodb_test.go index a61132590..4da4e14ff 100644 --- a/internal/dedupe/dynamodb_test.go +++ b/internal/dedupe/dynamodb_test.go @@ -554,7 +554,9 @@ func (h *throttledHTTP) Do(r *http.Request) (*http.Response, error) { // Through the real SDK stack at the default Timeout and MaxAttempts: a // throttled put is retried until its attempts run out, inside the call's -// deadline, so the Reserve fails with the throttle as its cause. +// deadline, so the Reserve fails with the throttle as its cause. The deadline +// is read from the error, never from the wall clock, which a loaded machine +// stretches past it around the call. func TestDynamo_ThrottledCallEndsOnItsLastAttempt(t *testing.T) { t.Parallel() h := &throttledHTTP{} @@ -566,16 +568,13 @@ func TestDynamo_ThrottledCallEndsOnItsLastAttempt(t *testing.T) { require.NoError(t, m.Apply(true)) calls := breakerTrips - 1 for range calls { - start := time.Now() _, err := m.Reserve(t.Context(), keys("a"), time.Minute) - took := time.Since(start) require.ErrorIs(t, err, ErrUnavailable) var maxed *retry.MaxAttemptsError require.ErrorAs(t, err, &maxed, "the attempts ran out, not the deadline") var throttled *types.ProvisionedThroughputExceededException require.ErrorAs(t, err, &throttled) require.NotErrorIs(t, err, context.DeadlineExceeded) - assert.Less(t, took, d.cfg.Timeout) } assert.Equal(t, int64(calls*d.cfg.MaxAttempts), h.puts.Load()) } From 8d47b2b905c1e1c173f78f06b2100ca75de5c8bc Mon Sep 17 00:00:00 2001 From: taitelee Date: Tue, 29 Sep 2026 15:43:06 -0400 Subject: [PATCH 3/4] fix(app): drop a moved tenant's registry before cache invalidation and wait for stopped loops on close --- internal/app/app_test.go | 86 +++++++++++++++++++++++++++++++++++++ internal/app/discoveries.go | 45 +++++++++++++++---- internal/app/wire.go | 17 +++++--- 3 files changed, 133 insertions(+), 15 deletions(-) diff --git a/internal/app/app_test.go b/internal/app/app_test.go index 9d761fc86..eff3e187f 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -1967,6 +1967,92 @@ func TestReload_MovedTenantFailedDiscoveryIsRetried(t *testing.T) { assert.Nil(t, moved.Get("events")) } +// registryAtInvalidation is a cache that records, at each InvalidateTenant, +// the tenants that still had a schema registry. +type registryAtInvalidation struct { + *testutil.MockCache + registry func(tenant.ID) *discovery.SchemaRegistry + + mu sync.Mutex + held []tenant.ID +} + +func (c *registryAtInvalidation) InvalidateTenant(ctx context.Context, id tenant.ID) error { + if c.registry(id) != nil { + c.mu.Lock() + c.held = append(c.held, id) + c.mu.Unlock() + } + return c.MockCache.InvalidateTenant(ctx, id) +} + +// A moved tenant's registry is gone before its cache is invalidated: the +// invalidation may wait on its backend with the tenant already on its new +// pool, and a request arriving meanwhile must not find the previous +// database's schema. +func TestReload_MovedTenantRegistryIsDroppedBeforeTheCacheInvalidation(t *testing.T) { + addr := closedAddr(t) + root := writeNestedSettings(t, map[string]map[string]any{"acme": databaseSettings(addr, "default")}) + a := newApp(t, testConfig(t, root), Options{}) + newFakeClickHouse(t, a, map[string][]string{"default": {"events"}, "moved_db": {"orders"}}) + recorder := ®istryAtInvalidation{MockCache: &testutil.MockCache{}, registry: a.discoveries.For} + a.cache = recorder + + rewriteSettings(t, filepath.Join(root, "acme"), databaseSettings(addr, "moved_db")) + _, adopted := a.tenants.Reload("test") + require.True(t, adopted) + require.Equal(t, []tenant.ID{"acme"}, recorder.GetTenants()) + recorder.mu.Lock() + defer recorder.mu.Unlock() + assert.Empty(t, recorder.held, "no registry while the cache is invalidated") +} + +// stuckConn is a connection whose first query, once entered, ignores its +// context and answers only when released. +type stuckConn struct { + driver.Conn + entered, release chan struct{} + once sync.Once +} + +func (c *stuckConn) QueryRow(context.Context, string, ...any) driver.Row { + c.once.Do(func() { close(c.entered) }) + <-c.release + return testutil.UTCRow{} +} + +func (c *stuckConn) Query(context.Context, string, ...any) (driver.Rows, error) { + return nil, errors.New("stuck connection") +} + +// A loop a reload stopped in the middle of a refresh is still running when +// the reload returns: Close waits for it with the loops of the served +// tenants, and names its tenant when the release budget ends first. +func TestClose_WaitsForALoopAReloadStopped(t *testing.T) { + conn := &stuckConn{entered: make(chan struct{}), release: make(chan struct{})} + d := newDiscoveries(t.Context(), nil, func(tenant.ID, error) {}, func(tenant.ID) {}) + d.adopt("acme", discovery.NewSchemaRegistry(func() (driver.Conn, string) { return conn, "default" }, "acme", + func(tenant.ID) time.Duration { return time.Hour })) + loop := (*d.cur.Load())["acme"] + <-conn.entered + + d.drop([]tenant.ID{"acme"}) + require.Nil(t, d.For("acme")) + budget, cancel := context.WithTimeout(t.Context(), 50*time.Millisecond) + defer cancel() + err := d.close(budget) + require.ErrorIs(t, err, context.DeadlineExceeded) + require.ErrorContains(t, err, "tenant acme not stopped") + + close(conn.release) + require.NoError(t, d.close(t.Context()), "a later close waits for it still") + select { + case <-loop.done: + default: + t.Fatal("close returned before the loop ended") + } +} + // Close stops every tenant's discovery loop within the release budget, and // the pools after them. func TestClose_StopsTheDiscoveryLoops(t *testing.T) { diff --git a/internal/app/discoveries.go b/internal/app/discoveries.go index d7f3a8e3e..f056b470b 100644 --- a/internal/app/discoveries.go +++ b/internal/app/discoveries.go @@ -26,7 +26,8 @@ import ( // #638): the schema it held describes tables it no longer reads, so its // lookups answer not loaded, never the previous database's schema, until the // first discovery of the new one. Every loop stops under App.Close within the -// release budget. A lookup is one lock-free load. +// release budget, the ones a reload stopped included. A lookup is one +// lock-free load. type discoveries struct { // ctx is the loops' parent: the App's stop context. ctx context.Context @@ -37,12 +38,17 @@ type discoveries struct { onAttempt func(tenant.ID, error) onLoaded func(tenant.ID) - mu sync.Mutex // serializes reconcile, adopt and close + mu sync.Mutex // serializes reconcile, drop, adopt and close cur atomic.Pointer[map[tenant.ID]*tenantDiscovery] + // retired is the loops a reload stopped that had not ended when last + // looked at: a refresh in flight ends on its own time, and close waits + // for it like the rest. Under mu. + retired []*tenantDiscovery } // tenantDiscovery is one tenant's registry and its loop. type tenantDiscovery struct { + id tenant.ID registry *discovery.SchemaRegistry cancel context.CancelFunc done chan struct{} @@ -79,13 +85,27 @@ func (d *discoveries) reconcile(tenants *settings.Registry) { } for id, td := range cur { if !served[id] { - td.cancel() + d.retire(td) delete(next, id) } } d.cur.Store(&next) } +// retire stops td's loop and keeps it for close to wait for, letting go of +// the retired loops that have ended since. Under mu. +func (d *discoveries) retire(td *tenantDiscovery) { + td.cancel() + d.retired = append(slices.DeleteFunc(d.retired, func(r *tenantDiscovery) bool { + select { + case <-r.done: + return true + default: + return false + } + }), td) +} + // drop stops the loop of each of ids and drops its registry, so the // reconcile that follows builds each a fresh one, whose loop discovers at // once and retries until it succeeds. No I/O: the hook calling it holds the @@ -97,7 +117,7 @@ func (d *discoveries) drop(ids []tenant.ID) { next := maps.Clone(*d.cur.Load()) for _, id := range ids { if td, ok := next[id]; ok { - td.cancel() + d.retire(td) delete(next, id) } } @@ -118,7 +138,7 @@ func (d *discoveries) adopt(id tenant.ID, reg *discovery.SchemaRegistry) { // for a registry already loaded, then the periodic refresh. func (d *discoveries) start(id tenant.ID, reg *discovery.SchemaRegistry) *tenantDiscovery { ctx, cancel := context.WithCancel(d.ctx) //nolint:gosec // G118: held on the tenantDiscovery, called by reconcile or close - td := &tenantDiscovery{registry: reg, cancel: cancel, done: make(chan struct{})} + td := &tenantDiscovery{id: id, registry: reg, cancel: cancel, done: make(chan struct{})} go func() { defer close(td.done) if !reg.Loaded() { @@ -135,7 +155,8 @@ func (d *discoveries) start(id tenant.ID, reg *discovery.SchemaRegistry) *tenant return td } -// close stops every loop and waits for them within ctx, the release budget. +// close stops every loop and waits for them within ctx, the release budget: +// the served tenants', then the ones a reload stopped earlier. func (d *discoveries) close(ctx context.Context) error { d.mu.Lock() defer d.mu.Unlock() @@ -143,12 +164,20 @@ func (d *discoveries) close(ctx context.Context) error { for _, td := range cur { td.cancel() } + loops := make([]*tenantDiscovery, 0, len(cur)+len(d.retired)) for _, id := range slices.Sorted(maps.Keys(cur)) { + loops = append(loops, cur[id]) + } + loops = append(loops, d.retired...) + for i, td := range loops { select { - case <-cur[id].done: + case <-td.done: case <-ctx.Done(): - return fmt.Errorf("schema discovery loop of tenant %s not stopped: %w", id, ctx.Err()) + // A later close waits for what this one did not see end. + d.retired = loops[i:] + return fmt.Errorf("schema discovery loop of tenant %s not stopped: %w", td.id, ctx.Err()) } } + d.retired = nil return nil } diff --git a/internal/app/wire.go b/internal/app/wire.go index 1893a1adb..06c274979 100644 --- a/internal/app/wire.go +++ b/internal/app/wire.go @@ -317,6 +317,16 @@ func (a *App) wireClickHouse() error { if err != nil { slog.Error("clickhouse pools reconciled in part; the next reload retries", "error", err) } + // The schema a moved tenant discovered is the previous database's + // (#638): its registry is dropped here, and the discovery hook, which + // runs after this one, builds it a fresh one over the pool it is on + // now. Dropped first, before the cache invalidation below, which may + // wait on its backend: the tenant is on its new pool already, and a + // request arriving meanwhile must find no schema rather than the + // previous database's. Only an API process discovers schemas. + if a.discoveries != nil { + a.discoveries.drop(stale) + } // A tenant back on a pool after an absence was out of the cache // fan-out (sharedTables) while away, and one moved to another // address or database now reads other tables: either way what it @@ -326,13 +336,6 @@ func (a *App) wireClickHouse() error { slog.Warn("cache invalidation of a stale tenant did not land; it may serve stale rows until it does", "tenant", id, "error", err) } } - // The schema a moved tenant discovered is the previous database's - // (#638): its registry is dropped here, and the discovery hook, which - // runs after this one, builds it a fresh one over the pool it is on - // now. Only an API process discovers schemas. - if a.discoveries != nil { - a.discoveries.drop(stale) - } }) return nil } From ddd489fe8cf5b044ed35cdc8838e668f0e591662 Mon Sep 17 00:00:00 2001 From: taitelee Date: Tue, 29 Sep 2026 16:34:38 -0400 Subject: [PATCH 4/4] test(app): release the stuck connection on every exit of the close test --- internal/app/app_test.go | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/internal/app/app_test.go b/internal/app/app_test.go index eff3e187f..7673a8988 100644 --- a/internal/app/app_test.go +++ b/internal/app/app_test.go @@ -2030,6 +2030,9 @@ func (c *stuckConn) Query(context.Context, string, ...any) (driver.Rows, error) // tenants, and names its tenant when the release budget ends first. func TestClose_WaitsForALoopAReloadStopped(t *testing.T) { conn := &stuckConn{entered: make(chan struct{}), release: make(chan struct{})} + // Released on every way out, so a failed assertion leaves no loop stuck. + release := sync.OnceFunc(func() { close(conn.release) }) + defer release() d := newDiscoveries(t.Context(), nil, func(tenant.ID, error) {}, func(tenant.ID) {}) d.adopt("acme", discovery.NewSchemaRegistry(func() (driver.Conn, string) { return conn, "default" }, "acme", func(tenant.ID) time.Duration { return time.Hour })) @@ -2044,7 +2047,7 @@ func TestClose_WaitsForALoopAReloadStopped(t *testing.T) { require.ErrorIs(t, err, context.DeadlineExceeded) require.ErrorContains(t, err, "tenant acme not stopped") - close(conn.release) + release() require.NoError(t, d.close(t.Context()), "a later close waits for it still") select { case <-loop.done: