From a867f837a4e99c7b06a695348cd30c09388d5192 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 24 Sep 2026 23:19:00 -0700 Subject: [PATCH] feat(migrations)!: align runMemoryMigrations with runMigrations runMemoryMigrations(config, { schema, ftsLanguage }) takes the same DBConfig and schema as Interchange runMigrations. schema names the host schema holding the tenant and principal tables, and the foreign-key references in the SQL are rewritten to it; memory's own tables stay in the memory schema. The runner replays every migration file on each run instead of keeping a memory._migrations ledger, which a new migration drops. Foreign keys and CHECK constraints are added only when missing, so a replay does not re-validate tables, a lock timeout keeps a replay from stalling readers, and the one-time temporal_class backfill runs only in the replay that adds the column. ftsLanguage is required, with no default, and the runner no longer reads FTS_LANGUAGE from the environment. --- CHANGELOG.md | 8 ++ CONTRIBUTING.md | 7 + IMPLEMENTATION.md | 12 +- PRODUCT.md | 2 +- README.md | 15 +- bun.lock | 2 + migrations/0002_memory_baseline.sql | 4 +- migrations/0003_claim_bearing.sql | 17 ++- migrations/0004_temporal_model.sql | 41 ++++-- migrations/0007_retention.sql | 17 ++- migrations/0008_tenant_principal_fks.sql | 152 +++++++++++++++++---- migrations/0009_drop_migrations_ledger.sql | 2 + package.json | 2 + scripts/db-setup.ts | 21 ++- src/core/fts-language.test.ts | 35 +++-- src/core/fts-language.ts | 11 +- src/migrations.ts | 99 ++++++++------ 17 files changed, 307 insertions(+), 140 deletions(-) create mode 100644 migrations/0009_drop_migrations_ledger.sql diff --git a/CHANGELOG.md b/CHANGELOG.md index 3db258e..fa175dc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 embed model registry, degrade metrics, FTS helpers), the test fakes, and `resolveGrantConfig` are no longer exported. The distiller stays at `@corbits/memory/distiller` and migrations at `@corbits/memory/migrations`. +- `runMemoryMigrations(config, { schema, ftsLanguage })` takes the same + `DBConfig` as Interchange `runMigrations` instead of a database URL. + `schema` names the host schema holding Interchange's `tenant` and + `principal` tables (the value passed to `runMigrations`, e.g. `"public"`); + memory's tables stay in the `memory` schema. `ftsLanguage` is required, and + the runner no longer reads `FTS_LANGUAGE` from the environment or accepts a `log` option. + Every migration file is idempotent and replayed on each run. The + `memory._migrations` ledger is dropped. ### Added diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 8a40ad4..fb6e9bd 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -44,6 +44,13 @@ bun run typecheck && bun run test `bun run typecheck` (`tsc --noEmit`) must be clean before any commit. +## Migrations + +`runMemoryMigrations` replays every file in `migrations/` on each run, so +every file must be idempotent. Never change a shipped file's effect on an +existing database: a changed constraint or a new column goes in a new +numbered file, because a guarded `ADD CONSTRAINT` keeps the old definition. + ## Branch and PR conventions - Branch off `main`; open PRs against `main`. diff --git a/IMPLEMENTATION.md b/IMPLEMENTATION.md index 7a5ddbb..26668e4 100644 --- a/IMPLEMENTATION.md +++ b/IMPLEMENTATION.md @@ -20,7 +20,7 @@ src/ http-client.ts # host-side HTTP client for tenant routes (imperative distill tick) log.ts # getLogger(["memory"]) from @intx/log - migrations.ts # runMemoryMigrations(url) + migrations.ts # runMemoryMigrations(dbConfig, { schema, ftsLanguage }) ports/ # DocumentStore / SourceProvider + fakes routes/ # the mounted tenant routes mount.ts # createMemoryRoutes (HTTP sub-app) @@ -49,7 +49,7 @@ src/ # @corbits/supermemory-memory-adapter → github.com/corbitsdev/corbits-supermemory-memory-adapter # @corbits/linear-tools → github.com/corbitsdev/corbits-linear-tools migrations/ # pgvector schema, applied in filename order by scripts/db-setup.ts -scripts/db-setup.ts # idempotent migration runner, tracked in `_migrations` +scripts/db-setup.ts # runs the idempotent migrations against DATABASE_URL compose.yml # pgvector + Ollama + reranker for local dev ``` @@ -150,8 +150,8 @@ egress control front the endpoints with an allowlisting proxy. ## Data model (`src/db/schema.ts` + `migrations/*.sql`) All tables are Drizzle-defined in `db/schema.ts`, DDL'd in `migrations/` -(applied by `scripts/db-setup.ts`, tracked in a `_migrations` ledger table so -re-running is a no-op). No memory table has a foreign key into any +(applied by `runMemoryMigrations`; every file is idempotent and replayed on +each run, so there is no ledger). No memory table has a foreign key into any control-plane table — `tenant_id`/`principal_id`/source refs are plain `text`. ### `memory_document` @@ -206,8 +206,8 @@ fresh full insert of its own chunks. **Changing `FTS_LANGUAGE` on an already-migrated database** (the mismatch `verifyFtsLanguage` throws on) requires rebuilding the generated column — -`runMemoryMigrations` only applies new files and will not retroactively -alter an existing one. One-time recipe (verified against a live +`runMemoryMigrations` adds the column only if it is missing and will not +retroactively alter an existing one. One-time recipe (verified against a live `postgres:16` instance): ```sql diff --git a/PRODUCT.md b/PRODUCT.md index b3336b5..83844e0 100644 --- a/PRODUCT.md +++ b/PRODUCT.md @@ -40,7 +40,7 @@ never creates one; it mounts onto yours. | `createMemoryRoutes({ memory, requireGrant })` | Hono sub-app the host mounts at `/api/tenants/:tenantId/memory` | | `mountWorkflowMemory(app, { memory, agentToken })` | Parallel run-scoped `/api/workflow-memory/*` for deployed agents | | `loadMemoryConfig()` | Config from env | -| `runMemoryMigrations(url)` | Apply pgvector schema | +| `runMemoryMigrations(dbConfig, { schema, ftsLanguage })` | Apply pgvector schema | | `@corbits/memory/sidecar-bundle` | Deployed-agent factory — no client code, no base URL, no token | | `@corbits/memory/distiller` | Optional process helpers: `runDistillTick`, `createResidentDistiller` | diff --git a/README.md b/README.md index 6fd10ab..c5ef414 100644 --- a/README.md +++ b/README.md @@ -65,20 +65,15 @@ Mount `installMemory` below the middleware that sets `principal`/`tenant` does). Identity comes from `c.get("principal")` — request bodies never carry tenant or principal. Missing principal → 401, missing grant → 403. -Apply migrations before serving traffic: +Apply migrations before serving traffic, with the same `DBConfig` and +`schema` your hub passes to Interchange `runMigrations`. Memory's tables land +in their own `memory` schema, with foreign keys into that schema's `tenant` +and `principal` tables: ```ts import { runMemoryMigrations } from "@corbits/memory/migrations"; -export async function migrateMemory(databaseUrl: string): Promise { - await runMemoryMigrations(databaseUrl); -} - -const databaseUrl = process.env.DATABASE_URL; -if (databaseUrl === undefined) { - throw new Error("DATABASE_URL is required to run memory migrations"); -} -await migrateMemory(databaseUrl); +await runMemoryMigrations(dbConfig, { schema: "public", ftsLanguage: "english" }); ``` `loadMemoryConfig()` reads `DATABASE_URL` (required — tables live in a diff --git a/bun.lock b/bun.lock index b9b0906..9976731 100644 --- a/bun.lock +++ b/bun.lock @@ -10,6 +10,7 @@ "devDependencies": { "@intx/agent": "0.4.0", "@intx/authz": "0.4.0", + "@intx/db": "0.4.0", "@intx/hub-api": "0.4.0", "@intx/log": "0.4.0", "@intx/types": "0.4.0", @@ -24,6 +25,7 @@ "peerDependencies": { "@intx/agent": "^0.4.0", "@intx/authz": "^0.4.0", + "@intx/db": "^0.4.0", "@intx/hub-api": "^0.4.0", "@intx/log": "^0.4.0", "@intx/types": "^0.4.0", diff --git a/migrations/0002_memory_baseline.sql b/migrations/0002_memory_baseline.sql index 60bf66b..54628b1 100644 --- a/migrations/0002_memory_baseline.sql +++ b/migrations/0002_memory_baseline.sql @@ -85,8 +85,8 @@ CREATE TABLE IF NOT EXISTS "memory"."chunk" ( CREATE UNIQUE INDEX IF NOT EXISTS "chunk_version_ordinal_uniq" ON "memory"."chunk" ("version_id", "ordinal"); --- {{FTS_LANGUAGE}} is substituted by runMemoryMigrations from FTS_LANGUAGE --- (or opts.ftsLanguage). Must match the language used at query time. +-- {{FTS_LANGUAGE}} is substituted by runMemoryMigrations from its +-- ftsLanguage option. Must match the language used at query time. ALTER TABLE "memory"."chunk" ADD COLUMN IF NOT EXISTS "text_fts" tsvector GENERATED ALWAYS AS (to_tsvector('{{FTS_LANGUAGE}}', "text")) STORED; diff --git a/migrations/0003_claim_bearing.sql b/migrations/0003_claim_bearing.sql index 94e38a4..9f9d912 100644 --- a/migrations/0003_claim_bearing.sql +++ b/migrations/0003_claim_bearing.sql @@ -5,9 +5,14 @@ ALTER TABLE "memory"."version" ADD COLUMN IF NOT EXISTS "provenance" text NOT NULL DEFAULT 'unknown'; -ALTER TABLE "memory"."version" - DROP CONSTRAINT IF EXISTS "version_provenance_check"; - -ALTER TABLE "memory"."version" - ADD CONSTRAINT "version_provenance_check" - CHECK ("provenance" IN ('stated', 'inferred', 'unknown')); +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'version_provenance_check' AND conrelid = '"memory"."version"'::regclass + ) THEN + ALTER TABLE "memory"."version" + ADD CONSTRAINT "version_provenance_check" + CHECK ("provenance" IN ('stated', 'inferred', 'unknown')); + END IF; +END $$; diff --git a/migrations/0004_temporal_model.sql b/migrations/0004_temporal_model.sql index 2e7170d..90e15cf 100644 --- a/migrations/0004_temporal_model.sql +++ b/migrations/0004_temporal_model.sql @@ -2,15 +2,35 @@ -- See docs/TEMPORAL.md. No asserted_at — occurred_at is effective time; -- ingested_at is when the memory plane learned the content. -ALTER TABLE "memory"."version" - ADD COLUMN IF NOT EXISTS "temporal_class" text NOT NULL DEFAULT 'event'; - -ALTER TABLE "memory"."version" - DROP CONSTRAINT IF EXISTS "version_temporal_class_check"; +-- The backfill runs only in the replay that adds the column, so distilled +-- claims (inferred provenance) written before this migration default to +-- state ranking exactly once; later inferred `event` rows are left alone. +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM information_schema.columns + WHERE table_schema = 'memory' AND table_name = 'version' + AND column_name = 'temporal_class' + ) THEN + ALTER TABLE "memory"."version" + ADD COLUMN "temporal_class" text NOT NULL DEFAULT 'event'; + UPDATE "memory"."version" + SET "temporal_class" = 'state' + WHERE "provenance" = 'inferred'; + END IF; +END $$; -ALTER TABLE "memory"."version" - ADD CONSTRAINT "version_temporal_class_check" - CHECK ("temporal_class" IN ('event', 'deadline', 'state', 'lesson')); +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'version_temporal_class_check' AND conrelid = '"memory"."version"'::regclass + ) THEN + ALTER TABLE "memory"."version" + ADD CONSTRAINT "version_temporal_class_check" + CHECK ("temporal_class" IN ('event', 'deadline', 'state', 'lesson')); + END IF; +END $$; ALTER TABLE "memory"."version" ADD COLUMN IF NOT EXISTS "valid_from" timestamp; @@ -18,8 +38,3 @@ ALTER TABLE "memory"."version" ALTER TABLE "memory"."version" ADD COLUMN IF NOT EXISTS "valid_until" timestamp; --- Distilled claims (inferred provenance) default to state ranking. -UPDATE "memory"."version" - SET "temporal_class" = 'state' - WHERE "provenance" = 'inferred' - AND "temporal_class" = 'event'; diff --git a/migrations/0007_retention.sql b/migrations/0007_retention.sql index 4366c4a..174f4a6 100644 --- a/migrations/0007_retention.sql +++ b/migrations/0007_retention.sql @@ -4,12 +4,17 @@ ALTER TABLE "memory"."version" ADD COLUMN IF NOT EXISTS "retention_class" text NOT NULL DEFAULT 'standard'; -ALTER TABLE "memory"."version" - DROP CONSTRAINT IF EXISTS "version_retention_class_check"; - -ALTER TABLE "memory"."version" - ADD CONSTRAINT "version_retention_class_check" - CHECK ("retention_class" IN ('durable', 'standard', 'ephemeral', 'source_only')); +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'version_retention_class_check' AND conrelid = '"memory"."version"'::regclass + ) THEN + ALTER TABLE "memory"."version" + ADD CONSTRAINT "version_retention_class_check" + CHECK ("retention_class" IN ('durable', 'standard', 'ephemeral', 'source_only')); + END IF; +END $$; CREATE INDEX IF NOT EXISTS "version_retention_ephemeral_idx" ON "memory"."version" ("tenant_id", "retention_class", "valid_until") diff --git a/migrations/0008_tenant_principal_fks.sql b/migrations/0008_tenant_principal_fks.sql index 9cd6b68..622be24 100644 --- a/migrations/0008_tenant_principal_fks.sql +++ b/migrations/0008_tenant_principal_fks.sql @@ -4,43 +4,135 @@ -- not exist. Tenant deletion cascades through this package's data; -- principal deletion cascades only the attribution column it owns -- (memory.version.created_by_principal_id), never a whole tenant's memory. +-- Each constraint is added only when missing, so re-running is a no-op; a +-- constraint that points at a different host schema fails the replay loudly. -ALTER TABLE "memory"."document" - ADD CONSTRAINT "document_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'document_tenant_id_fkey' AND conrelid = '"memory"."document"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."document" + ADD CONSTRAINT "document_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."version" - ADD CONSTRAINT "version_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'version_tenant_id_fkey' AND conrelid = '"memory"."version"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."version" + ADD CONSTRAINT "version_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."version" - ADD CONSTRAINT "version_created_by_principal_id_fkey" - FOREIGN KEY ("created_by_principal_id") REFERENCES "public"."principal"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'version_created_by_principal_id_fkey' AND conrelid = '"memory"."version"'::regclass + AND confrelid = '"public"."principal"'::regclass + ) THEN + ALTER TABLE "memory"."version" + ADD CONSTRAINT "version_created_by_principal_id_fkey" + FOREIGN KEY ("created_by_principal_id") REFERENCES "public"."principal"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."chunk" - ADD CONSTRAINT "chunk_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'chunk_tenant_id_fkey' AND conrelid = '"memory"."chunk"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."chunk" + ADD CONSTRAINT "chunk_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."entity" - ADD CONSTRAINT "entity_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'entity_tenant_id_fkey' AND conrelid = '"memory"."entity"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."entity" + ADD CONSTRAINT "entity_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."edge" - ADD CONSTRAINT "edge_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'edge_tenant_id_fkey' AND conrelid = '"memory"."edge"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."edge" + ADD CONSTRAINT "edge_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."raw_capture" - ADD CONSTRAINT "raw_capture_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'raw_capture_tenant_id_fkey' AND conrelid = '"memory"."raw_capture"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."raw_capture" + ADD CONSTRAINT "raw_capture_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."embed_model" - ADD CONSTRAINT "embed_model_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'embed_model_tenant_id_fkey' AND conrelid = '"memory"."embed_model"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."embed_model" + ADD CONSTRAINT "embed_model_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."transform_config" - ADD CONSTRAINT "transform_config_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'transform_config_tenant_id_fkey' AND conrelid = '"memory"."transform_config"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."transform_config" + ADD CONSTRAINT "transform_config_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; -ALTER TABLE "memory"."transform_run" - ADD CONSTRAINT "transform_run_tenant_id_fkey" - FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conname = 'transform_run_tenant_id_fkey' AND conrelid = '"memory"."transform_run"'::regclass + AND confrelid = '"public"."tenant"'::regclass + ) THEN + ALTER TABLE "memory"."transform_run" + ADD CONSTRAINT "transform_run_tenant_id_fkey" + FOREIGN KEY ("tenant_id") REFERENCES "public"."tenant"("id") ON DELETE CASCADE; + END IF; +END $$; diff --git a/migrations/0009_drop_migrations_ledger.sql b/migrations/0009_drop_migrations_ledger.sql new file mode 100644 index 0000000..2d6ad62 --- /dev/null +++ b/migrations/0009_drop_migrations_ledger.sql @@ -0,0 +1,2 @@ +-- Drops the ledger table earlier versions of runMemoryMigrations kept. +DROP TABLE IF EXISTS "memory"."_migrations"; diff --git a/package.json b/package.json index 351c5ca..c9ca4db 100644 --- a/package.json +++ b/package.json @@ -70,6 +70,7 @@ "peerDependencies": { "@intx/agent": "^0.4.0", "@intx/authz": "^0.4.0", + "@intx/db": "^0.4.0", "@intx/hub-api": "^0.4.0", "@intx/log": "^0.4.0", "@intx/types": "^0.4.0", @@ -82,6 +83,7 @@ "devDependencies": { "@intx/agent": "0.4.0", "@intx/authz": "0.4.0", + "@intx/db": "0.4.0", "@intx/hub-api": "0.4.0", "@intx/log": "0.4.0", "@intx/types": "0.4.0", diff --git a/scripts/db-setup.ts b/scripts/db-setup.ts index 045a26f..f7d64e4 100644 --- a/scripts/db-setup.ts +++ b/scripts/db-setup.ts @@ -1,10 +1,19 @@ +import { loadMemoryConfig } from "../src/mount-config.js"; import { runMemoryMigrations } from "../src/migrations.js"; -const url = - process.env["DATABASE_URL"]; -if (!url) throw new Error("DATABASE_URL is required"); +const { memory } = loadMemoryConfig(); +const url = new URL(memory.databaseUrl); +const sslmode = url.searchParams.get("sslmode"); -await runMemoryMigrations(url, { - log: (line) => console.log(` ${line}`), -}); +await runMemoryMigrations( + { + host: url.hostname, + port: Number(url.port || 5432), + user: decodeURIComponent(url.username), + password: decodeURIComponent(url.password), + database: url.pathname.slice(1), + ssl: sslmode === "require" || sslmode === "verify-ca" || sslmode === "verify-full", + }, + { schema: "public", ftsLanguage: memory.ftsLanguage }, +); console.log("Migrations complete."); diff --git a/src/core/fts-language.test.ts b/src/core/fts-language.test.ts index b1c295c..9c53b33 100644 --- a/src/core/fts-language.test.ts +++ b/src/core/fts-language.test.ts @@ -42,19 +42,30 @@ describe("the baseline migration language token", () => { }); }); -describe("runMemoryMigrations language boundary", () => { - it("falls back to the FTS_LANGUAGE env var when no option is passed", async () => { - // Pin the boundary contract without a live database: the runner must - // resolve exactly like the config loader, from the same env var. +describe("runMemoryMigrations boundary", () => { + const config = { + host: "unused", + port: 5432, + user: "unused", + password: "unused", + database: "unused", + }; + + it("rejects an invalid ftsLanguage before connecting", async () => { const { runMemoryMigrations } = await import("../migrations.js"); - process.env["FTS_LANGUAGE"] = "not a valid name"; - try { - await expect(runMemoryMigrations("postgres://unused")).rejects.toThrow( - "not a valid text search config name", - ); - } finally { - delete process.env["FTS_LANGUAGE"]; - } + await expect( + runMemoryMigrations(config, { + schema: "public", + ftsLanguage: "not a valid name", + }), + ).rejects.toThrow("not a valid text search config name"); + }); + + it("rejects an empty host schema before connecting", async () => { + const { runMemoryMigrations } = await import("../migrations.js"); + await expect( + runMemoryMigrations(config, { schema: "", ftsLanguage: "english" }), + ).rejects.toThrow("schema name must not be empty"); }); }); diff --git a/src/core/fts-language.ts b/src/core/fts-language.ts index 0a3ee68..35556bd 100644 --- a/src/core/fts-language.ts +++ b/src/core/fts-language.ts @@ -34,9 +34,14 @@ const APPLIED_REGCONFIG_RE = new RegExp( */ export function parseFtsLanguage(raw: string | undefined): string { if (raw === undefined || raw === "") return DEFAULT_FTS_LANGUAGE; + return assertFtsLanguage(raw); +} + +/** Validate a language with no default: an empty string is rejected. */ +export function assertFtsLanguage(raw: string): string { if (!FTS_LANGUAGE_PATTERN.test(raw)) { throw new Error( - `FTS_LANGUAGE "${raw}" is not a valid text search config name (expected ${FTS_LANGUAGE_PATTERN})`, + `FTS language "${raw}" is not a valid text search config name (expected ${FTS_LANGUAGE_PATTERN})`, ); } return raw; @@ -142,7 +147,7 @@ export async function verifyFtsLanguage( if (schema !== undefined) { throw new Error( `memory.chunk.text_fts was built with the schema-qualified text search config "${schema}.${applied}", ` + - `but FTS_LANGUAGE only supports unqualified pg_catalog configs. ` + + `but only unqualified pg_catalog configs are supported. ` + `Either drop the schema qualification (move/alias the config into pg_catalog), or rebuild the column ` + `under an unqualified config name:\n\n${rebuildColumnRecipe(ftsLanguage)}`, ); @@ -152,7 +157,7 @@ export async function verifyFtsLanguage( `FTS language mismatch: memory.chunk.text_fts was built with "${applied}" but the configuration says "${ftsLanguage}". ` + `Search would silently stem queries differently than the index.\n\n` + `To rebuild the column under the new language:\n\n${rebuildColumnRecipe(ftsLanguage)}\n\n` + - `Or, fix FTS_LANGUAGE back to "${applied}" instead.`, + `Or, set the FTS language (FTS_LANGUAGE / ftsLanguage) back to "${applied}" instead.`, ); } } diff --git a/src/migrations.ts b/src/migrations.ts index bc78a54..fa3c5a8 100644 --- a/src/migrations.ts +++ b/src/migrations.ts @@ -1,15 +1,16 @@ /** * Memory-plane (pgvector) schema migrations, callable by host apps. - * Applies every migrations/*.sql in filename order, each in its own - * transaction, tracked in memory._migrations so re-runs are idempotent - * and the ledger never collides with a host's public migration bookkeeping. + * Applies every migrations/*.sql in filename order on each run, the same + * way Interchange `runMigrations` does: every file is idempotent, so there + * is no ledger. */ import postgres from "postgres"; import { readdir, readFile } from "node:fs/promises"; import { join } from "node:path"; +import type { DBConfig } from "@intx/db"; import { FTS_LANGUAGE_TOKEN, - parseFtsLanguage, + assertFtsLanguage, verifyFtsLanguage, } from "./core/fts-language.js"; import { createRawSqlClient } from "./core/embed-sql.js"; @@ -19,55 +20,63 @@ import { MEMORY_SCHEMA } from "./db/schema.js"; // /migrations next to /dist, so this resolves in Node too. const MIGRATIONS_DIR = join(import.meta.dirname, "..", "migrations"); +// Advisory locks are namespaced by this integer alone; deliberately arbitrary +// and specific to @corbits/memory. +const LOCK_KEY = 0x3e30_7a11; + +function quoteIdentifier(name: string): string { + return `"${name.replace(/"/g, '""')}"`; +} + +/** + * Takes the same `config` and `schema` the host passes Interchange + * `runMigrations`: `schema` is where the host's `tenant` and `principal` + * tables live, and the `"public".` foreign-key references in the SQL are + * rewritten to it. Memory's own tables always live in the `memory` schema. + * `ftsLanguage` is fixed into the generated tsvector column; pass the same + * value `loadMemoryConfig` resolves for the query side. + */ export async function runMemoryMigrations( - databaseUrl: string, - opts: { log?: (line: string) => void; ftsLanguage?: string } = {}, + config: DBConfig, + options: { schema: string; ftsLanguage: string }, ): Promise { - const log = opts.log ?? (() => {}); - // This runner is an env-driven boundary like loadMemoryConfig: when the - // caller does not pass a language it reads the same FTS_LANGUAGE the query - // side will, so the two cannot diverge by defaulting differently. - const ftsLanguage = parseFtsLanguage( - opts.ftsLanguage ?? process.env["FTS_LANGUAGE"], - ); - const sql = postgres(databaseUrl, { max: 1 }); + if (options.schema.length === 0) { + throw new Error("runMemoryMigrations: schema name must not be empty"); + } + const hostSchema = quoteIdentifier(options.schema); + const ftsLanguage = assertFtsLanguage(options.ftsLanguage); + const sql = postgres({ + host: config.host, + port: config.port, + user: config.user, + password: config.password, + database: config.database, + ssl: config.ssl ?? false, + max: 1, + onnotice: () => undefined, + }); try { - // Schema first so the ledger and every later migration can land inside it - // even when 0001 has not been applied yet (fresh DB) or was skipped. - await sql.unsafe( - `CREATE SCHEMA IF NOT EXISTS "${MEMORY_SCHEMA}"`, - ); - await sql.unsafe( - `CREATE TABLE IF NOT EXISTS "${MEMORY_SCHEMA}"."_migrations" ( - "name" text PRIMARY KEY, - "applied_at" timestamp NOT NULL DEFAULT now() - )`, - ); - const appliedRows = (await sql.unsafe( - `SELECT name FROM "${MEMORY_SCHEMA}"."_migrations"`, - )) as unknown as { name: string }[]; - const applied = new Set(appliedRows.map((row) => row.name)); - const files = (await readdir(MIGRATIONS_DIR)) .filter((name) => name.endsWith(".sql")) .sort(); + const ddls = await Promise.all( + files.map(async (file) => + (await readFile(join(MIGRATIONS_DIR, file), "utf8")) + .replaceAll(FTS_LANGUAGE_TOKEN, ftsLanguage) + .replace(/"public"\.(?=")/g, `${hostSchema}.`), + ), + ); - for (const file of files) { - if (applied.has(file)) { - log(`(skip) ${file}`); - continue; - } - const raw = await readFile(join(MIGRATIONS_DIR, file), "utf8"); - const ddl = raw.replaceAll(FTS_LANGUAGE_TOKEN, ftsLanguage); - await sql.begin(async (tx) => { + // One transaction behind a transaction-scoped advisory lock: concurrent + // replicas serialize instead of racing CREATE ... IF NOT EXISTS or + // deadlocking on replayed ALTERs, and a failure releases the lock. + await sql.begin(async (tx) => { + await tx`SELECT pg_advisory_xact_lock(${LOCK_KEY})`; + await tx.unsafe(`CREATE SCHEMA IF NOT EXISTS "${MEMORY_SCHEMA}"`); + for (const ddl of ddls) { await tx.unsafe(ddl); - await tx.unsafe( - `INSERT INTO "${MEMORY_SCHEMA}"."_migrations" (name) VALUES ($1)`, - [file], - ); - }); - log(`applied ${file}`); - } + } + }); // The catalog is the authoritative record of which language the // generated column was actually built with; a previously-migrated