From c4bd138b9b119b0a152f2f37f0fa8714ee00b132 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 15:49:44 +0200 Subject: [PATCH 1/2] Add null and range predicates, variadic and/or, DISTINCT and row locks - isNull, isNotNull, between and notBetween on every expression; the ends of a range are encoded through the column, like a comparison. - CH.and / CH.or take any number of conditions, skip undefined ones, write one flat group, and return undefined when none are left. and keeps tenant evidence; or drops it, as .and / .or do. - distinct() and distinctOn(...aliases) on queries, both dialects. - forUpdate / forNoKeyUpdate / forShare / forKeyShare with skipLocked, noWait and of, written after LIMIT. DialectClauses.locking gates them; ClickHouse refuses at compile. - compile(query) takes params as optional, so a query without params no longer falls through to the write overload with a confusing error. Each feature runs in the shared core suite on Postgres and on both ClickHouse matrix servers. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 7 + design/gap-review.md | 6 +- docs/expressions.md | 30 +-- docs/queries.md | 32 ++++ docs/reference.md | 8 +- scripts/check-doc-examples.mjs | 2 +- src/ch/compile.ts | 37 +++- src/ch/dialect.ts | 3 + src/ch/expr.ts | 51 ++++++ src/ch/index.ts | 3 + src/ch/param.ts | 4 + src/ch/predicates.test.ts | 105 +++++++++++ src/ch/query.test-d.ts | 16 ++ src/ch/query.ts | 65 +++++++ src/database/database.test.ts | 50 +++++ src/expr.ts | 2 + src/pg/dialect.ts | 1 + src/sql/sql-query.ts | 18 +- tests/__snapshots__/core-sql.test.ts.snap | 213 ++++++++++++++++++++++ tests/core-cases.ts | 78 ++++++++ tests/database.clickhouse.test.ts | 33 ++++ tests/dialect-cases.ts | 16 ++ tests/dialect-coverage.test.ts | 6 +- 23 files changed, 763 insertions(+), 23 deletions(-) create mode 100644 src/ch/predicates.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index c99c1e1..6a8dc16 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,13 @@ ## Unreleased +- Add `isNull()`, `isNotNull()`, `between()` and `notBetween()` on every expression, and + variadic `CH.and(...)` / `CH.or(...)` that skip `undefined` and write one flat group. +- Add `distinct()` and `distinctOn(...aliases)` to queries, on both dialects. +- Add Postgres row locks: `forUpdate`, `forNoKeyUpdate`, `forShare`, `forKeyShare`, with + `skipLocked`, `noWait` and `of` (`LockOptions`). Add `DialectClauses.locking`; ClickHouse + refuses them. `SqlQuery` gains `distinct`, `distinctOn` and `lock`. +- `compile(query)` no longer needs a params argument when the query has no params. - Add `update(table).set(...).where(...)` and `deleteFrom(table).where(...)` (see `docs/updates-and-deletes.md`), with `returning` on Postgres and `settings` on ClickHouse, where they compile to an `ALTER TABLE ... UPDATE` mutation and a lightweight `DELETE`. A write diff --git a/design/gap-review.md b/design/gap-review.md index b15345e..c6ff5f2 100644 --- a/design/gap-review.md +++ b/design/gap-review.md @@ -28,10 +28,10 @@ builder; **P1** commonly used; **P2** niche. | ~~DELETE builder: WHERE, RETURNING~~ (built) | ~79 | S | | A typed, value-binding `sql` template usable inside expressions; `sql.join` / `raw` / `empty` on `Db.sql` | ~163 | M | | Postgres column types: `timestamptz` as `Date`, `timestamp`, `date`, `interval`, `varchar(n)`, serial / identity | 226 timestamp columns | S | -| DISTINCT (and DISTINCT ON) | ~10 | S | -| `FOR UPDATE` / `FOR SHARE` / `SKIP LOCKED` / `NOWAIT` | 7 | S | +| ~~DISTINCT (and DISTINCT ON)~~ (built) | ~10 | S | +| ~~`FOR UPDATE` / `FOR SHARE` / `SKIP LOCKED` / `NOWAIT`~~ (built) | 7 | S | | jsonb and array operators (`@>`, `->`, `?`, `&&`, `ANY`) | ~12 | M | -| `isNull` / `isNotNull` / `between`; variadic `and` / `or` that skip `undefined` | everywhere | S | +| ~~`isNull` / `isNotNull` / `between`; variadic `and` / `or` that skip `undefined`~~ (built) | everywhere | S | | Constraint error helpers (unique, foreign key, not null); keep ClickHouse's numeric error codes, which `sqlStateOf` drops today | all upserts | S | | Tenant-scope enforcement in `Database`, opt in, with an explicit cross-tenant entry point | safety | S | | Postgres `defineTable` (indexes, unique, FKs), Postgres migrations, a drizzle-kit importer | 68 tables, 90 indexes, 47 unique, 75 folders | L; can wait, drizzle-kit can keep migrating | diff --git a/docs/expressions.md b/docs/expressions.md index afdf6ea..bf76c01 100644 --- a/docs/expressions.md +++ b/docs/expressions.md @@ -19,24 +19,18 @@ $.Timestamp.gte(new Date(...)) // Timestamp >= '2026-01-01 00:00:00' ### Testing for NULL -`.eq(null)` emits `= NULL`; it does not test whether a value is missing. Use -`isNull` (or `isNotNull` for present values), declared with `defineCondFn`: +`.eq(null)` emits `= NULL`; it does not test whether a value is missing. Use `.isNull()` (or +`.isNotNull()` for present values), which write `IS NULL` and work on every dialect: ```ts title="null-filter.ts" import * as CH from "@maple-dev/effect-orm" import * as T from "@maple-dev/effect-orm/types" const Notes = CH.table("notes", { Note: T.nullable(T.string) }) -const isNull = CH.defineCondFn<[CH.Expr]>("isNull") -export const compiled = CH.compileUnsafe( - CH.from(Notes).select("Note").where(($) => [isNull($.Note)]), - {}, -) -console.log(compiled.sql) // SELECT Note AS Note FROM notes WHERE isNull(Note) +export const compiled = CH.compileUnsafe(CH.from(Notes).select("Note").where(($) => [$.Note.isNull()])) +console.log(compiled.sql) // SELECT Note AS Note FROM notes WHERE Note IS NULL ``` -See [ClickHouse NULL predicates](https://clickhouse.com/docs/reference/functions/regular-functions/functions-for-nulls#isNull). - ### Invalid literals A value the column cannot hold fails while the SQL is being built: @@ -62,6 +56,8 @@ Every `Expr` carries: | `.gt(x)` / `.gte(x)` | `> x` / `>= x` | | `.lt(x)` / `.lte(x)` | `< x` / `<= x` | | `.in_(...xs)` / `.notIn(...xs)` | `IN (…)` / `NOT IN (…)` | +| `.between(a, b)` / `.notBetween(a, b)` | `BETWEEN a AND b` / `NOT BETWEEN a AND b` | +| `.isNull()` / `.isNotNull()` | `IS NULL` / `IS NOT NULL` | Each accepts a raw value or another `Expr`. String literals are escaped; booleans emit as `1` / `0`. @@ -86,8 +82,20 @@ Each accepts a raw value or another `Expr`. String literals are escaped; bool `.and()` / `.or()` parenthesise their result, so precedence is explicit. `CH.not(condition)` wraps in `NOT (…)` and is available from the root and `/expr` subpath. +`CH.and(...)` and `CH.or(...)` take any number of conditions, skip `undefined` ones, and write +one flat group. With none left they return `undefined`, which `where` skips, so optional +filters combine without special cases: + +```ts +.where(($) => [ + $.OrgId.eq("org_123"), + CH.or(CH.when(name, (n) => $.Name.eq(n)), CH.when(minMs, (ms) => $.Ms.gte(ms))), +]) +// both given -> … AND (Name = 'checkout' OR Ms >= 100); neither -> only the OrgId test +``` + The `where` array is AND-joined. [Tenant scoping](./tenant-scoping.md) preserves evidence -through both separate entries and `.and()`; `.or()` discards it. +through both separate entries, `.and()` and `CH.and()`; `.or()` and `CH.or()` discard it. _(Backed by `docs/expressions.md > Combining conditions with and/or`.)_ diff --git a/docs/queries.md b/docs/queries.md index 427e5f1..39cff40 100644 --- a/docs/queries.md +++ b/docs/queries.md @@ -108,6 +108,19 @@ Untyped callers receive `QueryBuilderDefect`; use a tuple for each sort key. _(Backed by `docs/queries.md > orderBy takes tuples` and `> orderBy rejects a bare string`.)_ +## `distinct` / `distinctOn` + +```ts +.select("ServiceName").distinct() +// SELECT DISTINCT ServiceName … + +.select(($) => ({ org: $.OrgId, id: $.Id })).distinctOn("org").orderBy(["org", "asc"], ["id", "desc"]) +// SELECT DISTINCT ON (org) … — the newest row per org +``` + +`distinctOn` takes selected aliases and keeps the first row of each group in ORDER BY order; +Postgres wants those keys to lead the ORDER BY. Both ClickHouse and Postgres support it. + ## `limit` / `offset` ```ts @@ -128,6 +141,25 @@ your request boundary, and enforce an application maximum. Use a stable `orderBy Accepts `"JSON"` or `"JSONEachRow"`. Most clients set the format themselves; use this only when you are sending raw SQL somewhere that does not. +## Row locks + +On Postgres, `forUpdate`, `forNoKeyUpdate`, `forShare` and `forKeyShare` add a locking clause +after LIMIT. Each takes `{ skipLocked?, noWait?, of? }`. The usual job-queue claim: + +```ts +CH.from(Jobs) + .select("id") + .where(($) => [$.state.eq("queued")]) + .orderBy(["id", "asc"]) + .limit(1) + .forUpdate({ skipLocked: true }) +// … LIMIT 1 FOR UPDATE SKIP LOCKED +``` + +A lock lasts until the transaction ends, so run the query inside `Database.transaction`. +`skipLocked` and `noWait` together, and any lock on ClickHouse (which has no row locks), are a +`QueryBuilderDefect`. + ## `withCTE` See [Unions and CTEs](./unions-and-ctes.md#ctes). diff --git a/docs/reference.md b/docs/reference.md index c3cdd48..95bb0ce 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -77,6 +77,8 @@ Note `/sql` exports a `compile` (fragment → string) distinct from the root `co | `orderBy(...[col, dir])` | **Tuples**, not two strings | | `limit(n)` / `offset(n)` | Rounded before emission | | `format(fmt)` | `"JSON"` \| `"JSONEachRow"` | +| `distinct()` / `distinctOn(...aliases)` | `SELECT DISTINCT` / `SELECT DISTINCT ON (…)` | +| `forUpdate` / `forNoKeyUpdate` / `forShare` / `forKeyShare` | Postgres row locks; options `LockOptions` (`skipLocked`, `noWait`, `of`) | | `innerJoin` / `leftJoin` / `crossJoin` | `(table, alias, on?)` | | `innerJoinQuery` / `leftJoinQuery` / `crossJoinQuery` | `(query, alias, on?)` | | `withCTE(name, query)` / `withCTE(name, sql, options?)` | Typed query derives scope; SQL form can declare `options.tenantScope` | @@ -89,7 +91,7 @@ Note `/sql` exports a `compile` (fragment → string) distinct from the root `co | Export | Signature | | -------------------- | ------------------------------------------------------------------------------------ | -| `compile` | `(query, params, options?) => Effect, QueryBuilderError>`; also `(insert, params?, options?)`, whose options (`InsertCompileOptions`) are only `dialect` | +| `compile` | `(query, params?, options?) => Effect, QueryBuilderError>`; also `(insert, params?, options?)`, whose options (`InsertCompileOptions`) are only `dialect` | | `compileUnsafe` | The same, returning `CompiledQuery` and throwing instead | | `compileUnion` | `(union, params, options?) => Effect, QueryBuilderError>` | | `compileUnionUnsafe` | The same, throwing instead | @@ -125,13 +127,15 @@ time; see [Params and compilation](./params-and-compilation.md#what-each-kind-ac | `inExprList(expr, exprs)` | Same for expression lists | | `notInList(expr, values)` | `expr NOT IN ('a', 'b')` | | `not(condition)` | `NOT (…)` | +| `and(...conds)` / `or(...conds)` | One flat `(… AND …)` / `(… OR …)`; skips `undefined`, returns `undefined` when none are left | | `dynamicColumn(name, t?)` | An `Expr` from a runtime column name — a `GROUP BY` alias | | `exists(q)` | `EXISTS (…)` from a query or pre-compiled SQL | | `inSubquery(expr, q)` | `expr IN (…)` from a query or pre-compiled SQL | | `notInSubquery(expr, q)` | `expr NOT IN (…)`; note the NULL semantics | | `outerRef(name)` | Reference an outer column in a correlated subquery | -`Expr` methods: `eq`, `neq`, `gt`, `gte`, `lt`, `lte`, `in_`, `notIn`, `like`, `notLike`, +`Expr` methods: `eq`, `neq`, `gt`, `gte`, `lt`, `lte`, `in_`, `notIn`, `between`, `notBetween`, +`isNull`, `isNotNull`, `like`, `notLike`, `ilike` (string-only), and `add`, `sub`, `mul`, `div`, `mod` (number-only, **no parentheses**). `div` and `mod` decode as `number | null` — ClickHouse sends `inf`/`nan` as JSON `null` — except by a numeric literal of magnitude ≥ 1 (`Quotient`), which keeps the dividend's nullability; use diff --git a/scripts/check-doc-examples.mjs b/scripts/check-doc-examples.mjs index 5143aed..c566632 100644 --- a/scripts/check-doc-examples.mjs +++ b/scripts/check-doc-examples.mjs @@ -95,7 +95,7 @@ assert.equal(typedFunction.compiled.rowSchemaSource, "derived") assert.deepEqual(await Effect.runPromise(typedFunction.compiled.decodeRows([{ name: "checkout", durationMs: "42" }])), [{ name: "checkout", durationMs: 42 }]) await assert.rejects(() => Effect.runPromise(typedFunction.compiled.decodeRows([{ name: 42, durationMs: "oops" }]))) const nullFilter = await import("./null-filter") -assert.match(sql(nullFilter.compiled), /WHERE isNull\\(notes.Note\\)/) +assert.match(sql(nullFilter.compiled), /WHERE notes.Note IS NULL/) assert.deepEqual(await Effect.runPromise(nullFilter.compiled.decodeRows([{ Note: null }])), [{ Note: null }]) const escapedSql = await import("./escaped-sql") const fragments = await import("@maple-dev/effect-orm/sql") diff --git a/src/ch/compile.ts b/src/ch/compile.ts index 65a7439..817f561 100644 --- a/src/ch/compile.ts +++ b/src/ch/compile.ts @@ -49,6 +49,23 @@ export class CompiledQueryEncodeError extends Schema.TaggedError { + if (lock === undefined) return undefined + const dialect = currentDialect() + if (dialect.clauses.locking !== true) { + throw new QueryBuilderDefect({ + message: `CHQuery: FOR ${lock.strength} has no meaning for the ${dialect.name} dialect, which has no row locks`, + }) + } + if (lock.skipLocked === true && lock.noWait === true) { + throw new QueryBuilderDefect({ message: "CHQuery: a lock takes skipLocked or noWait, not both" }) + } + const of = lock.of !== undefined && lock.of.length > 0 ? ` OF ${lock.of.map(quoteIdentPath).join(", ")}` : "" + const wait = lock.skipLocked === true ? " SKIP LOCKED" : lock.noWait === true ? " NOWAIT" : "" + return `FOR ${lock.strength}${of}${wait}` +} + /** `orderBy` takes `[column, direction]` tuples. A bare string is the natural * mistake (`.orderBy("count", "desc")`), and it is invisible without types: * destructuring a string yields its first two characters, so `"count"` used to @@ -531,11 +548,12 @@ export function compileCH< Output extends Record, Joins extends Record, Route extends string | undefined, - Params extends Record, + Params extends Record = {}, Decoded extends Output = Output, >( query: CHQuery, - params: Params, + /** Values for the query's `param.*` markers. Optional when it has none. */ + params?: Params, options?: { skipFormat?: boolean rowSchema?: CompiledQueryRowSchema @@ -578,11 +596,12 @@ export function compileCHUnsafe< Output extends Record, Joins extends Record, Route extends string | undefined, - Params extends Record, + Params extends Record = {}, Decoded extends Output = Output, >( query: CHQuery, - params: Params, + /** Values for the query's `param.*` markers. Optional when it has none. */ + params?: Params, options?: { skipFormat?: boolean rowSchema?: CompiledQueryRowSchema @@ -827,6 +846,16 @@ function compileInner< }) const sqlQuery: SqlQuery = { + distinct: state.distinct !== undefined, + distinctOn: Array.isArray(state.distinct) + ? state.distinct.map((key) => { + if (!(options?.selectKeys ?? keys).includes(key)) { + throw new QueryBuilderDefect({ message: `CHQuery: distinctOn(${JSON.stringify(key)}) is not a selected alias` }) + } + return raw(quoteIdent(key)) + }) + : undefined, + lock: lockClause(state.lock), select: selectFragments, from: fromFragment, joins, diff --git a/src/ch/dialect.ts b/src/ch/dialect.ts index 00184ec..1574915 100644 --- a/src/ch/dialect.ts +++ b/src/ch/dialect.ts @@ -68,6 +68,9 @@ export interface DialectClauses { /** `SETTINGS` on an INSERT, UPDATE or DELETE (ClickHouse). Absent means * no: a write with `.settings()` fails to compile for the dialect. */ readonly writeSettings?: boolean + /** Row locking on a SELECT (`FOR UPDATE`, `FOR SHARE`, `SKIP LOCKED`). Absent + * means no: a query with `.forUpdate()` and the like fails to compile. */ + readonly locking?: boolean /** UPDATE is written `ALTER TABLE t UPDATE ... WHERE ...`, a ClickHouse * mutation, rather than `UPDATE t SET ... WHERE ...`. */ readonly alterTableUpdate?: boolean diff --git a/src/ch/expr.ts b/src/ch/expr.ts index b41c1b3..0b3420e 100644 --- a/src/ch/expr.ts +++ b/src/ch/expr.ts @@ -71,6 +71,22 @@ export interface Expr { notLike(this: Expr, pattern: string): Condition ilike(this: Expr, pattern: string): Condition + // NULL and ranges + /** `expr IS NULL`. */ + isNull(): Condition + /** `expr IS NOT NULL`. */ + isNotNull(): Condition + /** `expr BETWEEN low AND high`, both ends included. */ + between( + low: Comparable> | Expr | Expr>, + high: Comparable> | Expr | Expr>, + ): Condition + /** `expr NOT BETWEEN low AND high`. */ + notBetween( + low: Comparable> | Expr | Expr>, + high: Comparable> | Expr | Expr>, + ): Condition + // IN / NOT IN in_(...values: Array>>): Condition notIn(...values: Array>>): Condition @@ -248,6 +264,13 @@ export function makeExpr( lt: (other) => makeCond(lazy(() => `${compile(fragment)} < ${compile(operand(other))}`)), lte: (other) => makeCond(lazy(() => `${compile(fragment)} <= ${compile(operand(other))}`)), + isNull: () => makeCond(lazy(() => `${compile(fragment)} IS NULL`)), + isNotNull: () => makeCond(lazy(() => `${compile(fragment)} IS NOT NULL`)), + between: (low, high) => + makeCond(lazy(() => `${compile(fragment)} BETWEEN ${compile(operand(low))} AND ${compile(operand(high))}`)), + notBetween: (low, high) => + makeCond(lazy(() => `${compile(fragment)} NOT BETWEEN ${compile(operand(low))} AND ${compile(operand(high))}`)), + like: (pattern: string) => makeCond(lazy(() => `${compile(fragment)} LIKE ${compile(str(pattern))}`)), notLike: (pattern: string) => makeCond(lazy(() => `${compile(fragment)} NOT LIKE ${compile(str(pattern))}`)), ilike: (pattern: string) => makeCond(lazy(() => `${compile(fragment)} ILIKE ${compile(str(pattern))}`)), @@ -443,6 +466,34 @@ export function notInList(expr: Expr, values: readonly string[]): Condit return makeCond(lazy(() => `${compile(expr.toFragment())} NOT IN (${escaped()})`)) } +/** + * Conditions AND-joined, an `undefined` one skipped: `and(a, when(x, f), b)`. + * With none left it is `undefined`, which a `where` list skips in turn. Tenant + * evidence carries through, as with `.and`. + */ +export function and(...conditions: ReadonlyArray): Condition +export function and(...conditions: ReadonlyArray): Condition | undefined +export function and(...conditions: ReadonlyArray): Condition | undefined { + const present = conditions.filter((c): c is Condition => c !== undefined) + if (present.length <= 1) return present[0] + return markTenantPredicate( + makeCond(lazy(() => `(${present.map((c) => compile(c.toFragment())).join(" AND ")})`)), + present.flatMap((c) => tenantPredicatesOf(c)), + ) +} + +/** + * Conditions OR-joined, an `undefined` one skipped. With none left it is + * `undefined`. An OR proves no tenant, so it carries no tenant evidence. + */ +export function or(...conditions: ReadonlyArray): Condition +export function or(...conditions: ReadonlyArray): Condition | undefined +export function or(...conditions: ReadonlyArray): Condition | undefined { + const present = conditions.filter((c): c is Condition => c !== undefined) + if (present.length <= 1) return present[0] + return makeCond(lazy(() => `(${present.map((c) => compile(c.toFragment())).join(" OR ")})`)) +} + /** Wrap a condition in NOT (...). */ export function not(condition: Condition): Condition { return makeCond(lazy(() => `NOT (${compile(condition.toFragment())})`)) diff --git a/src/ch/index.ts b/src/ch/index.ts index 9e193b0..bd557be 100644 --- a/src/ch/index.ts +++ b/src/ch/index.ts @@ -68,6 +68,8 @@ export { notInList, // Negating a condition is table stakes; it was `/expr`-only. not, + and, + or, outerRef, // Reference an output alias (a GROUP BY key or aggregate) that isn't on the // column accessor — the usual way to write a `having()` body. @@ -234,6 +236,7 @@ export { type JoinOnCallback, type InferOutput, type InferQueryOutput, + type LockOptions, from, fromQuery, fromUnion, diff --git a/src/ch/param.ts b/src/ch/param.ts index 0a00432..bb6ff17 100644 --- a/src/ch/param.ts +++ b/src/ch/param.ts @@ -102,6 +102,10 @@ function makeParamMarker( like: raise, notLike: raise, ilike: raise, + isNull: raise, + isNotNull: raise, + between: raise, + notBetween: raise, div: raise, mul: raise, add: raise, diff --git a/src/ch/predicates.test.ts b/src/ch/predicates.test.ts new file mode 100644 index 0000000..a08533f --- /dev/null +++ b/src/ch/predicates.test.ts @@ -0,0 +1,105 @@ +import { describe, expect, it } from "@effect/vitest" +import { Effect, Exit } from "effect" +import * as CH from "./index" +import * as PG from "../postgres" +import { QueryBuilderDefect } from "./errors" + +const Events = CH.table( + "events", + { OrgId: CH.string, Name: CH.string, Ms: CH.uint64, Note: CH.nullable(CH.string), At: CH.dateTime }, + { tenantColumn: "OrgId" }, +) +const Jobs = CH.table("jobs", { id: PG.int4, org: PG.text, state: PG.text, done_at: PG.nullable(PG.timestamptz) }) + +const failure = (exit: Exit.Exit) => + Exit.isFailure(exit) ? exit.cause.reasons.map((r) => ("error" in r ? r.error : "defect" in r ? r.defect : r))[0] : undefined +const where = (sql: string) => sql.slice(sql.indexOf("WHERE")) + +describe("null and range predicates", () => { + it("isNull, isNotNull, between and notBetween, values encoded by the column", () => { + const compiled = CH.compileUnsafe( + CH.from(Events) + .select("Name") + .where(($) => [ + $.Note.isNull(), + $.Name.isNotNull(), + $.Ms.between(10, CH.param.int("hi")), + $.At.notBetween(new Date(0), "2026-01-01 00:00:00"), + ]), + { hi: 20 }, + ) + expect(where(compiled.sql)).toBe( + "WHERE events.Note IS NULL\n AND events.Name IS NOT NULL\n AND events.Ms BETWEEN 10 AND 20\n" + + " AND events.At NOT BETWEEN '1970-01-01 00:00:00' AND '2026-01-01 00:00:00'", + ) + }) + + it("binds between's ends on Postgres", () => { + const compiled = PG.compileUnsafe( + CH.from(Jobs).select("id").where(($) => [$.id.between(CH.param.int("lo"), 9), $.done_at.isNull()]), + { lo: 1 }, + ) + expect(where(compiled.sql)).toBe('WHERE "jobs"."id" BETWEEN $1 AND 9\n AND "jobs"."done_at" IS NULL') + expect(compiled.parameters).toEqual([1]) + }) +}) + +describe("and / or", () => { + it("skip undefined, render flat, and return undefined when nothing is left", () => { + const q = (build: (e: CH.ColumnAccessor) => CH.Condition | undefined) => + where(CH.compileUnsafe(CH.from(Events).select("Name").where(($) => [build($)])).sql) + expect(q(($) => CH.or($.Name.eq("a"), undefined, $.Name.eq("b"), $.Ms.gt(1)))).toBe( + "WHERE (events.Name = 'a' OR events.Name = 'b' OR events.Ms > 1)", + ) + expect(q(($) => CH.and($.Name.eq("a"), CH.or($.Ms.lt(1), $.Ms.gt(9))))).toBe( + "WHERE (events.Name = 'a' AND (events.Ms < 1 OR events.Ms > 9))", + ) + expect(q(($) => CH.and(undefined, $.Name.eq("a")))).toBe("WHERE events.Name = 'a'") + expect(CH.and(undefined, undefined)).toBeUndefined() + expect(CH.or()).toBeUndefined() + expect(CH.compileUnsafe(CH.from(Events).select("Name").where(() => [CH.or(undefined)])).sql).not.toContain("WHERE") + }) + + it("and carries tenant evidence; or does not", () => { + const scope = (build: (e: CH.ColumnAccessor) => CH.Condition | undefined) => + CH.compileUnsafe(CH.from(Events).select("Name").where(($) => [build($)])).tenantScope + expect(scope(($) => CH.and($.OrgId.eq("o"), $.Ms.gt(1)))).toBe("single-tenant") + expect(scope(($) => CH.or($.OrgId.eq("o"), $.OrgId.eq("p")))).toBe("cross-tenant") + }) +}) + +describe("distinct", () => { + it("SELECT DISTINCT and DISTINCT ON on both dialects", () => { + expect(CH.compileUnsafe(CH.from(Events).select("Name").distinct()).sql).toMatch(/^SELECT DISTINCT\n/) + const on = CH.from(Jobs) + .select(($) => ({ org: $.org, id: $.id })) + .distinctOn("org") + .orderBy(["org", "asc"], ["id", "desc"]) + expect(PG.compileUnsafe(on).sql).toMatch(/^SELECT DISTINCT ON \("org"\)\n/) + expect(CH.compileUnsafe(on).sql).toMatch(/^SELECT DISTINCT ON \(org\)\n/) + expect(CH.compileUnsafe(on.distinct()).sql).toMatch(/^SELECT DISTINCT\n/) + }) +}) + +describe("row locking", () => { + it("FOR UPDATE / NO KEY UPDATE / SHARE / KEY SHARE with OF, SKIP LOCKED and NOWAIT, after LIMIT", () => { + const base = CH.from(Jobs).select("id").where(($) => [$.state.eq("queued")]).limit(1) + const tail = (q: CH.CHQuery) => PG.compileUnsafe(q, {}).sql.split("\n").at(-1)!.trim() + expect(tail(base.forUpdate({ skipLocked: true }))).toBe("FOR UPDATE SKIP LOCKED") + expect(PG.compileUnsafe(base.forUpdate({ skipLocked: true }), {}).sql).toMatch(/LIMIT 1\n\s+FOR UPDATE SKIP LOCKED$/) + expect(tail(base.forNoKeyUpdate({ noWait: true }))).toBe("FOR NO KEY UPDATE NOWAIT") + expect(tail(base.forShare({ of: ["jobs"] }))).toBe('FOR SHARE OF "jobs"') + expect(tail(base.forKeyShare())).toBe("FOR KEY SHARE") + expect(tail(base.forShare().forUpdate())).toBe("FOR UPDATE") + }) + + it.effect("is a defect on ClickHouse and with skipLocked and noWait together", () => + Effect.gen(function* () { + const base = CH.from(Jobs).select("id") + expect(failure(yield* Effect.exit(CH.compile(base.forUpdate(), {})))).toBeInstanceOf(QueryBuilderDefect) + expect(failure(yield* Effect.exit(PG.compile(base.forUpdate({ skipLocked: true, noWait: true }), {})))).toBeInstanceOf( + QueryBuilderDefect, + ) + }), + ) +}) diff --git a/src/ch/query.test-d.ts b/src/ch/query.test-d.ts index 70bc8b0..457912e 100644 --- a/src/ch/query.test-d.ts +++ b/src/ch/query.test-d.ts @@ -305,3 +305,19 @@ expectTypeOf().toEqualTypeOf<{ readonly id: string readonly age: number }>() + +// Predicates, and/or, distinct and locking +{ + const T = CH.table("t", { a: CH.string, n: CH.uint64, o: CH.nullable(CH.string) }) + CH.from(T).select("a").where(($) => [$.o.isNull(), $.n.between(1, CH.param.int("hi")), $.a.notBetween("a", "m")]) + // @ts-expect-error between takes the column's type + CH.from(T).select("a").where(($) => [$.n.between("1", 2)]) + expectTypeOf(CH.and(CH.rawCond("1"), CH.rawCond("2"))).toEqualTypeOf() + expectTypeOf(CH.or(CH.rawCond("1"), undefined)).toEqualTypeOf() + CH.from(T).select("a", "n").distinctOn("a") + // @ts-expect-error distinctOn takes selected aliases + CH.from(T).select("a").distinctOn("n") + // @ts-expect-error distinctOn needs a key + CH.from(T).select("a").distinctOn() + CH.from(T).select("a").forUpdate({ skipLocked: true, of: ["t"] }) +} diff --git a/src/ch/query.ts b/src/ch/query.ts index 0a17f0a..351a77e 100644 --- a/src/ch/query.ts +++ b/src/ch/query.ts @@ -61,6 +61,21 @@ export type InferOutput = { type OrderBySpec = [keyof Output & string, "asc" | "desc"] +/** How a locking clause waits for rows another transaction holds. */ +export interface LockOptions { + /** `SKIP LOCKED`: leave out rows another transaction has locked. */ + readonly skipLocked?: boolean + /** `NOWAIT`: fail at once instead of waiting. Not with `skipLocked`. */ + readonly noWait?: boolean + /** `OF alias, ...`: lock only these tables' rows (the FROM alias or table name, or join aliases). */ + readonly of?: ReadonlyArray +} + +/** @internal — the locking clause of a query. */ +export interface LockClause extends LockOptions { + readonly strength: "UPDATE" | "NO KEY UPDATE" | "SHARE" | "KEY SHARE" +} + /** Callback for ON conditions — receives main and joined column accessors. */ export type JoinOnCallback = ( main: ColumnAccessor, @@ -98,6 +113,10 @@ export interface CHQueryState { readonly limitValue?: number readonly offsetValue?: number readonly formatValue?: string + /** Set by `distinct` (`true`) or `distinctOn` (the output aliases). */ + readonly distinct?: true | ReadonlyArray + /** Set by `forUpdate` and the other locking methods. */ + readonly lock?: LockClause /** Execution-route metadata carried onto the CompiledQuery (see compile.ts). */ readonly routeValue?: string /** Set by `.crossTenant()`. Forces `tenantScope: "cross-tenant"` (see compile.ts). */ @@ -175,6 +194,28 @@ export interface CHQuery< format(fmt: "JSON" | "JSONEachRow"): CHQuery + /** `SELECT DISTINCT`: drop duplicate output rows. */ + distinct(): CHQuery + + /** + * `SELECT DISTINCT ON (keys)`: keep the first row of each group of these + * output aliases, in ORDER BY order (Postgres wants the keys to lead the + * ORDER BY). Replaces `distinct()`. + */ + distinctOn(...keys: [keyof Output & string, ...Array]): CHQuery + + /** + * `FOR UPDATE`: lock the selected rows until the transaction ends. Run it + * inside `Database.transaction`. Postgres only; replaces any earlier lock. + */ + forUpdate(options?: LockOptions): CHQuery + /** `FOR NO KEY UPDATE`: as `forUpdate`, without blocking inserts that reference the rows. */ + forNoKeyUpdate(options?: LockOptions): CHQuery + /** `FOR SHARE`: a shared lock, which blocks writers but not other sharers. */ + forShare(options?: LockOptions): CHQuery + /** `FOR KEY SHARE`: the weakest lock, blocking only deletes and key updates. */ + forKeyShare(options?: LockOptions): CHQuery + /** * Tag this query with an execution route, carried through to the compiled * query as a type-level fact. The tag is opaque to the builder: what routes @@ -445,6 +486,30 @@ function makeQuery< return makeQuery({ ...state, formatValue: fmt }) }, + distinct() { + return makeQuery({ ...state, distinct: true }) + }, + + distinctOn(...keys) { + return makeQuery({ ...state, distinct: keys as ReadonlyArray }) + }, + + forUpdate(options = {}) { + return makeQuery({ ...state, lock: { strength: "UPDATE", ...options } }) + }, + + forNoKeyUpdate(options = {}) { + return makeQuery({ ...state, lock: { strength: "NO KEY UPDATE", ...options } }) + }, + + forShare(options = {}) { + return makeQuery({ ...state, lock: { strength: "SHARE", ...options } }) + }, + + forKeyShare(options = {}) { + return makeQuery({ ...state, lock: { strength: "KEY SHARE", ...options } }) + }, + route(route) { return makeQuery({ ...state, routeValue: route }) }, diff --git a/src/database/database.test.ts b/src/database/database.test.ts index db886f9..8f79c87 100644 --- a/src/database/database.test.ts +++ b/src/database/database.test.ts @@ -520,6 +520,56 @@ layer(Live, { excludeTestServices: true })("Database on PGlite", (it) => { }), ) + it.effect("claims a job with FOR UPDATE SKIP LOCKED, and reads with DISTINCT ON, IS NULL and BETWEEN", () => + Effect.gen(function* () { + yield* Db.execute(Db.sql`CREATE TABLE jobs (id int4 PRIMARY KEY, org text NOT NULL, state text NOT NULL, done_at timestamptz)`) + const Jobs = CH.table("jobs", { id: PG.int4, org: PG.text, state: PG.text, done_at: PG.nullable(PG.timestamptz) }) + yield* Db.run( + CH.insertInto(Jobs).values([ + { id: 1, org: "a", state: "queued" }, + { id: 2, org: "a", state: "queued" }, + { id: 3, org: "b", state: "queued" }, + ]), + ) + const claim = Db.transaction( + Effect.gen(function* () { + const [job] = yield* Db.run( + CH.from(Jobs) + .select("id") + .where(($) => [$.state.eq("queued"), $.done_at.isNull()]) + .orderBy(["id", "asc"]) + .limit(1) + .forUpdate({ skipLocked: true }), + ) + if (job === undefined) return undefined + yield* Db.run(CH.update(Jobs).set({ state: "running" }).where(($) => [$.id.eq(job.id)])) + return job.id + }), + ) + expect(yield* claim).toBe(1) + expect(yield* claim).toBe(2) + const firstPerOrg = yield* Db.run( + CH.from(Jobs) + .select(($) => ({ org: $.org, id: $.id })) + .distinctOn("org") + .orderBy(["org", "asc"], ["id", "desc"]), + ) + expect(firstPerOrg).toEqual([ + { org: "a", id: 2 }, + { org: "b", id: 3 }, + ]) + const inRange = yield* Db.run( + CH.from(Jobs).select("id").where(($) => [$.id.between(2, CH.param.int("hi")), CH.or($.org.eq("b"), $.state.eq("running"))]).orderBy(["id", "asc"]), + { hi: 3 }, + ) + expect(inRange).toEqual([{ id: 2 }, { id: 3 }]) + expect(yield* Db.run(CH.from(Jobs).select("state").distinct().orderBy(["state", "asc"]))).toEqual([ + { state: "queued" }, + { state: "running" }, + ]) + }), + ) + it.effect("an insert inside a failed transaction rolls back", () => Effect.gen(function* () { const table = yield* freshTable diff --git a/src/expr.ts b/src/expr.ts index 9a986c6..58f898c 100644 --- a/src/expr.ts +++ b/src/expr.ts @@ -6,6 +6,7 @@ export { type Condition, type Expr, type MapValueOf, + and, dynamicColumn, inExprList, inList, @@ -15,6 +16,7 @@ export { makeUntypedExpr, not, notInList, + or, outerRef, rawCond, rawExpr, diff --git a/src/pg/dialect.ts b/src/pg/dialect.ts index a6fd259..91d772e 100644 --- a/src/pg/dialect.ts +++ b/src/pg/dialect.ts @@ -104,6 +104,7 @@ export const postgresDialect: Dialect = { parenthesizeUnionBranches: true, returning: true, onConflict: true, + locking: true, }, paramCodecs: { bool: Schema.Boolean, diff --git a/src/sql/sql-query.ts b/src/sql/sql-query.ts index 1f20511..0076a80 100644 --- a/src/sql/sql-query.ts +++ b/src/sql/sql-query.ts @@ -11,6 +11,10 @@ interface SqlJoin { } export interface SqlQuery { + /** `SELECT DISTINCT`. */ + readonly distinct?: boolean + /** `SELECT DISTINCT ON (...)`; implies `distinct`. */ + readonly distinctOn?: ReadonlyArray readonly select: ReadonlyArray readonly from: SqlFragment readonly joins?: ReadonlyArray @@ -20,6 +24,8 @@ export interface SqlQuery { readonly orderBy: ReadonlyArray readonly limit?: SqlFragment readonly offset?: SqlFragment + /** A locking clause written after OFFSET, e.g. `FOR UPDATE SKIP LOCKED`. */ + readonly lock?: string readonly format?: string } @@ -30,7 +36,12 @@ export function compileQuery(q: SqlQuery): string { // SELECT const selectCols = q.select.map(compile).filter(Boolean) - parts.push(`SELECT\n ${selectCols.join(",\n ")}`) + const distinct = q.distinctOn?.length + ? ` DISTINCT ON (${q.distinctOn.map(compile).join(", ")})` + : q.distinct + ? " DISTINCT" + : "" + parts.push(`SELECT${distinct}\n ${selectCols.join(",\n ")}`) // FROM parts.push(`FROM ${compile(q.from)}`) @@ -76,6 +87,11 @@ export function compileQuery(q: SqlQuery): string { parts.push(`OFFSET ${compile(q.offset)}`) } + // Locking + if (q.lock) { + parts.push(q.lock) + } + // FORMAT if (q.format) { parts.push(`FORMAT ${q.format}`) diff --git a/tests/__snapshots__/core-sql.test.ts.snap b/tests/__snapshots__/core-sql.test.ts.snap index 41969fe..c5dbe8c 100644 --- a/tests/__snapshots__/core-sql.test.ts.snap +++ b/tests/__snapshots__/core-sql.test.ts.snap @@ -1,5 +1,23 @@ // Vitest Snapshot v1, https://vitest.dev/guide/snapshot.html +exports[`core SQL (clickhouse) > and-or 1`] = ` +{ + "parameters": [], + "sql": "WITH orders AS ( + +), +customers AS ( + +) +SELECT + orders.Id AS id + FROM orders + WHERE orders.OrgId = 'org_1' + AND ((orders.Status = 'paid' AND orders.Amount > 15) OR orders.Customer = 'globex') + ORDER BY id ASC", +} +`; + exports[`core SQL (clickhouse) > arithmetic 1`] = ` { "parameters": [], @@ -121,6 +139,41 @@ SELECT } `; +exports[`core SQL (clickhouse) > distinct 1`] = ` +{ + "parameters": [], + "sql": "WITH orders AS ( + +), +customers AS ( + +) +SELECT DISTINCT ON (customer) + orders.Customer AS customer, + orders.Id AS id + FROM orders + WHERE orders.OrgId = 'org_1' + ORDER BY customer ASC, id DESC", +} +`; + +exports[`core SQL (clickhouse) > distinct-plain 1`] = ` +{ + "parameters": [], + "sql": "WITH orders AS ( + +), +customers AS ( + +) +SELECT DISTINCT + orders.Status AS status + FROM orders + WHERE orders.OrgId = 'org_1' + ORDER BY status ASC", +} +`; + exports[`core SQL (clickhouse) > format 1`] = ` { "parameters": [], @@ -206,6 +259,24 @@ SELECT } `; +exports[`core SQL (clickhouse) > is-not-null 1`] = ` +{ + "parameters": [], + "sql": "WITH orders AS ( + +), +customers AS ( + +) +SELECT + orders.Id AS id + FROM orders + WHERE orders.OrgId = 'org_1' + AND orders.Note IS NOT NULL + ORDER BY id ASC", +} +`; + exports[`core SQL (clickhouse) > join-subqueries 1`] = ` { "parameters": [], @@ -288,6 +359,26 @@ SELECT } `; +exports[`core SQL (clickhouse) > null-and-range 1`] = ` +{ + "parameters": [], + "sql": "WITH orders AS ( + +), +customers AS ( + +) +SELECT + orders.Id AS id + FROM orders + WHERE orders.OrgId = 'org_1' + AND orders.Note IS NULL + AND orders.Amount BETWEEN 5 AND 20 + AND orders.Id NOT BETWEEN 3 AND 3 + ORDER BY id ASC", +} +`; + exports[`core SQL (clickhouse) > nulls-and-literals 1`] = ` { "parameters": [], @@ -488,6 +579,26 @@ FORMAT JSON", } `; +exports[`core SQL (postgres) > and-or 1`] = ` +{ + "parameters": [ + "org_1", + ], + "sql": "WITH "orders" AS ( + +), +"customers" AS ( + +) +SELECT + "orders"."Id" AS "id" + FROM "orders" + WHERE "orders"."OrgId" = $1 + AND (("orders"."Status" = 'paid' AND "orders"."Amount" > 15) OR "orders"."Customer" = 'globex') + ORDER BY "id" ASC", +} +`; + exports[`core SQL (postgres) > arithmetic 1`] = ` { "parameters": [ @@ -619,6 +730,45 @@ SELECT } `; +exports[`core SQL (postgres) > distinct 1`] = ` +{ + "parameters": [ + "org_1", + ], + "sql": "WITH "orders" AS ( + +), +"customers" AS ( + +) +SELECT DISTINCT ON ("customer") + "orders"."Customer" AS "customer", + "orders"."Id" AS "id" + FROM "orders" + WHERE "orders"."OrgId" = $1 + ORDER BY "customer" ASC, "id" DESC", +} +`; + +exports[`core SQL (postgres) > distinct-plain 1`] = ` +{ + "parameters": [ + "org_1", + ], + "sql": "WITH "orders" AS ( + +), +"customers" AS ( + +) +SELECT DISTINCT + "orders"."Status" AS "status" + FROM "orders" + WHERE "orders"."OrgId" = $1 + ORDER BY "status" ASC", +} +`; + exports[`core SQL (postgres) > from-subquery 1`] = ` { "parameters": [ @@ -692,6 +842,26 @@ SELECT } `; +exports[`core SQL (postgres) > is-not-null 1`] = ` +{ + "parameters": [ + "org_1", + ], + "sql": "WITH "orders" AS ( + +), +"customers" AS ( + +) +SELECT + "orders"."Id" AS "id" + FROM "orders" + WHERE "orders"."OrgId" = $1 + AND "orders"."Note" IS NOT NULL + ORDER BY "id" ASC", +} +`; + exports[`core SQL (postgres) > join-subqueries 1`] = ` { "parameters": [ @@ -780,6 +950,49 @@ SELECT } `; +exports[`core SQL (postgres) > locking 1`] = ` +{ + "parameters": [ + "org_1", + ], + "sql": "WITH "orders" AS ( + +), +"customers" AS ( + +) +SELECT + "orders"."Id" AS "id" + FROM "orders" + WHERE "orders"."OrgId" = $1 + AND "orders"."Id" = 1 + FOR UPDATE SKIP LOCKED", +} +`; + +exports[`core SQL (postgres) > null-and-range 1`] = ` +{ + "parameters": [ + "org_1", + 20, + ], + "sql": "WITH "orders" AS ( + +), +"customers" AS ( + +) +SELECT + "orders"."Id" AS "id" + FROM "orders" + WHERE "orders"."OrgId" = $1 + AND "orders"."Note" IS NULL + AND "orders"."Amount" BETWEEN 5 AND $2 + AND "orders"."Id" NOT BETWEEN 3 AND 3 + ORDER BY "id" ASC", +} +`; + exports[`core SQL (postgres) > nulls-and-literals 1`] = ` { "parameters": [ diff --git a/tests/core-cases.ts b/tests/core-cases.ts index 56bf88e..043d250 100644 --- a/tests/core-cases.ts +++ b/tests/core-cases.ts @@ -611,6 +611,84 @@ export const coreCases: readonly CoreCase[] = [ }, expected: [{ id: 1 }, { id: 2 }], }, + // Null and range predicates, variadic and/or, distinct, row locking. + { + id: "null-and-range", + covers: e("isNull", "between", "notBetween"), + build: (ctx) => + ctx.compile( + orgOrders(ctx) + .select(($) => ({ id: $.Id })) + .where(($) => [ + $.OrgId.eq(CH.param.string("orgId")), + $.Note.isNull(), + $.Amount.between(5, CH.param.int("hi")), + $.Id.notBetween(3, 3), + ]) + .orderBy(["id", "asc"]), + { ...org, hi: 20 }, + ), + expected: [{ id: 2 }], + }, + { + id: "is-not-null", + covers: e("isNotNull"), + build: (ctx) => + ctx.compile(orgOrders(ctx).select(($) => ({ id: $.Id })).where(($) => [$.OrgId.eq(CH.param.string("orgId")), $.Note.isNotNull()]).orderBy(["id", "asc"]), org), + expected: [{ id: 1 }, { id: 4 }], + }, + { + id: "and-or", + covers: ["function:and", "function:or"], + build: (ctx) => + ctx.compile( + orgOrders(ctx) + .select(($) => ({ id: $.Id })) + .where(($) => [ + $.OrgId.eq(CH.param.string("orgId")), + CH.or(CH.and($.Status.eq("paid"), $.Amount.gt(15)), undefined, $.Customer.eq("globex")), + ]) + .orderBy(["id", "asc"]), + org, + ), + expected: [{ id: 2 }, { id: 3 }], + }, + { + id: "distinct", + covers: q("distinct", "distinctOn"), + build: (ctx) => + ctx.compile( + orgOrders(ctx) + .select(($) => ({ customer: $.Customer, id: $.Id })) + .distinctOn("customer") + .orderBy(["customer", "asc"], ["id", "desc"]), + org, + ), + expected: [ + { customer: "acme", id: 2 }, + { customer: "globex", id: 3 }, + { customer: "initech", id: 4 }, + ], + }, + { + id: "distinct-plain", + covers: q("distinct"), + build: (ctx) => + ctx.compile(orgOrders(ctx).select(($) => ({ status: $.Status })).distinct().orderBy(["status", "asc"]), org), + expected: [{ status: "open" }, { status: "paid" }, { status: "void" }], + }, + { + id: "locking", + covers: q("forUpdate", "forNoKeyUpdate", "forShare", "forKeyShare"), + rejects: { clickhouse: /no row locks/ }, + build: (ctx) => { + // Each strength compiles; the one sent is FOR UPDATE SKIP LOCKED. + const base = orgOrders(ctx).select(($) => ({ id: $.Id })).where(($) => [$.OrgId.eq(CH.param.string("orgId")), $.Id.eq(1)]) + for (const locked of [base.forNoKeyUpdate({ noWait: true }), base.forShare(), base.forKeyShare()]) ctx.compile(locked, org) + return ctx.compile(base.forUpdate({ skipLocked: true }), org) + }, + expected: [{ id: 1 }], + }, ] /** Cases a dialect cannot run yet, each with the reason. Empty is the goal. */ diff --git a/tests/database.clickhouse.test.ts b/tests/database.clickhouse.test.ts index 68ba68f..43b6fd7 100644 --- a/tests/database.clickhouse.test.ts +++ b/tests/database.clickhouse.test.ts @@ -165,6 +165,39 @@ describe("database", () => { ]) }) + it("reads with DISTINCT, DISTINCT ON, IS NULL and BETWEEN", async () => { + const result = await Effect.runPromise( + withDatabase((db) => + Effect.gen(function* () { + yield* db.execute(Db.sql`CREATE TABLE ev (OrgId String, Id UInt32, Note Nullable(String)) ENGINE = MergeTree ORDER BY (OrgId, Id)`) + const Ev = CH.table("ev", { OrgId: CH.string, Id: CH.uint32, Note: CH.nullable(CH.string) }) + yield* db.run( + CH.insertInto(Ev).values([ + { OrgId: "o", Id: 1, Note: null }, + { OrgId: "o", Id: 2, Note: "n" }, + { OrgId: "p", Id: 3, Note: null }, + ]), + ) + const orgs = yield* db.run(CH.from(Ev).select("OrgId").distinct().orderBy(["OrgId", "asc"])) + const last = yield* db.run( + CH.from(Ev).select(($) => ({ OrgId: $.OrgId, Id: $.Id })).distinctOn("OrgId").orderBy(["OrgId", "asc"], ["Id", "desc"]), + ) + const nulls = yield* db.run( + CH.from(Ev).select("Id").where(($) => [$.Note.isNull(), $.Id.between(1, CH.param.int("hi"))]).orderBy(["Id", "asc"]), + { hi: 3 }, + ) + return { orgs, last, nulls } + }), + ), + ) + expect(result.orgs).toEqual([{ OrgId: "o" }, { OrgId: "p" }]) + expect(result.last).toEqual([ + { OrgId: "o", Id: 2 }, + { OrgId: "p", Id: 3 }, + ]) + expect(result.nulls).toEqual([{ Id: 1 }, { Id: 3 }]) + }) + it("refuses a transaction before sending anything", async () => { const result = await Effect.runPromise( withDatabase((db, sent) => diff --git a/tests/dialect-cases.ts b/tests/dialect-cases.ts index be184d7..9a6311c 100644 --- a/tests/dialect-cases.ts +++ b/tests/dialect-cases.ts @@ -211,6 +211,22 @@ export const dialectCases: readonly DialectCase[] = [ joinedList: "a,b", }, ), + { + id: "distinct-and-distinct-on", + covers: ["query:distinct", "query:distinctOn"], + build: () => + CH.compileUnsafe( + numbers() + .select(($) => ({ odd: CH.intDiv($.n, l(2)), n: $.n })) + .distinctOn("odd") + .orderBy(["odd", "asc"], ["n", "desc"]), + {}, + ), + expected: [ + { odd: 0, n: 1 }, + { odd: 1, n: 3 }, + ], + }, { id: "array-join-and-membership", covers: fn("arrayJoin", "has"), diff --git a/tests/dialect-coverage.test.ts b/tests/dialect-coverage.test.ts index ebf23f3..2a85aa4 100644 --- a/tests/dialect-coverage.test.ts +++ b/tests/dialect-coverage.test.ts @@ -31,6 +31,10 @@ const exemptions = { "Caller-defined SQL types and codecs; no finite dialect contract. Public tarball smoke exercises the factory.", "type:untyped": "Explicitly unvalidated escape hatch; no decoding guarantee. Public tarball smoke exercises the factory.", + "query:forUpdate": "Postgres row lock; ClickHouse refuses it at compile (core case `locking`), so there is no live ClickHouse run.", + "query:forNoKeyUpdate": "Postgres row lock; ClickHouse refuses it at compile (core case `locking`), so there is no live ClickHouse run.", + "query:forShare": "Postgres row lock; ClickHouse refuses it at compile (core case `locking`), so there is no live ClickHouse run.", + "query:forKeyShare": "Postgres row lock; ClickHouse refuses it at compile (core case `locking`), so there is no live ClickHouse run.", } export const dialectInventory = [ @@ -106,7 +110,7 @@ export const coreInventory = [ ...functionsOf(CH.lit(1), "expr"), ...functionsOf(CH.lit(1).eq(1), "condition"), ...Object.keys(CH.param).map((name) => `param:${name}`), - ...["from", "fromQuery", "unionAll", "lit", "not"].map((name) => `function:${name}`), + ...["from", "fromQuery", "unionAll", "lit", "not", "and", "or"].map((name) => `function:${name}`), ]), ].sort() From 63b3bc7c1448f4dfdde5060c6a564794a2b05ea4 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 15:54:52 +0200 Subject: [PATCH 2/2] Refuse lock forms Postgres rejects, and distinctOn without keys From review: a qualified name in a lock's OF (Postgres takes only unqualified names), a lock on a query with DISTINCT, GROUP BY or HAVING, and a lock on a unionAll branch were all compiled and then refused by the server; they are now QueryBuilderDefects. distinctOn() with no keys silently became a whole-row DISTINCT; it is now a defect too. Co-Authored-By: Claude Opus 5.5 --- docs/queries.md | 8 +++++--- src/ch/compile.ts | 30 +++++++++++++++++++++++++++--- src/ch/predicates.test.ts | 26 +++++++++++++++++++++----- src/ch/query.ts | 5 ++++- 4 files changed, 57 insertions(+), 12 deletions(-) diff --git a/docs/queries.md b/docs/queries.md index 39cff40..6b56b87 100644 --- a/docs/queries.md +++ b/docs/queries.md @@ -156,9 +156,11 @@ CH.from(Jobs) // … LIMIT 1 FOR UPDATE SKIP LOCKED ``` -A lock lasts until the transaction ends, so run the query inside `Database.transaction`. -`skipLocked` and `noWait` together, and any lock on ClickHouse (which has no row locks), are a -`QueryBuilderDefect`. +A lock lasts until the transaction ends, so run the query inside `Database.transaction`. These +are a `QueryBuilderDefect`, refused before anything is sent: `skipLocked` and `noWait` together; +a qualified name in `of` (use the alias or `jobs`, not `public.jobs`); a lock on a query with +DISTINCT, GROUP BY or HAVING, or on a `unionAll` branch, which Postgres refuses; and any lock on +ClickHouse, which has no row locks. ## `withCTE` diff --git a/src/ch/compile.ts b/src/ch/compile.ts index 817f561..cde91af 100644 --- a/src/ch/compile.ts +++ b/src/ch/compile.ts @@ -61,7 +61,15 @@ const lockClause = (lock: import("./query").LockClause | undefined): string | un if (lock.skipLocked === true && lock.noWait === true) { throw new QueryBuilderDefect({ message: "CHQuery: a lock takes skipLocked or noWait, not both" }) } - const of = lock.of !== undefined && lock.of.length > 0 ? ` OF ${lock.of.map(quoteIdentPath).join(", ")}` : "" + // Postgres takes only unqualified names here, so `public.jobs` is refused, not quoted. + for (const name of lock.of ?? []) { + if (name.includes(".")) { + throw new QueryBuilderDefect({ + message: `CHQuery: FOR ${lock.strength} OF ${JSON.stringify(name)}: name the table by its alias or unqualified name`, + }) + } + } + const of = lock.of !== undefined && lock.of.length > 0 ? ` OF ${lock.of.map(quoteIdent).join(", ")}` : "" const wait = lock.skipLocked === true ? " SKIP LOCKED" : lock.noWait === true ? " NOWAIT" : "" return `FOR ${lock.strength}${of}${wait}` } @@ -848,14 +856,27 @@ function compileInner< const sqlQuery: SqlQuery = { distinct: state.distinct !== undefined, distinctOn: Array.isArray(state.distinct) - ? state.distinct.map((key) => { + ? (state.distinct.length === 0 + ? (() => { + throw new QueryBuilderDefect({ message: "CHQuery: distinctOn() needs at least one key" }) + })() + : state.distinct + ).map((key: string) => { if (!(options?.selectKeys ?? keys).includes(key)) { throw new QueryBuilderDefect({ message: `CHQuery: distinctOn(${JSON.stringify(key)}) is not a selected alias` }) } return raw(quoteIdent(key)) }) : undefined, - lock: lockClause(state.lock), + lock: (() => { + // Postgres refuses a lock on rows that are no longer table rows; say so here. + if (state.lock !== undefined && (state.distinct !== undefined || state.groupByKeys.length > 0 || state.havingFn !== undefined)) { + throw new QueryBuilderDefect({ + message: `CHQuery: FOR ${state.lock.strength} cannot lock rows of a query with DISTINCT, GROUP BY or HAVING`, + }) + } + return lockClause(state.lock) + })(), select: selectFragments, from: fromFragment, joins, @@ -1166,6 +1187,9 @@ function compileUnionInner, Params extends Re const first = state.queries[0] if (first === undefined) throw new QueryBuilderDefect({ message: "unionAll requires at least one query" }) const selectKeys = Object.keys(selectExprsOf(first) ?? {}) + if (state.queries.some((q) => q._state.lock !== undefined)) { + throw new QueryBuilderDefect({ message: "unionAll: a branch cannot take a row lock; lock in a query over the union instead" }) + } const subQueries = state.queries.map((q) => compileInner(q, params, { skipFormat: true, deferParams, nested: true, selectKeys, enclosingCtes }), ) diff --git a/src/ch/predicates.test.ts b/src/ch/predicates.test.ts index a08533f..69beb94 100644 --- a/src/ch/predicates.test.ts +++ b/src/ch/predicates.test.ts @@ -81,6 +81,16 @@ describe("distinct", () => { }) }) +describe("distinctOn without keys", () => { + it.effect("is a defect rather than a whole-row DISTINCT", () => + Effect.gen(function* () { + const keys: ReadonlyArray<"org"> = [] + const q = CH.from(Jobs).select("org").distinctOn(...(keys as unknown as ["org"])) + expect(failure(yield* Effect.exit(PG.compile(q)))).toBeInstanceOf(QueryBuilderDefect) + }), + ) +}) + describe("row locking", () => { it("FOR UPDATE / NO KEY UPDATE / SHARE / KEY SHARE with OF, SKIP LOCKED and NOWAIT, after LIMIT", () => { const base = CH.from(Jobs).select("id").where(($) => [$.state.eq("queued")]).limit(1) @@ -93,13 +103,19 @@ describe("row locking", () => { expect(tail(base.forShare().forUpdate())).toBe("FOR UPDATE") }) - it.effect("is a defect on ClickHouse and with skipLocked and noWait together", () => + it.effect("is a defect on ClickHouse, with skipLocked and noWait together, and where Postgres refuses a lock", () => Effect.gen(function* () { const base = CH.from(Jobs).select("id") - expect(failure(yield* Effect.exit(CH.compile(base.forUpdate(), {})))).toBeInstanceOf(QueryBuilderDefect) - expect(failure(yield* Effect.exit(PG.compile(base.forUpdate({ skipLocked: true, noWait: true }), {})))).toBeInstanceOf( - QueryBuilderDefect, - ) + const defect = function* (q: Effect.Effect) { + expect(failure(yield* Effect.exit(q))).toBeInstanceOf(QueryBuilderDefect) + } + yield* defect(CH.compile(base.forUpdate(), {})) + yield* defect(PG.compile(base.forUpdate({ skipLocked: true, noWait: true }), {})) + yield* defect(PG.compile(base.forUpdate({ of: ["public.jobs"] }), {})) + yield* defect(PG.compile(base.distinct().forUpdate(), {})) + yield* defect(PG.compile(CH.from(Jobs).select("org").groupBy("org").forShare(), {})) + yield* defect(PG.compileUnion(CH.unionAll(base.forUpdate(), base), {})) + expect(PG.compileUnsafe(base.forUpdate({ of: ["jobs"] })).sql).toMatch(/FOR UPDATE OF "jobs"$/) }), ) }) diff --git a/src/ch/query.ts b/src/ch/query.ts index 351a77e..e5dc101 100644 --- a/src/ch/query.ts +++ b/src/ch/query.ts @@ -67,7 +67,10 @@ export interface LockOptions { readonly skipLocked?: boolean /** `NOWAIT`: fail at once instead of waiting. Not with `skipLocked`. */ readonly noWait?: boolean - /** `OF alias, ...`: lock only these tables' rows (the FROM alias or table name, or join aliases). */ + /** + * `OF name, ...`: lock only these tables' rows. Each is an alias or an + * unqualified table name (`jobs`, not `public.jobs`, which Postgres refuses). + */ readonly of?: ReadonlyArray }