diff --git a/AGENTS.md b/AGENTS.md index 5a55f6e1b..69d0d5bad 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -28,7 +28,7 @@ One binary: Twenty internal packages under `internal/` (plus `internal/testutil/` for shared test helpers): -- **`api/`** — Chi HTTP router, JWT/JWKS middleware (from `auth/`), ingest/query/structured-query/SSE/schema/DLQ/pipes handlers; `ch_errors.go` (`writeCHError`) is the one mapping from a failed ClickHouse query to status, `code` and `retryable` +- **`api/`** — Chi HTTP router, JWT/JWKS middleware (from `auth/`), ingest/query/structured-query/SSE/schema/DLQ/pipes handlers; `ch_errors.go` (`writeCHError`) is the one mapping from a failed ClickHouse query to status, `code` and `retryable` (`writeCHWriteError` for a write pipe: never retryable, no `Retry-After`, since the write may have run) - **`app/`** — the process wiring: `New` builds every component from the boot config and the settings directory (each one wired in one place — what it opens, what it loops, what it releases — with the settings registry handed to its wiring function whole, the injection point of the per-tenant registry of #583: store-keyed getters for the handlers, `perTenant` for the async paths (with the tenant each message's `mq.Topic` names for the stream hub and the ingest worker), the `chconn.Pools` and the per-tenant `discoveries` reconciled from `AfterAdopt`, `shortestKeepalive` for the one setting folded over every tenant served, `gapWindows` handing the sweeper each tenant's own gap window (a rejected tenant's as its folder last had it, unbounded for one rejected since boot) and the `mq.max_bytes_gb` reconcile each served tenant's byte budget, and `defaultPolicy` for the one setting that still follows tenant `0`, a flat directory's ops-gate admin role; the auth verifiers are per tenant, reconfigured (rebuilt only on changed wiring) and pruned from `AfterAdopt`, and the same hook's `Hub.Prune` ends the open streams of a tenant no longer served), `Run` drives the long-lived ones under one `errgroup` until the context is cancelled or one fails, `Close` releases them in reverse order. `New` wires only what the process's `roles` need (discovery, dedupe, auth verifiers, the hub bridge and keepalive per API process; the ingest worker per ingest process; the sweeper under its lease through `elected`); a process without `api` serves `api.NewOpsRouter` — probes, `/version`, metrics, and the settings reload behind the operator key alone. `cmd/wavehouse` and `tests/integration` both boot through it - **`auth/`** — JWT auth middleware: HMAC **or** JWKS verification with `alg` pinned to the active verifier, role extraction from a configurable claim path; always runs, never rejects (bad token → empty role + stashed reason). One verifier per tenant ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 9): `Authenticator` keys them by `tenant.ID` — the request store's `settings.Store.Tenant()`, through an injected `TenantSource`; `tenant.Default` on the tenant-exempt routes — built from each tenant's `auth` block by `Reconfigure`, dropped by `Prune` once the tenant stops being served (rejected or removed), released by `Close`; the secrets (`Config`) are boot-level and shared. A JWKS key set is fetched off the boot and reload paths: until one has been stored the verifier is pending and a token-bearing request gets `503` + `Retry-After` from `api.refuseUnverifiable` (`auth.ErrVerifierPending`), never a `default_role` evaluation; refresh is library-managed (Eric, 2026-09-22), response capped at 1 MiB; the operator key's admin role is the request tenant's - **`cache/`** — `Cache` interface → `LocalCache` (Ristretto: one pool for every tenant) + `VersionManager` (the invalidation index). Every key leads with the tenant ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 8) — `:query:` for the caller's query key and its singleflight (escaped whole as the lead field of the stored key, `|.||…`), `...
.` for a namespace, each field escaped by `keyenc` where the key is built (a `Namespace` carries the raw table and scope, so no caller escapes) — so no cached read or coalesced flight crosses tenants, a bump through `Invalidate` names one tenant's namespaces and no other's, and `InvalidateTenant` advances the tenant version folded into every namespace key and entry key of one tenant, orphaning its every cached result in one step, pipe results included (no insert reaches a pipe result until [#343](https://github.com/Wave-RF/WaveHouse/pull/343)); `Lookup` returns a `Snapshot` of the versions it read, taken before the handler chooses any input a bump invalidates — the tenant's connection included — and `Set` files the fill under it, so a write landing mid-query, or a reload moving the tenant to another address or database after the request took its connection, orphans the fill ([#382](https://github.com/Wave-RF/WaveHouse/issues/382)), and every backend runs the conformance suite `internal/testutil/cachetest`; the one crossing is the wiring's, above the package: `internal/app` hands the ingest worker the cache through `sharedTables`, which repeats each of the worker's bumps under every tenant on the same ClickHouse address and database (`chconn.Pools.SharingTables`, whatever their user or tls block — they read the same tables), and orphans the whole cache of a tenant back on a pool after an absence, since it was out of that fan-out while away, or moved to another address or database, since it now reads other tables (story 6) @@ -38,7 +38,7 @@ Twenty internal packages under `internal/` (plus `internal/testutil/` for shared - **`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 → `Embedded` (Pebble: every tenant's seen ids in one instance at `data_dir/pebble`, each key led by its tenant, 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), wrapped by `Managed` whose open/closed state follows the hot-reloadable `dedupe.enabled` in the settings directory's `config.json`; `Stores` holds one `Managed` per tenant, built through a `Factory` (`func(tenant.ID) *Managed`, `Embedded.Tenant` in production; `Managed` opens its store through a function, so every backend gets the same switch), and reconciled from the registry's `AfterAdopt` hook — open exactly when the tenant is served with its switch on, closed with its seen ids kept otherwise ([#583](https://github.com/Wave-RF/WaveHouse/issues/583) stories 7 and 3) - **`discovery/`** — `SchemaRegistry`, one per served tenant over a `Source` read once per refresh — the tenant's pool's connection and the database that pool was opened for, one snapshot, so a refused move keeps discovering the database the tenant's queries still use (`internal/app`'s `discoveries` builds, runs and stops them from `AfterAdopt` and `App.Close`: `RetryRefresh` until the first success, then `StartAutoRefresh` with a random first tick; `Lookup` answers `ErrNotLoaded` before the first success — the handlers' `503` with `Retry-After` — and `ErrUnknownTable` after; a failed loop attempt counts in `wavehouse_schema_refresh_failures_total{tenant}`), that introspects ClickHouse `system.columns` (name/type/nullability plus `default_expression` and 1-based `position`) and `system.tables` (each table's `create_table_query`, kept in-process and never serialized — an external-engine table renders its wiring there unconditionally — endpoint, bucket/host, database, username, S3 access key id; ClickHouse masks the password as `[HIDDEN]` from ~23.9, so the exposure is the topology, not the secret), records the server version, + `Validate()` for ingest payloads + `CanonicalizeTimestamps()` rewriting top-level `DateTime`/`DateTime64` column values to the canonical RFC 3339 UTC wire form pre-publish (Key Design Decision #19) -- **`ingest/`** — Ingest worker pipeline (`worker.go`: JetStream input → per-table batch INSERT with DLQ output). The pipeline is **insert-only**. The wire format `EventMessage` (`types.go`) carries `{table_name, scope, received_timestamp, format, columns, row}` and nothing else — `row` is one positional `JSONCompactEachRow` line and `columns` names its slots, the table's **insertable** columns (a `MATERIALIZED`/`ALIAS` column cannot be named in an `INSERT`); the worker batches per (tenant, table, column list), the tenant read off each message's `mq.Topic`, and inserts each batch into its tenant's own ClickHouse (`chconn.Pools.Target`); the worker accepts whatever table name the envelope carries (table existence was already checked by the HTTP ingest handler, which `404`s an unknown table before publish; the worker doesn't re-validate), then bulk-INSERTs. In the embedded-NATS deployment (the default), the server runs with `DontListen: true` (`internal/mq/embedded.go`), so the only Publishers reachable on the `ingest.>` subjects are in-process Go code — today, only the HTTP `/v1/ingest?table={table}` handler. Non-insert mutations (`DELETE`/`UPDATE`/`TRUNCATE`/…) must go through `POST /v1/ops/query` under the admin role (the same `RequireAdmin` gate as the rest of `/v1/ops/*`), so non-admin callers never reach the proxy. A request with no token (or an invalid one) resolves to the `default_role`, which in a production config is not the admin role (setting them equal is a loudly-warned dev-only setting), so it can't reach this endpoint. Plus `Sweeper` (Active Sweeper for NATS message lifecycle) + `EventMessage`/`BufferConsumerName` types (`types.go`) +- **`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 (`cache.VersionManager`) use it; changing what it keeps orphans every stored key - **`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`), `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 - **`observability/`** — OpenTelemetry pipeline: `InitProvider` wires trace/metric/log providers via OTLP gRPC (each signal independently gated). A top-level `Prometheus` config block drives an optional `/metrics` scrape endpoint that runs independently of OTLP push — standalone (Alloy/Mimir scrape, no collector), alongside OTLP, or off. `NewLogger` produces a slog handler that fans out to stdout AND OTLP (stdout always 100%, OTLP sample-rate-aware). `TraceHandler` injects trace_id/span_id from active spans. `tracer.go` provides W3C trace context propagation over message headers (`InjectHeaders`/`ExtractHeaders` on a plain header map; `internal/mq` injects on every publish and extracts on the `Subscribe` path, so this package never sees a NATS type). @@ -65,7 +65,7 @@ The invariant index — what must stay true. Full narrative and rationale live i 10. **Active Sweeper** — purges NATS messages that are both ACKed (written to CH) and older than the gap window; SSE gap-fill uses `DeliverByStartTime`, no in-process ring buffer. It runs only in the process holding the `sweeper` lease (`coord.RunElected`); that lease is not fenced: an overlap cannot lose ClickHouse data, since every sweep stops at the consumer's ack floor; it can only trim SSE replay history, and only when the two holders' settings views differ (one still reading a shorter `stream.gap_window_minutes`, or missing a tenant, after a reload the other has applied) — which fencing would not prevent either. Anything that does need exclusivity must check the term's `Token`. 11. **Hasura-style access control: fail-closed (security)** — `policy.IsAdmin` (role == `admin_role`, **exact case-sensitive**, default `"admin"`) is the single admin check, shared by `Evaluate`/`ResolveRole`/`Validate`/the `/v1/ops` gate/`RoleAllowed`. Empty/absent role matches nothing (no `"*"` wildcard); `Validate` rejects empty role keys; a `nil` policy (deleted) denies **everyone incl. admin** via a role — a total lockout for token-based callers, so recovery is writing `policies.json` and reloading, never an implicit admin grant (**exception:** the operator key's `auth.IsOperator` bit passes the `/v1/ops` gate even under a `nil` policy — a deliberate break-glass that can `POST /v1/ops/settings/reload` over HTTP, see #7). Over a nested settings directory the `/v1/ops` gate reads no policy at all — those routes reach every tenant, so the operator key alone passes and an admin-role token gets `403`; `api.NewRouter` decides that from the registry's shape, not from what was wired. `default_role` is the one sanctioned roleless exception (`ResolveRole` maps empty → it pre-eval); `default_role == admin_role` is permitted but dev-only and loudly warned (`policy.DefaultRoleGrantsAdmin`). Preserve when touching `internal/policy` (policy twin of #13; see #159). Detail: architecture.md § `policy/`. 12. **Structured queries: column authz fail-closed (security)** — `POST /v1/query?table={table}`: typed AST validated against schema, permission-enforced, timestamp-bucketed for cache, `DefaultMaxRows` (10,000) cap. Every column reference — projection, aggregation args, `filters`, `group_by`, `order_by`, `time_range` — is authorized inside `query.Build` (the single chokepoint that enumerates them all), so no clause can skip the role's `allow_columns`/`deny_columns` check (#223). A `select_all` read by a *column-restricted* role expands to its allowed columns via `policy.AllowedProjection`, never a bare `SELECT *`; *unrestricted*/admin roles keep `SELECT *` (`policy.RestrictsColumns` decides). Omitting `columns` selects nothing (`ErrEmptyProjection` → `200 []`); `["*"]` is the literal column `*` (schema-gated, not a wildcard); a table-granted role with no readable columns fails closed (`ErrNoReadableColumns` → `403`). Structured and live-stream (`stream.projectIndices`) reads share the one per-column decision `policy.IsColumnAllowed`, so column visibility can't drift. Row visibility has the same one-source guarantee (#319): `Evaluate` resolves a role's row-`filter` once (`resolvePredicates`), and both surfaces consume that single resolution — the query path renders it to SQL (`predicatesToSQL`), the stream evaluates it in memory per subscriber (`ResolvedPermissions.RowVisible`, whose type-aware comparison fails closed on anything it can't prove about the ingested payload — `policy.ColumnSpec`, with `DateTime`/`DateTime64` operands compared as instants through the ingest grammar (`discovery.Column.TimeParser`) and claim constants rendered canonically and digit-exact by the one shared rule `policy.CanonicalScalar` (#457 — which also refuses a float64 at/past 2^53 rather than match a neighboring ID, and whose ok=false — an absent claim, a structured value, no canonical form — makes the predicate match no rows on BOTH surfaces: `1 = 0` in SQL, every row withheld in memory); numeric comparison runs in the column's STORAGE domain (`policy.NumericSpec`, classified by `discovery.NumericStorageOf` — Float width rounding, Decimal scale truncation, integer exactness, both operands narrowed as ClickHouse narrows stored value and bound constant, out-of-range operands refused rather than modeled; the `tests/integration` differential oracle holds in-range verdicts equal to a live ClickHouse's and the never-admit-where-SQL-hides direction for the refused out-of-range ones); an event whose insert later fails into the DLQ is the one residual payload-vs-stored asymmetry, documented in the access-control enforcement caution) — so row visibility can't drift either. Preserve when touching `internal/query` or the structured-query handler. Detail: architecture.md § `query/`. -13. **Named query pipes: fail-closed (security)** — pre-defined SQL templates (Tinybird-style) with param binding + caching; `GET/POST /v1/pipes/{name}` sit outside `RequireAdmin`, so per-pipe `allowed_roles` is the *only* execute-path gate, via `policy.RoleAllowed`: exact allowlist membership (no `"*"`), admin always passes, empty/absent role and empty-string entries authorize nobody, and no `allowed_roles` → admin-only. Preserve and exercise via `testutil.RunRoleMatrix` / `StandardRoleMatrix` (see #159). Detail: architecture.md § `pipes/`. +13. **Named query pipes: fail-closed (security)** — pre-defined SQL templates (Tinybird-style) with param binding + caching — reads only: a pipe whose SQL `IsMutation` classifies as a write bypasses the cache and singleflight, since a cached or coalesced write is a dropped one (#386); `GET/POST /v1/pipes/{name}` sit outside `RequireAdmin`, so per-pipe `allowed_roles` is the *only* execute-path gate, via `policy.RoleAllowed`: exact allowlist membership (no `"*"`), admin always passes, empty/absent role and empty-string entries authorize nobody, and no `allowed_roles` → admin-only. Preserve and exercise via `testutil.RunRoleMatrix` / `StandardRoleMatrix` (see #159). Detail: architecture.md § `pipes/`. 14. **TypeScript SDK** — `@wavehouse/sdk`: typed query builder, real-time SSE over `fetch`, live queries (incrementable/decomposable/poll aggregation), codegen CLI. Exactly one runtime dependency — `eventsource-parser` (SSE framing, itself dependency-free); adding a second needs the same scrutiny the first got. The canonical client (see §SDK Sync). 15. **Observability invariants** — stdout always 100% (sampling is OTLP-push-only); WARN+ERROR always export at 100% (a non-configurable floor — don't expose it); gRPC OTel exporters dial lazily so an unreachable collector never blocks startup; the OTel Prometheus exporter uses a **private** `prometheus.Registry`. The OTLP endpoint/TLS/custom-CA/mTLS/headers are delegated to the OpenTelemetry SDK's standard `OTEL_EXPORTER_OTLP_*` env vars — `InitProvider` passes **no** endpoint/header options. Known gap, intentionally not patched in WaveHouse app code: the pinned gRPC logs exporter (`otlploggrpc` v0.19/v0.20) ignores the env TLS-cert vars, so a custom/private CA and mutual TLS apply to traces/metrics but **not** the logs signal (public-CA/system-roots TLS and plaintext still work for logs) — upstream bug open-telemetry/opentelemetry-go#6661. A malformed `OTEL_EXPORTER_OTLP_HEADERS` is logged and skipped by the SDK (fail-soft), not fatal. Preserve when touching the logger/sampler/provider. Detail: architecture.md § `observability/`. 16. **Bearer-token-only CORS posture (security)** — Bearer JWT on every request, no cookies/sessions; `corsMiddleware` deliberately **never** emits `Access-Control-Allow-Credentials` (not needed, and `*` + credentials is a spec violation browsers reject). `cors.allowed_origins` (settings directory, per tenant: a tenant route is decorated from the list of the tenant it names, everything else from tenant `0`'s — `corsOrigins`) controls who can *read* responses, not cookie scope; CSRF protection is structural. Don't reintroduce cookie auth or `Allow-Credentials` without a design discussion — answers GitHub #29/#30. Code: `internal/api/router.go`. @@ -130,6 +130,7 @@ Tooling notes (the non-obvious bits `make help` won't tell you): - **JWT helpers**: Use `testutil.MakeJWT(t, claims)` and `testutil.MakeExpiredJWT(t, claims)` for auth tests. See `testutil/jwt.go`. - **Schema helpers**: Use `testutil.NewTestSchemaRegistry(t, tables)` for schema-aware tests — it builds the registry through the real discovery path (`Refresh` against a mock ClickHouse connection), so timestamp specs are precomputed like production. - **Cache backends**: every `cache.Cache` backend runs `cachetest.Run` (`internal/testutil/cachetest`), the backend-agnostic conformance suite; a behavior the contract promises goes there, not in one backend's tests. +- **Write classifier cases**: `api.IsMutation`'s cases live in `internal/testutil/mutationtest`, shared by its unit test and `tests/integration/ismutation_test.go`, which checks each against ClickHouse's own parser; add a case there, not to either test. - **Policy helpers**: Use `policy.Static(p)` for a fixed `policy.Source` in tests. - **Pipes helpers**: Use `pipes.Static(queries...)` for a fixed `pipes.Source` in tests. - **Response assertions**: Use `testutil.AssertJSONResponse(t, rec, status, expected)` and `testutil.AssertJSONContains(t, rec, status, substring)`. @@ -447,7 +448,7 @@ internal/query/ → Structured query AST + SQL builder internal/settings/ → Settings directory (validate, adopted snapshot + reload, watcher, embedded seed) internal/stream/ → SSE fan-out (event Hub: project once per role, Subscriber outbound queue, Bucket fan-out, keepalive Heartbeater wheel) internal/tenant/ → Tenant id (type, grammar, reserved default, request header name) -internal/testutil/ → Shared test helpers (mocks, JWT + schema helpers; logtest/ captures or silences the default logger; cachetest/ is the conformance suite every cache.Cache backend runs) +internal/testutil/ → Shared test helpers (mocks, JWT + schema helpers; logtest/ captures or silences the default logger; cachetest/ is the conformance suite every cache.Cache backend runs; mutationtest/ holds the shared write-classifier cases) tests/ → Integration & E2E tests tests/integration/ → Go integration tests (//go:build integration; ClickHouse testcontainer) tests/e2e/ → E2E test stack (scripts/orchestrator boots a ClickHouse testcontainer + the wavehouse-cov binary) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8492d1576..b2a744936 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -85,6 +85,8 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Fixed +- **A pipe that writes runs on every call instead of being answered from the cache** (`internal/api/{pipes,ch_errors}.go` (+ tests), `docs/src/content/docs/{pipes.mdx,api.md,architecture.md,configuration.mdx,settings-directory.mdx,ingest-pipeline.md,sdk/pipes.md,sdk/reference.md}`, `clients/ts/src/pipes.ts` (doc comment), `internal/{settings/settings,app/wire}.go` (comments), `AGENTS.md`): fixes [#386](https://github.com/Wave-RF/WaveHouse/issues/386). `/v1/pipes/{name}` sent a write's SQL to ClickHouse through `Exec`, but still cached the `[]` it returned and coalesced identical calls in flight, so a repeat within the TTL answered `200` without executing and concurrent identical calls became one write — silently dropped writes, and with a shared cache ([#613](https://github.com/Wave-RF/WaveHouse/issues/613)) on every instance. A pipe whose bound SQL `IsMutation` classifies as a write — the same classifier that picks `Exec` — now skips the cache lookup, the fill and singleflight, and answers `X-Cache: BYPASS` with `Cache-Control: no-store`, so an HTTP cache in front of a `GET` cannot drop the write either. Classification stays automatic rather than a declared pipe property, so an operator cannot forget to mark one, and costs no ClickHouse round trip. A failed write answers with the status and `code` a failed read gets (see the ClickHouse-errors entry below), but always `retryable: false` and with no `Retry-After`, `503 clickhouse.unavailable` included: the statement may have run, so the SDK does not retry it. A write refused before it is sent, the tenant on no pool, keeps its `503` with `Retry-After: 30`. Read pipes are unchanged. Not in this fix: a write pipe still does not invalidate cached reads of the table it writes ([#394](https://github.com/Wave-RF/WaveHouse/issues/394)). +- **The write classifier skips whitespace, comments and quoted text the way ClickHouse's lexer does, classifies a `WITH`-led statement by `INSERT INTO` alone, and looks through `EXECUTE AS`** (`internal/api/clickhouse_exec.go` (+ tests), `internal/testutil/mutationtest` (new), `tests/integration/ismutation_test.go` (new), `docs/src/content/docs/pipes.mdx`, `AGENTS.md`): `IsMutation` picks `Exec` for a write, and since [#386](https://github.com/Wave-RF/WaveHouse/issues/386) keeps a write pipe out of the cache. It missed a write behind a backslash-escaped quote (`'it\'s'`, and the same inside `"…"` and `` `…` ``), a heredoc (`$$ ( $$`, `$tag$ … $tag$`), a curly-quoted literal or identifier (`‘(’`, `“c(d”`), a `//` line comment, a nested block comment (`/* a /* b */ SELECT */ INSERT …`), a number led by `.` with the verb glued to it (`WITH 1 AS a, .5INSERT INTO t …`, which ClickHouse reads as `.5` then `INSERT`), an `EXECUTE AS ` prefix (`EXECUTE AS u INSERT …`), or leading whitespace other than space, tab, CR and LF: `\v`, `\f`, a no-break space, a byte-order mark, and the other Unicode spaces ClickHouse skips. A missed write went through `Query`, which ran it and then failed the call with a `5xx` the TypeScript SDK retries, so one call could write three times. The same gaps, and a word led by `_` (`_delete`) whose tail was read as a verb, could make a read look like a write, which runs through `Exec` and answers `[]`. After a `WITH` list, which ClickHouse follows only with `SELECT`, a FROM-first `SELECT` or `INSERT INTO`, a name spelled like a keyword was taken for the statement: `WITH 'd' AS desc INSERT …` and `WITH 1 AS select INSERT …` ran as reads, and `WITH 1 AS set SELECT set` and `WITH 1 AS x FROM system.one SELECT x` as writes. A `WITH`-led statement is now a write exactly when it holds `INSERT INTO` outside parentheses. The classifier, exported as `IsMutation` for it, is now checked against the pinned ClickHouse's own parser (`EXPLAIN AST`) in the integration suite: every test case, and every ClickHouse keyword as a `WITH` list's name ahead of each statement a `WITH` list can lead. - **A failed ClickHouse query answers by what went wrong, not a flat `500`/`502`** (`internal/api/ch_errors.go` (new, + tests), `internal/api/{errors,query,structured_query,pipes,schema,ch_settings}.go`, `internal/chconn/errclass.go` (`HTTPStatus` exported), `clients/ts/src/errors.ts` (+ tests), `tests/integration/query_errors_test.go` (new), `tests/integration/query_limits_test.go`, `internal/app/app_test.go`, `tests/e2e/sdk/{admin,query}.test.ts`, `AGENTS.md`, `docs/src/content/docs/{api,architecture}.md`, `docs/src/content/docs/{access-control,configuration}.mdx`, `docs/src/content/docs/sdk/{reference.md,index.mdx}`): fixes [#403](https://github.com/Wave-RF/WaveHouse/issues/403) and [#271](https://github.com/Wave-RF/WaveHouse/issues/271), part of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). ClickHouse answers a syntax error, a missing grant and an overloaded server alike with HTTP `500`, so `/v1/ops/query` turned a bad statement into a `502` and `/v1/query` and pipes into a `500` the SDK retried. All three now class the failure with `chconn.Classify` through one helper, `writeCHError`: a statement ClickHouse refused is `400 clickhouse.rejected`; a query over a rows/bytes limit, the role's own memory cap, or its time cap where that is no longer than `query_timeout` is `400 clickhouse.limit_exceeded`; `ACCESS_DENIED` is `403 clickhouse.access_denied`; credentials, user or database refused, or a redirect or `4xx` with no exception code from whatever fronts ClickHouse, is `502 clickhouse.misconfigured`; ClickHouse down, unreachable or overloaded is `503 clickhouse.unavailable` with `Retry-After: 5`; a failure with no verdict stays `500` (`502` on the proxy) as `clickhouse.unknown`. The error envelope gains `code` and `retryable` next to `error` on these responses — additive. A role with `max_execution_time` now queries with no context deadline and a cancel two seconds past the cap instead: clickhouse-go overwrote the cap's `max_execution_time` with deadline+5s for any deadline over 1s, so an overrun came back as a bare deadline, indistinguishable from waiting for a pooled connection; ClickHouse now enforces the cap itself and reports `TIMEOUT_EXCEEDED`. `POST /v1/ops/schema/refresh` against an unreachable ClickHouse is a `503` with `Retry-After` instead of a `500`. **SDK:** `WaveHouseError.code` and `retryable` now take the server's `code`/`retryable` when the body has them (`HTTP_` and "5xx retries" otherwise), so a rejected query is `clickhouse.rejected` rather than `HTTP_500`, and is not retried. - **An unavailable ClickHouse is retried with backoff instead of dead-lettering every row** (`internal/chconn/errclass.go` (new, + tests), `internal/ingest/{worker,backoff}.go` (`backoff.go` new, + tests), `internal/mq/{mq,embedded}.go`, `internal/testutil/mocks.go`, `tests/integration/ingest_outage_test.go` (new), `AGENTS.md`, `README.md`, `docs/src/content/docs/{ingest-pipeline,architecture,api,deployment,why-wavehouse}.md`, `docs/src/content/docs/{settings-directory,index,access-control}.mdx`): workstream A of [#613](https://github.com/Wave-RF/WaveHouse/issues/613). A failed batch insert used to go through row-by-row isolation whatever the failure, so a ClickHouse that was down, overloaded or read-only failed every row twice and parked the whole batch on the DLQ. `chconn.Classify` now classes the failure first — `Rejected` (any ClickHouse exception code outside the availability and credential lists: the server read the row and refused it), `Unavailable` (connection refused/reset, timeouts, `TOO_MANY_SIMULTANEOUS_QUERIES`, `SERVER_OVERLOADED`, `MEMORY_LIMIT_EXCEEDED`, `TOO_MANY_PARTS`, `READONLY`, `TABLE_IS_READ_ONLY`, `KEEPER_EXCEPTION`, …), `Denied` (`AUTHENTICATION_FAILED`, `ACCESS_DENIED`, …) or `Unknown` (no code, no recognizable transport failure). Only `Rejected` is isolated and dead-lettered as before, and a multi-row batch refused with `TOO_MANY_PARTS` or `MEMORY_LIMIT_EXCEEDED` is split row by row first (`chconn.Splittable`), because a batch spanning too many partitions or too much memory can fail where each of its rows inserts; every other class hands the batch back to the queue with a delayed nak (`mq.Message.NakWithDelay`, new) under a jittered 1 s → 30 s backoff shared by every table on the same ClickHouse pool (a failure of one table — read-only, too many parts or mutations, a grant missing on it, `chconn.TableScoped` — backs off that table alone), which turns rows away without a request while it runs and probes once per window, and ClickHouse going away mid-isolation stops isolation and retries the rows it had not settled. Counted by the new `wavehouse_ingest_retries_total{table, reason}`; logged at `WARN` when an outage starts and at most every 30 s during it. A long outage now shows as a growing ingest stream and, at `mq.max_bytes_gb`, ingest `503`s — not as a full DLQ; a lasting failure of one table holds back its tenant's other tables once its waiting rows reach `maxAckPending`. Retried rows come back out of arrival order, which matters only to a `ReplacingMergeTree` without a version column or a `CollapsingMergeTree`. - **Schema discovery's retry loop jitters its backoff** (`internal/discovery/discovery.go` (+ tests), `internal/app/wire.go`, `internal/api/errors.go`, `AGENTS.md`, `docs/src/content/docs/{architecture,api,deployment}.md`): `RetryRefresh` slept exactly `2s * 2^n` capped at 60s, so instances retrying against one recovering ClickHouse fired in lockstep, every 60s on the same second. Each sleep is now drawn uniformly from below the backoff (full jitter), spreading the retries over the whole window and halving the mean wait — so a failing tenant's retries, their log lines and `wavehouse_schema_refresh_failures_total` come about twice as often ([#141](https://github.com/Wave-RF/WaveHouse/issues/141)). @@ -106,6 +108,8 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Security +- **The pipes page no longer says a parameter can never break out of its literal** (`docs/src/content/docs/pipes.mdx`, `internal/pipes/pipes.go`): that holds only for a placeholder written bare. A string value brings its own quotes, so inside a quoted placeholder they close the template's: the body `{"id": " OR 1=1 OR id = "}` turns `WHERE id = '{{id}}'` into `WHERE id = '' OR 1=1 OR id = ''`, which matches every row. The page now says to write each placeholder bare, never inside quotes. Check existing `pipes.json` templates for quoted placeholders (`'{{x}}'`) and write them bare ([#662](https://github.com/Wave-RF/WaveHouse/issues/662)). + - **An empty HMAC secret no longer verifies tokens signed with an empty key** (`internal/auth/auth.go` (+ tests), `SECURITY.md`): with `auth.jwt_secret` unset and no `auth.jwks_url` — the documented public-access posture, "no token can validate" — the key function handed `golang-jwt` an empty HMAC key, and the library verifies a token signed with one, so anyone could mint `{"role": "admin"}` and reach the whole data plane and `/v1/ops/*`. The verifier now refuses every token when it has neither a secret nor a JWKS URL, pinned by a test that signs with the empty key. Found by review on [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 9 ([#604](https://github.com/Wave-RF/WaveHouse/pull/604)), which carries the same fix. - **Policy validation now rejects the fail-open rule shapes strict decoding can't see** (`internal/policy/policy.go`, `docs/src/content/docs/access-control.mdx`; closes [#460](https://github.com/Wave-RF/WaveHouse/issues/460)): four new `validateRolePerms` rejections close the fail-open shapes strict decoding can't see because the document is syntactically innocent. A `filter` entry with no operator (`"tenant_id": {}`) resolved to zero predicates — no `WHERE` clause, row security silently off, the same shape a misspelled `"eq"` for `"_eq"` used to decode to before strict decoding closed that route; it is now rejected, as is its check-path twin (an operator-less `check` entry, skipped by `Evaluate`'s resolve switch — accepted but constraining nothing) and `filter:` under an `insert:` grant (resolved and then ignored by the ingest path — the same accept-but-ignore family as [#224](https://github.com/Wave-RF/WaveHouse/issues/224), and the pointed asymmetry #460 called out against the loud `check` `_neq`/`_gt`/`_lt` rejection) along with its mirror, `check:` under a `select:` grant — the likelier authoring slip and the fail-open direction: the author believes reads are row-scoped while `Evaluate` resolves the entry and nothing on the select or stream paths reads it. Because [#508](https://github.com/Wave-RF/WaveHouse/pull/508) funneled every adoption through the one `policy.Validate` path, the four checks land on boot, the directory watch, `SIGHUP`, `POST /v1/ops/settings/reload`, and `wavehouse validate` at once. #460's migration caveat (a stored policy hard-failing at boot) has evaporated with the settings directory being new and unreleased; no shipped seed, compose, or fixture policy carries any of the rejected shapes. diff --git a/clients/ts/src/pipes.ts b/clients/ts/src/pipes.ts index 0297e8511..7fba4cc3a 100644 --- a/clients/ts/src/pipes.ts +++ b/clients/ts/src/pipes.ts @@ -67,8 +67,8 @@ export class PipeRef> implements PromiseLike` names the [tenant](/deployment#the-nested-settings-directory) whose ClickHouse the SQL runs against — its own database, credentials and HTTP wiring; without it the SQL runs against tenant `0`'s, which is the whole settings directory unless it is nested. The parameter is parsed as strictly as on the [schema routes](#get-v1opsschema--list-all-table-schemas): `400` for a query string that does not parse or an empty, repeated or malformed id, `404` for an unknown tenant, `503` for one whose settings folder was rejected — all decided before the body is read. A tenant on no ClickHouse pool ([no pool could be opened for it](/settings-directory#clickhouse), such as one the connection ceiling refused) answers `503` `{"error":"no ClickHouse connection is open for this tenant"}` with `Retry-After: 30`. @@ -585,7 +587,7 @@ The inbound request body is capped at 1 MiB; a body over the cap is rejected wit ### `GET/POST /v1/pipes/{name}` — Execute Named Pipe -Executes a pre-defined named query (pipe) with parameter binding. Parameters can be supplied via query string and/or JSON body. Results are cached in the shared L1 (Ristretto) with singleflight coalescing — same machinery as the structured query endpoint, keyed by [tenant](/deployment#multi-tenant-deployments) like it, and again, unlike `/v1/ops/query`. +Executes a pre-defined named query (pipe) with parameter binding. Parameters can be supplied via query string and/or JSON body. A read's results are cached in the query cache with singleflight coalescing — same machinery as the structured query endpoint, keyed by [tenant](/deployment#multi-tenant-deployments) like it, and again, unlike `/v1/ops/query`; a [pipe that writes](/pipes#pipes-that-write) is neither cached nor coalesced (see Response). **Query Parameters:** Any key matching a pipe parameter name. @@ -600,7 +602,7 @@ Executes a pre-defined named query (pipe) with parameter binding. Parameters can **Response:** -JSON array of result rows, with `X-Cache: HIT` or `X-Cache: MISS` indicating whether the row came from the in-process L1. +JSON array of result rows, with `X-Cache: HIT` or `X-Cache: MISS` indicating whether the row came from the in-process L1. A pipe whose SQL is a write (`INSERT`, `ALTER`, `WITH … INSERT`, …) bypasses the cache and singleflight: it executes on every call, identical calls in flight are not coalesced, and the response is `[]` with `X-Cache: BYPASS` and `Cache-Control: no-store` (so an HTTP cache in front of a `GET` cannot answer a repeat) — see [Pipes that write](/pipes#pipes-that-write). The POST parameter body is capped at 1 MiB; a body over the cap is rejected with `413` (the same 1 MiB parameter/AST-body cap as [`POST /v1/query`](#post-v1querytabletable--structured-query) — see [reverse proxy → body limits](/reverse-proxy#request-body-size-limits)). A malformed-but-within-cap body is ignored rather than rejected, since parameters may legitimately come from the query string alone. @@ -615,7 +617,7 @@ The POST parameter body is capped at 1 MiB; a body over the cap is rejected with | 400 | `{"error":"parameter \"x\": unsupported parameter type object"}` | A non-scalar value with no SQL literal form — a JSON object, whether supplied directly or nested as an array element. A JSON **array** is valid and renders as an `IN`-style `(…)` list. | | 400 | `{"error":"parameter \"x\": array parameter must not be empty"}` | An empty array — it would render as the invalid `IN ()`. | | 413 | `{"error":"request body exceeded 1048576 bytes"}` | POST body over the 1 MiB cap | -| 400 / 403 / 500 / 502 / 503 | `{"error":"clickhouse query: …","code":"clickhouse.…","retryable":…}` | ClickHouse failed the pipe's query — for instance a parameter value it cannot use (`400 clickhouse.rejected`), or ClickHouse down (`503 clickhouse.unavailable`, `Retry-After: 5`); see [ClickHouse errors on the query paths](#clickhouse-errors-on-the-query-paths) | +| 400 / 403 / 500 / 502 / 503 | `{"error":"clickhouse query: …","code":"clickhouse.…","retryable":…}` | ClickHouse failed the pipe's query — for instance a parameter value it cannot use (`400 clickhouse.rejected`), or ClickHouse down (`503 clickhouse.unavailable`, `Retry-After: 5`); see [ClickHouse errors on the query paths](#clickhouse-errors-on-the-query-paths). A [pipe that writes](/pipes#pipes-that-write) answers `retryable: false` with no `Retry-After`, its message led by `clickhouse exec:` | | 503 | `{"error":"token verifier not ready: the tenant's JWKS has not been fetched yet"}` | A token was supplied, with no valid operator key, while the tenant's JWKS has not been fetched yet; refused before any policy runs, with a `Retry-After: 30` header — see [Authentication](#authentication) | --- diff --git a/docs/src/content/docs/architecture.md b/docs/src/content/docs/architecture.md index 5672ddd16..976e645fa 100644 --- a/docs/src/content/docs/architecture.md +++ b/docs/src/content/docs/architecture.md @@ -80,9 +80,9 @@ The API layer uses [Chi](https://github.com/go-chi/chi) for routing with Request - **router.go** — Route definitions. Public: `/livez`, `/readyz`, and the content-free `/v1/health` SDK ping (plus the permanent `/healthz` alias and the deprecated `/health`, `/ready` aliases). Policy-gated: `/v1/ingest?table={table}`, `/v1/query?table={table}` (structured), `/v1/pipes/{name}` (named pipes), `/v1/stream`. Admin-only (`RequireAdmin` — role == `policy.admin_role`, or a request bearing the operator key's operator bit, which passes even under a nil policy; over a nested settings directory `NewRouter` mounts the gate with no policy at all, whatever `Dependencies.PolicySource` was wired, so the operator key alone passes): `/v1/ops/schema/*`, `/v1/ops/dlq/stats`, `GET /v1/ops/pipes[/{name}]`, `/v1/ops/settings/reload`, `/v1/ops/query` (raw SQL — same gate as the rest of `/v1/ops/*`). `NewOpsRouter` is the router of a process without the `api` role: the probes and their aliases, `/version` and the same-port metrics path (the part it shares with `NewRouter`, `newProbeRouter`), and `POST /v1/ops/settings/reload` behind `RequireAdmin(nil)`, so only the operator key passes; every other route is a 404, under `/v1/ops` only once that gate has passed. - **auth middleware** — the JWT/JWKS authentication middleware is its own package, [`auth/`](#auth--authentication); the router runs it on every `/v1/*` route. - **tenant.go** — `TenantMW` resolves the request's tenant ahead of the auth middleware on every `/v1` route outside `/v1/ops/*`: the [`X-Tenant-ID`](/deployment#multi-tenant-deployments) header (absent means `tenant.Default`), validated by `tenant.Parse` (`400`), looked up in the `settings.Registry` (`resolveStore`: `404` for an id it does not hold, a bare `503` for a tenant whose folder was rejected — the findings stay out of a body answered before authentication), and the resolved `*settings.Store` stored in the request context (`WithStore` / `StoreFromContext` — here rather than in `tenant/`, because `settings` names `tenant.ID`). A handler reads the store once and passes it down as an argument — the per-tenant getters it holds take it as a parameter (`(*settings.Store).Policy`, `.DedupeFor`, `.DefaultMaxRows`, … in production) — and nothing below a handler reads the context; a tenant route reached without a resolved store answers `500` rather than fall back to a tenant. The probes, `/version`, the metrics path, and `/v1/ops/*` are tenant-exempt; an ops route that addresses one tenant — the admin pipe reads, the schema routes, the raw-SQL proxy, the settings reload and the DLQ stats — names it in `?tenant=` (`opsTenant`, and `opsStore` over it for the routes that need the tenant's store; the DLQ stats need none, since the MQ holds the queue), parsed strictly so that a query `url.ParseQuery` would half-read is a `400` rather than a read of the default tenant, which is what it means when absent on the reads (on the reload, absent is the whole directory). -- **pipes.go** — Named query pipe handlers: admin listing (`GET /v1/ops/pipes[/{name}]`, read per request from its `pipes.Source`) and execution with parameter binding. `pipes.json` is the only write path. +- **pipes.go** — Named query pipe handlers: admin listing (`GET /v1/ops/pipes[/{name}]`, read per request from its `pipes.Source`) and execution with parameter binding. A read is cached and coalesced; a write — bound SQL that `IsMutation` (`clickhouse_exec.go`) classifies as one — bypasses both and runs every call. `pipes.json` is the only way to define or change a pipe. - **structured_query.go** — Handler for `POST /v1/query?table={table}`: validates query AST, enforces permissions, builds and executes SQL. -- **ch_errors.go** — `writeCHError`, the one mapping from a failed ClickHouse query to a response, shared by `/v1/query`, pipes and `/v1/ops/query` so they cannot drift apart: `chconn.Classify` decides the class, and the class the status, `code` and `retryable` ([ClickHouse errors on the query paths](/api#clickhouse-errors-on-the-query-paths)). +- **ch_errors.go** — `writeCHError`, the one mapping from a failed ClickHouse query to a response, shared by `/v1/query`, pipes and `/v1/ops/query` so they cannot drift apart: `chconn.Classify` decides the class, and the class the status, `code` and `retryable` ([ClickHouse errors on the query paths](/api#clickhouse-errors-on-the-query-paths)). A write pipe answers through `writeCHWriteError`, the same mapping with `retryable` always `false` and no `Retry-After`, since the write may have run. - **ingest.go** — Accepts `POST /v1/ingest?table={table}` in three body shapes: one flat JSON object, a JSON array of them, or NDJSON. The **required** `Content-Type` chooses the format *family* — `application/json` versus the four NDJSON spellings — and within the JSON family the body's first non-whitespace byte picks array versus single object; the bytes never choose the family. Anything that is not exactly one readable media type is a `415`, decided before the body is read: the header is parsed per RFC 9110 §8.3, and because `Content-Type` is a singleton field, repeated header lines must all resolve to the same format and a value carrying a comma is refused unless the value as a whole parses as one media type — a comma inside a *quoted* parameter value is data, so `application/json; a=", application/x-ndjson; b="` is accepted. It then reads the whole (`MaxBytesReader`-capped) body into a pooled buffer and runs the per-format record readers over those bytes, so the `413` lands before any record is processed and peak memory per request is O(body) rather than O(record). Then it validates each record against the discovered schema, optional dedup, and publishes each row through `mq.Publisher` on `mq.Topic{Tenant, Table, Scope}` (the request's tenant, read off its resolved store — `store.Tenant()` — and raw names; the subject it becomes is `internal/mq`'s; a full queue comes back as `mq.ErrQueueFull`, which is the `503` + `Retry-After: 30`, and a broker that cannot be reached or does not answer in time as `mq.ErrUnavailable`, the `503` + `Retry-After: 5`). When dedup is on, a row missing the configured `id_field` can't be deduped: it is logged at `WARN` and counted by `wavehouse_ingest_dedupe_missing_id_total` (labeled by `table`), then published un-deduped — or rejected when `dedupe.require_id` is set ([#219](https://github.com/Wave-RF/WaveHouse/issues/219)). - **query.go** — Proxies raw SQL for `POST /v1/ops/query` straight to the `?tenant=`'s ClickHouse HTTP interface (`chconn.Pools.Target` by the resolved store's tenant; the zero target — no pool — is a `503` with `Retry-After`). **Not cached** — sets `Cache-Control: no-store` so every request hits ClickHouse; DateTime is rendered ISO-8601 via `date_time_output_format=iso` (the Go-side type conversion lives in the structured-query / pipes path, not here). - **stream.go** — Real-time streaming via SSE. Callers select a table with the `?table=` query parameter. Each connection registers one `Subscriber` (the `stream/` package) with both the event `Hub` (under its `(topic, role)`) and the shared keepalive wheel, then drains both from a single byte-pump — so idle streams keep emitting `:` keepalive comments (surviving reverse-proxy idle timeouts) while live events arrive already projected and serialized. Per-event projection/serialization happens **once per role** in the `Hub`, not once per subscriber ([#294](https://github.com/Wave-RF/WaveHouse/issues/294)); the handler also snapshots the connection's JWT claims onto the `Subscriber`, which the `Hub` evaluates per subscriber when the role carries a row-level `filter` ([#319](https://github.com/Wave-RF/WaveHouse/issues/319)). Gap-fill replay (`mq.Replayer.ReplaySince` on the connection's `mq.Topic` — a `DeliverByStartTime` consumer inside `internal/mq`) stays per-connection (low-volume, one-time on connect). A stream ends, a gap-fill in progress included, when the server begins shutting down (`Closing`) or its `Subscriber` is evicted because its tenant is no longer served (`Hub.Prune`); one admitted just before the reload that stopped serving its tenant, and registered just after the prune, is ended right after it registers (`Served`). @@ -149,7 +149,7 @@ The SSE fan-out, factored out of `api/` so the delivery hot path ([#294](https:/ ### `ingest/` — Ingest Pipeline, DLQ & Sweeping -- **worker.go** — `StartIngestWorker` launches an ingest pipeline: a durable `buffer-consumer` consumer of the ingest queue (created through `mq.ConsumerManager`) reads events, batches them per tenant table — the tenant read off each message's `mq.Topic` — and performs bulk INSERTs to ClickHouse. The pipeline is **insert-only**. The wire format `EventMessage` carries `{table_name, scope, received_timestamp, format, columns, row}` — the row positionally as one `JSONCompactEachRow` line, with `columns` naming its positions (the table's insertable columns — a computed one cannot be named in an `INSERT`); the worker batches per (tenant, table, column list) and writes `INSERT INTO … (cols) FORMAT JSONCompactEachRow`. It accepts any table name (events are addressed by `mq.Topic{Tenant, Table, Scope}` with raw names; `internal/mq` encodes them into subject tokens), then bulk-INSERTs. The embedded NATS server runs with `DontListen: true` (`internal/mq/embedded.go`), so the only publishers that can reach the ingest queue 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 (`policy.admin_role`) — see the Query Path section below; the `/v1/ops/*` `RequireAdmin` middleware enforces the check at the API layer, so a no/invalid-token request (resolved to `default_role`, not admin in a production config) never reaches the proxy. A batch whose tenant has no ClickHouse connection (no longer served, or no pool could be opened for it, such as by the connection ceiling) is never tried: no row of it could pass, so `parkBatch` takes it to the DLQ switch whole, logging once per batch rather than twice per row. Otherwise a bulk-insert failure is first classed by `chconn.Classify`: a ClickHouse that cannot take the insert (unavailable, denied, or no verdict at all) sends the batch back to the MQ for a delayed redelivery (`retryLater` → `mq.Message.NakWithDelay`), under a backoff shared by every table on the same pool (a failure of one table — read-only, too many parts — backs off that table alone), and never to the DLQ — the same when it stops answering mid-isolation. Only when ClickHouse rejects the batch, or refuses a multi-row batch for its size (`chconn.Splittable`: too many partitions for one INSERT, the memory limit), is it re-inserted row by row: rows that succeed are acked, and only the rows ClickHouse rejects again are routed to the DLQ (`sendToDLQ` → `mq.DeadLetterer.DeadLetter`), which parks the as-published `EventMessage` envelope under the topic it arrived on (`dlq.{tenant}.{table}` subjects inside `internal/mq`) with the failure context in `X-DLQ-*` headers when the tenant's `dlq.enabled` is on for the table — see [Ingest Pipeline](/ingest-pipeline) for the worker internals. +- **worker.go** — `StartIngestWorker` launches an ingest pipeline: a durable `buffer-consumer` consumer of the ingest queue (created through `mq.ConsumerManager`) reads events, batches them per tenant table — the tenant read off each message's `mq.Topic` — and performs bulk INSERTs to ClickHouse. The pipeline is **insert-only**. The wire format `EventMessage` carries `{table_name, scope, received_timestamp, format, columns, row}` — the row positionally as one `JSONCompactEachRow` line, with `columns` naming its positions (the table's insertable columns — a computed one cannot be named in an `INSERT`); the worker batches per (tenant, table, column list) and writes `INSERT INTO … (cols) FORMAT JSONCompactEachRow`. It accepts any table name (events are addressed by `mq.Topic{Tenant, Table, Scope}` with raw names; `internal/mq` encodes them into subject tokens), then bulk-INSERTs. The embedded NATS server runs with `DontListen: true` (`internal/mq/embedded.go`), so the only publishers that can reach the ingest queue 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 (`policy.admin_role`) — see the Query Path section below; the `/v1/ops/*` `RequireAdmin` middleware enforces the check at the API layer, so a no/invalid-token request (resolved to `default_role`, not admin in a production config) never reaches the proxy — or through an operator-authored [pipe that writes](/pipes#pipes-that-write), gated only by its `allowed_roles`. A batch whose tenant has no ClickHouse connection (no longer served, or no pool could be opened for it, such as by the connection ceiling) is never tried: no row of it could pass, so `parkBatch` takes it to the DLQ switch whole, logging once per batch rather than twice per row. Otherwise a bulk-insert failure is first classed by `chconn.Classify`: a ClickHouse that cannot take the insert (unavailable, denied, or no verdict at all) sends the batch back to the MQ for a delayed redelivery (`retryLater` → `mq.Message.NakWithDelay`), under a backoff shared by every table on the same pool (a failure of one table — read-only, too many parts — backs off that table alone), and never to the DLQ — the same when it stops answering mid-isolation. Only when ClickHouse rejects the batch, or refuses a multi-row batch for its size (`chconn.Splittable`: too many partitions for one INSERT, the memory limit), is it re-inserted row by row: rows that succeed are acked, and only the rows ClickHouse rejects again are routed to the DLQ (`sendToDLQ` → `mq.DeadLetterer.DeadLetter`), which parks the as-published `EventMessage` envelope under the topic it arrived on (`dlq.{tenant}.{table}` subjects inside `internal/mq`) with the failure context in `X-DLQ-*` headers when the tenant's `dlq.enabled` is on for the table — see [Ingest Pipeline](/ingest-pipeline) for the worker internals. - **backoff.go** — The retry backoff behind `retryLater`: a small circuit breaker per ClickHouse pool (the target's URL, user and database), and one per pool and table for a failure of one table (`chconn.TableScoped`). A failure opens it for 1 s, doubling to a 30 s cap, each window jittered down to half; while it is open, flushes and arriving rows are handed back without a request, and once it elapses one flush probes. Any answer that is not an outage closes it. - **types.go** — `EventMessage` struct (TableName, Scope — reserved, always empty today, ReceivedTimestamp, Format, Columns, Row; `Format` is `FormatJSONCompactEachRow` and `Row` is one positional line whose slots `Columns` names) and `BufferConsumerName` constant, shared across API handlers and the ingest pipeline. - **compact.go** — `EncodeCompactRow`, the positional row encoder every published row goes through, rendering one record over the table's **insertable** columns in declaration order. Serialization only: it validates nothing and judges no value. @@ -272,8 +272,9 @@ Ingest worker pipeline (StartIngestWorker): (Insert-only pipeline. The wire format `EventMessage` carries only {table_name, scope, received_timestamp, format, columns, row}; non-insert mutations - DELETE/UPDATE/TRUNCATE/DROP/etc. must go through POST /v1/ops/query — the - /v1/ops/* RequireAdmin gate rejects non-admin callers at the API layer, so + DELETE/UPDATE/TRUNCATE/DROP/etc. must go through POST /v1/ops/query (or an + operator-authored write pipe) — the /v1/ops/* RequireAdmin gate rejects + non-admin callers at the API layer, so a no/invalid-token request (resolved to default_role, not admin in a production config) cannot reach the proxy.) @@ -302,10 +303,11 @@ Client POST /v1/ops/query 401 when a stashed error shows the caller presented an invalid token, else 403. Raw SQL has no per-statement scope check (a full SQL parser would be needed to authorize predicates), so the role gate is the - entire authorization story. /v1/ops/query is the only sanctioned - surface for non-SELECT statements (DELETE/UPDATE/TRUNCATE/DROP/ALTER/…); - non-admin callers use `POST /v1/ingest?table={table}` for writes and - the structured query endpoint or named pipes for reads. + entire authorization story. /v1/ops/query is the only surface for + ad-hoc non-SELECT statements (DELETE/UPDATE/TRUNCATE/DROP/ALTER/…); + non-admin callers use `POST /v1/ingest?table={table}` or a write pipe + that lists their role for writes, and the structured query endpoint or + named pipes for reads. → Decode {"sql": "..."} from the request body. → POST the SQL verbatim to ClickHouse's HTTP interface at ://:/?default_format=JSON @@ -330,7 +332,7 @@ Client POST /v1/ops/query (browser, CDN, corp proxy) caches the result. ``` -The proxy-pattern wins are: zero classification logic on the WaveHouse side (no isMutation heuristic to maintain), and any ClickHouse statement type — including verbs added in future versions and inline FORMAT overrides — works without WaveHouse code changes. Multi-statement input (`SELECT 1; TRUNCATE t`) is supported when the upstream ClickHouse has multi-query enabled, which is the default on recent versions; older or restrictively-configured servers will return a clear error from ClickHouse itself for the second statement. The proxy buffers the response in memory with a 64 MiB cap (502 with `clickhouse response exceeded N bytes` on overflow, to keep a runaway `SELECT *` from pinning RAM on the API server), and passes ClickHouse's `Content-Type` through when an inline `FORMAT` directive overrides the default JSON envelope. The structured query endpoint and pipes still go through `clickhouse-go`'s native driver (Query/Exec) for performance and to keep the cached row-array shape consistent. +The proxy-pattern wins are: zero classification logic on the WaveHouse side (no `IsMutation` heuristic to maintain), and any ClickHouse statement type — including verbs added in future versions and inline FORMAT overrides — works without WaveHouse code changes. Multi-statement input (`SELECT 1; TRUNCATE t`) is supported when the upstream ClickHouse has multi-query enabled, which is the default on recent versions; older or restrictively-configured servers will return a clear error from ClickHouse itself for the second statement. The proxy buffers the response in memory with a 64 MiB cap (502 with `clickhouse response exceeded N bytes` on overflow, to keep a runaway `SELECT *` from pinning RAM on the API server), and passes ClickHouse's `Content-Type` through when an inline `FORMAT` directive overrides the default JSON envelope. The structured query endpoint and pipes still go through `clickhouse-go`'s native driver (Query/Exec) for performance and to keep the cached row-array shape consistent. ### Streaming Path diff --git a/docs/src/content/docs/configuration.mdx b/docs/src/content/docs/configuration.mdx index 5fac488b7..b4ceb2295 100644 --- a/docs/src/content/docs/configuration.mdx +++ b/docs/src/content/docs/configuration.mdx @@ -126,7 +126,7 @@ Set the backstop on the profile of the ClickHouse user WaveHouse connects as (th ``` -What a caller sees when one of these trips: a row or byte limit is `400 clickhouse.limit_exceeded`; a server-wide time, memory or quota limit is `503 clickhouse.unavailable` (retryable, so the SDK retries it), since WaveHouse cannot tell it from server pressure — except on `/v1/query` for a role that sets its own `max_memory_usage`, where every memory-limit error is read as that cap and answered `400 clickhouse.limit_exceeded`. See [ClickHouse errors on the query paths](/api#clickhouse-errors-on-the-query-paths). +What a caller sees when one of these trips: a row or byte limit is `400 clickhouse.limit_exceeded`; a server-wide time, memory or quota limit is `503 clickhouse.unavailable` (retryable, so the SDK retries it — except on a [pipe that writes](/pipes#pipes-that-write), which is never retryable), since WaveHouse cannot tell it from server pressure — except on `/v1/query` for a role that sets its own `max_memory_usage`, where every memory-limit error is read as that cap and answered `400 clickhouse.limit_exceeded`. See [ClickHouse errors on the query paths](/api#clickhouse-errors-on-the-query-paths). :::caution[How the two layers compose] WaveHouse's per-role caps are sent as per-query `SETTINGS` on its connection, so they **compose** with the ClickHouse profile — a per-role cap *tightens* within the profile's ceiling, and a `` block bounds how far any setting can move. But if the profile marks a setting `readonly` (or `` disallows changing it), ClickHouse will **reject** WaveHouse's per-query override and the query fails. So keep the settings WaveHouse manages (`max_memory_usage`, `max_execution_time`, `max_rows_to_read`, `max_result_rows`) **changeable** for its user — use a `` constraint, not `readonly`, if you want a hard ceiling. diff --git a/docs/src/content/docs/ingest-pipeline.md b/docs/src/content/docs/ingest-pipeline.md index f27356c0e..77bb97d24 100644 --- a/docs/src/content/docs/ingest-pipeline.md +++ b/docs/src/content/docs/ingest-pipeline.md @@ -19,7 +19,7 @@ It is deliberately detailed: this is a hot, concurrency-heavy path, and the goro | `sweeper.go` | The **Active Sweeper** — every minute, asks the MQ to purge the events that are both written to ClickHouse and past the SSE gap window (the purge arithmetic below lives in `internal/mq/purge.go`) | | `types.go` | `EventMessage` wire format and the `BufferConsumerName` constant | -The pipeline is **insert-only**. (Upgrading across the v2 envelope? [Drain the queue first](/deployment#upgrading-across-the-v2-ingest-envelope).) The wire format carries `{table_name, scope, received_timestamp, format, columns, row}`: `row` is one `JSONCompactEachRow` line — a positional JSON array — and `columns` names its positions — the table's insertable columns, in declaration order (a `MATERIALIZED` or `ALIAS` column cannot be named in an `INSERT`, so it is not part of the row's contract). (`scope` is reserved and always `""` today.) Each NATS message is its own envelope, so the names ride along per record; where they are carried once is the `INSERT` the worker emits per group. The worker parses the envelope, groups a batch by column list, and bulk-`INSERT`s each group as `INSERT INTO … (cols) FORMAT JSONCompactEachRow` — schema validation already happened at the HTTP ingest handler, before publish. Non-insert mutations go through `POST /v1/ops/query` (admin-only). +The pipeline is **insert-only**. (Upgrading across the v2 envelope? [Drain the queue first](/deployment#upgrading-across-the-v2-ingest-envelope).) The wire format carries `{table_name, scope, received_timestamp, format, columns, row}`: `row` is one `JSONCompactEachRow` line — a positional JSON array — and `columns` names its positions — the table's insertable columns, in declaration order (a `MATERIALIZED` or `ALIAS` column cannot be named in an `INSERT`, so it is not part of the row's contract). (`scope` is reserved and always `""` today.) Each NATS message is its own envelope, so the names ride along per record; where they are carried once is the `INSERT` the worker emits per group. The worker parses the envelope, groups a batch by column list, and bulk-`INSERT`s each group as `INSERT INTO … (cols) FORMAT JSONCompactEachRow` — schema validation already happened at the HTTP ingest handler, before publish. Non-insert mutations go through `POST /v1/ops/query` (admin-only) or an operator-authored [pipe that writes](/pipes#pipes-that-write). ## High-level shape diff --git a/docs/src/content/docs/pipes.mdx b/docs/src/content/docs/pipes.mdx index b8c7efa04..a6ce6417d 100644 --- a/docs/src/content/docs/pipes.mdx +++ b/docs/src/content/docs/pipes.mdx @@ -7,7 +7,7 @@ sidebar: A **named pipe** is a saved SQL query, registered under a name, that callers run by name with parameters — without ever sending raw SQL. They turn an ad-hoc query into a stable, cached, access-controlled endpoint: you write the SQL once as an operator in the settings directory's [`pipes.json`](/settings-directory#pipesjson), expose it at `GET/POST /v1/pipes/{name}`, and clients supply only the declared parameters. -Pipes are the right tool when a query is reusable and shouldn't live in client code — dashboards, reports, public APIs over curated slices of data. They sit on the **cached read path** (shared L1 + singleflight, same as structured queries), and authorize through a simple per-pipe allowlist rather than the full [policy engine](/access-control). +Pipes are the right tool when a query is reusable and shouldn't live in client code — dashboards, reports, public APIs over curated slices of data. They sit on the **cached read path** (shared L1 + singleflight, same as structured queries; a [pipe that writes](#pipes-that-write) bypasses both), and authorize through a simple per-pipe allowlist rather than the full [policy engine](/access-control). ## Anatomy of a pipe @@ -78,7 +78,7 @@ Independently, every `required` declared parameter must be supplied regardless o ### How a value becomes SQL -Bound values are **inlined directly into the SQL string** (not sent as positional driver parameters — that lets a parameter sit anywhere ClickHouse allows a literal, including `LIMIT`). Inlining is type-aware and escaped: +Bound values are **inlined directly into the SQL string** (not sent as positional driver parameters — that lets a parameter sit anywhere ClickHouse allows a literal, including `LIMIT`). Inlining is type-aware and escaped, and a string brings its own quotes, so **write each placeholder bare, never inside quotes** — `WHERE id = {{id}}`, not `WHERE id = '{{id}}'`: | Supplied value | Rendered as | Note | | -------------- | ----------- | ---- | @@ -89,7 +89,7 @@ Bound values are **inlined directly into the SQL string** (not sent as positiona | array | `('a', 'b')` | parenthesized list of escaped elements — for `IN` clauses (see below) | | null | `NULL` | | -Every leaf value is escaped the same way — including each element of an array — so a parameter value can't break out of its literal and inject SQL. The SQL *structure* still comes only from the operator-authored template. Values with no safe scalar form are **rejected** with a `400`: a JSON object, and an empty array (which would render as the invalid `IN ()`). +Every leaf value is escaped the same way — including each element of an array — so a value bound to a bare placeholder can't break out of its literal and inject SQL, and the SQL *structure* comes only from the operator-authored template. A quoted placeholder gives that up: the value's own quotes close the template's, so in `WHERE id = '{{id}}'` the body `{"id": " OR 1=1 OR id = "}` binds to `WHERE id = '' OR 1=1 OR id = ''`, which matches every row. Values with no safe scalar form are **rejected** with a `400`: a JSON object, and an empty array (which would render as the invalid `IN ()`). #### Array parameters and `IN` lists @@ -122,7 +122,7 @@ This is the *only* authorization check on the execute path — pipes deliberatel ## Creating and managing pipes -Pipes are defined in the settings directory's `pipes.json` — a `{"pipes": [...]}` list of the definitions above — and the files are the only write path: edit the file (standalone: on the host; on WaveHouse Cloud the control plane writes it) and the running server re-validates and adopts it on file change, `SIGHUP`, or `POST /v1/ops/settings/reload`, so a create, update, or delete applies without a restart. Validation (`wavehouse validate`, boot, and every reload) rejects an unknown key, a duplicate or empty name, empty SQL, an unknown parameter `type`, and an `allowed_roles` entry not declared in `roles.json`; a rejected reload keeps the previous pipes in effect. See [Settings Directory — `pipes.json`](/settings-directory#pipesjson) for the full rules. +Pipes are defined in the settings directory's `pipes.json` — a `{"pipes": [...]}` list of the definitions above — and the files are the only way to define or change one: edit the file (standalone: on the host; on WaveHouse Cloud the control plane writes it) and the running server re-validates and adopts it on file change, `SIGHUP`, or `POST /v1/ops/settings/reload`, so a create, update, or delete applies without a restart. Validation (`wavehouse validate`, boot, and every reload) rejects an unknown key, a duplicate or empty name, empty SQL, an unknown parameter `type`, and an `allowed_roles` entry not declared in `roles.json`; a rejected reload keeps the previous pipes in effect. See [Settings Directory — `pipes.json`](/settings-directory#pipesjson) for the full rules. ```json { @@ -169,7 +169,7 @@ curl -X POST http://localhost:8080/v1/pipes/top_pages \ -d '{"start_date": "2024-01-01", "limit": 20}' ``` -The response is a JSON array of rows. Results flow through the shared in-process L1 cache (Ristretto) with singleflight coalescing, so concurrent identical calls hit ClickHouse once; an `X-Cache: HIT` or `X-Cache: MISS` header tells you which path served the response. +The response is a JSON array of rows. Results flow through the shared in-process L1 cache (Ristretto) with singleflight coalescing, so concurrent identical calls hit ClickHouse once; an `X-Cache: HIT` or `X-Cache: MISS` header tells you which path served the response. A [pipe that writes](#pipes-that-write) skips both and answers `X-Cache: BYPASS`. | Status | Body | Cause | | ------ | ---- | ----- | @@ -178,6 +178,16 @@ The response is a JSON array of rows. Results flow through the shared in-process | 400 | `{"error":"missing required parameter: x"}` | A required parameter wasn't supplied | | 400 | `{"error":"parameter \"x\": unsupported parameter type object"}` | A non-scalar value with no SQL form — a JSON object (directly, or nested in an array). An empty array is likewise rejected (`array parameter must not be empty`). | +### Pipes that write + +A pipe's SQL may be a write: a statement led by a write verb WaveHouse recognizes — `INSERT`, `UPDATE`, `DELETE`, `ALTER` (so `ALTER … DELETE`), `CREATE`, `DROP`, `TRUNCATE`, `RENAME`, `EXCHANGE`, `REPLACE`, `OPTIMIZE`, `ATTACH`, `DETACH`, `GRANT`, `REVOKE`, `KILL`, `SET`, `USE` or `SYSTEM` — directly, or `INSERT INTO` after a `WITH` list (`WITH … INSERT INTO …`, the only write ClickHouse accepts there). An `EXECUTE AS ` prefix is looked through: the statement after it is the one classified. Such a pipe runs on **every** call: it never reads or fills the cache and is never coalesced with an identical call in flight, so ten identical calls are ten writes. The response is `[]` with `X-Cache: BYPASS` and `Cache-Control: no-store`, so an HTTP cache in front of a `GET` does not answer a repeat either. WaveHouse classifies the statement by its leading keyword — a `WITH`-led one by whether it holds `INSERT INTO` outside parentheses — with the same classifier that sends it to ClickHouse as a write, so no pipe property marks it. A statement led by any other keyword runs as a read: one that returns rows (`BACKUP`, `RESTORE`) is cached and coalesced, so a repeat within the TTL does not run; one that returns none (`UNDROP`, `MOVE`) fails the call after it has run, which the SDK may retry — don't put a write led by another verb in a pipe ([#666](https://github.com/Wave-RF/WaveHouse/issues/666)). + +A failed write is not retried automatically, because it may have run. It answers with the status and `code` a failed read would ([ClickHouse errors on the query paths](/api#clickhouse-errors-on-the-query-paths)), but always with `retryable: false` and no `Retry-After`, `503 clickhouse.unavailable` included: once the statement is on its way to ClickHouse, WaveHouse cannot tell whether it ran. The [SDK](/sdk/pipes) does not retry such an answer, so check whether the write landed before you send it again. A call refused before anything is sent — the tenant on no ClickHouse pool, `503` with `Retry-After: 30` — cannot have run, and the SDK retries it. The SDK also retries when WaveHouse's own answer never reaches it — a dropped connection, or a `502`/`503`/`504` from a proxy in front of WaveHouse that gave up waiting — so a write can still run twice that way; give a client that runs write pipes [`options.maxRetries`](/sdk#clientconfigdb) `0` if that matters. + +`allowed_roles` is a write pipe's only gate: the [policy engine](/access-control)'s insert rules do not apply to it, so any role you list — including a [`default_role`](/access-control#default_role--public-unauthenticated-access) that anonymous callers resolve to — can run the write. The operator fixes the statement and its predicate when authoring the pipe; callers supply only literal values, provided every placeholder is written bare ([how a value becomes SQL](#how-a-value-becomes-sql)). + +Two things a write pipe does not do yet: it does not invalidate cached reads of the table it writes — a structured query over that table can serve pre-write rows until its TTL ([#394](https://github.com/Wave-RF/WaveHouse/issues/394)); a read pipe's cached result expires only with its TTL whatever writes the table, ingest included ([#343](https://github.com/Wave-RF/WaveHouse/issues/343)) — and its rows do not reach [`/v1/stream`](/api#get-v1stream--server-sent-events-stream) subscribers, which only the [ingest pipeline](/ingest-pipeline) feeds ([#362](https://github.com/Wave-RF/WaveHouse/issues/362)). For writes that stream subscribers and structured queries should see at once, use [`POST /v1/ingest`](/api#post-v1ingesttabletable--ingest-data). + ## End-to-end example Ship a curated "top pages" endpoint that the public dashboard can call with no token. diff --git a/docs/src/content/docs/sdk/pipes.md b/docs/src/content/docs/sdk/pipes.md index 1bae86bf0..04bab0c98 100644 --- a/docs/src/content/docs/sdk/pipes.md +++ b/docs/src/content/docs/sdk/pipes.md @@ -17,12 +17,14 @@ const { data } = await wh.pipe('top_pages', { start_date: '2026-01-01', limit: 5 ### `.fetch(opts?)` -Execute and return results. Takes `PipeRequestOptions` — `{ signal }` only, narrower than the `.fetch(opts?)` on a [query builder](/sdk/queries), which also accepts `limit`. Passing a `limit` is a compile error rather than a silent no-op. +Execute and return results. A [pipe that writes](/pipes#pipes-that-write) returns `[]`. `.fetch()` takes `PipeRequestOptions` — `{ signal }` only, narrower than the `.fetch(opts?)` on a [query builder](/sdk/queries), which also accepts `limit`. Passing a `limit` is a compile error rather than a silent no-op. `limit` is typed `never` rather than left out, so the rejection also catches a value passed in a variable — leaving it out would only reject an inline object. That cuts both ways: a value *declared* as `RequestOptions` is rejected whether or not it actually carries a limit, since the type permits one. If you share one options object across calls, type it as `PipeRequestOptions` — the table and query-builder `.fetch()` accept that too — or inline `{ signal }` at the pipe call. There is no per-call row cap here: the endpoint binds your `params` as the pipe's parameters, so a limit has to be declared in the pipe's SQL as `{{limit}}` (see [Named Pipes](/pipes)) and passed as `wh.pipe(name, { limit })`, as in the example above. +A ClickHouse failure on a [pipe that writes](/pipes#pipes-that-write) comes back `retryable: false`, and the SDK does not retry it. The SDK does still retry when WaveHouse's own answer never reaches it — a dropped connection, or a `502`/`503`/`504` from a proxy in front of WaveHouse that gave up waiting — so a write can run twice that way. If that matters, give the client that runs write pipes [`options.maxRetries`](/sdk#clientconfigdb) `0`. + ### `.stream(opts?)` Open a live stream (see [Streaming](/sdk/streaming)). @@ -31,7 +33,7 @@ Open a live stream (see [Streaming](/sdk/streaming)). ## Pipes Admin — `wh.pipes` -Inspect the adopted named query pipes. Requires the admin gate — the admin role (`policy.admin_role`) or the [operator key](/api#authentication). Pipes are defined in the server's settings directory `pipes.json` — files are the only write path, so there is no `set` or `delete`: edit the file and let the server pick it up, or call [`wh.settings.reload()`](/sdk/admin#settings--whsettings). +Inspect the adopted named query pipes. Requires the admin gate — the admin role (`policy.admin_role`) or the [operator key](/api#authentication). Pipes are defined in the server's settings directory `pipes.json` — the files are the only way to define or change a pipe, so there is no `set` or `delete`: edit the file and let the server pick it up, or call [`wh.settings.reload()`](/sdk/admin#settings--whsettings). ```ts // List all pipes diff --git a/docs/src/content/docs/sdk/reference.md b/docs/src/content/docs/sdk/reference.md index fa6983da4..9db138f3f 100644 --- a/docs/src/content/docs/sdk/reference.md +++ b/docs/src/content/docs/sdk/reference.md @@ -25,7 +25,7 @@ if (error?.code === 'ABORTED') { The SDK **never throws** for anything the server returns — all API errors come back in `Result.error`. It does throw on caller and environment errors: a non-absolute `baseURL` (REST calls reject with a `TypeError`; streams report `SSE_CONNECT_ERROR` to the subscriber's `error` callback — see [Serving under a path prefix](/sdk#serving-under-a-path-prefix)), `.stream()` / `.liveQuery()` in a runtime with no global `fetch` and no `options.fetch` (see [Runtime support](/sdk#runtime-support)), and an `auth` callback that rejects — a token-refresh failure propagates out of the REST call, and on a stream is reported as a retryable `SSE_AUTH_ERROR`. One more exception escapes an SDK call synchronously, though it is yours rather than ours: your own `status` handler throwing on the first `.subscribe()` or `.liveQuery()`, described under *If your own callback throws* below. -`code` and `retryable` are the server's own when its error body carries them — a failed ClickHouse query does, with codes like `clickhouse.rejected` and `clickhouse.unavailable` ([the full list](/api#clickhouse-errors-on-the-query-paths)). Otherwise `code` is `HTTP_` and a `5xx` is retryable. +`code` and `retryable` are the server's own when its error body carries them — a failed ClickHouse query does, with codes like `clickhouse.rejected` and `clickhouse.unavailable` ([the full list](/api#clickhouse-errors-on-the-query-paths)). Otherwise `code` is `HTTP_` and a `5xx` is retryable. One exception to the table below: a [pipe that writes](/pipes#pipes-that-write) answers every ClickHouse failure `retryable: false` with no `Retry-After`, `clickhouse.unavailable` and `clickhouse.unknown` included, so the SDK returns it on the first attempt. | Status | Code | Retryable | Description | |--------|------|-----------|-------------| diff --git a/docs/src/content/docs/settings-directory.mdx b/docs/src/content/docs/settings-directory.mdx index 900e6a2c2..67992b6fb 100644 --- a/docs/src/content/docs/settings-directory.mdx +++ b/docs/src/content/docs/settings-directory.mdx @@ -108,7 +108,7 @@ The tenant tunables. Every key is required (a missing one is a validation error) | `clickhouse.http_scheme` | `http` | `http` or `https` for that HTTP hop — one of the two *outbound* TLS switches, with `tls.enabled` for the native hop; unrelated to your clients' TLS. | | `clickhouse.database` | `default` | Database tables are discovered from. | | `clickhouse.username` | `default` | Connection user; the password is boot config (`WH_CH_PASSWORD`). | -| `clickhouse.query_timeout` | `30` | Seconds (`>= 1`) a read may take. On `/v1/query` under a role's `max_execution_time`, the smaller of the two is sent to ClickHouse as `max_execution_time`; otherwise it bounds the client deadline, from which the driver derives a server-side `max_execution_time`. | +| `clickhouse.query_timeout` | `30` | Seconds (`>= 1`) a ClickHouse call may take on the query paths — structured queries, pipes (a write pipe included) and `/v1/ops/query`. On `/v1/query` under a role's `max_execution_time`, the smaller of the two is sent to ClickHouse as `max_execution_time`; otherwise, on the native paths, it bounds the client deadline, from which the driver derives a server-side `max_execution_time`; on `/v1/ops/query` it is the HTTP request's deadline. | | `clickhouse.tls.enabled` | `false` | Switches the native-protocol hop (`addr`) to TLS. The HTTP hop's switch stays `http_scheme`; the rest of the `tls` block applies to whichever hop uses TLS. See [ClickHouse](#clickhouse). | | `clickhouse.tls.ca_file` | `""` | PEM bundle the server certificate is verified against; empty uses the system roots. A path, read when the connection is built and re-read when the `tls` block changes — validation does not open it. | | `clickhouse.tls.cert_file` | `""` | Client certificate for mutual TLS, PEM; set together with `key_file` or not at all. | diff --git a/internal/api/ch_errors.go b/internal/api/ch_errors.go index e2d6c8e75..adb61aaf3 100644 --- a/internal/api/ch_errors.go +++ b/internal/api/ch_errors.go @@ -28,11 +28,13 @@ const ( // caller's. 502. codeCHMisconfigured = "clickhouse.misconfigured" // codeCHUnavailable: ClickHouse, or the way to it, could not take the - // query now. 503 with Retry-After. + // query now. 503 with Retry-After — without it for a write pipe, which + // may have run and is never retryable. codeCHUnavailable = "clickhouse.unavailable" // codeCHResponseTooLarge: the raw-SQL proxy's response cap. 502. codeCHResponseTooLarge = "clickhouse.response_too_large" - // codeCHUnknown: a failure with no verdict. 5xx, retryable. + // codeCHUnknown: a failure with no verdict. 5xx, retryable unless a + // write pipe's. codeCHUnknown = "clickhouse.unknown" ) @@ -119,10 +121,24 @@ func chFailureOf(err error, unknownStatus int, caps queryCaps) chFailure { // it is WaveHouse's configuration being refused, which an operator should // hear about even when the caller only sees a 403. func writeCHError(w http.ResponseWriter, r *http.Request, err error, message string, unknownStatus int, caps queryCaps) { - f := chFailureOf(err, unknownStatus, caps) + writeCHFailure(w, r, err, message, chFailureOf(err, unknownStatus, caps)) +} + +// writeCHWriteError answers a failed write pipe as writeCHError does, but +// never as retryable and with no Retry-After: the statement may have reached +// ClickHouse and run, so a client retrying would run it again. +func writeCHWriteError(w http.ResponseWriter, r *http.Request, err error, message string) { + f := chFailureOf(err, http.StatusInternalServerError, queryCaps{}) + f.retryable = false + writeCHFailure(w, r, err, message, f) +} + +func writeCHFailure(w http.ResponseWriter, r *http.Request, err error, message string, f chFailure) { switch f.code { case codeCHUnavailable: - w.Header().Set("Retry-After", retryAfterClickHouse) + if f.retryable { + w.Header().Set("Retry-After", retryAfterClickHouse) + } case codeCHAccessDenied, codeCHMisconfigured: exCode, _ := chconn.ExceptionCode(err) slog.WarnContext(r.Context(), "clickhouse refused WaveHouse's configuration", diff --git a/internal/api/ch_errors_test.go b/internal/api/ch_errors_test.go index 6daf11ea5..ee2df2db6 100644 --- a/internal/api/ch_errors_test.go +++ b/internal/api/ch_errors_test.go @@ -149,6 +149,32 @@ func TestPipes_ClickHouseErrors(t *testing.T) { } } +// TestPipes_WriteClickHouseErrors: a failed write pipe answers with its +// class's status and code, but never as retryable and with no Retry-After: +// the statement may have run, so a client that retried would run it again. +func TestPipes_WriteClickHouseErrors(t *testing.T) { + t.Parallel() + for _, tc := range chErrorCases(t) { + if tc.caps.MaxExecutionTime > 0 || tc.caps.MaxMemoryUsage > 0 { + continue + } + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + conn := &writeConn{err: tc.err} + h := writerPipesHandler(t, conn, nil, &pipes.NamedQuery{Name: "log", SQL: "INSERT INTO audit_log VALUES ({{msg}}, now())"}) + w := pipeCallAs(t, h, "log") + require.Equal(t, tc.wantStatus, w.Code, w.Body.String()) + var got errorBody + require.NoError(t, json.Unmarshal(w.Body.Bytes(), &got)) + assert.Equal(t, tc.wantCode, got.Code) + require.NotNil(t, got.Retryable) + assert.False(t, *got.Retryable) + assert.Empty(t, w.Header().Get("Retry-After")) + assert.Equal(t, int32(1), conn.execs.Load()) + }) + } +} + type errRow struct{ err error } func (r errRow) Err() error { return r.err } diff --git a/internal/api/clickhouse_exec.go b/internal/api/clickhouse_exec.go index 6a296fc70..51d41e583 100644 --- a/internal/api/clickhouse_exec.go +++ b/internal/api/clickhouse_exec.go @@ -6,6 +6,7 @@ import ( "reflect" "strings" "time" + "unicode/utf8" "github.com/ClickHouse/clickhouse-go/v2/lib/driver" "github.com/Wave-RF/WaveHouse/internal/settings" @@ -42,7 +43,7 @@ func timeoutOf(timeout func(*settings.Store) time.Duration, store *settings.Stor // The raw-SQL endpoint (/v1/ops/query) proxies straight to ClickHouse // over HTTP and never calls this; see internal/api/query.go. func executeCHQuery(ctx context.Context, conn driver.Conn, sql string, params []any) ([]map[string]any, error) { - if isMutation(sql) { + if IsMutation(sql) { if err := conn.Exec(ctx, sql, params...); err != nil { return nil, fmt.Errorf("clickhouse exec: %w", err) } @@ -113,27 +114,26 @@ var mutationVerbs = map[string]struct{}{ "SYSTEM": {}, } -// isMutation reports whether sql's leading statement is a non-SELECT — i.e. +// IsMutation reports whether sql's leading statement is a non-SELECT — i.e. // one that returns no result set and must go through Exec, not Query. -// Leading whitespace and SQL line/block comments are skipped, then the first -// alphabetic token is matched case-insensitively against mutationVerbs. A -// leading WITH clause (CTE) routes through a paren-aware scan because -// ClickHouse accepts `WITH cte AS (...) INSERT INTO t SELECT * FROM cte` as -// equivalent to `INSERT INTO t WITH cte AS (...) SELECT * FROM cte` (see -// https://clickhouse.com/docs/sql-reference/statements/insert-into). Without -// the skip, the WITH form would classify as a read, route through Query, -// silently succeed, and return `[]` — the same silent-success class that -// motivated the original cache-bypass guard. -func isMutation(sql string) bool { +// Leading whitespace and comments are skipped as ClickHouse's lexer skips +// them, then the first bareword is matched whole, case-insensitively, against +// mutationVerbs. After a WITH list ClickHouse parses only SELECT, a FROM-first +// SELECT or INSERT INTO, so a WITH-led statement is a write exactly when it +// holds INSERT INTO at the top level (hasTopLevelInsertInto). An +// `EXECUTE AS ` prefix is looked through to the statement it runs. A +// write classified as a read goes through Query, which runs it and then fails +// the call, so a client that retries the error writes again. +func IsMutation(sql string) bool { s := stripLeadingSQLComments(sql) - end := 0 - for end < len(s) { - c := s[end] - if (c < 'A' || c > 'Z') && (c < 'a' || c > 'z') { - break + if rest, ok := skipExecuteAs(s); ok { + // Bare, it switches the session's user and returns no result set. + if rest == "" || rest[0] == ';' { + return true } - end++ + s = rest } + end := skipWord(s, 0) if end == 0 { return false } @@ -142,52 +142,59 @@ func isMutation(sql string) bool { _, ok := mutationVerbs[first] return ok } - return containsMutationVerbAtTopLevel(s[end:]) + return hasTopLevelInsertInto(s[end:]) } -// nonMutationVerbs is the read/metadata-statement counterpart to -// mutationVerbs. Together they cover every ClickHouse statement-introducing -// keyword that can legally follow a CTE list. The CTE-aware scanner in -// containsMutationVerbAtTopLevel needs the union to identify *which* token -// is the statement keyword — without it, ordinary identifiers in the CTE -// list (table names, database names like the ClickHouse-built-in `system`) -// can collide with mutation-verb names and false-positive the classifier. -var nonMutationVerbs = map[string]struct{}{ - "SELECT": {}, - "SHOW": {}, - "DESCRIBE": {}, - "DESC": {}, - "EXPLAIN": {}, - "EXISTS": {}, - "CHECK": {}, +// skipExecuteAs returns what follows an `EXECUTE AS [@]` prefix +// leading s, past whitespace and comments, and true; or s and false if no such +// prefix leads it. +func skipExecuteAs(s string) (string, bool) { + i := skipWord(s, 0) + if !strings.EqualFold(s[:i], "EXECUTE") { + return s, false + } + i = skipSpaceAndComments(s, i) + j := skipWord(s, i) + if !strings.EqualFold(s[i:j], "AS") { + return s, false + } + i = skipSpaceAndComments(s, skipName(s, skipSpaceAndComments(s, j))) + if i < len(s) && s[i] == '@' { + i = skipSpaceAndComments(s, skipName(s, skipSpaceAndComments(s, i+1))) + } + return s[i:], true } -// containsMutationVerbAtTopLevel scans s for the statement-introducing -// keyword at paren-depth 0, stepping over SQL string literals (`'…'` with -// `”` escape), quoted identifiers (`"…"` and “ `…` “), parenthesized CTE -// subqueries, and SQL comments. The CTE list contains ordinary identifiers -// (CTE names, table/database names) that must not be matched as mutation -// verbs — `system` would otherwise pattern-match `SYSTEM` and route a -// `WITH … SELECT * FROM system.tables` read through `Exec` (silent empty- -// array result instead of the actual rows). Two-part fix: -// -// 1. Skip identifiers whose next non-whitespace, non-comment token is -// `AS` (case-insensitive) or `(` — those are CTE definition names -// (with optional column list before AS). This catches the harder -// class where the CTE alias is itself a mutation-verb name -// (`WITH set AS (…) SELECT …`, `WITH alter AS (…) …`, etc.). -// 2. Among the remaining identifiers, stop on the FIRST that's a -// known statement keyword (mutation OR read-class), and decide -// based on mutationVerbs membership. -// -// Tokens that aren't CTE names and aren't statement keywords (RECURSIVE, -// MATERIALIZED, scalar CTE aliases, etc.) are skipped silently. Returns -// false if no statement keyword is found — the SQL is syntactically -// incomplete or unrecognised; safer to treat as non-mutation than to -// silently route an unknown verb through Exec (an Exec'd SELECT returns -// `[]` with no error; a Query'd unrecognised statement surfaces a clear -// error). -func containsMutationVerbAtTopLevel(s string) bool { +// skipName returns the index just past the user or host name at s[i]: a +// bareword, a quoted identifier or string literal, or a heredoc. +func skipName(s string, i int) int { + if i >= len(s) { + return i + } + switch s[i] { + case '\'', '"', '`': + return skipQuoted(s, i) + case 0xE2: + return skipCurlyQuoted(s, i) + case '$': + if j := skipHeredoc(s, i); j > i { + return j + } + } + return skipWord(s, i) +} + +// hasTopLevelInsertInto reports whether s holds INSERT INTO outside +// parentheses, stepping over string literals and quoted identifiers +// (skipQuoted, skipCurlyQuoted), heredocs (skipHeredoc) and comments +// (skipComment). No other word is taken for the statement: a WITH list's +// names and aliases may be spelled like any keyword (`WITH 1 AS select`, +// `WITH desc AS (…)`, `WITH set -> 1 AS f`, `WITH t.from AS y`), but only the +// INSERT statement puts INTO after an insert. The exception, a read's +// `… AS insert INTO OUTFILE 'f'`, is classified as a write and answers `[]` +// uncached: harmless, and contrived. OUTFILE cannot tell the two apart, as +// `INSERT INTO outfile …` names a table. +func hasTopLevelInsertInto(s string) bool { depth := 0 i := 0 for i < len(s) { @@ -201,74 +208,44 @@ func containsMutationVerbAtTopLevel(s string) bool { depth-- } i++ - case c == '\'': - i++ - for i < len(s) { - if s[i] == '\'' { - if i+1 < len(s) && s[i+1] == '\'' { - i += 2 - continue - } - i++ - break - } - i++ - } - case c == '"' || c == '`': - q := c - i++ - for i < len(s) && s[i] != q { + case c == '\'' || c == '"' || c == '`': + i = skipQuoted(s, i) + case c == 0xE2: + // ‘…’ or “…”; any other character led by this byte is stepped + // over a byte at a time, like the default. + if j := skipCurlyQuoted(s, i); j > i { + i = j + } else { i++ } - if i < len(s) { - i++ + case c == '$': + // A heredoc, else a bareword led by `$` (never a keyword) or a + // lone `$`. + if j := skipHeredoc(s, i); j > i { + i = j + } else { + i = skipWord(s, i+1) } - case (c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z'): + case c == '.' && i+1 < len(s) && isDigit(s[i+1]): + // A number led by `.` ends with its digits, so a word glued to it + // is a word of its own: `.5INSERT` is `.5` then INSERT. After a + // name ClickHouse reads the `.` as a qualifier (`t.5insert`), which + // can only make a statement it rejects, or a read's INTO OUTFILE, + // look like a write. + i = skipDotNumber(s, i) + case isWordByte(c): + // A word led by a digit or `_` is read whole, so its tail is + // never taken for a keyword (`_insert`, `5insert`). start := i - for i < len(s) { - c2 := s[i] - if (c2 < 'A' || c2 > 'Z') && (c2 < 'a' || c2 > 'z') && (c2 < '0' || c2 > '9') && c2 != '_' { - break - } - i++ - } - if depth == 0 { - kw := strings.ToUpper(s[start:i]) - // Check non-mutation statement keywords (SELECT, SHOW, - // DESCRIBE, …) FIRST — these can legitimately be followed - // by `(` (e.g. `SELECT (1) FROM …`, `SELECT (a, b) FROM …` - // for tuple syntax), so we must not let the CTE-name - // lookahead below misclassify them as CTE aliases. - if _, ok := nonMutationVerbs[kw]; ok { - return false - } - // CTE name suppression: an identifier that ISN'T a - // non-mutation statement keyword and is followed by `AS` - // or `(` is a CTE definition name (with optional column - // list before AS). Skip without checking mutationVerbs - // — protects against CTE aliases that share a spelling - // with a mutation verb (`WITH set AS (...)`, - // `WITH alter AS (...)`, etc.). - if isCTENameLookahead(s, i) { - continue - } - if _, ok := mutationVerbs[kw]; ok { + i = skipWord(s, i) + if depth == 0 && strings.EqualFold(s[start:i], "INSERT") { + next := skipSpaceAndComments(s, i) + if strings.EqualFold(s[next:skipWord(s, next)], "INTO") { return true } } - case c == '-' && i+1 < len(s) && s[i+1] == '-', c == '#': - for i < len(s) && s[i] != '\n' { - i++ - } - case c == '/' && i+1 < len(s) && s[i+1] == '*': - i += 2 - for i+1 < len(s) { - if s[i] == '*' && s[i+1] == '/' { - i += 2 - break - } - i++ - } + case c == '-' && i+1 < len(s) && s[i+1] == '-', c == '#', c == '/' && i+1 < len(s) && (s[i+1] == '*' || s[i+1] == '/'): + i = skipComment(s, i) default: i++ } @@ -276,76 +253,183 @@ func containsMutationVerbAtTopLevel(s string) bool { return false } -// isCTENameLookahead returns true if the next non-whitespace, non-comment -// token at or after pos is `AS` (case-insensitive, word-boundary terminated) -// or `(` — signaling that whatever identifier just ended at pos is a CTE -// definition name (with optional column list before AS). Walks past space / -// tab / newline / `--` line comments / `#` line comments / `/* … */` block -// comments. Returns false on EOF or any other token. -func isCTENameLookahead(s string, pos int) bool { - i := pos - for i < len(s) { - c := s[i] - switch { - case c == ' ' || c == '\t' || c == '\r' || c == '\n': +// skipWord returns the index just past the bareword at s[i]: ClickHouse's +// barewords run over ASCII letters, digits, `_` and `$`. +func skipWord(s string, i int) int { + for i < len(s) && (isWordByte(s[i]) || s[i] == '$') { + i++ + } + return i +} + +func isWordByte(c byte) bool { + return (c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z') || isDigit(c) || c == '_' +} + +func isDigit(c byte) bool { return c >= '0' && c <= '9' } + +// skipDotNumber returns the index just past the number led by the `.` at s[i], +// read as ClickHouse's lexer reads one: digits, then an optional exponent (`e` +// or `E`, an optional sign, any digits), `_` allowed between two digits. +// Unlike a number led by a digit, it ends before any letters that follow it. +func skipDotNumber(s string, i int) int { + i = skipDigits(s, i+1) + if i < len(s) && (s[i] == 'e' || s[i] == 'E') { + i++ + if i < len(s) && (s[i] == '+' || s[i] == '-') { i++ - case c == '-' && i+1 < len(s) && s[i+1] == '-', c == '#': - for i < len(s) && s[i] != '\n' { - i++ - } - case c == '/' && i+1 < len(s) && s[i+1] == '*': - i += 2 - for i+1 < len(s) { - if s[i] == '*' && s[i+1] == '/' { - i += 2 - break - } - i++ - } - case c == '(': - return true - case (c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z'): - end := i - for end < len(s) { - c2 := s[end] - if (c2 < 'A' || c2 > 'Z') && (c2 < 'a' || c2 > 'z') && (c2 < '0' || c2 > '9') && c2 != '_' { - break - } - end++ - } - return strings.EqualFold(s[i:end], "AS") - default: - return false } + i = skipDigits(s, i) } - return false + return i +} + +// skipDigits returns the index just past the run of digits at s[i], in which +// each `_` stands between two digits. +func skipDigits(s string, i int) int { + for i < len(s) && (isDigit(s[i]) || s[i] == '_' && i > 0 && isDigit(s[i-1]) && i+1 < len(s) && isDigit(s[i+1])) { + i++ + } + return i +} + +// skipHeredoc returns the index just past the heredoc opening at s[i] — +// `$tag$ … $tag$`, the tag a possibly empty run of letters, digits and `_`, +// matched exactly — or i if none does, as an unclosed one is not a heredoc to +// ClickHouse either. +func skipHeredoc(s string, i int) int { + j := i + 1 + for j < len(s) && isWordByte(s[j]) { + j++ + } + if j >= len(s) || s[j] != '$' { + return i + } + tag := s[i : j+1] + if k := strings.Index(s[j+1:], tag); k >= 0 { + return j + 1 + k + len(tag) + } + return i } -// stripLeadingSQLComments trims whitespace plus line comments (`-- …` and -// MySQL-compat `# …`, both accepted by ClickHouse) and `/* block */` -// comments from the front of sql, returning the remainder with no leading -// whitespace. Unclosed block comments swallow the rest of the string — -// matches what ClickHouse itself would do at parse time. +// stripLeadingSQLComments trims whitespace and comments from the front of +// sql, the way ClickHouse's lexer skips them before the first token. func stripLeadingSQLComments(sql string) string { - s := strings.TrimLeft(sql, " \t\r\n") - for { - switch { - case strings.HasPrefix(s, "--"), strings.HasPrefix(s, "#"): - if i := strings.IndexByte(s, '\n'); i >= 0 { - s = strings.TrimLeft(s[i+1:], " \t\r\n") - } else { - return "" + return sql[skipSpaceAndComments(sql, 0):] +} + +// skipSpaceAndComments returns the index of the first byte at or after i that +// is neither whitespace nor inside a comment. +func skipSpaceAndComments(s string, i int) int { + for i < len(s) { + if n := sqlSpaceLen(s, i); n > 0 { + i += n + continue + } + j := skipComment(s, i) + if j == i { + return i + } + i = j + } + return i +} + +// sqlSpaceLen is the byte length of the whitespace character at s[i], or 0. +// The set is ClickHouse's lexer's: ASCII space, \t \n \v \f \r, and the +// Unicode spaces it skips so that SQL pasted from a word processor parses. A +// leading one the classifier did not skip would hide the verb behind it. +func sqlSpaceLen(s string, i int) int { + switch s[i] { + case ' ', '\t', '\n', '\v', '\f', '\r': + return 1 + } + if s[i] < utf8.RuneSelf { + return 0 + } + r, n := utf8.DecodeRuneInString(s[i:]) + switch { + case r == 0x85, r == 0xA0, r == 0x180E, r >= 0x2000 && r <= 0x200D, + r == 0x2028, r == 0x2029, r == 0x202F, r == 0x205F, r == 0x2060, + r == 0x3000, r == 0xFEFF: + return n + } + return 0 +} + +// skipComment returns the index just past the comment starting at s[i], or i +// if none starts there: `--`, `//` and MySQL-compat `#` to end of line, +// `/* … */` nesting as ClickHouse's do. An unclosed block comment runs to the +// end, as it does for ClickHouse, which then rejects the statement. +func skipComment(s string, i int) int { + switch { + case strings.HasPrefix(s[i:], "--"), strings.HasPrefix(s[i:], "//"), s[i] == '#': + if j := strings.IndexByte(s[i:], '\n'); j >= 0 { + return i + j + 1 + } + return len(s) + case strings.HasPrefix(s[i:], "/*"): + depth := 0 + for j := i; j+1 < len(s); { + switch { + case s[j] == '/' && s[j+1] == '*': + depth++ + j += 2 + case s[j] == '*' && s[j+1] == '/': + depth-- + j += 2 + if depth == 0 { + return j + } + default: + j++ } - case strings.HasPrefix(s, "/*"): - if i := strings.Index(s[2:], "*/"); i >= 0 { - s = strings.TrimLeft(s[2+i+2:], " \t\r\n") - } else { - return "" + } + return len(s) + } + return i +} + +// skipCurlyQuoted returns the index just past a string literal in ‘…’ or a +// quoted identifier in “…”, which ClickHouse reads so that SQL pasted from a +// word processor parses, or i if none opens at s[i]. Nothing escapes inside +// them; an unclosed one runs to the end. +func skipCurlyQuoted(s string, i int) int { + var closer string + switch { + case strings.HasPrefix(s[i:], "\u2018"): + closer = "\u2019" + case strings.HasPrefix(s[i:], "\u201c"): + closer = "\u201d" + default: + return i + } + start := i + len(closer) // the opener is as long as its closer + if k := strings.Index(s[start:], closer); k >= 0 { + return start + k + len(closer) + } + return len(s) +} + +// skipQuoted returns the index just past the string literal or quoted +// identifier opening at s[i] (`'`, `"` or backtick). As in ClickHouse's lexer, +// a doubled quote or a backslash escapes the next byte; an unclosed one runs +// to the end. +func skipQuoted(s string, i int) int { + q := s[i] + for i++; i < len(s); i++ { + switch s[i] { + case '\\': + i++ + case q: + if i+1 < len(s) && s[i+1] == q { + i++ + continue } - default: - return s + return i + 1 } } + return len(s) } // transformRow converts ClickHouse-specific types to JSON-friendly values. diff --git a/internal/api/clickhouse_exec_test.go b/internal/api/clickhouse_exec_test.go index 947413b70..1ead574f3 100644 --- a/internal/api/clickhouse_exec_test.go +++ b/internal/api/clickhouse_exec_test.go @@ -7,6 +7,7 @@ import ( "time" "github.com/ClickHouse/clickhouse-go/v2/lib/driver" + "github.com/Wave-RF/WaveHouse/internal/testutil/mutationtest" "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -39,89 +40,38 @@ func (c *stubConn) Query(_ context.Context, _ string, _ ...any) (driver.Rows, er return &chainEmptyRows{}, nil } +// TestIsMutation runs the shared cases; the integration suite checks the +// same cases against ClickHouse's parser. func TestIsMutation(t *testing.T) { t.Parallel() - tests := []struct { - name string - sql string - want bool - }{ - {"select", "SELECT 1", false}, - {"select lower", "select 1", false}, - {"with cte", "WITH x AS (SELECT 1) SELECT * FROM x", false}, - {"show", "SHOW TABLES", false}, - {"describe", "DESCRIBE clicks", false}, - {"explain", "EXPLAIN SELECT 1", false}, - {"exists", "EXISTS TABLE clicks", false}, - - {"insert", "INSERT INTO t VALUES (1)", true}, - {"update", "UPDATE t SET a=1 WHERE b=2", true}, - {"delete", "DELETE FROM t WHERE id=1", true}, - {"truncate", "TRUNCATE TABLE t", true}, - {"truncate lower", "truncate table t", true}, - {"drop", "DROP TABLE t", true}, - {"alter", "ALTER TABLE t ADD COLUMN c String", true}, - {"create", "CREATE TABLE t (a Int)", true}, - {"rename", "RENAME TABLE a TO b", true}, - {"exchange", "EXCHANGE TABLES t1 AND t2", true}, - {"optimize", "OPTIMIZE TABLE t", true}, - {"replace", "REPLACE INTO t VALUES (1)", true}, - {"grant", "GRANT SELECT ON t TO u", true}, - {"revoke", "REVOKE SELECT ON t FROM u", true}, - {"system", "SYSTEM RELOAD CONFIG", true}, - {"attach", "ATTACH TABLE t FROM '/path'", true}, - {"detach", "DETACH TABLE t", true}, - {"kill", "KILL QUERY WHERE query_id = 'abc'", true}, - {"set", "SET max_threads = 4", true}, - {"use", "USE mydb", true}, - - {"leading whitespace", " \n\tTRUNCATE TABLE t", true}, - {"line comment then mutation", "-- drop guard\nDROP TABLE t", true}, - {"hash line comment then mutation", "# audit\nDROP TABLE t", true}, - {"block comment then mutation", "/* admin */ ALTER TABLE t ADD COLUMN c Int", true}, - {"mixed comments then select", "-- foo\n# bar\n/* baz */ SELECT 1", false}, - {"with insert", "WITH cte AS (SELECT 1) INSERT INTO t SELECT * FROM cte", true}, - {"with insert lower", "with cte as (select 1) insert into t select * from cte", true}, - {"with delete", "WITH cte AS (SELECT id FROM x) DELETE FROM t WHERE id IN (SELECT id FROM cte)", true}, - {"with update", "WITH cte AS (SELECT 1) ALTER TABLE t UPDATE a=1 WHERE id IN (SELECT id FROM cte)", true}, - {"with truncate", "WITH cte AS (SELECT 1) TRUNCATE TABLE t", true}, - {"with multi-cte insert", "WITH a AS (SELECT 1), b AS (SELECT 2) INSERT INTO t SELECT * FROM a JOIN b", true}, - {"with nested parens insert", "WITH cte AS (SELECT id FROM t WHERE id IN (1,2,3)) INSERT INTO t2 SELECT * FROM cte", true}, - {"with paren-in-string insert", "WITH cte AS (SELECT ')' AS x) INSERT INTO t2 SELECT * FROM cte", true}, - {"with materialized insert", "WITH cte AS MATERIALIZED (SELECT 1) INSERT INTO t SELECT * FROM cte", true}, - {"with recursive select", "WITH RECURSIVE x AS (SELECT 1 UNION ALL SELECT * FROM x) SELECT * FROM x", false}, - {"with nested select", "WITH x AS (SELECT 1) SELECT * FROM (SELECT * FROM x)", false}, - {"with scalar insert", "WITH '/path' AS p INSERT INTO files VALUES (p)", true}, - {"with line comment containing DELETE then select", "WITH cte AS (SELECT 1) -- old DELETE approach\nSELECT * FROM cte", false}, - {"with hash comment containing TRUNCATE then select", "WITH cte AS (SELECT 1) # was TRUNCATE\nSELECT * FROM cte", false}, - {"with block comment containing INSERT then select", "WITH cte AS (SELECT 1) /* INSERT reminder */ SELECT * FROM cte", false}, - {"with comment then real mutation", "WITH cte AS (SELECT 1) -- explanatory\nINSERT INTO t SELECT * FROM cte", true}, - {"with unclosed block comment", "WITH cte AS (SELECT 1) /* unterminated comment DELETE", false}, - {"with select from system tables (collision regression)", "WITH x AS (SELECT 1) SELECT * FROM system.tables", false}, - {"with select from system columns lower (collision regression)", "with x as (select 1) select name from system.columns", false}, - {"with select aliased as set (false positive regression)", "WITH cte AS (SELECT 1) SELECT * FROM cte AS set", false}, - {"with select from system tables then real insert", "WITH x AS (SELECT * FROM system.tables) INSERT INTO snapshot SELECT * FROM x", true}, - {"with CTE alias named set (read)", "WITH set AS (SELECT 1) SELECT * FROM set", false}, - {"with CTE alias named alter (read)", "WITH alter AS (SELECT 1) SELECT id FROM alter", false}, - {"with CTE alias named drop lowercase (read)", "with drop as (select 1) select * from drop", false}, - {"with CTE alias named update then real update", "WITH update AS (SELECT id FROM x) ALTER TABLE other UPDATE c=1 WHERE id IN (SELECT id FROM update)", true}, - {"with CTE name with column list (read)", "WITH cte (a, b) AS (SELECT 1, 2) SELECT * FROM cte", false}, - {"with multi-CTE both with verb-name aliases (read)", "WITH set AS (SELECT 1), kill AS (SELECT 2) SELECT * FROM set JOIN kill", false}, - {"with parenthesized SELECT then system table (CTE-lookahead ordering regression)", "WITH x AS (SELECT 1) SELECT (1) FROM system.tables", false}, - {"with tuple-shape SELECT then system table", "WITH x AS (SELECT 1) SELECT (a, b) FROM system.parts", false}, - - {"empty", "", false}, - {"comment only", "-- just a comment", false}, - {"unclosed block comment", "/* never closed", false}, - } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { + for _, tc := range mutationtest.Cases { + t.Run(tc.Name, func(t *testing.T) { t.Parallel() - assert.Equal(t, tt.want, isMutation(tt.sql)) + assert.Equal(t, tc.Mutation, IsMutation(tc.SQL)) }) } } +// TestIsMutation_ClickHouseWhitespace pins every character ClickHouse 26.6's +// lexer accepts as whitespace, each checked against a live server: ahead of a +// write it must not hide the verb, and ahead of a read it must not make one. +func TestIsMutation_ClickHouseWhitespace(t *testing.T) { + t.Parallel() + spaces := []rune{' ', '\t', '\n', '\v', '\f', '\r', 0x85, 0xA0, 0x180E, 0x2028, 0x2029, 0x202F, 0x205F, 0x2060, 0x3000, 0xFEFF} + for r := rune(0x2000); r <= 0x200D; r++ { + spaces = append(spaces, r) + } + for _, r := range spaces { + ws := string(r) + assert.True(t, IsMutation(ws+"INSERT INTO t VALUES (1)"), "U+%04X before INSERT", r) + assert.False(t, IsMutation(ws+"SELECT 1"), "U+%04X before SELECT", r) + assert.True(t, IsMutation("WITH x AS (SELECT 1)"+ws+"INSERT INTO t SELECT * FROM x"), "U+%04X before a WITH's INSERT", r) + assert.True(t, IsMutation("WITH 1 AS x INSERT"+ws+"INTO t SELECT x"), "U+%04X between a WITH's INSERT and INTO", r) + } + // Not whitespace to ClickHouse (it rejects the statement), so not skipped. + assert.False(t, IsMutation("\u1680INSERT INTO t VALUES (1)")) +} + func TestExecuteCHQuery_MutationRoutesToExec(t *testing.T) { t.Parallel() // Mutations route through driver.Exec because clickhouse-go's @@ -137,6 +87,7 @@ func TestExecuteCHQuery_MutationRoutesToExec(t *testing.T) { "ALTER TABLE clicks ADD COLUMN c String", "INSERT INTO clicks VALUES (1)", " -- audit log\n UPDATE clicks SET v = 1 WHERE id = 2", + "EXECUTE AS writer INSERT INTO clicks VALUES (1)", } { t.Run(sql, func(t *testing.T) { t.Parallel() @@ -152,12 +103,17 @@ func TestExecuteCHQuery_MutationRoutesToExec(t *testing.T) { func TestExecuteCHQuery_SelectRoutesToQuery(t *testing.T) { t.Parallel() - conn := &stubConn{} - rows, err := executeCHQuery(context.Background(), conn, "SELECT 1", nil) - require.NoError(t, err) - assert.Zero(t, conn.execCount, "Exec must not be used for SELECT") - assert.Equal(t, 1, conn.queryCount, "Query must be used for SELECT") - assert.Equal(t, []map[string]any{}, rows, "zero-row SELECT must marshal to [] not null") + for _, sql := range []string{"SELECT 1", "EXECUTE AS reader SELECT 1"} { + t.Run(sql, func(t *testing.T) { + t.Parallel() + conn := &stubConn{} + rows, err := executeCHQuery(context.Background(), conn, sql, nil) + require.NoError(t, err) + assert.Zero(t, conn.execCount, "Exec must not be used for SELECT") + assert.Equal(t, 1, conn.queryCount, "Query must be used for SELECT") + assert.Equal(t, []map[string]any{}, rows, "zero-row SELECT must marshal to [] not null") + }) + } } // TestExecuteCHQuery_TransformsClickHouseTypes pins transformRow's contract diff --git a/internal/api/pipes.go b/internal/api/pipes.go index 428d3b8d0..3a850d3d5 100644 --- a/internal/api/pipes.go +++ b/internal/api/pipes.go @@ -152,6 +152,11 @@ func (h *PipesHandler) Execute(w http.ResponseWriter, r *http.Request) { return } + if IsMutation(sql) { + h.executeWrite(w, r, store, sql, params) + return + } + // Cache. A pipe can read several tables, but the current pipe impl doesn't // expose its table/scope dependencies, so we pass no deps: the result folds // the tenant's version alone, so InvalidateTenant orphans it but no insert @@ -185,28 +190,12 @@ func (h *PipesHandler) Execute(w http.ResponseWriter, r *http.Request) { // Execute with singleflight. v, err, _ := h.sf.Do(cacheKey, func() (interface{}, error) { - queryCtx, cancel := context.WithTimeout(r.Context(), timeoutOf(h.queryTimeout, store)) - defer cancel() - - start := time.Now() - - rows, err := executeCHQuery(queryCtx, conn, sql, params) - queryDuration := time.Since(start) + data, queryDuration, err := h.run(r.Context(), store, conn, sql, params) if err != nil { - // TODO: depending on the error, we may actually want to cache it return nil, err } - - data, err := json.Marshal(rows) - if err != nil { - // TODO: eventually we want CSV support etc - return nil, err - } - - ttl := cache.QueryTimeToTTL(queryDuration) - if h.Cache != nil { - _ = h.Cache.Set(r.Context(), snap, data, ttl) + _ = h.Cache.Set(r.Context(), snap, data, cache.QueryTimeToTTL(queryDuration)) } return data, nil }) @@ -219,3 +208,41 @@ func (h *PipesHandler) Execute(w http.ResponseWriter, r *http.Request) { w.Header().Set("X-Cache", "MISS") _, _ = w.Write(v.([]byte)) //nolint:gosec // G705: the tenant id on the key only selects the entry; the bytes are JSON the handler marshalled from ClickHouse rows } + +// executeWrite runs a pipe that writes, on every call: a cached or coalesced +// response would answer a repeat without executing it, silently dropping the +// write (#386) — on every instance once the cache is shared. IsMutation is the +// classifier executeCHQuery routes Exec by, so what bypasses here is exactly +// what runs as a write. no-store keeps an HTTP cache in front of a GET from +// answering a repeat the same way. +func (h *PipesHandler) executeWrite(w http.ResponseWriter, r *http.Request, store *settings.Store, sql string, params []any) { + conn := connOf(h.CHConn, store) + if conn == nil { + writeUnavailable(w, noConnectionMessage, retryAfterPool) + return + } + data, _, err := h.run(r.Context(), store, conn, sql, params) + if err != nil { + writeCHWriteError(w, r, err, err.Error()) + return + } + w.Header().Set("Content-Type", "application/json") + w.Header().Set("X-Cache", "BYPASS") + w.Header().Set("Cache-Control", "no-store") + _, _ = w.Write(data) //nolint:gosec // G705: JSON the handler marshalled from the exec result +} + +// run executes a pipe's bound SQL under the tenant's query timeout and +// returns the rows as JSON with how long ClickHouse took. +func (h *PipesHandler) run(ctx context.Context, store *settings.Store, conn driver.Conn, sql string, params []any) ([]byte, time.Duration, error) { + queryCtx, cancel := context.WithTimeout(ctx, timeoutOf(h.queryTimeout, store)) + defer cancel() + start := time.Now() + rows, err := executeCHQuery(queryCtx, conn, sql, params) + queryDuration := time.Since(start) + if err != nil { + return nil, 0, err + } + data, err := json.Marshal(rows) + return data, queryDuration, err +} diff --git a/internal/api/pipes_test.go b/internal/api/pipes_test.go index db7ff2c41..f3a163d44 100644 --- a/internal/api/pipes_test.go +++ b/internal/api/pipes_test.go @@ -7,10 +7,15 @@ import ( "net/http" "net/http/httptest" "strings" + "sync" + "sync/atomic" "testing" + "testing/synctest" "time" + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" "github.com/Wave-RF/WaveHouse/internal/auth" + "github.com/Wave-RF/WaveHouse/internal/cache" "github.com/Wave-RF/WaveHouse/internal/pipes" "github.com/Wave-RF/WaveHouse/internal/policy" "github.com/Wave-RF/WaveHouse/internal/settings" @@ -511,3 +516,120 @@ func TestPipesHandler_Execute_NoAllowedRoles_AdminAllowed(t *testing.T) { "admin bypasses the allowlist on a pipe with no allowed_roles") assert.NotEqual(t, http.StatusNotFound, w.Code) } + +// writeConn counts Exec and Query calls, and every Exec returns err. With +// gate set, every Exec reports itself on entered and holds until gate is +// closed, so a test can hold requests in flight together. +type writeConn struct { + driver.Conn + execs, queries atomic.Int32 + entered, gate chan struct{} + err error +} + +func (c *writeConn) Exec(context.Context, string, ...any) error { + c.execs.Add(1) + if c.gate != nil { + c.entered <- struct{}{} + <-c.gate + } + return c.err +} + +func (c *writeConn) Query(context.Context, string, ...any) (driver.Rows, error) { + c.queries.Add(1) + return &chainEmptyRows{}, nil +} + +// pipeCallAs runs the pipe name as the writer role and returns the recorder. +func pipeCallAs(t *testing.T, h *PipesHandler, name string) *httptest.ResponseRecorder { + t.Helper() + w := httptest.NewRecorder() + r := pipesRequest(t, http.MethodPost, "/v1/pipes/"+name, name, map[string]any{"msg": "hello"}) + h.Execute(w, withTenant(r.WithContext(auth.WithRole(r.Context(), "writer")))) + return w +} + +func writerPipesHandler(t *testing.T, conn driver.Conn, c cache.Cache, queries ...*pipes.NamedQuery) *PipesHandler { + t.Helper() + for _, q := range queries { + q.AllowedRoles = []string{"writer"} + } + timeout := func(*settings.Store) time.Duration { return 5 * time.Second } + return NewPipesHandler(staticPipes(queries...), staticPolicy(&policy.Policy{}), fixedConn(conn), c, timeout) +} + +// #386: a pipe that writes executes on every call. Served from the cache, a +// repeat would answer 200 with the first call's `[]` and never reach +// ClickHouse — the write silently dropped. +func TestPipesHandler_Execute_MutationRunsEveryCall(t *testing.T) { + t.Parallel() + for name, sql := range map[string]string{ + "insert": "INSERT INTO audit_log VALUES ({{msg}}, now())", + "insert after cte": "WITH m AS (SELECT {{msg}} AS msg) INSERT INTO audit_log SELECT msg, now() FROM m", + "alter delete": "ALTER TABLE audit_log DELETE WHERE msg = {{msg}}", + } { + t.Run(name, func(t *testing.T) { + t.Parallel() + l1, err := cache.NewLocal(1 << 20) + require.NoError(t, err) + t.Cleanup(func() { _ = l1.Close() }) + conn := &writeConn{} + h := writerPipesHandler(t, conn, l1, &pipes.NamedQuery{Name: "log", SQL: sql}) + + for range 3 { + w := pipeCallAs(t, h, "log") + require.Equal(t, http.StatusOK, w.Code, "body: %s", w.Body.String()) + assert.Equal(t, "BYPASS", w.Header().Get("X-Cache")) + assert.Equal(t, "no-store", w.Header().Get("Cache-Control")) + assert.JSONEq(t, `[]`, w.Body.String()) + l1.Wait() + } + assert.Equal(t, int32(3), conn.execs.Load(), "every call must reach ClickHouse") + assert.Zero(t, conn.queries.Load()) + }) + } +} + +// Identical mutation calls in flight together are each executed: coalescing +// them would run one write for all of them. Under synctest, Wait returns once +// every request is inside Exec or parked on another's flight. +func TestPipesHandler_Execute_ConcurrentMutationsNotCoalesced(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + const calls = 3 + conn := &writeConn{entered: make(chan struct{}, calls), gate: make(chan struct{})} + h := writerPipesHandler(t, conn, nil, &pipes.NamedQuery{Name: "log", SQL: "INSERT INTO audit_log VALUES ({{msg}}, now())"}) + var wg sync.WaitGroup + for range calls { + wg.Go(func() { + w := pipeCallAs(t, h, "log") + assert.Equal(t, http.StatusOK, w.Code, "body: %s", w.Body.String()) + }) + } + synctest.Wait() + assert.Len(t, conn.entered, calls, "writes in flight once every request is blocked") + close(conn.gate) + wg.Wait() + assert.Equal(t, int32(calls), conn.execs.Load()) + }) +} + +// A read pipe keeps its cache, including one whose table name starts with a +// write verb: the classifier reads the statement, not the words in it. +func TestPipesHandler_Execute_ReadPipeStaysCached(t *testing.T) { + t.Parallel() + l1, err := cache.NewLocal(1 << 20) + require.NoError(t, err) + t.Cleanup(func() { _ = l1.Close() }) + conn := &writeConn{} + h := writerPipesHandler(t, conn, l1, &pipes.NamedQuery{Name: "recent", SQL: "SELECT * FROM insert_log WHERE msg = {{msg}}"}) + + for _, want := range []string{"MISS", "HIT", "HIT"} { + w := pipeCallAs(t, h, "recent") + require.Equal(t, http.StatusOK, w.Code, "body: %s", w.Body.String()) + assert.Equal(t, want, w.Header().Get("X-Cache")) + l1.Wait() + } + assert.Equal(t, int32(1), conn.queries.Load()) + assert.Zero(t, conn.execs.Load()) +} diff --git a/internal/api/query.go b/internal/api/query.go index 0959bf670..fa782b0c8 100644 --- a/internal/api/query.go +++ b/internal/api/query.go @@ -37,7 +37,7 @@ import ( // the upstream ClickHouse has multi-query enabled, which is the // default in recent versions; older or restrictively-configured // servers may reject the second statement with a clear error. -// - There is no isMutation heuristic to maintain — no leading-verb table, +// - There is no IsMutation heuristic to maintain — no leading-verb table, // no comment stripper, no CTE-aware paren scanner, no class of bug // where a future ClickHouse verb routes the wrong way. // - ClickHouse's own error messages reach the admin verbatim, which is diff --git a/internal/api/tenant_clickhouse_test.go b/internal/api/tenant_clickhouse_test.go index 8181baa2d..781b3e4a3 100644 --- a/internal/api/tenant_clickhouse_test.go +++ b/internal/api/tenant_clickhouse_test.go @@ -135,6 +135,15 @@ func TestClickHouseRoutes_NoPoolIs503(t *testing.T) { h.Execute(w, withTenant(pipesRequest(t, http.MethodGet, "/v1/pipes/top_pages", "top_pages", nil))) assertUnavailable(t, w, noConnectionMessage, retryAfterPool) }) + // A write that never reached ClickHouse cannot have run, so this 503 + // keeps its Retry-After where a failed write's answer drops it. + t.Run("write pipe execute", func(t *testing.T) { + t.Parallel() + h := NewPipesHandler(staticPipes(&pipes.NamedQuery{Name: "log", SQL: "INSERT INTO audit_log VALUES (1)", AllowedRoles: []string{"viewer"}}), allowAll, noConn, nil, noTimeout) + w := httptest.NewRecorder() + h.Execute(w, withTenant(pipesRequest(t, http.MethodGet, "/v1/pipes/log", "log", nil))) + assertUnavailable(t, w, noConnectionMessage, retryAfterPool) + }) t.Run("raw-SQL proxy", func(t *testing.T) { t.Parallel() h := newTestQueryHandler(func(*settings.Store) chconn.Target { return chconn.Target{} }, noTimeout) diff --git a/internal/app/wire.go b/internal/app/wire.go index efcb44fc9..c6f7c7f97 100644 --- a/internal/app/wire.go +++ b/internal/app/wire.go @@ -367,8 +367,8 @@ func (a *App) registryFor(s *settings.Store) *discovery.SchemaRegistry { return a.discoveries.For(s.Tenant()) } -// queryTimeout is the tenant's read deadline, a per-call setting rather -// than a property of the pool it shares. +// queryTimeout is the tenant's deadline for a call on the query paths, a +// per-call setting rather than a property of the pool it shares. func queryTimeout(s *settings.Store) time.Duration { return s.ClickHouse().QueryTimeout } // wireDiscovery builds one schema registry per served tenant, each with a diff --git a/internal/pipes/pipes.go b/internal/pipes/pipes.go index e146a2012..28e8b44db 100644 --- a/internal/pipes/pipes.go +++ b/internal/pipes/pipes.go @@ -143,7 +143,9 @@ func BindParams(q *NamedQuery, supplied map[string]any) (string, []any, error) { // comma-separated list of recursively formatted elements — the `(v1, v2, …)` // shape ClickHouse expects on the right of `IN`, matching how the // structured-query builder renders an IN clause. Because every scalar leaf is -// escaped, no value — or array element — can break out of its literal. +// escaped, no value — or array element — can break out of its literal, as long +// as the template writes the placeholder bare: inside quotes (`'{{id}}'`) the +// value's own quotes close the template's. // // Values with no scalar SQL representation are refused rather than emitted as // Go's `%v` text: a JSON object has no meaning here, and an empty array would diff --git a/internal/settings/settings.go b/internal/settings/settings.go index c1e439238..9ad25b2e7 100644 --- a/internal/settings/settings.go +++ b/internal/settings/settings.go @@ -98,7 +98,8 @@ type ClickHouseConfig struct { HTTPScheme *string `json:"http_scheme"` Database *string `json:"database"` Username *string `json:"username"` - // QueryTimeout is the read deadline in seconds (>= 1). + // QueryTimeout is the deadline in seconds (>= 1) of a call on the query + // paths: structured queries, pipes (writes included) and the raw-SQL proxy. QueryTimeout *int `json:"query_timeout"` // TLS is the TLS wiring of both hops: `enabled` switches the native // protocol to TLS, `http_scheme` stays the HTTP hop's switch, and the diff --git a/internal/testutil/mutationtest/cases.go b/internal/testutil/mutationtest/cases.go new file mode 100644 index 000000000..0a94be19d --- /dev/null +++ b/internal/testutil/mutationtest/cases.go @@ -0,0 +1,208 @@ +// Package mutationtest holds the statements api.IsMutation is tested on, +// shared by its unit test and by the integration test that checks each one +// against ClickHouse's own parser. Add a case here, not to either test. +package mutationtest + +// Case is a statement and how api.IsMutation classifies it. +type Case struct { + Name string + SQL string + // Mutation is true for a statement that goes through Exec. + Mutation bool + // Unparsed marks a statement ClickHouse rejects as a syntax error, so its + // parser has no answer to check Mutation against; the integration test + // fails if one starts to parse. + Unparsed bool +} + +// Cases is every statement api.IsMutation is tested on. +var Cases = []Case{ + {"select", "SELECT 1", false, false}, + {"select lower", "select 1", false, false}, + {"with cte", "WITH x AS (SELECT 1) SELECT * FROM x", false, false}, + {"show", "SHOW TABLES", false, false}, + {"describe", "DESCRIBE clicks", false, false}, + {"explain", "EXPLAIN SELECT 1", false, false}, + {"exists", "EXISTS TABLE clicks", false, false}, + + {"insert", "INSERT INTO t VALUES (1)", true, false}, + {"update", "UPDATE t SET a=1 WHERE b=2", true, false}, + {"delete", "DELETE FROM t WHERE id=1", true, false}, + {"truncate", "TRUNCATE TABLE t", true, false}, + {"truncate lower", "truncate table t", true, false}, + {"drop", "DROP TABLE t", true, false}, + {"alter", "ALTER TABLE t ADD COLUMN c String", true, false}, + {"create", "CREATE TABLE t (a Int)", true, false}, + {"rename", "RENAME TABLE a TO b", true, false}, + {"exchange", "EXCHANGE TABLES t1 AND t2", true, false}, + {"optimize", "OPTIMIZE TABLE t", true, false}, + {"replace into", "REPLACE INTO t VALUES (1)", true, true}, + {"replace table", "REPLACE TABLE t (a Int) ENGINE = Memory", true, false}, + {"grant", "GRANT SELECT ON t TO u", true, false}, + {"revoke", "REVOKE SELECT ON t FROM u", true, false}, + {"system", "SYSTEM RELOAD CONFIG", true, false}, + {"attach", "ATTACH TABLE t FROM '/path'", true, false}, + {"detach", "DETACH TABLE t", true, false}, + {"kill", "KILL QUERY WHERE query_id = 'abc'", true, false}, + {"set", "SET max_threads = 4", true, false}, + {"use", "USE mydb", true, false}, + + {"leading whitespace", " \n\tTRUNCATE TABLE t", true, false}, + {"line comment then mutation", "-- drop guard\nDROP TABLE t", true, false}, + {"hash line comment then mutation", "# audit\nDROP TABLE t", true, false}, + {"block comment then mutation", "/* admin */ ALTER TABLE t ADD COLUMN c Int", true, false}, + {"mixed comments then select", "-- foo\n# bar\n/* baz */ SELECT 1", false, false}, + {"with insert", "WITH cte AS (SELECT 1) INSERT INTO t SELECT * FROM cte", true, false}, + {"with insert lower", "with cte as (select 1) insert into t select * from cte", true, false}, + {"with multi-cte insert", "WITH a AS (SELECT 1), b AS (SELECT 2) INSERT INTO t SELECT * FROM a CROSS JOIN b", true, false}, + {"with nested parens insert", "WITH cte AS (SELECT id FROM t WHERE id IN (1,2,3)) INSERT INTO t2 SELECT * FROM cte", true, false}, + {"with paren-in-string insert", "WITH cte AS (SELECT ')' AS x) INSERT INTO t2 SELECT * FROM cte", true, false}, + {"with materialized insert", "WITH cte AS MATERIALIZED (SELECT 1) INSERT INTO t SELECT * FROM cte", true, false}, + {"with recursive select", "WITH RECURSIVE x AS (SELECT 1 UNION ALL SELECT * FROM x) SELECT * FROM x", false, false}, + {"with nested select", "WITH x AS (SELECT 1) SELECT * FROM (SELECT * FROM x)", false, false}, + {"with scalar insert", "WITH '/path' AS p INSERT INTO files VALUES (p)", true, false}, + {"with line comment containing DELETE then select", "WITH cte AS (SELECT 1) -- old DELETE approach\nSELECT * FROM cte", false, false}, + {"with hash comment containing TRUNCATE then select", "WITH cte AS (SELECT 1) # was TRUNCATE\nSELECT * FROM cte", false, false}, + {"with block comment containing INSERT then select", "WITH cte AS (SELECT 1) /* INSERT reminder */ SELECT * FROM cte", false, false}, + {"with comment then real mutation", "WITH cte AS (SELECT 1) -- explanatory\nINSERT INTO t SELECT * FROM cte", true, false}, + {"with unclosed block comment", "WITH cte AS (SELECT 1) /* unterminated comment DELETE", false, true}, + {"with select from system tables (collision regression)", "WITH x AS (SELECT 1) SELECT * FROM system.tables", false, false}, + {"with select from system columns lower (collision regression)", "with x as (select 1) select name from system.columns", false, false}, + {"with select aliased as set (false positive regression)", "WITH cte AS (SELECT 1) SELECT * FROM cte AS set", false, false}, + {"with select from system tables then real insert", "WITH x AS (SELECT * FROM system.tables) INSERT INTO snapshot SELECT * FROM x", true, false}, + {"with CTE alias named set (read)", "WITH set AS (SELECT 1) SELECT * FROM set", false, false}, + {"with CTE alias named alter (read)", "WITH alter AS (SELECT 1) SELECT id FROM alter", false, false}, + {"with CTE alias named drop lowercase (read)", "with drop as (select 1) select * from drop", false, false}, + {"with CTE name with column list (read)", "WITH cte (a, b) AS (SELECT 1, 2) SELECT * FROM cte", false, false}, + {"with multi-CTE both with verb-name aliases (read)", "WITH set AS (SELECT 1), kill AS (SELECT 2) SELECT * FROM set CROSS JOIN kill", false, false}, + {"with parenthesized SELECT then system table (CTE-lookahead ordering regression)", "WITH x AS (SELECT 1) SELECT (1) FROM system.tables", false, false}, + {"with tuple-shape SELECT then system table", "WITH x AS (SELECT 1) SELECT (a, b) FROM system.parts", false, false}, + + // A backslash escapes the next byte inside all three quote kinds, so an + // escaped quote does not end the literal or identifier. + {"with backslash-escaped quote in literal then insert", `WITH m AS (SELECT 'it\'s' AS s) INSERT INTO t SELECT s FROM m`, true, false}, + {"with backslash-escaped quote in literal then select", `WITH m AS (SELECT 'a\'b' AS s) SELECT 'x) INSERT' FROM m`, false, false}, + {"with backslash-escaped double quote then insert", `WITH m AS (SELECT 'x' AS "a\"(b") INSERT INTO t SELECT * FROM m`, true, false}, + {"with backslash-escaped double quote then select", `WITH m AS (SELECT 1 AS "a\"b") SELECT 2 AS "x) INSERT" FROM m`, false, false}, + {"with backslash-escaped backtick then insert", "WITH m AS (SELECT 'x' AS `a\\`(b`) INSERT INTO t SELECT * FROM m", true, false}, + {"with backslash-escaped backtick then select", "WITH m AS (SELECT 1 AS `a\\`b`) SELECT 2 AS `x) INSERT` FROM m", false, false}, + + // A heredoc ($$…$$, $tag$…$tag$) is a literal: its parens, quotes + // and words are not the statement's. + {"with heredoc holding a paren then insert", "WITH $$ ( $$ AS s INSERT INTO t SELECT s", true, false}, + {"with tagged heredoc holding a quote then insert", "WITH $x$ it's $x$ AS s INSERT INTO t SELECT s", true, false}, + {"with tagged heredoc holding a paren and another tag then insert", "WITH $x$ ( $y$ $x$ AS s INSERT INTO t SELECT s", true, false}, + {"with heredoc holding a verb then select", "WITH $$INSERT$$ AS s SELECT s", false, false}, + {"with tagged heredoc holding a paren and a verb then select", "WITH $x$ ) INSERT $x$ AS s SELECT s", false, false}, + {"with CTE alias set$ (read)", "WITH set$ AS (SELECT 1 AS v) SELECT * FROM set$", false, false}, + + // A word led by `_` is one bareword, never a keyword's tail, and a + // leading bareword is matched whole, never by its first letters. + {"leading bareword insert_log", "insert_log VALUES (1)", false, true}, + {"leading bareword insert2", "insert2 INTO t VALUES (1)", false, true}, + {"with alias _delete (read)", "WITH 1 AS _delete SELECT _delete", false, false}, + {"with alias _set (read)", "WITH [1,2] AS _set SELECT has(_set, 1)", false, false}, + + // `//` starts a line comment. + {"slash comment then insert", "// note\nINSERT INTO t VALUES (1)", true, false}, + {"slash comment hiding insert then select", "// INSERT\nSELECT 1", false, false}, + {"with slash comment holding a paren then insert", "WITH x AS (SELECT 'a' AS s) // (\nINSERT INTO t SELECT * FROM x", true, false}, + + // ‘…’ is a string literal and “…” a quoted identifier; nothing escapes + // inside them. + {"with curly-quoted literal holding a paren then insert", "WITH x AS (SELECT \u2018(\u2019 AS s) INSERT INTO t SELECT s FROM x", true, false}, + {"with curly-quoted literal holding a verb then select", "WITH x AS (SELECT \u2018) INSERT\u2019 AS s) SELECT s FROM x", false, false}, + {"with curly-quoted identifier holding a paren then insert", "WITH x AS (SELECT 'q' AS \u201cc(d\u201d) INSERT INTO t SELECT * FROM x", true, false}, + + // ClickHouse's lexer skips \v, \f and Unicode spaces as whitespace + // (TestIsMutation_ClickHouseWhitespace covers the whole set). + {"leading form feed then insert", "\fINSERT INTO t VALUES (1)", true, false}, + {"leading vertical tab then insert", "\vINSERT INTO t VALUES (1)", true, false}, + {"leading NBSP then insert", "\u00a0INSERT INTO t VALUES (1)", true, false}, + {"leading BOM then insert", "\ufeffINSERT INTO t VALUES (1)", true, false}, + {"leading NBSP then select", "\u00a0SELECT 1", false, false}, + {"with CTE alias named set before form feed AS (read)", "WITH set\fAS (SELECT 1) SELECT * FROM set", false, false}, + {"with CTE alias named set before NBSP AS (read)", "WITH set\u00a0AS (SELECT 1) SELECT * FROM set", false, false}, + + // ClickHouse block comments nest. + {"nested block comment hiding select then insert", "/* a /* b */ SELECT */ INSERT INTO t VALUES (1)", true, false}, + {"nested block comment hiding insert then select", "/* a /* b */ INSERT */ SELECT 1", false, false}, + {"unclosed nested block comment", "/* a /* b */ INSERT INTO t VALUES (1)", false, true}, + {"with nested block comment hiding select then insert", "WITH x AS (SELECT 1) /* a /* b */ SELECT */ INSERT INTO t SELECT * FROM x", true, false}, + {"with nested block comment hiding insert then select", "WITH x AS (SELECT 1) /* a /* b */ INSERT */ SELECT * FROM x", false, false}, + {"with nested block comment before CTE AS (read)", "WITH set /* a /* b */ c */ AS (SELECT 1) SELECT * FROM set", false, false}, + + // After a WITH list ClickHouse parses only SELECT, a FROM-first SELECT and + // INSERT INTO; any other statement is a syntax error, so nothing runs. + {"with delete", "WITH cte AS (SELECT id FROM x) DELETE FROM t WHERE id IN (SELECT id FROM cte)", false, true}, + {"with alter update", "WITH cte AS (SELECT 1) ALTER TABLE t UPDATE a=1 WHERE id IN (SELECT id FROM cte)", false, true}, + {"with truncate", "WITH cte AS (SELECT 1) TRUNCATE TABLE t", false, true}, + {"with CTE alias named update then alter update", "WITH update AS (SELECT id FROM x) ALTER TABLE other UPDATE c=1 WHERE id IN (SELECT id FROM update)", false, true}, + {"with from-first select", "WITH 1 AS x FROM system.one SELECT x", false, false}, + {"with from-first select from a subquery", "WITH 1 AS x FROM (SELECT 1) SELECT x", false, false}, + + // A WITH list's names may be spelled like any keyword: a CTE name, an + // alias, a function, a lambda parameter, an operand, a qualified name's + // part, an array element or a bare element. + {"with alias desc then insert", "WITH 'd' AS desc INSERT INTO t SELECT length(desc)", true, false}, + {"with CTE named desc then insert", "WITH desc AS (SELECT 1 AS x) INSERT INTO t SELECT * FROM desc", true, false}, + {"with CTE named check then insert", "WITH check AS (SELECT 1 AS x) INSERT INTO t SELECT * FROM check", true, false}, + {"with CTE named explain then insert", "WITH explain AS (SELECT 1 AS x) INSERT INTO t SELECT * FROM explain", true, false}, + {"with alias show then insert", "WITH 1 AS show INSERT INTO t SELECT show", true, false}, + {"with alias describe then insert", "WITH 1 AS describe INSERT INTO t SELECT describe", true, false}, + {"with EXISTS expression then insert", "WITH EXISTS(SELECT 1 FROM t WHERE x = 7) AS seen INSERT INTO t2 SELECT seen", true, false}, + {"with alias select then insert", "WITH 1 AS select INSERT INTO t SELECT 1", true, false}, + {"with operand select then insert", "WITH 1 + select AS y INSERT INTO t SELECT y", true, false}, + {"with qualified select then insert", "WITH t.select AS y INSERT INTO t SELECT y", true, false}, + {"with array of select then insert", "WITH [select] AS a INSERT INTO t SELECT a", true, false}, + {"with bare element from then insert", "WITH from INSERT INTO t SELECT 1", true, false}, + {"with comment between INSERT and INTO", "WITH 1 AS x INSERT /* c */ INTO t SELECT x", true, false}, + {"with alias set (read)", "WITH 1 AS set SELECT set", false, false}, + {"with alias use (read)", "WITH 1 AS use SELECT use", false, false}, + {"with alias kill (read)", "WITH 1 AS kill SELECT kill", false, false}, + {"with alias system (read)", "WITH 1 AS system SELECT system", false, false}, + {"with alias insert (read)", "WITH 1 AS insert SELECT insert", false, false}, + {"with alias alter then from-first select", "WITH 1 AS alter FROM system.one SELECT alter", false, false}, + {"with lambda parameter set (read)", "WITH set -> 1 AS f SELECT f(2)", false, false}, + {"with lambda parameter insert (read)", "WITH insert -> 1 AS f SELECT f(2)", false, false}, + {"with operand insert (read)", "WITH insert + 1 AS y SELECT y", false, false}, + {"with bare element insert (read)", "WITH insert SELECT 1", false, false}, + {"with function insert (read)", "WITH insert(1) AS y SELECT y", false, false}, + + // A number led by `.` ends with its digits and exponent, unlike one led + // by a digit, so a word glued to it is a word of its own. + {"with dot-led number glued to insert", "WITH 1 AS a, .5INSERT INTO t SELECT a", true, false}, + {"with signed dot-led exponent glued to insert", "WITH -.5e-3INSERT INTO t SELECT 1", true, false}, + {"with dot-led number and digit separator glued to insert", "WITH 1 AS a, .5_0INSERT INTO t SELECT a", true, false}, + {"with dot-led exponent in a sum glued to insert", "WITH 1+.5e3INSERT INTO t SELECT 1", true, false}, + {"with Unicode minus and dot-led number glued to insert", "WITH −.5INSERT INTO t SELECT 1", true, false}, + + // EXECUTE AS runs the statement after the user as that user, so that + // statement is the one classified; bare, it switches the session's user + // and returns no result set. + {"execute as insert", "EXECUTE AS default INSERT INTO t SELECT 1", true, false}, + {"execute as with insert", "EXECUTE AS default WITH 1 AS a INSERT INTO t SELECT a", true, false}, + {"execute as select", "EXECUTE AS default SELECT 1", false, false}, + {"execute as with select", "EXECUTE AS default WITH 1 AS a SELECT a", false, false}, + {"execute as drop", "EXECUTE AS u1 DROP TABLE t", true, false}, + {"execute as show", "EXECUTE AS u1 SHOW TABLES", false, false}, + {"execute as backticked user then insert", "EXECUTE AS `u 1` INSERT INTO t SELECT 1", true, false}, + {"execute as double-quoted user then select", `EXECUTE AS "u1" SELECT 1`, false, false}, + {"execute as string user glued to insert", "EXECUTE AS 'u1'INSERT INTO t SELECT 1", true, false}, + {"execute as heredoc user then insert", "EXECUTE AS $$u1$$ INSERT INTO t SELECT 1", true, false}, + {"execute as curly-quoted user then insert", "EXECUTE AS “u1” INSERT INTO t SELECT 1", true, false}, + {"execute as user at host then insert", "EXECUTE AS u1@'localhost' INSERT INTO t SELECT 1", true, false}, + {"execute as user at host then select", "EXECUTE AS u1 @ `h` SELECT 1", false, false}, + {"execute as lower with comments then insert", "execute /* c */ as u1 -- c\ninsert into t select 1", true, false}, + {"execute as user named insert then select", "EXECUTE AS insert SELECT 1", false, false}, + {"execute as user named select then insert", "EXECUTE AS select INSERT INTO t SELECT 1", true, false}, + {"execute as bare", "EXECUTE AS u1", true, false}, + {"execute as bare with semicolon", "EXECUTE AS u1;", true, false}, + {"execute as user glued to insert", "EXECUTE AS u1INSERT INTO t SELECT 1", false, true}, + {"execute as nested", "EXECUTE AS u1 EXECUTE AS u2 INSERT INTO t SELECT 1", false, true}, + {"execute without as", "EXECUTE u1 INSERT INTO t SELECT 1", false, true}, + + {"empty", "", false, true}, + {"comment only", "-- just a comment", false, true}, + {"unclosed block comment", "/* never closed", false, true}, +} diff --git a/tests/integration/ismutation_test.go b/tests/integration/ismutation_test.go new file mode 100644 index 000000000..1458a9337 --- /dev/null +++ b/tests/integration/ismutation_test.go @@ -0,0 +1,222 @@ +//go:build integration + +package tests + +import ( + "context" + "errors" + "fmt" + "sort" + "strings" + "sync" + "testing" + + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/Wave-RF/WaveHouse/internal/api" + "github.com/Wave-RF/WaveHouse/internal/chconn" + "github.com/Wave-RF/WaveHouse/internal/testutil/mutationtest" +) + +// codeSyntaxError is ClickHouse's SYNTAX_ERROR. +const codeSyntaxError = 62 + +// astMutation is api.IsMutation's answer for each statement kind astRoot +// names. +var astMutation = map[string]bool{ + // A bare EXECUTE AS: it switches the session's user and returns no + // result set. + "ExecuteAsQuery": true, + "SelectWithUnionQuery": false, + "ShowTables": false, + "DescribeQuery": false, + "Explain": false, + "ExistsTableQuery": false, + "InsertQuery": true, + "UpdateQuery": true, + "DeleteQuery": true, + "TruncateQuery": true, + "DropQuery": true, + "AlterQuery": true, + "CreateQuery": true, + "Rename": true, + "OptimizeQuery": true, + "GrantQuery": true, + "SYSTEM": true, + "AttachQuery": true, + "DetachQuery": true, + "KillQueryQuery": true, + "Set": true, + "UseQuery": true, +} + +// astRoot is the statement kind ClickHouse's parser gives sql — the root node +// of its tree, or for an EXECUTE AS that leads a statement, that statement's +// — or the error it rejects sql with. Parsing only: nothing runs, and no table +// need exist. +func astRoot(ctx context.Context, conn driver.Conn, sql string) (string, error) { + rows, err := conn.Query(ctx, "EXPLAIN AST "+sql) + if err != nil { + return "", err + } + defer func() { _ = rows.Close() }() + var lines []string + for rows.Next() { + var line string + if err := rows.Scan(&line); err != nil { + return "", err + } + lines = append(lines, line) + } + if err := rows.Err(); err != nil { + return "", err + } + if len(lines) == 0 { + return "", errors.New("EXPLAIN AST returned no rows") + } + root := strings.Fields(lines[0]) + if len(root) == 0 { + return "", fmt.Errorf("EXPLAIN AST returned %q", lines[0]) + } + if root[0] != "ExecuteAsQuery" { + return root[0], nil + } + // Its children, indented one space: the user, then any statement. + var children []string + for _, line := range lines[1:] { + if strings.HasPrefix(line, " ") && !strings.HasPrefix(line, " ") { + children = append(children, strings.Fields(line)[0]) + } + } + if len(children) < 2 { + return root[0], nil + } + return children[1], nil +} + +func isSyntaxError(err error) bool { + code, ok := chconn.ExceptionCode(err) + return ok && code == codeSyntaxError +} + +// TestIsMutation_AgreesWithClickHouseParser checks every shared IsMutation +// case against the parser of the pinned ClickHouse: a case that parses is +// classified as its statement kind, and one marked unparsed still fails to. +func TestIsMutation_AgreesWithClickHouseParser(t *testing.T) { + e := env(t) + ctx := context.Background() + for _, tc := range mutationtest.Cases { + t.Run(tc.Name, func(t *testing.T) { + root, err := astRoot(ctx, e.chConn, tc.SQL) + if tc.Unparsed { + require.Error(t, err, "ClickHouse now parses this case as %s: set Mutation to match and drop Unparsed", root) + assert.True(t, isSyntaxError(err), "want a syntax error, got %v", err) + return + } + require.NoError(t, err) + want, known := astMutation[root] + require.True(t, known, "no IsMutation answer for statement kind %s: add it to astMutation", root) + assert.Equal(t, want, tc.Mutation, "ClickHouse parses this as %s", root) + }) + } +} + +// TestIsMutation_KeywordNamesInWithList spells a WITH list's names — a CTE, +// an alias, a function, a lambda parameter, an operand, a qualified name's +// part, an array element, a bare element — as every ClickHouse keyword, ahead +// of each statement a WITH list can lead, and checks IsMutation against the +// parser on each combination that parses. +func TestIsMutation_KeywordNamesInWithList(t *testing.T) { + e := env(t) + ctx := context.Background() + + rows, err := e.chConn.Query(ctx, "SELECT keyword FROM system.keywords") + require.NoError(t, err) + seen := map[string]bool{} + for rows.Next() { + var k string + require.NoError(t, rows.Scan(&k)) + for _, w := range strings.Fields(k) { + seen[strings.ToLower(w)] = true + } + } + require.NoError(t, rows.Err()) + _ = rows.Close() + words := make([]string, 0, len(seen)) + for w := range seen { + words = append(words, w) + } + sort.Strings(words) + require.NotEmpty(t, words) + + shapes := []string{ + "WITH %s AS (SELECT 1 AS x) %s", + "WITH 1 AS a, %s AS (SELECT 1) %s", + "WITH 1 AS %s %s", + "WITH (SELECT 1) AS a, 2 AS %s %s", + "WITH %s AS y %s", + "WITH %s(1) AS y %s", + "WITH %s -> 1 AS f %s", + "WITH 1 + %s AS y %s", + "WITH %s + 1 AS y %s", + "WITH t.%s AS y %s", + "WITH [%s] AS a %s", + "WITH %s %s", + } + statements := []string{"SELECT 1", "FROM system.one SELECT 1", "INSERT INTO t SELECT 1"} + queries := make(chan string) + go func() { + defer close(queries) + for _, w := range words { + for _, shape := range shapes { + for _, stmt := range statements { + queries <- fmt.Sprintf(shape, w, stmt) + } + } + } + }() + + var ( + mu sync.Mutex + wg sync.WaitGroup + parsed int + failures []string + ) + for range 8 { + wg.Add(1) + go func() { + defer wg.Done() + for sql := range queries { + root, err := astRoot(ctx, e.chConn, sql) + var failure string + switch want, known := astMutation[root]; { + case err != nil && !isSyntaxError(err): + failure = fmt.Sprintf("%q: %v", sql, err) + case err != nil: + case !known: + failure = fmt.Sprintf("%q: no IsMutation answer for statement kind %s", sql, root) + case api.IsMutation(sql) != want: + failure = fmt.Sprintf("%q: ClickHouse parses it as %s, IsMutation says %v", sql, root, !want) + } + mu.Lock() + if err == nil { + parsed++ + } + if failure != "" { + failures = append(failures, failure) + } + mu.Unlock() + } + }() + } + wg.Wait() + + sort.Strings(failures) + for _, f := range failures { + t.Error(f) + } + // Most combinations parse; far fewer means the check stopped checking. + assert.Greater(t, parsed, len(words)*len(shapes)*len(statements)/2) +}