Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
c3709bc
refactor(ingest)!: lock pipeline to inserts only; mutations require /…
taitelee May 19, 2026
cd2a8ba
Merge remote-tracking branch 'origin/main' into lock-ingest-inserts-only
taitelee May 19, 2026
50af71d
fix(api): bypass cache for mutations on /v1/query
taitelee May 19, 2026
d51342a
test: bump integration timeout to 120s to prevent flakes on loaded ru…
taitelee May 19, 2026
2b174b5
test: added mutation verb tests
taitelee May 19, 2026
a23b72e
fix: route EXCHANGE through Exec; drop dead inFlight counter
taitelee May 19, 2026
f044a11
Merge remote-tracking branch 'origin/main' into lock-ingest-inserts-only
taitelee May 19, 2026
4075a2d
fix(query): parse CTE clauses in isMutation to correctly classify WIT…
taitelee May 19, 2026
a65aebe
fix(query): add comment opacity to isMutation scanner to prevent fals…
taitelee May 19, 2026
adfc2ec
refactor(ingest): drop action field from wire format
taitelee May 19, 2026
45bc840
refactor(api)!: lock raw SQL to /v1/admin/query under admin/service
EricAndrechek May 19, 2026
421c3ce
refactor(api): strip cache + singleflight from /v1/admin/query
EricAndrechek May 19, 2026
7e93528
fix(integration): drop stale NewQueryHandler args in test harness
EricAndrechek May 19, 2026
0fe15b0
docs: fix query-path mermaid in why-wavehouse.md
EricAndrechek May 19, 2026
38365c7
refactor(api)!: proxy /v1/admin/query to ClickHouse HTTP, drop verb c…
EricAndrechek May 19, 2026
72aeb37
refactor(api): tighten /v1/admin/query proxy per PR review feedback
EricAndrechek May 19, 2026
1d04493
refactor(api): close self-review findings on /v1/admin/query proxy
EricAndrechek May 19, 2026
e816607
refactor(api): address CodeRabbit review on /v1/admin/query proxy
EricAndrechek May 20, 2026
3a4d8dd
docs/api: address CodeRabbit round-2 on /v1/admin/query proxy
EricAndrechek May 20, 2026
cfa0795
fix(api): isMutation no longer false-positives CTE-prefixed reads
EricAndrechek May 20, 2026
0486bc4
fix(ingest)+docs: restore non-insert envelope rejection + CR round-4 …
EricAndrechek May 20, 2026
0e8c07b
docs: qualify API-layer admin gate as auth-enabled-only
EricAndrechek May 20, 2026
b6bcb8e
refactor(api): CR round-6 — auth qualifiers + integration test consol…
EricAndrechek May 20, 2026
6d37646
Merge remote-tracking branch 'origin/main' into lock-ingest-inserts-only
EricAndrechek May 20, 2026
adbb10b
fix(api): reject unknown fields on /v1/admin/query requests
EricAndrechek May 20, 2026
cebca98
test(api): bound post-cancel wait in ContextCancelPropagates
EricAndrechek May 20, 2026
f16b15b
fix(api): correct CTE-lookahead ordering for SELECT-with-parens
EricAndrechek May 20, 2026
6927f2a
refactor(api): CR round-11 — defensive guards + test consolidation
EricAndrechek May 20, 2026
0b84563
refactor: CR round-12 — scope claims to current deployment + 405 guards
EricAndrechek May 20, 2026
e916c45
test(integration): split shared ctx per Eventually phase
EricAndrechek May 20, 2026
b1c048e
removing unneeded check for non ingest data in nats
EricAndrechek May 20, 2026
949715d
final tweaks to remove excessive check for action type in nats
EricAndrechek May 20, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .gemini/styleguide.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ End every review with a one-line verdict: **`Ship it`** (no MUSTs, few/no SHOULD
Walk every diff against:

- **SQL injection** in ClickHouse paths (`BindParams`, dynamic table names, user-supplied filters). The `safeIdentifierRe` regex exists for a reason; flag any SQL built without it.
- **Broken auth / authz**: JWT claim handling, role extraction (`auth.role_claim`), policy templating (`{{ jwt.path }}`), raw-SQL access without `raw_sql: true`.
- **Broken auth / authz**: JWT claim handling, role extraction (`auth.role_claim`), policy templating (`{{ jwt.path }}`). When authentication is enabled, `/v1/admin/query` requires the `admin` or `service` role; with `auth.enabled=false` (or `auth.dev_mode=true`) the endpoint is intentionally open, so flag bypasses only in the auth-enabled path.
- **Sensitive data exposure**: secrets in logs, full error messages to clients, credentials or JWT payload echoed into responses.
- **Security misconfiguration**: CORS allowlist bypass, TLS downgrade, default credentials, permissive-by-default flags.
- **Input validation gaps**: unvalidated JSON reaching ClickHouse, unbounded request sizes, missing rate limits.
Expand Down
2 changes: 1 addition & 1 deletion .github/prompts/pr-review.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ Review against each of these, in this order:
- SQL injection in any ClickHouse-bound path (`BindParams`, query builders, dynamic table names)
- Broken authentication / authorization (JWT claim handling, role extraction, policy templating)
- Sensitive data exposure (secrets in logs, error messages leaking internal state)
- Broken access control (policy bypass, raw-SQL without `raw_sql: true` permission)
- Broken access control (policy bypass, raw-SQL outside the `admin` / `service` role on `/v1/admin/query`)
- Security misconfiguration (CORS, TLS, default credentials, permissive defaults)
- Insufficient logging / monitoring
- SSRF, XXE, deserialization flaws if touched
Expand Down
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ Eleven internal packages under `internal/` (plus `internal/testutil/` for shared
- **`config/`** — YAML + env var config loading (cleanenv)
- **`dedupe/`** — `Deduplicator` interface → `Embedded` (Pebble) — optional, controlled by `dedupe.enabled`
- **`discovery/`** — `SchemaRegistry` that introspects ClickHouse `system.columns` + `Validate()` for ingest payloads
- **`ingest/`** — Bento-based ingest pipeline (`bento.go`: JetStream input → per-table batch INSERT with DLQ output, plus inline delete handling — failed deletes route to `dlq.<table>` and `DoubleAck` rather than `Nak`, to break the infinite-redelivery loop that a deterministic delete error would otherwise produce) + `Sweeper` (Active Sweeper for NATS message lifecycle) + `EventMessage`/`BufferConsumerName` types (`types.go`)
- **`ingest/`** — Bento-based ingest pipeline (`bento.go`: JetStream input → per-table batch INSERT with DLQ output). The pipeline is **insert-only**. The wire format `EventMessage` (`types.go`) carries `{table_name, received_timestamp, data}` and nothing else; the worker validates the table name and the payload's presence, 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}` handler. Non-insert mutations (`DELETE`/`UPDATE`/`TRUNCATE`/…) must go through `POST /v1/admin/query` under the admin/service role (the same gate as the rest of `/v1/admin/*`); when authentication is enabled, the `/v1/admin/*` middleware enforces that role check at the API layer, so non-admin callers never reach the proxy. With `auth.enabled=false` (or `auth.dev_mode=true`) the endpoint is intentionally open — the dev/test posture. Plus `Sweeper` (Active Sweeper for NATS message lifecycle) + `EventMessage`/`BufferConsumerName` types (`types.go`)
- **`mq/`** — `Publisher`/`Subscriber` interfaces → `EmbeddedNATS` + `RemoteNATS`
- **`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 NATS headers.
- **`pipes/`** — Named query pipes: `NamedQuery` type + NATS KV store (`WAVEHOUSE_PIPES`) + `.sql` file bootstrap
Expand Down
22 changes: 21 additions & 1 deletion CHANGELOG.md

Large diffs are not rendered by default.

6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -70,8 +70,10 @@ curl -s -X POST http://localhost:8080/v1/ingest/clicks \
# Check discovered schemas
curl -s http://localhost:8080/v1/schema | jq

# Query data (wait ~5s for the batch flush to ClickHouse)
curl -s -X POST http://localhost:8080/v1/query \
# Query data (wait ~5s for the batch flush to ClickHouse).
# /v1/admin/query requires the admin or service role when auth is on. With
# auth.enabled=false (the default for the quickstart) it's open.
curl -s -X POST http://localhost:8080/v1/admin/query \
-H "Content-Type: application/json" \
-d '{"sql": "SELECT * FROM clicks LIMIT 10"}'

Expand Down
2 changes: 1 addition & 1 deletion SECURITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ WaveHouse handles data and enforces strict isolation:
- **JWT validation**: When auth is enabled, all `/v1/*` endpoints require valid JWTs. Signing supports either an HMAC shared secret or a remote JWKS endpoint (`auth.jwks_url`).
- **Role-based access control**: Roles are extracted from a configurable JWT claim path. Non-admin/service roles have per-table, per-column, row-level policies enforced on ingest and query.
- **Input validation**: JSON payloads are validated against ClickHouse schemas before processing.
- **Query passthrough**: Raw SQL via `POST /v1/query` is restricted to roles with `raw_sql: true` in their policy. Structured queries (`POST /v1/tables/{table}/query`) are validated against schema with permission injection.
- **Query passthrough**: When authentication is enabled, raw SQL via `POST /v1/admin/query` is restricted to the `admin` / `service` role — the same gate as the rest of `/v1/admin/*`. With `auth.enabled=false` or `auth.dev_mode=true` (the dev/test postures) the endpoint is intentionally open; production deployments must enable authentication. 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; service tokens already hold admin-scoped powers across the admin tree, so carving out a tighter gate for raw SQL alone would be inconsistency without a real authorization win. Non-admin callers use structured queries (`POST /v1/tables/{table}/query`, validated against schema with permission injection) or named pipes (`GET/POST /v1/pipes/{name}`); raw-SQL grants to non-admin roles via the policy engine are no longer supported (the `raw_sql` field on policies has been removed).
- **Supply chain**: Third-party GitHub Actions are pinned to full commit SHAs. `govulncheck` runs on every push/PR. Dependabot opens weekly grouped PRs for Go modules and Actions.

## Disclosure Policy
Expand Down
21 changes: 19 additions & 2 deletions clients/ts/src/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ describe('WaveHouseClient.pipe()', () => {
});

describe('WaveHouseClient.sql()', () => {
it('delegates to sql() and POSTs to /v1/query', async () => {
it('delegates to sql() and POSTs to /v1/admin/query', async () => {
fetchSpy.mockResolvedValue(
new Response(JSON.stringify([{ count: 42 }]), { status: 200 }),
);
Expand All @@ -114,7 +114,24 @@ describe('WaveHouseClient.sql()', () => {
const result = await client.sql('SELECT count() FROM clicks');

expect(result.data).toEqual([{ count: 42 }]);
expect(fetchSpy.mock.calls[0][0]).toContain('/v1/query');
expect(fetchSpy.mock.calls[0][0]).toContain('/v1/admin/query');
});

it('throws a migration-clear error when called with a legacy params array', () => {
// The second argument used to be a positional-`?` params array.
// TS callers get a compile-time error; JS callers (or `any`-typed
// call sites) would silently pass the array as `opts` and get
// confusing downstream SQL errors. The runtime guard catches
// that case and points at the migration.
const client = createClient({ baseURL: 'http://localhost:8080' });
// Force-cast so the test compiles under strict TS.
const callWithLegacyParams = () =>
(client.sql as unknown as (q: string, p: unknown) => unknown)(
'SELECT * FROM clicks WHERE id = ?',
['some-id'],
);

expect(callWithLegacyParams).toThrow(/client\.sql\(sql, params\) was removed/);
});
});

Expand Down
21 changes: 18 additions & 3 deletions clients/ts/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,13 +72,28 @@ export class WaveHouseClient<DB extends Database = Database> {
);
}

/** Execute a raw SQL query against ClickHouse. Requires admin/service role when policy is active. */
/**
* Execute a raw SQL query against ClickHouse. Requires admin/service role
* when auth is active. The endpoint proxies straight to ClickHouse's HTTP
* interface so any ClickHouse-accepted SQL works; positional `?` param
* binding is NOT supported — inline literals or use the structured query
* builder for safe binding. See sql.ts for details.
*/
sql<Row = Record<string, unknown>>(
query: string,
params?: unknown[],
opts?: { signal?: AbortSignal },
): Promise<Result<Row[]>> {
return sql<Row>(this._ctx, query, params, opts);
// Migration guard: the second argument used to be a positional-`?`
// params array. TS callers get a compile-time error from the type
// signature, but JS callers (or `any`-typed callsites) would silently
// pass an array as `opts` and only discover the break via downstream
// SQL errors. Throw a clear runtime error pointing at the migration.
if (Array.isArray(opts)) {
throw new Error(
'[WaveHouse SDK] client.sql(sql, params) was removed. The /v1/admin/query endpoint does not accept positional `?` params. Inline literals into the SQL, or use the structured query builder (wh.from(table)…) for safe binding from user input.',
);
}
return sql<Row>(this._ctx, query, opts);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

/** @internal Create a stream for the given table. */
Expand Down
27 changes: 16 additions & 11 deletions clients/ts/src/namespaces.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,26 +22,31 @@ describe('sql', () => {

afterEach(() => vi.restoreAllMocks());

it('POSTs to /v1/query with sql field', async () => {
it('POSTs to /v1/admin/query with sql field', async () => {
const result = await sql(makeCtx(), 'SELECT count() FROM clicks');

expect(result.data).toEqual([{ count: 10 }]);
const [url, init] = fetchSpy.mock.calls[0];
expect(url).toContain('/v1/query');
expect(url).toContain('/v1/admin/query');
expect(JSON.parse(init.body)).toEqual({ sql: 'SELECT count() FROM clicks' });
});

it('includes params when provided', async () => {
await sql(makeCtx(), 'SELECT * FROM clicks WHERE page = ?', ['/home']);

const body = JSON.parse(fetchSpy.mock.calls[0][1].body);
expect(body.params).toEqual(['/home']);
});

it('omits params when empty', async () => {
await sql(makeCtx(), 'SELECT 1');
// The proxy endpoint doesn't accept positional `?` params — ClickHouse's
// HTTP interface uses a different binding model entirely. The SDK no
// longer accepts a params array; callers inline literals or use
// ClickHouse's named-param syntax in the SQL string. This test pins that
// contract: the request body never carries a `params` field, even if the
// SQL string itself contains a `?` (which would now be passed to
// ClickHouse verbatim and trigger a parse error there — the admin's
// responsibility).
it('never sends a params field', async () => {
// Use SQL that actually contains `?` so the test exercises what its
// comment claims — that the SDK passes the `?` through to ClickHouse
// verbatim instead of trying to bind it to a missing params array.
await sql(makeCtx(), 'SELECT ?');

const body = JSON.parse(fetchSpy.mock.calls[0][1].body);
expect(body).toEqual({ sql: 'SELECT ?' });
expect(body.params).toBeUndefined();
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});

Expand Down
41 changes: 34 additions & 7 deletions clients/ts/src/sql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,20 +2,47 @@ import type { Result, HttpContext } from './types.js';
import { request } from './http.js';
import { ok, err } from './errors.js';

/** Execute a raw SQL query against ClickHouse. */
/**
* Execute a raw SQL query against ClickHouse.
*
* Backed by `POST /v1/admin/query`, which is gated on `admin` / `service` —
* the same role set as the rest of `/v1/admin/*`. When authentication is
* enabled, callers must hold a JWT with one of those roles. With
* `auth.enabled=false` or `auth.dev_mode=true` (both dev/test postures)
* the endpoint is open. Non-admin use cases should use the structured
* query builder (`wh.from(table)...`) instead.
*
* The server proxies the SQL string verbatim to ClickHouse's HTTP interface,
* so any ClickHouse-accepted statement works — including multi-statement
* input (`SELECT 1; TRUNCATE t`) and arbitrary DDL/DML/SYSTEM verbs.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
*
* **JSON-row contract.** This helper returns `Result<Row[]>` and assumes the
* response is the standard `FORMAT JSON` envelope (or an empty body for
* no-result mutations — coerced to `[]`). Inline non-JSON `FORMAT` overrides
* (`SELECT 1 FORMAT CSV` / `FORMAT TSV` / `FORMAT Pretty` / etc.) are NOT
* compatible with `sql()` — the proxy passes the upstream Content-Type
* through, and the SDK's JSON decoder throws on a non-JSON body, which the
* retry layer surfaces as a `NETWORK_ERROR` result (not a structured
* format-mismatch error). For CSV/TSV exports, hit `/v1/admin/query`
* directly with `fetch()` and read the body as text.
*
* **No parameter binding.** Positional `?` substitution is not supported.
* The SDK has no way to forward ClickHouse-style named params
* (`WHERE id = {id:UInt32}` with `param_id=42` on the query string) —
* the proxy doesn't forward arbitrary query-string params to ClickHouse
* and the SDK doesn't expose a hook to add them. Inline literals into
* the SQL string, or — for safe binding from user-supplied input — use
* the structured query builder (`wh.from(...).select(...).where(...)`).
*/
export async function sql<Row = Record<string, unknown>>(
ctx: HttpContext,
query: string,
params?: unknown[],
opts?: { signal?: AbortSignal },
): Promise<Result<Row[]>> {
const body: { sql: string; params?: unknown[] } = { sql: query };
if (params && params.length > 0) body.params = params;

const { data, error } = await request<Row[]>(ctx, {
method: 'POST',
path: '/v1/query',
body,
path: '/v1/admin/query',
body: { sql: query },
signal: opts?.signal,
});
if (error) return err(error);
Expand Down
1 change: 0 additions & 1 deletion clients/ts/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,6 @@ export interface RolePermissions {
denied_aggregations?: string[];
max_rows?: number;
max_execution_time_ms?: number;
raw_sql?: boolean;
}

export interface PolicyFilter {
Expand Down
18 changes: 15 additions & 3 deletions cmd/wavehouse/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"errors"
"fmt"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
Expand Down Expand Up @@ -280,7 +281,6 @@ func run() int {
ingestStream, err := ingest.StartIngestWorker(
ctx,
embeddedMQ.NatsConn(),
chConn,
cfg.ClickHouse.Addr,
cfg.ClickHouse.HTTPPort, // Uses 8123 by default
cfg.ClickHouse.HTTPScheme,
Expand Down Expand Up @@ -332,8 +332,20 @@ func run() int {
dlqHandler = api.NewDLQHandler(js, logger)
}

queryHandler := api.NewQueryHandler(chConn, tiered, time.Duration(cfg.Cache.DefaultTTL)*time.Second)
queryHandler.PolicyStore = policyStore
// TODO: is this really the best/right way to do this?
// /v1/admin/query proxies straight to ClickHouse over HTTP — no native
// driver involvement. Construct the base URL from the same fields the
// ingest worker uses, defaulting the scheme to http if blank.
queryHost, _, err := net.SplitHostPort(cfg.ClickHouse.Addr)
if err != nil {
queryHost = cfg.ClickHouse.Addr
}
queryScheme := cfg.ClickHouse.HTTPScheme
if queryScheme == "" {
queryScheme = "http"
}
queryEndpoint := fmt.Sprintf("%s://%s", queryScheme, net.JoinHostPort(queryHost, cfg.ClickHouse.HTTPPort))
queryHandler := api.NewQueryHandler(queryEndpoint, cfg.ClickHouse.Username, cfg.ClickHouse.Password, cfg.ClickHouse.Database)

healthHandler := api.NewHealthHandler(chConn)
healthHandler.Boot = bootState
Expand Down
Loading
Loading