From 0aa9f8f8e3322a1f85c8645692ec785a88a9812d Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 23:32:08 +0200 Subject: [PATCH 1/2] Add Postgres schema definitions and migrations Schema-as-code and migrations were ClickHouse-only: the entity model, DDL, diff, ledger and drift check were written against ClickHouse with no dialect seam, so a Postgres consumer (Maple) could not replace drizzle-kit. The dialect-neutral parts stay shared (snapshot envelope, branch graph, folder loading, MigrationDriver); each dialect now owns its entities, definitions, diff, DDL, ledger and drift check. - S.pg.table with column defaults, identity, composite primary keys, partial/expression indexes and foreign keys - dialect: "postgres" for generate/check/migrate/status/verify - one transaction per migration under an advisory lock - verify builds the expected schema in a rolled-back scratch schema and compares catalogs, so Postgres normalizes both sides - adopt a drizzle-kit folder: generate --baseline [--from-drizzle] and effect-orm baseline / Migrate.baseline Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 15 ++ design/gap-review.md | 2 +- design/migrations.md | 40 +++- docs/README.md | 4 +- docs/migrations.md | 104 +++++++++- src/kit/cli.ts | 41 +++- src/kit/generate.ts | 154 ++++++++++---- src/kit/graph.ts | 43 ++-- src/kit/kit.test.ts | 89 +++++++++ src/migrate.ts | 10 +- src/migrate/driver.ts | 11 +- src/migrate/pg-ledger.ts | 68 +++++++ src/migrate/pg-migrate.test.ts | 201 +++++++++++++++++++ src/migrate/pg-verify.ts | 232 +++++++++++++++++++++ src/migrate/run.ts | 119 ++++++++++- src/migrate/source.ts | 48 ++++- src/migrate/verify.ts | 42 +++- src/schema.ts | 45 ++++- src/schema/diff.ts | 4 +- src/schema/drizzle.ts | 175 ++++++++++++++++ src/schema/entities.ts | 58 ++++-- src/schema/ops.ts | 12 +- src/schema/pg-define.ts | 355 +++++++++++++++++++++++++++++++++ src/schema/pg-diff.ts | 161 +++++++++++++++ src/schema/pg-entities.ts | 128 ++++++++++++ src/schema/pg-ops.ts | 150 ++++++++++++++ src/schema/pg-schema.test.ts | 277 +++++++++++++++++++++++++ src/schema/snapshot.ts | 106 ++++++++-- 28 files changed, 2568 insertions(+), 126 deletions(-) create mode 100644 src/migrate/pg-ledger.ts create mode 100644 src/migrate/pg-migrate.test.ts create mode 100644 src/migrate/pg-verify.ts create mode 100644 src/schema/drizzle.ts create mode 100644 src/schema/pg-define.ts create mode 100644 src/schema/pg-diff.ts create mode 100644 src/schema/pg-entities.ts create mode 100644 src/schema/pg-ops.ts create mode 100644 src/schema/pg-schema.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index f31eaec..0bb5eb9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,21 @@ ## Unreleased +- Add Postgres schema definitions and migrations. `S.pg.table` (with `S.pg.column`, `S.pg.index`, + `S.pg.uniqueIndex`, `S.pg.foreignKey`) defines tables with primary keys, partial and expression + indexes, foreign keys, defaults and identity columns; `dialect: "postgres"` in the kit config + makes `generate`, `check`, `migrate`, `status` and `verify` work on Postgres. Each migration + runs in one transaction under an advisory lock, and `verify` compares the catalog with the + snapshot built in a rolled-back scratch schema. + - Adopt a drizzle-kit folder with `generate --baseline --from-drizzle` (or `--baseline` from + the definitions) and record an existing database with `effect-orm baseline ` / + `Migrate.baseline`. drizzle-kit `snapshot.json` files are recognized; migrations before the + first effect-orm snapshot are legacy and run as they are. + - `Snapshot` is a union over `dialect` (`ClickHouseSnapshot`, `PgSnapshot`); `MigrationFile` + over its generated files (`ClickHouseMigrationFile`, `PgMigrationFile`). `MigrationDriver` + gains an optional `transaction`, which `fromSqlClient` provides. `Drift.problem` adds + `not_null`, `identity` and `foreign_key`. + - **Breaking:** invalid queries are refused before any SQL is sent: as type errors where the type can see them, otherwise as a `QueryBuilderError` / `QueryBuilderDefect` from `compile`. - Params are in the query's type. `compile`, `compileUnion` and `Database.run` require every diff --git a/design/gap-review.md b/design/gap-review.md index 8586e00..852c5eb 100644 --- a/design/gap-review.md +++ b/design/gap-review.md @@ -34,7 +34,7 @@ builder; **P1** commonly used; **P2** niche. | ~~`isNull` / `isNotNull` / `between`; variadic `and` / `or` that skip `undefined`~~ (built) | everywhere | S | | Constraint error helpers (unique, foreign key, not null); keep ClickHouse's numeric error codes, which `sqlStateOf` drops today | all upserts | S | | Tenant-scope enforcement in `Database`, opt in, with an explicit cross-tenant entry point | safety | S | -| Postgres `defineTable` (indexes, unique, FKs), Postgres migrations, a drizzle-kit importer | 68 tables, 90 indexes, 47 unique, 75 folders | L; can wait, drizzle-kit can keep migrating | +| ~~Postgres `defineTable` (indexes, unique, FKs), Postgres migrations, a drizzle-kit importer~~ Done: `S.pg.table`, `dialect: "postgres"`, `--baseline --from-drizzle` (design/migrations.md section 8) | 68 tables, 90 indexes, 47 unique, 75 folders | L | ## P1 diff --git a/design/migrations.md b/design/migrations.md index 20f7e75..19656fd 100644 --- a/design/migrations.md +++ b/design/migrations.md @@ -1,6 +1,7 @@ # Migrations: schema-as-code, snapshots, and a migrator -Status: phases 0 to 3 implemented on 2026-10-03 (branch `feat/migrations`); phases 4 to 7 open. +Status: phases 0 to 3 implemented on 2026-10-03 (branch `feat/migrations`); phase 6 (Postgres) +on 2026-10-04, see section 8; phases 4, 5 and 7 open. Section 7 lists what was built and where it departs from this plan. User docs: [`docs/migrations.md`](../docs/migrations.md). @@ -483,3 +484,40 @@ rendering, diff, snapshots, the CLI in a temp folder, branch analysis). additive change with a recreated view checked with real inserts, resume after a failed statement, hash mismatch under `strict`, the lease, and drift. Passing on 26.2.19.43 and 26.8.2.7. `tests/package-consumer.mts` imports all three entry points from the packed tarball under Node. + +--- + +## 8. Postgres (phase 6) + +Open question 1 is answered with full authoring, not runtime-only: Maple's Postgres schema +(70 tables, 76 drizzle-kit migrations) was the consumer, and keeping drizzle-kit for authoring +would have meant two definitions per table. + +**Where the dialect seam is.** Shared: the snapshot envelope (`Snapshot` is a union over +`dialect`), `entityKey` / `sortEntities` / hashing, the branch graph (`kit/graph.ts`), folder +loading, `MigrationDriver`. Per dialect: entities (`pg-entities.ts`), definitions +(`pg-define.ts`, exported as `S.pg`), the diff (`pg-diff.ts`), ops and DDL (`pg-ops.ts`), the +ledger (`pg-ledger.ts`), drift (`pg-verify.ts`). `run`, `status`, `verify` and `generate` pick +the dialect from the config or the snapshots and dispatch; nothing ClickHouse-specific moved. + +**Runtime.** One transaction per migration, ledger row included, under +`pg_advisory_xact_lock`: transaction-scoped, so a pooled connection cannot leak the lock. The +ledger is a plain table with a primary key. The ClickHouse step journal and lease are not used. + +**Drift without normalizing SQL.** The catalog stores `'open'` as `'open'::text` and +`status in ('a', 'b')` as `(status = ANY (ARRAY[...]))`; any text normalizer would be a +heuristic. `verify` instead renders the snapshot into a scratch schema inside a transaction it +always rolls back, and reads both catalogs with the same queries. Checked against Maple: its 76 +migrations replayed on PGlite, baselined from drizzle-kit's last snapshot, verify clean except +two real orphans (a table and a column the SQL created and no migration dropped, absent from +drizzle-kit's snapshot). + +**Adoption.** A drizzle-kit `snapshot.json` is recognized (it has `ddl`, not `entities`) and +kept aside as `foreignSnapshot`. Migrations sorting before the first effect-orm snapshot are +legacy: `check` accepts them, `generate` refuses to diff until `--baseline` (from the +definitions, or `--from-drizzle` from drizzle-kit's last snapshot) starts the history. +`Migrate.baseline` records already-applied migrations without running them. + +**Not done.** Check and unique constraints, enums, views, sequences beyond identity defaults, +non-`public` schemas, `CREATE INDEX CONCURRENTLY` (needs a migration outside a transaction), +rename detection, and `pull` / `push`. diff --git a/docs/README.md b/docs/README.md index 6d246ba..34cf2c9 100644 --- a/docs/README.md +++ b/docs/README.md @@ -56,7 +56,7 @@ Roughly in reading order. | [Tenant scoping](./tenant-scoping.md) | `tenantScope`, what marks a query scoped, `crossTenant()` | | [Extending the DSL](./extending.md) | `defineFn`, raw escape hatches, handwritten SQL | | [Postgres](./postgres.md) | The Postgres dialect, its column types and functions | -| [Schema and migrations](./migrations.md) | `defineTable`, `materializedView`, `effect-orm generate`, applying migrations | +| [Schema and migrations](./migrations.md) | `defineTable`, `materializedView`, `S.pg.table`, `effect-orm generate`, applying migrations, adopting drizzle-kit | | [Statements and transactions](./database.md) | `Database` over your `SqlClient`: `run`, `execute`, `transaction`, retry | ## Reference @@ -79,7 +79,7 @@ Roughly in reading order. | `@maple-dev/effect-orm/benchmark/cli` | `runCli(args)` for embedding the bundled `ch-bench` commands | | `@maple-dev/effect-orm/schema` | `defineTable`, `materializedView`, DDL rendering, snapshots, and the schema diff. Pure | | `@maple-dev/effect-orm/kit` | `generate` and `check` over a migrations folder, `defineConfig`, and `runCli` for the bundled `effect-orm` command. Node or Bun | -| `@maple-dev/effect-orm/migrate` | `run`, `status`, `verify`, and `MigrationDriver`: applies migrations through a driver you provide | +| `@maple-dev/effect-orm/migrate` | `run`, `status`, `verify`, `baseline`, and `MigrationDriver`: applies ClickHouse or Postgres migrations through a driver you provide | | `@maple-dev/effect-orm/database` | `Database` over your `SqlClient`: `run` compiled queries, `execute` statements, `transaction` with settings and contention retry | The root barrel is curated, not exhaustive — see diff --git a/docs/migrations.md b/docs/migrations.md index b56cdfd..d4e4d09 100644 --- a/docs/migrations.md +++ b/docs/migrations.md @@ -5,13 +5,16 @@ in TypeScript too, three entry points add that, all opt-in: | Entry | Runs where | What it does | | ---------------------------------- | ------------------ | ----------------------------------------------------------------------------- | -| `@maple-dev/effect-orm/schema` | anywhere, pure | `defineTable` / `materializedView`, DDL rendering, snapshots, the diff | +| `@maple-dev/effect-orm/schema` | anywhere, pure | `defineTable` / `materializedView` / `pg.table`, DDL rendering, snapshots, the diff | | `@maple-dev/effect-orm/kit` | Node or Bun | `generate` and `check` over a migrations folder; the `effect-orm` command | | `@maple-dev/effect-orm/migrate` | anywhere Effect runs | applies migrations through a driver you provide, `status`, `verify` | -ClickHouse only, for now. The model follows drizzle-kit (a committed snapshot per migration, -an offline `generate`, data-loss confirmations by prompt or by hints), and the runtime is -built for a database without transactions. +Both ClickHouse and Postgres. The model follows drizzle-kit (a committed snapshot per migration, +an offline `generate`, data-loss confirmations by prompt or by hints). Snapshots, the branch +check and folder loading are shared; table definitions, the diff, the DDL, the runtime and +drift detection are per database, because the two differ where it matters: ClickHouse DDL is +not transactional and cannot change most things in place, and Postgres DDL is and can. The +sections below describe ClickHouse first; [Postgres](#postgres) covers what differs. ## Defining tables @@ -176,3 +179,96 @@ or comments yet. `effect-orm verify` exits 3 when it finds drift. _(Effect's own `ClickhouseMigrator` creates its ledger with a statement ClickHouse 26.8 rejects, and inserts ledger rows before running each migration. That is why this package has its own runner.)_ + +## Postgres + +Set `dialect: "postgres"` in the config and define tables with `S.pg.table`. `generate`, `check`, +`migrate`, `status` and `verify` then work as above, with the differences below. + +```ts title="migrations-postgres.ts" +import * as CH from "@maple-dev/effect-orm" +import * as PG from "@maple-dev/effect-orm/postgres" +import * as S from "@maple-dev/effect-orm/schema" + +export const Dashboards = S.pg.table("dashboards", { + columns: { + org_id: PG.text, + id: PG.text, + status: S.pg.column(PG.text, { default: "open" }), + created_at: S.pg.column(PG.timestamptz, { defaultExpr: "now()" }), + archived_at: PG.nullable(PG.timestamptz), + }, + primaryKey: ["org_id", "id"], + indexes: [S.pg.index("dashboards_open_idx", ["org_id"], { where: ($) => $.archived_at.isNull() })], + tenantColumn: "org_id", +}) + +export const Shares = S.pg.table("dashboard_shares", { + columns: { org_id: PG.text, id: PG.text, dashboard_id: PG.text, widget_id: PG.nullable(PG.text), revoked_at: PG.nullable(PG.timestamptz) }, + primaryKey: ["org_id", "id"], + indexes: [ + // At most one live share per dashboard and widget: a partial unique index on an expression. + S.pg.uniqueIndex("dashboard_shares_live_unq", ($) => [$.org_id, $.dashboard_id, CH.coalesce($.widget_id, CH.lit(""))], { + where: "revoked_at is null", + }), + ], + foreignKeys: [ + S.pg.foreignKey({ columns: ["org_id", "dashboard_id"], references: Dashboards, foreignColumns: ["org_id", "id"], onDelete: "cascade" }), + ], +}) + +export const ddl = S.renderPgSchema(S.pgEntitiesOf([Dashboards, Shares])) +``` + +**Definitions.** A column is `NOT NULL` unless its type is `PG.nullable(...)`. `S.pg.column(type, +options)` adds a `default` (a value of the column's type), a `defaultExpr` (SQL or a DSL +expression) or an `identity` (`"always"` or `"by default"`); any of them makes the column +optional on insert. `primaryKey` takes column names, or `{ columns, name }`; the default name is +`_pkey`. Indexes are `S.pg.index` / `S.pg.uniqueIndex` over column names or expressions, +with `where` for a partial index and `using` for the access method. A foreign key without a +`name` gets drizzle-kit's, `
____fk`. Types are +stored as Postgres names them (`int4` is `integer`), so snapshots compare with the catalog and +with drizzle-kit. Check and unique constraints, enums, views, sequences and other schemas are +not modeled yet; write them in a `--custom` migration. + +**Generating.** Postgres changes a column's type, nullability, default or identity in place +(`ALTER COLUMN`), and a primary key, an index or a foreign key by dropping and re-creating it, +so `generate` reports nothing as unsupported. A type change is labeled `rewrite`: Postgres +rewrites the table under an exclusive lock. Drops still need confirmation, and renames still +read as a drop and an add. Generated files carry `"dialect": "postgres"`. + +**Applying.** Each migration runs in one transaction with its ledger row, under a +transaction-scoped advisory lock, so concurrent deploys wait for each other rather than +racing. A failed statement rolls the whole migration back and the next run starts it from the +top: there is no step journal, no lease, no `partial` or `uncertain` state, and nothing for +`resolve` to do. The driver needs a transaction, which `Migrate.fromSqlClient` provides from +`SqlClient.withTransaction`. A statement Postgres refuses inside a transaction (`CREATE INDEX +CONCURRENTLY`, `ALTER TYPE ... ADD VALUE` before Postgres 12) cannot be in a migration yet. + +**Drift.** `verify` builds the expected schema in a scratch schema, inside a transaction it +always rolls back, and reads both catalogs with the same queries, so Postgres deparses both +sides and `'open'` matches the stored `'open'::text`. It checks tables, columns (type, `NOT NULL`, +default, identity), primary keys, indexes (uniqueness, method, keys, predicate) and foreign keys +in the current schema. `Migrate.verify(migrations, { ignoreTables })` skips tables another tool +owns. + +### Adopting a drizzle-kit folder + +A drizzle-kit (v1) folder already has the layout `migrate` reads: `_/migration.sql` +split on `--> statement-breakpoint`. Its `snapshot.json` files are recognized as drizzle-kit's +and set aside, so the folder runs as it is. Adoption is two steps: + +1. `effect-orm generate --baseline --from-drizzle` writes a migration that runs nothing, whose + snapshot is drizzle-kit's last one converted to entities. Anything the conversion cannot model + is listed and nothing is written. Without `--from-drizzle` the snapshot comes from your + `S.pg.table` definitions instead. Migrations before the baseline are legacy: they run, but + nothing diffs against them, and a plain `generate` refuses to run until a baseline exists. +2. On a database drizzle-kit (or anything else) already migrated, `effect-orm baseline ` + (`Migrate.baseline`) records the baseline and every migration before it as applied, without + running them. A fresh database, such as a test's, simply runs everything. + +The first `generate` after the baseline diffs your `S.pg.table` definitions against what +drizzle-kit recorded, so every place they disagree (a constraint name, a default) shows up as an +op to accept or fix. `verify` against the baseline also finds objects the database has and +drizzle-kit's snapshot does not, such as a table a hand-written migration created and nothing +dropped. diff --git a/src/kit/cli.ts b/src/kit/cli.ts index d3b9937..353783e 100644 --- a/src/kit/cli.ts +++ b/src/kit/cli.ts @@ -8,14 +8,17 @@ import * as Migrate from "../migrate" import { Hints, type Hint } from "../schema/diff" import { check, generate, KitError, readMigrations, type KitConfig } from "./generate" -const help = `effect-orm: schema migrations for ClickHouse +const help = `effect-orm: schema migrations for ClickHouse and Postgres generate [--name x] [--custom] Diff the schema against the migrations folder and write the next migration [--hints ] [--hints-file ] + [--baseline [--from-drizzle]] Start the history of a folder another tool wrote: a migration that runs + nothing, with the current schema (or drizzle-kit's last snapshot) as its snapshot check Validate the migrations folder (snapshot chain, branch conflicts) migrate [--strict] Apply pending migrations (needs config.driver) status Applied, pending, partial, or changed, per migration verify Compare the database with the last applied snapshot + baseline Record and every one before it as applied, without running them resolve Record what a step left "started" did, after checking the database: --ran | --not-ran --ran skips it on the next migrate, --not-ran runs it @@ -28,6 +31,8 @@ const options = { config: { type: "string" }, name: { type: "string" }, custom: { type: "boolean" }, + baseline: { type: "boolean" }, + "from-drizzle": { type: "boolean" }, hints: { type: "string" }, "hints-file": { type: "string" }, strict: { type: "boolean" }, @@ -95,6 +100,7 @@ const program = (args: ReadonlyArray, print: (line: string) => void) => hints, ...(parsed.values.name !== undefined ? { name: parsed.values.name } : undefined), ...(parsed.values.custom === true ? { custom: true } : undefined), + ...(parsed.values.baseline === true ? { baseline: parsed.values["from-drizzle"] === true ? "drizzle" : "schema" } : undefined), ...(interactive ? { confirm: promptConfirm } : undefined), }) out(result, () => @@ -107,7 +113,7 @@ const program = (args: ReadonlyArray, print: (line: string) => void) => case "check": { const result = yield* check(config, cwd) out(result, () => [ - `${result.migrations} migrations, ok.`, + `${result.migrations} migrations, ok.${result.legacy > 0 ? ` ${result.legacy} predate the first snapshot.` : ""}`, ...(result.leaves.length > 1 ? [`Independent branches (merged by the next generate): ${result.leaves.join(", ")}`] : []), ]) return 0 @@ -124,29 +130,52 @@ const program = (args: ReadonlyArray, print: (line: string) => void) => if (migration === undefined) return yield* new KitError({ code: "config", message: `no migration named ${migrationName}` }) yield* withDriver( config, - Migrate.resolveStep({ migration, step: stepId, outcome: ran ? "ran" : "not-ran", render: config.render ?? {} }), + Migrate.resolveStep({ + migration, + step: stepId, + outcome: ran ? "ran" : "not-ran", + render: config.render ?? {}, + ...(config.dialect !== undefined ? { dialect: config.dialect } : undefined), + }), ) out({ migration: migrationName, step: stepId, outcome: ran ? "ran" : "not-ran" }, () => [ `Recorded ${migrationName} step ${stepId} as ${ran ? "done; migrate will skip it" : "not run; migrate will run it"}.`, ]) return 0 } + case "baseline": { + const [, upTo] = parsed.positionals + if (upTo === undefined) return yield* new KitError({ code: "config", message: "usage: effect-orm baseline " }) + const migrations = yield* readMigrations(resolve(cwd, config.out)) + const recorded = yield* withDriver( + config, + Migrate.baseline({ + migrations, + upTo, + render: config.render ?? {}, + ...(config.dialect !== undefined ? { dialect: config.dialect } : undefined), + }), + ) + out(recorded, () => (recorded.length === 0 ? ["Already recorded; nothing to do."] : recorded.map((name) => `recorded ${name} as applied`))) + return 0 + } case "migrate": case "status": case "verify": { const migrations = yield* readMigrations(resolve(cwd, config.out)) const render = config.render ?? {} + const dialect = config.dialect !== undefined ? { dialect: config.dialect } : undefined if (command === "migrate") { - const ran = yield* withDriver(config, Migrate.run({ migrations, render, strict: parsed.values.strict === true })) + const ran = yield* withDriver(config, Migrate.run({ migrations, render, strict: parsed.values.strict === true, ...dialect })) out(ran, () => (ran.length === 0 ? ["Nothing to apply."] : ran.map((m) => `applied ${m.name} (${m.steps} statements${m.resumedSteps > 0 ? `, ${m.resumedSteps} resumed` : ""})`))) return 0 } if (command === "status") { - const rows = yield* withDriver(config, Migrate.status(migrations, render)) + const rows = yield* withDriver(config, Migrate.status(migrations, render, config.dialect)) out(rows, () => rows.map((r) => `${r.state.padEnd(8)} ${r.name}${r.appliedAt !== undefined ? ` ${r.appliedAt}` : ""}`)) return 0 } - const result = yield* withDriver(config, Migrate.verify(migrations)) + const result = yield* withDriver(config, Migrate.verify(migrations, { ...dialect })) out(result, () => result.against === undefined ? ["No applied migration with a snapshot to compare against."] diff --git a/src/kit/generate.ts b/src/kit/generate.ts index 24a204c..b3b6f50 100644 --- a/src/kit/generate.ts +++ b/src/kit/generate.ts @@ -8,18 +8,33 @@ import { Effect, Schema } from "effect" import type { Layer } from "effect" import { fromRecord, isStagingName, type LoadedMigration, type MigrationInput } from "../migrate/source" import type { MigrationDriver } from "../migrate/driver" -import { diffSchemas, type Hint } from "../schema/diff" -import { labelOf, renderOp, type MigrationFile } from "../schema/ops" +import { diffSchemas, type DiffResult, type Hint } from "../schema/diff" +import { fromDrizzleSnapshot } from "../schema/drizzle" +import { ORIGIN_ID, type AnySchemaEntity, type SchemaDialect, type SchemaEntity } from "../schema/entities" +import { labelOf, renderOp, type MigrationFile, type MigrationOp } from "../schema/ops" +import { diffPgSchemas } from "../schema/pg-diff" +import type { PgSchemaEntity } from "../schema/pg-entities" +import { labelOfPg, renderPgOp, type PgMigrationOp } from "../schema/pg-ops" import type { RenderOptions } from "../schema/render" -import { entitiesOf, isSchemaObject, makeSnapshot, serializeSnapshot, type SchemaObject } from "../schema/snapshot" +import { + dialectOfObject, + entitiesOf, + isSchemaObject, + makeSnapshot, + pgEntitiesOf, + serializeSnapshot, + type SchemaObject, +} from "../schema/snapshot" import { analyze, type GraphProblem } from "./graph" export interface KitConfig { - /** Modules whose exports include `defineTable` / `materializedView` values. */ + /** The database this folder migrates. Default `clickhouse`. */ + readonly dialect?: SchemaDialect + /** Modules whose exports include `defineTable` / `materializedView` (or `S.pg.table`) values. */ readonly schema: string | ReadonlyArray /** The migrations folder. */ readonly out: string - /** How `migrate`, `status`, and `verify` render statements for this deployment. */ + /** How `migrate`, `status`, and `verify` render ClickHouse statements for this deployment. */ readonly render?: RenderOptions /** Needed by `migrate`, `status`, and `verify`. Build it from your `SqlClient`. */ readonly driver?: Layer.Layer @@ -50,7 +65,15 @@ export const loadSchema = (config: KitConfig, cwd: string): Effect.Effect dialectOfObject(object) === dialect) + if (mine.length === 0 && objects.size > 0) { + return yield* new KitError({ + code: "config", + message: `the schema modules define no ${dialect} tables; set \`dialect\` in the config to the database they are for`, + }) + } + return mine }) /** Read `//...` with node:fs. */ @@ -86,7 +109,7 @@ export const check = (config: KitConfig, cwd: string) => if (analysis.problems.length > 0) { return yield* new KitError({ code: "check", message: "the migrations folder is inconsistent", details: describeProblems(analysis.problems) }) } - return { migrations: migrations.length, leaves: analysis.leaves.map((l) => l.name) } + return { migrations: migrations.length, leaves: analysis.leaves.map((l) => l.name), legacy: analysis.legacy.length } }) const ADJECTIVES = ["amber", "brisk", "calm", "deft", "eager", "fond", "gentle", "hardy", "keen", "lucid", "merry", "nimble", "quiet", "rapid", "steady", "tidy", "vivid", "witty"] @@ -113,6 +136,14 @@ export interface GenerateOptions { /** Asked for each data-loss confirmation the hints do not cover. Absent: report them as missing. */ readonly confirm?: (hint: Hint) => Effect.Effect readonly now?: Date + /** + * Adopt a folder another tool wrote, or a database that already exists: + * write a migration with no statements whose snapshot is the schema as it + * stands. `schema` takes it from the definitions; `drizzle` from the newest + * drizzle-kit snapshot in the folder, so the next `generate` shows every + * place the definitions and the database disagree. + */ + readonly baseline?: "schema" | "drizzle" } export interface GenerateResult { @@ -132,44 +163,83 @@ export const generate = (config: KitConfig, cwd: string, options: GenerateOption if (options.name !== undefined && !/^[a-z0-9_]+$/.test(options.name)) { return yield* new KitError({ code: "config", message: "--name takes lowercase letters, digits, and underscores" }) } - const objects = yield* loadSchema(config, cwd) - const next = yield* Effect.try({ - try: () => entitiesOf(objects), - catch: (cause) => new KitError({ code: "config", message: cause instanceof Error ? cause.message : String(cause) }), - }) + const dialect = config.dialect ?? "clickhouse" + if (analysis.dialect !== undefined && analysis.dialect !== dialect) { + return yield* new KitError({ code: "config", message: `the migrations folder holds ${analysis.dialect} snapshots, but the config's dialect is ${dialect}` }) + } + const adopting = analysis.leaves.length === 0 && analysis.legacy.length > 0 + if (options.baseline !== undefined && analysis.leaves.length > 0) { + return yield* new KitError({ code: "config", message: "--baseline starts a folder's history; this one already has effect-orm snapshots" }) + } + if (options.baseline === undefined && adopting && options.custom !== true) { + return yield* new KitError({ + code: "config", + message: `the folder has ${analysis.legacy.length} migrations from another tool and no effect-orm snapshot yet; run generate --baseline first (add --from-drizzle for a drizzle-kit folder)`, + }) + } let file: MigrationFile | undefined - let snapshotEntities = next - if (options.custom === true) { - snapshotEntities = analysis.base - } else { - const hints = [...(options.hints ?? [])] - let diff = diffSchemas(analysis.base, next, hints) - if (diff.unsupported.length > 0) { + let snapshotEntities: ReadonlyArray + let prevIds = analysis.baseIds + if (options.baseline === "drizzle") { + if (dialect !== "postgres") return yield* new KitError({ code: "config", message: "--from-drizzle needs dialect: \"postgres\"" }) + const source = [...analysis.legacy].reverse().find((m) => m.foreignSnapshot?.tool === "drizzle-kit") + if (source === undefined) return yield* new KitError({ code: "config", message: "no drizzle-kit snapshot.json in the folder to start from" }) + const imported = fromDrizzleSnapshot(source.foreignSnapshot!.json) + if (imported.unsupported.length > 0) { return yield* new KitError({ code: "unsupported", - message: "these changes need a table rebuild or a data rewrite, which generate does not write yet", - details: diff.unsupported.map((u) => `${u.entity}: ${u.message}`), + message: `${source.name}/snapshot.json holds objects effect-orm cannot model yet`, + details: imported.unsupported, }) } - if (diff.missingHints.length > 0 && options.confirm !== undefined) { - for (const hint of diff.missingHints) { - if (yield* options.confirm(hint)) hints.push(hint) + snapshotEntities = imported.entities + prevIds = [ORIGIN_ID] + } else { + const objects = yield* loadSchema(config, cwd) + const next = yield* Effect.try({ + try: (): ReadonlyArray => (dialect === "postgres" ? pgEntitiesOf(objects) : entitiesOf(objects)), + catch: (cause) => new KitError({ code: "config", message: cause instanceof Error ? cause.message : String(cause) }), + }) + snapshotEntities = next + if (options.baseline === "schema") prevIds = [ORIGIN_ID] + else if (options.custom === true) snapshotEntities = analysis.base + else { + const hints = [...(options.hints ?? [])] + const diff = (): DiffResult => + dialect === "postgres" + ? diffPgSchemas(analysis.base as ReadonlyArray, next as ReadonlyArray, hints) + : diffSchemas(analysis.base as ReadonlyArray, next as ReadonlyArray, hints) + let result = diff() + if (result.unsupported.length > 0) { + return yield* new KitError({ + code: "unsupported", + message: "these changes need a table rebuild or a data rewrite, which generate does not write yet", + details: result.unsupported.map((u) => `${u.entity}: ${u.message}`), + }) } - diff = diffSchemas(analysis.base, next, hints) - } - if (diff.missingHints.length > 0) { - return yield* new KitError({ - code: "missing_hints", - message: "confirm data loss with --hints, or run in a terminal to be asked", - details: [JSON.stringify(diff.missingHints)], - }) + if (result.missingHints.length > 0 && options.confirm !== undefined) { + for (const hint of result.missingHints) { + if (yield* options.confirm(hint)) hints.push(hint) + } + result = diff() + } + if (result.missingHints.length > 0) { + return yield* new KitError({ + code: "missing_hints", + message: "confirm data loss with --hints, or run in a terminal to be asked", + details: [JSON.stringify(result.missingHints)], + }) + } + if (result.ops.length === 0) return { written: undefined, plan: [] } satisfies GenerateResult + file = + dialect === "postgres" + ? { version: "1", dialect: "postgres", ops: result.ops as ReadonlyArray } + : { version: "1", ops: result.ops as ReadonlyArray } } - if (diff.ops.length === 0) return { written: undefined, plan: [] } satisfies GenerateResult - file = { version: "1", ops: diff.ops } } - const snapshot = yield* makeSnapshot(snapshotEntities, analysis.baseIds) + const snapshot = yield* makeSnapshot(snapshotEntities, prevIds, dialect) const taken = new Set(migrations.map((m) => m.name.slice(0, 14))) const name = `${timestamp(options.now ?? new Date(), taken)}_${options.name ?? `${pick(ADJECTIVES)}_${pick(NOUNS)}`}` const dir = join(out, name) @@ -180,7 +250,12 @@ export const generate = (config: KitConfig, cwd: string, options: GenerateOption yield* io( () => file === undefined - ? writeFile(join(staging, "migration.sql"), "-- Custom SQL migration. Separate statements with a line holding only:\n-- --> statement-breakpoint\n") + ? writeFile( + join(staging, "migration.sql"), + options.baseline !== undefined + ? "-- Baseline: the schema as it stood when effect-orm took over this folder. Runs nothing.\n" + : "-- Custom SQL migration. Separate statements with a line holding only:\n-- --> statement-breakpoint\n", + ) : writeFile(join(staging, "migration.json"), `${JSON.stringify(file, null, "\t")}\n`), `Cannot write ${dir}`, ) @@ -188,6 +263,11 @@ export const generate = (config: KitConfig, cwd: string, options: GenerateOption yield* io(() => rename(staging, dir), `Cannot move ${staging} to ${dir}`) }).pipe(Effect.onError(() => Effect.promise(() => rm(staging, { recursive: true, force: true }).catch(() => undefined)))) - const plan = (file?.ops ?? []).map((op) => ({ label: labelOf(op), sql: renderOp(op, config.render) })) + const plan = + file === undefined + ? [] + : "dialect" in file + ? file.ops.map((op) => ({ label: labelOfPg(op), sql: renderPgOp(op) })) + : file.ops.map((op) => ({ label: labelOf(op), sql: renderOp(op, config.render) })) return { written: dir, plan } satisfies GenerateResult }) diff --git a/src/kit/graph.ts b/src/kit/graph.ts index 16f7a5f..6c6dc9f 100644 --- a/src/kit/graph.ts +++ b/src/kit/graph.ts @@ -7,7 +7,7 @@ // to be regenerated on top of the other. import { Effect } from "effect" -import { canonicalJson, entityKey, ORIGIN_ID, sha256Hex, sortEntities, type SchemaEntity } from "../schema/entities" +import { canonicalJson, entityKey, ORIGIN_ID, sha256Hex, sortEntities, type AnySchemaEntity, type SchemaDialect } from "../schema/entities" import { migrationParents, type LoadedMigration } from "../migrate/source" export interface GraphProblem { @@ -20,9 +20,17 @@ export interface GraphAnalysis { /** Migrations nothing builds on yet. More than one means unmerged branches. */ readonly leaves: ReadonlyArray /** The schema the next migration starts from. */ - readonly base: ReadonlyArray + readonly base: ReadonlyArray /** Parents for the next migration's snapshot. */ readonly baseIds: ReadonlyArray + /** The dialect the snapshots are written for; `undefined` when there are none yet. */ + readonly dialect: SchemaDialect | undefined + /** + * Migrations from before the first snapshot: written by another tool (a + * drizzle-kit folder being adopted) or by hand. They run, but nothing diffs + * against them. + */ + readonly legacy: ReadonlyArray } const ancestorsOf = ( @@ -42,12 +50,12 @@ const ancestorsOf = ( /** Entity changes from `from` to `to`: key -> new entity, or null for removed. */ const changesBetween = ( - from: ReadonlyArray, - to: ReadonlyArray, -): Map => { + from: ReadonlyArray, + to: ReadonlyArray, +): Map => { const before = new Map(from.map((e) => [entityKey(e), canonicalJson(e)])) const after = new Map(to.map((e) => [entityKey(e), e])) - const out = new Map() + const out = new Map() for (const [key, entity] of after) if (before.get(key) !== canonicalJson(entity)) out.set(key, entity) for (const key of before.keys()) if (!after.has(key)) out.set(key, null) return out @@ -56,19 +64,28 @@ const changesBetween = ( /** The table an entity belongs to, so two branches editing one table conflict. */ const ownerOf = (key: string): string => { const [kind, rest = ""] = key.split(":") - return kind === "column" || kind === "index" ? `table:${rest.split(".")[0]}` : key + return kind === "column" || kind === "index" || kind === "foreign_key" ? `table:${rest.split(".")[0]}` : key } export const analyze = (migrations: ReadonlyArray): Effect.Effect => Effect.gen(function* () { const problems: Array = [] const withSnapshots = migrations.filter((m) => m.snapshot !== undefined) + const first = withSnapshots.reduce((min, m) => (min === undefined || m.name < min ? m.name : min), undefined) + const legacy = migrations.filter((m) => m.snapshot === undefined && (first === undefined || m.name < first)) for (const m of migrations) { - if (m.snapshot === undefined) problems.push({ migration: m.name, message: "has no snapshot.json" }) + if (m.snapshot === undefined && !legacy.includes(m)) { + problems.push({ migration: m.name, message: "has no snapshot.json, and sorts after the first migration that has one" }) + } + } + const dialects = [...new Set(withSnapshots.map((m) => m.snapshot!.dialect))] + if (dialects.length > 1) { + problems.push({ migration: withSnapshots.at(-1)!.name, message: `the folder mixes ${dialects.join(" and ")} snapshots; keep one folder per database` }) } + const dialect = dialects[0] const knownIds = new Set(withSnapshots.map((m) => m.snapshot!.id)) for (const m of withSnapshots) { - const id = yield* Effect.promise(() => sha256Hex(canonicalJson(sortEntities(m.snapshot!.entities)))) + const id = yield* Effect.promise(() => sha256Hex(canonicalJson(sortEntities(m.snapshot!.entities)))) if (id !== m.snapshot!.id) { problems.push({ migration: m.name, message: "snapshot id does not match its entities; the snapshot was edited by hand" }) } @@ -88,16 +105,16 @@ export const analyze = (migrations: ReadonlyArray): Effect.Effe } // A broken chain makes leaves and ancestors meaningless; report it alone. - if (problems.length > 0) return { problems, leaves: [], base: [], baseIds: [] } + if (problems.length > 0) return { problems, leaves: [], base: [], baseIds: [], dialect, legacy } const hasChild = new Set() for (const list of parents.values()) for (const p of list) hasChild.add(p) const leaves = withSnapshots.filter((m) => !hasChild.has(m)) - if (leaves.length === 0) return { problems, leaves, base: [], baseIds: [ORIGIN_ID] } + if (leaves.length === 0) return { problems, leaves, base: [], baseIds: [ORIGIN_ID], dialect, legacy } if (leaves.length === 1) { const leaf = leaves[0]! - return { problems, leaves, base: leaf.snapshot!.entities, baseIds: [leaf.snapshot!.id] } + return { problems, leaves, base: leaf.snapshot!.entities, baseIds: [leaf.snapshot!.id], dialect, legacy } } // Several leaves: merge them on their nearest common ancestor. @@ -125,5 +142,5 @@ export const analyze = (migrations: ReadonlyArray): Effect.Effe else merged.set(key, entity) } } - return { problems, leaves, base: [...merged.values()], baseIds: leaves.map((l) => l.snapshot!.id) } + return { problems, leaves, base: [...merged.values()], baseIds: leaves.map((l) => l.snapshot!.id), dialect, legacy } }) diff --git a/src/kit/kit.test.ts b/src/kit/kit.test.ts index dcf27e6..d8901d2 100644 --- a/src/kit/kit.test.ts +++ b/src/kit/kit.test.ts @@ -98,4 +98,93 @@ describe("effect-orm CLI", () => { writeFileSync(path, readFileSync(path, "utf8").replace('"String"', '"UInt8"')) expect(await cli("check")).toBe(1) }) + + describe("Postgres", () => { + const pgModule = (extra: { owner?: boolean } = {}) => ` +import * as PG from "${src}/postgres" +import * as S from "${src}/schema" + +export const Dashboards = S.pg.table("dashboards", { + columns: { + org_id: PG.text, + id: PG.text, + status: S.pg.column(PG.text, { default: "open" }), + ${extra.owner === true ? "owner: PG.nullable(PG.text)," : ""} + }, + primaryKey: { columns: ["org_id", "id"], name: "dashboards_org_id_id_pk" }, + indexes: [S.pg.index("dashboards_open_idx", ["org_id"], { where: "status = 'open'" })], +}) +` + const pgConfig = (schema: string, name = "effect-orm.config.ts") => + writeFileSync(join(dir, name), `export default { dialect: "postgres", schema: "./${schema}", out: "./migrations" }\n`) + + it("generates Postgres migrations and refuses a ClickHouse config on the folder", async () => { + writeFileSync(join(dir, "pg1.ts"), pgModule()) + pgConfig("pg1.ts") + expect(await cli("generate", "--name", "init")).toBe(0) + const [first] = folders() + const file = JSON.parse(readFileSync(join(dir, "migrations", first!, "migration.json"), "utf8")) + expect(file.dialect).toBe("postgres") + expect(file.ops.map((op: { op: string }) => op.op)).toEqual(["create_table", "create_index"]) + expect(JSON.parse(readFileSync(join(dir, "migrations", first!, "snapshot.json"), "utf8")).dialect).toBe("postgres") + + writeFileSync(join(dir, "pg2.ts"), pgModule({ owner: true })) + pgConfig("pg2.ts", "v2.config.ts") + expect(await cli("generate", "--name", "owner", "--config", "v2.config.ts")).toBe(0) + expect(lines.join("\n")).toContain('[metadata] ALTER TABLE "dashboards" ADD COLUMN IF NOT EXISTS "owner" text') + expect(await cli("check")).toBe(0) + + writeFileSync(join(dir, "ch.config.ts"), `export default { schema: "./pg2.ts", out: "./migrations" }\n`) + expect(await cli("generate", "--config", "ch.config.ts")).toBe(1) + }) + + it("adopts a drizzle-kit folder: baseline from its last snapshot, then diff the definitions against it", async () => { + const legacy = join(dir, "migrations", "20260101000000_drizzle_init") + mkdirSync(legacy, { recursive: true }) + writeFileSync(join(legacy, "migration.sql"), `CREATE TABLE "dashboards" ("org_id" text NOT NULL, "id" text NOT NULL);`) + writeFileSync( + join(legacy, "snapshot.json"), + JSON.stringify({ + version: "8", + dialect: "postgres", + id: "a", + prevIds: [], + renames: [], + ddl: [ + { isRlsEnabled: false, name: "dashboards", entityType: "tables", schema: "public" }, + { type: "text", typeSchema: null, notNull: true, dimensions: 0, default: null, generated: null, identity: null, name: "org_id", entityType: "columns", schema: "public", table: "dashboards" }, + { type: "text", typeSchema: null, notNull: true, dimensions: 0, default: null, generated: null, identity: null, name: "id", entityType: "columns", schema: "public", table: "dashboards" }, + { type: "text", typeSchema: null, notNull: true, dimensions: 0, default: "'open'", generated: null, identity: null, name: "status", entityType: "columns", schema: "public", table: "dashboards" }, + { columns: ["org_id", "id"], nameExplicit: false, name: "dashboards_org_id_id_pk", entityType: "pks", schema: "public", table: "dashboards" }, + { + nameExplicit: true, + columns: [{ value: "org_id", isExpression: false, asc: true, nullsFirst: false, opclass: null }], + isUnique: false, + where: "status = 'open'", + with: "", + method: "btree", + concurrently: false, + name: "dashboards_open_idx", + entityType: "indexes", + schema: "public", + table: "dashboards", + }, + ], + }), + ) + writeFileSync(join(dir, "pg3.ts"), pgModule({ owner: true })) + pgConfig("pg3.ts") + expect(await cli("generate")).toBe(1) + expect(await cli("generate", "--baseline", "--from-drizzle", "--name", "effect_orm_baseline")).toBe(0) + const baseline = folders().at(-1)! + expect(baseline).toMatch(/_effect_orm_baseline$/) + expect(readFileSync(join(dir, "migrations", baseline, "migration.sql"), "utf8")).toMatch(/^-- Baseline/) + expect(await cli("check")).toBe(0) + expect(lines.at(-1)).toBe("2 migrations, ok. 1 predate the first snapshot.") + // The definitions add `owner`; everything else matches what drizzle-kit recorded. + expect(await cli("generate", "--json", "--name", "owner")).toBe(0) + const plan = JSON.parse(lines.at(-1)!).plan + expect(plan).toEqual([{ label: "metadata", sql: ['ALTER TABLE "dashboards" ADD COLUMN IF NOT EXISTS "owner" text'] }]) + }) + }) }) diff --git a/src/migrate.ts b/src/migrate.ts index b3e066e..f46c81e 100644 --- a/src/migrate.ts +++ b/src/migrate.ts @@ -1,8 +1,9 @@ // @maple-dev/effect-orm/migrate // // Applies migrations written by `effect-orm generate` (or by hand) to a -// ClickHouse database through a `MigrationDriver` you provide. Statements are -// journaled one by one, so a failed run resumes where it stopped. See +// ClickHouse or Postgres database through a `MigrationDriver` you provide. On +// ClickHouse statements are journaled one by one, so a failed run resumes +// where it stopped; on Postgres each migration is one transaction. See // docs/migrations.md. export { MigrationDriver, fromSqlClient, layerSqlClient, type FromSqlClientOptions, type MigrationDriverApi } from "./migrate/driver" @@ -19,6 +20,7 @@ export { export { LEDGER_TABLES } from "./migrate/ledger" export { STATEMENT_BREAKPOINT, + dialectOf, fromFileSystem, isStagingName, fromRecord, @@ -29,13 +31,15 @@ export { type MigrationStep, } from "./migrate/source" export { + baseline, resolveStep, run, status, type AppliedMigration, + type BaselineOptions, type MigrationState, type MigrationStatus, type ResolveStepOptions, type RunOptions, } from "./migrate/run" -export { verify, type Drift, type VerifyResult } from "./migrate/verify" +export { verify, type Drift, type VerifyOptions, type VerifyResult } from "./migrate/verify" diff --git a/src/migrate/driver.ts b/src/migrate/driver.ts index a01b6d2..e4e6723 100644 --- a/src/migrate/driver.ts +++ b/src/migrate/driver.ts @@ -16,6 +16,11 @@ export interface MigrationDriverApi { readonly execute: (sql: string) => Effect.Effect /** Run a SELECT and return its rows as plain records. */ readonly query: (sql: string) => Effect.Effect>, MigrateSqlError> + /** + * Run an effect in one transaction: commit when it succeeds, roll back when + * it fails. Postgres migrations need it; ClickHouse has none to offer. + */ + readonly transaction?: (effect: Effect.Effect) => Effect.Effect } export class MigrationDriver extends Context.Service()( @@ -32,10 +37,14 @@ export interface FromSqlClientOptions { readonly command?: (effect: Effect.Effect) => Effect.Effect } -/** A driver over an Effect `SqlClient`. */ +/** + * A driver over an Effect `SqlClient`. Its `transaction` is the client's + * `withTransaction`, which every statement the driver runs inside it joins. + */ export const fromSqlClient = (sql: SqlClient.SqlClient, options: FromSqlClientOptions = {}): MigrationDriverApi => { const command = options.command ?? ((effect) => effect) return { + transaction: (effect) => sql.withTransaction(effect).pipe(Effect.catchTag("SqlError", (cause) => Effect.fail(sqlError("BEGIN / COMMIT")(cause)))), execute: (text) => command(sql.unsafe(text)).pipe(Effect.asVoid, Effect.mapError(sqlError(text))), query: (text) => sql.unsafe>(text).pipe( diff --git a/src/migrate/pg-ledger.ts b/src/migrate/pg-ledger.ts new file mode 100644 index 0000000..dba7d29 --- /dev/null +++ b/src/migrate/pg-ledger.ts @@ -0,0 +1,68 @@ +// The Postgres migration ledger. +// +// One table, `_effect_orm_migrations`, with a row per applied migration. A +// Postgres migration runs in one transaction together with its ledger row, so +// the ClickHouse ledger's step journal and lease have nothing to do here: a +// migration is applied or it is not. Concurrent runs are serialized with a +// transaction-scoped advisory lock, which a pooled connection cannot leak. + +import { Effect } from "effect" +import { MigrationDriver } from "./driver" +import { MigrateSourceError, type MigrateSqlError } from "./errors" +import { LEDGER_TABLES, type AppliedRow } from "./ledger" + +/** The advisory lock every effect-orm migration run takes. Arbitrary, fixed. */ +export const PG_MIGRATION_LOCK = 7_243_567_198 + +const table = `"${LEDGER_TABLES.migrations}"` +const q = (value: string): string => `'${value.replace(/'/g, "''")}'` + +/** The driver's transaction, which Postgres migrations cannot run without. */ +export const transactionOf = (driver: typeof MigrationDriver.Service, migration: string) => + driver.transaction === undefined + ? Effect.fail( + new MigrateSourceError({ + migration, + message: "Postgres migrations run in a transaction, and this MigrationDriver has none; build it with fromSqlClient", + }), + ) + : Effect.succeed(driver.transaction) + +/** Wait for the migration lock, inside the current transaction. */ +export const lockPg = Effect.gen(function* () { + const driver = yield* MigrationDriver + yield* driver.query(`SELECT pg_advisory_xact_lock(${PG_MIGRATION_LOCK})`) +}) + +export const ensurePgLedger: Effect.Effect = Effect.gen(function* () { + const driver = yield* MigrationDriver + const transaction = yield* transactionOf(driver, LEDGER_TABLES.migrations) + // Under the lock: two runs creating the table at once would collide in pg_type. + yield* transaction( + Effect.gen(function* () { + yield* lockPg + yield* driver.execute( + `CREATE TABLE IF NOT EXISTS ${table} (name text PRIMARY KEY, hash text NOT NULL, applied_at timestamptz NOT NULL DEFAULT now())`, + ) + }), + ) +}) + +export const readPgApplied = Effect.gen(function* () { + const driver = yield* MigrationDriver + const rows = yield* driver.query(`SELECT name, hash, applied_at::text AS applied FROM ${table} ORDER BY applied_at, name`) + return rows.map((row): AppliedRow => ({ name: String(row.name), hash: String(row.hash), appliedAt: String(row.applied) })) +}) + +export const isPgApplied = (name: string) => + Effect.gen(function* () { + const driver = yield* MigrationDriver + const rows = yield* driver.query(`SELECT 1 AS applied FROM ${table} WHERE name = ${q(name)}`) + return rows.length > 0 + }) + +export const recordPgMigration = (name: string, hash: string) => + Effect.gen(function* () { + const driver = yield* MigrationDriver + yield* driver.execute(`INSERT INTO ${table} (name, hash) VALUES (${q(name)}, ${q(hash)}) ON CONFLICT (name) DO NOTHING`) + }) diff --git a/src/migrate/pg-migrate.test.ts b/src/migrate/pg-migrate.test.ts new file mode 100644 index 0000000..7a5d498 --- /dev/null +++ b/src/migrate/pg-migrate.test.ts @@ -0,0 +1,201 @@ +import { PgliteClient } from "@effect/sql-pglite" +import { assert, describe, expect, it } from "@effect/vitest" +import { Effect, Exit, Layer } from "effect" +import * as SqlClient from "effect/sql/SqlClient" +import * as CH from "../ch/index" +import * as Migrate from "../migrate" +import * as PG from "../postgres" +import * as S from "../schema" + +const Dashboards = S.pg.table("dashboards", { + columns: { + org_id: PG.text, + id: PG.text, + status: S.pg.column(PG.text, { default: "open" }), + tags: S.pg.column(PG.array(PG.text), { default: [] }), + layout: S.pg.column(PG.jsonb(), { default: {} }), + archived: S.pg.column(PG.bool, { default: false }), + created_at: S.pg.column(PG.timestamptz, { defaultExpr: "now()" }), + archived_at: PG.nullable(PG.timestamptz), + embedding: PG.nullable(PG.array(PG.float4)), + }, + primaryKey: { columns: ["org_id", "id"], name: "dashboards_org_id_id_pk" }, + indexes: [ + S.pg.index("dashboards_open_idx", ["org_id"], { where: ($) => $.status.in_("open", "waiting") }), + S.pg.index("dashboards_created_idx", ($) => [$.org_id, `"created_at" DESC`]), + ], +}) + +const Shares = S.pg.table("dashboard_shares", { + columns: { + org_id: PG.text, + id: PG.text, + dashboard_id: PG.text, + widget_id: PG.nullable(PG.text), + revoked_at: PG.nullable(PG.timestamptz), + }, + primaryKey: ["org_id", "id"], + indexes: [ + S.pg.uniqueIndex("dashboard_shares_live_unq", ($) => [$.org_id, $.dashboard_id, CH.coalesce($.widget_id, CH.lit(""))], { + where: "revoked_at is null", + }), + ], + foreignKeys: [ + S.pg.foreignKey({ columns: ["org_id", "dashboard_id"], references: Dashboards, foreignColumns: ["org_id", "id"], onDelete: "cascade" }), + ], +}) + +/** A generated migration from `prev` to `next`, as `generate` would write it. */ +const generated = (prev: ReadonlyArray, next: ReadonlyArray, prevIds: ReadonlyArray) => + Effect.gen(function* () { + const { ops, missingHints } = S.diffPgSchemas(prev, next) + assert.deepStrictEqual(missingHints, []) + const snapshot = yield* S.makeSnapshot(next, prevIds, "postgres") + return { + snapshot, + input: { + kind: "ops" as const, + migration: JSON.stringify({ version: "1", dialect: "postgres", ops }), + snapshot: S.serializeSnapshot(snapshot), + }, + } + }) + +const Live = Migrate.layerSqlClient().pipe(Layer.provideMerge(PgliteClient.layer({ postgresqlconf: "timezone = 'UTC'" }))) + +const tables = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient + const rows = yield* sql.unsafe<{ name: string }>( + "SELECT tablename AS name FROM pg_tables WHERE schemaname = current_schema() ORDER BY tablename", + ) + return rows.map((r) => r.name) +}) + +describe("Postgres migrations", () => { + it.effect("apply generated migrations once, record them, and verify with no drift", () => + Effect.gen(function* () { + const v1 = S.pgEntitiesOf([Dashboards]) + const v2 = S.pgEntitiesOf([Dashboards, Shares]) + const first = yield* generated([], v1, [S.ORIGIN_ID]) + const second = yield* generated(v1, v2, [first.snapshot.id]) + const migrations = yield* Migrate.fromRecord({ "20261004000000_init": first.input, "20261004000001_shares": second.input }) + + const ran = yield* Migrate.run({ migrations }) + expect(ran.map((m) => m.name)).toEqual(["20261004000000_init", "20261004000001_shares"]) + expect(yield* tables).toEqual(["_effect_orm_migrations", "dashboard_shares", "dashboards"]) + expect(yield* Migrate.run({ migrations })).toEqual([]) + expect((yield* Migrate.status(migrations)).map((s) => s.state)).toEqual(["applied", "applied"]) + + const result = yield* Migrate.verify(migrations) + expect(result).toEqual({ against: "20261004000001_shares", drift: [] }) + + // The constraints the snapshot describes are the ones the database enforces. + const sql = yield* SqlClient.SqlClient + yield* sql.unsafe(`INSERT INTO dashboards (org_id, id) VALUES ('o', 'd')`) + yield* sql.unsafe(`INSERT INTO dashboard_shares (org_id, id, dashboard_id) VALUES ('o', 's1', 'd')`) + const duplicate = yield* Effect.exit(sql.unsafe(`INSERT INTO dashboard_shares (org_id, id, dashboard_id) VALUES ('o', 's2', 'd')`)) + expect(Exit.isFailure(duplicate)).toBe(true) + yield* sql.unsafe(`DELETE FROM dashboards`) + expect(yield* sql.unsafe(`SELECT * FROM dashboard_shares`)).toEqual([]) + }).pipe(Effect.provide(Live)), + ) + + it.effect("verify reports drift in columns, defaults, indexes, keys, and extra tables", () => + Effect.gen(function* () { + const v1 = S.pgEntitiesOf([Dashboards, Shares]) + const first = yield* generated([], v1, [S.ORIGIN_ID]) + const migrations = yield* Migrate.fromRecord({ "20261004000000_init": first.input }) + yield* Migrate.run({ migrations }) + + const sql = yield* SqlClient.SqlClient + yield* sql.unsafe(`ALTER TABLE dashboards ALTER COLUMN status SET DEFAULT 'closed'`) + yield* sql.unsafe(`ALTER TABLE dashboards ALTER COLUMN archived_at SET NOT NULL`) + yield* sql.unsafe(`ALTER TABLE dashboards ADD COLUMN extra int4`) + yield* sql.unsafe(`DROP INDEX dashboards_open_idx`) + yield* sql.unsafe(`CREATE INDEX dashboards_open_idx ON dashboards (org_id) WHERE status = 'open'`) + yield* sql.unsafe(`ALTER TABLE dashboard_shares DROP CONSTRAINT dashboard_shares_org_id_dashboard_id_dashboards_org_id_id_fk`) + yield* sql.unsafe(`CREATE TABLE stray (a int)`) + + const { drift } = yield* Migrate.verify(migrations) + expect(drift.map((d) => `${d.problem} ${d.entity}`).sort()).toEqual([ + "default dashboards.status", + "index dashboards.dashboards_open_idx", + "missing dashboard_shares.dashboard_shares_org_id_dashboard_id_dashboards_org_id_id_fk", + "not_null dashboards.archived_at", + "unexpected dashboards.extra", + "unexpected stray", + ]) + expect(drift.find((d) => d.problem === "default")).toMatchObject({ expected: "'open'::text", actual: "'closed'::text" }) + expect((yield* Migrate.verify(migrations, { ignoreTables: ["stray"] })).drift).toHaveLength(5) + // The scratch schema verify builds the expectation in is rolled back. + const schemas = yield* sql.unsafe<{ n: string }>(`SELECT nspname AS n FROM pg_namespace WHERE nspname LIKE '_effect_orm_verify%'`) + expect(schemas).toEqual([]) + }).pipe(Effect.provide(Live)), + ) + + it.effect("a failed statement rolls the whole migration back, and the next run starts it again", () => + Effect.gen(function* () { + const broken = { kind: "sql" as const, migration: "CREATE TABLE a (x int)\n--> statement-breakpoint\nCREATE TABLE a (x int)" } + const failing = yield* Migrate.fromRecord({ "20261004000000_first": broken }) + const exit = yield* Effect.exit(Migrate.run({ migrations: failing, dialect: "postgres" })) + assert(Exit.isFailure(exit)) + const error = exit.cause.reasons.find((r) => r._tag === "Fail") + expect(error?._tag === "Fail" && error.error._tag).toBe("@maple-dev/effect-orm/MigrateStepFailed") + expect(yield* tables).toEqual(["_effect_orm_migrations"]) + expect((yield* Migrate.status(failing, {}, "postgres")).map((s) => s.state)).toEqual(["pending"]) + + const fixed = yield* Migrate.fromRecord({ "20261004000000_first": { kind: "sql", migration: "CREATE TABLE a (x int)" } }) + expect((yield* Migrate.run({ migrations: fixed, dialect: "postgres" })).map((m) => m.steps)).toEqual([1]) + expect(yield* tables).toEqual(["_effect_orm_migrations", "a"]) + }).pipe(Effect.provide(Live)), + ) + + it.effect("an applied migration whose file changed fails under strict", () => + Effect.gen(function* () { + const v1 = yield* Migrate.fromRecord({ "20261004000000_a": { kind: "sql", migration: "CREATE TABLE a (x int)" } }) + yield* Migrate.run({ migrations: v1, dialect: "postgres" }) + const edited = yield* Migrate.fromRecord({ "20261004000000_a": { kind: "sql", migration: "CREATE TABLE a (y int)" } }) + const exit = yield* Effect.exit(Migrate.run({ migrations: edited, dialect: "postgres", strict: true })) + assert(Exit.isFailure(exit)) + expect((yield* Migrate.status(edited, {}, "postgres"))[0]?.state).toBe("changed") + }).pipe(Effect.provide(Live)), + ) + + it.effect("baseline adopts a database another tool built: legacy SQL is recorded, not run", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient + // What drizzle-kit (or a deploy pipeline) already applied. + yield* sql.unsafe(`CREATE TABLE legacy (id text PRIMARY KEY)`) + const base = yield* S.makeSnapshot(S.pgEntitiesOf([S.pg.table("legacy", { columns: { id: PG.text }, primaryKey: ["id"] })]), [S.ORIGIN_ID], "postgres") + const migrations = yield* Migrate.fromRecord({ + "20260101000000_drizzle_init": { + kind: "sql", + migration: `CREATE TABLE "legacy" ("id" text PRIMARY KEY)`, + snapshot: JSON.stringify({ version: "8", dialect: "postgres", id: "d", prevIds: [], ddl: [], renames: [] }), + }, + "20261004000000_baseline": { kind: "sql", migration: "-- Baseline: runs nothing.\n", snapshot: S.serializeSnapshot(base) }, + "20261004000001_next": { kind: "sql", migration: "ALTER TABLE legacy ADD COLUMN name text" }, + }) + expect(migrations[0]?.foreignSnapshot?.tool).toBe("drizzle-kit") + expect(Migrate.dialectOf(migrations)).toBe("postgres") + + expect(yield* Migrate.baseline({ migrations, upTo: "20261004000000_baseline" })).toEqual([ + "20260101000000_drizzle_init", + "20261004000000_baseline", + ]) + expect((yield* Migrate.run({ migrations })).map((m) => m.name)).toEqual(["20261004000001_next"]) + // Verify compares against the baseline's snapshot, the last applied one with a snapshot. + const { against, drift } = yield* Migrate.verify(migrations) + expect(against).toBe("20261004000000_baseline") + expect(drift).toEqual([{ entity: "legacy.name", problem: "unexpected" }]) + }).pipe(Effect.provide(Live)), + ) + + it.effect("resolveStep has nothing to resolve on Postgres", () => + Effect.gen(function* () { + const [migration] = yield* Migrate.fromRecord({ "20261004000000_a": { kind: "sql", migration: "SELECT 1" } }) + const exit = yield* Effect.exit(Migrate.resolveStep({ migration: migration!, step: "0", outcome: "ran", dialect: "postgres" })) + expect(Exit.isFailure(exit)).toBe(true) + }).pipe(Effect.provide(Live)), + ) +}) diff --git a/src/migrate/pg-verify.ts b/src/migrate/pg-verify.ts new file mode 100644 index 0000000..0abbf49 --- /dev/null +++ b/src/migrate/pg-verify.ts @@ -0,0 +1,232 @@ +// Postgres drift: the live schema against the snapshot of the last applied +// migration. +// +// Comparing SQL text is hopeless in Postgres: the catalog stores a default of +// `'open'` as `'open'::text`, a predicate `status in ('a', 'b')` as +// `(status = ANY (ARRAY['a'::text, 'b'::text]))`. So nothing is normalized +// here. The snapshot is built in a scratch schema, inside a transaction that +// is always rolled back, and the two catalogs are read with the same queries: +// the server deparses both sides, and equal definitions read back equal. +// +// Checked: tables, columns (type, NOT NULL, default), primary keys, indexes +// (uniqueness, method, keys, predicate), foreign keys (columns, target, +// actions). Not checked: anything the entity model does not hold (triggers, +// grants, check constraints), and tables outside the current schema. + +import { Effect, Schema } from "effect" +import type { PgSnapshot } from "../schema/entities" +import { renderPgSchema } from "../schema/pg-ops" +import type { PgReferentialAction } from "../schema/pg-entities" +import { MigrationDriver } from "./driver" +import type { MigrateSourceError, MigrateSqlError } from "./errors" +import { LEDGER_TABLES } from "./ledger" +import { transactionOf } from "./pg-ledger" +import type { Drift } from "./verify" + +const q = (value: string): string => `'${value.replace(/'/g, "''")}'` + +interface CatalogTable { + readonly primaryKey: { readonly name: string; readonly columns: ReadonlyArray } | null + readonly columns: Map + readonly indexes: Map; readonly where: string | null }> + readonly foreignKeys: Map< + string, + { + readonly columns: ReadonlyArray + readonly foreignTable: string + readonly foreignColumns: ReadonlyArray + readonly onDelete: PgReferentialAction + readonly onUpdate: PgReferentialAction + } + > +} + +const ACTIONS: Readonly> = { + a: "NO ACTION", + r: "RESTRICT", + c: "CASCADE", + n: "SET NULL", + d: "SET DEFAULT", +} + +const jsonList = (value: unknown): ReadonlyArray => + value === null || value === undefined ? [] : (JSON.parse(String(value)) as Array).map(String) + +/** Every table of `schema`, read from the catalog. */ +const readCatalog = (schema: string) => + Effect.gen(function* () { + const driver = yield* MigrationDriver + const ns = `n.nspname = ${q(schema)}` + const tables = yield* driver.query( + `SELECT c.relname AS name FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace WHERE ${ns} AND c.relkind IN ('r', 'p')`, + ) + const columns = yield* driver.query( + `SELECT c.relname AS tbl, a.attname AS name, format_type(a.atttypid, a.atttypmod) AS type, a.attnotnull AS not_null, pg_get_expr(d.adbin, d.adrelid) AS dflt, a.attidentity::text AS ident + FROM pg_attribute a + JOIN pg_class c ON c.oid = a.attrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + LEFT JOIN pg_attrdef d ON d.adrelid = a.attrelid AND d.adnum = a.attnum + WHERE ${ns} AND c.relkind IN ('r', 'p') AND a.attnum > 0 AND NOT a.attisdropped`, + ) + const constraints = yield* driver.query( + `SELECT con.conname AS name, c.relname AS tbl, con.contype AS type, fc.relname AS foreign_table, + con.confdeltype AS on_delete, con.confupdtype AS on_update, + (SELECT json_agg(att.attname::text ORDER BY k.ord) FROM unnest(con.conkey) WITH ORDINALITY k(attnum, ord) + JOIN pg_attribute att ON att.attrelid = con.conrelid AND att.attnum = k.attnum)::text AS cols, + (SELECT json_agg(att.attname::text ORDER BY k.ord) FROM unnest(con.confkey) WITH ORDINALITY k(attnum, ord) + JOIN pg_attribute att ON att.attrelid = con.confrelid AND att.attnum = k.attnum)::text AS foreign_cols + FROM pg_constraint con + JOIN pg_class c ON c.oid = con.conrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + LEFT JOIN pg_class fc ON fc.oid = con.confrelid + WHERE ${ns} AND con.contype IN ('p', 'f')`, + ) + // Key parts one by one, with their sort order: pg_get_indexdef(oid, k, true) + // writes a key without DESC / NULLS, which live in indoption. + const indexes = yield* driver.query( + `SELECT i.relname AS name, t.relname AS tbl, ix.indisunique AS is_unique, am.amname AS method, + pg_get_expr(ix.indpred, ix.indrelid, true) AS predicate, + (SELECT json_agg( + pg_get_indexdef(ix.indexrelid, k, true) + || CASE WHEN ix.indoption[k - 1] & 1 = 1 THEN ' DESC' ELSE '' END + || CASE + WHEN ix.indoption[k - 1] & 1 = 1 AND ix.indoption[k - 1] & 2 = 0 THEN ' NULLS LAST' + WHEN ix.indoption[k - 1] & 1 = 0 AND ix.indoption[k - 1] & 2 = 2 THEN ' NULLS FIRST' + ELSE '' END + ORDER BY k) FROM generate_series(1, ix.indnkeyatts) k)::text AS keys + FROM pg_index ix + JOIN pg_class i ON i.oid = ix.indexrelid + JOIN pg_class t ON t.oid = ix.indrelid + JOIN pg_namespace n ON n.oid = t.relnamespace + JOIN pg_am am ON am.oid = i.relam + WHERE ${ns} AND NOT EXISTS (SELECT 1 FROM pg_constraint con WHERE con.conindid = ix.indexrelid AND con.contype IN ('p', 'u', 'x'))`, + ) + + const out = new Map() + for (const row of tables) { + out.set(String(row.name), { primaryKey: null, columns: new Map(), indexes: new Map(), foreignKeys: new Map() }) + } + for (const row of columns) { + out.get(String(row.tbl))?.columns.set(String(row.name), { + type: String(row.type), + notNull: row.not_null === true || row.not_null === "t", + default: row.dflt === null || row.dflt === undefined ? null : String(row.dflt), + identity: row.ident === "a" ? "always" : row.ident === "d" ? "by default" : "none", + }) + } + for (const row of constraints) { + const table = out.get(String(row.tbl)) + if (table === undefined) continue + if (row.type === "p") { + out.set(String(row.tbl), { ...table, primaryKey: { name: String(row.name), columns: jsonList(row.cols) } }) + } else { + table.foreignKeys.set(String(row.name), { + columns: jsonList(row.cols), + foreignTable: String(row.foreign_table), + foreignColumns: jsonList(row.foreign_cols), + onDelete: ACTIONS[String(row.on_delete)] ?? "NO ACTION", + onUpdate: ACTIONS[String(row.on_update)] ?? "NO ACTION", + }) + } + } + for (const row of indexes) { + out.get(String(row.tbl))?.indexes.set(String(row.name), { + unique: row.is_unique === true || row.is_unique === "t", + method: String(row.method), + keys: jsonList(row.keys), + where: row.predicate === null || row.predicate === undefined ? null : String(row.predicate), + }) + } + return out + }) + +/** Rolls the scratch transaction back once its catalog is read. */ +class Rollback extends Schema.TaggedError()("@maple-dev/effect-orm/VerifyRollback", { + expected: Schema.Unknown, +}) {} + +/** The snapshot as Postgres would store it: built in a scratch schema and rolled back. */ +const expectedCatalog = (snapshot: PgSnapshot, migration: string) => + Effect.gen(function* () { + const driver = yield* MigrationDriver + const transaction = yield* transactionOf(driver, migration) + const scratch = `_effect_orm_verify_${globalThis.crypto.randomUUID().replace(/-/g, "").slice(0, 16)}` + const result = yield* transaction( + Effect.gen(function* () { + yield* driver.execute(`CREATE SCHEMA "${scratch}"`) + // pg_catalog stays implicitly first, so built-ins resolve as they do in the real schema. + yield* driver.execute(`SET LOCAL search_path TO "${scratch}"`) + for (const statement of renderPgSchema(snapshot.entities)) yield* driver.execute(statement) + const expected = yield* readCatalog(scratch) + return yield* new Rollback({ expected }) + }), + ).pipe(Effect.flip) + if (result._tag !== "@maple-dev/effect-orm/VerifyRollback") return yield* Effect.fail(result) + return result.expected as Map + }) + +const same = (a: unknown, b: unknown): boolean => JSON.stringify(a) === JSON.stringify(b) +const show = (value: unknown): string => (typeof value === "string" ? value : JSON.stringify(value)) + +export interface PgVerifyOptions { + /** Tables in the schema that are not this schema's business (another tool's ledger, say). */ + readonly ignoreTables?: ReadonlyArray +} + +export const verifyPg = ( + snapshot: PgSnapshot, + migration: string, + options: PgVerifyOptions = {}, +): Effect.Effect, MigrateSqlError | MigrateSourceError, MigrationDriver> => + Effect.gen(function* () { + const driver = yield* MigrationDriver + const schema = String((yield* driver.query("SELECT current_schema() AS s"))[0]?.s ?? "public") + const actual = yield* readCatalog(schema) + const expected = yield* expectedCatalog(snapshot, migration) + const ignored = new Set([...Object.values(LEDGER_TABLES), ...(options.ignoreTables ?? [])]) + const drift: Array = [] + const push = (entity: string, problem: Drift["problem"], expectedValue?: unknown, actualValue?: unknown) => + drift.push({ + entity, + problem, + ...(expectedValue !== undefined ? { expected: show(expectedValue) } : undefined), + ...(actualValue !== undefined ? { actual: show(actualValue) } : undefined), + }) + + for (const [name, want] of expected) { + const have = actual.get(name) + if (have === undefined) { + push(name, "missing") + continue + } + if (!same(want.primaryKey, have.primaryKey)) push(name, "primary_key", want.primaryKey ?? "none", have.primaryKey ?? "none") + for (const [column, w] of want.columns) { + const h = have.columns.get(column) + const key = `${name}.${column}` + if (h === undefined) push(key, "missing") + else { + if (w.type !== h.type) push(key, "type", w.type, h.type) + if (w.notNull !== h.notNull) push(key, "not_null", w.notNull ? "NOT NULL" : "NULL", h.notNull ? "NOT NULL" : "NULL") + if (w.default !== h.default) push(key, "default", w.default ?? "none", h.default ?? "none") + if (w.identity !== h.identity) push(key, "identity", w.identity, h.identity) + } + } + for (const column of have.columns.keys()) if (!want.columns.has(column)) push(`${name}.${column}`, "unexpected") + for (const [index, w] of want.indexes) { + const h = have.indexes.get(index) + if (h === undefined) push(`${name}.${index}`, "missing") + else if (!same(w, h)) push(`${name}.${index}`, "index", w, h) + } + for (const index of have.indexes.keys()) if (!want.indexes.has(index)) push(`${name}.${index}`, "unexpected") + for (const [fk, w] of want.foreignKeys) { + const h = have.foreignKeys.get(fk) + if (h === undefined) push(`${name}.${fk}`, "missing") + else if (!same(w, h)) push(`${name}.${fk}`, "foreign_key", w, h) + } + for (const fk of have.foreignKeys.keys()) if (!want.foreignKeys.has(fk)) push(`${name}.${fk}`, "unexpected") + } + for (const name of actual.keys()) { + if (!expected.has(name) && !ignored.has(name)) push(name, "unexpected") + } + return drift + }) diff --git a/src/migrate/run.ts b/src/migrate/run.ts index 8104e62..f227171 100644 --- a/src/migrate/run.ts +++ b/src/migrate/run.ts @@ -1,5 +1,9 @@ // Running migrations. // +// ClickHouse and Postgres differ here more than anywhere else. Postgres DDL is +// transactional, so a Postgres migration runs whole or not at all (`runPg`). +// ClickHouse's is not, which is what the rest of this comment is about. +// // Unlike Drizzle's migrator (one transaction, apply by name) and Effect's // (insert the ledger rows first, then run, inside a transaction), this one // assumes nothing is atomic. Each statement is journaled after it finishes and @@ -7,7 +11,7 @@ // rerun resumes at the statement that failed. import { Effect, Option } from "effect" -import { sha256Hex } from "../schema/entities" +import { sha256Hex, type SchemaDialect } from "../schema/entities" import type { RenderOptions } from "../schema/render" import { MigrationDriver } from "./driver" import { MigrateHashMismatch, MigrateLeaseHeld, MigrateSourceError, MigrateStepChanged, MigrateStepFailed, MigrateStepUncertain, type MigrateError } from "./errors" @@ -21,19 +25,29 @@ import { writeLease, type AppliedRow, } from "./ledger" -import { stepsOf, type LoadedMigration } from "./source" +import { ensurePgLedger, isPgApplied, lockPg, readPgApplied, recordPgMigration, transactionOf } from "./pg-ledger" +import { dialectOf, stepsOf, type LoadedMigration } from "./source" export interface RunOptions { readonly migrations: ReadonlyArray + /** + * The database being migrated. Read from the migrations' snapshots when + * omitted; a folder of hand-written SQL alone needs it said. Default + * `clickhouse`. + */ + readonly dialect?: SchemaDialect readonly render?: RenderOptions /** Fail when an applied migration's file has changed. Otherwise log a warning. */ readonly strict?: boolean /** Names this run in the lease. Defaults to a random id. */ readonly owner?: string - /** How long the lease lasts without renewal. Each finished statement renews it. */ + /** How long the lease lasts without renewal. Each finished statement renews it. ClickHouse only. */ readonly leaseSeconds?: number } +const resolveDialect = (migrations: ReadonlyArray, dialect: SchemaDialect | undefined): SchemaDialect => + dialect ?? dialectOf(migrations) ?? "clickhouse" + export interface AppliedMigration { readonly name: string readonly steps: number @@ -148,6 +162,9 @@ const applyOne = (migration: LoadedMigration, render: RenderOptions, renew: Effe * have their hash checked. Returns what ran. */ export const run = (options: RunOptions): Effect.Effect, MigrateError, MigrationDriver> => + resolveDialect(options.migrations, options.dialect) === "postgres" ? runPg(options) : runClickHouse(options) + +const runClickHouse = (options: RunOptions): Effect.Effect, MigrateError, MigrationDriver> => Effect.gen(function* () { const render = options.render ?? {} const owner = options.owner ?? `effect-orm-${globalThis.crypto.randomUUID()}` @@ -170,12 +187,101 @@ export const run = (options: RunOptions): Effect.Effect, MigrateError, MigrationDriver> => + Effect.gen(function* () { + const driver = yield* MigrationDriver + yield* ensurePgLedger + const applied = new Map((yield* readPgApplied).map((row) => [row.name, row])) + const ran: Array = [] + for (const migration of options.migrations) { + const row = applied.get(migration.name) + if (row !== undefined) { + yield* checkHash(migration, row, options.strict ?? false) + continue + } + const transaction = yield* transactionOf(driver, migration.name) + const steps = stepsOf(migration) + const result = yield* transaction( + Effect.gen(function* () { + yield* lockPg + if (yield* isPgApplied(migration.name)) return undefined + for (const step of steps) { + yield* driver.execute(step.sql).pipe( + Effect.mapError( + (cause) => + new MigrateStepFailed({ + migration: migration.name, + step: step.id, + sql: step.sql, + message: `${migration.name} step ${step.id} failed, and the migration was rolled back: ${cause.message}`, + cause, + }), + ), + Effect.withSpan("effect_orm.migrate.step", { attributes: { "effect_orm.migration.step": step.id } }), + ) + } + yield* recordPgMigration(migration.name, migration.hash) + const done: AppliedMigration = { name: migration.name, steps: steps.length, resumedSteps: 0 } + return done + }), + ).pipe(Effect.withSpan("effect_orm.migrate.migration", { attributes: { "effect_orm.migration.name": migration.name } })) + if (result !== undefined) ran.push(result) + } + return ran + }).pipe(Effect.withSpan("effect_orm.migrate")) + +export interface BaselineOptions { + readonly migrations: ReadonlyArray + /** The last migration the database already has. It and every migration before it are recorded as applied. */ + readonly upTo: string + readonly dialect?: SchemaDialect + readonly render?: RenderOptions +} + +/** + * Record migrations as applied without running them, for a database whose + * schema another tool (drizzle-kit, a deploy pipeline) already built. Returns + * the names it recorded; ones already in the ledger are left alone. + */ +export const baseline = (options: BaselineOptions): Effect.Effect, MigrateError, MigrationDriver> => + Effect.gen(function* () { + const index = options.migrations.findIndex((m) => m.name === options.upTo) + if (index === -1) return yield* new MigrateSourceError({ migration: options.upTo, message: "is not one of the migrations" }) + const postgres = resolveDialect(options.migrations, options.dialect) === "postgres" + if (postgres) yield* ensurePgLedger + else yield* ensureLedger(options.render ?? {}) + const applied = new Set((yield* postgres ? readPgApplied : readApplied).map((row) => row.name)) + const recorded: Array = [] + for (const migration of options.migrations.slice(0, index + 1)) { + if (applied.has(migration.name)) continue + yield* postgres ? recordPgMigration(migration.name, migration.hash) : recordMigration(migration.name, migration.hash) + recorded.push(migration.name) + } + return recorded + }) + /** Where each migration stands, without running anything. */ export const status = ( migrations: ReadonlyArray, render: RenderOptions = {}, + dialect?: SchemaDialect, ): Effect.Effect, MigrateError, MigrationDriver> => Effect.gen(function* () { + if (resolveDialect(migrations, dialect) === "postgres") { + yield* ensurePgLedger + const applied = new Map((yield* readPgApplied).map((row) => [row.name, row])) + return migrations.map((migration): MigrationStatus => { + const row = applied.get(migration.name) + if (row === undefined) return { name: migration.name, state: "pending", appliedAt: undefined } + return { name: migration.name, state: row.hash === migration.hash ? "applied" : "changed", appliedAt: row.appliedAt } + }) + } yield* ensureLedger(render) const applied = new Map((yield* readApplied).map((row) => [row.name, row])) return yield* Effect.forEach(migrations, (migration) => @@ -203,6 +309,7 @@ export interface ResolveStepOptions { /** What checking the database showed: the statement took effect, or it did not. */ readonly outcome: "ran" | "not-ran" readonly render?: RenderOptions + readonly dialect?: SchemaDialect } /** @@ -213,6 +320,12 @@ export interface ResolveStepOptions { export const resolveStep = (options: ResolveStepOptions): Effect.Effect => Effect.gen(function* () { const { migration } = options + if (resolveDialect([migration], options.dialect) === "postgres") { + return yield* new MigrateSourceError({ + migration: migration.name, + message: "a Postgres migration runs in one transaction, so no step of it is ever uncertain", + }) + } const step = stepsOf(migration, options.render ?? {}).find((s) => s.id === options.step) if (step === undefined) { return yield* new MigrateSourceError({ migration: migration.name, message: `has no step ${options.step}` }) diff --git a/src/migrate/source.ts b/src/migrate/source.ts index 02b229c..11e1e16 100644 --- a/src/migrate/source.ts +++ b/src/migrate/source.ts @@ -3,11 +3,14 @@ // A migration is a folder name plus its file: `migration.json` (generated ops, // rendered when it runs) or `migration.sql` (hand-written, statements split on // `--> statement-breakpoint`, the drizzle-kit separator). Its `snapshot.json` -// is optional at run time; `verify` needs it, and ordering prefers it. +// is optional at run time; `verify` needs it, and ordering prefers it. A +// drizzle-kit `snapshot.json` is recognized and kept aside, so a drizzle folder +// loads as it is. import { Effect, FileSystem, Path, Schema } from "effect" -import { canonicalJson, sha256Hex, Snapshot } from "../schema/entities" +import { canonicalJson, sha256Hex, Snapshot, type SchemaDialect } from "../schema/entities" import { MigrationFile, renderOp } from "../schema/ops" +import { renderPgOp } from "../schema/pg-ops" import type { RenderOptions } from "../schema/render" import { MigrateSourceError } from "./errors" @@ -32,6 +35,11 @@ export interface LoadedMigration { readonly file: MigrationFile | undefined readonly sql: ReadonlyArray readonly snapshot: Snapshot | undefined + /** + * A `snapshot.json` another tool wrote (drizzle-kit's, in a folder being + * adopted), parsed but not interpreted. `snapshot` is `undefined` then. + */ + readonly foreignSnapshot: { readonly tool: "drizzle-kit"; readonly json: unknown } | undefined } export interface MigrationStep { @@ -47,8 +55,9 @@ const load = (name: string, input: MigrationInput): Effect.Effect (cause: unknown) => new MigrateSourceError({ migration: name, message, cause }) const file = input.kind === "ops" ? yield* decodeFile(input.migration).pipe(Effect.mapError(fail("migration.json does not decode"))) : undefined + const foreignSnapshot = input.snapshot === undefined ? undefined : foreignSnapshotOf(input.snapshot) const snapshot = - input.snapshot === undefined + input.snapshot === undefined || foreignSnapshot !== undefined ? undefined : yield* decodeSnapshot(input.snapshot).pipe(Effect.mapError(fail("snapshot.json does not decode"))) // Hash the decoded content for ops, so reformatting the JSON is not an edit. @@ -61,9 +70,22 @@ const load = (name: string, input: MigrationInput): Effect.Effect statement.trim()) .filter((statement) => statement.replace(/^\s*--.*$/gm, "").trim().length > 0) : [] - return { name, kind: input.kind, hash, file, sql, snapshot } + return { name, kind: input.kind, hash, file, sql, snapshot, foreignSnapshot } }) +/** drizzle-kit's snapshot: `ddl` rather than `entities`, and its own `version`. */ +const foreignSnapshotOf = (text: string): LoadedMigration["foreignSnapshot"] => { + let json: unknown + try { + json = JSON.parse(text) + } catch { + return undefined + } + return typeof json === "object" && json !== null && "ddl" in json && Array.isArray(json.ddl) && !("entities" in json) + ? { tool: "drizzle-kit", json } + : undefined +} + /** * Each migration's parent migrations, from its snapshot's `prevIds`. A custom * migration copies its parent's entities, so several migrations can share one @@ -177,4 +199,20 @@ export const fromFileSystem = ( export const stepsOf = (migration: LoadedMigration, render: RenderOptions = {}): ReadonlyArray => migration.file === undefined ? migration.sql.map((sql, i) => ({ id: String(i), sql })) - : migration.file.ops.flatMap((op, i) => renderOp(op, render).map((sql, j) => ({ id: `${i}.${j}`, sql }))) + : "dialect" in migration.file + ? migration.file.ops.flatMap((op, i) => renderPgOp(op).map((sql, j) => ({ id: `${i}.${j}`, sql }))) + : migration.file.ops.flatMap((op, i) => renderOp(op, render).map((sql, j) => ({ id: `${i}.${j}`, sql }))) + +/** + * The dialect a set of migrations is for, read from their snapshots and + * generated files. `undefined` for a folder of hand-written SQL alone, which + * says nothing about its database. + */ +export const dialectOf = (migrations: ReadonlyArray): SchemaDialect | undefined => { + for (const m of migrations) { + if (m.snapshot !== undefined) return m.snapshot.dialect + if (m.file !== undefined) return "dialect" in m.file ? "postgres" : "clickhouse" + if (m.foreignSnapshot !== undefined) return "postgres" + } + return undefined +} diff --git a/src/migrate/verify.ts b/src/migrate/verify.ts index c4768d8..354cbb2 100644 --- a/src/migrate/verify.ts +++ b/src/migrate/verify.ts @@ -12,11 +12,13 @@ import { Effect } from "effect" import { quoteClickHouseString } from "../sql/sql-fragment" -import type { ColumnEntity, IndexEntity, MaterializedViewEntity, Snapshot, TableEntity } from "../schema/entities" +import type { SchemaDialect, ClickHouseSnapshot, ColumnEntity, IndexEntity, MaterializedViewEntity, TableEntity } from "../schema/entities" import { MigrationDriver } from "./driver" -import type { MigrateSqlError } from "./errors" +import type { MigrateSourceError, MigrateSqlError } from "./errors" import { LEDGER_TABLES, readApplied } from "./ledger" -import type { LoadedMigration } from "./source" +import { readPgApplied } from "./pg-ledger" +import { verifyPg, type PgVerifyOptions } from "./pg-verify" +import { dialectOf, type LoadedMigration } from "./source" export interface Drift { readonly entity: string @@ -28,8 +30,11 @@ export interface Drift { | "partition_key" | "primary_key" | "type" + | "not_null" + | "identity" | "default" | "index" + | "foreign_key" | "view_target" | "view_body" readonly expected?: string @@ -80,15 +85,34 @@ const canonicalTypes = (driver: typeof MigrationDriver.Service, types: ReadonlyS return out }) +export interface VerifyOptions extends PgVerifyOptions { + /** Read from the migrations when omitted. */ + readonly dialect?: SchemaDialect +} + +/** + * Compare the database with the snapshot of the last applied migration. The + * dialect comes from the snapshots; see `pg-verify.ts` for how Postgres is + * compared. + */ export const verify = ( migrations: ReadonlyArray, -): Effect.Effect => + options: VerifyOptions = {}, +): Effect.Effect => Effect.gen(function* () { - const driver = yield* MigrationDriver - const appliedNames = new Set((yield* readApplied.pipe(Effect.orElseSucceed(() => []))).map((row) => row.name)) + const postgres = (options.dialect ?? dialectOf(migrations)) === "postgres" + const appliedRows = yield* (postgres ? readPgApplied : readApplied).pipe(Effect.orElseSucceed(() => [])) + const appliedNames = new Set(appliedRows.map((row) => row.name)) const against = [...migrations].reverse().find((m) => appliedNames.has(m.name) && m.snapshot !== undefined) - const snapshot: Snapshot | undefined = against?.snapshot - if (snapshot === undefined) return { against: undefined, drift: [] } + const snapshot = against?.snapshot + if (against === undefined || snapshot === undefined) return { against: undefined, drift: [] } + if (snapshot.dialect === "postgres") return { against: against.name, drift: yield* verifyPg(snapshot, against.name, options) } + return { against: against.name, drift: yield* verifyClickHouse(snapshot) } + }) + +const verifyClickHouse = (snapshot: ClickHouseSnapshot): Effect.Effect, MigrateSqlError, MigrationDriver> => + Effect.gen(function* () { + const driver = yield* MigrationDriver const db = String((yield* driver.query("SELECT currentDatabase() AS db"))[0]?.db ?? "default") const tables = yield* driver.query( @@ -204,5 +228,5 @@ export const verify = ( } } - return { against: against?.name, drift } + return drift }) diff --git a/src/schema.ts b/src/schema.ts index 3c62fc9..da4da4d 100644 --- a/src/schema.ts +++ b/src/schema.ts @@ -1,8 +1,12 @@ // @maple-dev/effect-orm/schema // // Tables and materialized views that carry their DDL, snapshots of them, and -// the offline diff that turns two snapshots into migration ops. Pure: nothing -// here reads files or opens a connection. See docs/migrations.md. +// the offline diff that turns two snapshots into migration ops. ClickHouse +// definitions are the top-level exports; Postgres ones live under `pg` +// (`S.pg.table`). Pure: nothing here reads files or opens a connection. See +// docs/migrations.md. + +export * as pg from "./schema/pg-define" export { column, @@ -28,8 +32,12 @@ export { type TableDefinition, } from "./schema/define" export { + ClickHouseSnapshot, ColumnDefault, ColumnEntity, + PgSnapshot, + type AnySchemaEntity, + type SchemaDialect, EngineSpec, IndexEntity, MaterializedViewEntity, @@ -55,6 +63,35 @@ export { renderSchema, type RenderOptions, } from "./schema/render" -export { MigrationFile, MigrationOp, labelOf, renderOp, type OpLabel } from "./schema/ops" -export { entitiesOf, isSchemaObject, makeSnapshot, serializeSnapshot, type SchemaObject } from "./schema/snapshot" +export { ClickHouseMigrationFile, MigrationFile, MigrationOp, labelOf, renderOp, type OpLabel } from "./schema/ops" +export { + dialectOfObject, + entitiesOf, + isSchemaObject, + makeSnapshot, + pgEntitiesOf, + serializeSnapshot, + type SchemaObject, +} from "./schema/snapshot" +export { + PgColumnEntity, + PgForeignKeyEntity, + PgIdentity, + PgIndexEntity, + PgReferentialAction, + PgSchemaEntity, + PgTableEntity, + canonicalPgType, +} from "./schema/pg-entities" +export { + PgMigrationFile, + PgMigrationOp, + labelOfPg, + pgIdent, + renderPgOp, + renderPgSchema, + type PgOpLabel, +} from "./schema/pg-ops" +export { diffPgSchemas } from "./schema/pg-diff" +export { fromDrizzleSnapshot, type DrizzleImport } from "./schema/drizzle" export { Hint, Hints, diffSchemas, type DiffResult, type UnsupportedChange } from "./schema/diff" diff --git a/src/schema/diff.ts b/src/schema/diff.ts index 6de748e..c59556f 100644 --- a/src/schema/diff.ts +++ b/src/schema/diff.ts @@ -42,8 +42,8 @@ export interface UnsupportedChange { readonly message: string } -export interface DiffResult { - readonly ops: ReadonlyArray +export interface DiffResult { + readonly ops: ReadonlyArray readonly missingHints: ReadonlyArray readonly unsupported: ReadonlyArray } diff --git a/src/schema/drizzle.ts b/src/schema/drizzle.ts new file mode 100644 index 0000000..7072726 --- /dev/null +++ b/src/schema/drizzle.ts @@ -0,0 +1,175 @@ +// Reading a drizzle-kit (v1, snapshot version 8) Postgres snapshot as entities. +// +// Adopting effect-orm in a folder drizzle-kit wrote starts from a baseline: a +// snapshot of the schema the database already has. drizzle-kit's last +// snapshot is exactly that, so `generate --baseline --from-drizzle` converts +// it rather than trusting the TypeScript definitions to match the database. +// The first `generate` after it then shows, as ordinary ops, every place the +// definitions and the database disagree. +// +// Only what the entity model holds converts: tables, columns, primary keys, +// indexes, foreign keys in the `public` schema. Anything else (enums, +// sequences, views, policies, check and unique constraints, identity and +// generated columns, other schemas) is reported, and nothing is written. + +import { canonicalPgType, type PgIdentity, type PgReferentialAction, type PgSchemaEntity, type PgTableEntity } from "./pg-entities" +import { sortEntities } from "./entities" + +export interface DrizzleImport { + readonly entities: ReadonlyArray + /** Objects the snapshot holds that have no entity here. Empty when the import is complete. */ + readonly unsupported: ReadonlyArray +} + +type Json = Record + +const isRecord = (value: unknown): value is Json => typeof value === "object" && value !== null && !Array.isArray(value) +const str = (value: unknown): string => (typeof value === "string" ? value : String(value)) +const strings = (value: unknown): ReadonlyArray => (Array.isArray(value) ? value.map(str) : []) + +const quote = (name: string): string => `"${name.replace(/"/g, '""')}"` + +const ACTIONS: ReadonlyArray = ["NO ACTION", "RESTRICT", "CASCADE", "SET NULL", "SET DEFAULT"] +const action = (value: unknown): PgReferentialAction => { + const upper = typeof value === "string" ? value.toUpperCase() : "NO ACTION" + return ACTIONS.find((a) => a === upper) ?? "NO ACTION" +} + +/** drizzle-kit writes a default as SQL text, or (in some versions) as `{ value, type }`. */ +const defaultOf = (value: unknown): string | null => { + if (value === null || value === undefined) return null + if (typeof value === "string") return value + if (isRecord(value) && "value" in value) return str(value.value) + return str(value) +} + +const indexPart = (part: unknown): string => { + if (typeof part === "string") return quote(part) + if (!isRecord(part)) return str(part) + let sql = part.isExpression === true ? str(part.value) : quote(str(part.value)) + if (typeof part.opclass === "string" && part.opclass.length > 0) sql += ` ${part.opclass}` + const asc = part.asc !== false + if (!asc) sql += " DESC" + // Postgres's defaults: NULLS LAST ascending, NULLS FIRST descending. + if (asc && part.nullsFirst === true) sql += " NULLS FIRST" + if (!asc && part.nullsFirst === false) sql += " NULLS LAST" + return sql +} + +const MAX: Readonly> = { smallint: "32767", integer: "2147483647", bigint: "9223372036854775807" } + +/** + * An identity's kind, or `"!"` when its sequence has options the + * entity does not hold (anything but Postgres's defaults and the default name). + */ +const identityOf = (value: unknown, table: string, column: string, type: string): PgIdentity | null | `!${string}` => { + if (value === null || value === undefined) return null + if (!isRecord(value)) return "!an identity column drizzle-kit wrote in an unknown form" + const kind = value.type === "always" ? "always" : value.type === "byDefault" ? "by default" : undefined + if (kind === undefined) return `!an identity of type ${str(value.type)}` + const defaults: Record = { + name: `${table}_${column}_seq`, + increment: "1", + startWith: "1", + minValue: "1", + maxValue: MAX[canonicalPgType(type)], + cache: 1, + cycle: false, + } + const custom = Object.entries(defaults).filter(([key, expected]) => key in value && value[key] !== undefined && value[key] !== null && String(value[key]) !== String(expected)) + return custom.length > 0 ? `!an identity whose sequence sets ${custom.map(([k]) => k).join(", ")}` : kind +} + +/** Convert a drizzle-kit Postgres snapshot (`snapshot.json`, version 8). */ +export const fromDrizzleSnapshot = (json: unknown): DrizzleImport => { + const unsupported: Array = [] + if (!isRecord(json) || !Array.isArray(json.ddl)) { + return { entities: [], unsupported: ["not a drizzle-kit snapshot: it has no ddl list"] } + } + if (json.dialect !== undefined && json.dialect !== "postgres" && json.dialect !== "postgresql") { + unsupported.push(`dialect ${str(json.dialect)}: only Postgres snapshots convert`) + } + const tables = new Map() + const entities: Array = [] + const positions = new Map() + + for (const raw of json.ddl as ReadonlyArray) { + if (!isRecord(raw)) continue + const type = str(raw.entityType) + const schema = raw.schema === undefined ? "public" : str(raw.schema) + const where = `${type} ${raw.table !== undefined ? `${str(raw.table)}.` : ""}${str(raw.name)}` + if (type === "schemas") { + if (str(raw.name) !== "public") unsupported.push(`schema ${str(raw.name)}`) + continue + } + if (schema !== "public") { + unsupported.push(`${where} (schema ${schema})`) + continue + } + switch (type) { + case "tables": { + if (raw.isRlsEnabled === true) unsupported.push(`${where}: row-level security`) + tables.set(str(raw.name), { name: str(raw.name), primaryKey: null }) + break + } + case "columns": { + const table = str(raw.table) + if (raw.typeSchema !== null && raw.typeSchema !== undefined) unsupported.push(`${where}: an enum or custom type`) + if (raw.generated !== null && raw.generated !== undefined) unsupported.push(`${where}: a generated column`) + const identity = identityOf(raw.identity, str(raw.table), str(raw.name), str(raw.type)) + if (typeof identity === "string" && identity.startsWith("!")) unsupported.push(`${where}: ${identity.slice(1)}`) + const position = positions.get(table) ?? 0 + positions.set(table, position + 1) + const dimensions = typeof raw.dimensions === "number" ? raw.dimensions : 0 + entities.push({ + kind: "column", + table, + name: str(raw.name), + position, + type: canonicalPgType(`${str(raw.type)}${"[]".repeat(dimensions)}`), + notNull: raw.notNull === true, + default: defaultOf(raw.default), + identity: identity === "always" || identity === "by default" ? identity : null, + }) + break + } + case "pks": { + const table = tables.get(str(raw.table)) + if (table === undefined) unsupported.push(`${where}: its table comes later in the snapshot`) + else table.primaryKey = { name: str(raw.name), columns: strings(raw.columns) } + break + } + case "indexes": { + if (typeof raw.with === "string" && raw.with.length > 0) unsupported.push(`${where}: WITH (${raw.with})`) + entities.push({ + kind: "index", + table: str(raw.table), + name: str(raw.name), + unique: raw.isUnique === true, + method: typeof raw.method === "string" && raw.method.length > 0 ? raw.method.toLowerCase() : "btree", + columns: Array.isArray(raw.columns) ? raw.columns.map(indexPart) : [], + where: typeof raw.where === "string" && raw.where.length > 0 ? raw.where : null, + }) + break + } + case "fks": { + if (raw.schemaTo !== undefined && raw.schemaTo !== "public") unsupported.push(`${where}: references schema ${str(raw.schemaTo)}`) + entities.push({ + kind: "foreign_key", + table: str(raw.table), + name: str(raw.name), + columns: strings(raw.columns), + foreignTable: str(raw.tableTo), + foreignColumns: strings(raw.columnsTo), + onDelete: action(raw.onDelete), + onUpdate: action(raw.onUpdate), + }) + break + } + default: + unsupported.push(where) + } + } + for (const table of tables.values()) entities.push({ kind: "table", ...table }) + return { entities: sortEntities(entities), unsupported } +} diff --git a/src/schema/entities.ts b/src/schema/entities.ts index 9c18f80..e5b087a 100644 --- a/src/schema/entities.ts +++ b/src/schema/entities.ts @@ -1,4 +1,6 @@ -// Schema entities: the normalized, serializable form of a ClickHouse schema. +// Schema entities: the normalized, serializable form of a ClickHouse schema, +// and the snapshot envelope both dialects share (Postgres entities live in +// `pg-entities.ts`). // // A `defineTable` value is code; a snapshot is data. Everything downstream of // the definitions (DDL rendering, diffing, the migrator's drift check) reads @@ -6,6 +8,7 @@ // and diffs exactly as it did when it was written. import { Schema } from "effect" +import { PgSchemaEntity } from "./pg-entities" /** `DEFAULT`, `MATERIALIZED`, or `ALIAS`, with its SQL expression. */ export const ColumnDefault = Schema.Struct({ @@ -78,34 +81,51 @@ export const SNAPSHOT_VERSION = "1" /** The parent id of a first migration. */ export const ORIGIN_ID = "0000000000000000000000000000000000000000000000000000000000000000" -export const Snapshot = Schema.Struct({ +const snapshotFields = { version: Schema.Literal(SNAPSHOT_VERSION), - dialect: Schema.Literal("clickhouse"), /** sha256 of the canonical entity list; two branches reaching one schema agree. */ id: Schema.String, prevIds: Schema.Array(Schema.String), +} + +export const ClickHouseSnapshot = Schema.Struct({ + ...snapshotFields, + dialect: Schema.Literal("clickhouse"), entities: Schema.Array(SchemaEntity), }) +export type ClickHouseSnapshot = typeof ClickHouseSnapshot.Type + +export const PgSnapshot = Schema.Struct({ + ...snapshotFields, + dialect: Schema.Literal("postgres"), + entities: Schema.Array(PgSchemaEntity), +}) +export type PgSnapshot = typeof PgSnapshot.Type + +/** The schema at one point in history. `dialect` says which entity set it holds. */ +export const Snapshot = Schema.Union([ClickHouseSnapshot, PgSnapshot]) export type Snapshot = typeof Snapshot.Type -/** A stable key per entity, unique within a snapshot. */ -export const entityKey = (entity: SchemaEntity): string => { - switch (entity.kind) { - case "table": - return `table:${entity.name}` - case "materialized_view": - return `materialized_view:${entity.name}` - case "column": - return `column:${entity.table}.${entity.name}` - case "index": - return `index:${entity.table}.${entity.name}` - } -} +/** The dialects a schema can be written for. */ +export type SchemaDialect = Snapshot["dialect"] + +/** An entity of either dialect. Everything dialect-neutral (keys, ordering, hashing, the branch graph) takes this. */ +export type AnySchemaEntity = SchemaEntity | PgSchemaEntity -const kindOrder: Record = { table: 0, column: 1, index: 2, materialized_view: 3 } +/** A stable key per entity, unique within a snapshot: `:`, or `:
.` for a table's parts. */ +export const entityKey = (entity: AnySchemaEntity): string => + "table" in entity ? `${entity.kind}:${entity.table}.${entity.name}` : `${entity.kind}:${entity.name}` + +const kindOrder: Record = { + table: 0, + column: 1, + index: 2, + foreign_key: 3, + materialized_view: 4, +} -/** Tables, then columns in declaration order, then indexes, then views. Deterministic. */ -export const sortEntities = (entities: ReadonlyArray): ReadonlyArray => +/** Tables, then columns in declaration order, then indexes, foreign keys, views. Deterministic. */ +export const sortEntities = (entities: ReadonlyArray): ReadonlyArray => [...entities].sort((a, b) => { const byKind = kindOrder[a.kind] - kindOrder[b.kind] if (byKind !== 0) return byKind diff --git a/src/schema/ops.ts b/src/schema/ops.ts index e7a312a..f431a3f 100644 --- a/src/schema/ops.ts +++ b/src/schema/ops.ts @@ -6,6 +6,7 @@ import { Schema } from "effect" import { ColumnEntity, IndexEntity, MaterializedViewEntity, TableEntity } from "./entities" +import { PgMigrationFile } from "./pg-ops" import { ident, renderAlter, @@ -46,11 +47,18 @@ export const MigrationOp = Schema.Union([ ]) export type MigrationOp = typeof MigrationOp.Type -/** The file a generated migration is written to. */ -export const MigrationFile = Schema.Struct({ +/** The file a generated ClickHouse migration is written to. */ +export const ClickHouseMigrationFile = Schema.Struct({ version: Schema.Literal("1"), ops: Schema.Array(MigrationOp), }) +export type ClickHouseMigrationFile = typeof ClickHouseMigrationFile.Type + +/** + * A generated `migration.json`, of either dialect. A Postgres file says + * `"dialect": "postgres"`; a ClickHouse one predates the field and has none. + */ +export const MigrationFile = Schema.Union([PgMigrationFile, ClickHouseMigrationFile]) export type MigrationFile = typeof MigrationFile.Type /** Labels the plan prints, so the expensive lines stand out. */ diff --git a/src/schema/pg-define.ts b/src/schema/pg-define.ts new file mode 100644 index 0000000..23cd392 --- /dev/null +++ b/src/schema/pg-define.ts @@ -0,0 +1,355 @@ +// Postgres schema definitions: `S.pg.table` and its column, index and foreign +// key helpers. +// +// The Postgres counterpart of `defineTable`: the value IS a `Table`, so every +// query API accepts it, and its DDL rides beside it on `ddl` as Postgres +// entities. Expressions are written with the query DSL (or as SQL strings) and +// rendered once, here, with the Postgres dialect. A definition that cannot +// become DDL dies as a `SchemaDefinitionDefect` while the module loads. + +import { compile as compileFragment } from "../sql/sql-fragment" +import { withDialect } from "../ch/dialect" +import type { Condition, Expr } from "../ch/expr" +import { encodeColumnLiteral } from "../ch/literal" +import { createColumnAccessor, type ColumnAccessor } from "../ch/query" +import type { Table } from "../ch/table" +import type { CHType, ColumnDefs, InferTS } from "../ch/types" +import { postgresDialect } from "../pg/dialect" +import { SchemaDefinitionDefect, type DdlExpr, type DdlKey } from "./define" +import { + canonicalPgType, + PG_MAX_IDENTIFIER, + type PgColumnEntity, + type PgForeignKeyEntity, + type PgIdentity, + type PgIndexEntity, + type PgReferentialAction, + type PgTableEntity, +} from "./pg-entities" + +const quoteIdent = (name: string): string => `"${name.replace(/"/g, '""')}"` + +const renderValue = (value: Expr | Condition | string): string => + typeof value === "string" ? value : withDialect(postgresDialect, () => compileFragment(value.toFragment())) + +/** A partial index's predicate: SQL, or a condition built with the DSL. */ +export type DdlPredicate = string | (($: ColumnAccessor) => Condition | Expr | string) + +const renderExpr = (expr: DdlExpr | DdlPredicate, columns: Cols): string => + typeof expr === "string" ? expr : renderValue(expr(createColumnAccessor(columns))) + +/** Key parts: a column name is written quoted, an expression as it renders. */ +const renderKeyParts = (key: DdlKey, columns: Cols): ReadonlyArray => + typeof key === "function" ? key(createColumnAccessor(columns)).map(renderValue) : key.map(quoteIdent) + +// Columns + +export interface ColumnOptions> { + /** A literal default, encoded through the column's own type. */ + readonly default?: InferTS + /** `DEFAULT `, for a computed default such as `"now()"`. */ + readonly defaultExpr?: DdlExpr + /** + * An identity column, numbered by its own sequence. An insert may leave it + * out; with `"always"`, Postgres also rejects an insert that gives it. + */ + readonly identity?: PgIdentity +} + +/** `Given` records which options were set, so the table knows which columns an insert may leave out. */ +export interface ColumnSpec< + T extends CHType, + Given extends keyof ColumnOptions = keyof ColumnOptions, +> { + readonly _tag: "PgColumnSpec" + readonly type: T + readonly options: ColumnOptions + readonly _given?: Given +} + +/** A column with a default. A bare column type works too where none is needed. */ +export const column = , Given extends keyof ColumnOptions = never>( + type: T, + options: ColumnOptions & { readonly [K in Given]: unknown } = {} as ColumnOptions & { + readonly [K in Given]: unknown + }, +): ColumnSpec => ({ _tag: "PgColumnSpec", type, options }) + +export type ColumnInput = CHType | ColumnSpec, any> + +/** The query-side column types of a `columns` record. */ +export type ColumnsOf> = { + readonly [K in keyof I]: I[K] extends ColumnSpec ? T : I[K] extends CHType ? I[K] : never +} + +/** Columns declared with `default`, `defaultExpr` or `identity`: an insert may leave them out. */ +export type DefaultedColumnsOf> = { + [K in keyof I]: I[K] extends ColumnSpec ? ([Given] extends [never] ? never : K) : never +}[keyof I] & + string + +const isColumnSpec = (input: ColumnInput): input is ColumnSpec, any> => + "_tag" in input && input._tag === "PgColumnSpec" + +// Indexes and foreign keys + +export interface IndexOptions { + /** A partial index: only rows this predicate holds for are indexed. */ + readonly where?: DdlPredicate + /** The access method. Default `btree`. */ + readonly using?: string +} + +export interface IndexSpec { + readonly _tag: "PgIndexSpec" + readonly name: string + readonly unique: boolean + readonly on: DdlKey + readonly options: IndexOptions +} + +/** An index on columns or expressions: `index("t_org_idx", ["org_id"])`, or a callback building expressions. */ +export const index = ( + name: string, + on: DdlKey, + options: IndexOptions = {}, +): IndexSpec => ({ _tag: "PgIndexSpec", name, unique: false, on, options }) + +/** A unique index. With `where`, a partial unique index: at most one matching row per key. */ +export const uniqueIndex = ( + name: string, + on: DdlKey, + options: IndexOptions = {}, +): IndexSpec => ({ _tag: "PgIndexSpec", name, unique: true, on, options }) + +/** Referential actions, written as drizzle writes them or as the catalog does. */ +export type ReferentialAction = + | PgReferentialAction + | "no action" + | "restrict" + | "cascade" + | "set null" + | "set default" + +export interface ForeignKeySpec { + readonly _tag: "PgForeignKeySpec" + readonly columns: ReadonlyArray + readonly references: string + readonly foreignColumns: ReadonlyArray + readonly onDelete: PgReferentialAction + readonly onUpdate: PgReferentialAction + readonly name: string | undefined +} + +/** + * A foreign key. `references` is the referenced table (a `Table`, so its + * columns are checked) or its name, for a table that references itself. + */ +export const foreignKey = (spec: { + readonly columns: ReadonlyArray + readonly references: Table | string + readonly foreignColumns: ReadonlyArray + readonly onDelete?: ReferentialAction + readonly onUpdate?: ReferentialAction + /** Default `
____fk`, the name drizzle-kit gives. */ + readonly name?: string +}): ForeignKeySpec => ({ + _tag: "PgForeignKeySpec", + columns: spec.columns, + references: typeof spec.references === "string" ? spec.references : spec.references.name, + foreignColumns: spec.foreignColumns, + onDelete: (spec.onDelete?.toUpperCase() ?? "NO ACTION") as PgReferentialAction, + onUpdate: (spec.onUpdate?.toUpperCase() ?? "NO ACTION") as PgReferentialAction, + name: spec.name, +}) + +// Tables + +export interface TableDefinition> { + readonly columns: Columns + /** Column names, or with a constraint name. The default name is `
_pkey`, Postgres's own. */ + readonly primaryKey?: + | ReadonlyArray + | { readonly columns: ReadonlyArray; readonly name?: string } + readonly indexes?: ReadonlyArray>> + readonly foreignKeys?: ReadonlyArray> + /** As for `table()`: the column carrying row-level tenancy. */ + readonly tenantColumn?: keyof Columns & string +} + +/** The DDL a Postgres table carries, as entities. */ +export interface TableDdl { + readonly dialect: "postgres" + readonly table: PgTableEntity + readonly columns: ReadonlyArray + readonly indexes: ReadonlyArray + readonly foreignKeys: ReadonlyArray +} + +export interface PgSchemaTable + extends Table { + readonly ddl: TableDdl +} + +const IDENTIFIER = /^[A-Za-z_][A-Za-z0-9_$]*$/ + +const assertIdentifier = (object: string, name: string): void => { + if (!IDENTIFIER.test(name)) { + throw new SchemaDefinitionDefect({ + object, + message: `${JSON.stringify(name)} is not a plain identifier ([A-Za-z_][A-Za-z0-9_$]*)`, + }) + } + if (name.length > PG_MAX_IDENTIFIER) { + throw new SchemaDefinitionDefect({ + object, + message: `${JSON.stringify(name)} is longer than ${PG_MAX_IDENTIFIER} characters, which Postgres would truncate`, + }) + } +} + +/** `TRUE` is how the dialect writes a literal; the catalog and drizzle-kit write `true`. */ +const lowerKeywords = (sql: string): string => (/^(TRUE|FALSE|NULL)$/.test(sql) ? sql.toLowerCase() : sql) + +const columnDefault = ( + table: string, + name: string, + spec: ColumnSpec>, + columns: ColumnDefs, +): string | null => { + const { options } = spec + const given = [options.default !== undefined, options.defaultExpr !== undefined, options.identity !== undefined].filter(Boolean) + if (given.length > 1) { + throw new SchemaDefinitionDefect({ object: `${table}.${name}`, message: "a column takes one of default, defaultExpr, identity" }) + } + if (options.defaultExpr !== undefined) return renderExpr(options.defaultExpr, columns) + if (options.default === undefined || options.default === null) return null + const literal = withDialect(postgresDialect, () => encodeColumnLiteral(spec.type, options.default, name)) + return lowerKeywords(literal) +} + +/** + * A Postgres table with its DDL. Usable everywhere a `table()` is; `generate` + * reads its `ddl` when the config's dialect is `postgres`. + * + * A column is `NOT NULL` unless its type is `PG.nullable(...)`. + */ +export function table>( + name: Name, + definition: TableDefinition, +): PgSchemaTable, DefaultedColumnsOf> { + assertIdentifier(name, name) + const inputs = Object.entries(definition.columns) + if (inputs.length === 0) throw new SchemaDefinitionDefect({ object: name, message: "a table needs columns" }) + const types = Object.fromEntries( + inputs.map(([column, input]) => [column, isColumnSpec(input) ? input.type : input]), + ) as ColumnsOf + const has = (column: string) => Object.hasOwn(types, column) + + const columnEntities = inputs.map(([column, input], position): PgColumnEntity => { + assertIdentifier(`${name}.${column}`, column) + const spec: ColumnSpec> = isColumnSpec(input) + ? input + : { _tag: "PgColumnSpec", type: input as CHType, options: {} } + return { + kind: "column", + table: name, + name: column, + position, + type: canonicalPgType(spec.type.sql), + notNull: spec.type._tag !== "Nullable", + default: columnDefault(name, column, spec, types), + identity: spec.options.identity ?? null, + } + }) + for (const column of columnEntities) { + if (column.identity !== null && !["smallint", "integer", "bigint"].includes(column.type)) { + throw new SchemaDefinitionDefect({ object: `${name}.${column.name}`, message: `an identity column must be smallint, integer or bigint, not ${column.type}` }) + } + if (column.identity !== null && !column.notNull) { + throw new SchemaDefinitionDefect({ object: `${name}.${column.name}`, message: "an identity column cannot be nullable" }) + } + } + + const pk = definition.primaryKey + const pkColumns = pk === undefined ? [] : "columns" in pk ? pk.columns : pk + const pkName = pk !== undefined && "columns" in pk && pk.name !== undefined ? pk.name : `${name}_pkey` + if (pk !== undefined) { + assertIdentifier(`${name} primary key`, pkName) + if (pkColumns.length === 0) throw new SchemaDefinitionDefect({ object: name, message: "a primary key needs columns" }) + for (const column of pkColumns) { + if (!has(column)) throw new SchemaDefinitionDefect({ object: `${name} primary key`, message: `${column} is not a column` }) + if (columnEntities.find((c) => c.name === column)?.notNull === false) { + throw new SchemaDefinitionDefect({ + object: `${name}.${column}`, + message: "a primary key column cannot be nullable; Postgres would make it NOT NULL anyway", + }) + } + } + } + + const indexEntities = (definition.indexes ?? []).map((spec): PgIndexEntity => { + assertIdentifier(`${name} index`, spec.name) + if (Array.isArray(spec.on)) { + for (const column of spec.on) { + if (!has(column)) throw new SchemaDefinitionDefect({ object: `${name} index ${spec.name}`, message: `${column} is not a column` }) + } + } + const columns = renderKeyParts(spec.on, types) + if (columns.length === 0) throw new SchemaDefinitionDefect({ object: `${name} index ${spec.name}`, message: "an index needs columns" }) + return { + kind: "index", + table: name, + name: spec.name, + unique: spec.unique, + method: (spec.options.using ?? "btree").toLowerCase(), + columns, + where: spec.options.where === undefined ? null : renderExpr(spec.options.where, types), + } + }) + + const foreignKeyEntities = (definition.foreignKeys ?? []).map((spec): PgForeignKeyEntity => { + const fkName = spec.name ?? `${name}_${spec.columns.join("_")}_${spec.references}_${spec.foreignColumns.join("_")}_fk` + assertIdentifier(`${name} foreign key`, fkName) + for (const column of spec.columns) { + if (!has(column)) throw new SchemaDefinitionDefect({ object: `${name} foreign key ${fkName}`, message: `${column} is not a column` }) + } + if (spec.columns.length === 0 || spec.columns.length !== spec.foreignColumns.length) { + throw new SchemaDefinitionDefect({ + object: `${name} foreign key ${fkName}`, + message: "columns and foreignColumns need the same, non-zero, length", + }) + } + return { + kind: "foreign_key", + table: name, + name: fkName, + columns: [...spec.columns], + foreignTable: spec.references, + foreignColumns: [...spec.foreignColumns], + onDelete: spec.onDelete, + onUpdate: spec.onUpdate, + } + }) + + const tableEntity: PgTableEntity = { + kind: "table", + name, + primaryKey: pk === undefined ? null : { name: pkName, columns: [...pkColumns] }, + } + const defaults = columnEntities.filter((c) => c.default !== null || c.identity !== null).map((c) => c.name) + return { + _tag: "Table", + name, + columns: types, + ...(definition.tenantColumn !== undefined ? { tenantColumn: definition.tenantColumn } : undefined), + ...(defaults.length > 0 ? { defaults: defaults as unknown as Array> } : undefined), + ddl: { + dialect: "postgres", + table: tableEntity, + columns: columnEntities, + indexes: indexEntities, + foreignKeys: foreignKeyEntities, + }, + } +} diff --git a/src/schema/pg-diff.ts b/src/schema/pg-diff.ts new file mode 100644 index 0000000..97075be --- /dev/null +++ b/src/schema/pg-diff.ts @@ -0,0 +1,161 @@ +// Diff two Postgres schemas into migration ops. +// +// Pure and offline, like `diffSchemas` for ClickHouse, with the same data-loss +// rule: dropping a table or a column needs a `confirm_data_loss` hint. Postgres +// can change almost anything else in place (a column's type, a key, a default), +// so nothing is reported as unsupported; a change the data cannot take (a type +// cast that fails, NOT NULL over existing nulls) fails when the migration runs, +// inside its transaction, and leaves the database as it was. +// +// Renames are not detected: a renamed column reads as a drop plus an add, and +// the drop asks for confirmation. + +import { canonicalJson, entityKey } from "./entities" +import type { DiffResult, Hint } from "./diff" +import type { PgColumnEntity, PgForeignKeyEntity, PgIndexEntity, PgSchemaEntity, PgTableEntity } from "./pg-entities" +import type { PgMigrationOp } from "./pg-ops" + +const hintKey = (hint: Hint): string => `${hint.type}:${hint.kind}:${hint.entity}` + +const same = (a: unknown, b: unknown): boolean => canonicalJson(a) === canonicalJson(b) + +const byKind = (entities: ReadonlyArray, kind: K) => + new Map( + entities + .filter((e): e is Extract => e.kind === kind) + .map((e) => [entityKey(e), e] as const), + ) + +// Constraints and indexes go first, so a dropped column or table is not still +// referenced; new tables exist before the indexes and keys that need them. +const opOrder: Record = { + drop_foreign_key: 0, + drop_index: 1, + drop_column: 2, + drop_table: 3, + create_table: 4, + add_column: 5, + alter_column: 6, + set_primary_key: 7, + create_index: 8, + add_foreign_key: 9, +} + +export const diffPgSchemas = ( + prev: ReadonlyArray, + next: ReadonlyArray, + hints: ReadonlyArray = [], +): DiffResult => { + const given = new Set(hints.map(hintKey)) + const ops: Array = [] + const missingHints: Array = [] + + const confirm = (hint: Hint, op: PgMigrationOp): void => { + if (given.has(hintKey(hint))) ops.push(op) + else missingHints.push(hint) + } + + const prevTables = byKind(prev, "table") + const nextTables = byKind(next, "table") + const prevColumns = byKind(prev, "column") + const nextColumns = byKind(next, "column") + + const columnsOf = (columns: Map, table: string): ReadonlyArray => + [...columns.values()].filter((c) => c.table === table).sort((a, b) => a.position - b.position) + + for (const [key, table] of nextTables) { + if (!prevTables.has(key)) ops.push({ op: "create_table", table, columns: columnsOf(nextColumns, table.name) }) + } + for (const [key, table] of prevTables) { + if (!nextTables.has(key)) { + confirm({ type: "confirm_data_loss", kind: "table", entity: table.name }, { op: "drop_table", name: table.name }) + } + } + + for (const [key, after] of nextTables) { + const before = prevTables.get(key) + if (before === undefined) continue + diffColumns(columnsOf(prevColumns, after.name), columnsOf(nextColumns, after.name), after.name, ops, confirm) + diffPrimaryKey(before, after, ops) + } + + // Indexes and foreign keys are diffed whole, across tables: one that moved + // tables, or belongs to a table that was dropped and re-created, is a drop + // and a create like any other change. + diffNamed(byKind(prev, "index"), byKind(next, "index"), prevTables, ops, { + drop: (i: PgIndexEntity) => ({ op: "drop_index", table: i.table, name: i.name }), + create: (i: PgIndexEntity) => ({ op: "create_index", index: i }), + }) + diffNamed(byKind(prev, "foreign_key"), byKind(next, "foreign_key"), prevTables, ops, { + drop: (fk: PgForeignKeyEntity) => ({ op: "drop_foreign_key", table: fk.table, name: fk.name }), + create: (fk: PgForeignKeyEntity) => ({ op: "add_foreign_key", foreignKey: fk }), + }) + + return { + ops: ops + .map((op, i) => [op, i] as const) + .sort(([a, i], [b, j]) => opOrder[a.op] - opOrder[b.op] || i - j) + .map(([op]) => op), + missingHints, + unsupported: [], + } +} + +const diffColumns = ( + before: ReadonlyArray, + after: ReadonlyArray, + table: string, + ops: Array, + confirm: (hint: Hint, op: PgMigrationOp) => void, +): void => { + const beforeByName = new Map(before.map((c) => [c.name, c])) + const afterNames = new Set(after.map((c) => c.name)) + for (const column of after) { + const old = beforeByName.get(column.name) + if (old === undefined) ops.push({ op: "add_column", column }) + // Position alone is not a change: Postgres cannot reorder columns, and the + // order only matters to a CREATE TABLE. + else if (old.type !== column.type || old.notNull !== column.notNull || old.default !== column.default || old.identity !== column.identity) { + ops.push({ op: "alter_column", from: old, to: column }) + } + } + for (const column of before) { + if (!afterNames.has(column.name)) { + confirm( + { type: "confirm_data_loss", kind: "column", entity: `${table}.${column.name}` }, + { op: "drop_column", table, name: column.name }, + ) + } + } +} + +const diffPrimaryKey = (before: PgTableEntity, after: PgTableEntity, ops: Array): void => { + if (!same(before.primaryKey, after.primaryKey)) { + ops.push({ op: "set_primary_key", table: after.name, from: before.primaryKey, to: after.primaryKey }) + } +} + +const diffNamed = ( + before: Map, + after: Map, + prevTables: Map, + ops: Array, + make: { readonly drop: (e: E) => PgMigrationOp; readonly create: (e: E) => PgMigrationOp }, +): void => { + const dropped = new Set(ops.flatMap((op) => (op.op === "drop_table" ? [op.name] : []))) + // A table that is gone takes its indexes and keys with it; dropping them first is still + // correct, and keeps a foreign key from another table from blocking the DROP TABLE. + for (const [key, entity] of after) { + const old = before.get(key) + if (old !== undefined && same(old, entity)) continue + if (old !== undefined) ops.push(make.drop(old)) + ops.push(make.create(entity)) + } + for (const [key, entity] of before) { + if (after.has(key)) continue + // Only an index of a dropped table can be left to the DROP TABLE; a foreign + // key is dropped either way, because one pointing *at* the table blocks it. + if (entity.kind === "index" && dropped.has(entity.table) && prevTables.has(`table:${entity.table}`)) continue + ops.push(make.drop(entity)) + } +} diff --git a/src/schema/pg-entities.ts b/src/schema/pg-entities.ts new file mode 100644 index 0000000..b1de457 --- /dev/null +++ b/src/schema/pg-entities.ts @@ -0,0 +1,128 @@ +// Postgres schema entities: the normalized, serializable form of a Postgres +// schema, as `entities.ts` is for ClickHouse. +// +// Types are stored in the spelling `format_type` reports (`integer`, +// `timestamp with time zone`, `real[]`), so a snapshot, a drizzle-kit snapshot +// and the live catalog all compare as text. Expressions (defaults, index keys, +// partial-index predicates) are SQL as written; `verify` normalizes both sides +// before comparing them. + +import { Schema } from "effect" + +export const PgTableEntity = Schema.Struct({ + kind: Schema.Literal("table"), + name: Schema.String, + /** The primary key constraint, or `null` for a table without one. */ + primaryKey: Schema.NullOr(Schema.Struct({ name: Schema.String, columns: Schema.Array(Schema.String) })), +}) +export type PgTableEntity = typeof PgTableEntity.Type + +/** An identity column's kind. */ +export const PgIdentity = Schema.Literals(["always", "by default"]) +export type PgIdentity = typeof PgIdentity.Type + +export const PgColumnEntity = Schema.Struct({ + kind: Schema.Literal("column"), + table: Schema.String, + name: Schema.String, + /** Declaration order, the order `CREATE TABLE` writes the columns in. */ + position: Schema.Number, + /** As `format_type` spells it: `integer`, `timestamp with time zone`, `text[]`. */ + type: Schema.String, + notNull: Schema.Boolean, + /** The `DEFAULT` expression as SQL, or `null`. */ + default: Schema.NullOr(Schema.String), + /** `GENERATED ALWAYS | BY DEFAULT AS IDENTITY`, with the sequence's default options; `null` for none. */ + identity: Schema.NullOr(PgIdentity), +}) +export type PgColumnEntity = typeof PgColumnEntity.Type + +export const PgIndexEntity = Schema.Struct({ + kind: Schema.Literal("index"), + table: Schema.String, + name: Schema.String, + unique: Schema.Boolean, + /** `btree`, `gin`, `hash`, ... */ + method: Schema.String, + /** Key parts as SQL: a quoted column (`"org_id"`) or an expression. */ + columns: Schema.Array(Schema.String), + /** A partial index's predicate, or `null`. */ + where: Schema.NullOr(Schema.String), +}) +export type PgIndexEntity = typeof PgIndexEntity.Type + +/** `ON DELETE` / `ON UPDATE` actions, as the catalog names them. */ +export const PgReferentialAction = Schema.Literals(["NO ACTION", "RESTRICT", "CASCADE", "SET NULL", "SET DEFAULT"]) +export type PgReferentialAction = typeof PgReferentialAction.Type + +export const PgForeignKeyEntity = Schema.Struct({ + kind: Schema.Literal("foreign_key"), + table: Schema.String, + name: Schema.String, + columns: Schema.Array(Schema.String), + foreignTable: Schema.String, + foreignColumns: Schema.Array(Schema.String), + onDelete: PgReferentialAction, + onUpdate: PgReferentialAction, +}) +export type PgForeignKeyEntity = typeof PgForeignKeyEntity.Type + +export const PgSchemaEntity = Schema.Union([PgTableEntity, PgColumnEntity, PgIndexEntity, PgForeignKeyEntity]) +export type PgSchemaEntity = typeof PgSchemaEntity.Type + +/** Postgres truncates longer identifiers silently, which would make every later comparison miss. */ +export const PG_MAX_IDENTIFIER = 63 + +const TYPE_ALIASES: Readonly> = { + int2: "smallint", + smallint: "smallint", + int: "integer", + int4: "integer", + integer: "integer", + int8: "bigint", + bigint: "bigint", + float4: "real", + real: "real", + float8: "double precision", + "double precision": "double precision", + bool: "boolean", + boolean: "boolean", + timestamptz: "timestamp with time zone", + "timestamp with time zone": "timestamp with time zone", + timestamp: "timestamp without time zone", + "timestamp without time zone": "timestamp without time zone", + timetz: "time with time zone", + "time with time zone": "time with time zone", + time: "time without time zone", + "time without time zone": "time without time zone", + varchar: "character varying", + "character varying": "character varying", + char: "character", + character: "character", + decimal: "numeric", + numeric: "numeric", +} + +/** + * A type name as `format_type` writes it: `int4` is `integer`, `timestamptz[]` + * is `timestamp with time zone[]`, `varchar(20)` is `character varying(20)`. + * Names it does not know (`text`, `jsonb`, `uuid`, a custom type) are kept. + */ +export const canonicalPgType = (type: string): string => { + const trimmed = type.trim().replace(/\s+/g, " ") + const array = /^(.*?)((?:\s*\[\s*\d*\s*\])+)$/.exec(trimmed) + if (array !== null) { + const dimensions = (array[2]!.match(/\[/g) ?? []).length + return `${canonicalPgType(array[1]!)}${"[]".repeat(dimensions)}` + } + const match = /^([a-z0-9_ ]+?)\s*(\(.*\))?$/i.exec(trimmed) + if (match === null) return trimmed + const base = match[1]!.toLowerCase() + const args = match[2]?.replace(/\s+/g, "") ?? "" + const alias = TYPE_ALIASES[base] + if (alias === undefined) return `${base}${args}` + // `timestamp(3) with time zone`: format_type puts the precision after `timestamp`. + if (args.length > 0 && alias.startsWith("timestamp ")) return alias.replace("timestamp", `timestamp${args}`) + if (args.length > 0 && alias.startsWith("time ")) return alias.replace("time", `time${args}`) + return `${alias}${args}` +} diff --git a/src/schema/pg-ops.ts b/src/schema/pg-ops.ts new file mode 100644 index 0000000..75071c0 --- /dev/null +++ b/src/schema/pg-ops.ts @@ -0,0 +1,150 @@ +// Postgres migration operations and the DDL they render to. +// +// The Postgres counterpart of `ops.ts`. A Postgres migration runs inside one +// transaction, so its statements do not have to be safe to repeat; they still +// use `IF EXISTS` / `IF NOT EXISTS` where Postgres has it, so a migration +// replayed against a hand-fixed database does not trip over its own objects. + +import { Schema } from "effect" +import { PgColumnEntity, PgForeignKeyEntity, PgIndexEntity, PgTableEntity, type PgSchemaEntity } from "./pg-entities" + +const PrimaryKey = Schema.NullOr(Schema.Struct({ name: Schema.String, columns: Schema.Array(Schema.String) })) + +export const PgMigrationOp = Schema.Union([ + Schema.Struct({ op: Schema.Literal("create_table"), table: PgTableEntity, columns: Schema.Array(PgColumnEntity) }), + Schema.Struct({ op: Schema.Literal("drop_table"), name: Schema.String }), + Schema.Struct({ op: Schema.Literal("add_column"), column: PgColumnEntity }), + Schema.Struct({ op: Schema.Literal("drop_column"), table: Schema.String, name: Schema.String }), + /** Type, nullability, or default changed. */ + Schema.Struct({ op: Schema.Literal("alter_column"), from: PgColumnEntity, to: PgColumnEntity }), + Schema.Struct({ op: Schema.Literal("set_primary_key"), table: Schema.String, from: PrimaryKey, to: PrimaryKey }), + Schema.Struct({ op: Schema.Literal("create_index"), index: PgIndexEntity }), + Schema.Struct({ op: Schema.Literal("drop_index"), table: Schema.String, name: Schema.String }), + Schema.Struct({ op: Schema.Literal("add_foreign_key"), foreignKey: PgForeignKeyEntity }), + Schema.Struct({ op: Schema.Literal("drop_foreign_key"), table: Schema.String, name: Schema.String }), +]) +export type PgMigrationOp = typeof PgMigrationOp.Type + +/** A generated Postgres migration. `dialect` tells it apart from a ClickHouse `migration.json`. */ +export const PgMigrationFile = Schema.Struct({ + version: Schema.Literal("1"), + dialect: Schema.Literal("postgres"), + ops: Schema.Array(PgMigrationOp), +}) +export type PgMigrationFile = typeof PgMigrationFile.Type + +/** An identifier, always double-quoted, so mixed case and keywords survive. */ +export const pgIdent = (name: string): string => `"${name.replace(/"/g, '""')}"` + +const list = (names: ReadonlyArray): string => names.map(pgIdent).join(", ") + +export const renderPgColumnDefinition = (column: PgColumnEntity): string => { + const parts = [pgIdent(column.name), column.type] + if (column.notNull) parts.push("NOT NULL") + if (column.default !== null) parts.push(`DEFAULT ${column.default}`) + if (column.identity !== null) parts.push(`GENERATED ${column.identity.toUpperCase()} AS IDENTITY`) + return parts.join(" ") +} + +export const renderPgCreateTable = (table: PgTableEntity, columns: ReadonlyArray): string => { + const body = [...columns].sort((a, b) => a.position - b.position).map(renderPgColumnDefinition) + if (table.primaryKey !== null) { + body.push(`CONSTRAINT ${pgIdent(table.primaryKey.name)} PRIMARY KEY (${list(table.primaryKey.columns)})`) + } + return `CREATE TABLE IF NOT EXISTS ${pgIdent(table.name)} (\n\t${body.join(",\n\t")}\n)` +} + +export const renderPgCreateIndex = (index: PgIndexEntity): string => + `CREATE ${index.unique ? "UNIQUE " : ""}INDEX IF NOT EXISTS ${pgIdent(index.name)} ON ${pgIdent(index.table)} USING ${index.method} (${index.columns.join(", ")})${index.where === null ? "" : ` WHERE ${index.where}`}` + +export const renderPgAddForeignKey = (fk: PgForeignKeyEntity): string => { + const actions = [ + ...(fk.onDelete === "NO ACTION" ? [] : [`ON DELETE ${fk.onDelete}`]), + ...(fk.onUpdate === "NO ACTION" ? [] : [`ON UPDATE ${fk.onUpdate}`]), + ] + return `ALTER TABLE ${pgIdent(fk.table)} ADD CONSTRAINT ${pgIdent(fk.name)} FOREIGN KEY (${list(fk.columns)}) REFERENCES ${pgIdent(fk.foreignTable)} (${list(fk.foreignColumns)})${actions.map((a) => ` ${a}`).join("")}` +} + +const alter = (table: string, action: string): string => `ALTER TABLE ${pgIdent(table)} ${action}` + +/** + * Every CREATE statement for a schema: tables, then indexes, then foreign + * keys, so every table a key references exists before the key. + */ +export const renderPgSchema = (entities: ReadonlyArray): ReadonlyArray => { + const pick = (kind: K) => + entities.filter((e): e is Extract => e.kind === kind) + const columns = pick("column") + return [ + ...pick("table").map((table) => renderPgCreateTable(table, columns.filter((c) => c.table === table.name))), + ...pick("index").map(renderPgCreateIndex), + ...pick("foreign_key").map(renderPgAddForeignKey), + ] +} + +/** Labels the plan prints. `rewrite`: a column type change, which rewrites the table under a lock. */ +export type PgOpLabel = "metadata" | "destructive" | "rewrite" + +export const labelOfPg = (op: PgMigrationOp): PgOpLabel => { + switch (op.op) { + case "drop_table": + case "drop_column": + return "destructive" + case "alter_column": + return op.from.type === op.to.type ? "metadata" : "rewrite" + default: + return "metadata" + } +} + +/** The statements one op runs, in order. */ +export const renderPgOp = (op: PgMigrationOp): ReadonlyArray => { + switch (op.op) { + case "create_table": + return [renderPgCreateTable(op.table, op.columns)] + case "drop_table": + return [`DROP TABLE IF EXISTS ${pgIdent(op.name)}`] + case "add_column": + return [alter(op.column.table, `ADD COLUMN IF NOT EXISTS ${renderPgColumnDefinition(op.column)}`)] + case "drop_column": + return [alter(op.table, `DROP COLUMN IF EXISTS ${pgIdent(op.name)}`)] + case "alter_column": { + const { from, to } = op + const name = pgIdent(to.name) + const out: Array = [] + const typeChanged = from.type !== to.type + // The old default may not cast to the new type, so it goes before the type changes. + if (from.default !== null && (to.default === null || (typeChanged && from.default !== to.default))) { + out.push(alter(to.table, `ALTER COLUMN ${name} DROP DEFAULT`)) + } + if (typeChanged) out.push(alter(to.table, `ALTER COLUMN ${name} SET DATA TYPE ${to.type} USING ${name}::${to.type}`)) + if (to.default !== null && (from.default !== to.default || typeChanged)) { + out.push(alter(to.table, `ALTER COLUMN ${name} SET DEFAULT ${to.default}`)) + } + if (from.notNull !== to.notNull) out.push(alter(to.table, `ALTER COLUMN ${name} ${to.notNull ? "SET" : "DROP"} NOT NULL`)) + if (from.identity !== to.identity) { + out.push( + from.identity === null + ? alter(to.table, `ALTER COLUMN ${name} ADD GENERATED ${to.identity!.toUpperCase()} AS IDENTITY`) + : to.identity === null + ? alter(to.table, `ALTER COLUMN ${name} DROP IDENTITY IF EXISTS`) + : alter(to.table, `ALTER COLUMN ${name} SET GENERATED ${to.identity.toUpperCase()}`), + ) + } + return out + } + case "set_primary_key": + return [ + ...(op.from === null ? [] : [alter(op.table, `DROP CONSTRAINT IF EXISTS ${pgIdent(op.from.name)}`)]), + ...(op.to === null ? [] : [alter(op.table, `ADD CONSTRAINT ${pgIdent(op.to.name)} PRIMARY KEY (${list(op.to.columns)})`)]), + ] + case "create_index": + return [renderPgCreateIndex(op.index)] + case "drop_index": + return [`DROP INDEX IF EXISTS ${pgIdent(op.name)}`] + case "add_foreign_key": + return [renderPgAddForeignKey(op.foreignKey)] + case "drop_foreign_key": + return [alter(op.table, `DROP CONSTRAINT IF EXISTS ${pgIdent(op.name)}`)] + } +} diff --git a/src/schema/pg-schema.test.ts b/src/schema/pg-schema.test.ts new file mode 100644 index 0000000..fc47c1e --- /dev/null +++ b/src/schema/pg-schema.test.ts @@ -0,0 +1,277 @@ +import { Effect } from "effect" +import { describe, expect, it } from "vitest" +import * as CH from "../ch/index" +import * as PG from "../postgres" +import * as S from "../schema" + +const Dashboards = S.pg.table("dashboards", { + columns: { + org_id: PG.text, + id: PG.text, + name: S.pg.column(PG.text, { default: "Untitled" }), + tags: S.pg.column(PG.array(PG.text), { default: [] }), + layout: S.pg.column(PG.jsonb(), { default: {} }), + widgets: PG.int4, + archived: S.pg.column(PG.bool, { default: false }), + created_at: S.pg.column(PG.timestamptz, { defaultExpr: "now()" }), + archived_at: PG.nullable(PG.timestamptz), + }, + primaryKey: { columns: ["org_id", "id"], name: "dashboards_org_id_id_pk" }, + indexes: [ + S.pg.index("dashboards_org_idx", ["org_id"]), + S.pg.index("dashboards_live_idx", ["org_id", "created_at"], { where: ($) => $.archived_at.isNull() }), + ], + tenantColumn: "org_id", +}) + +const Shares = S.pg.table("dashboard_shares", { + columns: { + org_id: PG.text, + id: PG.text, + dashboard_id: PG.text, + widget_id: PG.nullable(PG.text), + embedding: PG.nullable(PG.array(PG.float4)), + revoked_at: PG.nullable(PG.timestamptz), + }, + primaryKey: ["org_id", "id"], + indexes: [ + S.pg.uniqueIndex("dashboard_shares_live_unq", ($) => [$.org_id, $.dashboard_id, CH.coalesce($.widget_id, CH.lit(""))], { + where: "revoked_at is null", + }), + ], + foreignKeys: [ + S.pg.foreignKey({ + columns: ["org_id", "dashboard_id"], + references: Dashboards, + foreignColumns: ["org_id", "id"], + onDelete: "cascade", + name: "dashboard_shares_dashboard_fk", + }), + ], +}) + +describe("S.pg.table", () => { + it("is a Table the query builder accepts, with defaults optional on insert", () => { + const { sql } = PG.compileUnsafe(CH.from(Dashboards).select("name").where(($) => [$.org_id.eq("o")]), {}) + expect(sql).toContain('FROM "dashboards"') + const insert = PG.compileUnsafe(CH.insertInto(Dashboards).values({ org_id: "o", id: "d", widgets: 0 }), {}) + expect(insert.sql).toContain('INSERT INTO "dashboards"') + expect(Dashboards.defaults).toEqual(["name", "tags", "layout", "archived", "created_at"]) + }) + + it("normalizes types and nullability into entities", () => { + const columns = Object.fromEntries(Shares.ddl.columns.map((c) => [c.name, [c.type, c.notNull]])) + expect(columns).toEqual({ + org_id: ["text", true], + id: ["text", true], + dashboard_id: ["text", true], + widget_id: ["text", false], + embedding: ["real[]", false], + revoked_at: ["timestamp with time zone", false], + }) + expect(Shares.ddl.table.primaryKey).toEqual({ name: "dashboard_shares_pkey", columns: ["org_id", "id"] }) + }) + + it("renders its DDL", () => { + expect(S.renderPgSchema(S.pgEntitiesOf([Dashboards, Shares]))).toEqual([ + [ + 'CREATE TABLE IF NOT EXISTS "dashboard_shares" (', + '\t"org_id" text NOT NULL,', + '\t"id" text NOT NULL,', + '\t"dashboard_id" text NOT NULL,', + '\t"widget_id" text,', + '\t"embedding" real[],', + '\t"revoked_at" timestamp with time zone,', + '\tCONSTRAINT "dashboard_shares_pkey" PRIMARY KEY ("org_id", "id")', + ")", + ].join("\n"), + [ + 'CREATE TABLE IF NOT EXISTS "dashboards" (', + '\t"org_id" text NOT NULL,', + '\t"id" text NOT NULL,', + "\t\"name\" text NOT NULL DEFAULT 'Untitled',", + "\t\"tags\" text[] NOT NULL DEFAULT '{}',", + "\t\"layout\" jsonb NOT NULL DEFAULT '{}',", + '\t"widgets" integer NOT NULL,', + '\t"archived" boolean NOT NULL DEFAULT false,', + '\t"created_at" timestamp with time zone NOT NULL DEFAULT now(),', + '\t"archived_at" timestamp with time zone,', + '\tCONSTRAINT "dashboards_org_id_id_pk" PRIMARY KEY ("org_id", "id")', + ")", + ].join("\n"), + `CREATE UNIQUE INDEX IF NOT EXISTS "dashboard_shares_live_unq" ON "dashboard_shares" USING btree ("org_id", "dashboard_id", coalesce("widget_id", '')) WHERE revoked_at is null`, + 'CREATE INDEX IF NOT EXISTS "dashboards_live_idx" ON "dashboards" USING btree ("org_id", "created_at") WHERE "archived_at" IS NULL', + 'CREATE INDEX IF NOT EXISTS "dashboards_org_idx" ON "dashboards" USING btree ("org_id")', + 'ALTER TABLE "dashboard_shares" ADD CONSTRAINT "dashboard_shares_dashboard_fk" FOREIGN KEY ("org_id", "dashboard_id") REFERENCES "dashboards" ("org_id", "id") ON DELETE CASCADE', + ]) + }) + + it("names a foreign key as drizzle-kit does when no name is given", () => { + const Checks = S.pg.table("checks", { + columns: { id: PG.text, target_id: PG.text }, + foreignKeys: [S.pg.foreignKey({ columns: ["target_id"], references: "targets", foreignColumns: ["id"] })], + }) + expect(Checks.ddl.foreignKeys[0]?.name).toBe("checks_target_id_targets_id_fk") + }) + + it("rejects definitions Postgres would not take as written", () => { + expect(() => S.pg.table("t", { columns: { a: PG.nullable(PG.text) }, primaryKey: ["a"] })).toThrow(/cannot be nullable/) + expect(() => S.pg.table("t".repeat(64), { columns: { a: PG.text } })).toThrow(/longer than 63/) + expect(() => + S.pg.table("t", { columns: { a: PG.text }, foreignKeys: [S.pg.foreignKey({ columns: ["a"], references: "u", foreignColumns: ["x", "y"] })] }), + ).toThrow(/same, non-zero, length/) + }) + + it("validates a schema as a whole", () => { + const Other = S.pg.table("other", { + columns: { a: PG.text }, + indexes: [S.pg.index("dashboards_org_idx", ["a"])], + }) + expect(() => S.pgEntitiesOf([Dashboards, Other])).toThrow(/one namespace/) + expect(() => S.pgEntitiesOf([Shares])).toThrow(/not a table in this schema/) + expect(() => S.entitiesOf([Shares])).toThrow(/Postgres table/) + }) +}) + +describe("canonicalPgType", () => { + it("spells types the way format_type does", () => { + expect(S.canonicalPgType("int4")).toBe("integer") + expect(S.canonicalPgType("float8")).toBe("double precision") + expect(S.canonicalPgType("timestamptz[]")).toBe("timestamp with time zone[]") + expect(S.canonicalPgType("timestamptz(3)")).toBe("timestamp(3) with time zone") + expect(S.canonicalPgType("varchar(20)")).toBe("character varying(20)") + expect(S.canonicalPgType("jsonb")).toBe("jsonb") + }) +}) + +describe("diffPgSchemas", () => { + const v1 = S.pgEntitiesOf([Dashboards]) + const v2 = S.pgEntitiesOf([ + S.pg.table("dashboards", { + columns: { + org_id: PG.text, + id: PG.text, + name: S.pg.column(PG.text, { default: "New dashboard" }), + widgets: PG.int8, + archived: PG.bool, + created_at: S.pg.column(PG.timestamptz, { defaultExpr: "now()" }), + archived_at: PG.nullable(PG.timestamptz), + owner: PG.nullable(PG.text), + }, + primaryKey: ["org_id", "id"], + indexes: [S.pg.index("dashboards_org_idx", ["org_id", "owner"])], + }), + ]) + + it("alters columns, keys and indexes in place, and asks before dropping", () => { + const first = S.diffPgSchemas(v1, v2) + expect(first.missingHints).toEqual([ + { type: "confirm_data_loss", kind: "column", entity: "dashboards.tags" }, + { type: "confirm_data_loss", kind: "column", entity: "dashboards.layout" }, + ]) + const { ops, unsupported } = S.diffPgSchemas(v1, v2, first.missingHints) + expect(unsupported).toEqual([]) + expect(ops.map((op) => op.op)).toEqual([ + "drop_index", + "drop_index", + "drop_column", + "drop_column", + "add_column", + "alter_column", + "alter_column", + "alter_column", + "set_primary_key", + "create_index", + ]) + expect(ops.flatMap((op) => S.renderPgOp(op))).toEqual([ + 'DROP INDEX IF EXISTS "dashboards_org_idx"', + 'DROP INDEX IF EXISTS "dashboards_live_idx"', + 'ALTER TABLE "dashboards" DROP COLUMN IF EXISTS "tags"', + 'ALTER TABLE "dashboards" DROP COLUMN IF EXISTS "layout"', + 'ALTER TABLE "dashboards" ADD COLUMN IF NOT EXISTS "owner" text', + `ALTER TABLE "dashboards" ALTER COLUMN "name" SET DEFAULT 'New dashboard'`, + 'ALTER TABLE "dashboards" ALTER COLUMN "widgets" SET DATA TYPE bigint USING "widgets"::bigint', + 'ALTER TABLE "dashboards" ALTER COLUMN "archived" DROP DEFAULT', + 'ALTER TABLE "dashboards" DROP CONSTRAINT IF EXISTS "dashboards_org_id_id_pk"', + 'ALTER TABLE "dashboards" ADD CONSTRAINT "dashboards_pkey" PRIMARY KEY ("org_id", "id")', + 'CREATE INDEX IF NOT EXISTS "dashboards_org_idx" ON "dashboards" USING btree ("org_id", "owner")', + ]) + expect(ops.map(S.labelOfPg)).toContain("rewrite") + }) + + it("drops foreign keys before the tables they point at, and nothing for an unchanged schema", () => { + const both = S.pgEntitiesOf([Dashboards, Shares]) + expect(S.diffPgSchemas(both, both).ops).toEqual([]) + const hint = { type: "confirm_data_loss", kind: "table", entity: "dashboards" } as const + const hint2 = { type: "confirm_data_loss", kind: "table", entity: "dashboard_shares" } as const + const { ops } = S.diffPgSchemas(both, [], [hint, hint2]) + expect(ops.map((op) => op.op)).toEqual(["drop_foreign_key", "drop_table", "drop_table"]) + }) + + it("snapshots carry their dialect and hash entities alone", async () => { + const snapshot = await Effect.runPromise(S.makeSnapshot(v1, [S.ORIGIN_ID], "postgres")) + expect(snapshot.dialect).toBe("postgres") + const again = await Effect.runPromise(S.makeSnapshot([...v1].reverse(), [S.ORIGIN_ID], "postgres")) + expect(again.id).toBe(snapshot.id) + }) +}) + +describe("fromDrizzleSnapshot", () => { + const drizzle = { + version: "8", + dialect: "postgres", + id: "x", + prevIds: [], + renames: [], + ddl: [ + { isRlsEnabled: false, name: "dashboard_shares", entityType: "tables", schema: "public" }, + { type: "text", typeSchema: null, notNull: true, dimensions: 0, default: null, generated: null, identity: null, name: "org_id", entityType: "columns", schema: "public", table: "dashboard_shares" }, + { type: "real", typeSchema: null, notNull: false, dimensions: 1, default: null, generated: null, identity: null, name: "embedding", entityType: "columns", schema: "public", table: "dashboard_shares" }, + { type: "text", typeSchema: null, notNull: true, dimensions: 0, default: "'open'", generated: null, identity: null, name: "status", entityType: "columns", schema: "public", table: "dashboard_shares" }, + { + nameExplicit: true, + columns: [ + { value: "org_id", isExpression: false, asc: true, nullsFirst: false, opclass: null }, + { value: "coalesce(widget_id, '')", isExpression: true, asc: true, nullsFirst: false, opclass: null }, + ], + isUnique: true, + where: "revoked_at is null", + with: "", + method: "btree", + concurrently: false, + name: "dashboard_shares_live_unq", + entityType: "indexes", + schema: "public", + table: "dashboard_shares", + }, + { nameExplicit: true, columns: ["org_id"], schemaTo: "public", tableTo: "dashboards", columnsTo: ["org_id"], onUpdate: "NO ACTION", onDelete: "CASCADE", name: "fk", entityType: "fks", schema: "public", table: "dashboard_shares" }, + { columns: ["org_id"], nameExplicit: false, name: "dashboard_shares_org_id_pk", entityType: "pks", schema: "public", table: "dashboard_shares" }, + ], + } + + it("converts tables, columns, keys and indexes", () => { + const { entities, unsupported } = S.fromDrizzleSnapshot(drizzle) + expect(unsupported).toEqual([]) + expect(entities).toContainEqual({ kind: "table", name: "dashboard_shares", primaryKey: { name: "dashboard_shares_org_id_pk", columns: ["org_id"] } }) + expect(entities).toContainEqual({ kind: "column", table: "dashboard_shares", name: "embedding", position: 1, type: "real[]", notNull: false, default: null, identity: null }) + expect(entities).toContainEqual({ kind: "column", table: "dashboard_shares", name: "status", position: 2, type: "text", notNull: true, default: "'open'", identity: null }) + expect(entities).toContainEqual({ + kind: "index", + table: "dashboard_shares", + name: "dashboard_shares_live_unq", + unique: true, + method: "btree", + columns: ['"org_id"', "coalesce(widget_id, '')"], + where: "revoked_at is null", + }) + expect(entities.find((e) => e.kind === "foreign_key")).toMatchObject({ foreignTable: "dashboards", onDelete: "CASCADE" }) + }) + + it("reports what it cannot model instead of dropping it", () => { + const { unsupported } = S.fromDrizzleSnapshot({ + ...drizzle, + ddl: [...drizzle.ddl, { name: "mood", values: ["a"], entityType: "enums", schema: "public" }], + }) + expect(unsupported).toEqual(["enums mood"]) + }) +}) diff --git a/src/schema/snapshot.ts b/src/schema/snapshot.ts index 56209c6..9504514 100644 --- a/src/schema/snapshot.ts +++ b/src/schema/snapshot.ts @@ -9,12 +9,18 @@ import { sha256Hex, sortEntities, SNAPSHOT_VERSION, + type AnySchemaEntity, + type PgSnapshot, + type ClickHouseSnapshot, + type SchemaDialect, type SchemaEntity, type Snapshot, } from "./entities" +import type { PgSchemaTable } from "./pg-define" +import type { PgSchemaEntity } from "./pg-entities" /** Anything a schema module may export that `generate` collects. */ -export type SchemaObject = SchemaTable | MaterializedView +export type SchemaObject = SchemaTable | MaterializedView | PgSchemaTable export const isSchemaObject = (value: unknown): value is SchemaObject => typeof value === "object" && @@ -23,12 +29,21 @@ export const isSchemaObject = (value: unknown): value is SchemaObject => "_tag" in value && (value._tag === "Table" || value._tag === "MaterializedView") +/** The dialect a definition was written for. */ +export const dialectOfObject = (object: SchemaObject): SchemaDialect => + object._tag === "Table" && "dialect" in object.ddl ? object.ddl.dialect : "clickhouse" + /** The entities of a set of definitions, validated as one schema. */ export const entitiesOf = (objects: ReadonlyArray): ReadonlyArray => { const entities: Array = [] for (const object of objects) { - if (object._tag === "Table") entities.push(object.ddl.table, ...object.ddl.columns, ...object.ddl.indexes) - else entities.push(object.ddl) + if (dialectOfObject(object) !== "clickhouse") { + throw new SchemaDefinitionDefect({ object: object.name, message: "is a Postgres table; use pgEntitiesOf" }) + } + if (object._tag === "Table") { + const ddl = object.ddl as SchemaTable["ddl"] + entities.push(ddl.table, ...ddl.columns, ...ddl.indexes) + } else entities.push(object.ddl) } const seen = new Set() const names = new Map() @@ -62,22 +77,85 @@ export const entitiesOf = (objects: ReadonlyArray): ReadonlyArray< return sortEntities(entities) } +/** + * The entities of a set of Postgres tables, validated as one schema: names are + * unique, index and primary-key names do not collide (Postgres keeps them in + * one namespace per schema), and every foreign key references a table and + * columns of this schema. + */ +export const pgEntitiesOf = (objects: ReadonlyArray): ReadonlyArray => { + const tables: Array> = [] + for (const object of objects) { + if (dialectOfObject(object) !== "postgres") { + throw new SchemaDefinitionDefect({ object: object.name, message: "is a ClickHouse definition; use entitiesOf" }) + } + tables.push(object as PgSchemaTable) + } + const entities: Array = [] + const relations = new Map() + const claim = (name: string, what: string) => { + const other = relations.get(name) + if (other !== undefined) { + throw new SchemaDefinitionDefect({ object: name, message: `names both ${other} and ${what}; Postgres keeps them in one namespace` }) + } + relations.set(name, what) + } + for (const { ddl } of tables) { + claim(ddl.table.name, `table ${ddl.table.name}`) + if (ddl.table.primaryKey !== null) claim(ddl.table.primaryKey.name, `the primary key of ${ddl.table.name}`) + for (const index of ddl.indexes) claim(index.name, `an index on ${ddl.table.name}`) + entities.push(ddl.table, ...ddl.columns, ...ddl.indexes, ...ddl.foreignKeys) + } + const columnsOf = new Map(tables.map(({ ddl }) => [ddl.table.name, new Set(ddl.columns.map((c) => c.name))])) + const seenConstraints = new Set() + for (const { ddl } of tables) { + for (const fk of ddl.foreignKeys) { + const key = `${fk.table}.${fk.name}` + if (seenConstraints.has(key)) throw new SchemaDefinitionDefect({ object: key, message: "defined twice" }) + seenConstraints.add(key) + const target = columnsOf.get(fk.foreignTable) + if (target === undefined) { + throw new SchemaDefinitionDefect({ object: key, message: `references ${fk.foreignTable}, which is not a table in this schema` }) + } + const missing = fk.foreignColumns.filter((c) => !target.has(c)) + if (missing.length > 0) { + throw new SchemaDefinitionDefect({ object: key, message: `references ${missing.join(", ")}, not columns of ${fk.foreignTable}` }) + } + } + } + return sortEntities(entities) +} + /** A snapshot of `entities` with the given parents. Its id is the hash of the entities alone. */ -export const makeSnapshot = ( - entities: ReadonlyArray, +export function makeSnapshot(entities: ReadonlyArray, prevIds: ReadonlyArray): Effect.Effect +export function makeSnapshot( + entities: ReadonlyArray, + prevIds: ReadonlyArray, + dialect: "postgres", +): Effect.Effect +export function makeSnapshot( + entities: ReadonlyArray, prevIds: ReadonlyArray, -): Effect.Effect => - Effect.promise(() => sha256Hex(canonicalJson(sortEntities(entities)))).pipe( + dialect?: SchemaDialect, +): Effect.Effect +export function makeSnapshot( + entities: ReadonlyArray, + prevIds: ReadonlyArray, + dialect: SchemaDialect = "clickhouse", +): Effect.Effect { + return Effect.promise(() => sha256Hex(canonicalJson(sortEntities(entities)))).pipe( Effect.map( - (id): Snapshot => ({ - version: SNAPSHOT_VERSION, - dialect: "clickhouse", - id, - prevIds: [...prevIds], - entities: sortEntities(entities), - }), + (id) => + ({ + version: SNAPSHOT_VERSION, + dialect, + id, + prevIds: [...prevIds], + entities: sortEntities(entities), + }) as Snapshot, ), ) +} /** Pretty, stable snapshot JSON for committing. */ export const serializeSnapshot = (snapshot: Snapshot): string => `${JSON.stringify(snapshot, null, "\t")}\n` From 7bad019cba65c874778b1725bd290d7af0129d74 Mon Sep 17 00:00:00 2001 From: Makisuo Date: Sun, 4 Oct 2026 23:37:21 +0200 Subject: [PATCH 2/2] Shorten default foreign key names past 63 characters like drizzle-kit A long composite key's default name passed Postgres's identifier limit and failed at module load. Hash it to
__fk with drizzle-kit's hash, so names stay valid and deterministic. Co-Authored-By: Claude Opus 5.5 --- docs/migrations.md | 3 ++- src/schema/pg-define.ts | 39 +++++++++++++++++++++++++++++++++++- src/schema/pg-schema.test.ts | 15 ++++++++++++++ 3 files changed, 55 insertions(+), 2 deletions(-) diff --git a/docs/migrations.md b/docs/migrations.md index d4e4d09..83bb04a 100644 --- a/docs/migrations.md +++ b/docs/migrations.md @@ -226,7 +226,8 @@ expression) or an `identity` (`"always"` or `"by default"`); any of them makes t optional on insert. `primaryKey` takes column names, or `{ columns, name }`; the default name is `
_pkey`. Indexes are `S.pg.index` / `S.pg.uniqueIndex` over column names or expressions, with `where` for a partial index and `using` for the access method. A foreign key without a -`name` gets drizzle-kit's, `
____fk`. Types are +`name` gets drizzle-orm's, `
____fk`, shortened +with drizzle-kit's hash to `
__fk` when it would pass 63 characters. Types are stored as Postgres names them (`int4` is `integer`), so snapshots compare with the catalog and with drizzle-kit. Check and unique constraints, enums, views, sequences and other schemas are not modeled yet; write them in a `--custom` migration. diff --git a/src/schema/pg-define.ts b/src/schema/pg-define.ts index 23cd392..39fab54 100644 --- a/src/schema/pg-define.ts +++ b/src/schema/pg-define.ts @@ -208,6 +208,43 @@ const assertIdentifier = (object: string, name: string): void => { } } +const HASH_DICTIONARY = "0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ" + +/** drizzle-kit's name hash (`dialects/utils.ts`), so a shortened name matches the one it writes. */ +const drizzleHash = (input: string, length = 12): string => { + const base = BigInt(HASH_DICTIONARY.length) + const modulus = base ** BigInt(length) + let power = 1n + let hash = 0n + for (const ch of input) { + hash = (hash + BigInt(ch.codePointAt(0) ?? 0) * power) % modulus + power = (power * 53n) % modulus + } + const out: Array = [] + for (let i = 0; i < length; i++) { + out.unshift(HASH_DICTIONARY[Number(hash % base)]!) + hash /= base + } + return out.join("") +} + +/** + * `
____fk`, drizzle-orm's name. + * Over 63 characters, Postgres would truncate it, so it is shortened the way + * drizzle-kit shortens its own: `
__fk`, or `_fk` for a + * table name of 45 characters or more. + */ +export const defaultForeignKeyName = ( + table: string, + columns: ReadonlyArray, + foreignTable: string, + foreignColumns: ReadonlyArray, +): string => { + const desired = `${table}_${columns.join("_")}_${foreignTable}_${foreignColumns.join("_")}_fk` + if (desired.length <= PG_MAX_IDENTIFIER) return desired + return table.length < 45 ? `${table}_${drizzleHash(desired)}_fk` : `${drizzleHash(desired)}_fk` +} + /** `TRUE` is how the dialect writes a literal; the catalog and drizzle-kit write `true`. */ const lowerKeywords = (sql: string): string => (/^(TRUE|FALSE|NULL)$/.test(sql) ? sql.toLowerCase() : sql) @@ -309,7 +346,7 @@ export function table { - const fkName = spec.name ?? `${name}_${spec.columns.join("_")}_${spec.references}_${spec.foreignColumns.join("_")}_fk` + const fkName = spec.name ?? defaultForeignKeyName(name, spec.columns, spec.references, spec.foreignColumns) assertIdentifier(`${name} foreign key`, fkName) for (const column of spec.columns) { if (!has(column)) throw new SchemaDefinitionDefect({ object: `${name} foreign key ${fkName}`, message: `${column} is not a column` }) diff --git a/src/schema/pg-schema.test.ts b/src/schema/pg-schema.test.ts index fc47c1e..e742b0f 100644 --- a/src/schema/pg-schema.test.ts +++ b/src/schema/pg-schema.test.ts @@ -114,6 +114,21 @@ describe("S.pg.table", () => { expect(Checks.ddl.foreignKeys[0]?.name).toBe("checks_target_id_targets_id_fk") }) + it("shortens a default foreign key name past 63 characters with drizzle-kit's hash", () => { + const long = S.pg.table("organization_membership_invitations", { + columns: { organization_id: PG.text, invited_by_user_id: PG.text }, + foreignKeys: [ + S.pg.foreignKey({ columns: ["organization_id", "invited_by_user_id"], references: "organization_members", foreignColumns: ["organization_id", "user_id"] }), + ], + }) + const name = long.ddl.foreignKeys[0]!.name + expect(name).toMatch(/^organization_membership_invitations_[0-9A-Za-z]{12}_fk$/) + expect(name.length).toBeLessThanOrEqual(63) + // Deterministic, so a snapshot and the next generate agree. + expect(S.pg.defaultForeignKeyName("organization_membership_invitations", ["organization_id", "invited_by_user_id"], "organization_members", ["organization_id", "user_id"])).toBe(name) + expect(S.pg.defaultForeignKeyName("t".repeat(60), ["a"], "u", ["b"])).toMatch(/^[0-9A-Za-z]{12}_fk$/) + }) + it("rejects definitions Postgres would not take as written", () => { expect(() => S.pg.table("t", { columns: { a: PG.nullable(PG.text) }, primaryKey: ["a"] })).toThrow(/cannot be nullable/) expect(() => S.pg.table("t".repeat(64), { columns: { a: PG.text } })).toThrow(/longer than 63/)