From fd2b270b0d63cc7f1c9f248498bb126fc6034d7f Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 15:26:31 +0200 Subject: [PATCH 1/2] Add update() and deleteFrom() for Postgres and ClickHouse update(table).set(...).where(...) and deleteFrom(table).where(...), with returning on Postgres and settings on ClickHouse. They share the insert's value encoding, SET record, RETURNING and settings, now factored into valueCells, setAssignments, returningOf and writeSettingsClause. A write with no where() is a defect unless allRows() says so; a where() whose conditions all came out undefined fails, since that comes from data and would widen a filtered write to every row. ClickHouse compiles UPDATE to an ALTER TABLE ... UPDATE mutation and DELETE to a lightweight DELETE, both with WHERE 1 for allRows() and settings last; checked on 26.2 and 26.8. Tenant scope is derived from the WHERE, and an update that moves rows to another tenant is cross-tenant. CompiledQuery.kind gains update and delete, and Database.run sends any write without RETURNING through command. DialectClauses.insertSettings (unreleased) becomes writeSettings; alterTableUpdate is new. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 9 +- README.md | 1 + design/gap-review.md | 4 +- design/writes.md | 24 ++- docs/README.md | 3 +- docs/database.md | 14 +- docs/reference.md | 4 +- docs/updates-and-deletes.md | 83 +++++++++ src/ch/compile.ts | 291 ++++++++++++++++++++++-------- src/ch/dialect.ts | 12 +- src/ch/index.ts | 12 ++ src/ch/insert.test-d.ts | 18 ++ src/ch/update.test.ts | 105 +++++++++++ src/ch/update.ts | 156 ++++++++++++++++ src/database/database.test.ts | 29 +++ src/database/database.ts | 19 +- tests/database.clickhouse.test.ts | 31 ++++ 17 files changed, 717 insertions(+), 98 deletions(-) create mode 100644 docs/updates-and-deletes.md create mode 100644 src/ch/update.test.ts create mode 100644 src/ch/update.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index d9fe1cc..c99c1e1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,13 @@ ## Unreleased +- 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 + with no `where()` is refused unless `allRows()` says so, and one whose conditions all came out + `undefined` fails. Tenant scope is derived from the WHERE. `CompiledQuery.kind` gains `update` + and `delete`; `DialectClauses.insertSettings` is renamed `writeSettings` (unreleased), and + `DialectClauses.alterTableUpdate` is added. - `returning()` with no arguments returns every column, as in Drizzle. `insertInto(table)` now returns `CHInsertStart`, which offers only `values` and `select`, so an insert without rows no longer type-checks. Add `TableOptions.computed` for generated columns. `INSERT ... SELECT` @@ -24,7 +31,7 @@ against the table at the type level. Tenant scope follows the read and where the written tenant comes from. - Add `settings(record)` to an insert (ClickHouse): `INSERT ... SETTINGS name = value`. Add - `DialectClauses.insertSettings`, optional. + `DialectClauses.writeSettings`, optional. - Add `CompiledQuery.kind` (`"select"` or `"insert"`); `rawCompiledQuery` takes it as an option. - Add `ParamStyle.maxParameters`; Postgres sets 65535, and a statement over it fails to compile. - Add `@maple-dev/effect-orm/database`, opt-in (see `docs/database.md`): a `Database` over the diff --git a/README.md b/README.md index 5b316f8..1b428a5 100644 --- a/README.md +++ b/README.md @@ -150,6 +150,7 @@ Full guides live in [`docs/`](./docs/README.md): | [Joins and subqueries](./docs/joins-and-subqueries.md) | The join family, `fromQuery`, correlated subqueries | | [Unions and CTEs](./docs/unions-and-ctes.md) | `unionAll`, `fromUnion`, `withCTE` | | [Inserting rows](./docs/inserts.md) | `insertInto`, the insert row type, `DEFAULT`, binding | +| [Updating and deleting](./docs/updates-and-deletes.md) | `update`, `deleteFrom`, `allRows`, ClickHouse mutations | | [Params and compilation](./docs/params-and-compilation.md) | `param.*`, how values reach the SQL, `CompiledQuery` | | [Decoding results](./docs/decoding-results.md) | `rowSchema`, `decodeRows`, decode errors | | [Running a query](./docs/running-queries.md) | Executing the SQL with a real client, wire settings, `SETTINGS` | diff --git a/design/gap-review.md b/design/gap-review.md index dea7823..b15345e 100644 --- a/design/gap-review.md +++ b/design/gap-review.md @@ -24,8 +24,8 @@ builder; **P1** commonly used; **P2** niche. | Gap | Maple | Effort | | --- | --- | --- | -| UPDATE builder: SET values and expressions, WHERE, RETURNING | ~120 | M | -| DELETE builder: WHERE, RETURNING | ~79 | S | +| ~~UPDATE builder: SET values and expressions, WHERE, RETURNING~~ (built) | ~120 | M | +| ~~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 | diff --git a/design/writes.md b/design/writes.md index 4ca6381..abef65b 100644 --- a/design/writes.md +++ b/design/writes.md @@ -1,7 +1,8 @@ # Writes: INSERT Status: phases 1 to 4 built (`values`, `returning`, `onConflictDoNothing` / -`onConflictDoUpdate`, `select`, `settings`; see `docs/inserts.md`). `encodeInsertRows` (§7) is +`onConflictDoUpdate`, `select`, `settings`; see `docs/inserts.md`), and UPDATE and DELETE (§12, +`docs/updates-and-deletes.md`). `encodeInsertRows` (§7) is not built: no consumer has asked for it. Section 11 lists where the build differs from the plan. UPDATE and DELETE come later and will reuse what this note sets up (the write-statement state, the `RETURNING` path, value encoding). @@ -257,4 +258,23 @@ note, reusing `kind: "write"`, the returning path and the `set` record type from Its tenant scope is the SELECT's for an untenanted target (the read), and for a tenant target single-tenant only when the read is and each row takes its tenant from a source tenant column or the same param. -- **Settings are a dialect clause**, `DialectClauses.insertSettings`, like the others. +- **Settings are a dialect clause**, `DialectClauses.writeSettings`, like the others. + +## 12. UPDATE and DELETE + +`update(table).set(...).where(...)` and `deleteFrom(table).where(...)`, in `src/ch/update.ts`, +compiled beside INSERT and sharing its value encoding (`valueCells`), SET record +(`setAssignments`, also used by `onConflictDoUpdate`), RETURNING (`returningOf`) and settings. + +- **No WHERE is refused.** No `where()` is a defect; a `where()` whose conditions are all + `undefined` is a failure, since optional filters come from data and that case would widen the + write to every row. `allRows()` opts in. +- **ClickHouse**: UPDATE is an `ALTER TABLE ... UPDATE` mutation (`DialectClauses.alterTableUpdate`), + the one form every supported server takes; DELETE is the lightweight `DELETE FROM`. Both need + a WHERE, so `allRows()` writes `WHERE 1`. Settings go last (`mutations_sync`, + `lightweight_deletes_sync`), and `DialectClauses.insertSettings` became `writeSettings`. + Checked on 26.2.19.43 and 26.8.2.7. +- **Tenant scope** is derived from the WHERE as for a query over the table; an UPDATE that sets + the tenant column to anything but the pinned value is cross-tenant. +- Not built: UPDATE ... FROM / joins, DELETE USING, ORDER BY / LIMIT, a VALUES source for bulk + updates. See `design/gap-review.md`. diff --git a/docs/README.md b/docs/README.md index 1915615..6d246ba 100644 --- a/docs/README.md +++ b/docs/README.md @@ -16,7 +16,7 @@ You do not need a Maple account, Maple's schema, or tenant columns. Tenant analy optional feature for applications that share tables between tenants. The root builder does not manage connections, create tables, or run migrations. It builds -SELECTs and [INSERTs](./inserts.md); UPDATE and DELETE are not built yet. Opt-in [schema and migration entry points](./migrations.md) add DDL and migrations for +SELECTs, [INSERTs](./inserts.md), and [UPDATEs and DELETEs](./updates-and-deletes.md). Opt-in [schema and migration entry points](./migrations.md) add DDL and migrations for ClickHouse. It does not validate SQL against a live server, choose query plans, enforce authorization, or supply retries. Existing ClickHouse tables and your executor own those responsibilities. [Getting started](./getting-started.md) covers npm installation and building from source. @@ -47,6 +47,7 @@ Roughly in reading order. | [Joins and subqueries](./joins-and-subqueries.md) | The join family, `fromQuery`, correlated subqueries | | [Unions and CTEs](./unions-and-ctes.md) | `unionAll`, `fromUnion`, `withCTE` | | [Inserting rows](./inserts.md) | `insertInto`, the insert row type, `DEFAULT`, binding | +| [Updating and deleting](./updates-and-deletes.md) | `update`, `deleteFrom`, `allRows`, ClickHouse mutations | | [Params and compilation](./params-and-compilation.md) | `param.*`, how values reach the SQL, `CompiledQuery` | | [Decoding results](./decoding-results.md) | `rowSchema`, `decodeRows`, `decodeFirstRow`, decode errors | | [Running a query](./running-queries.md) | Executing the SQL with a real client, wire settings, `SETTINGS` | diff --git a/docs/database.md b/docs/database.md index 573f919..7e6b42f 100644 --- a/docs/database.md +++ b/docs/database.md @@ -99,12 +99,13 @@ Calling `withdraw(1, 30)` outside `transfer` does not compile: `requireTransacti ### Queries and statements -`run` takes the query you built, a `unionAll`, an `insertInto`, or a query compiled elsewhere. It compiles with +`run` takes the query you built, a `unionAll`, an `insertInto`, `update` or `deleteFrom`, or a +query compiled elsewhere. It compiles with the database's dialect, so you never pick a `compile`; `params` fills the query's `param.*` markers, and a missing one fails with `QueryBuilderError`. A query compiled elsewhere must have been compiled for the same dialect, or `run` dies: the root `compile` is ClickHouse's. -`sql` writes the statements the builder does not have yet (UPDATE, DELETE, DDL, advisory +`sql` writes the statements the builder does not have yet (DDL, bulk `UPDATE ... FROM`, advisory locks). Each `${value}` is bound, as `$1, $2, ...` on Postgres and as an escaped literal on ClickHouse, so nothing in a value becomes SQL. A `sql` inside another is spliced, so statements compose. Names go through `sql.identifier`, which accepts only plain identifiers @@ -247,7 +248,8 @@ fails at BEGIN today: through the query path with a syntax error, and through `a ## Writes -The builder compiles SELECTs and [INSERTs](./inserts.md). `run` runs an insert and returns its -`returning` rows, decoded, or none without `returning`; an insert without it goes through -`command`, as `execute` does. Write UPDATE and DELETE with `sql`, as above, and read `RETURNING` -with `query` and a schema. +The builder compiles SELECTs, [INSERTs](./inserts.md) and +[UPDATEs and DELETEs](./updates-and-deletes.md). `run` runs a write and returns its `returning` +rows, decoded, or none without `returning`; a write without it goes through `command`, as +`execute` does. Write other statements with `sql`, as above, and read `RETURNING` with `query` +and a schema. diff --git a/docs/reference.md b/docs/reference.md index 974aec6..c3cdd48 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -62,6 +62,8 @@ Note `/sql` exports a `compile` (fragment → string) distinct from the root `co | `fromQuery` | `(query, alias) => CHQuery` | | `fromUnion` | `(union, alias) => CHQuery` | | `unionAll` | `(...queries) => CHUnionQuery` | +| `update` | `(table) => CHUpdateStart`, then `CHUpdate`: `.set(record \| fn)`, `.where(fn)` or `.allRows()`, `.returning(...)`, `.settings(record)`. See [Updating and deleting](./updates-and-deletes.md) | +| `deleteFrom` | `(table) => CHDelete`: `.where(fn)` or `.allRows()`, `.returning(...)`, `.settings(record)` | | `insertInto` | `(table) => CHInsertStart`, then `CHInsert`; `.values(row \| rows)` or `.select(query)` sets its rows, `.settings(record)` ClickHouse `SETTINGS`, `.returning(...)` the RETURNING list, `.onConflictDoNothing(options?)` / `.onConflictDoUpdate(options)` the ON CONFLICT clause (Postgres). See [Inserting rows](./inserts.md) | ### `CHQuery` methods @@ -279,7 +281,7 @@ Types: `WindowSpec`, `CompiledWindowSpec`, `WindowFrameBound`, `WindowRowsFrame` **Everything else** — `Table`, `TableOptions`, `Expr`, `ColumnRef`, `Condition`, `Comparable` (what a value of a type may be compared against), `MapValueOf`, `Subquery`, `ParamMarker`, -`ParamKind`, `CHQuery`, `CHUnionQuery`, `CHInsert`, `CHInsertStart`, `InsertRow`, `InsertRowOf`, `InsertValue`, `InsertSelectMisfits`, `InsertSelectMissing`, `InsertSettingValue`, `ConflictTarget`, `ConflictSet`, `OnConflictDoNothing`, `OnConflictDoUpdate`, `ColumnAccessor`, `JoinedColumnAccessor`, +`ParamKind`, `CHQuery`, `CHUnionQuery`, `CHInsert`, `CHInsertStart`, `CHUpdate`, `CHUpdateStart`, `CHDelete`, `CHWrite`, `UpdateSet`, `UpdateSetOf`, `InsertRow`, `InsertRowOf`, `InsertValue`, `InsertSelectMisfits`, `InsertSelectMissing`, `InsertSettingValue`, `ConflictTarget`, `ConflictSet`, `OnConflictDoNothing`, `OnConflictDoUpdate`, `ColumnAccessor`, `JoinedColumnAccessor`, `JoinOnCallback`, `CompiledQuery`, `CompiledQueryInput`, `CompiledQueryRowSchema`, `RowSchemaMismatch`, `TenantScope`, `Dialect`, `DialectClauses`, `DialectTransactions`, `IsolationLevel`, `TransactionSettings`, `ParamStyle`, `FnResult`, `WindowFunnelMode`, `WindowSpec`, `WindowRowsFrame`, `WindowFrameBound`, `WindowOrderDirection`, `CompiledWindowSpec`. diff --git a/docs/updates-and-deletes.md b/docs/updates-and-deletes.md new file mode 100644 index 0000000..a8f052e --- /dev/null +++ b/docs/updates-and-deletes.md @@ -0,0 +1,83 @@ +# Updating and deleting rows + +`update(table)` and `deleteFrom(table)` build UPDATE and DELETE from the same table definitions, +like [`insertInto`](./inserts.md). They are immutable values; `Database.run` compiles them for +its database's dialect and runs them. + +```ts +import * as CH from "@maple-dev/effect-orm" +import * as PG from "@maple-dev/effect-orm/postgres" + +const Tickets = CH.table( + "tickets", + { id: PG.int4, org: PG.text, seats: PG.int4, tags: PG.array(PG.text) }, + { tenantColumn: "org" }, +) + +const bump = CH.update(Tickets) + .set(($) => ({ seats: $.seats.add(1) })) + .where(($) => [$.org.eq(CH.param.string("org")), $.seats.lt(5)]) + .returning("id", "seats") + +const revoke = CH.deleteFrom(Tickets).where(($) => [$.id.eq(CH.param.int("id"))]) + +// yield* Db.run(bump, { org }) // [{ id, seats }] +// yield* Db.run(revoke, { id }) // [] +``` + +## SET + +`set` takes a record of values, params or expressions, or a callback that gets the row's +columns as `$`. A key left out (or `undefined`) keeps the existing value. Values are encoded +through the column's codec and bound on Postgres, as in an insert. A computed column (see +[`computed`](./inserts.md#which-columns-have-defaults)) cannot be set. `UpdateSetOf` +names the record type. + +`update(table)` offers only `set` until it has one (its type is `CHUpdateStart`). + +## WHERE, and writing every row + +`where` works as in a query: a list of conditions, AND-joined, with an `undefined` one skipped, +so optional filters compose. A write with no `where` would change every row, so: + +- compiling an UPDATE or DELETE with no `where()` is a `QueryBuilderDefect`; +- a `where()` whose conditions all came out `undefined` is a `QueryBuilderError`, because that + happens with data (every optional filter absent) and would otherwise widen a filtered write to + the whole table; +- `allRows()` says a write over every row is meant. + +## RETURNING + +On Postgres, `returning` works as on an insert: no arguments for every column, column names, or +a callback. `Database.run` returns the changed or deleted rows, decoded. Without it, a write +returns no rows. + +## ClickHouse + +ClickHouse has no `UPDATE ... SET` on every server, so `update` compiles to a mutation, and +`deleteFrom` to a lightweight delete. Both need a WHERE, so `allRows()` writes `WHERE 1`: + +```sql +ALTER TABLE spans UPDATE Name = 'x' +WHERE OrgId = 'o' SETTINGS mutations_sync = 2 + +DELETE FROM spans +WHERE OrgId = 'o' SETTINGS lightweight_deletes_sync = 2 +``` + +A mutation runs in the background: without `mutations_sync`, `Database.run` returns before the +rows change. Pass `settings({ mutations_sync: 2 })` when the caller reads its own write. +Mutations rewrite whole parts, so they suit occasional corrections, not per-request updates; +model frequently changing state with a `ReplacingMergeTree` and inserts instead. ClickHouse +cannot update a column of the sorting key. `returning` is refused on ClickHouse +(`QueryBuilderDefect`), and `settings` on Postgres. + +## Tenant scope + +An UPDATE or DELETE has the scope a query over the table with the same WHERE would have: +`"single-tenant"` when the WHERE pins the tenant column, `"cross-tenant"` otherwise (including +`allRows()`). An UPDATE that sets the tenant column to another value moves rows out of the +tenant, so it is `"cross-tenant"` too. + +_(Backed by `src/ch/update.test.ts`, `src/database/database.test.ts` and +`tests/database.clickhouse.test.ts`.)_ diff --git a/src/ch/compile.ts b/src/ch/compile.ts index 3d1d288..9c3d06a 100644 --- a/src/ch/compile.ts +++ b/src/ch/compile.ts @@ -11,6 +11,8 @@ import { custom, dateTime, dateTime64, type CHType, type ColumnDefs } from "./ty import type { CHQuery, CHQueryState } from "./query" import type { CHUnionQuery } from "./union" import { isInsert, type CHInsert } from "./insert" +import { isDelete, isUpdate, type CHDelete, type CHUpdate } from "./update" +import type { Table } from "./table" import { createColumnAccessor, createQualifiedColumnAccessor, createJoinedColumnAccessor, sourceAlias } from "./query" import { aliased, columnTypeOf, isExprLike, type Expr } from "./expr" import { raw, identPath, quoteIdent, quoteIdentPath, compile as compileSqlFragment } from "../sql/sql-fragment" @@ -133,9 +135,9 @@ interface CompiledQueryBase { * so an executor runs it the way it runs DDL (`Database.run` does), not * through a client path that expects a result set. */ - readonly kind: "select" | "insert" + readonly kind: "select" | "insert" | "update" | "delete" /** - * The aliases of an insert's RETURNING list, when it has one. An insert + * The aliases of a write's RETURNING list, when it has one. A write * without it sends back no rows; one with it is read like a query. */ readonly returning?: ReadonlyArray @@ -347,7 +349,7 @@ const makeCompiledQuery = ( rawSql?: { readonly reason: string; readonly justification: string }, rowSchemaMismatch?: RowSchemaMismatch, dialect?: string, - kind: "select" | "insert" = "select", + kind: CompiledQueryBase["kind"] = "select", returning?: ReadonlyArray, ): CompiledQuery => { let cachedDecodeRow: ((row: unknown) => Effect.Effect) | undefined @@ -473,8 +475,8 @@ export const rawCompiledQuery = < readonly route?: Route /** The `name` of the dialect the SQL is written for, so an executor can check it. */ readonly dialect?: string - /** `insert` for a write that returns no rows. Default `select`. */ - readonly kind?: "select" | "insert" + /** `insert`, `update` or `delete` for a write that returns no rows. Default `select`. */ + readonly kind?: "select" | "insert" | "update" | "delete" }): CompiledQuery => makeCompiledQuery( args.sql, @@ -541,21 +543,24 @@ export function compileCH< dialect?: Dialect }, ): Effect.Effect, QueryBuilderError> -/** An INSERT. `params` fills the `param.*` markers among its values. */ +/** An INSERT, UPDATE or DELETE. `params` fills the `param.*` markers among its values. */ export function compileCH( - insert: CHInsert, + insert: CHWrite, params?: Record, options?: InsertCompileOptions, ): Effect.Effect, QueryBuilderError> export function compileCH( - query: CHQuery | CHInsert, + query: CHQuery | CHWrite, params?: Record, options?: any, ): Effect.Effect, QueryBuilderError> { return asEffect(() => compileCHUnsafe(query as CHQuery, params ?? {}, options)) } -/** What compiling an INSERT takes: only the dialect. */ +/** A write statement: what `compile` takes besides a query. */ +export type CHWrite = CHInsert | CHUpdate | CHDelete + +/** What compiling a write takes: only the dialect. */ export interface InsertCompileOptions { readonly dialect?: Dialect } @@ -589,19 +594,23 @@ export function compileCHUnsafe< dialect?: Dialect }, ): CompiledQuery -/** An INSERT. `params` fills the `param.*` markers among its values. */ +/** An INSERT, UPDATE or DELETE. `params` fills the `param.*` markers among its values. */ export function compileCHUnsafe( - insert: CHInsert, + insert: CHWrite, params?: Record, options?: InsertCompileOptions, ): CompiledQuery export function compileCHUnsafe( - query: CHQuery | CHInsert, + query: CHQuery | CHWrite, params?: Record, options?: any, ): CompiledQuery { return withDialect(options?.dialect ?? currentDialect(), () => - isInsert(query) ? compileInsert(query, params ?? {}) : compileInner(query, params ?? {}, options), + isInsert(query) + ? compileInsert(query, params ?? {}) + : isUpdate(query) || isDelete(query) + ? compileUpdateOrDelete(query, params ?? {}) + : compileInner(query as CHQuery, params ?? {}, options), ) } @@ -1378,23 +1387,7 @@ const onConflictClause = ( const existing = createQualifiedColumnAccessor(table.name, undefined, table.columns) const excluded = createQualifiedColumnAccessor("excluded", undefined, table.columns) const set = typeof conflict.set === "function" ? conflict.set(existing, excluded) : conflict.set - const computed = new Set(table.computed ?? []) - // The SET record may come from data, so a bad key is a failure, as in a row. - const assignments = Object.entries(set as Record).flatMap(([column, value]) => { - if (value === undefined) return [] - if (!Object.hasOwn(table.columns, column) || computed.has(column)) { - throw new QueryBuilderError({ - code: "InvalidArguments", - message: `${where}: onConflictDoUpdate sets ${JSON.stringify(column)}, which is not an insertable column of the table`, - }) - } - const sql = cell(column, value, "onConflictDoUpdate set") - wrote(column, value, sql) - return [`${quoteIdent(column)} = ${sql}`] - }) - if (assignments.length === 0) { - throw new QueryBuilderError({ code: "InvalidArguments", message: `${where}: onConflictDoUpdate sets no columns` }) - } + const assignments = setAssignments(table, set as Record, cell, wrote, where, "onConflictDoUpdate") const condition = conflict.where?.(existing, excluded) return `\nON CONFLICT${target} DO UPDATE SET ${assignments.join(", ")}${ condition === undefined ? "" : ` WHERE ${compileSqlFragment(condition.toFragment())}` @@ -1402,21 +1395,101 @@ const onConflictClause = ( } /** The RETURNING clause and its row schema, or none without `returning`. */ -const returningOf = (insert: CHInsert) => { - const { table, returningFn } = insert._state +const returningOf = ( + table: Table, + returningFn: (($: any) => Record>) | undefined, + where: string, +) => { if (returningFn === undefined) return undefined const dialect = currentDialect() if (dialect.clauses.returning !== true) { throw new QueryBuilderDefect({ - message: `insertInto(${table.name}): returning() has no meaning for the ${dialect.name} dialect, which has no RETURNING clause`, + message: `${where}: returning() has no meaning for the ${dialect.name} dialect, which has no RETURNING clause`, }) } const exprs = returningFn(createColumnAccessor(table.columns)) const aliases = Object.keys(exprs) if (aliases.length === 0) { - throw new QueryBuilderDefect({ message: `insertInto(${table.name}): returning() needs at least one column` }) + throw new QueryBuilderDefect({ message: `${where}: returning() needs at least one column` }) } - return { exprs, aliases, derived: deriveRowSchema(exprs) } + const sql = `\nRETURNING ${aliases.map((alias) => compileSqlFragment(aliased(exprs[alias]!, alias))).join(", ")}` + return { exprs, aliases, sql, derived: deriveRowSchema(exprs) } +} + +/** A write's compiled query, with the row schema its RETURNING list derives. */ +const writeCompiledQuery = ( + rendered: { readonly sql: string; readonly parameters: ReadonlyArray }, + kind: "insert" | "update" | "delete", + tenantScope: TenantScope, + tenantBound: string | undefined, + returning: ReturnType, +): CompiledQuery => { + const returnedSchema = returning !== undefined && "schema" in returning.derived ? returning.derived.schema : undefined + return withTenantBound( + makeCompiledQuery( + rendered.sql, + rendered.parameters, + tenantScope, + returning === undefined || returnedSchema !== undefined ? "derived" : "none", + () => (returning === undefined ? EMPTY_ROW : returnedSchema), + undefined, + returning !== undefined && "untyped" in returning.derived ? returning.derived.untyped : [], + undefined, + undefined, + currentDialect().name, + kind, + returning?.aliases, + ), + tenantBound, + ) +} + +/** + * Literal values of a write as params of their column's type: encoded by the + * column's codec (so a failure names the column), then sent the way the + * dialect sends params. `values` is the params bag to render with. + */ +const valueCells = (table: Table, params: Record, where: string) => { + const values: Record = { ...params } + let next = 0 + const cell = (column: string, value: unknown, context: string): string => { + if (value === undefined) return "DEFAULT" + if (isExprLike(value)) return compileSqlFragment(value.toFragment()) + const wire = encodeValue(table.columns[column]!.literalSchema, value, `${where}: ${context}, column ${column}`) + let name = `${VALUE_PARAM}${next++}` + while (Object.hasOwn(params, name)) name = `${VALUE_PARAM}${next++}` + values[name] = wire + return compileSqlFragment(param.of(insertWireValue, name).toFragment()) + } + return { values, cell } +} + +/** `column = value, ...` for a SET record, which may come from data, so a bad key is a failure. */ +const setAssignments = ( + table: Table, + set: Record, + cell: (column: string, value: unknown, context: string) => string, + wrote: (column: string, value: unknown, sql: string) => void, + where: string, + context: string, +): ReadonlyArray => { + const computed = new Set(table.computed ?? []) + const assignments = Object.entries(set).flatMap(([column, value]) => { + if (value === undefined) return [] + if (!Object.hasOwn(table.columns, column) || computed.has(column)) { + throw new QueryBuilderError({ + code: "InvalidArguments", + message: `${where}: ${context} sets ${JSON.stringify(column)}, which is not an insertable column of the table`, + }) + } + const sql = cell(column, value, `${context} set`) + wrote(column, value, sql) + return [`${quoteIdent(column)} = ${sql}`] + }) + if (assignments.length === 0) { + throw new QueryBuilderError({ code: "InvalidArguments", message: `${where}: ${context} sets no columns` }) + } + return assignments } /** @@ -1451,20 +1524,23 @@ const insertSelectSource = ( } /** `SETTINGS a = 1, b = 'x'`, or `""` without settings. */ -const insertSettingsClause = (insert: CHInsert): string => { - const { table, settings } = insert._state +const writeSettingsClause = ( + table: Table, + settings: Readonly> | undefined, + where: string, +): string => { const entries = Object.entries(settings ?? {}) if (entries.length === 0) return "" const dialect = currentDialect() - if (dialect.clauses.insertSettings !== true) { + if (dialect.clauses.writeSettings !== true) { throw new QueryBuilderDefect({ - message: `insertInto(${table.name}): settings() has no meaning for the ${dialect.name} dialect, which has no INSERT SETTINGS`, + message: `${where}: settings() has no meaning for the ${dialect.name} dialect, which has no SETTINGS on writes`, }) } return ` SETTINGS ${entries .map(([name, value]) => { if (!/^[A-Za-z_][A-Za-z0-9_]*$/.test(name)) { - throw new QueryBuilderDefect({ message: `insertInto(${table.name}): ${JSON.stringify(name)} is not a setting name` }) + throw new QueryBuilderDefect({ message: `${where}: ${JSON.stringify(name)} is not a setting name` }) } return `${name} = ${checkedLiteral(dialect, value, `setting ${name}`)}` }) @@ -1484,17 +1560,7 @@ function compileInsert(insert: CHInsert, params: Record = { ...params } - let next = 0 - const cell = (column: string, value: unknown, context: string): string => { - if (value === undefined) return "DEFAULT" - if (isExprLike(value)) return compileSqlFragment(value.toFragment()) - const wire = encodeValue(table.columns[column]!.literalSchema, value, `${where}: ${context}, column ${column}`) - let name = `${VALUE_PARAM}${next++}` - while (Object.hasOwn(params, name)) name = `${VALUE_PARAM}${next++}` - values[name] = wire - return compileSqlFragment(param.of(insertWireValue, name).toFragment()) - } + const { values, cell } = valueCells(table, params, where) // A tenant table's insert is single-tenant when every row it writes (and an // upsert's SET, if it writes the column) pins the tenant column to the same @@ -1557,13 +1623,9 @@ function compileInsert(insert: CHInsert, params: Record compileSqlFragment(aliased(returning.exprs[alias]!, alias))).join(", ")}` + const returning = returningOf(table, insert._state.returningFn, where) const rendered = renderParams( - `INSERT INTO ${quoteIdentPath(table.name)} (${columns.map(quoteIdent).join(", ")})${insertSettingsClause(insert)}\n${source}${conflictSql}${returningSql}`, + `INSERT INTO ${quoteIdentPath(table.name)} (${columns.map(quoteIdent).join(", ")})${writeSettingsClause(table, insert._state.settings, where)}\n${source}${conflictSql}${returning?.sql ?? ""}`, values, dialect, ) @@ -1575,25 +1637,7 @@ function compileInsert(insert: CHInsert, params: Record( - rendered.sql, - rendered.parameters, - tenantScope, - returning === undefined || returnedSchema !== undefined ? "derived" : "none", - () => (returning === undefined ? EMPTY_ROW : returnedSchema), - undefined, - returning !== undefined && "untyped" in returning.derived ? returning.derived.untyped : [], - undefined, - undefined, - dialect.name, - "insert", - returning?.aliases, - ), - tenantBound, - ) + return writeCompiledQuery(rendered, "insert", tenantScope, tenantBound, returning) } /** The VALUES rows' columns, in table order, after checking every key. */ @@ -1632,3 +1676,98 @@ const valuesColumns = ( } return columns } + +// UPDATE and DELETE + +/** + * An UPDATE or DELETE. Postgres writes `UPDATE t SET ... WHERE ...` and + * `DELETE FROM t WHERE ...`; ClickHouse an `ALTER TABLE t UPDATE ... WHERE ...` + * mutation and a lightweight `DELETE FROM t WHERE ...`, both of which need a + * WHERE, so `allRows()` writes `WHERE 1` there. + */ +function compileUpdateOrDelete( + write: CHUpdate | CHDelete, + params: Record, +): CompiledQuery { + const state = write._state + const { table } = state + const kind = write._tag === "CHUpdate" ? "update" : "delete" + const where = `${kind === "update" ? "update" : "deleteFrom"}(${table.name})` + const dialect = currentDialect() + const { values, cell } = valueCells(table, params, where) + const $ = createColumnAccessor(table.columns, table.tenantColumn) + + // No WHERE at all is a mistake in the source; a WHERE whose conditions all + // came out undefined is data, and would otherwise turn a filtered write into + // one over every row. + if (state.whereFn === undefined && state.allRows !== true) { + throw new QueryBuilderDefect({ message: `${where}: no where(); call allRows() to write every row` }) + } + const conditions = (state.whereFn?.($) ?? []).filter((c): c is NonNullable => c != null) + if (state.whereFn !== undefined && conditions.length === 0 && state.allRows !== true) { + throw new QueryBuilderError({ + code: "InvalidArguments", + message: `${where}: every where() condition was undefined, which would write every row; call allRows() if that is meant`, + }) + } + + const tenant = table.tenantColumn + const setWrites: Array<{ readonly value: unknown; readonly sql: string }> = [] + + const [assignments, whereSql] = withSubqueryCompiler( + (subquery) => + typeof subquery === "string" ? subquery : compileInner(subquery, values, { skipFormat: true, nested: true }).sql, + () => { + let assignments: ReadonlyArray = [] + if (write._tag === "CHUpdate") { + const set = (write as CHUpdate)._state.set + if (set === undefined) throw new QueryBuilderDefect({ message: `${where}: set() is required` }) + const record = typeof set === "function" ? set($) : set + assignments = setAssignments(table, record as Record, cell, (column, value, sql) => { + if (column === tenant) setWrites.push({ value, sql }) + }, where, "update") + } + const rendered = conditions.map((c) => compileSqlFragment(c.toFragment())) + const whereSql = + rendered.length > 0 + ? `\nWHERE ${rendered.join("\n AND ")}` + : dialect.clauses.alterTableUpdate === true + ? "\nWHERE 1" + : "" + return [assignments, whereSql] as const + }, + ) + + // Scope as for a query over the table, from the WHERE; an UPDATE that moves + // rows to another tenant reaches past it. + let tenantScope: TenantScope = "untenanted" + let tenantBound: string | undefined + if (tenant !== undefined) { + const derived = deriveTenantScope( + [{ column: tenant, scope: "cross-tenant" }], + [{ predicates: conditions.flatMap((c) => tenantPredicatesOf(c)) }], + (value) => inlineParams(compileSqlFragment(value), values), + ) + tenantScope = derived.scope + tenantBound = derived.bound + for (const { value, sql } of setWrites) { + const pinned = value !== null && (!isExprLike(value) || "_paramName" in value) + if (!pinned || inlineParams(sql, values) !== tenantBound) { + tenantScope = "cross-tenant" + tenantBound = undefined + } + } + } + + const returning = returningOf(table, state.returningFn, where) + const settings = writeSettingsClause(table, state.settings, where) + const target = quoteIdentPath(table.name) + const head = + kind === "delete" + ? `DELETE FROM ${target}` + : dialect.clauses.alterTableUpdate === true + ? `ALTER TABLE ${target} UPDATE ${assignments.join(", ")}` + : `UPDATE ${target} SET ${assignments.join(", ")}` + const rendered = renderParams(`${head}${whereSql}${returning?.sql ?? ""}${settings}`, values, dialect) + return writeCompiledQuery(rendered, kind, tenantScope, tenantScope === "single-tenant" ? tenantBound : undefined, returning) +} diff --git a/src/ch/dialect.ts b/src/ch/dialect.ts index 70229f4..00184ec 100644 --- a/src/ch/dialect.ts +++ b/src/ch/dialect.ts @@ -65,9 +65,12 @@ export interface DialectClauses { /** `ON CONFLICT ... DO NOTHING / DO UPDATE` after an INSERT. Absent means * no: an insert with `onConflict*` fails to compile for the dialect. */ readonly onConflict?: boolean - /** `SETTINGS` on an INSERT (ClickHouse). Absent means no: an insert with - * `.settings()` fails to compile for the dialect. */ - readonly insertSettings?: boolean + /** `SETTINGS` on an INSERT, UPDATE or DELETE (ClickHouse). Absent means + * no: a write with `.settings()` fails to compile for the dialect. */ + readonly writeSettings?: boolean + /** UPDATE is written `ALTER TABLE t UPDATE ... WHERE ...`, a ClickHouse + * mutation, rather than `UPDATE t SET ... WHERE ...`. */ + readonly alterTableUpdate?: boolean } /** A transaction isolation level. `read uncommitted` is left out: Postgres runs it as `read committed`. */ @@ -154,7 +157,8 @@ export const clickhouseDialect: Dialect = { parenthesizeUnionBranches: false, returning: false, onConflict: false, - insertSettings: true, + writeSettings: true, + alterTableUpdate: true, }, transactions: noTransactions, } diff --git a/src/ch/index.ts b/src/ch/index.ts index c181212..9e193b0 100644 --- a/src/ch/index.ts +++ b/src/ch/index.ts @@ -256,6 +256,17 @@ export { insertInto, } from "./insert" +// Update and delete builders +export { + type CHDelete, + type CHUpdate, + type CHUpdateStart, + type UpdateSet, + type UpdateSetOf, + deleteFrom, + update, +} from "./update" + // Compilation export { // `compileCH` / `compileCHUnsafe` are the internal names; the public API is @@ -268,6 +279,7 @@ export { type CompiledQuery, type CompiledQueryInput, type CompiledQueryRowSchema, + type CHWrite, type InsertCompileOptions, type RowSchemaMismatch, type TenantScope, diff --git a/src/ch/insert.test-d.ts b/src/ch/insert.test-d.ts index f3d88be..9fac5cb 100644 --- a/src/ch/insert.test-d.ts +++ b/src/ch/insert.test-d.ts @@ -162,3 +162,21 @@ expectTypeOf(PG.compileUnsafe(insert, {})).toEqualTypeOf >() + +// UPDATE / DELETE +const Ctr = CH.table("ctr", { key: PG.text, count: PG.int8, gen: PG.text }, { computed: ["gen"] }) +CH.update(Ctr).set({ count: 1 }).where(($) => [$.key.eq("k")]) +CH.update(Ctr).set(($) => ({ count: $.count.add(1) })).allRows() +// @ts-expect-error a value of another type +CH.update(Ctr).set({ count: "1" }) +// @ts-expect-error gen is generated +CH.update(Ctr).set({ gen: "x" }) +// @ts-expect-error set() first +CH.update(Ctr).where(($) => [$.key.eq("k")]) +const updated = CH.update(Ctr).set({ count: 1 }).allRows().returning("count") +expectTypeOf>().toEqualTypeOf<{ readonly count: number }>() +const deleted = CH.deleteFrom(Ctr).where(($) => [$.key.eq("k")]).returning() +expectTypeOf>().toEqualTypeOf<{ readonly key: string; readonly count: number; readonly gen: string }>() +expectTypeOf>>>().toEqualTypeOf() +expectTypeOf(PG.compileUnsafe(updated)).toEqualTypeOf>() +expectTypeOf>().toEqualTypeOf>() diff --git a/src/ch/update.test.ts b/src/ch/update.test.ts new file mode 100644 index 0000000..199cc02 --- /dev/null +++ b/src/ch/update.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, QueryBuilderError } from "./errors" + +const Counters = CH.table( + "counters", + { org: PG.text, key: PG.text, count: PG.int8, meta: PG.jsonb(), search: PG.text }, + { tenantColumn: "org", computed: ["search"] }, +) +const Spans = CH.table("spans", { OrgId: CH.string, Name: CH.string, Ms: CH.uint64 }, { tenantColumn: "OrgId" }) + +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 + +describe("update", () => { + it("writes UPDATE ... SET ... WHERE ... RETURNING on Postgres, values bound", () => { + const compiled = PG.compileUnsafe( + CH.update(Counters) + .set(($) => ({ count: $.count.add(1), meta: { a: 1 } })) + .where(($) => [$.org.eq(CH.param.string("org")), $.key.eq("k")]) + .returning("count"), + { org: "o" }, + ) + expect(compiled.sql).toBe( + 'UPDATE "counters" SET "count" = "count" + 1, "meta" = $1\nWHERE "org" = $2\n AND "key" = \'k\'\nRETURNING "count" AS "count"', + ) + expect(compiled.parameters).toEqual(['{"a":1}', "o"]) + expect(compiled.kind).toBe("update") + expect(compiled.returning).toEqual(["count"]) + expect(compiled.tenantScope).toBe("single-tenant") + }) + + it("writes an ALTER TABLE ... UPDATE mutation on ClickHouse, with settings last", () => { + const compiled = CH.compileUnsafe( + CH.update(Spans) + .set({ Name: "x" }) + .where(($) => [$.OrgId.eq("o")]) + .settings({ mutations_sync: 2 }), + ) + expect(compiled.sql).toBe("ALTER TABLE spans UPDATE Name = 'x'\nWHERE OrgId = 'o' SETTINGS mutations_sync = 2") + expect(CH.compileUnsafe(CH.update(Spans).set({ Ms: 0 }).allRows()).sql).toBe("ALTER TABLE spans UPDATE Ms = 0\nWHERE 1") + }) + + it("allRows() writes no WHERE on Postgres and makes the write cross-tenant", () => { + const compiled = PG.compileUnsafe(CH.update(Counters).set({ count: 0 }).allRows()) + expect(compiled.sql).toBe('UPDATE "counters" SET "count" = $1') + expect(compiled.tenantScope).toBe("cross-tenant") + }) + + it("an update that moves rows to another tenant is cross-tenant", () => { + const scope = (org: string) => + PG.compileUnsafe(CH.update(Counters).set({ org }).where(($) => [$.org.eq("o")])).tenantScope + expect(scope("o")).toBe("single-tenant") + expect(scope("p")).toBe("cross-tenant") + }) + + it.effect("refuses writes over every row unless asked, and bad SETs", () => + Effect.gen(function* () { + const noWhere = yield* Effect.exit(PG.compile(CH.update(Counters).set({ count: 0 }) as unknown as CH.CHUpdate)) + expect(failure(noWhere)).toBeInstanceOf(QueryBuilderDefect) + // Optional predicates that all came out undefined are data, so a failure. + const allUndefined = yield* Effect.flip( + PG.compile(CH.update(Counters).set({ count: 0 }).where(($) => [CH.when(undefined, (v: string) => $.key.eq(v))])), + ) + expect(allUndefined).toBeInstanceOf(QueryBuilderError) + expect(allUndefined.message).toContain("would write every row") + const computed = yield* Effect.flip( + PG.compile(CH.update(Counters).set({ search: "x" } as any).where(($) => [$.key.eq("k")])), + ) + expect(computed.message).toContain('sets "search"') + const empty = yield* Effect.flip(PG.compile(CH.update(Counters).set({}).where(($) => [$.key.eq("k")]))) + expect(empty.message).toContain("sets no columns") + const returning = yield* Effect.exit(CH.compile(CH.update(Spans).set({ Ms: 0 }).allRows().returning())) + expect(failure(returning)).toBeInstanceOf(QueryBuilderDefect) + const settings = yield* Effect.exit(PG.compile(CH.update(Counters).set({ count: 0 }).allRows().settings({ a: 1 }))) + expect(failure(settings)).toBeInstanceOf(QueryBuilderDefect) + }), + ) +}) + +describe("deleteFrom", () => { + it("writes DELETE ... WHERE ... RETURNING on Postgres and a lightweight DELETE on ClickHouse", () => { + const pg = PG.compileUnsafe( + CH.deleteFrom(Counters).where(($) => [$.org.eq(CH.param.string("org"))]).returning(), + { org: "o" }, + ) + expect(pg.sql).toBe( + 'DELETE FROM "counters"\nWHERE "org" = $1\nRETURNING "org" AS "org", "key" AS "key", "count" AS "count", "meta" AS "meta", "search" AS "search"', + ) + expect(pg.kind).toBe("delete") + expect(pg.tenantScope).toBe("single-tenant") + const ch = CH.compileUnsafe(CH.deleteFrom(Spans).where(($) => [$.OrgId.eq("o")]).settings({ lightweight_deletes_sync: 2 })) + expect(ch.sql).toBe("DELETE FROM spans\nWHERE OrgId = 'o' SETTINGS lightweight_deletes_sync = 2") + expect(CH.compileUnsafe(CH.deleteFrom(Spans).allRows()).sql).toBe("DELETE FROM spans\nWHERE 1") + expect(CH.compileUnsafe(CH.deleteFrom(Spans).allRows()).tenantScope).toBe("cross-tenant") + }) + + it.effect("refuses a delete with no where()", () => + Effect.gen(function* () { + expect(failure(yield* Effect.exit(CH.compile(CH.deleteFrom(Spans))))).toBeInstanceOf(QueryBuilderDefect) + }), + ) +}) diff --git a/src/ch/update.ts b/src/ch/update.ts new file mode 100644 index 0000000..a370ac5 --- /dev/null +++ b/src/ch/update.ts @@ -0,0 +1,156 @@ +// Update and Delete Builders +// +// `update(table).set(...).where(...)` and `deleteFrom(table).where(...)` +// describe writes the way `insertInto` does: immutable values read by +// `compile`, which writes them for the dialect it is given. SET values are +// encoded through each column's own codec. See `design/writes.md`. +// +// A write with no WHERE changes every row, so it has to say so with +// `allRows()`; compiling one that does not is a defect. +// +// Usage: +// CH.update(Counters) +// .set(($) => ({ count: $.count.add(1) })) +// .where(($) => [$.key.eq(CH.param.string("key"))]) +// .returning("count") +// +// CH.deleteFrom(ApiKeys).where(($) => [$.orgId.eq(CH.param.string("orgId"))]) + +import type { Condition, Expr } from "./expr" +import type { ConflictSet, InsertSettingValue } from "./insert" +import type { ColumnAccessor, InferOutput } from "./query" +import type { Table } from "./table" +import type { ColumnDefs, InferTS } from "./types" + +/** + * What an UPDATE writes: any insertable column, as a value, param or + * expression. A key left out (or `undefined`) keeps the existing value. + */ +export type UpdateSet = ConflictSet + +/** The SET record of a table value: `UpdateSetOf`. */ +export type UpdateSetOf = T extends Table ? UpdateSet : never + +type WhereFn = ($: ColumnAccessor) => Array + +/** @internal — what UPDATE and DELETE share. */ +interface WriteState { + readonly table: Table + readonly whereFn?: ($: any) => Array + /** Set by `allRows()`: the write is meant to touch every row. */ + readonly allRows?: boolean + readonly returningFn?: ($: any) => Record> + readonly settings?: Readonly> +} + +/** @internal — runtime update state */ +export interface CHUpdateState extends WriteState { + readonly set?: ConflictSet | (($: any) => ConflictSet) +} + +/** @internal — runtime delete state */ +export type CHDeleteState = WriteState + +type AllColumns = { readonly [P in keyof Cols & string]: InferTS } + +/** The clauses UPDATE and DELETE share, returning `Self` with `Output` replaced. */ +interface WriteClauses { + /** + * The rows to change, as in a query's `where`: conditions AND-joined, an + * `undefined` one skipped. Calling it again replaces them. + */ + where(fn: WhereFn): Self + /** Change every row. Without it, compiling a write with no WHERE is a defect. */ + allRows(): Self + /** + * ClickHouse `SETTINGS` for this write, such as `{ mutations_sync: 2 }` so an + * `ALTER TABLE ... UPDATE` waits for the mutation. Postgres refuses them. + */ + settings(settings: Readonly>): Self +} + +export interface CHUpdate + extends WriteClauses> { + readonly _tag: "CHUpdate" + /** @internal — runtime update state */ + readonly _state: CHUpdateState + /** phantom. `output` is the row `Database.run` returns: none without RETURNING. */ + readonly _phantom?: { readonly cols: Cols; readonly output: Output } + + /** Replace the SET record. */ + set(set: UpdateSet | (($: ColumnAccessor) => UpdateSet)): CHUpdate + + /** The changed rows, as for an insert: every column, the named ones, or a callback. Postgres only. */ + returning(): CHUpdate> + returning( + ...columns: [K, ...Array] + ): CHUpdate }> + returning>>(fn: ($: ColumnAccessor) => S): CHUpdate> +} + +export interface CHDelete + extends WriteClauses> { + readonly _tag: "CHDelete" + /** @internal — runtime delete state */ + readonly _state: CHDeleteState + readonly _phantom?: { readonly cols: Cols; readonly output: Output } + + /** The deleted rows, as for an insert: every column, the named ones, or a callback. Postgres only. */ + returning(): CHDelete> + returning(...columns: [K, ...Array]): CHDelete }> + returning>>(fn: ($: ColumnAccessor) => S): CHDelete> +} + +/** An update with no SET yet: only `set`, so it cannot be compiled before it says what to write. */ +export type CHUpdateStart = Pick< + CHUpdate, + "set" +> + +/** `returning(...)`'s arguments as a select callback, shared with inserts. */ +export const returningFnOf = + (table: Table) => + (args: ReadonlyArray): (($: any) => Record>) => { + const [first] = args + if (typeof first === "function") return first as ($: any) => Record> + const columns = args.length === 0 ? Object.keys(table.columns) : (args as ReadonlyArray) + return ($: any) => Object.fromEntries(columns.map((column) => [column, $[column]])) + } + +const writeClauses = (state: State, make: (state: State) => Self) => ({ + where: (whereFn: ($: any) => Array) => make({ ...state, whereFn }), + allRows: () => make({ ...state, allRows: true }), + settings: (settings: Readonly>) => make({ ...state, settings: { ...settings } }), + returning: (...args: ReadonlyArray) => make({ ...state, returningFn: returningFnOf(state.table)(args) }), +}) + +const makeUpdate = (state: CHUpdateState): CHUpdate => ({ + _tag: "CHUpdate", + _state: state, + ...writeClauses(state, makeUpdate), + set: (set) => makeUpdate({ ...state, set }), +}) + +const makeDelete = (state: CHDeleteState): CHDelete => ({ + _tag: "CHDelete", + _state: state, + ...writeClauses(state, makeDelete), +}) + +/** Start an UPDATE of `table`. Give what to write with `set`, and the rows with `where` or `allRows`. */ +export function update( + table: Table, +): CHUpdateStart { + return makeUpdate({ table: table as Table }) as CHUpdateStart +} + +/** Start a DELETE from `table`. Give the rows with `where` or `allRows`. */ +export function deleteFrom(table: Table): CHDelete { + return makeDelete({ table: table as Table }) as CHDelete +} + +export const isUpdate = (value: unknown): value is CHUpdate => + typeof value === "object" && value !== null && (value as { readonly _tag?: unknown })._tag === "CHUpdate" + +export const isDelete = (value: unknown): value is CHDelete => + typeof value === "object" && value !== null && (value as { readonly _tag?: unknown })._tag === "CHDelete" diff --git a/src/database/database.test.ts b/src/database/database.test.ts index bcd85ad..db886f9 100644 --- a/src/database/database.test.ts +++ b/src/database/database.test.ts @@ -491,6 +491,35 @@ layer(Live, { excludeTestServices: true })("Database on PGlite", (it) => { }), ) + it.effect("update and deleteFrom change rows and return them", () => + Effect.gen(function* () { + yield* Db.execute(Db.sql`CREATE TABLE tickets (id int4 PRIMARY KEY, org text NOT NULL, seats int4 NOT NULL, tags text[] NOT NULL)`) + const Tickets = CH.table("tickets", { id: PG.int4, org: PG.text, seats: PG.int4, tags: PG.array(PG.text) }, { tenantColumn: "org" }) + yield* Db.run( + CH.insertInto(Tickets).values([ + { id: 1, org: "a", seats: 1, tags: [] }, + { id: 2, org: "a", seats: 5, tags: [] }, + { id: 3, org: "b", seats: 9, tags: [] }, + ]), + ) + const bumped = yield* Db.run( + CH.update(Tickets) + .set(($) => ({ seats: $.seats.add(1), tags: ["x", "it's"] })) + .where(($) => [$.org.eq(CH.param.string("org")), $.seats.lt(5)]) + .returning("id", "seats", "tags"), + { org: "a" }, + ) + expect(bumped).toEqual([{ id: 1, seats: 2, tags: ["x", "it's"] }]) + // Without RETURNING, a write runs through `execute` and returns nothing. + expect(yield* Db.run(CH.update(Tickets).set({ seats: 0 }).where(($) => [$.id.eq(2)]))).toEqual([]) + const removed = yield* Db.run(CH.deleteFrom(Tickets).where(($) => [$.org.eq("a")]).returning("id")) + expect([...removed].map((r) => r.id).sort()).toEqual([1, 2]) + expect(yield* Db.run(CH.from(Tickets).select("id", "seats"))).toEqual([{ id: 3, seats: 9 }]) + yield* Db.run(CH.deleteFrom(Tickets).allRows()) + expect(yield* Db.run(CH.from(Tickets).select("id"))).toEqual([]) + }), + ) + it.effect("an insert inside a failed transaction rolls back", () => Effect.gen(function* () { const table = yield* freshTable diff --git a/src/database/database.ts b/src/database/database.ts index 1c367c5..3f04d7f 100644 --- a/src/database/database.ts +++ b/src/database/database.ts @@ -15,6 +15,7 @@ import { compileCH, compileUnion, CompiledQueryDecodeError, type CompiledQuery } import { noTransactions, type Dialect, type IsolationLevel, type TransactionSettings } from "../ch/dialect" import type { QueryBuilderError } from "../ch/errors" import type { CHInsert } from "../ch/insert" +import type { CHDelete, CHUpdate } from "../ch/update" import type { CHQuery } from "../ch/query" import type { CHUnionQuery } from "../ch/union" import { @@ -39,7 +40,13 @@ export interface Statement { export type StatementInput = SqlTemplate | Statement /** What `run` takes: a built query, a union, an insert, or one already compiled. */ -export type Runnable = CHQuery | CHUnionQuery | CHInsert | CompiledQuery +export type Runnable = + | CHQuery + | CHUnionQuery + | CHInsert + | CHUpdate + | CHDelete + | CompiledQuery /** The decoded row of a `Runnable`. */ export type RowOf = @@ -85,7 +92,7 @@ export interface DatabaseApi { readonly dialect: Dialect /** * Compile a query for this database's dialect, run it, and decode its rows. - * An insert returns its RETURNING rows, or none without `returning`. + * A write returns its RETURNING rows, or none without `returning`. * `params` fills the query's `param.*` markers. A query compiled elsewhere * runs as it is, if it was compiled for this dialect. */ @@ -259,16 +266,18 @@ export const fromSqlClient = (sql: SqlClient.SqlClient, options: FromSqlClientOp : Effect.succeed(runnable) } if ("_tag" in runnable && runnable._tag === "CHUnionQuery") return compileUnion(runnable, params, { dialect }) - if ("_tag" in runnable && runnable._tag === "CHInsert") return compileCH(runnable, params, { dialect }) + if ("_tag" in runnable && (runnable._tag === "CHInsert" || runnable._tag === "CHUpdate" || runnable._tag === "CHDelete")) { + return compileCH(runnable, params, { dialect }) + } return compileCH(runnable as CHQuery, params, { dialect }) } - // An insert without RETURNING sends back no rows, so it runs the way `execute` + // A write without RETURNING sends back no rows, so it runs the way `execute` // does: through `command`, which a ClickHouse client needs for a statement // with no result set. const run: DatabaseApi["run"] = (runnable, params = {}) => Effect.flatMap(compileFor(runnable, params), (compiled) => - compiled.kind === "insert" && compiled.returning === undefined + compiled.kind !== "select" && compiled.returning === undefined ? Effect.as(execute(compiled), []) : Effect.flatMap(rows(compiled), (wire) => compiled.decodeRows(wire)), ) diff --git a/tests/database.clickhouse.test.ts b/tests/database.clickhouse.test.ts index 1aa89b9..68ba68f 100644 --- a/tests/database.clickhouse.test.ts +++ b/tests/database.clickhouse.test.ts @@ -134,6 +134,37 @@ describe("database", () => { expect(rows).toEqual([{ OrgId: "o", Name: "a", Total: 3 }]) }) + it("runs update as an ALTER TABLE mutation and deleteFrom as a lightweight delete", async () => { + const rows = await Effect.runPromise( + withDatabase((db) => + Effect.gen(function* () { + yield* db.execute(Db.sql`CREATE TABLE jobs (OrgId String, Id UInt32, State String) ENGINE = MergeTree ORDER BY (OrgId, Id)`) + const Jobs = CH.table("jobs", { OrgId: CH.string, Id: CH.uint32, State: CH.string }, { tenantColumn: "OrgId" }) + yield* db.run( + CH.insertInto(Jobs).values([ + { OrgId: "o", Id: 1, State: "queued" }, + { OrgId: "o", Id: 2, State: "queued" }, + { OrgId: "p", Id: 3, State: "queued" }, + ]), + ) + yield* db.run( + CH.update(Jobs) + .set({ State: "done" }) + .where(($) => [$.OrgId.eq(CH.param.string("org")), $.Id.eq(1)]) + .settings({ mutations_sync: 2 }), + { org: "o" }, + ) + yield* db.run(CH.deleteFrom(Jobs).where(($) => [$.OrgId.eq("p")]).settings({ lightweight_deletes_sync: 2 })) + return yield* db.run(CH.from(Jobs).select("OrgId", "Id", "State").orderBy(["Id", "asc"])) + }), + ), + ) + expect(rows).toEqual([ + { OrgId: "o", Id: 1, State: "done" }, + { OrgId: "o", Id: 2, State: "queued" }, + ]) + }) + it("refuses a transaction before sending anything", async () => { const result = await Effect.runPromise( withDatabase((db, sent) => From 816265dc088dfc89788ddb86e5ef7b8d2f70b626 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 15:34:00 +0200 Subject: [PATCH 2/2] Count subqueries in a write's tenant scope; ignore empty conditions A subquery in an UPDATE or DELETE's SET or WHERE, or in an insert's VALUES or onConflictDoUpdate, was compiled for its SQL and its scope dropped, so a pinned write that read every tenant through a subquery reported single-tenant, and a write into an untenanted table reported untenanted. Writes now record each subquery's scope (a string subquery is cross-tenant) and combine it with their own, as queries do. A where() condition that renders to nothing no longer leaves a dangling WHERE, and does not count as a filter: a write left with none fails unless allRows() says so. Co-Authored-By: Claude Opus 5.5 --- docs/updates-and-deletes.md | 11 +++--- src/ch/compile.ts | 68 +++++++++++++++++++++++++++++-------- src/ch/insert.ts | 4 ++- src/ch/update.test.ts | 52 +++++++++++++++++++++++++++- 4 files changed, 114 insertions(+), 21 deletions(-) diff --git a/docs/updates-and-deletes.md b/docs/updates-and-deletes.md index a8f052e..a1b53c7 100644 --- a/docs/updates-and-deletes.md +++ b/docs/updates-and-deletes.md @@ -41,9 +41,9 @@ names the record type. so optional filters compose. A write with no `where` would change every row, so: - compiling an UPDATE or DELETE with no `where()` is a `QueryBuilderDefect`; -- a `where()` whose conditions all came out `undefined` is a `QueryBuilderError`, because that - happens with data (every optional filter absent) and would otherwise widen a filtered write to - the whole table; +- a `where()` whose conditions all came out `undefined` (or render to nothing) is a + `QueryBuilderError`, because that happens with data (every optional filter absent) and would + otherwise widen a filtered write to the whole table; - `allRows()` says a write over every row is meant. ## RETURNING @@ -77,7 +77,10 @@ cannot update a column of the sorting key. `returning` is refused on ClickHouse An UPDATE or DELETE has the scope a query over the table with the same WHERE would have: `"single-tenant"` when the WHERE pins the tenant column, `"cross-tenant"` otherwise (including `allRows()`). An UPDATE that sets the tenant column to another value moves rows out of the -tenant, so it is `"cross-tenant"` too. +tenant, so it is `"cross-tenant"` too. A subquery in the SET or WHERE counts as it would in a +query: one that reads another tenant (or every tenant) makes the write `"cross-tenant"`, and a +write into a table without a tenant column takes the scope of what its subqueries read. The +same holds for a subquery in an insert's values or `onConflictDoUpdate`. _(Backed by `src/ch/update.test.ts`, `src/database/database.test.ts` and `tests/database.clickhouse.test.ts`.)_ diff --git a/src/ch/compile.ts b/src/ch/compile.ts index 9c3d06a..65a7439 100644 --- a/src/ch/compile.ts +++ b/src/ch/compile.ts @@ -1464,6 +1464,35 @@ const valueCells = (table: Table, params: Record, reads: Array) => + (subquery: Parameters[0]>[0]): string => { + if (typeof subquery === "string") { + reads.push({ scope: "cross-tenant" }) + return subquery + } + const compiled = compileInner(subquery, values, { skipFormat: true, nested: true }) + reads.push({ scope: compiled.tenantScope, bound: tenantBoundOf(compiled) }) + return compiled.sql + } + +/** A write's own scope combined with what its subqueries read. */ +const withReads = ( + own: { readonly scope: TenantScope; readonly bound: string | undefined }, + reads: ReadonlyArray, +): { readonly scope: TenantScope; readonly bound: string | undefined } => { + if (reads.length === 0) return own + const derived = deriveTenantScope([{ scope: own.scope, bound: own.bound }, ...reads], [], (value) => + compileSqlFragment(value), + ) + return { scope: derived.scope, bound: derived.scope === "single-tenant" ? derived.bound : undefined } +} + /** `column = value, ...` for a SET record, which may come from data, so a bad key is a failure. */ const setAssignments = ( table: Table, @@ -1574,9 +1603,10 @@ function compileInsert(insert: CHInsert, params: Record = [] + let readBound: string | undefined const [columns, source, conflictSql, readScope] = withSubqueryCompiler( - (subquery) => - typeof subquery === "string" ? subquery : compileInner(subquery, values, { skipFormat: true, nested: true }).sql, + recordingSubqueries(values, reads), () => { let columns: ReadonlyArray let source: string @@ -1602,6 +1632,7 @@ function compileInsert(insert: CHInsert, params: Record, params: Record => c != null) - if (state.whereFn !== undefined && conditions.length === 0 && state.allRows !== true) { - throw new QueryBuilderError({ - code: "InvalidArguments", - message: `${where}: every where() condition was undefined, which would write every row; call allRows() if that is meant`, - }) - } const tenant = table.tenantColumn const setWrites: Array<{ readonly value: unknown; readonly sql: string }> = [] + const reads: Array = [] const [assignments, whereSql] = withSubqueryCompiler( - (subquery) => - typeof subquery === "string" ? subquery : compileInner(subquery, values, { skipFormat: true, nested: true }).sql, + recordingSubqueries(values, reads), () => { let assignments: ReadonlyArray = [] if (write._tag === "CHUpdate") { @@ -1727,7 +1757,14 @@ function compileUpdateOrDelete( if (column === tenant) setWrites.push({ value, sql }) }, where, "update") } - const rendered = conditions.map((c) => compileSqlFragment(c.toFragment())) + // An empty rendering (a `rawCond("")`) filters nothing, so it does not count. + const rendered = conditions.map((c) => compileSqlFragment(c.toFragment())).filter((sql) => sql.trim() !== "") + if (state.whereFn !== undefined && rendered.length === 0 && state.allRows !== true) { + throw new QueryBuilderError({ + code: "InvalidArguments", + message: `${where}: where() gave no conditions (each was undefined or empty), which would write every row; call allRows() if that is meant`, + }) + } const whereSql = rendered.length > 0 ? `\nWHERE ${rendered.join("\n AND ")}` @@ -1769,5 +1806,6 @@ function compileUpdateOrDelete( ? `ALTER TABLE ${target} UPDATE ${assignments.join(", ")}` : `UPDATE ${target} SET ${assignments.join(", ")}` const rendered = renderParams(`${head}${whereSql}${returning?.sql ?? ""}${settings}`, values, dialect) - return writeCompiledQuery(rendered, kind, tenantScope, tenantScope === "single-tenant" ? tenantBound : undefined, returning) + const scope = withReads({ scope: tenantScope, bound: tenantScope === "single-tenant" ? tenantBound : undefined }, reads) + return writeCompiledQuery(rendered, kind, scope.scope, scope.bound, returning) } diff --git a/src/ch/insert.ts b/src/ch/insert.ts index 26ec888..e7b9dea 100644 --- a/src/ch/insert.ts +++ b/src/ch/insert.ts @@ -72,7 +72,9 @@ type SelectedRow = Q extends { readonly _phantom?: { readonly output: infer O /** Selected columns the table cannot take: not an insertable column, or of another type. */ export type InsertSelectMisfits = { - // The rule a comparison uses: a branded column takes the plain primitive. + // The rule `values` and comparisons use: a column takes its widened + // primitive, so a branded column takes a plain string. That also lets a plain + // string into a literal-union column; the server checks those values. [K in keyof Output]: K extends Exclude> ? [Output[K]] extends [InferTS | Widen>] ? never diff --git a/src/ch/update.test.ts b/src/ch/update.test.ts index 199cc02..cc2edce 100644 --- a/src/ch/update.test.ts +++ b/src/ch/update.test.ts @@ -80,6 +80,49 @@ describe("update", () => { ) }) +describe("subqueries in a write count toward its tenant scope", () => { + const Other = CH.table("other", { org: PG.text, key: PG.text }, { tenantColumn: "org" }) + const Plain = CH.table("plain", { id: PG.text }) + + it("a WHERE subquery over every tenant makes a pinned update cross-tenant", () => { + const update = (inner: CH.CHQuery) => + PG.compileUnsafe( + CH.update(Counters) + .set({ count: 1 }) + .where(($) => [$.org.eq(CH.param.string("org")), CH.inSubquery($.key, inner)]), + { org: "o", other: "p" }, + ).tenantScope + expect(update(CH.from(Other).select(($) => ({ k: $.key })))).toBe("cross-tenant") + expect(update(CH.from(Other).select(($) => ({ k: $.key })).where(($) => [$.org.eq(CH.param.string("org"))]))).toBe( + "single-tenant", + ) + expect(update(CH.from(Other).select(($) => ({ k: $.key })).where(($) => [$.org.eq(CH.param.string("other"))]))).toBe( + "cross-tenant", + ) + }) + + it("an untenanted target takes the scope of what its subqueries read", () => { + const scope = (inner: CH.CHQuery) => + PG.compileUnsafe(CH.deleteFrom(Plain).where(($) => [CH.inSubquery($.id, inner)]), { org: "o" }).tenantScope + expect(scope(CH.from(Other).select(($) => ({ k: $.key })))).toBe("cross-tenant") + expect(scope(CH.from(Other).select(($) => ({ k: $.key })).where(($) => [$.org.eq(CH.param.string("org"))]))).toBe( + "single-tenant", + ) + expect(PG.compileUnsafe(CH.deleteFrom(Plain).where(($) => [$.id.eq("x")])).tenantScope).toBe("untenanted") + }) + + it("a subquery in an insert's VALUES or ON CONFLICT SET counts too", () => { + const anyKey = CH.subqueryExpr(CH.from(Other).select(($) => ({ k: $.key })).limit(1), PG.text) + const pinned = { org: "o", key: "k", count: 1, meta: {} } + expect(PG.compileUnsafe(CH.insertInto(Counters).values({ ...pinned, key: anyKey })).tenantScope).toBe("cross-tenant") + expect( + PG.compileUnsafe(CH.insertInto(Counters).values(pinned).onConflictDoUpdate({ target: ["key"], set: { key: anyKey } })) + .tenantScope, + ).toBe("cross-tenant") + expect(PG.compileUnsafe(CH.insertInto(Counters).values(pinned)).tenantScope).toBe("single-tenant") + }) +}) + describe("deleteFrom", () => { it("writes DELETE ... WHERE ... RETURNING on Postgres and a lightweight DELETE on ClickHouse", () => { const pg = PG.compileUnsafe( @@ -97,9 +140,16 @@ describe("deleteFrom", () => { expect(CH.compileUnsafe(CH.deleteFrom(Spans).allRows()).tenantScope).toBe("cross-tenant") }) - it.effect("refuses a delete with no where()", () => + it.effect("refuses a delete with no where(), or whose conditions filter nothing", () => Effect.gen(function* () { expect(failure(yield* Effect.exit(CH.compile(CH.deleteFrom(Spans))))).toBeInstanceOf(QueryBuilderDefect) + for (const conditions of [[], [CH.rawCond("")], [undefined, CH.rawCond(" ")]]) { + const error = yield* Effect.flip(PG.compile(CH.deleteFrom(Counters).where(() => conditions))) + expect(error.message).toContain("would write every row") + } + expect(PG.compileUnsafe(CH.deleteFrom(Counters).where(($) => [CH.rawCond(""), $.key.eq("k")])).sql).toBe( + 'DELETE FROM "counters"\nWHERE "key" = \'k\'', + ) }), ) })