Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
e8d6a66
feat(discovery): per-tenant registry lookups over a connection getter
taitelee Sep 23, 2026
8f6b885
feat(clickhouse): one pool per tuple and one schema registry per tenant
taitelee Sep 23, 2026
2f586ea
feat(cache): fan an insert out to the tenants sharing its tables
taitelee Sep 23, 2026
0c8b285
feat(sdk): tenant option on the schema reads, the refresh and sql
taitelee Sep 23, 2026
3b74734
docs: per-tenant clickhouse pools, discovery, probes and the ceiling
taitelee Sep 23, 2026
4de967a
docs: sweep the last shared-clickhouse wording
taitelee Sep 23, 2026
464dda2
Merge origin/main into ch-per-tenant
taitelee Sep 24, 2026
c11413f
fix(app): flip the nested boot state under one lock; review doc fixes
taitelee Sep 24, 2026
9a51e50
fix(discovery): read the database from the pool the connection came from
taitelee Sep 24, 2026
b39e7f7
fix(cache): orphan a tenant moved to another clickhouse address or da…
taitelee Sep 24, 2026
13aadd3
docs: certificate reads per pool, the cache interface, the chconn tre…
taitelee Sep 24, 2026
821d365
docs(api): what a tenant answers when its clickhouse goes down after …
taitelee Sep 24, 2026
946f46e
fix(chconn): open a new pool as the walk places it, so a failed open …
taitelee Sep 24, 2026
ff7d904
fix(chconn): format a tuple by its String, never the struct
taitelee Sep 24, 2026
9b58aea
fix(chconn): name a tuple in errors without a String call on it
taitelee Sep 24, 2026
d588ed6
docs: shared-pool grace, driver refusals, the no-pool 503 on schema r…
taitelee Sep 24, 2026
3474320
docs: a tenant on no pool is any pool that could not be opened, not o…
taitelee Sep 24, 2026
c987860
docs: the remaining ceiling-only wording for a refused or poolless te…
taitelee Sep 24, 2026
0867ec6
docs: every refusal cause in the chconn comments and the agents bullet
taitelee Sep 24, 2026
7412465
fix(chconn): never open a new pool above the ceiling
taitelee Sep 24, 2026
66ee513
docs: a new pool's opening size, and a rewrapped cache comment
taitelee Sep 24, 2026
d85621b
docs(changelog): a flat directory's address or database move orphans …
taitelee Sep 24, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 10 additions & 10 deletions AGENTS.md

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),

### Added

