From eab33d455be3e25f1ebac050998f4cfbb11ad583 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 01:01:59 +0200 Subject: [PATCH] Add onConflictDoNothing and onConflictDoUpdate to inserts ON CONFLICT on Postgres, with Drizzle's options so call sites move over with renames: target (columns or a named constraint), targetWhere for a partial index, set as a record or a callback over the existing row and excluded, and where to limit the update. The existing row is qualified with the table name, since an unqualified column is ambiguous with excluded. A SET that writes the tenant column counts toward tenant scope. DialectClauses.onConflict is optional; ClickHouse refuses at compile. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 3 + design/writes.md | 9 ++- docs/inserts.md | 45 +++++++++++++- docs/reference.md | 4 +- src/ch/compile.ts | 109 ++++++++++++++++++++++++++++------ src/ch/dialect.ts | 4 ++ src/ch/index.ts | 12 +++- src/ch/insert.test-d.ts | 16 +++++ src/ch/insert.test.ts | 68 +++++++++++++++++++++ src/ch/insert.ts | 64 +++++++++++++++++++- src/database/database.test.ts | 36 +++++++++++ src/pg/dialect.ts | 1 + 12 files changed, 346 insertions(+), 25 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a19ede9..d26c317 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,9 @@ - Add `returning` to an insert (Postgres): column names or a callback, as in `select`. `run` returns the inserted rows decoded through the derived row schema; `CompiledQuery.returning` lists the aliases. Add `DialectClauses.returning`, optional, absent meaning no. +- Add `onConflictDoNothing` and `onConflictDoUpdate` to an insert (Postgres), with Drizzle's + options: `target` (columns or `{ constraint }`), `targetWhere`, `set` (a record, or a callback + over the existing row and `excluded`) and `where`. Add `DialectClauses.onConflict`, 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/design/writes.md b/design/writes.md index 64a2738..6a33e6f 100644 --- a/design/writes.md +++ b/design/writes.md @@ -1,7 +1,7 @@ # Writes: INSERT -Status: phases 1 and 2 built (`insertInto(...).values(...).returning(...)`, see -`docs/inserts.md`); phases 3 and 4 not started. Section 11 lists where the build differs from the plan. UPDATE and DELETE come later and +Status: phases 1 to 3 built (`values`, `returning`, `onConflictDoNothing` / +`onConflictDoUpdate`, see `docs/inserts.md`); phase 4 not started. 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). @@ -245,3 +245,8 @@ note, reusing `kind: "write"`, the returning path and the `set` record type from optional so a dialect written against the earlier interface still compiles. - **`CompiledQuery.returning`** lists the RETURNING aliases; `run` reads it to choose between the command path and the row path, so a precompiled insert runs correctly too. +- **`ON CONFLICT` takes Drizzle's shape**, `onConflictDoNothing({ target?, targetWhere? })` and + `onConflictDoUpdate({ target, targetWhere?, set, where? })`, instead of the chained + `.onConflict(target).doUpdate(...)` in §6: Maple's 68 call sites then move over with renames. + `$` in `set` and `where` is qualified with the table name (`"counters"."count"`), because an + unqualified column there is ambiguous with `excluded`. diff --git a/docs/inserts.md b/docs/inserts.md index 26eb51b..07389fa 100644 --- a/docs/inserts.md +++ b/docs/inserts.md @@ -109,11 +109,52 @@ The row schema is derived from the list, as it is from a SELECT: an untyped expr `untypedColumns`. `CompiledQuery.returning` lists the aliases. ClickHouse has no RETURNING, so compiling an insert with `returning` for it is a `QueryBuilderDefect`. +## On conflict + +On Postgres, `onConflictDoNothing` and `onConflictDoUpdate` add an `ON CONFLICT` clause. Their +options follow Drizzle's, so code moving from Drizzle changes little. + +```ts +const Counters = CH.table("counters", { key: PG.text, count: PG.int8, locked: PG.bool }, { defaults: ["locked"] }) + +// Skip a row whose key exists. Without `target`, any unique index or constraint counts. +CH.insertInto(Counters).values({ key: "a", count: 1 }).onConflictDoNothing({ target: ["key"] }) + +// Upsert: add to the existing count, unless the row is locked. +CH.insertInto(Counters) + .values({ key: "a", count: 1 }) + .onConflictDoUpdate({ + target: ["key"], + set: ($, excluded) => ({ count: $.count.add(excluded.count) }), + where: ($) => $.locked.eq(false), + }) + .returning("key", "count") +``` + +```sql +INSERT INTO "counters" ("key", "count") +VALUES ($1, $2) +ON CONFLICT ("key") DO UPDATE SET "count" = "counters"."count" + "excluded"."count" WHERE "counters"."locked" = FALSE +RETURNING "key" AS "key", "count" AS "count" +``` + +- `target` is column names, or `{ constraint: "name" }`. `targetWhere` gives a partial unique + index's predicate. `onConflictDoUpdate` requires a `target`; `onConflictDoNothing` does not. +- `set` is a record of values, params or expressions, or a callback that gets `$` (the existing + row) and `excluded` (the row proposed for insertion). `$` is qualified with the table name, + because an unqualified column would be ambiguous with `excluded`. A key left out keeps the + existing value. +- `where` limits the update to existing rows it holds for. A row it skips is not updated and, + with `returning`, returns nothing; the same goes for a row `onConflictDoNothing` skips. +- Calling either again replaces the clause. ClickHouse has no `ON CONFLICT` (deduplicate with a + `ReplacingMergeTree` instead), so compiling one for it is a `QueryBuilderDefect`. + ## Tenant scope An insert has a `tenantScope` like a query, worked out the same way. On a table with a `tenantColumn`, it is `"single-tenant"` when every row gives that column the same value or the -same param, and `"cross-tenant"` when rows differ or a row uses another expression. A table +same param, and `"cross-tenant"` when rows differ or a row uses another expression. An +`onConflictDoUpdate` that sets the tenant column counts as one more row. A table without a tenant column gives `"untenanted"`. ## Failures @@ -127,6 +168,8 @@ without a tenant column gives `"untenanted"`. | Over the dialect's bound-value limit | `QueryBuilderError` `InvalidArguments` | | Compiling without `values` | `QueryBuilderDefect` | | `returning` for a dialect without RETURNING | `QueryBuilderDefect` | +| `onConflictDoUpdate` setting no or unknown columns | `QueryBuilderError` `InvalidArguments` | +| `onConflict*` without ON CONFLICT, a bad target | `QueryBuilderDefect` | _(Backed by `src/ch/insert.test.ts`, `src/database/database.test.ts` and `tests/database.clickhouse.test.ts`.)_ diff --git a/docs/reference.md b/docs/reference.md index ff2c72e..1d96d21 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -62,7 +62,7 @@ Note `/sql` exports a `compile` (fragment → string) distinct from the root `co | `fromQuery` | `(query, alias) => CHQuery` | | `fromUnion` | `(union, alias) => CHQuery` | | `unionAll` | `(...queries) => CHUnionQuery` | -| `insertInto` | `(table) => CHInsert`; `.values(row \| rows)` sets its rows, `.returning(...)` the RETURNING list (Postgres). See [Inserting rows](./inserts.md) | +| `insertInto` | `(table) => CHInsert`; `.values(row \| rows)` sets its rows, `.returning(...)` the RETURNING list, `.onConflictDoNothing(options?)` / `.onConflictDoUpdate(options)` the ON CONFLICT clause (Postgres). See [Inserting rows](./inserts.md) | ### `CHQuery` methods @@ -279,7 +279,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`, `InsertRow`, `InsertRowOf`, `InsertValue`, `ColumnAccessor`, `JoinedColumnAccessor`, +`ParamKind`, `CHQuery`, `CHUnionQuery`, `CHInsert`, `InsertRow`, `InsertRowOf`, `InsertValue`, `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/src/ch/compile.ts b/src/ch/compile.ts index d1bd5a4..585749d 100644 --- a/src/ch/compile.ts +++ b/src/ch/compile.ts @@ -1333,6 +1333,74 @@ const VALUE_PARAM = "$$v" const EMPTY_ROW = Schema.Struct({}) as unknown as CompiledQueryRowSchema +/** + * The ON CONFLICT clause, or `""` without one. `$` in SET and WHERE is the + * existing row, qualified with the table name: an unqualified column there + * would be ambiguous with `excluded`. + */ +const onConflictClause = ( + insert: CHInsert, + cell: (column: string, value: unknown, context: string) => string, + wrote: (column: string, value: unknown, sql: string) => void, +): string => { + const { table, conflict } = insert._state + if (conflict === undefined) return "" + const where = `insertInto(${table.name})` + const dialect = currentDialect() + if (dialect.clauses.onConflict !== true) { + throw new QueryBuilderDefect({ + message: `${where}: onConflict has no meaning for the ${dialect.name} dialect, which has no ON CONFLICT clause`, + }) + } + + let target = "" + if (conflict.target !== undefined) { + if ("constraint" in conflict.target) { + if (conflict.targetWhere !== undefined) { + throw new QueryBuilderDefect({ message: `${where}: targetWhere needs a column target, not a constraint` }) + } + target = ` ON CONSTRAINT ${quoteIdent(conflict.target.constraint)}` + } else { + const targetColumns = conflict.target + if (targetColumns.length === 0 || targetColumns.some((column) => !Object.hasOwn(table.columns, column))) { + throw new QueryBuilderDefect({ message: `${where}: the conflict target must name columns of the table` }) + } + const predicate = conflict.targetWhere?.(createColumnAccessor(table.columns)) + target = ` (${targetColumns.map(quoteIdent).join(", ")})${ + predicate === undefined ? "" : ` WHERE ${compileSqlFragment(predicate.toFragment())}` + }` + } + } else if (conflict.targetWhere !== undefined) { + throw new QueryBuilderDefect({ message: `${where}: targetWhere needs a target` }) + } + if (conflict.action === "nothing") return `\nON CONFLICT${target} DO NOTHING` + + 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 condition = conflict.where?.(existing, excluded) + return `\nON CONFLICT${target} DO UPDATE SET ${assignments.join(", ")}${ + condition === undefined ? "" : ` WHERE ${compileSqlFragment(condition.toFragment())}` + }` +} + /** The RETURNING clause and its row schema, or none without `returning`. */ const returningOf = (insert: CHInsert) => { const { table, returningFn } = insert._state @@ -1389,52 +1457,57 @@ function compileInsert(insert: CHInsert, params: Record = { ...params } let next = 0 - const cell = (column: string, value: unknown, index: number): string => { + 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}: row ${index}, column ${column}`) + 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()) } - // A tenant table's insert is single-tenant when every row pins the tenant - // column to the same value or param. Any other expression, a NULL or a - // default could be anything. + // A tenant table's insert is single-tenant when every row (and an upsert's + // SET, if it writes the column) pins the tenant column to the same value or + // param. Any other expression, a NULL or a default could be anything. const tenant = table.tenantColumn const bounds = new Set() - let pinned = tenant !== undefined + let pinned = tenant !== undefined && present.has(tenant) + const pin = (value: unknown, sql: string): void => { + if (!pinned) return + if (value === undefined || value === null || (isExprLike(value) && !("_paramName" in value))) pinned = false + else bounds.add(inlineParams(sql, values)) + } - const tuples = withSubqueryCompiler( + const [tuples, conflictSql] = withSubqueryCompiler( (subquery) => typeof subquery === "string" ? subquery : compileInner(subquery, values, { skipFormat: true, nested: true }).sql, - () => - rows.map((row, index) => { + () => { + const tuples = rows.map((row, index) => { const cells = columns.map((column) => { const value = row[column] - const sql = cell(column, value, index) - if (column === tenant && pinned) { - if (value === undefined || value === null || (isExprLike(value) && !("_paramName" in value))) pinned = false - else bounds.add(inlineParams(sql, values)) - } + const sql = cell(column, value, `row ${index}`) + if (column === tenant) pin(value, sql) return sql }) - if (tenant !== undefined && !present.has(tenant)) pinned = false return `(${cells.join(", ")})` - }), + }) + return [tuples, onConflictClause(insert, cell, (column, value, sql) => { + if (column === tenant) pin(value, sql) + })] as const + }, ) - const dialect = currentDialect() const returning = returningOf(insert) const returningSql = returning === undefined ? "" : `\nRETURNING ${returning.aliases.map((alias) => compileSqlFragment(aliased(returning.exprs[alias]!, alias))).join(", ")}` const rendered = renderParams( - `INSERT INTO ${quoteIdentPath(table.name)} (${columns.map(quoteIdent).join(", ")})\nVALUES ${tuples.join(", ")}${returningSql}`, + `INSERT INTO ${quoteIdentPath(table.name)} (${columns.map(quoteIdent).join(", ")})\nVALUES ${tuples.join(", ")}${conflictSql}${returningSql}`, values, dialect, ) diff --git a/src/ch/dialect.ts b/src/ch/dialect.ts index e117c0f..089f65d 100644 --- a/src/ch/dialect.ts +++ b/src/ch/dialect.ts @@ -62,6 +62,9 @@ export interface DialectClauses { /** `RETURNING` after an INSERT. Absent means no: an insert with * `.returning()` fails to compile for the dialect. */ readonly returning?: boolean + /** `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 } /** A transaction isolation level. `read uncommitted` is left out: Postgres runs it as `read committed`. */ @@ -147,6 +150,7 @@ export const clickhouseDialect: Dialect = { groupByAlias: true, parenthesizeUnionBranches: false, returning: false, + onConflict: false, }, transactions: noTransactions, } diff --git a/src/ch/index.ts b/src/ch/index.ts index 0bca824..82e8e1f 100644 --- a/src/ch/index.ts +++ b/src/ch/index.ts @@ -240,7 +240,17 @@ export { } from "./query" // Insert builder -export { type CHInsert, type InsertRow, type InsertRowOf, type InsertValue, insertInto } from "./insert" +export { + type CHInsert, + type ConflictSet, + type ConflictTarget, + type InsertRow, + type InsertRowOf, + type InsertValue, + type OnConflictDoNothing, + type OnConflictDoUpdate, + insertInto, +} from "./insert" // Compilation export { diff --git a/src/ch/insert.test-d.ts b/src/ch/insert.test-d.ts index 40014be..4fee379 100644 --- a/src/ch/insert.test-d.ts +++ b/src/ch/insert.test-d.ts @@ -94,6 +94,22 @@ expectTypeOf(PG.compileUnsafe(returningExprs)).toEqualTypeOf< // @ts-expect-error not a column CH.insertInto(Keys).values({ id: "k" }).returning("nope") +// ON CONFLICT: targets are column names, SET takes values or expressions. +const Counters = CH.table("counters", { key: PG.text, count: PG.int8 }) +CH.insertInto(Counters).values({ key: "k", count: 1 }).onConflictDoNothing({ target: ["key"] }) +CH.insertInto(Counters) + .values({ key: "k", count: 1 }) + .onConflictDoUpdate({ target: ["key"], set: ($, excluded) => ({ count: $.count.add(excluded.count) }) }) +CH.insertInto(Counters).values({ key: "k", count: 1 }).onConflictDoUpdate({ target: { constraint: "c" }, set: { count: 0 } }) +// @ts-expect-error not a column +CH.insertInto(Counters).values({ key: "k", count: 1 }).onConflictDoNothing({ target: ["nope"] }) +// @ts-expect-error a value of another type +CH.insertInto(Counters).values({ key: "k", count: 1 }).onConflictDoUpdate({ target: ["key"], set: { count: "1" } }) +// @ts-expect-error DO UPDATE needs a target +CH.insertInto(Counters).values({ key: "k", count: 1 }).onConflictDoUpdate({ set: { count: 0 } }) +// @ts-expect-error a computed column cannot be set +CH.insertInto(Spans).values({ OrgId: "o", Label: "l" }).onConflictDoUpdate({ target: ["OrgId"], set: { Day: "x" } }) + // An insert without RETURNING runs to no rows; compile gives a CompiledQuery. const insert = CH.insertInto(Plain).values({ A: "a", B: 1 }) expectTypeOf>().toEqualTypeOf() diff --git a/src/ch/insert.test.ts b/src/ch/insert.test.ts index cc36dac..c8864ad 100644 --- a/src/ch/insert.test.ts +++ b/src/ch/insert.test.ts @@ -144,6 +144,74 @@ describe("insertInto", () => { ) }) + describe("on conflict", () => { + const Counters = CH.table("counters", { org: PG.text, key: PG.text, count: PG.int8, locked: PG.bool }, { tenantColumn: "org" }) + const row = { org: "o", key: "k", count: 1, locked: false } + + it("DO NOTHING, with and without a target", () => { + const insert = CH.insertInto(Counters).values(row) + expect(PG.compileUnsafe(insert.onConflictDoNothing()).sql).toMatch(/\nON CONFLICT DO NOTHING$/) + expect(PG.compileUnsafe(insert.onConflictDoNothing({ target: ["org", "key"] })).sql).toMatch( + /\nON CONFLICT \("org", "key"\) DO NOTHING$/, + ) + expect(PG.compileUnsafe(insert.onConflictDoNothing({ target: { constraint: "counters_pkey" } })).sql).toMatch( + /\nON CONFLICT ON CONSTRAINT "counters_pkey" DO NOTHING$/, + ) + expect( + PG.compileUnsafe(insert.onConflictDoNothing({ target: ["key"], targetWhere: ($) => $.locked.eq(false) })).sql, + ).toMatch(/\nON CONFLICT \("key"\) WHERE "locked" = FALSE DO NOTHING$/) + }) + + it("DO UPDATE with excluded, a qualified existing row, values and a WHERE", () => { + const compiled = PG.compileUnsafe( + CH.insertInto(Counters) + .values(row) + .onConflictDoUpdate({ + target: ["org", "key"], + set: ($, excluded) => ({ count: $.count.add(excluded.count), locked: true }), + where: ($) => $.locked.eq(false), + }) + .returning("count"), + ) + expect(compiled.sql).toBe( + 'INSERT INTO "counters" ("org", "key", "count", "locked")\nVALUES ($1, $2, $3, $4)\n' + + 'ON CONFLICT ("org", "key") DO UPDATE SET "count" = "counters"."count" + "excluded"."count", "locked" = $5 ' + + 'WHERE "counters"."locked" = FALSE\nRETURNING "count" AS "count"', + ) + expect(compiled.parameters).toEqual(["o", "k", 1, false, true]) + }) + + it("a SET that writes another tenant makes the insert cross-tenant", () => { + const insert = CH.insertInto(Counters).values(row) + const scope = (org: string) => + PG.compileUnsafe(insert.onConflictDoUpdate({ target: ["key"], set: { org } })).tenantScope + expect(scope("o")).toBe("single-tenant") + expect(scope("p")).toBe("cross-tenant") + }) + + it.effect("a bad SET fails; misuse and ClickHouse are defects", () => + Effect.gen(function* () { + const insert = CH.insertInto(Counters).values(row) + const empty = yield* Effect.flip( + PG.compile(insert.onConflictDoUpdate({ target: ["key"], set: { count: undefined } })), + ) + expect(empty.message).toContain("sets no columns") + const unknown = yield* Effect.flip( + PG.compile(insert.onConflictDoUpdate({ target: ["key"], set: { nope: 1 } as any })), + ) + expect(unknown.code).toBe("InvalidArguments") + for (const bad of [ + PG.compile(insert.onConflictDoNothing({ target: [] })), + PG.compile(insert.onConflictDoNothing({ targetWhere: ($) => $.locked.eq(false) })), + PG.compile(insert.onConflictDoNothing({ target: { constraint: "c" }, targetWhere: ($) => $.locked.eq(false) })), + CH.compile(insert.onConflictDoNothing()), + ]) { + expect(failure(yield* Effect.exit(bad))).toBeInstanceOf(QueryBuilderDefect) + } + }), + ) + }) + describe("tenant scope", () => { const scope = (rows: ReadonlyArray>, params: Record = {}) => CH.compileUnsafe(CH.insertInto(Events).values(rows as any), params).tenantScope diff --git a/src/ch/insert.ts b/src/ch/insert.ts index 7107b9c..932ae2a 100644 --- a/src/ch/insert.ts +++ b/src/ch/insert.ts @@ -13,7 +13,7 @@ // }) // yield* Database.run(insert, { id, orgId }) -import type { Comparable, Expr, Widen } from "./expr" +import type { Comparable, Condition, Expr, Widen } from "./expr" import type { ColumnAccessor, InferOutput } from "./query" import type { Table } from "./table" import type { CHType, ColumnDefs, InferTS } from "./types" @@ -65,6 +65,51 @@ export type InsertRow = T extends Table ? InsertRow : never +/** + * The unique index or constraint a conflict is detected on: column names + * (Postgres infers the index from them), or a constraint by name. + */ +export type ConflictTarget = + | ReadonlyArray + | { readonly constraint: string } + +/** + * What `ON CONFLICT DO UPDATE` writes into the existing row: any insertable + * column, as a value, param or expression. A key left out (or `undefined`) + * keeps the existing value. + */ +export type ConflictSet = { + readonly [K in Exclude>]?: InsertValue | undefined +} + +export interface OnConflictDoNothing { + /** Omit to skip a row that conflicts on any unique index or constraint. */ + readonly target?: ConflictTarget + /** The partial index's predicate, for a target on a partial unique index. */ + readonly targetWhere?: ($: ColumnAccessor) => Condition +} + +export interface OnConflictDoUpdate { + /** Required: Postgres must know which index the update is for. */ + readonly target: ConflictTarget + readonly targetWhere?: ($: ColumnAccessor) => Condition + /** + * The columns to write into the existing row. As a callback, `$` is the + * existing row and `excluded` the row that was proposed for insertion: + * `(($, excluded) => ({ count: $.count.add(excluded.count) }))`. + */ + readonly set: + | ConflictSet + | (($: ColumnAccessor, excluded: ColumnAccessor) => ConflictSet) + /** Update only the existing rows this holds for; the others are skipped. */ + readonly where?: ($: ColumnAccessor, excluded: ColumnAccessor) => Condition +} + +/** @internal — what an insert does on conflict. */ +export type ConflictClause = + | ({ readonly action: "nothing" } & OnConflictDoNothing) + | ({ readonly action: "update" } & OnConflictDoUpdate) + /** @internal — runtime insert state */ export interface CHInsertState { readonly table: Table @@ -72,6 +117,8 @@ export interface CHInsertState { readonly rows?: ReadonlyArray>> /** Set by `returning`: the RETURNING list, as a select callback. */ readonly returningFn?: ($: any) => Record> + /** Set by `onConflictDoNothing` / `onConflictDoUpdate`. */ + readonly conflict?: ConflictClause } export interface CHInsert< @@ -108,6 +155,19 @@ export interface CHInsert< returning>>( fn: ($: ColumnAccessor) => S, ): CHInsert> + + /** + * `ON CONFLICT DO NOTHING`: skip a row that conflicts. With `returning`, a + * skipped row returns nothing. Postgres only; replaces any earlier + * `onConflict*`. + */ + onConflictDoNothing(options?: OnConflictDoNothing): CHInsert + + /** + * `ON CONFLICT (target) DO UPDATE SET ...`: an upsert. Postgres only; + * replaces any earlier `onConflict*`. + */ + onConflictDoUpdate(options: OnConflictDoUpdate): CHInsert } const makeInsert = ( @@ -129,6 +189,8 @@ const makeInsert = Object.fromEntries((args as ReadonlyArray).map((column) => [column, $[column]])) return makeInsert({ ...state, returningFn }) }) as CHInsert["returning"], + onConflictDoNothing: (options = {}) => makeInsert({ ...state, conflict: { action: "nothing", ...options } }), + onConflictDoUpdate: (options) => makeInsert({ ...state, conflict: { action: "update", ...options } }), }) /** Start an INSERT into `table`. Give its rows with `values`. */ diff --git a/src/database/database.test.ts b/src/database/database.test.ts index 631e33d..010ffc7 100644 --- a/src/database/database.test.ts +++ b/src/database/database.test.ts @@ -419,6 +419,42 @@ layer(Live, { excludeTestServices: true })("Database on PGlite", (it) => { }), ) + it.effect("upserts with onConflictDoUpdate and skips with onConflictDoNothing", () => + Effect.gen(function* () { + yield* Db.execute( + Db.sql`CREATE TABLE counters (key text PRIMARY KEY, count int8 NOT NULL, locked boolean NOT NULL DEFAULT false)`, + ) + const Counters = CH.table("counters", { key: PG.text, count: PG.int8, locked: PG.bool }, { defaults: ["locked"] }) + const bump = (key: string, by: number) => + Db.run( + CH.insertInto(Counters) + .values({ key, count: by }) + .onConflictDoUpdate({ + target: ["key"], + set: ($, excluded) => ({ count: $.count.add(excluded.count) }), + where: ($) => $.locked.eq(false), + }) + .returning("key", "count"), + ) + expect(yield* bump("a", 1)).toEqual([{ key: "a", count: 1 }]) + expect(yield* bump("a", 2)).toEqual([{ key: "a", count: 3 }]) + yield* Db.execute(Db.sql`UPDATE counters SET locked = true WHERE key = 'a'`) + // The WHERE skips a locked row: nothing is updated, so nothing returns. + expect(yield* bump("a", 5)).toEqual([]) + const skipped = yield* Db.run( + CH.insertInto(Counters) + .values([{ key: "a", count: 100 }, { key: "b", count: 1 }]) + .onConflictDoNothing({ target: ["key"] }) + .returning("key"), + ) + expect(skipped).toEqual([{ key: "b" }]) + expect(yield* Db.run(CH.from(Counters).select("key", "count").orderBy(["key", "asc"]))).toEqual([ + { key: "a", count: 3 }, + { key: "b", count: 1 }, + ]) + }), + ) + it.effect("an insert inside a failed transaction rolls back", () => Effect.gen(function* () { const table = yield* freshTable diff --git a/src/pg/dialect.ts b/src/pg/dialect.ts index b8d497e..a6fd259 100644 --- a/src/pg/dialect.ts +++ b/src/pg/dialect.ts @@ -103,6 +103,7 @@ export const postgresDialect: Dialect = { groupByAlias: false, parenthesizeUnionBranches: true, returning: true, + onConflict: true, }, paramCodecs: { bool: Schema.Boolean,