From 2e24ac3f2ff2b94aa9f251796e2864f8d7a5653c Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 24 Sep 2026 23:48:02 -0700 Subject: [PATCH 1/2] feat(migrations)!: ship SQL migrations applied like runMigrations runArtifactMigrations now takes the same (config, { schema }) arguments as @intx/db's runMigrations and applies the idempotent SQL files under migrations/, rewriting the "public". foreign-key references to the host schema that holds tenant and principal. The package's own tables stay in the artifacts schema. One advisory-locked transaction still serializes concurrent boots. The embedded TypeScript DDL, the checksum ledger, the adopt path and their errors are gone. 0001 creates the final 0.1.0 shape, so a database 0.1.0 migrated no-ops, and 0002 drops the old ledger table. --- CHANGELOG.md | 7 + CONTRIBUTING.md | 92 +- README.md | 8 +- bun.lock | 2 + examples/reference-host/src/index.ts | 26 +- .../reference-host/test/acceptance.test.ts | 11 +- migrations/0001_artifacts.sql | 111 +++ migrations/0002_drop_migration_ledger.sql | 1 + package.json | 7 +- src/index.ts | 3 +- src/migrations.test.ts | 876 +++--------------- src/migrations.ts | 636 ++----------- src/test-helpers.ts | 15 +- 13 files changed, 382 insertions(+), 1413 deletions(-) create mode 100644 migrations/0001_artifacts.sql create mode 100644 migrations/0002_drop_migration_ledger.sql diff --git a/CHANGELOG.md b/CHANGELOG.md index 2bac9d0..ee2661a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -124,6 +124,13 @@ always called out under their own heading. ### Breaking +- `runArtifactMigrations(config, { schema })` takes the same arguments as + Interchange's `runMigrations`: a `DBConfig` and the host schema holding + `tenant` and `principal`. It applies the SQL files shipped under + `migrations/`, all idempotent, with no ledger. The `adopt` option, + `RunArtifactMigrationsOptions`, `MigrationChecksumError` and + `MigrationAdoptError` are removed, and the `artifacts.migrations` ledger + table is dropped on the next boot. - The drizzle tables (`artifact`, `artifactVersion`, `upload`, `mailAttachmentRef`) are no longer exported from the package entry. Hosts reach artifacts through the routes and functions; `ARTIFACTS_SCHEMA` and the diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 584fbf7..557e745 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -70,12 +70,12 @@ make this package uninstallable outside the project that defines it. ## Migrations -Shipped migrations are immutable. Each ledger row records a checksum of the migration's -rendered SQL, so editing one that has already been applied fails with -`MigrationChecksumError` on the next boot rather than letting fresh and existing -databases diverge. Add a new migration instead. +Migrations are SQL files under `migrations/`, applied in filename order on every +boot, so every statement must be idempotent (`IF NOT EXISTS`, `IF EXISTS`). There +is no ledger: a schema change is a new file whose statements are safe to re-run, +never an edit that assumes it runs once. -`schema.ts` and `migrations.ts` must agree — every query goes through the drizzle +`schema.ts` and `migrations/` must agree — every query goes through the drizzle table objects, and a test asserts the migrations create exactly the tables `schema.ts` declares, no more and no less. Change one, change the other, in the same commit. @@ -235,12 +235,12 @@ which store is installed. ### Data model Four physical tables — `artifact`, `artifact_version`, `upload`, -`mail_attachment_ref` — plus this package's own migration ledger. +`mail_attachment_ref`. **Hard control-plane foreign keys, by design.** `tenant_id` is `NOT NULL` and -references `public.tenant(id)` (`ON DELETE CASCADE` — a deleted tenant takes its +references the host's `tenant(id)` (`ON DELETE CASCADE` — a deleted tenant takes its artifacts with it) and `principal_id` / `owner_principal_id` reference -`public.principal(id)` (`ON DELETE SET NULL` — a removed principal detaches its +the host's `principal(id)` (`ON DELETE SET NULL` — a removed principal detaches its artifacts rather than destroying them). This package is coupled to Interchange: it mounts on Interchange-shaped hosts only, and the host's own migrations must have run before `runArtifactMigrations`. The internal key — @@ -248,7 +248,7 @@ have run before `runArtifactMigrations`. The internal key — **Cheap row-local CHECKs.** `artifact.version` and `artifact_version.version` must be ≥ 1; `upload.size` and `mail_attachment_ref.size` must be ≥ 0. These are -single-column constraints applied by a ledgered migration — free at write time. +single-column constraints — free at write time. **Principal↔tenant alignment is host-owned.** The package FKs each column into the control plane independently; it does **not** enforce that `principal_id` (or @@ -302,9 +302,7 @@ path that promises it. Separately, this package also has no way to confirm existing tenants are already free of duplicate `(title, kind)` rows, which would make even a scoped constraint risky to backfill. That is not the main reason for -rejecting the constraint, and it is not by itself decisive. See -`0003_schema_invariants` for this repo's own pattern for guarding a -migration against exactly that kind of bad existing data. +rejecting the constraint, and it is not by itself decisive. Instead, `findOrVersionArtifact(db, args)` (in `artifacts.ts`) closes the race with a transaction-scoped advisory lock keyed by @@ -336,55 +334,35 @@ behind it. ### Migration runner -`runArtifactMigrations(db)` is idempotent and safe to call unconditionally on -every boot of every replica. +`runArtifactMigrations(config, { schema })` takes the same arguments as +Interchange's `runMigrations`, and a host calls it right after that, with the +same values. `schema` is where the host's `tenant` and `principal` tables live; +the runner rewrites the `"public".` foreign-key references in the SQL files to +it. It is idempotent and safe to call on every boot of every replica. -- The whole run is one transaction whose first statements are - `SET LOCAL client_min_messages = warning` and a **transaction-scoped** - advisory lock. A transaction pins one pooled connection, so the lock, the - ledger read and the DDL are the same session; the lock releases on commit or - rollback, so there is no unlock call to lose on an error path. +- The whole run is one transaction whose first statement takes a + **transaction-scoped** advisory lock, so the lock releases on commit or + rollback and there is no unlock call to lose on an error path. `CREATE TABLE IF NOT EXISTS` is not itself race-safe, so the lock — not the `IF NOT EXISTS` — is what makes concurrent cold starts safe. -- Lowering `client_min_messages` is why a re-run prints **nothing**: every - statement is `IF NOT EXISTS`, and on the second boot Postgres answers each with - a NOTICE that postgres.js would otherwise dump to the console, making a clean - re-boot look like a wall of errors. `SET LOCAL` scopes it to the transaction - and stops at NOTICE — WARNING and above still reach the host. -- Each migration applies inside a nested transaction (a savepoint) together with - its ledger row, so a migration can never be recorded as applied with only some - of its statements run. -- The ledger is this package's own table, `artifacts.migrations`, - never shared with a host's. Each row records a **checksum of the migration's - rendered SQL**, so editing a shipped migration fails with - `MigrationChecksumError` on the next boot instead of letting existing and - fresh databases diverge silently. Ship a new migration instead. The column is - `NOT NULL`, so the guarantee is unconditional: there is no unrecorded row for - the runner to adopt and wave through. +- The runner opens its own single-connection client and discards NOTICEs, so a + re-run, where Postgres answers every `IF NOT EXISTS` with a NOTICE, prints + nothing. - Event timestamps (`created_at`, `updated_at`, `archived_at`) are - **`timestamptz`**. The initial create migration still lays them down as - zoneless `timestamp`; a follow-on migration retypes them with - `USING col AT TIME ZONE 'UTC'`, treating existing walls as the UTC clocks the - package always assumed. List keyset cursors project through - `AT TIME ZONE 'UTC'` and compare with `::timestamptz`, so paging and date - filters stay on the absolute instant under any session `TimeZone`. Rollback is - the reverse cast (`TYPE timestamp USING col AT TIME ZONE 'UTC'`) plus a new - ledgered migration — never edit a shipped one. -- A later ledgered migration sets `artifact.tenant_id NOT NULL` and adds the - version/size CHECKs. If null-tenant rows still exist, that migration raises - before altering the column so the operator can clean them up first. -- Empty ledger + pre-existing package objects fails closed - (`MigrationAdoptError`). `{ adopt: true }` records checksums without re-DDL - only after shape validation: tables, column types, required nullability, and - the named CHECK constraints. Column presence alone is not enough. - -**The package owns its own Postgres schema.** Every table, index and the ledger -live in `artifacts`, created by the runner and qualified in every -DDL statement and every query — nothing resolves through `search_path`, so the -package shares a database with the host's control plane without ever being able -to collide with (or silently adopt) a host table of the same name. The coupling -to the host is explicit instead: `tenant_id` and the principal columns are hard -FKs into `public.tenant` / `public.principal` (see the data model). + **`timestamptz`**. List keyset cursors project through `AT TIME ZONE 'UTC'` + and compare with `::timestamptz`, so paging and date filters stay on the + absolute instant under any session `TimeZone`. +- Releases up to 0.1.0 kept a checksum ledger in `artifacts.migrations`. + `0002_drop_migration_ledger.sql` removes it; a database migrated by 0.1.0 + already has the shape `0001_artifacts.sql` creates, so its statements no-op. + +**The package owns its own Postgres schema.** Every table and index lives in +`artifacts`, created by the runner and qualified in every DDL statement and +every query — nothing resolves through `search_path`, so the package shares a +database with the host's control plane without ever being able to collide with +(or silently adopt) a host table of the same name. The coupling to the host is +explicit instead: `tenant_id` and the principal columns are hard FKs into the +host schema's `tenant` / `principal` (see the data model). ### Boundaries diff --git a/README.md b/README.md index c32c33f..0c3fadf 100644 --- a/README.md +++ b/README.md @@ -10,13 +10,17 @@ npm add @corbits/artifacts Requires Node 24 or newer and `@intx/*` 0.4.0 or newer. -Open a database handle and apply this package's migrations at boot, before mounting any routes. A host that already has a drizzle handle passes that instead of calling `createArtifactDb`. +At boot, right after Interchange's `runMigrations`, apply this package's migrations with the same `config` and `schema`. The tables go in their own `artifacts` Postgres schema, with tenant and principal foreign keys pointing into `schema`. Then open a database handle for the routes; a host that already has a drizzle handle passes that instead of calling `createArtifactDb`. ```ts +import { runMigrations } from "@intx/db"; import { createArtifactDb, runArtifactMigrations } from "@corbits/artifacts"; +// `config` is the host's `DBConfig` from `@intx/db`. +await runMigrations(config, { schema: "public" }); +await runArtifactMigrations(config, { schema: "public" }); + const { db, close } = createArtifactDb(process.env.DATABASE_URL!); -await runArtifactMigrations(db); // on shutdown await close(); diff --git a/bun.lock b/bun.lock index e21314e..cfa6ef3 100644 --- a/bun.lock +++ b/bun.lock @@ -11,6 +11,7 @@ "devDependencies": { "@intx/agent": "0.4.0", "@intx/authz": "0.4.0", + "@intx/db": "0.4.0", "@intx/hub-api": "0.4.0", "@intx/types": "0.4.0", "@types/bun": "1.1.14", @@ -23,6 +24,7 @@ }, "peerDependencies": { "@intx/agent": "^0.4.0", + "@intx/db": "^0.4.0", "@intx/hub-api": "^0.4.0", "@intx/types": "^0.4.0", "drizzle-orm": "^0.45.2", diff --git a/examples/reference-host/src/index.ts b/examples/reference-host/src/index.ts index 99ba762..ec869cb 100644 --- a/examples/reference-host/src/index.ts +++ b/examples/reference-host/src/index.ts @@ -22,7 +22,13 @@ import { type RequireGrant, type TenantEnv, } from "@intx/hub-api"; -import { createDB, createGrantStore, runMigrations, schema as intxSchema } from "@intx/db"; +import { + createDB, + createGrantStore, + runMigrations, + schema as intxSchema, + type DBConfig, +} from "@intx/db"; // Interchange owns its id scheme; the host mints its OWN control-plane rows // with it rather than inventing a second one. import { generateId } from "@intx/hub-common"; @@ -50,7 +56,7 @@ export const DATABASE_URL = const EPOCH = new Date(0); -function parsePostgresUrl(raw: string) { +function parsePostgresUrl(raw: string): DBConfig { const url = new URL(raw); return { host: url.hostname, @@ -103,6 +109,8 @@ export type Session = { userId: string } | null; export type ReferenceHost = { db: ArtifactDb; + /** The `DBConfig` the host migrates with. */ + config: DBConfig; /** Interchange tenant id every artifact in this host is scoped to. */ tenantId: string; /** Principal id of the agent Alice owns. */ @@ -130,13 +138,14 @@ export async function createReferenceHost(): Promise { // ONE pool. The artifact module mounts on the handle the host already has // from `createDB` — the seam takes any drizzle postgres-js instance, so there // is no second connection to the same database. - const hub = createDB(parsePostgresUrl(DATABASE_URL)); + const config = parsePostgresUrl(DATABASE_URL); + const hub = createDB(config); const db: ArtifactDb = hub.db; // This host resets and truncates its database on boot — refuse to run // against anything that doesn't look like a throwaway database unless // explicitly opted in. - const { database } = parsePostgresUrl(DATABASE_URL); + const { database } = config; if ( !database.startsWith("artifact_") && process.env.ARTIFACT_REFERENCE_ALLOW_RESET !== "1" @@ -156,8 +165,8 @@ export async function createReferenceHost(): Promise { // principal stand-ins as FK targets; on a shared dev database those look // present by name but lack Interchange's columns, so detect by shape // (`tenant.slug`), drop the stand-ins, and migrate for real. - // `runArtifactMigrations` needs no such guard — carrying its own ledger is - // precisely why it can be called unconditionally on every boot. + // `runArtifactMigrations` needs no such guard: every statement is idempotent, + // so it runs unconditionally on every boot. const [hostSchema] = await db.execute<{ present: boolean }>(sql` SELECT EXISTS ( SELECT 1 FROM information_schema.columns @@ -168,9 +177,9 @@ export async function createReferenceHost(): Promise { await db.execute(sql`DROP TABLE IF EXISTS "public"."principal" CASCADE`); await db.execute(sql`DROP TABLE IF EXISTS "public"."tenant" CASCADE`); await db.execute(sql`DROP SCHEMA IF EXISTS "artifacts" CASCADE`); - await runMigrations(parsePostgresUrl(DATABASE_URL), { schema: "public" }); + await runMigrations(config, { schema: "public" }); } - await runArtifactMigrations(db); + await runArtifactMigrations(config, { schema: "public" }); await db.execute( sql`TRUNCATE TABLE "artifacts"."artifact", "artifacts"."artifact_version", "artifacts"."upload", "artifacts"."mail_attachment_ref" CASCADE`, ); @@ -362,6 +371,7 @@ export async function createReferenceHost(): Promise { return { db, + config, tenantId: tenant.id, agentPrincipal, scope: () => ({ diff --git a/examples/reference-host/test/acceptance.test.ts b/examples/reference-host/test/acceptance.test.ts index 176e239..784e444 100644 --- a/examples/reference-host/test/acceptance.test.ts +++ b/examples/reference-host/test/acceptance.test.ts @@ -679,15 +679,8 @@ describe("a skill-draft is invisible over the mounted host", () => { }); describe("the migration runner is re-runnable", () => { - test("re-running applies nothing new and destroys no data", async () => { - const before = await host.db.execute<{ id: string }>( - sql`SELECT "id" FROM "artifacts"."migrations"`, - ); - await runArtifactMigrations(host.db); - const after = await host.db.execute<{ id: string }>( - sql`SELECT "id" FROM "artifacts"."migrations"`, - ); - expect(after.length).toBe(before.length); + test("re-running destroys no data", async () => { + await runArtifactMigrations(host.config, { schema: "public" }); const survived = await json<{ artifacts: unknown[] }>( await host.request("/api/artifacts?limit=100"), diff --git a/migrations/0001_artifacts.sql b/migrations/0001_artifacts.sql new file mode 100644 index 0000000..dcce632 --- /dev/null +++ b/migrations/0001_artifacts.sql @@ -0,0 +1,111 @@ +CREATE SCHEMA IF NOT EXISTS "artifacts"; +--> statement-breakpoint +CREATE TABLE IF NOT EXISTS "artifacts"."artifact" ( + "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, + "tenant_id" text NOT NULL REFERENCES "public"."tenant"("id") ON DELETE CASCADE, + "principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, + "owner_principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, + "kind" text NOT NULL, + "title" text NOT NULL, + "content" text NOT NULL, + "source" jsonb, + "version" integer NOT NULL DEFAULT 1, + "metadata" jsonb, + "content_sha256" text, + "archived_at" timestamptz, + "created_at" timestamptz NOT NULL DEFAULT now(), + "updated_at" timestamptz NOT NULL DEFAULT now(), + CONSTRAINT "artifact_version_gte_1" CHECK ("version" >= 1) +); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "artifact_tenant_updated_id_idx" + ON "artifacts"."artifact" ("tenant_id", "updated_at", "id"); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "artifact_principal_idx" + ON "artifacts"."artifact" ("principal_id"); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "artifact_owner_principal_idx" + ON "artifacts"."artifact" ("owner_principal_id"); +--> statement-breakpoint +CREATE TABLE IF NOT EXISTS "artifacts"."artifact_version" ( + "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, + "artifact_id" text NOT NULL REFERENCES "artifacts"."artifact"("id") ON DELETE CASCADE, + "version" integer NOT NULL, + "title" text NOT NULL, + "content" text NOT NULL, + "author_id" text NOT NULL, + "metadata" jsonb, + "parent_version_ids" text[], + "content_sha256" text, + "created_at" timestamptz NOT NULL DEFAULT now(), + CONSTRAINT "artifact_version_artifact_id_version" UNIQUE ("artifact_id", "version"), + CONSTRAINT "artifact_version_version_gte_1" CHECK ("version" >= 1) +); +--> statement-breakpoint +CREATE TABLE IF NOT EXISTS "artifacts"."upload" ( + "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, + "tenant_id" text NOT NULL REFERENCES "public"."tenant"("id") ON DELETE CASCADE, + "principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, + "filename" text NOT NULL, + "mime_type" text NOT NULL, + "content" bytea NOT NULL, + "size" integer NOT NULL, + "created_at" timestamptz NOT NULL DEFAULT now(), + CONSTRAINT "upload_size_gte_0" CHECK ("size" >= 0) +); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "upload_tenant_idx" ON "artifacts"."upload" ("tenant_id"); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "upload_principal_idx" ON "artifacts"."upload" ("principal_id"); +--> statement-breakpoint +CREATE TABLE IF NOT EXISTS "artifacts"."mail_attachment_ref" ( + "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, + "tenant_id" text NOT NULL REFERENCES "public"."tenant"("id") ON DELETE CASCADE, + "principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, + "instance_id" text NOT NULL, + "mail_id" text NOT NULL, + "artifact_id" text NOT NULL, + "name" text NOT NULL, + "mime_type" text NOT NULL, + "size" integer NOT NULL, + "created_at" timestamptz NOT NULL DEFAULT now(), + CONSTRAINT "mail_attachment_ref_mail_id_artifact_id" UNIQUE ("mail_id", "artifact_id"), + CONSTRAINT "mail_attachment_ref_size_gte_0" CHECK ("size" >= 0) +); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "mail_attachment_ref_instance_idx" + ON "artifacts"."mail_attachment_ref" ("instance_id"); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "mail_attachment_ref_tenant_idx" + ON "artifacts"."mail_attachment_ref" ("tenant_id"); +--> statement-breakpoint +CREATE INDEX IF NOT EXISTS "mail_attachment_ref_principal_idx" + ON "artifacts"."mail_attachment_ref" ("principal_id"); +--> statement-breakpoint +-- A database from before 0.1.0 that never booted on 0.1.0 still lacks the +-- columns 0.1.0 added and keeps zoneless timestamps; bring it to 0.1.0's shape. +ALTER TABLE "artifacts"."artifact" + ADD COLUMN IF NOT EXISTS "metadata" jsonb, + ADD COLUMN IF NOT EXISTS "content_sha256" text; +--> statement-breakpoint +ALTER TABLE "artifacts"."artifact_version" + ADD COLUMN IF NOT EXISTS "metadata" jsonb, + ADD COLUMN IF NOT EXISTS "parent_version_ids" text[], + ADD COLUMN IF NOT EXISTS "content_sha256" text; +--> statement-breakpoint +-- 0.1.0 treated existing zoneless values as UTC; only columns still zoneless +-- are rewritten, so a current database is untouched. +DO $$ +DECLARE + "col" record; +BEGIN + FOR "col" IN + SELECT "table_name", "column_name" FROM information_schema.columns + WHERE "table_schema" = 'artifacts' AND "data_type" = 'timestamp without time zone' + LOOP + EXECUTE format( + 'ALTER TABLE "artifacts".%I ALTER COLUMN %I TYPE timestamptz USING %I AT TIME ZONE ''UTC''', + "col"."table_name", "col"."column_name", "col"."column_name" + ); + END LOOP; +END $$; diff --git a/migrations/0002_drop_migration_ledger.sql b/migrations/0002_drop_migration_ledger.sql new file mode 100644 index 0000000..d2efaed --- /dev/null +++ b/migrations/0002_drop_migration_ledger.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS "artifacts"."migrations"; diff --git a/package.json b/package.json index ae9b78a..29956b2 100644 --- a/package.json +++ b/package.json @@ -45,6 +45,7 @@ }, "files": [ "dist", + "migrations", "LICENSE", "README.md" ], @@ -70,7 +71,8 @@ "hono-openapi": "^1.2.0", "postgres": "^3.4.9", "@intx/hub-api": "^0.4.0", - "@intx/agent": "^0.4.0" + "@intx/agent": "^0.4.0", + "@intx/db": "^0.4.0" }, "peerDependenciesMeta": { "@intx/agent": { @@ -88,7 +90,8 @@ "typescript": "5.7.2", "@intx/hub-api": "0.4.0", "@intx/authz": "0.4.0", - "@intx/agent": "0.4.0" + "@intx/agent": "0.4.0", + "@intx/db": "0.4.0" }, "overrides": { "drizzle-orm": "0.45.2" diff --git a/src/index.ts b/src/index.ts index b5235f8..bcc64c9 100644 --- a/src/index.ts +++ b/src/index.ts @@ -19,8 +19,7 @@ export type { ArtifactCounts, ArtifactCountSegments } from "./counts.js"; export { artifactPreviewHeaders, resolveArtifactPreview } from "./preview.js"; export type { ArtifactPreviewResult } from "./preview.js"; -export { runArtifactMigrations, MigrationChecksumError, MigrationAdoptError } from "./migrations.js"; -export type { RunArtifactMigrationsOptions } from "./migrations.js"; +export { runArtifactMigrations } from "./migrations.js"; export { createArtifactDb } from "./db.js"; export type { ArtifactDb, ArtifactTx } from "./db.js"; diff --git a/src/migrations.test.ts b/src/migrations.test.ts index a8c07e0..1d0572f 100644 --- a/src/migrations.test.ts +++ b/src/migrations.test.ts @@ -1,810 +1,188 @@ -import { describe, expect, test } from "bun:test"; +import { beforeAll, describe, expect, test } from "bun:test"; import { getTableName, is, sql } from "drizzle-orm"; -import { drizzle } from "drizzle-orm/postgres-js"; import { PgTable } from "drizzle-orm/pg-core"; -import postgres from "postgres"; import { createArtifactDb } from "./db.js"; -import { - migrationChecksum, - MigrationAdoptError, - MigrationChecksumError, - MIGRATIONS, - runArtifactMigrations, -} from "./migrations.js"; +import { runArtifactMigrations } from "./migrations.js"; import * as schema from "./schema.js"; import { assertDestructiveArtifactTestsAllowed, + databaseConfig, DATABASE_URL, + ensureControlPlane, } from "./test-helpers.js"; const SCHEMA = "artifacts"; -const LEDGER = "migrations"; +const config = databaseConfig(DATABASE_URL); /** - * DERIVED FROM `schema.ts`, never restated. A hardcoded list is not a coverage - * guard: adding a table to the schema without a migration would leave the list - * describing the old world and the suite green, which is the precise failure - * this test exists to catch. Reading the declarations back off the drizzle - * objects means the guard grows itself the moment the schema does. + * DERIVED FROM `schema.ts`, never restated, so a table added to the schema + * without a migration (or the reverse) turns the equality check red. */ const DECLARED_TABLES = (Object.values(schema) as unknown[]) .filter((value): value is PgTable => is(value, PgTable)) .map(getTableName) .sort(); -// A shared handle for the assertions that only inspect catalogue state. const { db } = createArtifactDb(DATABASE_URL); -describe("migrations", () => { - test("creates every declared table and records each migration once", async () => { - await runArtifactMigrations(db); - - const tables = await db.execute<{ table_name: string }>(sql` - SELECT table_name FROM information_schema.tables - WHERE table_schema = ${SCHEMA} - `); - const names = new Set(tables.map((t) => t.table_name)); - expect(DECLARED_TABLES.length).toBeGreaterThan(0); - for (const table of [...DECLARED_TABLES, LEDGER]) { - expect({ table, created: names.has(table) }).toEqual({ table, created: true }); - } - - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toEqual(MIGRATIONS.map((m) => m.id)); - }); +async function packageTables(): Promise { + const rows = await db.execute<{ table_name: string }>(sql` + SELECT table_name FROM information_schema.tables + WHERE table_schema = ${SCHEMA} + `); + return rows.map((r) => r.table_name).sort(); +} + +async function rejection(query: ReturnType): Promise { + return await db.execute(query).then( + () => "", + (error: unknown) => String((error as { cause?: unknown }).cause), + ); +} + +beforeAll(async () => { + assertDestructiveArtifactTestsAllowed(DATABASE_URL); + await ensureControlPlane(db); +}); - test("re-running is a no-op: no duplicate ledger rows, no error", async () => { - await runArtifactMigrations(db); - await runArtifactMigrations(db); - await runArtifactMigrations(db); +describe("runArtifactMigrations", () => { + test("creates exactly the tables schema.ts declares, and re-running is a no-op", async () => { + await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); + await runArtifactMigrations(config, { schema: "public" }); + await runArtifactMigrations(config, { schema: "public" }); - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)}`, - ); - expect(ledger.length).toBe(MIGRATIONS.length); + expect(DECLARED_TABLES.length).toBeGreaterThan(0); + expect(await packageTables()).toEqual(DECLARED_TABLES); }); - // Four cold starts against a database where the package schema does not - // exist at all — the real first-boot race. Without the advisory lock this - // fails: `CREATE SCHEMA/TABLE IF NOT EXISTS` checks the catalogue and - // inserts non-atomically, so the losers get a 23505 duplicate key. It is - // the lock, not the IF NOT EXISTS, that makes the runner safe to call from - // every instance at once. + // It is the advisory lock, not IF NOT EXISTS, that makes concurrent cold + // starts safe: the catalogue check and insert are not atomic. test("concurrent first boots against a fresh database all succeed", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - // Each handle is its own pool, so all four genuinely believe they are the - // first boot. - const handles = [0, 1, 2, 3].map(() => createArtifactDb(DATABASE_URL)); - try { - await Promise.all(handles.map((h) => runArtifactMigrations(h.db))); - - const ledger = await handles[0]!.db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toEqual(MIGRATIONS.map((m) => m.id)); - - const tables = await db.execute<{ table_name: string }>(sql` - SELECT table_name FROM information_schema.tables - WHERE table_schema = ${SCHEMA} - `); - const names = new Set(tables.map((t) => t.table_name)); - for (const table of [...DECLARED_TABLES, LEDGER]) { - expect({ table, created: names.has(table) }).toEqual({ table, created: true }); - } - } finally { - await Promise.all(handles.map((h) => h.close())); - } - }); - - // The lock is transaction-scoped, so a failed run releases it by rolling - // back — there is no unlock statement that could be issued on a different - // pooled connection, return false, and leak the lock forever. - test("the advisory lock is not held after the run, success or failure", async () => { - // pg_locks is cluster-wide, so this observes every backend, not just ours. - const heldNow = async () => { - const rows = await db.execute<{ n: number }>(sql` - SELECT count(*)::int AS n FROM pg_locks - WHERE locktype = 'advisory' AND objid = ${0x0a27_1f04} - `); - return rows[0]!.n; - }; - - const { db: handle, close } = createArtifactDb(DATABASE_URL); - try { - // SUCCESS path. - await runArtifactMigrations(handle); - expect({ path: "success", held: await heldNow() }).toEqual({ - path: "success", - held: 0, - }); - - // FAILURE path — the half the name claimed and the body never ran. A - // stale checksum makes the run throw AFTER the lock has been taken, which - // is the only interesting case: a lock leaked here would block every - // later boot forever, and it is exactly what the session-scoped - // `pg_advisory_lock`/`pg_advisory_unlock` pair used to risk. - await db.execute( - sql`UPDATE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} SET "checksum" = 'stale' - WHERE "id" = ${MIGRATIONS[0]!.id}`, - ); - await expect(runArtifactMigrations(handle)).rejects.toThrow( - MigrationChecksumError, - ); - expect({ path: "failure", held: await heldNow() }).toEqual({ - path: "failure", - held: 0, - }); - - // And the proof that "released" means usable, not merely absent from - // pg_locks: the very next run takes the lock again instead of hanging. - await db.execute( - sql`UPDATE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} SET "checksum" = ${migrationChecksum(MIGRATIONS[0]!)} - WHERE "id" = ${MIGRATIONS[0]!.id}`, - ); - await runArtifactMigrations(handle); - expect(await heldNow()).toBe(0); - } finally { - await close(); - } - }); - - test("the standalone handle can be closed, so a migration script exits", async () => { - const { db: standalone, close } = createArtifactDb(DATABASE_URL); - await runArtifactMigrations(standalone); - await close(); - await expect(runArtifactMigrations(standalone)).rejects.toThrow(); - }); - - // Every DDL in the runner is `IF NOT EXISTS`, so on the second and every - // later boot Postgres answers each one with a NOTICE. postgres.js has no - // notice handler by default and dumps the raw object to the console, so a - // clean re-boot of a runner documented as "safe to call on every boot" read - // as a wall of errors on every replica start. - test("a re-run emits no NOTICE, on a handle this package did not construct", async () => { - const notices: { severity?: string; code?: string; message?: string }[] = []; - const client = postgres(DATABASE_URL, { - onnotice: (notice) => notices.push(notice), - }); - const handle = drizzle(client); - try { - // First run creates whatever is missing; second run is the boot that used - // to be noisy — every object already exists. - await runArtifactMigrations(handle); - notices.length = 0; - await runArtifactMigrations(handle); - expect(notices).toEqual([]); - } finally { - await client.end(); - } - }); - - // The silencing must stop at NOTICE. A host that raises a WARNING on the same - // connection still hears it — this is not a blanket mute. - test("WARNING and above still reach the host", async () => { - const notices: { severity?: string; message?: string }[] = []; - const client = postgres(DATABASE_URL, { - onnotice: (notice) => notices.push(notice), - }); - const handle = drizzle(client); - try { - await runArtifactMigrations(handle); - await handle.transaction(async (tx) => { - await tx.execute(sql`SET LOCAL client_min_messages = warning`); - await tx.execute(sql`DO $$ BEGIN RAISE NOTICE 'quiet'; END $$`); - await tx.execute(sql`DO $$ BEGIN RAISE WARNING 'loud'; END $$`); - }); - expect(notices.map((n) => n.severity)).toEqual(["WARNING"]); - expect(notices[0]!.message).toBe("loud"); - } finally { - await client.end(); - } - }); - - test("keeps its ledger inside its own schema, not a shared journal", async () => { - await runArtifactMigrations(db); - const rows = await db.execute<{ n: number }>(sql` - SELECT count(*)::int AS n FROM information_schema.tables - WHERE table_schema = ${SCHEMA} AND table_name = ${LEDGER} - `); - expect(rows[0]!.n).toBe(1); - }); - - test("every migration id is unique and every migration has statements", () => { - const ids = MIGRATIONS.map((m) => m.id); - expect(new Set(ids).size).toBe(ids.length); - for (const migration of MIGRATIONS) { - expect(migration.statements.length).toBeGreaterThan(0); - } - }); - - test("records a checksum of each migration BODY, not just its id", async () => { - await runArtifactMigrations(db); - const ledger = await db.execute<{ id: string; checksum: string | null }>( - sql`SELECT "id", "checksum" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, + await Promise.all( + [0, 1, 2, 3].map(() => runArtifactMigrations(config, { schema: "public" })), ); - for (const row of ledger) { - const migration = MIGRATIONS.find((m) => m.id === row.id)!; - expect(row.checksum).toBe(migrationChecksum(migration)); - } - }); - - test("an edited shipped migration fails loudly instead of silently diverging", async () => { - await runArtifactMigrations(db); - // Simulate the ledger having recorded a DIFFERENT body for 0001: exactly - // what a deployed database looks like after someone edits a shipped - // migration. Without a checksum this database would skip it forever while - // a fresh one gets the new DDL, and nothing would notice. - await db.execute( - sql`UPDATE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} SET "checksum" = 'stale' WHERE "id" = '0001_artifacts'`, - ); - await expect(runArtifactMigrations(db)).rejects.toThrow(MigrationChecksumError); - - // The advisory lock is released even on the failure path, so the next boot - // can still run rather than hanging forever. - await db.execute( - sql`UPDATE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} SET "checksum" = ${migrationChecksum(MIGRATIONS[0]!)} WHERE "id" = '0001_artifacts'`, - ); - await runArtifactMigrations(db); - }); - - test("the ledger has no nullable-checksum escape hatch", async () => { - // The adopt-silently branch for "ledgers written before checksums existed" - // is gone: 0.1.0 is the first public release, so no such ledger can exist, - // and while the column was nullable the runner would accept exactly one - // edit to a shipped migration without complaint. NOT NULL is what makes the - // documented immutability guarantee unconditional. - await runArtifactMigrations(db); - // Drizzle wraps driver errors, so the NOT NULL violation is on `.cause`, - // not on the message `toThrow` would match. - const failure = await db - .execute(sql`UPDATE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} SET "checksum" = NULL`) - .then( - () => null, - (error: unknown) => error, - ); - expect(failure).not.toBeNull(); - expect(String((failure as { cause?: unknown }).cause)).toMatch( - /null value in column "checksum"/, - ); - }); - - test("the keyset index carries the id tie-break", async () => { - await runArtifactMigrations(db); - const indexes = await db.execute<{ indexname: string; indexdef: string }>(sql` - SELECT indexname, indexdef FROM pg_indexes - WHERE schemaname = ${SCHEMA} AND tablename = 'artifact' - `); - const keyset = indexes.find( - (i) => i.indexname === "artifact_tenant_updated_id_idx", - ); - expect(keyset?.indexdef).toContain("(tenant_id, updated_at, id)"); - }); - - test("the (artifactId, version) uniqueness backstop is really in the database", async () => { - await runArtifactMigrations(db); - const constraints = await db.execute<{ conname: string }>(sql` - SELECT conname FROM pg_constraint - WHERE conname = 'artifact_version_artifact_id_version' - `); - expect(constraints.length).toBe(1); + expect(await packageTables()).toEqual(DECLARED_TABLES); }); - /** - * The coverage guard proper, and the reason it runs against a VIRGIN schema: - * the shared `public` schema accumulates whatever else has ever been created - * there, so it can only support a "created everything declared" check. On an - * empty schema the migrations are the only writer, which makes the comparison - * an EQUALITY — the guard then catches drift in both directions: - * - * - a table added to `schema.ts` with no migration (nothing creates it), and - * - a table created by a migration that `schema.ts` never declares (dead - * schema no code can address). - */ - test("the migrations create exactly the tables schema.ts declares, no more, no less", async () => { - // The package schema is dropped and rebuilt from empty, so the migrations - // are its only writer and the comparison is an EQUALITY. - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - const tables = await db.execute<{ table_name: string }>(sql` - SELECT table_name FROM information_schema.tables - WHERE table_schema = ${SCHEMA} + test("drops the migration ledger earlier releases kept", async () => { + await db.execute(sql` + CREATE TABLE IF NOT EXISTS "artifacts"."migrations" ("id" text PRIMARY KEY) `); - const created = tables - .map((t) => t.table_name) - .filter((name) => name !== LEDGER) - .sort(); - expect(created).toEqual(DECLARED_TABLES); - }); - - /** - * Empty ledger + pre-existing package objects is the AUDIT-013 footgun: - * `CREATE IF NOT EXISTS` no-ops and a naive runner would still stamp the - * current checksum, so boot "succeeds" while the live shape is wrong. - * These three cases pin the adopt-or-fail contract. - */ - test("clean install still migrates when the package schema is empty", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - - const ledger = await db.execute<{ id: string; checksum: string }>( - sql`SELECT "id", "checksum" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toEqual(MIGRATIONS.map((m) => m.id)); - for (const row of ledger) { - const migration = MIGRATIONS.find((m) => m.id === row.id)!; - expect(row.checksum).toBe(migrationChecksum(migration)); - } + await runArtifactMigrations(config, { schema: "public" }); + expect(await packageTables()).toEqual(DECLARED_TABLES); }); - test("empty ledger + wrong shape fails closed and does not write a checksum", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - // Partial / wrong object: schema + a table that is not the package shape. - await db.execute(sql`CREATE SCHEMA ${sql.identifier(SCHEMA)}`); + test("brings a database from before 0.1.0 up to 0.1.0's columns", async () => { await db.execute(sql` - CREATE TABLE ${sql.identifier(SCHEMA)}."artifact" ( - "id" text PRIMARY KEY, - "not_the_real_shape" text - ) + ALTER TABLE "artifacts"."artifact" + DROP COLUMN "metadata", DROP COLUMN "content_sha256", + ALTER COLUMN "created_at" TYPE timestamp `); - - await expect(runArtifactMigrations(db)).rejects.toThrow(MigrationAdoptError); - await expect( - runArtifactMigrations(db, { adopt: true }), - ).rejects.toThrow(MigrationAdoptError); - - // No ledger row may have been recorded — including a ledger that was - // created mid-run and then rolled back with the failed transaction. - const ledgerTables = await db.execute<{ n: number }>(sql` - SELECT count(*)::int AS n FROM information_schema.tables - WHERE table_schema = ${SCHEMA} AND table_name = ${LEDGER} + await db.execute(sql` + ALTER TABLE "artifacts"."artifact_version" + DROP COLUMN "metadata", DROP COLUMN "parent_version_ids", DROP COLUMN "content_sha256" `); - expect(ledgerTables[0]!.n).toBe(0); - - // And the wrong object is still the only package table — runner must not - // have half-applied real DDL around it. - const tables = await db.execute<{ table_name: string }>(sql` - SELECT table_name FROM information_schema.tables + await runArtifactMigrations(config, { schema: "public" }); + const columns = await db.execute<{ column: string; type: string }>(sql` + SELECT table_name || '.' || column_name AS column, udt_name AS type + FROM information_schema.columns WHERE table_schema = ${SCHEMA} - `); - expect(tables.map((t) => t.table_name).sort()).toEqual(["artifact"]); + AND column_name IN ('metadata', 'parent_version_ids', 'content_sha256', 'created_at') + AND table_name IN ('artifact', 'artifact_version') + ORDER BY 1 + `); + expect(columns.map((c) => `${c.column} ${c.type}`)).toEqual([ + "artifact.content_sha256 text", + "artifact.created_at timestamptz", + "artifact.metadata jsonb", + "artifact_version.content_sha256 text", + "artifact_version.created_at timestamptz", + "artifact_version.metadata jsonb", + "artifact_version.parent_version_ids _text", + ]); }); - test("empty ledger + compatible shape requires explicit adopt; adopt records checksums without re-DDL", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); + test("points the tenant and principal foreign keys at the host schema", async () => { + const host = "artifact_host_test"; await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - // Build a correct shape the honest way, then erase the ledger so the next - // boot sees "objects present, ledger empty" — the adopt path's input. - await runArtifactMigrations(db); + await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(host)} CASCADE`); + await db.execute(sql`CREATE SCHEMA ${sql.identifier(host)}`); + await db.execute(sql`CREATE TABLE ${sql.identifier(host)}."tenant" ("id" text PRIMARY KEY)`); await db.execute( - sql`DROP TABLE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)}`, - ); - - // Without the flag: fail closed even though the shape is right. Silent - // adoption is how a wrong-but-lucky shape used to get a checksum stamp. - await expect(runArtifactMigrations(db)).rejects.toThrow(MigrationAdoptError); - - const stillEmpty = await db.execute<{ n: number }>(sql` - SELECT count(*)::int AS n FROM information_schema.tables - WHERE table_schema = ${SCHEMA} AND table_name = ${LEDGER} - `); - expect(stillEmpty[0]!.n).toBe(0); - - // With the flag: shape validates, ledger is written, objects stay put. - await runArtifactMigrations(db, { adopt: true }); - - const ledger = await db.execute<{ id: string; checksum: string }>( - sql`SELECT "id", "checksum" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, + sql`CREATE TABLE ${sql.identifier(host)}."principal" ("id" text PRIMARY KEY)`, ); - expect(ledger.map((r) => r.id)).toEqual(MIGRATIONS.map((m) => m.id)); - for (const row of ledger) { - const migration = MIGRATIONS.find((m) => m.id === row.id)!; - expect(row.checksum).toBe(migrationChecksum(migration)); + try { + await runArtifactMigrations(config, { schema: host }); + const targets = await db.execute<{ target: string }>(sql` + SELECT DISTINCT tn.nspname AS target + FROM pg_constraint con + JOIN pg_class c ON c.oid = con.conrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + JOIN pg_class t ON t.oid = con.confrelid + JOIN pg_namespace tn ON tn.oid = t.relnamespace + WHERE con.contype = 'f' AND n.nspname = ${SCHEMA} AND t.relname IN ('tenant', 'principal') + `); + expect(targets.map((r) => r.target)).toEqual([host]); + } finally { + await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); + await db.execute(sql`DROP SCHEMA ${sql.identifier(host)} CASCADE`); + await runArtifactMigrations(config, { schema: "public" }); } - - // Re-run (no adopt needed once ledgered) remains a quiet no-op. - await runArtifactMigrations(db); - const again = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)}`, - ); - expect(again.length).toBe(MIGRATIONS.length); }); - test("empty ledger + columns without 0003 CHECKs refuses adopt", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - // Honest migrate, drop ledger, then strip a row-local CHECK so columns and - // types still match but 0003 invariants are gone — adopt must not stamp. - await runArtifactMigrations(db); - await db.execute( - sql`DROP TABLE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)}`, + test("refuses an empty schema name", async () => { + await expect(runArtifactMigrations(config, { schema: "" })).rejects.toThrow( + "schema name must not be empty", ); - await db.execute(sql` - ALTER TABLE ${sql.identifier(SCHEMA)}."artifact" - DROP CONSTRAINT "artifact_version_gte_1" - `); - - await expect( - runArtifactMigrations(db, { adopt: true }), - ).rejects.toThrow(MigrationAdoptError); - await expect( - runArtifactMigrations(db, { adopt: true }), - ).rejects.toThrow(/artifact_version_gte_1/); - - const ledgerTables = await db.execute<{ n: number }>(sql` - SELECT count(*)::int AS n FROM information_schema.tables - WHERE table_schema = ${SCHEMA} AND table_name = ${LEDGER} - `); - expect(ledgerTables[0]!.n).toBe(0); - - // Later tests call runArtifactMigrations without a reset; restore a - // fully-ledgered schema so they are not stranded on an empty ledger. - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); }); - test("empty ledger + nullable tenant_id refuses adopt", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - await db.execute( - sql`DROP TABLE ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)}`, - ); - await db.execute(sql` - ALTER TABLE ${sql.identifier(SCHEMA)}."artifact" - ALTER COLUMN "tenant_id" DROP NOT NULL - `); - - await expect( - runArtifactMigrations(db, { adopt: true }), - ).rejects.toThrow(MigrationAdoptError); - await expect( - runArtifactMigrations(db, { adopt: true }), - ).rejects.toThrow(/tenant_id.*NOT NULL/i); - - const ledgerTables = await db.execute<{ n: number }>(sql` - SELECT count(*)::int AS n FROM information_schema.tables - WHERE table_schema = ${SCHEMA} AND table_name = ${LEDGER} + test("the keyset index carries the id tie-break", async () => { + const [keyset] = await db.execute<{ indexdef: string }>(sql` + SELECT indexdef FROM pg_indexes + WHERE schemaname = ${SCHEMA} AND indexname = 'artifact_tenant_updated_id_idx' `); - expect(ledgerTables[0]!.n).toBe(0); - - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - }); - - /** - * DB invariants: tenant_id is required on every artifact row; version and size - * stay non-negative. Principal↔tenant alignment is host middleware/context - * (TenantEnv) — no multi-table trigger here. - */ - test("null tenant_id on artifact is rejected after migrations", async () => { - await runArtifactMigrations(db); - const failure = await db - .execute(sql` - INSERT INTO "artifacts"."artifact" - ("tenant_id", "kind", "title", "content", "version") - VALUES (NULL, 'document', 'orphan', 'body', 1) - `) - .then( - () => null, - (error: unknown) => error, - ); - expect(failure).not.toBeNull(); - expect(String((failure as { cause?: unknown }).cause)).toMatch( - /null value in column "tenant_id"/, - ); + expect(keyset?.indexdef).toContain("(tenant_id, updated_at, id)"); }); - test("version and size CHECK constraints reject impossible values", async () => { - await runArtifactMigrations(db); + test("tenant_id NOT NULL and the version and size CHECKs are enforced", async () => { + expect( + await rejection(sql` + INSERT INTO "artifacts"."artifact" ("tenant_id", "kind", "title", "content") + VALUES (NULL, 'document', 'orphan', 'body') + `), + ).toMatch(/null value in column "tenant_id"/); + expect( + await rejection(sql` + INSERT INTO "artifacts"."artifact" ("tenant_id", "kind", "title", "content", "version") + VALUES ('acme', 'document', 'bad', 'body', 0) + `), + ).toMatch(/artifact_version_gte_1/); - const badArtifactVersion = await db - .execute(sql` - INSERT INTO "artifacts"."artifact" - ("tenant_id", "kind", "title", "content", "version") - VALUES ('acme', 'document', 'bad-ver', 'body', 0) - `) - .then( - () => null, - (error: unknown) => error, - ); - expect(badArtifactVersion).not.toBeNull(); - expect(String((badArtifactVersion as { cause?: unknown }).cause)).toMatch( - /artifact_version_gte_1|check constraint/i, - ); - - // A legal artifact so we can try a bad history row and a bad upload. const [row] = await db.execute<{ id: string }>(sql` - INSERT INTO "artifacts"."artifact" - ("tenant_id", "kind", "title", "content", "version") - VALUES ('acme', 'document', 'ok', 'body', 1) + INSERT INTO "artifacts"."artifact" ("tenant_id", "kind", "title", "content") + VALUES ('acme', 'document', 'ok', 'body') RETURNING "id" `); - - const badHistory = await db - .execute(sql` + expect( + await rejection(sql` INSERT INTO "artifacts"."artifact_version" ("artifact_id", "version", "title", "content", "author_id") VALUES (${row!.id}, 0, 'ok', 'body', 'user-1') - `) - .then( - () => null, - (error: unknown) => error, - ); - expect(badHistory).not.toBeNull(); - expect(String((badHistory as { cause?: unknown }).cause)).toMatch( - /artifact_version_version_gte_1|check constraint/i, - ); - - const badUpload = await db - .execute(sql` - INSERT INTO "artifacts"."upload" - ("tenant_id", "filename", "mime_type", "content", "size") + `), + ).toMatch(/artifact_version_version_gte_1/); + expect( + await rejection(sql` + INSERT INTO "artifacts"."upload" ("tenant_id", "filename", "mime_type", "content", "size") VALUES ('acme', 'x.bin', 'application/octet-stream', decode('00', 'hex'), -1) - `) - .then( - () => null, - (error: unknown) => error, - ); - expect(badUpload).not.toBeNull(); - expect(String((badUpload as { cause?: unknown }).cause)).toMatch( - /upload_size_gte_0|check constraint/i, - ); - - const badRef = await db - .execute(sql` + `), + ).toMatch(/upload_size_gte_0/); + expect( + await rejection(sql` INSERT INTO "artifacts"."mail_attachment_ref" - ("tenant_id", "instance_id", "mail_id", "artifact_id", - "name", "mime_type", "size") - VALUES ('acme', 'inst', 'mail', ${row!.id}, 'a.bin', - 'application/octet-stream', -1) - `) - .then( - () => null, - (error: unknown) => error, - ); - expect(badRef).not.toBeNull(); - expect(String((badRef as { cause?: unknown }).cause)).toMatch( - /mail_attachment_ref_size_gte_0|check constraint/i, - ); - }); - - test("migration fails clearly when null tenant_id rows already exist", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - - // Apply only the baseline + timestamptz migrations so tenant_id is still - // nullable, plant a null-tenant row, then attempt the full runner (which - // must include the invariants migration). - const preInvariant = MIGRATIONS.filter( - (m) => m.id === "0001_artifacts" || m.id === "0002_timestamptz", - ); - expect(preInvariant.length).toBe(2); - expect(MIGRATIONS.some((m) => m.id === "0003_schema_invariants")).toBe(true); - - await db.transaction(async (tx) => { - await tx.execute(sql`CREATE SCHEMA IF NOT EXISTS ${sql.identifier(SCHEMA)}`); - await tx.execute(sql` - CREATE TABLE IF NOT EXISTS ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ( - "id" text PRIMARY KEY, - "checksum" text NOT NULL, - "applied_at" timestamptz NOT NULL DEFAULT now() - ) - `); - for (const migration of preInvariant) { - for (const statement of migration.statements) { - await tx.execute(statement); - } - await tx.execute(sql` - INSERT INTO ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} - ("id", "checksum") - VALUES (${migration.id}, ${migrationChecksum(migration)}) - `); - } - }); - - await db.execute(sql` - INSERT INTO "artifacts"."artifact" - ("tenant_id", "kind", "title", "content", "version") - VALUES (NULL, 'document', 'orphan', 'body', 1) - `); - - const failure = await runArtifactMigrations(db).then( - () => null, - (error: unknown) => error, - ); - expect(failure).not.toBeNull(); - expect(String(failure)).toMatch(/null tenant_id/i); - - // The invariants migration must not be stamped when the guard fails. - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toEqual([ - "0001_artifacts", - "0002_timestamptz", - ]); - - // Full suite (and re-runs of this file) must not inherit null-tenant rows - // and a half-applied ledger. - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - }); - - test("0004_version_metadata adds metadata and parent_version_ids columns", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - - const columns = await db.execute<{ - table_name: string; - column_name: string; - udt_name: string; - }>(sql` - SELECT table_name, column_name, udt_name - FROM information_schema.columns - WHERE table_schema = ${SCHEMA} - AND table_name IN ('artifact', 'artifact_version') - AND column_name IN ('metadata', 'parent_version_ids') - ORDER BY table_name, column_name - `); - expect([...columns]).toEqual([ - { table_name: "artifact", column_name: "metadata", udt_name: "jsonb" }, - { - table_name: "artifact_version", - column_name: "metadata", - udt_name: "jsonb", - }, - { - table_name: "artifact_version", - column_name: "parent_version_ids", - udt_name: "_text", - }, - ]); - - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toContain("0004_version_metadata"); - }); - - test("adopting a 0003-only database applies 0004's new columns forward", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - - // Build a database that only ever saw migrations through 0003 — the shape - // an existing production database has before this change ships. - const upTo0003 = MIGRATIONS.filter((m) => m.id !== "0004_version_metadata"); - await db.transaction(async (tx) => { - await tx.execute(sql`CREATE SCHEMA IF NOT EXISTS ${sql.identifier(SCHEMA)}`); - await tx.execute(sql` - CREATE TABLE IF NOT EXISTS ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ( - "id" text PRIMARY KEY, - "checksum" text NOT NULL, - "applied_at" timestamptz NOT NULL DEFAULT now() - ) - `); - for (const migration of upTo0003) { - for (const statement of migration.statements) { - await tx.execute(statement); - } - await tx.execute(sql` - INSERT INTO ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} - ("id", "checksum") - VALUES (${migration.id}, ${migrationChecksum(migration)}) - `); - } - }); - - // Running the full migration set adopts the new migration cleanly: no - // adopt flag needed, since the ledger already has real rows for 0001-0003 - // and only 0004 is missing — the ordinary "apply what's new" path. - await runArtifactMigrations(db); - - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toEqual(MIGRATIONS.map((m) => m.id)); - - const columns = await db.execute<{ column_name: string }>(sql` - SELECT column_name FROM information_schema.columns - WHERE table_schema = ${SCHEMA} AND table_name = 'artifact_version' - AND column_name IN ('metadata', 'parent_version_ids') - `); - expect(columns.map((c) => c.column_name).sort()).toEqual([ - "metadata", - "parent_version_ids", - ]); - - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - }); - - test("0005_version_content_digest adds a nullable content_sha256 column to both tables", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); - - const columns = await db.execute<{ - table_name: string; - column_name: string; - udt_name: string; - is_nullable: string; - }>(sql` - SELECT table_name, column_name, udt_name, is_nullable - FROM information_schema.columns - WHERE table_schema = ${SCHEMA} - AND table_name IN ('artifact', 'artifact_version') - AND column_name = 'content_sha256' - ORDER BY table_name - `); - expect([...columns]).toEqual([ - { - table_name: "artifact", - column_name: "content_sha256", - udt_name: "text", - is_nullable: "YES", - }, - { - table_name: "artifact_version", - column_name: "content_sha256", - udt_name: "text", - is_nullable: "YES", - }, - ]); - - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toContain("0005_version_content_digest"); - }); - - test("adopting a 0004-only database applies 0005's new column forward", async () => { - assertDestructiveArtifactTestsAllowed(DATABASE_URL); - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - - const upTo0004 = MIGRATIONS.filter((m) => m.id !== "0005_version_content_digest"); - await db.transaction(async (tx) => { - await tx.execute(sql`CREATE SCHEMA IF NOT EXISTS ${sql.identifier(SCHEMA)}`); - await tx.execute(sql` - CREATE TABLE IF NOT EXISTS ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ( - "id" text PRIMARY KEY, - "checksum" text NOT NULL, - "applied_at" timestamptz NOT NULL DEFAULT now() - ) - `); - for (const migration of upTo0004) { - for (const statement of migration.statements) { - await tx.execute(statement); - } - await tx.execute(sql` - INSERT INTO ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} - ("id", "checksum") - VALUES (${migration.id}, ${migrationChecksum(migration)}) - `); - } - }); - - await runArtifactMigrations(db); - - const ledger = await db.execute<{ id: string }>( - sql`SELECT "id" FROM ${sql.identifier(SCHEMA)}.${sql.identifier(LEDGER)} ORDER BY "id"`, - ); - expect(ledger.map((r) => r.id)).toEqual(MIGRATIONS.map((m) => m.id)); - - const columns = await db.execute<{ column_name: string }>(sql` - SELECT column_name FROM information_schema.columns - WHERE table_schema = ${SCHEMA} AND table_name = 'artifact_version' - AND column_name = 'content_sha256' - `); - expect(columns.map((c) => c.column_name)).toEqual(["content_sha256"]); - - await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`); - await runArtifactMigrations(db); + ("tenant_id", "instance_id", "mail_id", "artifact_id", "name", "mime_type", "size") + VALUES ('acme', 'inst', 'mail', ${row!.id}, 'a.bin', 'application/octet-stream', -1) + `), + ).toMatch(/mail_attachment_ref_size_gte_0/); }); }); diff --git a/src/migrations.ts b/src/migrations.ts index 795035c..2f68a00 100644 --- a/src/migrations.ts +++ b/src/migrations.ts @@ -1,601 +1,71 @@ -import { createHash } from "node:crypto"; -import { sql, type SQL } from "drizzle-orm"; -import { PgDialect } from "drizzle-orm/pg-core"; -import type { ArtifactDb, ArtifactTx } from "./db.js"; +// Applies migrations/*.sql, shipped next to dist/, into the `artifacts` schema. +import { readdir, readFile } from "node:fs/promises"; +import { join } from "node:path"; +import type { DBConfig } from "@intx/db"; +import postgres from "postgres"; -import { ARTIFACTS_SCHEMA } from "./schema.js"; - -// Own migration ledger, inside the package's own schema, so it never collides -// with the host's migration bookkeeping. -const LEDGER_TABLE = "migrations"; -const LEDGER = sql`${sql.identifier(ARTIFACTS_SCHEMA)}.${sql.identifier(LEDGER_TABLE)}`; - -export type Migration = { id: string; statements: SQL[] }; - -// Every statement is schema-qualified: this package owns the `artifacts` -// schema outright and never writes into the host's search_path. The tenant and -// principal columns are hard FKs into Interchange's control plane in `public` -// — the host's own migrations must have run first. -export const MIGRATIONS: Migration[] = [ - { - id: "0001_artifacts", - statements: [ - sql` - CREATE TABLE IF NOT EXISTS "artifacts"."artifact" ( - "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, - "tenant_id" text REFERENCES "public"."tenant"("id") ON DELETE CASCADE, - "principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, - "owner_principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, - "kind" text NOT NULL, - "title" text NOT NULL, - "content" text NOT NULL, - "source" jsonb, - "version" integer NOT NULL DEFAULT 1, - "archived_at" timestamp, - "created_at" timestamp NOT NULL DEFAULT now(), - "updated_at" timestamp NOT NULL DEFAULT now() - ) - `, - // The list is a keyset scan ordered by (updated_at, id) — the id - // tie-break must be IN the index or the cursor predicate degrades from - // an Index Cond to a Filter plus a sort. - sql` - CREATE INDEX IF NOT EXISTS "artifact_tenant_updated_id_idx" - ON "artifacts"."artifact" ("tenant_id", "updated_at", "id") - `, - // FK-support indexes, so ON DELETE CASCADE / SET NULL on the control - // plane never seq-scans these tables. - sql` - CREATE INDEX IF NOT EXISTS "artifact_principal_idx" - ON "artifacts"."artifact" ("principal_id") - `, - sql` - CREATE INDEX IF NOT EXISTS "artifact_owner_principal_idx" - ON "artifacts"."artifact" ("owner_principal_id") - `, - sql` - CREATE TABLE IF NOT EXISTS "artifacts"."artifact_version" ( - "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, - "artifact_id" text NOT NULL REFERENCES "artifacts"."artifact"("id") ON DELETE CASCADE, - "version" integer NOT NULL, - "title" text NOT NULL, - "content" text NOT NULL, - "author_id" text NOT NULL, - "created_at" timestamp NOT NULL DEFAULT now(), - CONSTRAINT "artifact_version_artifact_id_version" - UNIQUE ("artifact_id", "version") - ) - `, - sql` - CREATE TABLE IF NOT EXISTS "artifacts"."upload" ( - "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, - "tenant_id" text NOT NULL REFERENCES "public"."tenant"("id") ON DELETE CASCADE, - "principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, - "filename" text NOT NULL, - "mime_type" text NOT NULL, - "content" bytea NOT NULL, - "size" integer NOT NULL, - "created_at" timestamp NOT NULL DEFAULT now() - ) - `, - sql` - CREATE INDEX IF NOT EXISTS "upload_tenant_idx" - ON "artifacts"."upload" ("tenant_id") - `, - sql` - CREATE INDEX IF NOT EXISTS "upload_principal_idx" - ON "artifacts"."upload" ("principal_id") - `, - sql` - CREATE TABLE IF NOT EXISTS "artifacts"."mail_attachment_ref" ( - "id" text PRIMARY KEY DEFAULT gen_random_uuid()::text, - "tenant_id" text NOT NULL REFERENCES "public"."tenant"("id") ON DELETE CASCADE, - "principal_id" text REFERENCES "public"."principal"("id") ON DELETE SET NULL, - "instance_id" text NOT NULL, - "mail_id" text NOT NULL, - "artifact_id" text NOT NULL, - "name" text NOT NULL, - "mime_type" text NOT NULL, - "size" integer NOT NULL, - "created_at" timestamp NOT NULL DEFAULT now(), - CONSTRAINT "mail_attachment_ref_mail_id_artifact_id" - UNIQUE ("mail_id", "artifact_id") - ) - `, - sql` - CREATE INDEX IF NOT EXISTS "mail_attachment_ref_instance_idx" - ON "artifacts"."mail_attachment_ref" ("instance_id") - `, - sql` - CREATE INDEX IF NOT EXISTS "mail_attachment_ref_tenant_idx" - ON "artifacts"."mail_attachment_ref" ("tenant_id") - `, - sql` - CREATE INDEX IF NOT EXISTS "mail_attachment_ref_principal_idx" - ON "artifacts"."mail_attachment_ref" ("principal_id") - `, - ], - }, - { - // Zoneless `timestamp` stores a wall clock. Absolute instants written through - // a non-UTC session TimeZone land as session-local walls, so list cursors - // (which stamp a literal Z) and date filters compare the wrong instant. - // `timestamptz` keeps absolute time; existing walls are treated as UTC — the - // package always documented these columns as UTC wall clocks. - id: "0002_timestamptz", - statements: [ - sql` - ALTER TABLE "artifacts"."artifact" - ALTER COLUMN "archived_at" TYPE timestamptz - USING "archived_at" AT TIME ZONE 'UTC', - ALTER COLUMN "created_at" TYPE timestamptz - USING "created_at" AT TIME ZONE 'UTC', - ALTER COLUMN "updated_at" TYPE timestamptz - USING "updated_at" AT TIME ZONE 'UTC' - `, - sql` - ALTER TABLE "artifacts"."artifact_version" - ALTER COLUMN "created_at" TYPE timestamptz - USING "created_at" AT TIME ZONE 'UTC' - `, - sql` - ALTER TABLE "artifacts"."upload" - ALTER COLUMN "created_at" TYPE timestamptz - USING "created_at" AT TIME ZONE 'UTC' - `, - sql` - ALTER TABLE "artifacts"."mail_attachment_ref" - ALTER COLUMN "created_at" TYPE timestamptz - USING "created_at" AT TIME ZONE 'UTC' - `, - ], - }, - { - // Cheap row-local invariants. tenant_id was nullable in 0001 so a restored - // dump can still hold orphans; refuse to SET NOT NULL over them and tell the - // operator to assign a tenant or delete the rows first. Version/size CHECKs - // are single-column and free at write time. Principal↔tenant alignment is - // host middleware/context (TenantEnv) — a multi-table trigger into - // public.principal is out of scope and would couple write path latency to - // the control plane. - id: "0003_schema_invariants", - statements: [ - sql` - DO $guard$ - BEGIN - IF EXISTS ( - SELECT 1 FROM "artifacts"."artifact" WHERE "tenant_id" IS NULL - ) THEN - RAISE EXCEPTION - 'Cannot apply 0003_schema_invariants: artifacts.artifact has rows with null tenant_id. Assign a valid tenant or delete those rows, then re-run migrations.'; - END IF; - END - $guard$ - `, - sql` - ALTER TABLE "artifacts"."artifact" - ALTER COLUMN "tenant_id" SET NOT NULL - `, - sql` - ALTER TABLE "artifacts"."artifact" - ADD CONSTRAINT "artifact_version_gte_1" CHECK ("version" >= 1) - `, - sql` - ALTER TABLE "artifacts"."artifact_version" - ADD CONSTRAINT "artifact_version_version_gte_1" CHECK ("version" >= 1) - `, - sql` - ALTER TABLE "artifacts"."upload" - ADD CONSTRAINT "upload_size_gte_0" CHECK ("size" >= 0) - `, - sql` - ALTER TABLE "artifacts"."mail_attachment_ref" - ADD CONSTRAINT "mail_attachment_ref_size_gte_0" CHECK ("size" >= 0) - `, - ], - }, - { - // `metadata` is opaque-to-the-package jsonb: `artifact_version` carries it - // per version, and `artifact` mirrors the current version's value the same - // way it already mirrors `title`/`content`/`version`. `parent_version_ids` - // is explicit lineage set by the writer, never inferred from version - // order, so it lives only on `artifact_version` — there is no "current - // parents" concept to mirror onto `artifact`. Both columns are nullable - // and additive: `ADD COLUMN IF NOT EXISTS` is safe to re-run and an - // existing database adopts cleanly with no backfill. - id: "0004_version_metadata", - statements: [ - sql` - ALTER TABLE "artifacts"."artifact" - ADD COLUMN IF NOT EXISTS "metadata" jsonb - `, - sql` - ALTER TABLE "artifacts"."artifact_version" - ADD COLUMN IF NOT EXISTS "metadata" jsonb - `, - sql` - ALTER TABLE "artifacts"."artifact_version" - ADD COLUMN IF NOT EXISTS "parent_version_ids" text[] - `, - ], - }, - { - // `content_sha256` is computed at write time (sha256 hex over the UTF-8 - // bytes of `content` for text/URL artifacts, or over the uploaded bytes - // for a blob-backed file artifact's first version) and mirrored onto - // `artifact` the same way `metadata` already is. Both columns are - // nullable and additive: `ADD COLUMN IF NOT EXISTS` is safe to re-run, and - // existing rows stay null — "written before digests existed" — with no - // backfill. - id: "0005_version_content_digest", - statements: [ - sql` - ALTER TABLE "artifacts"."artifact" - ADD COLUMN IF NOT EXISTS "content_sha256" text - `, - sql` - ALTER TABLE "artifacts"."artifact_version" - ADD COLUMN IF NOT EXISTS "content_sha256" text - `, - ], - }, -]; +const MIGRATIONS_DIR = join(import.meta.dirname, "..", "migrations"); // Advisory locks are namespaced by this integer alone; deliberately arbitrary // and specific to @corbits/artifacts. const LOCK_KEY = 0x0a27_1f04; -const DIALECT = new PgDialect(); - -/** - * Checksum of a migration's BODY, rendered through the executing dialect, with - * whitespace runs collapsed so reindenting is not a schema change. Editing a - * shipped migration therefore fails loudly on the next boot instead of letting - * existing and fresh databases diverge silently. - */ -export function migrationChecksum(migration: Migration): string { - const rendered = migration.statements - .map((statement) => - DIALECT.sqlToQuery(statement).sql.replace(/\s+/g, " ").trim(), - ) - .join(";\n"); - return createHash("sha256").update(rendered).digest("hex"); -} - -export class MigrationChecksumError extends Error { - constructor(id: string, recorded: string, current: string) { - super( - `Migration "${id}" has changed since it was applied to this database ` + - `(recorded ${recorded}, now ${current}). A shipped migration must never ` + - `be edited — add a new one instead.`, - ); - this.name = "MigrationChecksumError"; - } +function quoteIdentifier(name: string): string { + return `"${name.replace(/"/g, '""')}"`; } /** - * Thrown when the migration ledger is empty but package-owned objects already - * exist in the `artifacts` schema. Without an explicit `{ adopt: true }`, the - * runner fails closed rather than letting `CREATE IF NOT EXISTS` no-op and - * stamp a checksum over a shape it never verified. With `adopt: true`, the - * same error is thrown when the live shape does not match what the migrations - * would create. + * Takes the same `config` and `schema` the host passes Interchange's + * `runMigrations`. `schema` is where the host's `tenant` and `principal` + * tables live; the package's own tables always go in the `artifacts` schema, + * and the `"public".` foreign-key references in the SQL are rewritten to + * `schema`. Every statement is idempotent, and the run is one transaction + * behind an advisory lock so concurrent hub replicas cannot race the DDL. */ -export class MigrationAdoptError extends Error { - constructor(message: string) { - super(message); - this.name = "MigrationAdoptError"; +export async function runArtifactMigrations( + config: DBConfig, + options: { schema: string }, +): Promise { + if (options.schema.length === 0) { + throw new Error("runArtifactMigrations: schema name must not be empty"); } -} - -/** - * Options for {@link runArtifactMigrations}. - * - * `adopt` is an operator escape hatch for the rare case where package tables - * already exist (restored dump, manual DDL, ledger dropped) and the operator - * has confirmed they match the expected shape. It is never set by default and - * must not be passed on ordinary boots. Adopt validates tables, column types, - * required nullability (`artifact.tenant_id`), and the named CHECK constraints - * from `0003_schema_invariants` before writing ledger rows — it does not - * re-run DDL. - */ -export type RunArtifactMigrationsOptions = { - adopt?: boolean; -}; - -/** Package-owned tables (excluding the ledger) and the columns each must have. */ -type ExpectedColumn = { - name: string; - udt: string; - /** When set, live `is_nullable` must be `NO`. */ - notNull?: true; -}; - -const EXPECTED_OWNED_SHAPE: Readonly> = - { - artifact: [ - { name: "id", udt: "text" }, - { name: "tenant_id", udt: "text", notNull: true }, - { name: "principal_id", udt: "text" }, - { name: "owner_principal_id", udt: "text" }, - { name: "kind", udt: "text" }, - { name: "title", udt: "text" }, - { name: "content", udt: "text" }, - { name: "source", udt: "jsonb" }, - { name: "version", udt: "int4" }, - { name: "archived_at", udt: "timestamptz" }, - { name: "created_at", udt: "timestamptz" }, - { name: "updated_at", udt: "timestamptz" }, - { name: "metadata", udt: "jsonb" }, - { name: "content_sha256", udt: "text" }, - ], - artifact_version: [ - { name: "id", udt: "text" }, - { name: "artifact_id", udt: "text" }, - { name: "version", udt: "int4" }, - { name: "title", udt: "text" }, - { name: "content", udt: "text" }, - { name: "author_id", udt: "text" }, - { name: "created_at", udt: "timestamptz" }, - { name: "metadata", udt: "jsonb" }, - { name: "parent_version_ids", udt: "_text" }, - { name: "content_sha256", udt: "text" }, - ], - upload: [ - { name: "id", udt: "text" }, - { name: "tenant_id", udt: "text" }, - { name: "principal_id", udt: "text" }, - { name: "filename", udt: "text" }, - { name: "mime_type", udt: "text" }, - { name: "content", udt: "bytea" }, - { name: "size", udt: "int4" }, - { name: "created_at", udt: "timestamptz" }, - ], - mail_attachment_ref: [ - { name: "id", udt: "text" }, - { name: "tenant_id", udt: "text" }, - { name: "principal_id", udt: "text" }, - { name: "instance_id", udt: "text" }, - { name: "mail_id", udt: "text" }, - { name: "artifact_id", udt: "text" }, - { name: "name", udt: "text" }, - { name: "mime_type", udt: "text" }, - { name: "size", udt: "int4" }, - { name: "created_at", udt: "timestamptz" }, - ], - }; - -/** - * Row-local CHECKs applied by `0003_schema_invariants`. Adopt must see these - * names so stamping 0003 without the constraints cannot pass shape validation. - */ -const EXPECTED_CHECK_CONSTRAINTS: ReadonlyArray<{ - table: string; - name: string; -}> = [ - { table: "artifact", name: "artifact_version_gte_1" }, - { table: "artifact_version", name: "artifact_version_version_gte_1" }, - { table: "upload", name: "upload_size_gte_0" }, - { table: "mail_attachment_ref", name: "mail_attachment_ref_size_gte_0" }, -]; - -async function listOwnedTables(tx: ArtifactTx): Promise { - const rows = await tx.execute<{ table_name: string }>(sql` - SELECT table_name FROM information_schema.tables - WHERE table_schema = ${ARTIFACTS_SCHEMA} - AND table_type = 'BASE TABLE' - AND table_name <> ${LEDGER_TABLE} - ORDER BY table_name - `); - return rows.map((row) => row.table_name); -} - -/** - * Compare live catalogue against the shape the migrations create: tables, - * column types, required nullability, and named CHECK constraints. Returns a - * human-readable mismatch list (empty when compatible). - */ -async function shapeMismatches(tx: ArtifactTx): Promise { - const owned = await listOwnedTables(tx); - const expectedTables = Object.keys(EXPECTED_OWNED_SHAPE).sort(); - const mismatches: string[] = []; - - const ownedSet = new Set(owned); - const expectedSet = new Set(expectedTables); - for (const table of expectedTables) { - if (!ownedSet.has(table)) { - mismatches.push(`missing table ${ARTIFACTS_SCHEMA}.${table}`); - } + const schemaIdent = quoteIdentifier(options.schema); + const files = (await readdir(MIGRATIONS_DIR)).filter((f) => f.endsWith(".sql")).sort(); + if (files.length === 0) { + throw new Error(`runArtifactMigrations: no .sql files found in ${MIGRATIONS_DIR}`); } - for (const table of owned) { - if (!expectedSet.has(table)) { - mismatches.push(`unexpected table ${ARTIFACTS_SCHEMA}.${table}`); + const statements: string[] = []; + for (const file of files) { + const raw = await readFile(join(MIGRATIONS_DIR, file), "utf8"); + for (const stmt of raw + .replace(/"public"\.(?=")/g, `${schemaIdent}.`) + .split("--> statement-breakpoint")) { + if (stmt.trim().length > 0) statements.push(stmt); } } - if (mismatches.length > 0) return mismatches; - - const columns = await tx.execute<{ - table_name: string; - column_name: string; - udt_name: string; - is_nullable: string; - }>(sql` - SELECT table_name, column_name, udt_name, is_nullable - FROM information_schema.columns - WHERE table_schema = ${ARTIFACTS_SCHEMA} - AND table_name <> ${LEDGER_TABLE} - `); - - const byTable = new Map>(); - for (const col of columns) { - let cols = byTable.get(col.table_name); - if (!cols) { - cols = new Map(); - byTable.set(col.table_name, cols); - } - cols.set(col.column_name, { - udt: col.udt_name, - nullable: col.is_nullable, + const clientOptions: postgres.Options<{}> = { + host: config.host, + port: config.port, + user: config.user, + password: config.password, + database: config.database, + max: 1, + onnotice: () => undefined, + }; + if (config.ssl !== undefined) clientOptions.ssl = config.ssl; + const client = postgres(clientOptions); + try { + await client.begin(async (tx) => { + await tx.unsafe(`SELECT pg_advisory_xact_lock(${LOCK_KEY})`); + for (const stmt of statements) await tx.unsafe(stmt); }); - } - - for (const table of expectedTables) { - const expectedCols = EXPECTED_OWNED_SHAPE[table]!; - const live = byTable.get(table) ?? new Map(); - for (const expected of expectedCols) { - const liveCol = live.get(expected.name); - if (liveCol === undefined) { - mismatches.push( - `missing column ${ARTIFACTS_SCHEMA}.${table}.${expected.name}`, - ); - } else if (liveCol.udt !== expected.udt) { - mismatches.push( - `column ${ARTIFACTS_SCHEMA}.${table}.${expected.name} has type ${liveCol.udt}, expected ${expected.udt}`, - ); - } else if (expected.notNull && liveCol.nullable !== "NO") { - mismatches.push( - `column ${ARTIFACTS_SCHEMA}.${table}.${expected.name} is nullable, expected NOT NULL`, - ); - } - } - } - - if (mismatches.length > 0) return mismatches; - - // Named CHECKs from 0003 — without these, adopt would stamp the ledger over a - // schema that never gained the row-local invariants. - const checks = await tx.execute<{ table_name: string; constraint_name: string }>( - sql` - SELECT c.relname AS table_name, con.conname AS constraint_name - FROM pg_catalog.pg_constraint con - JOIN pg_catalog.pg_class c ON c.oid = con.conrelid - JOIN pg_catalog.pg_namespace n ON n.oid = c.relnamespace - WHERE n.nspname = ${ARTIFACTS_SCHEMA} - AND con.contype = 'c' - `, - ); - const checkByTable = new Map>(); - for (const row of checks) { - let names = checkByTable.get(row.table_name); - if (!names) { - names = new Set(); - checkByTable.set(row.table_name, names); - } - names.add(row.constraint_name); - } - for (const { table, name } of EXPECTED_CHECK_CONSTRAINTS) { - const live = checkByTable.get(table); - if (!live?.has(name)) { - mismatches.push( - `missing CHECK constraint ${name} on ${ARTIFACTS_SCHEMA}.${table}`, - ); - } - } - - return mismatches; -} - -/** - * Idempotent migration runner. Safe to call on every boot, from every - * instance: the whole run is ONE transaction whose first act is taking a - * TRANSACTION-scoped advisory lock, so concurrent cold starts serialize on the - * same session and the lock releases on commit or rollback. `SET LOCAL - * client_min_messages = warning` silences the re-boot NOTICE chatter from the - * IF NOT EXISTS statements without muting real warnings. Each migration - * applies inside a savepoint together with its ledger row, so it can never be - * recorded as applied with only some statements run. - * - * When the ledger is empty but package-owned objects already exist, the runner - * fails closed with {@link MigrationAdoptError} unless `{ adopt: true }` is - * passed and the live shape matches what the migrations would create (tables, - * column types, required nullability, and named CHECK constraints). That path - * records checksums without re-running DDL. Ledger checksum drift still throws - * {@link MigrationChecksumError}. - */ -export async function runArtifactMigrations( - db: ArtifactDb, - options: RunArtifactMigrationsOptions = {}, -): Promise { - const adopt = options.adopt === true; - - await db.transaction(async (tx) => { - await tx.execute(sql`SET LOCAL client_min_messages = warning`); - await tx.execute(sql`SELECT pg_advisory_xact_lock(${LOCK_KEY})`); - - await tx.execute( - sql`CREATE SCHEMA IF NOT EXISTS ${sql.identifier(ARTIFACTS_SCHEMA)}`, + } catch (error) { + throw new Error( + `@corbits/artifacts migration failed: ${error instanceof Error ? error.message : String(error)}`, + { cause: error }, ); - - await tx.execute(sql` - CREATE TABLE IF NOT EXISTS ${LEDGER} ( - "id" text PRIMARY KEY, - "checksum" text NOT NULL, - "applied_at" timestamptz NOT NULL DEFAULT now() - ) - `); - - const applied = await tx.execute<{ id: string; checksum: string }>( - sql`SELECT "id", "checksum" FROM ${LEDGER}`, - ); - const appliedChecksums = new Map(applied.map((row) => [row.id, row.checksum])); - - // Empty ledger + pre-existing owned objects: never let IF NOT EXISTS no-op - // and stamp a checksum. Fail closed, or adopt only after shape validation. - if (appliedChecksums.size === 0) { - const owned = await listOwnedTables(tx); - if (owned.length > 0) { - const mismatches = await shapeMismatches(tx); - if (mismatches.length > 0) { - throw new MigrationAdoptError( - `Package schema ${ARTIFACTS_SCHEMA} already has objects but the ` + - `migration ledger is empty, and the live shape is incompatible: ` + - `${mismatches.join("; ")}. Refusing to record a checksum over an ` + - `unverified schema. Drop the incompatible objects and re-run, or ` + - `repair the shape to match the package migrations before adopting.`, - ); - } - if (!adopt) { - throw new MigrationAdoptError( - `Package schema ${ARTIFACTS_SCHEMA} already has objects ` + - `(${owned.join(", ")}) but the migration ledger is empty. ` + - `Refusing to silently adopt them. Re-run with ` + - `{ adopt: true } only after confirming the live shape matches ` + - `what the package migrations create, or drop the schema and let ` + - `the runner create it cleanly.`, - ); - } - // Compatible shape + explicit adopt: record every migration checksum - // without re-running DDL (IF NOT EXISTS would only hide drift). - for (const migration of MIGRATIONS) { - const checksum = migrationChecksum(migration); - await tx.execute(sql` - INSERT INTO ${LEDGER} ("id", "checksum") - VALUES (${migration.id}, ${checksum}) - `); - } - return; - } - } - - for (const migration of MIGRATIONS) { - const checksum = migrationChecksum(migration); - const recorded = appliedChecksums.get(migration.id); - if (recorded !== undefined) { - if (recorded !== checksum) { - throw new MigrationChecksumError(migration.id, recorded, checksum); - } - continue; - } - await tx.transaction(async (step) => { - for (const statement of migration.statements) { - await step.execute(statement); - } - await step.execute(sql` - INSERT INTO ${LEDGER} ("id", "checksum") - VALUES (${migration.id}, ${checksum}) - `); - }); - } - }); + } finally { + await client.end({ timeout: 5 }); + } } diff --git a/src/test-helpers.ts b/src/test-helpers.ts index 6406b4d..03a6b02 100644 --- a/src/test-helpers.ts +++ b/src/test-helpers.ts @@ -1,4 +1,5 @@ import { sql } from "drizzle-orm"; +import type { DBConfig } from "@intx/db"; import { createArtifactDb, type ArtifactDb } from "../src/db.js"; import { runArtifactMigrations } from "../src/migrations.js"; import { createArtifact } from "../src/artifacts.js"; @@ -8,6 +9,18 @@ export const DATABASE_URL = process.env.ARTIFACT_DATABASE_URL ?? "postgres://postgres:postgres@localhost:5457/artifact_core"; +/** `DATABASE_URL` in the shape Interchange's `runMigrations` takes. */ +export function databaseConfig(connectionString: string): DBConfig { + const url = new URL(connectionString); + return { + host: url.hostname, + port: Number(url.port || 5432), + user: decodeURIComponent(url.username), + password: decodeURIComponent(url.password), + database: databaseNameFromConnectionString(connectionString), + }; +} + /** * Explicit opt-in required before the harness runs TRUNCATE or DROP SCHEMA. * Must be the string `"1"` — any other value (including `"true"`) is refused. @@ -150,7 +163,7 @@ export async function testDb(): Promise { db = createArtifactDb(DATABASE_URL).db; shared = db; await ensureControlPlane(db); - await runArtifactMigrations(db); + await runArtifactMigrations(databaseConfig(DATABASE_URL), { schema: "public" }); } await db.execute( sql`TRUNCATE TABLE "artifacts"."artifact", "artifacts"."artifact_version", "artifacts"."upload", "artifacts"."mail_attachment_ref" CASCADE`, From 00fce3a327be425f579c54245724ebab8adb4ac0 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 25 Sep 2026 07:24:45 -0700 Subject: [PATCH 2/2] test(migrations): close the database handle after the run --- src/migrations.test.ts | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/src/migrations.test.ts b/src/migrations.test.ts index 1d0572f..cee514c 100644 --- a/src/migrations.test.ts +++ b/src/migrations.test.ts @@ -1,4 +1,4 @@ -import { beforeAll, describe, expect, test } from "bun:test"; +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; import { getTableName, is, sql } from "drizzle-orm"; import { PgTable } from "drizzle-orm/pg-core"; import { createArtifactDb } from "./db.js"; @@ -23,7 +23,7 @@ const DECLARED_TABLES = (Object.values(schema) as unknown[]) .map(getTableName) .sort(); -const { db } = createArtifactDb(DATABASE_URL); +const { db, close } = createArtifactDb(DATABASE_URL); async function packageTables(): Promise { const rows = await db.execute<{ table_name: string }>(sql` @@ -45,6 +45,8 @@ beforeAll(async () => { await ensureControlPlane(db); }); +afterAll(close); + describe("runArtifactMigrations", () => { test("creates exactly the tables schema.ts declares, and re-running is a no-op", async () => { await db.execute(sql`DROP SCHEMA IF EXISTS ${sql.identifier(SCHEMA)} CASCADE`);