diff --git a/AGENTS.md b/AGENTS.md index 03bafa62..b29e203e 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; `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 - **`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; over a `Sharded` queue, `claims.go`'s `ClaimShards` narrows the worker to the units this process is assigned — membership leases, capped rendezvous, halt-drain-then-release handover and stop, reset at takeover from a dead owner, each unit's share of a 10,000-row budget of unsettled rows). 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). A broker whose ingest queue is split into units one consumer at a time owns implements `Sharded` too (`IngestUnits`, `ResetOrphaned`, `Unowned`; `ConsumerConfig.Units` narrows a consumer to some of them, and its consumer is a `Releaser` and a `Halter`, capped per unit by `ConsumerConfig.MaxHeld`, with `ErrConsumerMismatch` when an operator durable no longer fits), and `Message.OnSettled` runs a hook once, at the first ack or nak attempt, confirmed or not. 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 implementations, whose subject tokens are escaped by the shared `internal/keyenc`: `EmbeddedNATS` (`embedded.go`, `subject.go`, `purge.go`, `deadletter.go`), which `internal/app` constructs and hands everything else as a `mq.Broker`, and `ExternalNATS` (`external.go`, `subject_nats.go`, `nats_topology.go`: an operator-owned cluster whose streams, durables and lease bucket it never creates, changes, purges or deletes; `lease.go` holds `coord.backend: nats`'s leases in that bucket), which `internal/app` constructs from the `mq.nats` block when `mq.backend` is `nats` ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)). 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 305b77cc..10229066 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -95,6 +95,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. - **A shard's owner keeps it while its ClickHouse is slow, one stuck shard no longer stalls the process's others, and a handover or clean stop keeps each table's rows in order** (`internal/mq/external_consumer.go` (new, from `external.go`), `internal/mq/{external,mq}.go`, `internal/ingest/{claims,worker}.go`, `.testcoverage.yml`, tests in `internal/mq/external_test.go`, `internal/ingest/{claims,worker}_test.go` and `tests/integration/shard_order_test.go` (new), `docs/src/content/docs/{deployment,ingest-pipeline,architecture}.md`, `AGENTS.md`): part of the external-NATS workstream of [#613](https://github.com/Wave-RF/WaveHouse/issues/613), fixing the shard ownership of the entries above before they ship. The pull consumer ran the worker's handler on its delivery goroutine and pulled again only when it returned, and the server renews a shard's pin only on a pull, so a process at its cap of held rows for longer than the pinned TTL (10 seconds) lost its pins: `wavehouse_ingest_shards_unowned` counted its shards, and a peer that judged it dead reset them, redelivering rows it still held — two writers per table. Each shard now has a puller that never runs the handler, fetches only what the shard's cap leaves room for, and at the cap renews the pin every 5 seconds with a one-row fetch, at most one `ack_wait` of such rows past the cap, 12 at the generated 1 minute (past that, and once halted, with a pull of `max_bytes` 1, which delivers nothing); measured, an owner blocked 13 seconds on a hung insert kept one pin throughout and took 2 rows past its share, and a unit drains as fast as before (160,000–192,000 rows a second, against 158,000–172,000). `ResetOrphaned` and `Unowned` also require that the shard delivered nothing and had no ack for its pinned TTL. The process-wide cap of 10,000 held rows let one stuck table take every slot and stall every shard the process owned; each shard now holds its own share, plus at most the 12 rows its renewals take (see the entry above), and with one shard's share full a healthy shard's rows kept flowing. A clean stop now halts every shard first (`mq.Halter`), keeping the pins, so the worker writes what it holds and what the shards still deliver, and releases them after its final flush; rows delivered during that flush were dropped before, and came back only once they had been waited out. A handover keeps the pin until the rows delivered are written, bounded by the worker's 60-second ack wait instead of 15 seconds, and a row the worker drops while stopping is NAKed before the release. Checked under continuous publishing, with the first process's ClickHouse slower than a batch window: every table's rows arrived exactly once and in order across a handover and a clean stop, and the next owner wrote its first rows about 2 seconds after the stop, against 9 seconds when the stop leaves rows behind. A shard fetches in 1-second pulls, so a shard stops fetching within a second of a halt, and an idle shard costs one pull a second. A shard durable deleted while no pull of it was waiting stalled silently until the five-minute topology check, since the next pull only got no responders: two such pulls in a row, or one renewal, now look the durable up, and one that is gone, or whose stream is, ends the worker (measured: 5 seconds after the delete, with the handler busy). A bind whose durable or stream is gone, or whose durable no longer fits (`mq.ErrConsumerMismatch`: `ack_wait` shorter than asked, or no `max_ack_pending`), ends the worker instead of retrying every tick with a warning, and a bind failure that lasts a minute logs an error. `wavehouse_ingest_rows_held_waits_total` is gone: nothing waits at the cap any more. A renewal at the cap takes a row rather than holding it back: a held-back redelivery goes to the back of the server's redelivery queue, and with max_bytes-1 renewals rows NAKed as 0, 1, 2, 3 came back as 2, 3, 0, 1 (measured); now in order, until a stuck shard is 12 rows past its cap. The server requeues the same way while another process's pull waits for the pin, so rows redelivered across a handover or a stop may come back out of order among themselves, though still ahead of newer rows. The unpin on release names no pin id, so it is preceded by a check that the pin is still this process's; the unit renews its pin until the unpin returns (for up to `ack_wait` and one 5-second renewal after its halt has handed on what it fetched, longer than a handover waits), so the pin cannot lapse between the two, and only a consumer leader change there could move it. - **External NATS: the shipped manifests fit the shipped volume, boot refuses permissions narrower than the shards, and the tooling covers `js_domain`** (`internal/mq/{nats_manifests,nats_topology,external}.go`, `internal/mq/nats_permprobe.go` (new), `internal/mq/natstest/natstest.go`, `cmd/wavehouse/mq.go` (+ tests; `internal/mq/{nats_manifests,external_perms}_test.go` new), `deployments/nats/{jetstream.yaml,values.yaml}`, `docs/src/content/docs/{deployment,architecture}.md`, `configuration.mdx`): part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). The shipped streams reserved 225 GiB against the shipped 100 Gi volume, whose size the NATS Helm chart makes each server's `max_file_store`, so the server refused the third partition. `wavehouse mq manifests --file-store` (default `100Gi`) now sizes each partition at 15% of it, the history at 10% and the dead-letter stream at 5%, so the shipped four partitions reserve 75 GiB; a partition's size does not follow N, and the generator refuses streams that reserve more than the store (`--partition-max-bytes` resizes them). The test server takes its file store from the shipped values as the chart does, where it had a limit no test could reach. Boot now probes every shard durable as the connecting user, with a pull and an unpin the server rejects on their merits, and a `required` finding names each durable the user may not pull, where permissions generated for fewer shards passed boot and left those shards' rows unpulled; a consumer request the server refuses later sets `wavehouse_mq_topology_ok` to `0` at once. `wavehouse mq permissions --js-domain` allows and denies the `$JS..API` subjects a client in a JetStream domain sends, beside the plain ones a server in the domain checks after mapping them; `wavehouse mq manifests` takes `--ingest-consumer`, `--history-stream` and `--publish-timeout`, the duplicate window following the timeout; the `wavehouse` user may ack only for its shard durables, under either ack subject layout, where it could ack for any consumer; and an `mq.nats.ingest_consumer` or `history_stream` outside `[a-zA-Z0-9_-]` refuses boot. - **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. diff --git a/docs/src/content/docs/architecture.md b/docs/src/content/docs/architecture.md index 5938a6fb..94ead446 100644 --- a/docs/src/content/docs/architecture.md +++ b/docs/src/content/docs/architecture.md @@ -156,7 +156,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 dc4be85c..49975c2c 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 80a2e8c5..7673a898 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,259 @@ 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")) +} + +// 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{})} + // 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 })) + 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") + + 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 60fbee41..f056b470 100644 --- a/internal/app/discoveries.go +++ b/internal/app/discoveries.go @@ -21,8 +21,13 @@ 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, 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 @@ -33,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{} @@ -75,7 +85,39 @@ 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 +// 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 { + d.retire(td) delete(next, id) } } @@ -96,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() { @@ -113,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() @@ -121,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/roles_test.go b/internal/app/roles_test.go index c50eb9a3..55eceaac 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 4286341a..70ada1bb 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 @@ -342,11 +352,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/dedupe/dynamodb_test.go b/internal/dedupe/dynamodb_test.go index a6113259..4da4e14f 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()) } diff --git a/internal/discovery/discovery.go b/internal/discovery/discovery.go index 90fe9b0e..0a7282bf 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 8b3cefb4..0a8f0021 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.