diff --git a/CHANGELOG.md b/CHANGELOG.md index d26c317..7ca95a5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,11 @@ - 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 `select(query)` to an insert: `INSERT ... SELECT` from a query or union, its row checked + 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. - 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 6a33e6f..4ca6381 100644 --- a/design/writes.md +++ b/design/writes.md @@ -1,7 +1,8 @@ # Writes: INSERT -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 +Status: phases 1 to 4 built (`values`, `returning`, `onConflictDoNothing` / +`onConflictDoUpdate`, `select`, `settings`; see `docs/inserts.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). @@ -250,3 +251,10 @@ note, reusing `kind: "write"`, the returning path and the `set` record type from `.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`. +- **`INSERT ... SELECT` checks types with its own `InsertSelectMisfits` / `InsertSelectMissing`** + rather than moving `MisfitColumns` out of `schema/define.ts`: an insert also has to exclude + computed columns and require the ones without a default, which a materialized view does not. + 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. diff --git a/docs/inserts.md b/docs/inserts.md index 07389fa..b15918d 100644 --- a/docs/inserts.md +++ b/docs/inserts.md @@ -90,6 +90,37 @@ the codec rejects fails to compile with a `QueryBuilderError` that names the row several rows is bound once. A statement over 65535 bound values (Postgres's limit) fails to compile instead of being split: send fewer rows per statement. +## Insert ... select + +`select(query)` inserts the rows a query (or a `unionAll`) selects, instead of `values`. Each +selected alias names the column it goes into, so select under the target's column names: + +```ts +const Spans = CH.table("spans", { OrgId: CH.string, Name: CH.string, Ms: CH.uint64 }, { tenantColumn: "OrgId" }) +const Daily = CH.table("daily", { OrgId: CH.string, Name: CH.string, Total: CH.uint64 }, { tenantColumn: "OrgId" }) + +CH.insertInto(Daily).select( + CH.from(Spans) + .select(($) => ({ OrgId: $.OrgId, Name: $.Name, Total: CH.sum($.Ms) })) + .where(($) => [$.OrgId.eq(CH.param.string("orgId"))]) + .groupBy("OrgId", "Name"), +) +// INSERT INTO daily (OrgId, Name, Total) +// SELECT ... FROM spans WHERE spans.OrgId = 'o1' GROUP BY OrgId, Name +``` + +The selected row is checked against the table: selecting a column the table does not have (or +a computed one), selecting a value of another type, or leaving out a required column is a type +error naming the columns (`targetCannotTake`, `missingColumns`). A nullable result, such as a +Postgres `sum`, does not fit a NOT NULL column; wrap it in `coalesce`. + +The column list is the aliases in select order, which is the order the SELECT writes them, so +ClickHouse and Postgres agree. `returning` and `onConflict*` work with `select` as with +`values`. `select` and `values` replace each other. + +A long `INSERT ... SELECT` on ClickHouse (a backfill over a big table) can outlast an HTTP +timeout; run those in slices. + ## Returning On Postgres, `returning` adds a RETURNING list and `Database.run` returns the inserted rows, @@ -149,12 +180,31 @@ RETURNING "key" AS "key", "count" AS "count" - Calling either again replaces the clause. ClickHouse has no `ON CONFLICT` (deduplicate with a `ReplacingMergeTree` instead), so compiling one for it is a `QueryBuilderDefect`. +## Settings + +On ClickHouse, `settings` adds a `SETTINGS` clause to the insert, before `VALUES` or the +`SELECT`: + +```ts +CH.insertInto(Daily).values(rows).settings({ async_insert: 1, wait_for_async_insert: 1 }) +// INSERT INTO daily (OrgId, Name, Total) SETTINGS async_insert = 1, wait_for_async_insert = 1 +// VALUES ... +``` + +Names must be plain identifiers; values (strings, numbers, booleans) are written as literals. +Calling it again replaces them. Postgres has no insert settings, so compiling one with +`settings` 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. An -`onConflictDoUpdate` that sets the tenant column counts as one more row. A table +`onConflictDoUpdate` that sets the tenant column counts as one more row. + +An `INSERT ... SELECT` into a tenant table is `"single-tenant"` when the SELECT is, and each row +takes its tenant from a tenant column of the source or from the same param that pins the SELECT. +Into a table without a tenant column, the insert has the SELECT's scope: what it reads. A table without a tenant column gives `"untenanted"`. ## Failures @@ -170,6 +220,7 @@ without a tenant column gives `"untenanted"`. | `returning` for a dialect without RETURNING | `QueryBuilderDefect` | | `onConflictDoUpdate` setting no or unknown columns | `QueryBuilderError` `InvalidArguments` | | `onConflict*` without ON CONFLICT, a bad target | `QueryBuilderDefect` | +| `settings` without insert settings, a bad name | `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 1d96d21..f6f1085 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, `.onConflictDoNothing(options?)` / `.onConflictDoUpdate(options)` the ON CONFLICT clause (Postgres). See [Inserting rows](./inserts.md) | +| `insertInto` | `(table) => 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 +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`, `ConflictTarget`, `ConflictSet`, `OnConflictDoNothing`, `OnConflictDoUpdate`, `ColumnAccessor`, `JoinedColumnAccessor`, +`ParamKind`, `CHQuery`, `CHUnionQuery`, `CHInsert`, `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/src/ch/compile.ts b/src/ch/compile.ts index 585749d..3d1d288 100644 --- a/src/ch/compile.ts +++ b/src/ch/compile.ts @@ -12,7 +12,7 @@ import type { CHQuery, CHQueryState } from "./query" import type { CHUnionQuery } from "./union" import { isInsert, type CHInsert } from "./insert" import { createColumnAccessor, createQualifiedColumnAccessor, createJoinedColumnAccessor, sourceAlias } from "./query" -import { aliased, columnTypeOf, isExprLike } from "./expr" +import { aliased, columnTypeOf, isExprLike, type Expr } from "./expr" import { raw, identPath, quoteIdent, quoteIdentPath, compile as compileSqlFragment } from "../sql/sql-fragment" import { splitTerminalClauses } from "../sql/terminal-clauses" import { compileQuery, type SqlQuery } from "../sql/sql-query" @@ -23,7 +23,7 @@ import { checkedLiteral, clickhouseDialect, currentDialect, withDialect, type Di import { Effect, Option, Schema } from "effect" import { QueryBuilderDefect, QueryBuilderError } from "./errors" import { withSubqueryCompiler } from "./subquery-context" -import { tenantBoundOf, tenantPredicatesOf, withTenantBound, type TenantPredicate } from "./tenant" +import { tenantBoundOf, tenantColumnOf, tenantPredicatesOf, withTenantBound, type TenantPredicate } from "./tenant" // `QueryBuilderError` moved to ./errors so `expr.ts` can raise it too; still // exported from here, which is where every caller imports it from. @@ -1420,43 +1420,69 @@ const returningOf = (insert: CHInsert) => { } /** - * An INSERT ... VALUES. Columns are written in table order, so two rows with - * their keys in different orders cannot swap values; a column some rows leave - * out is `DEFAULT` in those rows, which both dialects accept. + * The column list and SQL of an `INSERT ... SELECT` source. The selected + * aliases name the columns, in select order, which is also the order the + * SELECT writes them in, so positional and named matching agree. */ -function compileInsert(insert: CHInsert, params: Record): CompiledQuery { - const { table, rows } = insert._state - const where = `insertInto(${table.name})` - if (rows === undefined) throw new QueryBuilderDefect({ message: `${where}: values() is required` }) - // The rows usually come from data, so their number and keys are failures, not defects. - if (rows.length === 0) { - throw new QueryBuilderError({ code: "InvalidArguments", message: `${where}: values() was given no rows` }) - } +const insertSelectSource = ( + insert: CHInsert, + query: CHQuery | CHUnionQuery, + values: Record, +) => { + const { table } = insert._state + const isUnion = "_tag" in query && query._tag === "CHUnionQuery" + const exprs = selectExprsOf(isUnion ? (query as CHUnionQuery)._state.queries[0]! : (query as CHQuery)) ?? {} + const columns = Object.keys(exprs) const computed = new Set(table.computed ?? []) - const present = new Set() - rows.forEach((row, index) => { - for (const [column, value] of Object.entries(row)) { - if (value === undefined) continue - if (!Object.hasOwn(table.columns, column)) { - throw new QueryBuilderError({ - code: "InvalidArguments", - message: `${where}: row ${index} has ${JSON.stringify(column)}, which is not a column of the table`, - }) - } - if (computed.has(column)) { - throw new QueryBuilderError({ - code: "InvalidArguments", - message: `${where}: row ${index} writes ${column}, which the database computes (MATERIALIZED or ALIAS)`, - }) - } - present.add(column) + // The query is written in source, so a column the table cannot take is a + // defect, as the type error on `select` says. + for (const column of columns) { + if (!Object.hasOwn(table.columns, column) || computed.has(column)) { + throw new QueryBuilderDefect({ + message: `insertInto(${table.name}): the query selects ${JSON.stringify(column)}, which is not an insertable column of the table`, + }) } - }) - const columns = Object.keys(table.columns).filter((column) => present.has(column)) - if (columns.length === 0) { - throw new QueryBuilderError({ code: "InvalidArguments", message: `${where}: every row is empty; give at least one column` }) } + const inner = isUnion + ? compileUnionInner(query as CHUnionQuery, values, { nested: true }) + : compileInner(query as CHQuery, values, { skipFormat: true, nested: true }) + const sql = isUnion && currentDialect().clauses.format ? splitTerminalClauses(inner.sql).body : inner.sql + return { columns, exprs, inner, sql } +} +/** `SETTINGS a = 1, b = 'x'`, or `""` without settings. */ +const insertSettingsClause = (insert: CHInsert): string => { + const { table, settings } = insert._state + const entries = Object.entries(settings ?? {}) + if (entries.length === 0) return "" + const dialect = currentDialect() + if (dialect.clauses.insertSettings !== true) { + throw new QueryBuilderDefect({ + message: `insertInto(${table.name}): settings() has no meaning for the ${dialect.name} dialect, which has no INSERT SETTINGS`, + }) + } + 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` }) + } + return `${name} = ${checkedLiteral(dialect, value, `setting ${name}`)}` + }) + .join(", ")}` +} + +/** + * An INSERT ... VALUES or INSERT ... SELECT. For VALUES, columns are written in + * table order, so two rows with their keys in different orders cannot swap + * values, and a column some rows leave out is `DEFAULT` in those rows, which + * both dialects accept. + */ +function compileInsert(insert: CHInsert, params: Record): CompiledQuery { + const { table, rows, selectQuery } = insert._state + const where = `insertInto(${table.name})` + if (rows === undefined && selectQuery === undefined) { + throw new QueryBuilderDefect({ message: `${where}: values() or select() is required` }) + } const dialect = currentDialect() const values: Record = { ...params } let next = 0 @@ -1470,34 +1496,64 @@ function compileInsert(insert: CHInsert, params: Record() - let pinned = tenant !== undefined && present.has(tenant) + let pinned = tenant !== undefined 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, conflictSql] = withSubqueryCompiler( + const [columns, source, conflictSql, readScope] = withSubqueryCompiler( (subquery) => typeof subquery === "string" ? subquery : compileInner(subquery, values, { skipFormat: true, nested: true }).sql, () => { - const tuples = rows.map((row, index) => { - const cells = columns.map((column) => { - const value = row[column] - const sql = cell(column, value, `row ${index}`) - if (column === tenant) pin(value, sql) - return sql + let columns: ReadonlyArray + let source: string + // What the statement reads, besides what it writes: an INSERT ... SELECT + // into any table reads with the SELECT's scope. + let readScope: TenantScope | undefined + if (rows !== undefined) { + columns = valuesColumns(table, rows) + if (tenant !== undefined && !columns.includes(tenant)) pinned = false + const tuples = rows.map((row, index) => { + const cells = columns.map((column) => { + const value = row[column] + const sql = cell(column, value, `row ${index}`) + if (column === tenant) pin(value, sql) + return sql + }) + return `(${cells.join(", ")})` }) - return `(${cells.join(", ")})` - }) - return [tuples, onConflictClause(insert, cell, (column, value, sql) => { + source = `VALUES ${tuples.join(", ")}` + } else { + const selected = insertSelectSource(insert, selectQuery!, values) + columns = selected.columns + source = selected.sql + readScope = selected.inner.tenantScope + const bound = tenantBoundOf(selected.inner) + if (tenant !== undefined) { + // The written tenant is pinned when the read is, and the row takes its + // tenant from a tenant column of the source or from that same value. + const expr = selected.exprs[tenant] as Expr | undefined + const fromSource = expr !== undefined && isExprLike(expr) && tenantColumnOf(expr) !== undefined + const sameValue = + expr !== undefined && + isExprLike(expr) && + "_paramName" in expr && + inlineParams(compileSqlFragment(expr.toFragment()), values) === bound + if (readScope !== "single-tenant" || bound === undefined || !(fromSource || sameValue)) pinned = false + else bounds.add(bound) + } + } + const conflictSql = onConflictClause(insert, cell, (column, value, sql) => { if (column === tenant) pin(value, sql) - })] as const + }) + return [columns, source, conflictSql, readScope] as const }, ) @@ -1507,13 +1563,19 @@ function compileInsert(insert: CHInsert, params: Record compileSqlFragment(aliased(returning.exprs[alias]!, alias))).join(", ")}` const rendered = renderParams( - `INSERT INTO ${quoteIdentPath(table.name)} (${columns.map(quoteIdent).join(", ")})\nVALUES ${tuples.join(", ")}${conflictSql}${returningSql}`, + `INSERT INTO ${quoteIdentPath(table.name)} (${columns.map(quoteIdent).join(", ")})${insertSettingsClause(insert)}\n${source}${conflictSql}${returningSql}`, values, dialect, ) - const returnedSchema = returning !== undefined && "schema" in returning.derived ? returning.derived.schema : undefined const tenantScope: TenantScope = - tenant === undefined ? "untenanted" : pinned && bounds.size === 1 ? "single-tenant" : "cross-tenant" + tenant === undefined + ? (readScope ?? "untenanted") + : pinned && bounds.size === 1 + ? "single-tenant" + : "cross-tenant" + const tenantBound = + tenantScope !== "single-tenant" ? undefined : tenant === undefined ? undefined : [...bounds][0] + const returnedSchema = returning !== undefined && "schema" in returning.derived ? returning.derived.schema : undefined return withTenantBound( makeCompiledQuery( @@ -1530,6 +1592,43 @@ function compileInsert(insert: CHInsert, params: Record["_state"]["table"], + rows: ReadonlyArray>>, +): ReadonlyArray => { + const where = `insertInto(${table.name})` + // The rows usually come from data, so their number and keys are failures, not defects. + if (rows.length === 0) { + throw new QueryBuilderError({ code: "InvalidArguments", message: `${where}: values() was given no rows` }) + } + const computed = new Set(table.computed ?? []) + const present = new Set() + rows.forEach((row, index) => { + for (const [column, value] of Object.entries(row)) { + if (value === undefined) continue + if (!Object.hasOwn(table.columns, column)) { + throw new QueryBuilderError({ + code: "InvalidArguments", + message: `${where}: row ${index} has ${JSON.stringify(column)}, which is not a column of the table`, + }) + } + if (computed.has(column)) { + throw new QueryBuilderError({ + code: "InvalidArguments", + message: `${where}: row ${index} writes ${column}, which the database computes (MATERIALIZED or ALIAS)`, + }) + } + present.add(column) + } + }) + const columns = Object.keys(table.columns).filter((column) => present.has(column)) + if (columns.length === 0) { + throw new QueryBuilderError({ code: "InvalidArguments", message: `${where}: every row is empty; give at least one column` }) + } + return columns +} diff --git a/src/ch/dialect.ts b/src/ch/dialect.ts index 089f65d..70229f4 100644 --- a/src/ch/dialect.ts +++ b/src/ch/dialect.ts @@ -65,6 +65,9 @@ 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 } /** A transaction isolation level. `read uncommitted` is left out: Postgres runs it as `read committed`. */ @@ -151,6 +154,7 @@ export const clickhouseDialect: Dialect = { parenthesizeUnionBranches: false, returning: false, onConflict: false, + insertSettings: true, }, transactions: noTransactions, } diff --git a/src/ch/index.ts b/src/ch/index.ts index 82e8e1f..92183b2 100644 --- a/src/ch/index.ts +++ b/src/ch/index.ts @@ -246,6 +246,9 @@ export { type ConflictTarget, type InsertRow, type InsertRowOf, + type InsertSelectMisfits, + type InsertSelectMissing, + type InsertSettingValue, type InsertValue, type OnConflictDoNothing, type OnConflictDoUpdate, diff --git a/src/ch/insert.test-d.ts b/src/ch/insert.test-d.ts index 4fee379..2d5e20d 100644 --- a/src/ch/insert.test-d.ts +++ b/src/ch/insert.test-d.ts @@ -110,6 +110,28 @@ CH.insertInto(Counters).values({ key: "k", count: 1 }).onConflictDoUpdate({ set: // @ts-expect-error a computed column cannot be set CH.insertInto(Spans).values({ OrgId: "o", Label: "l" }).onConflictDoUpdate({ target: ["OrgId"], set: { Day: "x" } }) +// INSERT ... SELECT: the selected row must fit the table. +const Daily = CH.table("daily", { OrgId: CH.string, Total: CH.uint64, Note: CH.nullable(CH.string) }) +const Source = CH.table("source", { OrgId: CH.string, Ms: CH.uint64, Label: CH.string }) +CH.insertInto(Daily).select(CH.from(Source).select(($) => ({ OrgId: $.OrgId, Total: $.Ms }))) +CH.insertInto(Daily).select(CH.from(Source).select(($) => ({ OrgId: $.OrgId, Total: $.Ms, Note: $.Label }))) +CH.insertInto(Daily).select( + CH.unionAll(CH.from(Source).select(($) => ({ OrgId: $.OrgId, Total: $.Ms })), CH.from(Source).select(($) => ({ OrgId: $.OrgId, Total: $.Ms }))), +) +// @ts-expect-error Total is required and not selected +CH.insertInto(Daily).select(CH.from(Source).select("OrgId")) +// @ts-expect-error Label is not a column of daily +CH.insertInto(Daily).select(CH.from(Source).select(($) => ({ OrgId: $.OrgId, Total: $.Ms, Label: $.Label }))) +// @ts-expect-error Total is a number, not a string +CH.insertInto(Daily).select(CH.from(Source).select(($) => ({ OrgId: $.OrgId, Total: $.Label }))) +// @ts-expect-error a computed column cannot be selected into +CH.insertInto(Spans).select(CH.from(Spans).select("OrgId", "Label", "Day")) + +// settings take plain values. +CH.insertInto(Daily).values({ OrgId: "o", Total: 1 }).settings({ async_insert: 1, wait_for_async_insert: true }) +// @ts-expect-error not a setting value +CH.insertInto(Daily).values({ OrgId: "o", Total: 1 }).settings({ async_insert: [1] }) + // 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 c8864ad..1df0daf 100644 --- a/src/ch/insert.test.ts +++ b/src/ch/insert.test.ts @@ -212,6 +212,94 @@ describe("insertInto", () => { ) }) + describe("insert ... select", () => { + const Spans = CH.table("spans", { OrgId: CH.string, Name: CH.string, Ms: CH.uint64 }, { tenantColumn: "OrgId" }) + const Daily = CH.table("daily", { OrgId: CH.string, Name: CH.string, Total: CH.uint64 }, { tenantColumn: "OrgId" }) + + it("names the columns from the selected aliases, in select order", () => { + const compiled = CH.compileUnsafe( + CH.insertInto(Daily).select( + CH.from(Spans) + .select(($) => ({ Total: CH.sum($.Ms), OrgId: $.OrgId, Name: $.Name })) + .where(($) => [$.OrgId.eq(CH.param.string("org"))]) + .groupBy("OrgId", "Name"), + ), + { org: "o" }, + ) + expect(compiled.sql).toMatch(/^INSERT INTO daily \(Total, OrgId, Name\)\nSELECT\s+sum\(spans\.Ms\) AS Total,\s+spans\.OrgId AS OrgId,/) + expect(compiled.sql).toMatch(/WHERE spans\.OrgId = 'o'\s+GROUP BY OrgId, Name$/) + expect(compiled.tenantScope).toBe("single-tenant") + }) + + it("binds the query's params on Postgres, with values and select sharing the numbering", () => { + const Src = CH.table("src", { org: PG.text, n: PG.int4 }) + const Dst = CH.table("dst", { org: PG.text, n: PG.int4 }) + const compiled = PG.compileUnsafe( + CH.insertInto(Dst) + .select(CH.from(Src).select("org", "n").where(($) => [$.org.eq(CH.param.string("org"))])) + .onConflictDoUpdate({ target: ["org"], set: { n: CH.param.int("n") } }), + { org: "o", n: 3 }, + ) + expect(compiled.sql).toMatch(/^INSERT INTO "dst" \("org", "n"\)\nSELECT[\s\S]*WHERE "src"\."org" = \$1\nON CONFLICT \("org"\) DO UPDATE SET "n" = \$2$/) + expect(compiled.parameters).toEqual(["o", 3]) + }) + + it("takes a union", () => { + const branch = (org: string) => + CH.from(Spans).select(($) => ({ OrgId: $.OrgId, Name: $.Name, Total: $.Ms })).where(($) => [$.OrgId.eq(org)]) + const compiled = CH.compileUnsafe(CH.insertInto(Daily).select(CH.unionAll(branch("a"), branch("b")))) + expect(compiled.sql).toMatch(/^INSERT INTO daily \(OrgId, Name, Total\)\nSELECT[\s\S]*UNION ALL[\s\S]*'b'$/) + expect(compiled.tenantScope).toBe("cross-tenant") + }) + + it("tenant scope follows the read and where the tenant column comes from", () => { + const scope = (query: CH.CHQuery, params: Record = { org: "o" }) => + CH.compileUnsafe(CH.insertInto(Daily).select(query as any), params).tenantScope + const pinned = CH.from(Spans).where(($) => [$.OrgId.eq(CH.param.string("org"))]) + expect(scope(pinned.select(($) => ({ OrgId: $.OrgId, Name: $.Name, Total: $.Ms })))).toBe("single-tenant") + expect(scope(pinned.select(($) => ({ OrgId: CH.param.string("org"), Name: $.Name, Total: $.Ms })))).toBe("single-tenant") + // Pinned read, but the rows are written to another tenant. + expect(scope(pinned.select(($) => ({ OrgId: CH.param.string("other"), Name: $.Name, Total: $.Ms })), { org: "o", other: "p" })).toBe( + "cross-tenant", + ) + expect(scope(CH.from(Spans).select(($) => ({ OrgId: $.OrgId, Name: $.Name, Total: $.Ms })))).toBe("cross-tenant") + // An untenanted target reads with the query's scope. + const Names = CH.table("names", { Name: CH.string }) + expect(CH.compileUnsafe(CH.insertInto(Names).select(pinned.select("Name")), { org: "o" }).tenantScope).toBe("single-tenant") + expect(CH.compileUnsafe(CH.insertInto(Names).select(CH.from(Spans).select("Name"))).tenantScope).toBe("cross-tenant") + }) + + it("select and values replace each other", () => { + const fromQuery = CH.insertInto(Daily).select(CH.from(Daily).select("OrgId", "Name", "Total")) + expect(CH.compileUnsafe(fromQuery.values({ OrgId: "o", Name: "n", Total: 1 })).sql).toContain("VALUES") + expect(CH.compileUnsafe(fromQuery.values({ OrgId: "o", Name: "n", Total: 1 }).select(CH.from(Daily).select("OrgId", "Name", "Total"))).sql).toContain( + "SELECT", + ) + }) + }) + + describe("settings", () => { + const Notes = CH.table("notes", { Body: CH.string }) + + it("writes SETTINGS before VALUES on ClickHouse", () => { + const compiled = CH.compileUnsafe( + CH.insertInto(Notes).values({ Body: "x" }).settings({ async_insert: 1, wait_for_async_insert: true, insert_deduplication_token: "t'1" }), + ) + expect(compiled.sql).toBe( + "INSERT INTO notes (Body) SETTINGS async_insert = 1, wait_for_async_insert = 1, insert_deduplication_token = 't\\'1'\nVALUES ('x')", + ) + }) + + it.effect("a bad name and Postgres are defects", () => + Effect.gen(function* () { + const insert = CH.insertInto(Notes).values({ Body: "x" }) + expect(failure(yield* Effect.exit(CH.compile(insert.settings({ "a b": 1 }))))).toBeInstanceOf(QueryBuilderDefect) + expect(failure(yield* Effect.exit(PG.compile(insert.settings({ a: 1 }))))).toBeInstanceOf(QueryBuilderDefect) + expect(PG.compileUnsafe(insert.settings({})).sql).not.toContain("SETTINGS") + }), + ) + }) + 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 932ae2a..f4a4b4c 100644 --- a/src/ch/insert.ts +++ b/src/ch/insert.ts @@ -14,8 +14,9 @@ // yield* Database.run(insert, { id, orgId }) import type { Comparable, Condition, Expr, Widen } from "./expr" -import type { ColumnAccessor, InferOutput } from "./query" +import type { CHQuery, ColumnAccessor, InferOutput } from "./query" import type { Table } from "./table" +import type { CHUnionQuery } from "./union" import type { CHType, ColumnDefs, InferTS } from "./types" /** @@ -61,6 +62,47 @@ export type InsertRow +type RequiredKeys = Exclude< + keyof Cols & string, + OptionalKeys | Known +> + +/** The row a query or union selects. */ +type SelectedRow = Q extends { readonly _phantom?: { readonly output: infer Output } } ? Output : never + +/** Selected columns the table cannot take: not an insertable column, or of another type. */ +export type InsertSelectMisfits = { + [K in keyof Output]: K extends Exclude> + ? [Output[K]] extends [InferTS] + ? never + : K + : K +}[keyof Output] + +/** Required columns of the table the query does not select. */ +export type InsertSelectMissing< + Output, + Cols extends ColumnDefs, + Defaulted extends string = never, + Computed extends string = never, +> = Exclude, keyof Output> + +/** + * `unknown` when a query's row fits the table, otherwise a property naming + * what does not, so the error says which columns to fix. + */ +type SelectFits = ([ + InsertSelectMisfits, +] extends [never] + ? unknown + : { readonly targetCannotTake: InsertSelectMisfits }) & + ([InsertSelectMissing] extends [never] + ? unknown + : { readonly missingColumns: InsertSelectMissing }) + +/** A ClickHouse setting's value, as `SETTINGS name = value` writes it. */ +export type InsertSettingValue = string | number | boolean + /** The insert row of a table value: `InsertRowOf`. */ export type InsertRowOf = T extends Table ? InsertRow : never @@ -113,8 +155,12 @@ export type ConflictClause = /** @internal — runtime insert state */ export interface CHInsertState { readonly table: Table - /** Set by `values`. Compiling without it is a defect. */ + /** Set by `values`. Compiling without it or `selectQuery` is a defect. */ readonly rows?: ReadonlyArray>> + /** Set by `select`, which clears `rows` (and `values` clears it). */ + readonly selectQuery?: CHQuery | CHUnionQuery + /** Set by `settings`: ClickHouse `SETTINGS` for this insert. */ + readonly settings?: Readonly> /** Set by `returning`: the RETURNING list, as a select callback. */ readonly returningFn?: ($: any) => Record> /** Set by `onConflictDoNothing` / `onConflictDoUpdate`. */ @@ -142,6 +188,24 @@ export interface CHInsert< rows: InsertRow | ReadonlyArray>, ): CHInsert + /** + * `INSERT ... SELECT`: insert the rows a query (or union) selects. Each + * selected alias names the column it is written to, so the query must + * select every required column, and only columns the table can take, of + * their types. Replaces any `values`. + */ + select | CHUnionQuery>( + query: Q & SelectFits, Cols, Defaulted, Computed>, + ): CHInsert + + /** + * ClickHouse `SETTINGS` for this insert, such as + * `{ async_insert: 1, wait_for_async_insert: 1 }`. Names must be plain + * identifiers; values are written as literals. Calling it again replaces + * them. On a dialect without insert settings (Postgres) compiling is a defect. + */ + settings(settings: Readonly>): CHInsert + /** * Return the inserted rows: column names, or a callback building an * expression per alias, as in `select`. `Database.run` then decodes them @@ -178,9 +242,12 @@ const makeInsert = makeInsert({ ...state, + selectQuery: undefined, // Copied, so a caller pushing to its array later does not change the insert. rows: Array.isArray(rows) ? [...rows] : [rows as Readonly>], }), + select: (query) => makeInsert({ ...state, rows: undefined, selectQuery: query }), + settings: (settings) => makeInsert({ ...state, settings: { ...settings } }), returning: ((...args: ReadonlyArray) => { const [first] = args const returningFn = diff --git a/src/database/database.test.ts b/src/database/database.test.ts index 010ffc7..625ff00 100644 --- a/src/database/database.test.ts +++ b/src/database/database.test.ts @@ -455,6 +455,26 @@ layer(Live, { excludeTestServices: true })("Database on PGlite", (it) => { }), ) + it.effect("insert ... select copies rows, with ON CONFLICT and RETURNING", () => + Effect.gen(function* () { + yield* Db.execute(Db.sql`CREATE TABLE src (org text NOT NULL, n int4 NOT NULL)`) + yield* Db.execute(Db.sql`CREATE TABLE dst (org text PRIMARY KEY, total int8 NOT NULL)`) + yield* Db.execute(Db.sql`INSERT INTO src VALUES ('a', 1), ('a', 2), ('b', 5)`) + const Src = CH.table("src", { org: PG.text, n: PG.int4 }) + const Dst = CH.table("dst", { org: PG.text, total: PG.int8 }) + const rollup = CH.insertInto(Dst) + .select(CH.from(Src).select(($) => ({ org: $.org, total: CH.coalesce(PG.sum($.n), CH.lit(0)) })).groupBy("org")) + .onConflictDoUpdate({ target: ["org"], set: ($, excluded) => ({ total: $.total.add(excluded.total) }) }) + .returning("org", "total") + expect(yield* Db.run(rollup)).toHaveLength(2) + const again = yield* Db.run(rollup) + expect([...again].sort((x, y) => x.org.localeCompare(y.org))).toEqual([ + { org: "a", total: 6 }, + { org: "b", total: 10 }, + ]) + }), + ) + it.effect("an insert inside a failed transaction rolls back", () => Effect.gen(function* () { const table = yield* freshTable diff --git a/tests/database.clickhouse.test.ts b/tests/database.clickhouse.test.ts index 1d0a2eb..1aa89b9 100644 --- a/tests/database.clickhouse.test.ts +++ b/tests/database.clickhouse.test.ts @@ -99,6 +99,41 @@ describe("database", () => { expect(sent[1]).toMatch(/^INSERT INTO events \(OrgId, Id, At, Attrs, Note, Tags\)\nVALUES \('o1', DEFAULT, /) }) + it("runs insert ... select with settings", async () => { + const rows = await Effect.runPromise( + withDatabase((db) => + Effect.gen(function* () { + yield* db.execute(Db.sql`CREATE TABLE spans (OrgId String, Name String, Ms UInt64) ENGINE = MergeTree ORDER BY OrgId`) + yield* db.execute(Db.sql`CREATE TABLE daily (OrgId String, Name String, Total UInt64) ENGINE = MergeTree ORDER BY OrgId`) + const Spans = CH.table("spans", { OrgId: CH.string, Name: CH.string, Ms: CH.uint64 }, { tenantColumn: "OrgId" }) + const Daily = CH.table("daily", { OrgId: CH.string, Name: CH.string, Total: CH.uint64 }, { tenantColumn: "OrgId" }) + yield* db.run( + CH.insertInto(Spans) + .values([ + { OrgId: "o", Name: "a", Ms: 1 }, + { OrgId: "o", Name: "a", Ms: 2 }, + { OrgId: "p", Name: "b", Ms: 9 }, + ]) + .settings({ async_insert: 0 }), + ) + yield* db.run( + CH.insertInto(Daily) + .select( + CH.from(Spans) + .select(($) => ({ Total: CH.sum($.Ms), OrgId: $.OrgId, Name: $.Name })) + .where(($) => [$.OrgId.eq(CH.param.string("org"))]) + .groupBy("OrgId", "Name"), + ) + .settings({ max_threads: 1 }), + { org: "o" }, + ) + return yield* db.run(CH.from(Daily).select("OrgId", "Name", "Total")) + }), + ), + ) + expect(rows).toEqual([{ OrgId: "o", Name: "a", Total: 3 }]) + }) + it("refuses a transaction before sending anything", async () => { const result = await Effect.runPromise( withDatabase((db, sent) =>