- **One ClickHouse pool per tuple and one schema registry per tenant** (`internal/chconn/chconn.go` (+ tests), `internal/discovery/discovery.go` (+ tests), `internal/app/discoveries.go` (new), `internal/app/{app,wire}.go` (+ tests), `internal/api/{schema,ingest,structured_query,pipes,query,health,errors,clickhouse_exec}.go` (+ tests), `internal/stream/hub.go`, `internal/ingest/worker.go`, `internal/cache/{cache,local,version_manager}.go` (+ tests), `internal/testutil/{testutil,mocks}.go`, `tests/integration/{setup,tenants,boot_resilience,query_limits}_test.go`, `clients/ts/src/{schema,table,sql,client,types}.ts` (+ tests), `tests/e2e/sdk/admin.test.ts`, `docs/src/content/docs/{api,deployment,architecture,ingest-pipeline}.md`, `docs/src/content/docs/{settings-directory,configuration,access-control,reverse-proxy}.mdx`, `docs/src/content/docs/sdk/{admin,reference,queries}.md`, `AGENTS.md`): the second slice of story 6 of the multi-tenant epic ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)), with no behavior change for a settings directory that holds the four files beyond the three noted at the end. The process opens one native pool per distinct `clickhouse.addr` / `database` / `username` / password / `tls` tuple among the served tenants (`chconn.Identity`, `chconn.Pools`), shared by the tenants naming it and sized to their largest `max_open_conns` and `max_idle_conns`; `http_port`, `http_scheme`, `headers` and `query_timeout` stay each tenant's own. Every reload reconciles the pools: a new tuple opens (never dials), a tenant whose tuple changed is repointed, a tuple no tenant names closes after the longest `query_timeout` among the tenants it had, and a pool whose largest ask changed is resized with the same grace. The boot config's `clickhouse.max_total_conns` now bounds the open pools together: boot is refused naming the sum and the ceiling; at a reload a resize above it is refused with the pool kept at its size, and a tuple that cannot be opened — the ceiling, a certificate file that cannot be read, or options the driver refuses — leaves its tenants on the pool they had (the keep-previous-wiring rule of the connection ceiling) or on none when they had none; both are logged and the next reload retries. A tenant on no pool fails closed: `503` with `Retry-After: 30` on `POST /v1/query`, `GET/POST /v1/pipes/{name}`, `POST /v1/ops/query` and `POST /v1/ops/schema/refresh`, ahead of the cache. Each served tenant gets a `discovery.SchemaRegistry` of its own over its pool, kept fresh by its own loop — the boot retry until the first success, then `schema.refresh_interval` with the first refresh at a random point within the interval so tenants adopted together do not refresh together — created and stopped from the settings reload and stopped under `App.Close`; `wavehouse_schema_refresh_failures_total{tenant}` counts a loop's failed attempts. `GET /v1/ops/schema`, `POST /v1/ops/schema/refresh` and `POST /v1/ops/query` take the strict `?tenant=` the pipe reads take (absent is tenant `0`, `400` malformed, `404` unknown, `503` rejected); the SDK sends it as the `tenant` option of `wh.schema.list()`, `wh.schema.refresh()`, `wh.from(t).schema()` and `wh.sql()`. Over a nested directory `/livez` (and `/v1/health`) is `503` with the latest discovery failure, naming its tenant, while no tenant has completed a first discovery, then `200` for the rest of the process lifetime; `/readyz` pings every open pool at once, is ready at the first answer, and names every pool that did not answer when none does — one tenant's ClickHouse outage is its log line and counter, never a probe failure. The ingest worker inserts each batch into its own tenant's ClickHouse — the tenant the message's topic names, through that tenant's HTTP wiring (`chconn.Pools.Target`); a tenant on no pool takes the failure path an unreachable ClickHouse takes — and its cache invalidation fans out to the tenants on the same address and database as the batch's tenant, whatever their user or `tls` block (they read the same tables), rather than to every known tenant, and a tenant adopted after an absence — rejected or removed, so out of that fan-out — or moved to another address or database has its cached structured-query results orphaned in one step (`Cache.InvalidateTenant`, a tenant generation in every version key), so a repaired folder never serves query rows cached before the inserts it missed (a pipe result names no table, so neither this nor any insert invalidates it: it stays until its TTL expires). Three changes reach the single-tenant directory too: a table lookup before the first discovery is a `503` with `Retry-After: 5` (`schema not loaded yet`) rather than a `404`, on `POST /v1/ingest`, `POST /v1/query` and `GET /v1/ops/schema` (the list included, where `[]` would read as no tables); the first periodic refresh fires at a random point within the interval rather than a full interval after boot; and a reload that moves `clickhouse.addr` or `clickhouse.database` now orphans the structured-query results cached before it, where they were served until their TTL. Until [#529](https://github.com/Wave-RF/WaveHouse/issues/529) every tenant's user authenticates with `WH_CH_PASSWORD`, so the tuple is in effect the address, database, username and `tls` block.
- **One token verifier per tenant, built off the boot and reload paths** (`internal/auth/auth.go` (+ tests), `internal/api/router.go` (+ tests), `internal/app/{app,wire}.go` (+ tests), `go.mod`, `docs/src/content/docs/sdk/{reference,streaming}.md`, `docs/src/content/docs/{architecture,deployment,api}.md`, `docs/src/content/docs/{settings-directory,configuration}.mdx`, `SECURITY.md`): story 9 of the multi-tenant epic ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)). Over a nested settings directory each tenant's folder now wires that tenant's verifier (`auth.jwks_url`, `auth.role_claim`), so a JWKS-issued token verifies only under the tenants whose `jwks_url` names its identity provider's key set — under any other tenant's header it is refused as invalid — where before one verifier, tenant `0`'s, accepted a token under any header. `auth.Config` is now the boot-config half alone (the HMAC secret and the operator key, shared by every tenant) and the new `auth.Wiring` a tenant's half; `NewAuthenticator` builds no verifier, `Reconfigure(id, wiring)` gives a tenant one — swapped atomically when its wiring changed, kept when it did not — `Prune` drops the verifiers of the tenants a reload stopped serving, rejected or removed alike (no work runs for a tenant that is not served; a folder adopted again gets a fresh verifier), and `Close` stops every JWKS refresh as a component of `App.Close`. This holds on every reload because the registry's `AfterAdopt` hooks run after every reload it applied, empty list included, so a per-tenant reload that rejects a folder drops its verifier then rather than at the next adoption. The middleware reads the request tenant's verifier through an injected `auth.TenantSource` — the store `api.TenantMW` resolved names its tenant with `settings.Store.Tenant()`, one read of the context — a tenant-exempt route (the ops tree) verifies as tenant `0`, and a tenant with no verifier fails closed. The operator key stamps the request tenant's `admin_role`, read from that tenant's policy, rather than tenant `0`'s. A JWKS key set is fetched on its own goroutine: the verifier is in place at once and *pending* until a set has been stored — the first fetch retried with backoff from one second to a minute, a stored set then kept fresh by the library hourly and, rate-limited, on an unknown key id — so neither boot nor a reload (which holds the registry's lock) waits on the endpoint, and **an unreachable JWKS no longer refuses boot** — it logs (`jwks refresh failed; no token validates until it succeeds`) and that tenant alone is affected. A token checked against a pending verifier is a new outcome, `auth.ErrVerifierPending`: every `/v1` route answers it `503 {"error": "token verifier not ready: …"}` with `Retry-After: 30` (`api.refuseUnverifiable`) rather than evaluate the request under the `default_role`, which could accept its data under a lesser role while another pod holding the keys would have served it as its own; a request without a token, and the operator key, are unaffected. Every fetch goes through one client that caps the response at 1 MiB, refusing a larger one as unreachable with `jwks response exceeds 1048576 bytes` as the logged cause. Nothing changes for a flat directory beyond that boot rule: its one tenant gets exactly the verifier it had. The boot warning for the secretless posture (no `auth.jwt_secret` and no `jwks_url`, whose token refusal landed in [#607](https://github.com/Wave-RF/WaveHouse/pull/607)) is now per tenant — each served tenant with no `jwks_url` while the boot secret is unset — where one line that any JWKS tenant silenced used to stand for the whole directory.

- **ClickHouse TLS, HTTP-interface headers, pool sizes and a connection ceiling** (`internal/settings/{settings,validate,store}.go` + seed, `internal/chconn/chconn.go`, `internal/ingest/worker.go`, `internal/api/query.go`, `internal/config/config.go`, `internal/app/wire.go`, `deployments/compose/settings/config.json`, `deployments/compose/standalone.yaml`, `config.yaml`, `docs/src/content/docs/{settings-directory,configuration,reverse-proxy}.mdx`, `docs/src/content/docs/{architecture,deployment}.md`): the tenant-agnostic first slice of story 6 of the multi-tenant epic ([#583](https://github.com/Wave-RF/WaveHouse/issues/583)). `config.json`'s `clickhouse` block gains a `tls` block (`enabled`, `ca_file`, `cert_file`, `key_file`, `insecure_skip_verify`, `server_name`), a `headers` map for the HTTP interface, and `max_open_conns` / `max_idle_conns`. **Every key is required, so an existing `config.json` must add them**; the seed values (`tls.enabled: false` with the other `tls` keys empty, `headers: {}`, `10` / `5`) change nothing. `tls.enabled` switches the native hop to TLS, `http_scheme` stays the HTTP hop's switch, and the material applies to whichever hop uses TLS: the driver gets the TLS config and the pool sizes, the ingest worker and the raw-SQL proxy get the TLS config and the headers, set ahead of their own so the credentials win (naming `X-ClickHouse-User`, `X-ClickHouse-Key` or `Authorization` is a validation error, and so are two spellings of one name). Validation checks shape only — the paths are not opened, so `wavehouse validate` runs anywhere — warns when `insecure_skip_verify` is on or when only one of the two hops is on TLS (each carries the credentials in the clear without it), and a certificate file that cannot be read or parsed refuses boot or leaves a reload's connection unchanged; the files are read when the connection is built, so a file replaced in place needs a restart. Boot config gains the optional `clickhouse.max_total_conns` (`WH_CH_MAX_TOTAL_CONNS`, `0` = no ceiling): a settings pool above it refuses boot, naming both numbers in the error, and a reload that raises the pool above it is refused and logged (the reload still reports adopted), leaving the connection as it was.
Expand Down
15 changes: 12 additions & 3 deletions clients/ts/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,14 @@ import { StreamController } from "./stream/controller.js";
import { SSETransport } from "./stream/sse.js";
import { SysNamespace } from "./sys.js";
import { TableRef } from "./table.js";
import type { ClientConfig, Database, HttpContext, Result, StreamOptions } from "./types.js";
import type {
ClientConfig,
Database,
HttpContext,
OpsRequestOptions,
Result,
StreamOptions,
} from "./types.js";

type TableName<DB> = DB extends Database ? Extract<keyof DB, string> : string;
type RowType<DB, T extends string> = DB extends Database
Expand Down Expand Up @@ -71,11 +78,13 @@ export class WaveHouseClient<DB extends Database = Database> {
* `service` role). The endpoint proxies straight to ClickHouse's HTTP
* interface so any ClickHouse-accepted SQL works; positional `?` param
* binding is NOT supported — inline literals or use the structured query
* builder for safe binding. See sql.ts for details.
* builder for safe binding. See sql.ts for details. `opts.tenant` names the
* tenant whose ClickHouse the SQL runs against, the default tenant without
* it.
*/
sql<Row = Record<string, unknown>>(
query: string,
opts?: { signal?: AbortSignal },
opts?: OpsRequestOptions,
): Promise<Result<Row[]>> {
// Migration guard: the second argument used to be a positional-`?`
// params array. TS callers get a compile-time error from the type
Expand Down
24 changes: 24 additions & 0 deletions clients/ts/src/namespaces.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,12 @@ describe("sql", () => {

it("POSTs to /v1/ops/query with sql field", async () => {
const result = await sql(makeCtx(), "SELECT count() FROM clicks");
expect(new URL(fetchSpy.mock.calls[0][0]).search).toBe("");

await sql(makeCtx(), "SELECT 1", { tenant: "acme" });
expect(
new URL(fetchSpy.mock.calls[1][0]).pathname + new URL(fetchSpy.mock.calls[1][0]).search,
).toBe("/v1/ops/query?tenant=acme");

expect(result.data).toEqual([{ count: 10 }]);
const [url, init] = fetchSpy.mock.calls[0];
Expand Down Expand Up @@ -79,6 +85,24 @@ describe("SchemaNamespace", () => {
expect(fetchSpy.mock.calls[0][0]).toContain("/v1/ops/schema");
});

it("list() and refresh() send opts.tenant as ?tenant=, and nothing without it", async () => {
fetchSpy.mockImplementation(async () => new Response("[]", { status: 200 }));
const ns = new SchemaNamespace(makeCtx());

await ns.list({ tenant: "acme" });
await ns.refresh({ tenant: "acme" });
await ns.list();
// An empty id is the caller's bug: it is sent for the server to refuse,
// never dropped into a read of the default tenant.
await ns.refresh({ tenant: "" });

const urls = fetchSpy.mock.calls.map((call) => new URL(call[0]));
expect(urls[0].pathname + urls[0].search).toBe("/v1/ops/schema?tenant=acme");
expect(urls[1].pathname + urls[1].search).toBe("/v1/ops/schema/refresh?tenant=acme");
expect(urls[2].search).toBe("");
expect(urls[3].search).toBe("?tenant=");
});

it("refresh() POSTs to /v1/ops/schema/refresh", async () => {
fetchSpy.mockResolvedValue(new Response(JSON.stringify({}), { status: 200 }));

Expand Down
21 changes: 15 additions & 6 deletions clients/ts/src/schema.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { err, ok } from "./errors.js";
import { request } from "./http.js";
import type { HttpContext, Result, Schemas } from "./types.js";
import { request, tenantParam } from "./http.js";
import type { HttpContext, OpsRequestOptions, Result, Schemas } from "./types.js";

/** Namespace for schema introspection. */
export class SchemaNamespace {
Expand All @@ -10,12 +10,17 @@ export class SchemaNamespace {
this._ctx = ctx;
}

/** List all table schemas discovered from ClickHouse. */
async list(opts?: { signal?: AbortSignal }): Promise<Result<Schemas>> {
/**
* List all table schemas discovered from ClickHouse — `opts.tenant`'s, the
* default tenant's without it. A `503` with `Retry-After` is a tenant whose
* first discovery has not succeeded yet.
*/
async list(opts?: OpsRequestOptions): Promise<Result<Schemas>> {
// The backend returns TableSchema[] — transform to Record<string, TableSchema>.
const { data, error } = await request<unknown>(this._ctx, {
method: "GET",
path: "/v1/ops/schema",
params: tenantParam(opts),
signal: opts?.signal,
});
if (error) return err(error);
Expand All @@ -34,11 +39,15 @@ export class SchemaNamespace {
return ok(schemas);
}

/** Force a schema refresh from ClickHouse system.columns. */
async refresh(opts?: { signal?: AbortSignal }): Promise<Result<void>> {
/**
* Force a schema refresh from ClickHouse system.columns — of `opts.tenant`,
* the default tenant without it.
*/
async refresh(opts?: OpsRequestOptions): Promise<Result<void>> {
const { error } = await request<Schemas>(this._ctx, {
method: "POST",
path: "/v1/ops/schema/refresh",
params: tenantParam(opts),
signal: opts?.signal,
});
if (error) return err<void>(error);
Expand Down
7 changes: 4 additions & 3 deletions clients/ts/src/sql.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { err, ok } from "./errors.js";
import { request } from "./http.js";
import type { HttpContext, Result } from "./types.js";
import { request, tenantParam } from "./http.js";
import type { HttpContext, OpsRequestOptions, Result } from "./types.js";

/**
* Execute a raw SQL query against ClickHouse.
Expand Down Expand Up @@ -37,11 +37,12 @@ import type { HttpContext, Result } from "./types.js";
export async function sql<Row = Record<string, unknown>>(
ctx: HttpContext,
query: string,
opts?: { signal?: AbortSignal },
opts?: OpsRequestOptions,
): Promise<Result<Row[]>> {
const { data, error } = await request<Row[]>(ctx, {
method: "POST",
path: "/v1/ops/query",
params: tenantParam(opts),
body: { sql: query },
signal: opts?.signal,
});
Expand Down
7 changes: 7 additions & 0 deletions clients/ts/src/table.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,13 @@ describe("TableRef", () => {

expect(result.data).toEqual(schema);
expect(fetchSpy.mock.calls[0][0]).toContain("/v1/ops/schema?table=clicks");

// The tenant rides beside the table, as on every admin route.
await table().schema({ tenant: "acme" });
const url = new URL(fetchSpy.mock.calls[1][0]);
expect(url.pathname).toBe("/v1/ops/schema");
expect(url.searchParams.get("table")).toBe("clicks");
expect(url.searchParams.get("tenant")).toBe("acme");
});

// --- stream ---
Expand Down
10 changes: 6 additions & 4 deletions clients/ts/src/table.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
import { err, ok } from "./errors.js";
import { request } from "./http.js";
import { request, tenantParam } from "./http.js";
import { QueryBuilder } from "./query-builder.js";
import type { StreamController } from "./stream/controller.js";
import type {
HttpContext,
InsertRecordResult,
InsertResult,
OpsRequestOptions,
RequestOptions,
Result,
StreamOptions,
Expand Down Expand Up @@ -190,11 +191,12 @@ export class TableRef<Row = Record<string, unknown>> {
return ok(result);
}

/** Fetch the schema for this table. */
async schema(opts?: { signal?: AbortSignal }): Promise<Result<TableSchema>> {
/** Fetch the schema for this table — under `opts.tenant`, the default tenant without it. */
async schema(opts?: OpsRequestOptions): Promise<Result<TableSchema>> {
const { data, error } = await request<TableSchema>(this._ctx, {
method: "GET",
path: `/v1/ops/schema?table=${encodeURIComponent(this._table)}`,
path: "/v1/ops/schema",
params: { table: this._table, ...tenantParam(opts) },
signal: opts?.signal,
});
if (error) return err(error);
Expand Down
11 changes: 7 additions & 4 deletions clients/ts/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -470,16 +470,19 @@ export interface PipeRequestOptions {

/**
* Options for a call to one of the admin routes that address a tenant:
* `wh.pipes.list()`, `wh.pipes.get()`, and `wh.settings.reload()`.
* `wh.pipes.list()`, `wh.pipes.get()`, `wh.settings.reload()`,
* `wh.schema.list()`, `wh.schema.refresh()`, `wh.from(t).schema()` and
* `wh.sql()`.
*/
export interface OpsRequestOptions {
signal?: AbortSignal;
/**
* The tenant the call addresses, sent as `?tenant=`. The admin routes ignore
* the `X-Tenant-ID` header, so `options.headers` cannot select one. Omitted,
* the reads serve the default tenant (`0`) and `reload()` reloads every
* tenant. An id the server does not accept — the empty string included — is
* a `400`, never a silent fallback to the default.
* the reads, the refresh and `sql()` serve the default tenant (`0`) and
* `reload()` reloads every tenant. An id the server does not accept — the
* empty string included — is a `400`, never a silent fallback to the
* default.
*/
tenant?: string;
}
Expand Down
Loading
Loading