diff --git a/bun.lock b/bun.lock index e94fd8011..48e18e22e 100644 --- a/bun.lock +++ b/bun.lock @@ -1185,6 +1185,7 @@ "version": "0.0.1", "dependencies": { "@corbits/folded-run-one-shot": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@corbits/slug": "workspace:*", "@corbits/workflow-catalog": "workspace:*", "@intx/db": "workspace:*", @@ -1468,6 +1469,7 @@ "dependencies": { "@corbits/error-sink": "workspace:*", "@corbits/folded-runs": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", "@intx/hub-common": "0.3.0", @@ -1513,6 +1515,7 @@ "version": "0.0.1", "dependencies": { "@corbits/error-sink": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/hub-sessions": "workspace:*", "@intx/types": "workspace:*", "arktype": "catalog:", @@ -3587,18 +3590,8 @@ "@babel/helper-compilation-targets/semver": ["semver@6.3.1", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-BR7VvDCVHO+q2xBEWskxS6DJE1qRnb7DxzUrogb71CWoSficBxYsiAGd+Kl0mmq/MprG9yArRkyrQxTO6XjMzA=="], - "@corbits/artifact-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], - - "@corbits/chat-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], - - "@corbits/context-menu/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], - "@corbits/memory-hub/@corbits/memory": ["@corbits/memory@github:corbitsdev/corbits-memory#9e6f213", { "dependencies": { "@intx/agent": "0.2.2", "@intx/authz": "0.2.2", "@intx/hub-api": "0.2.2", "@intx/log": "0.2.2", "@intx/workflow": "0.2.2", "arktype": "^2.1.29", "drizzle-orm": "^0.45.1", "hono": "^4.9.0", "hono-openapi": "^1.3.1", "postgres": "^3.4.7" } }, "corbitsdev-corbits-memory-9e6f213", "sha512-utnM4ZT2zmslcPXYWAAqxlDNLcpGsXFiTOtj8h7+OXnhCP0Eaw8yl25+yCTyHpvt3jcdeG4h5uFsSj7ou0BZCA=="], - "@corbits/plugins-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], - - "@corbits/settings-ui/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], - "@esbuild-kit/core-utils/esbuild": ["esbuild@0.18.20", "", { "optionalDependencies": { "@esbuild/android-arm": "0.18.20", "@esbuild/android-arm64": "0.18.20", "@esbuild/android-x64": "0.18.20", "@esbuild/darwin-arm64": "0.18.20", "@esbuild/darwin-x64": "0.18.20", "@esbuild/freebsd-arm64": "0.18.20", "@esbuild/freebsd-x64": "0.18.20", "@esbuild/linux-arm": "0.18.20", "@esbuild/linux-arm64": "0.18.20", "@esbuild/linux-ia32": "0.18.20", "@esbuild/linux-loong64": "0.18.20", "@esbuild/linux-mips64el": "0.18.20", "@esbuild/linux-ppc64": "0.18.20", "@esbuild/linux-riscv64": "0.18.20", "@esbuild/linux-s390x": "0.18.20", "@esbuild/linux-x64": "0.18.20", "@esbuild/netbsd-x64": "0.18.20", "@esbuild/openbsd-x64": "0.18.20", "@esbuild/sunos-x64": "0.18.20", "@esbuild/win32-arm64": "0.18.20", "@esbuild/win32-ia32": "0.18.20", "@esbuild/win32-x64": "0.18.20" }, "bin": { "esbuild": "bin/esbuild" } }, "sha512-ceqxoedUrcayh7Y7ZX6NdbbDzGROiyVBgC4PriJThBKSVPWnnFHZAkfI1lJT8QFkOwH4qOS2SJkS4wvpGl8BpA=="], "@eslint-community/eslint-utils/eslint-visitor-keys": ["eslint-visitor-keys@3.4.3", "", {}, "sha512-wpc+LXeiyiisxPlEkUzU6svyS1frIO3Mgxj1fdy7Pm8Ygzguax2N3Fa/D/ag1WqbOprdI+uY6wMUl8/a2G+iag=="], @@ -3623,8 +3616,6 @@ "@workbench/hub/@corbits/memory": ["@corbits/memory@github:corbitsdev/corbits-memory#9e6f213", { "dependencies": { "@intx/agent": "0.2.2", "@intx/authz": "0.2.2", "@intx/hub-api": "0.2.2", "@intx/log": "0.2.2", "@intx/workflow": "0.2.2", "arktype": "^2.1.29", "drizzle-orm": "^0.45.1", "hono": "^4.9.0", "hono-openapi": "^1.3.1", "postgres": "^3.4.7" } }, "corbitsdev-corbits-memory-9e6f213", "sha512-utnM4ZT2zmslcPXYWAAqxlDNLcpGsXFiTOtj8h7+OXnhCP0Eaw8yl25+yCTyHpvt3jcdeG4h5uFsSj7ou0BZCA=="], - "@workbench/web/@corbits/react-ui": ["@corbits/react-ui@github:corbitsdev/react-ui#3b12281", { "dependencies": { "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "tailwind-merge": "^3.3.1" }, "peerDependencies": { "@radix-ui/react-dialog": "^1.1.15", "@radix-ui/react-dropdown-menu": "^2.1.16", "@radix-ui/react-slot": "^1.2.3", "@radix-ui/react-tooltip": "^1.2.8", "@tanstack/react-query": "^5.90.2", "lucide-react": "^0.545.0 || ^1.0.0", "react": "^18.2.0 || ^19.0.0", "react-dom": "^18.2.0 || ^19.0.0", "sonner": "^2.0.7" }, "optionalPeers": ["@tanstack/react-query"] }, "corbitsdev-react-ui-3b12281", "sha512-Abvm/DO0Gqg0ITHGT9355ZxyKRPMVJLSSQSjpd3a8qt4JPrSMOLIOS4sX8ZMNNaArIbnY9F+VKrOWkUJUyO4Nw=="], - "ajv-formats/ajv": ["ajv@8.20.0", "", { "dependencies": { "fast-deep-equal": "^3.1.3", "fast-uri": "^3.0.1", "json-schema-traverse": "^1.0.0", "require-from-string": "^2.0.2" } }, "sha512-Thbli+OlOj+iMPYFBVBfJ3OmCAnaSyNn4M1vz9T6Gka5Jt9ba/HIR56joy65tY6kx/FCF5VXNB819Y7/GUrBGA=="], "better-call/@better-auth/utils": ["@better-auth/utils@0.5.0", "", { "dependencies": { "@noble/hashes": "^2.0.1" } }, "sha512-BL8W4EfIZFwlu0r54m3v1ztjDhu6dDe/amLTm0xybmbZaNgYUqhD3SjpAsnq0q8YD6/ki4iwIgxJNLP/N3TxiA=="], diff --git a/docs/package-migrations.md b/docs/package-migrations.md index 5dada978c..54406c44b 100644 --- a/docs/package-migrations.md +++ b/docs/package-migrations.md @@ -72,7 +72,8 @@ package onto the transactional pattern — see below): 1. **Self-contained, transactional** — `@corbits/chat`, `@corbits/notify`, `@corbits/webhook-triggers`, `@corbits/routines`, `@corbits/insights`, `@corbits/skills`, `@corbits/bench`, `@corbits/preferences`, - `@corbits/inference-catalog`, `@corbits/evals`, `@corbits/access-policy`. + `@corbits/inference-catalog`, `@corbits/evals`, `@corbits/access-policy`, + `@corbits/workflow-deploy-source`. The package's `src/migrations.ts` owns only a literal `{ name, sql }[]` array and a thin `applyXMigrations(databaseUrl)` wrapper; the mechanics — schema/ledger bootstrap, the transactional apply loop, and the advisory diff --git a/packages/migration-runner/README.md b/packages/migration-runner/README.md index 2b844a47c..6dfe3ad8f 100644 --- a/packages/migration-runner/README.md +++ b/packages/migration-runner/README.md @@ -25,11 +25,13 @@ name rather than the name itself, so two distinct table names could in principle hash to the same key — two packages would then serialize their boot-time migrations against each other instead of running in parallel, a liveness cost (one waits its turn) and never a correctness one (each still -applies to its own schema and ledger). The six current ledger tables +applies to its own schema and ledger). The nine current ledger tables (`access_policy_migrations`, `bench_migrations`, `evals_migrations`, `inference_catalog_migrations`, `insights_migrations`, -`preferences_migrations`) do not collide — verified against a live -`hashtext()`. Confirm the same before naming a seventh. +`preferences_migrations`, `webhook_triggers_migrations`, +`routine_migrations`, `workflow_deploy_source_migrations`) do not collide — +verified pairwise against a live PostgreSQL 17.11 `hashtext()`. Confirm the +same before naming a tenth. ## What it does not change diff --git a/packages/routines/package.json b/packages/routines/package.json index 282e85a42..5f33664b3 100644 --- a/packages/routines/package.json +++ b/packages/routines/package.json @@ -17,6 +17,7 @@ }, "dependencies": { "@corbits/folded-run-one-shot": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@corbits/slug": "workspace:*", "cronstrue": "^3.24.0", "@corbits/workflow-catalog": "workspace:*", diff --git a/packages/routines/src/migrations.ts b/packages/routines/src/migrations.ts index 62f45717a..3df66c0f2 100644 --- a/packages/routines/src/migrations.ts +++ b/packages/routines/src/migrations.ts @@ -4,13 +4,17 @@ // this package's half of the "mount + migrations is the entire install // story" install contract. Bookkeeping is its own ledger table, never // the platform's drizzle journal, so this package's migration history -// stays extractable on its own. -import postgres from "postgres"; +// stays extractable on its own. Mechanics (schema/ledger bootstrap, +// transactional apply, the advisory lock across concurrent hub +// replicas) live in @corbits/migration-runner — this file owns only +// the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; -export interface RoutineMigration { - name: string; - sql: string; -} +export type RoutineMigration = PackageMigration; export const routineMigrations: readonly RoutineMigration[] = [ { @@ -119,18 +123,7 @@ export const routineMigrations: readonly RoutineMigration[] = [ const SCHEMA = "routines"; const LEDGER_TABLE = "routine_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyRoutineMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyRoutineMigrationsReport = ApplyPackageMigrationsReport; /** * Apply `routineMigrations` against `databaseUrl`, idempotently: a @@ -141,40 +134,11 @@ export interface ApplyRoutineMigrationsReport { export async function applyRoutineMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - const rows = await sql.unsafe( - `SELECT name FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)}`, - ); - const alreadyApplied = new Set(rows.map((row) => String(row["name"]))); - const applied: string[] = []; - for (const migration of routineMigrations) { - if (alreadyApplied.has(migration.name)) continue; - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (error) { - throw new Error( - `@corbits/routines migration ${JSON.stringify(migration.name)} failed: ` + - `${error instanceof Error ? error.message : String(error)}`, - { cause: error }, - ); - } - } - return { applied, alreadyApplied: [...alreadyApplied] }; - } finally { - await sql.end(); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: routineMigrations, + packageLabel: "@corbits/routines", + }); } diff --git a/packages/routines/test/migrations.test.ts b/packages/routines/test/migrations.test.ts index f14fd8c0d..21a9a28d9 100644 --- a/packages/routines/test/migrations.test.ts +++ b/packages/routines/test/migrations.test.ts @@ -126,3 +126,71 @@ describeIfDb("applyRoutineMigrations", () => { } }); }); + +// Separate database from the suites above: two replicas racing the same +// ledger must not collide with the idempotency test's own already-applied +// rows, and must start from a schema that has never seen this migration +// set before. +describeIfDb("applyRoutineMigrations concurrency", () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ).replace("_routine_migrations_test", "_routine_migrations_concurrent_test"); + const scratchDatabase = new URL(scratchUrl).pathname.replace(/^\//, ""); + + beforeAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }, 20000); + + afterAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }, 20000); + + test("two replicas booting concurrently both complete without either crashing on a duplicate ledger insert", async () => { + const [first, second] = await Promise.all([ + applyRoutineMigrations(scratchUrl), + applyRoutineMigrations(scratchUrl), + ]); + + const appliedNames = [...first.applied, ...second.applied].sort(); + expect(new Set(appliedNames).size).toBe(appliedNames.length); + + const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + try { + const ledgerRows = await sql.unsafe( + `SELECT name FROM "routines"."routine_migrations" ORDER BY name`, + ); + const ledgerNames = ledgerRows.map((row) => String(row["name"])); + expect(new Set(ledgerNames).size).toBe(ledgerNames.length); + expect( + [ + ...appliedNames, + ...first.alreadyApplied, + ...second.alreadyApplied, + ].sort(), + ).toEqual([...ledgerNames, ...ledgerNames].sort()); + } finally { + await sql.end(); + } + }, 10000); +}); diff --git a/packages/webhook-triggers/package.json b/packages/webhook-triggers/package.json index d3e8234f7..bc1fe4620 100644 --- a/packages/webhook-triggers/package.json +++ b/packages/webhook-triggers/package.json @@ -16,6 +16,7 @@ "dependencies": { "@corbits/error-sink": "workspace:*", "@corbits/folded-runs": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/db": "workspace:*", "@intx/hub-api": "workspace:*", "@intx/hub-common": "0.3.0", diff --git a/packages/webhook-triggers/src/migrations.ts b/packages/webhook-triggers/src/migrations.ts index ff12bc634..d2765e4eb 100644 --- a/packages/webhook-triggers/src/migrations.ts +++ b/packages/webhook-triggers/src/migrations.ts @@ -5,13 +5,17 @@ // install story, mirroring `@corbits/chat`'s `migrations.ts`. // Bookkeeping is deliberately its own table, never the platform's // drizzle journal, so this package's migration history stays -// extractable on its own. -import postgres from "postgres"; +// extractable on its own. Mechanics (schema/ledger bootstrap, +// transactional apply, the advisory lock across concurrent hub +// replicas) live in @corbits/migration-runner — this file owns only +// the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; -export interface WebhookTriggersMigration { - name: string; - sql: string; -} +export type WebhookTriggersMigration = PackageMigration; /** * Ordered, explicit migration set for the table declared in @@ -47,26 +51,15 @@ export const webhookTriggersMigrations: readonly WebhookTriggersMigration[] = [ }, ]; -// Bookkeeping table for this package's own migrations. Named -// distinctly from the platform's setup ledger and from any drizzle -// journal, so extracting `@corbits/webhook-triggers` out of this repo -// never has to disentangle its history from the platform's. Lives in -// the package's own `webhook_triggers` schema, like the table it owns. +// Named distinctly from the platform's setup ledger and from any +// drizzle journal, so extracting `@corbits/webhook-triggers` out of +// this repo never has to disentangle its history from the platform's. +// Lives in the package's own `webhook_triggers` schema, like the +// table it owns. const SCHEMA = "webhook_triggers"; const LEDGER_TABLE = "webhook_triggers_migrations"; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - -export interface ApplyWebhookTriggersMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyWebhookTriggersMigrationsReport = ApplyPackageMigrationsReport; /** * Apply `webhookTriggersMigrations` against `databaseUrl`, @@ -78,40 +71,11 @@ export interface ApplyWebhookTriggersMigrationsReport { export async function applyWebhookTriggersMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - const rows = await sql.unsafe( - `SELECT name FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)}`, - ); - const alreadyApplied = new Set(rows.map((row) => String(row["name"]))); - const applied: string[] = []; - for (const migration of webhookTriggersMigrations) { - if (alreadyApplied.has(migration.name)) continue; - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (error) { - throw new Error( - `@corbits/webhook-triggers migration ${JSON.stringify(migration.name)} failed: ` + - `${error instanceof Error ? error.message : String(error)}`, - { cause: error }, - ); - } - } - return { applied, alreadyApplied: [...alreadyApplied] }; - } finally { - await sql.end(); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: webhookTriggersMigrations, + packageLabel: "@corbits/webhook-triggers", + }); } diff --git a/packages/webhook-triggers/test/migrations.test.ts b/packages/webhook-triggers/test/migrations.test.ts index affa5e34a..addd51948 100644 --- a/packages/webhook-triggers/test/migrations.test.ts +++ b/packages/webhook-triggers/test/migrations.test.ts @@ -89,3 +89,80 @@ describeIfDb("applyWebhookTriggersMigrations", () => { } }); }); + +// Separate database from the suite above: two replicas racing the same +// ledger must not collide with the idempotency test's own already-applied +// rows, and must start from a schema that has never seen this migration +// set before. +describeIfDb("applyWebhookTriggersMigrations concurrency", () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ).replace( + "_webhook_triggers_migrations_test", + "_webhook_triggers_migrations_concurrent_test", + ); + const scratchDatabase = new URL(scratchUrl).pathname.replace(/^\//, ""); + + beforeAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }, 20000); + + afterAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + }, 20000); + + test("two replicas booting concurrently both complete without either crashing on a duplicate ledger insert", async () => { + const [first, second] = await Promise.all([ + applyWebhookTriggersMigrations(scratchUrl), + applyWebhookTriggersMigrations(scratchUrl), + ]); + + const appliedNames = [...first.applied, ...second.applied].sort(); + expect(new Set(appliedNames).size).toBe(appliedNames.length); + expect( + [ + ...appliedNames, + ...first.alreadyApplied, + ...second.alreadyApplied, + ].sort(), + ).toEqual( + ["0001_webhook_trigger", "0002_webhook_trigger_tenant_index"] + .flatMap((name) => [name, name]) + .sort(), + ); + + const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + try { + const ledgerRows = await sql.unsafe( + `SELECT name FROM "webhook_triggers"."webhook_triggers_migrations" ORDER BY name`, + ); + expect(ledgerRows.map((row) => String(row["name"]))).toEqual([ + "0001_webhook_trigger", + "0002_webhook_trigger_tenant_index", + ]); + } finally { + await sql.end(); + } + }, 10000); +}); diff --git a/packages/workflow-deploy-source/package.json b/packages/workflow-deploy-source/package.json index aa3eec059..a1fb9111b 100644 --- a/packages/workflow-deploy-source/package.json +++ b/packages/workflow-deploy-source/package.json @@ -15,6 +15,7 @@ }, "dependencies": { "@corbits/error-sink": "workspace:*", + "@corbits/migration-runner": "workspace:*", "@intx/hub-sessions": "workspace:*", "@intx/types": "workspace:*", "arktype": "catalog:", diff --git a/packages/workflow-deploy-source/src/migrations.ts b/packages/workflow-deploy-source/src/migrations.ts index c5739b6fd..38c2a7aea 100644 --- a/packages/workflow-deploy-source/src/migrations.ts +++ b/packages/workflow-deploy-source/src/migrations.ts @@ -3,12 +3,16 @@ // without disentangling history from the platform drizzle journal. The // table this package owns lives in its own `workflow_deploy_source` // Postgres schema, never `public`; see docs/package-migrations.md. -import postgres from "postgres"; +// Mechanics (schema/ledger bootstrap, transactional apply, the +// advisory lock across concurrent hub replicas) live in +// @corbits/migration-runner — this file owns only the domain SQL. +import { + applyPackageMigrations, + type ApplyPackageMigrationsReport, + type PackageMigration, +} from "@corbits/migration-runner"; -export interface WorkflowDeploySourceMigration { - name: string; - sql: string; -} +export type WorkflowDeploySourceMigration = PackageMigration; const SCHEMA = "workflow_deploy_source"; @@ -35,65 +39,19 @@ export const workflowDeploySourceMigrations: readonly WorkflowDeploySourceMigrat }, ]; -function quoteIdentifier(name: string): string { - return `"${name.replace(/"/g, '""')}"`; -} - -function quoteQualified(schema: string, name: string): string { - return `${quoteIdentifier(schema)}.${quoteIdentifier(name)}`; -} - const LEDGER_TABLE = "workflow_deploy_source_migrations"; -export interface ApplyWorkflowDeploySourceMigrationsReport { - applied: string[]; - alreadyApplied: string[]; -} +export type ApplyWorkflowDeploySourceMigrationsReport = + ApplyPackageMigrationsReport; export async function applyWorkflowDeploySourceMigrations( databaseUrl: string, ): Promise { - const sql = postgres(databaseUrl, { max: 1, onnotice: () => undefined }); - try { - await sql.unsafe(`CREATE SCHEMA IF NOT EXISTS ${quoteIdentifier(SCHEMA)}`); - - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS ${quoteQualified(SCHEMA, LEDGER_TABLE)} (` + - `name text PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())`, - ); - - const applied: string[] = []; - const alreadyApplied: string[] = []; - - for (const migration of workflowDeploySourceMigrations) { - const existing = await sql.unsafe( - `SELECT 1 FROM ${quoteQualified(SCHEMA, LEDGER_TABLE)} WHERE name = $1`, - [migration.name], - ); - if (existing.length > 0) { - alreadyApplied.push(migration.name); - continue; - } - try { - await sql.begin(async (tx) => { - await tx.unsafe(migration.sql); - await tx.unsafe( - `INSERT INTO ${quoteQualified(SCHEMA, LEDGER_TABLE)} (name) VALUES ($1)`, - [migration.name], - ); - }); - applied.push(migration.name); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - throw new Error( - `workflow_deploy_source migration ${migration.name} failed: ${message}`, - { cause: err }, - ); - } - } - - return { applied, alreadyApplied }; - } finally { - await sql.end({ timeout: 5 }); - } + return applyPackageMigrations({ + databaseUrl, + schema: SCHEMA, + ledgerTable: LEDGER_TABLE, + migrations: workflowDeploySourceMigrations, + packageLabel: "workflow_deploy_source", + }); } diff --git a/packages/workflow-deploy-source/test/migrations.test.ts b/packages/workflow-deploy-source/test/migrations.test.ts index 3c5babe74..284a93805 100644 --- a/packages/workflow-deploy-source/test/migrations.test.ts +++ b/packages/workflow-deploy-source/test/migrations.test.ts @@ -77,3 +77,73 @@ describeIfDb("applyWorkflowDeploySourceMigrations", () => { expect(report.alreadyApplied).toEqual(["0001_workflow_deploy_source"]); }); }); + +// Separate database from the suites above: two replicas racing the same +// ledger must not collide with the earlier tests' own already-applied +// rows, and must start from a schema that has never seen this migration +// set before. +describeIfDb("applyWorkflowDeploySourceMigrations concurrency", () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ).replace( + "_workflow_deploy_source_migrations_test", + "_workflow_deploy_source_migrations_concurrent_test", + ); + const scratchDatabase = new URL(scratchUrl).pathname.replace(/^\//, ""); + + async function withMaintenance(run: (sql: postgres.Sql) => Promise) { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await run(maintenance); + } finally { + await maintenance.end(); + } + } + + beforeAll(async () => { + await withMaintenance(async (sql) => { + await sql.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + await sql.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + }); + }, 20000); + + afterAll(async () => { + await withMaintenance(async (sql) => { + await sql.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`); + }); + }, 20000); + + test("two replicas booting concurrently both complete without either crashing on a duplicate ledger insert", async () => { + const [first, second] = await Promise.all([ + applyWorkflowDeploySourceMigrations(scratchUrl), + applyWorkflowDeploySourceMigrations(scratchUrl), + ]); + + const appliedNames = [...first.applied, ...second.applied].sort(); + expect(new Set(appliedNames).size).toBe(appliedNames.length); + expect( + [ + ...appliedNames, + ...first.alreadyApplied, + ...second.alreadyApplied, + ].sort(), + ).toEqual(["0001_workflow_deploy_source", "0001_workflow_deploy_source"]); + + const sql = postgres(scratchUrl, { max: 1, onnotice: () => undefined }); + try { + const ledgerRows = await sql.unsafe( + `SELECT name FROM "workflow_deploy_source"."workflow_deploy_source_migrations" ORDER BY name`, + ); + expect(ledgerRows.map((row) => String(row["name"]))).toEqual([ + "0001_workflow_deploy_source", + ]); + } finally { + await sql.end(); + } + }, 10000); +});