diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index e9ef921..e242b52 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -11,7 +11,8 @@ A **library, not a service**. It creates no HTTP server, opens no connection pool by default, owns no configuration, and starts no background work. A host calls two functions: -- `runMailboxMigrations(db)` — once at boot, before serving. +- `runMailboxMigrations(config, { schema })` — once at boot, before serving, + with the same arguments the host passes Interchange's `runMigrations`. - `createMailboxRoutes(deps)` — returns a `Hono` sub-app serving `/me/inbox*`, which the host mounts with `app.route`. @@ -48,9 +49,9 @@ is reached for. What it does **not** require: no session library, no logger configuration, no UI. What it *does* require of the database is an -Interchange-shaped control plane: `public.tenant` and `public.principal` in the -same database, in place before `runMailboxMigrations` runs, because the mailbox -tables foreign-key to both. Nothing changed in Interchange to make that work — +Interchange-shaped control plane: `tenant` and `principal` in the host schema +of the same database, in place before `runMailboxMigrations` runs, because the +mailbox tables foreign-key to both. Nothing changed in Interchange to make that work — the coupling lives entirely on this side. One further seam lives outside `createMailboxRoutes`, on the write side: @@ -171,7 +172,7 @@ own offboarding transaction; neither is scoped by view, because an offboarded tenant's trash is as much their data as their inbox. **Hard control-plane foreign keys.** `tenant_id` and `principal_id` on both -tables reference the host's `public.tenant` and `public.principal`, both +tables reference the host schema's `tenant` and `principal`, both `ON DELETE CASCADE` — the same posture as Interchange's own `session_mail.tenant_id`, extended to the principal. Consequences, stated rather than hidden: the control plane and the mail plane must share one @@ -431,8 +432,14 @@ two would query columns or rely on indexes the migrations never created. ## Migrations -`runMailboxMigrations(db)` is idempotent and safe to call unconditionally on -every boot of every replica. +`runMailboxMigrations(config, { schema })` is idempotent and safe to call +unconditionally on every boot of every replica. It opens one connection from +`config` and closes it when done. `schema` names the host schema holding +`tenant` and `principal`; the FKs are pointed there at execution time, while +ledger checksums hash the statements as shipped, so the same migration has the +same checksum whatever the host schema. The FKs are fixed by the run that +first applies each migration; passing a different `schema` later does not move +them. - The whole run is one transaction whose first statements are `SET LOCAL client_min_messages = warning` and a **transaction-scoped** @@ -474,8 +481,8 @@ every boot of every replica. **Everything lands in the `mailbox` schema, fully qualified.** Nothing resolves through `search_path`, so the host's own setting cannot redirect or shadow where the mailbox tables live. The one ordering constraint mounting imposes is -the control plane's: the DDL's foreign keys reference `public.tenant` and -`public.principal`, so those tables must exist before the first run. +the control plane's: the DDL's foreign keys reference the host schema's +`tenant` and `principal`, so those tables must exist before the first run. ## Boundaries diff --git a/CHANGELOG.md b/CHANGELOG.md index 44de4f5..8912ac7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -214,6 +214,12 @@ always called out under their own heading. ### Breaking +- **`runMailboxMigrations(config, { schema })` replaces `runMailboxMigrations(db)`.** + It takes the same `DBConfig` (`@intx/db`, now a peer) and `schema` the host + passes Interchange's `runMigrations`. `schema` is the host schema holding + `tenant` and `principal`; the mailbox FKs point there. The mailbox tables + stay in the `mailbox` schema. `createMailboxDb` is no longer exported; hosts + pass the handle they already have (e.g. `@intx/db`'s `createDB`). - **`mountMailbox` is replaced by `createMailboxRoutes(deps): Hono`.** The host mounts the returned sub-app with `app.route`. `deps.requireGrant` (`@intx/hub-api`'s `RequireGrant`, now a peer) gates reads on `mailbox:*` diff --git a/README.md b/README.md index 39e362a..91104e8 100644 --- a/README.md +++ b/README.md @@ -4,24 +4,24 @@ Give a **person** in an Interchange hub an inbox: list, read, flag, send, and li ## Runtime support -Node >= 24 consumes built `dist/`. Bun >= 1.2 runs TypeScript source. Peers: `@intx/hub-api`, `@intx/log`, `@intx/mailbox`, `@intx/mime`, `@intx/types`, `drizzle-orm`, `hono`, `postgres`. +Node >= 24 consumes built `dist/`. Bun >= 1.2 runs TypeScript source. Peers: `@intx/db`, `@intx/hub-api`, `@intx/log`, `@intx/mailbox`, `@intx/mime`, `@intx/types`, `drizzle-orm`, `hono`, `postgres`. ## Quickstart ```bash -npm add @corbits/mailbox @intx/hub-api @intx/log @intx/mailbox @intx/mime @intx/types drizzle-orm hono postgres +npm add @corbits/mailbox @intx/db @intx/hub-api @intx/log @intx/mailbox @intx/mime @intx/types drizzle-orm hono postgres ``` ```ts +import { createDB } from "@intx/db"; import { createInMemoryMailboxEventBus, - createMailboxDb, createMailboxRoutes, runMailboxMigrations, } from "@corbits/mailbox"; -const { db, close } = createMailboxDb(databaseUrl); -await runMailboxMigrations(db); +await runMailboxMigrations(dbConfig, { schema: "public" }); +const { db, close } = createDB(dbConfig); app.route( "/api/tenants/:tenantId/mailbox", @@ -37,33 +37,33 @@ app.route( process.once("SIGTERM", close); ``` -`app` is the host's `Hono` behind its tenant middleware, `requireGrant` comes from `@intx/hub-api`'s `createRequireGrant`, and `addressOf` and `transport` are the host's own directory and mail transport. `runMailboxMigrations` creates the `mailbox` schema; it needs the host's `tenant` and `principal` tables to exist first. +`app` is the host's `Hono` behind its tenant middleware, `requireGrant` comes from `@intx/hub-api`'s `createRequireGrant`, and `addressOf` and `transport` are the host's own directory and mail transport. `dbConfig` is the same config and `schema` the host passes Interchange's `runMigrations`: the schema holding its `tenant` and `principal` tables, which must exist first. The mailbox's own tables always live in the `mailbox` schema. -| `deps` | Type | What the host provides | -| --------------------- | ----------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `db` | `MailboxDb` | Mail lives there (schema `mailbox`). `createMailboxDb` opens a handle; a hub that already has one passes it instead. | -| `bus` | `MailboxEventBus` | SSE fan-out only; mail itself is Postgres. `createInMemoryMailboxEventBus()` for a single process; a shared bus when several processes must fan the same inbox events. | -| `requireGrant` | `RequireGrant` | Gates reads on `mailbox:*` `read`, send on `create`, and the flag and move verbs on `manage`. | -| `senderAddressFor` | `(principal: ResolvedPrincipal) => string` | That person's From: address, from the host's own directory. | -| `deliver` | `(message: OutgoingMailboxMessage) => void` | The host's mail transport. Called once per send with `{ raw, from, to, messageId }` after the message is filed in `Sent`. This package builds MIME; transmission is the host's job. | -| `heartbeatIntervalMs` | `number` (optional) | SSE keep-alive period. Defaults to 25s. | +| `deps` | Type | What the host provides | +| --------------------- | ----------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `db` | `MailboxDb` | Mail lives there (schema `mailbox`). A hub passes the handle it already has. | +| `bus` | `MailboxEventBus` | SSE fan-out only; mail itself is Postgres. `createInMemoryMailboxEventBus()` for a single process; a shared bus when several processes must fan the same inbox events. | +| `requireGrant` | `RequireGrant` | Gates reads on `mailbox:*` `read`, send on `create`, and the flag and move verbs on `manage`. | +| `senderAddressFor` | `(principal: ResolvedPrincipal) => string` | That person's From: address, from the host's own directory. | +| `deliver` | `(message: OutgoingMailboxMessage) => void` | The host's mail transport. Called once per send with `{ raw, from, to, messageId }` after the message is filed in `Sent`. This package builds MIME; transmission is the host's job. | +| `heartbeatIntervalMs` | `number` (optional) | SSE keep-alive period. Defaults to 25s. | ### Routes Paths are relative to where the host mounts the sub-app. Every route reads or writes only the caller's own mailbox. -| Method | Path | Purpose | -| ------ | ---------------------------- | ------------------------------------------------------------------------------------------- | +| Method | Path | Purpose | +| ------ | ---------------------------- | ---------------------------------------------------------------------------------------------------- | | GET | `/me/inbox` | List a folder newest first. `?folder=` INBOX (default), Sent, Archive, Trash; `?limit=`, `?cursor=`. | -| GET | `/me/inbox/threads` | A folder as threads (REFERENCES algorithm). `?folder=` as above. | -| GET | `/me/inbox/threads/:rootUid` | One thread, rooted at `rootUid`. `?folder=` as above. | -| POST | `/me/inbox/send` | Build a message from `{ to, subject?, body, inReplyTo? }`, file it in `Sent`, call `deliver`. | -| GET | `/me/inbox/events` | Server-sent `mailbox` events for the caller, with a heartbeat. | -| POST | `/me/inbox/:uid/read` | Set `\Seen` on a message in `?folder=` (INBOX by default). | -| POST | `/me/inbox/:uid/unread` | Clear `\Seen` on a message in `?folder=` (INBOX by default). | -| POST | `/me/inbox/:uid/archive` | Move from INBOX to Archive. | -| POST | `/me/inbox/:uid/trash` | Move from INBOX to Trash. | -| POST | `/me/inbox/:uid/restore` | Move back to INBOX from `?folder=` (Archive by default). | +| GET | `/me/inbox/threads` | A folder as threads (REFERENCES algorithm). `?folder=` as above. | +| GET | `/me/inbox/threads/:rootUid` | One thread, rooted at `rootUid`. `?folder=` as above. | +| POST | `/me/inbox/send` | Build a message from `{ to, subject?, body, inReplyTo? }`, file it in `Sent`, call `deliver`. | +| GET | `/me/inbox/events` | Server-sent `mailbox` events for the caller, with a heartbeat. | +| POST | `/me/inbox/:uid/read` | Set `\Seen` on a message in `?folder=` (INBOX by default). | +| POST | `/me/inbox/:uid/unread` | Clear `\Seen` on a message in `?folder=` (INBOX by default). | +| POST | `/me/inbox/:uid/archive` | Move from INBOX to Archive. | +| POST | `/me/inbox/:uid/trash` | Move from INBOX to Trash. | +| POST | `/me/inbox/:uid/restore` | Move back to INBOX from `?folder=` (Archive by default). | ### Agent-originated mail diff --git a/bun.lock b/bun.lock index 1ddd13c..048bf68 100644 --- a/bun.lock +++ b/bun.lock @@ -13,6 +13,7 @@ "hono-openapi": "1.3.1", }, "devDependencies": { + "@intx/db": "0.4.0", "@intx/hub-api": "0.4.0", "@intx/log": "0.4.0", "@intx/mailbox": "0.4.0", @@ -28,6 +29,7 @@ "typescript": "5.7.2", }, "peerDependencies": { + "@intx/db": "^0.4.0", "@intx/hub-api": "^0.4.0", "@intx/log": "^0.4.0", "@intx/mailbox": "^0.4.0", diff --git a/package.json b/package.json index 6589b77..6a316b9 100644 --- a/package.json +++ b/package.json @@ -61,6 +61,7 @@ "hono-openapi": "1.3.1" }, "peerDependencies": { + "@intx/db": "^0.4.0", "@intx/hub-api": "^0.4.0", "@intx/log": "^0.4.0", "@intx/mailbox": "^0.4.0", @@ -71,6 +72,7 @@ "postgres": "^3.4.0" }, "devDependencies": { + "@intx/db": "0.4.0", "@intx/hub-api": "0.4.0", "@intx/log": "0.4.0", "@intx/mailbox": "0.4.0", diff --git a/src/db.ts b/src/db.ts index b699e0a..a88c5fe 100644 --- a/src/db.ts +++ b/src/db.ts @@ -1,5 +1,4 @@ -import { drizzle, type PostgresJsDatabase } from "drizzle-orm/postgres-js"; -import postgres from "postgres"; +import type { PostgresJsDatabase } from "drizzle-orm/postgres-js"; /** * The db handle this package expects a host to hand in — the drizzle instance @@ -13,15 +12,3 @@ import postgres from "postgres"; // handle bound to its own (e.g. `createDB`'s). Nothing here reads `db.query`. export type MailboxDb = PostgresJsDatabase; -/** - * Opens a standalone handle, for hosts and scripts that don't already have one. - * `close` drains the pool — without it a migrate-only script keeps an open - * socket and never exits. - */ -export function createMailboxDb(connectionString: string): { - db: MailboxDb; - close: () => Promise; -} { - const client = postgres(connectionString); - return { db: drizzle(client), close: () => client.end() }; -} diff --git a/src/index.ts b/src/index.ts index 836ff49..0480f81 100644 --- a/src/index.ts +++ b/src/index.ts @@ -13,7 +13,6 @@ export { runMailboxMigrations, MigrationChecksumError } from "./migrations.js"; export { SchemaTypeMismatchError } from "./schema-check.js"; -export { createMailboxDb } from "./db.js"; export type { MailboxDb } from "./db.js"; // The native `MailboxStore` over `mailbox.principal_mail` / diff --git a/src/migrations.test.ts b/src/migrations.test.ts index f04ffe1..147a928 100644 --- a/src/migrations.test.ts +++ b/src/migrations.test.ts @@ -10,15 +10,17 @@ import { afterAll, beforeAll, describe, expect, test } from "bun:test"; import postgres from "postgres"; import { drizzle } from "drizzle-orm/postgres-js"; import { sql } from "drizzle-orm"; -import { createMailboxDb } from "./db.js"; import { MIGRATIONS, MigrationChecksumError, + applyMailboxMigrations, runMailboxMigrations, } from "./migrations.js"; import { buildMailFrame } from "./frame.js"; import { createHostControlPlane, + createMailboxDb, + dbConfigFromUrl, seedScope, TEST_DATABASE_URL, } from "./test-helpers.js"; @@ -33,7 +35,8 @@ beforeAll(async () => { afterAll(async () => { // Leave the schema the way every other suite expects to find it, whatever // the last case here did to it. - await runMailboxMigrations(adminDb); + await dropMailboxSchema(); + await applyMailboxMigrations(adminDb, "public"); await admin.end(); }); @@ -61,9 +64,43 @@ async function fromEmpty( } describe("runMailboxMigrations", () => { + test("points the FKs at the host schema it is given", async () => { + await fromEmpty(async ({ client }) => { + await admin.unsafe(`DROP SCHEMA IF EXISTS "host_cp" CASCADE`); + await admin.unsafe(`CREATE SCHEMA "host_cp"`); + await admin.unsafe(`CREATE TABLE "host_cp"."tenant" ("id" text PRIMARY KEY)`); + await admin.unsafe(`CREATE TABLE "host_cp"."principal" ("id" text PRIMARY KEY)`); + try { + await runMailboxMigrations(dbConfigFromUrl(TEST_DATABASE_URL), { + schema: "host_cp", + }); + const targets = await client<{ target: string }[]>` + SELECT DISTINCT confrelid::regclass::text AS target + FROM pg_constraint + WHERE contype = 'f' + AND connamespace = 'mailbox'::regnamespace + AND confrelid::regclass::text NOT LIKE 'mailbox.%' + ORDER BY target`; + expect(targets.map((row) => row.target)).toEqual([ + "host_cp.principal", + "host_cp.tenant", + ]); + } finally { + await dropMailboxSchema(); + await admin.unsafe(`DROP SCHEMA "host_cp" CASCADE`); + } + }); + }); + + test("refuses an empty schema name", async () => { + await expect( + runMailboxMigrations(dbConfigFromUrl(TEST_DATABASE_URL), { schema: "" }), + ).rejects.toThrow("schema name must not be empty"); + }); + test("builds the full schema from an empty database", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); // The mail plane reads 1-1 with Interchange's `session_mail`: the message // as delivered, plus the cached header columns and this package's scope. @@ -133,7 +170,7 @@ describe("runMailboxMigrations", () => { test("0005 leaves uid and modseq NOT NULL: every write path is the native store now", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ column_name: string; is_nullable: string }>( sql`SELECT column_name, is_nullable FROM information_schema.columns WHERE table_schema = 'mailbox' AND table_name = 'principal_mail' @@ -149,7 +186,7 @@ describe("runMailboxMigrations", () => { test("the keyset index matches the list query's ORDER BY exactly", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const [row] = await db.execute<{ indexdef: string }>( sql`SELECT indexdef FROM pg_indexes WHERE schemaname = 'mailbox' AND indexname = 'principal_mail_tenant_id_principal_id_created_at_id_idx'`, @@ -169,7 +206,7 @@ describe("runMailboxMigrations", () => { // the upgrade and every older message would project no parent. await fromEmpty(async ({ db }) => { // Build the pre-0002 schema, then seed through it. - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); await db.execute( sql`ALTER TABLE "mailbox"."principal_mail" DROP COLUMN "message_id", DROP COLUMN "in_reply_to"`, @@ -215,7 +252,7 @@ describe("runMailboxMigrations", () => { `); } - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ message_key: string; @@ -241,7 +278,7 @@ describe("runMailboxMigrations", () => { // entire `raw`, NUL included) and GREEN once only the NUL-stripped header // slice reaches `convert_from`. await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); await db.execute( sql`ALTER TABLE "mailbox"."principal_mail" DROP COLUMN "message_id", DROP COLUMN "in_reply_to"`, @@ -276,7 +313,7 @@ describe("runMailboxMigrations", () => { `); } - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const ledger = await db.execute<{ id: string }>( sql`SELECT "id" FROM "mailbox"."corbits_mailbox_migrations" ORDER BY "id"`, @@ -311,7 +348,7 @@ describe("runMailboxMigrations", () => { // only the first fragment, and every older message would then link to the // wrong ancestor — worse than linking to none. await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); await db.execute( sql`ALTER TABLE "mailbox"."principal_mail" DROP COLUMN "references"`, ); @@ -356,7 +393,7 @@ describe("runMailboxMigrations", () => { `); } - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ message_key: string; @@ -381,7 +418,7 @@ describe("runMailboxMigrations", () => { // runtime path now uses too, so a frame decoded before or after the // upgrade projects the same cached `in_reply_to`. await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); await db.execute( sql`ALTER TABLE "mailbox"."principal_mail" DROP COLUMN "message_id", DROP COLUMN "in_reply_to"`, @@ -418,7 +455,7 @@ describe("runMailboxMigrations", () => { ${Buffer.from(enc.encode(text))}, ${key}, ${i + 1}, ${i + 1}) `); } - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ message_key: string; message_id: string | null; @@ -440,8 +477,8 @@ describe("runMailboxMigrations", () => { test("is idempotent: running twice does not error and applies once", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ id: string; count: string }>( sql`SELECT "id", count(*)::text AS count @@ -460,7 +497,7 @@ describe("runMailboxMigrations", () => { test("records a checksum per applied migration", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ id: string; checksum: string | null }>( sql`SELECT "id", "checksum" FROM "mailbox"."corbits_mailbox_migrations"`, ); @@ -473,7 +510,7 @@ describe("runMailboxMigrations", () => { test("refuses to boot when a shipped migration was edited after it applied", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); // Stand in for someone editing MIGRATIONS[0].statements in place: the // ledger now disagrees with the code, which is exactly the divergence // that used to be invisible — old environments skip the edit forever @@ -485,10 +522,10 @@ describe("runMailboxMigrations", () => { // A NAMED error, matching both sibling cores: a host catching this to // tell "someone edited a migration" apart from "the database is down" // should not have to regex-match a message string. - await expect(runMailboxMigrations(db)).rejects.toThrow( + await expect(applyMailboxMigrations(db, "public")).rejects.toThrow( MigrationChecksumError, ); - await expect(runMailboxMigrations(db)).rejects.toThrow( + await expect(applyMailboxMigrations(db, "public")).rejects.toThrow( /has changed since it was applied/, ); }); @@ -501,7 +538,7 @@ describe("runMailboxMigrations", () => { // 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 runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); // Drizzle wraps driver errors, so the NOT NULL violation is on `.cause`, // not on the message `toThrow` would match. const failure = await db @@ -521,7 +558,7 @@ describe("runMailboxMigrations", () => { test("creates its own ledger table distinct from any host table", async () => { await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ exists: boolean }>( sql`SELECT EXISTS (SELECT 1 FROM information_schema.tables WHERE table_schema = 'mailbox' @@ -536,7 +573,7 @@ describe("runMailboxMigrations", () => { // only belong to a tenant and principal the host knows, and offboarding // either carries the mailbox rows out with it. await fromEmpty(async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); const rows = await db.execute<{ constraint_name: string; table_name: string; @@ -575,7 +612,7 @@ describe("runMailboxMigrations", () => { connection: { search_path: "mbx_elsewhere" }, }); try { - await runMailboxMigrations(drizzle(client)); + await applyMailboxMigrations(drizzle(client), "public"); const found = await admin.unsafe( `SELECT to_regclass('mailbox.principal_mail') AS t, to_regclass('mailbox.corbits_mailbox_migrations') AS l, @@ -592,7 +629,7 @@ describe("runMailboxMigrations", () => { test("closing the handle it opened drains the pool", async () => { const { db, close } = createMailboxDb(TEST_DATABASE_URL); - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); await close(); // A closed pool refuses further work rather than hanging the process. expect(async () => { @@ -601,7 +638,7 @@ describe("runMailboxMigrations", () => { }); }); -describe("runMailboxMigrations under concurrent cold start", () => { +describe("applyMailboxMigrations under concurrent cold start", () => { // `CREATE TABLE IF NOT EXISTS` is NOT race-safe: the existence check and the // pg_type insert are not atomic, so without an advisory lock the losers crash // with 23505 on (typname, typnamespace). Pre-creating the ledger only moves @@ -610,7 +647,7 @@ describe("runMailboxMigrations under concurrent cold start", () => { await dropMailboxSchema(); const runners = Array.from({ length: 4 }, () => handle()); const results = await Promise.allSettled( - runners.map((r) => runMailboxMigrations(r.db)), + runners.map((r) => applyMailboxMigrations(r.db, "public")), ); await Promise.all(runners.map((r) => r.client.end())); @@ -635,12 +672,12 @@ describe("runMailboxMigrations under concurrent cold start", () => { test("a second wave against an already-migrated schema is a no-op for all", async () => { await dropMailboxSchema(); const first = handle(); - await runMailboxMigrations(first.db); + await applyMailboxMigrations(first.db, "public"); await first.client.end(); const runners = Array.from({ length: 3 }, () => handle()); const results = await Promise.allSettled( - runners.map((r) => runMailboxMigrations(r.db)), + runners.map((r) => applyMailboxMigrations(r.db, "public")), ); await Promise.all(runners.map((r) => r.client.end())); expect(results.every((r) => r.status === "fulfilled")).toBe(true); diff --git a/src/migrations.ts b/src/migrations.ts index 74f995a..c9e77ff 100644 --- a/src/migrations.ts +++ b/src/migrations.ts @@ -1,6 +1,9 @@ import { createHash } from "node:crypto"; import { sql, type SQL } from "drizzle-orm"; import { PgDialect } from "drizzle-orm/pg-core"; +import { drizzle } from "drizzle-orm/postgres-js"; +import postgres from "postgres"; +import type { DBConfig } from "@intx/db"; import type { MailboxDb } from "./db.js"; import { assertExpectedColumnTypes } from "./schema-check.js"; @@ -534,8 +537,8 @@ export class MigrationChecksumError extends Error { * every boot" printed a wall of what looked like errors on every replica start. * `SET LOCAL` scopes the change to this transaction and stops at NOTICE: * WARNING and above still reach the host untouched. It is set on the connection - * rather than via a client option so it holds for ANY handle a host hands in, - * including one this package did not construct. + * rather than via a client option so it holds for every handle the tests pass + * in, not only the one `runMailboxMigrations` builds. * * Each migration applies inside its own nested transaction (a savepoint under * the outer one), so a migration is all-or-nothing with its ledger row and can @@ -544,7 +547,11 @@ export class MigrationChecksumError extends Error { * stops being equivalent, and a runner that rolls back inconsistently is not * something to discover then. */ -export async function runMailboxMigrations(db: MailboxDb): Promise { +export async function applyMailboxMigrations( + db: MailboxDb, + hostSchema: string, +): Promise { + const hostIdent = `"${hostSchema.replace(/"/g, '""')}".`; 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})`); @@ -587,7 +594,12 @@ export async function runMailboxMigrations(db: MailboxDb): Promise { if (migration.assertColumnsBeforeStatement === index) { await assertExpectedColumnTypes(step); } - await step.execute(statement); + // Checksums hash the statements as shipped, so the ledger is the + // same whichever host schema the FKs are pointed at. + const rendered = DIALECT.sqlToQuery(statement).sql; + await step.execute( + sql.raw(rendered.replace(/"public"\.(?=")/g, hostIdent)), + ); } await step.execute( sql`INSERT INTO "mailbox".${sql.identifier(LEDGER_TABLE)} ("id", "checksum") VALUES (${migration.id}, ${expected})`, @@ -603,3 +615,33 @@ export async function runMailboxMigrations(db: MailboxDb): Promise { await assertExpectedColumnTypes(tx); }); } + +/** + * Takes the same `config` and `schema` the host passes Interchange's + * `runMigrations`: `schema` is where the host's `tenant` and `principal` + * tables live, and the mailbox FKs point there. The mailbox's own tables + * always live in the `mailbox` schema. + */ +export async function runMailboxMigrations( + config: DBConfig, + options: { schema: string }, +): Promise { + if (options.schema.length === 0) { + throw new Error("runMailboxMigrations: schema name must not be empty"); + } + const client = postgres({ + host: config.host, + port: config.port, + user: config.user, + password: config.password, + database: config.database, + ssl: config.ssl, + max: 1, + onnotice: () => undefined, + }); + try { + await applyMailboxMigrations(drizzle(client), options.schema); + } finally { + await client.end({ timeout: 5 }); + } +} diff --git a/src/schema-check.test.ts b/src/schema-check.test.ts index 374fe93..affd0ef 100644 --- a/src/schema-check.test.ts +++ b/src/schema-check.test.ts @@ -12,7 +12,7 @@ import { afterAll, beforeAll, describe, expect, test } from "bun:test"; import postgres from "postgres"; import { drizzle } from "drizzle-orm/postgres-js"; import { sql } from "drizzle-orm"; -import { MIGRATIONS, runMailboxMigrations } from "./migrations.js"; +import { MIGRATIONS, applyMailboxMigrations } from "./migrations.js"; import { expectedColumnTypes, SchemaTypeMismatchError, @@ -29,7 +29,7 @@ afterAll(async () => { // Leave the schema rebuilt for whatever suite runs after this one — the // planted conflicting table has to go first, or the rebuild rejects too. await admin.unsafe(`DROP SCHEMA IF EXISTS "mailbox" CASCADE`); - await runMailboxMigrations(drizzle(admin)); + await applyMailboxMigrations(drizzle(admin), "public"); await admin.end(); }); @@ -74,7 +74,7 @@ async function ledgerRows(schema: string): Promise { } /** - * The error `runMailboxMigrations` rejected with. A boot that SUCCEEDS is the + * The error `applyMailboxMigrations` rejected with. A boot that SUCCEEDS is the * failure every case below is written to catch, so it must not slip through as * an `undefined` that the assertions then read properties off. */ @@ -110,7 +110,7 @@ describe("expectedColumnTypes", () => { describe("boot against a host table this package did not create", () => { test("a fresh, correct database boots and records the migration", async () => { await inFreshSchema("mbx_check_ok", async ({ db }) => { - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); expect(await ledgerRows("mailbox")).toBe(MIGRATIONS.length); }); }); @@ -136,7 +136,7 @@ describe("boot against a host table this package did not create", () => { "refs" jsonb, "created_at" timestamptz NOT NULL DEFAULT now() )`); - const failure = await bootFailure(runMailboxMigrations(db)); + const failure = await bootFailure(applyMailboxMigrations(db, "public")); expect(failure).toBeInstanceOf(SchemaTypeMismatchError); expect((failure as SchemaTypeMismatchError).mismatches).toEqual([ "principal_mail.created_at is timestamp with time zone, " + @@ -174,7 +174,7 @@ describe("boot against a host table this package did not create", () => { "refs" jsonb, "created_at" timestamp NOT NULL DEFAULT now() )`); - const failure = await bootFailure(runMailboxMigrations(db)); + const failure = await bootFailure(applyMailboxMigrations(db, "public")); expect(failure).toBeInstanceOf(SchemaTypeMismatchError); expect((failure as SchemaTypeMismatchError).mismatches).toEqual([ "principal_mail.subject is missing (expected text)", @@ -199,13 +199,13 @@ describe("boot against a host table this package did not create", () => { "refs" jsonb, "created_at" timestamptz NOT NULL DEFAULT now() )`); - await expect(runMailboxMigrations(db)).rejects.toThrow( + await expect(applyMailboxMigrations(db, "public")).rejects.toThrow( SchemaTypeMismatchError, ); // No "it already applied, skip it" shortcut on the second attempt: the // ledger is empty, so the check runs again and fails again. A guard that // only fires on the first boot is one a restart disables. - await expect(runMailboxMigrations(db)).rejects.toThrow( + await expect(applyMailboxMigrations(db, "public")).rejects.toThrow( SchemaTypeMismatchError, ); expect(await ledgerRows(schema)).toBe(0); diff --git a/src/schema-ddl-parity.test.ts b/src/schema-ddl-parity.test.ts index 1351ae5..5abd3d6 100644 --- a/src/schema-ddl-parity.test.ts +++ b/src/schema-ddl-parity.test.ts @@ -9,7 +9,7 @@ import { drizzle } from "drizzle-orm/postgres-js"; import { sql } from "drizzle-orm"; import { getTableConfig } from "drizzle-orm/pg-core"; import { principalMail } from "./schema.js"; -import { runMailboxMigrations } from "./migrations.js"; +import { applyMailboxMigrations } from "./migrations.js"; import { createHostControlPlane, TEST_DATABASE_URL } from "./test-helpers.js"; // The tables live in this package's own `mailbox` schema, so what is compared @@ -22,7 +22,7 @@ const client = postgres(TEST_DATABASE_URL, { onnotice: () => {} }); beforeAll(async () => { const db = drizzle(client); await createHostControlPlane(db); - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); }); afterAll(async () => { @@ -111,7 +111,7 @@ async function liveIndexes(table: string): Promise { // hold to this parity. const TABLES = [{ name: "principal_mail", declared: principalMail }] as const; -describe("schema.ts vs. the DDL runMailboxMigrations actually creates", () => { +describe("schema.ts vs. the DDL applyMailboxMigrations actually creates", () => { for (const { name, declared } of TABLES) { it(`${name}: declares exactly the indexes the live table has, in the same column order`, async () => { expect(declaredIndexes(declared)).toEqual(await liveIndexes(name)); diff --git a/src/test-helpers.ts b/src/test-helpers.ts index 8ebafe8..4531f80 100644 --- a/src/test-helpers.ts +++ b/src/test-helpers.ts @@ -1,7 +1,10 @@ -import { createMailboxDb, type MailboxDb } from "./db.js"; -import { runMailboxMigrations } from "./migrations.js"; +import { drizzle } from "drizzle-orm/postgres-js"; +import postgres from "postgres"; +import type { MailboxDb } from "./db.js"; +import { applyMailboxMigrations } from "./migrations.js"; import { sql } from "drizzle-orm"; import { Hono } from "hono"; +import type { DBConfig } from "@intx/db"; import type { RequireGrant, TenantEnv } from "@intx/hub-api"; import type { ResolvedPrincipal } from "./mount.js"; @@ -9,6 +12,30 @@ export const TEST_DATABASE_URL = process.env.MAILBOX_TEST_DATABASE_URL ?? "postgres://postgres:postgres@localhost:5433/mailbox_core"; +/** + * Opens a standalone handle. `close` drains the pool — without it a suite + * keeps an open socket and never exits. + */ +export function createMailboxDb(connectionString: string): { + db: MailboxDb; + close: () => Promise; +} { + const client = postgres(connectionString); + return { db: drizzle(client), close: () => client.end() }; +} + +/** `url` as the `DBConfig` Interchange's migration runners take. */ +export function dbConfigFromUrl(url: string): DBConfig { + const parsed = new URL(url); + return { + host: parsed.hostname, + port: Number(parsed.port || 5432), + user: decodeURIComponent(parsed.username), + password: decodeURIComponent(parsed.password), + database: parsed.pathname.slice(1), + }; +} + /** * The minimum control plane the FKs require: the host's `tenant` and * `principal` tables, with only the columns the mailbox references. Real @@ -67,7 +94,7 @@ export async function withTestDb(): Promise { // The control plane must exist before the mailbox migrations can FK to it // — same order a real host boots in. await createHostControlPlane(db); - await runMailboxMigrations(db); + await applyMailboxMigrations(db, "public"); return db; })(); const db = await shared; diff --git a/tests/lib/db-harness.ts b/tests/lib/db-harness.ts index 0af809b..35b1633 100644 --- a/tests/lib/db-harness.ts +++ b/tests/lib/db-harness.ts @@ -12,6 +12,7 @@ import { import { allowAllGrants, createHostControlPlane, + dbConfigFromUrl, TEST_DATABASE_URL, } from "../../src/test-helpers.js"; @@ -41,13 +42,14 @@ export async function createTestDb(): Promise { url.pathname = `/${name}`; const client = postgres(url.toString(), { onnotice: () => {} }); const db = drizzle(client); + const config = dbConfigFromUrl(url.toString()); const close = async () => { await client.end(); await admin((sql) => sql.unsafe(`DROP DATABASE "${name}" WITH (FORCE)`)); }; try { await createHostControlPlane(db); - await runMailboxMigrations(db); + await runMailboxMigrations(config, { schema: "public" }); } catch (err) { await close(); throw err;