From fde06c7e1ed5fa2a97908d805255159ba4197a07 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 24 Sep 2026 23:48:43 -0700 Subject: [PATCH 1/3] test(migrations): fold the schema parity and column-type checks in The column-type rejection and the schema.ts-vs-DDL index parity are migration behavior, so they live with the migration suite. Closes CL-9270. --- ARCHITECTURE.md | 2 +- CONTRIBUTING.md | 2 +- src/migrations.test.ts | 200 ++++++++++++++++++++++++++++++- src/native-store.ts | 2 +- src/schema-check.test.ts | 214 ---------------------------------- src/schema-ddl-parity.test.ts | 136 --------------------- src/schema.ts | 2 +- 7 files changed, 203 insertions(+), 355 deletions(-) delete mode 100644 src/schema-check.test.ts delete mode 100644 src/schema-ddl-parity.test.ts diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index e242b52..d39f505 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -428,7 +428,7 @@ index-servable, ~3.8x its pre-split cost. `schema.ts` and `migrations.ts` must agree statement for statement: the runtime queries read through the drizzle table object, so a drift between the two would query columns or rely on indexes the migrations never created. -`src/schema-ddl-parity.test.ts` diffs the two against a live database. +`src/migrations.test.ts` diffs the two against a live database. ## Migrations diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 9d2d1ba..926516e 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -52,7 +52,7 @@ statements, so editing one that has already been applied fails loudly on the nex rather than letting fresh and existing databases diverge. Add a new migration instead. `schema.ts` and `migrations.ts` must agree statement for statement — the runtime -queries read through the drizzle table object, and `src/schema-ddl-parity.test.ts` +queries read through the drizzle table object, and `src/migrations.test.ts` diffs the two against a live database. Change one, change the other, in the same commit. ## Agent-originated mail diff --git a/src/migrations.test.ts b/src/migrations.test.ts index 147a928..c20606c 100644 --- a/src/migrations.test.ts +++ b/src/migrations.test.ts @@ -10,6 +10,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 { getTableConfig } from "drizzle-orm/pg-core"; import { MIGRATIONS, MigrationChecksumError, @@ -17,6 +18,11 @@ import { runMailboxMigrations, } from "./migrations.js"; import { buildMailFrame } from "./frame.js"; +import { principalMail } from "./schema.js"; +import { + expectedColumnTypes, + SchemaTypeMismatchError, +} from "./schema-check.js"; import { createHostControlPlane, createMailboxDb, @@ -149,7 +155,7 @@ describe("runMailboxMigrations", () => { // the keyset the default page seeks on, and the three the thread read // adds — the msg-id lookup, the GIN index serving the `refs` containment // filter, and the thread's own oldest-first keyset. - // `schema-ddl-parity.test.ts` holds schema.ts to this same list. + // The schema.ts parity suite below holds schema.ts to this same list. expect(mailIndexes.map((i) => i.indexname)).toEqual([ "principal_mail_pkey", "principal_mail_refs_idx", @@ -683,3 +689,195 @@ describe("applyMailboxMigrations under concurrent cold start", () => { expect(results.every((r) => r.status === "fulfilled")).toBe(true); }); }); + +// The drizzle table object is a public export, so a host can point +// `drizzle-kit push`/`generate` at it. If it declares an index the migrations +// do not create, or with a different column order, that host's schema silently +// diverges from the one this package's queries were planned against. +describe("schema.ts vs. the DDL applyMailboxMigrations actually creates", () => { + /** `name USING method(col asc, col desc)` plus `unique`/`partial` markers. */ + type IndexDescriptor = string; + + function declaredIndexes(): IndexDescriptor[] { + return getTableConfig(principalMail) + .indexes.map((index) => { + const config = index.config; + const columns = config.columns + .map((column) => { + // An expression index has no `.name`; fail rather than compare it + // as blank. + const name = (column as { name?: string }).name; + if (name === undefined) { + throw new Error( + `index ${config.name} uses an expression column this parity check cannot canonicalize`, + ); + } + const order = + (column as { indexConfig?: { order?: string } }).indexConfig + ?.order ?? "asc"; + return `${name} ${order}`; + }) + .join(", "); + const flags = [ + config.unique === true ? "unique" : null, + config.where !== undefined ? "partial" : null, + ].filter((flag) => flag !== null); + const suffix = flags.length > 0 ? ` [${flags.join(" ")}]` : ""; + // The access method matters: `refs @> …` is only servable by GIN. + const method = (config as { method?: string }).method ?? "btree"; + return `${config.name} USING ${method}(${columns})${suffix}`; + }) + .sort(); + } + + // `pg_get_indexdef` renders `CREATE [UNIQUE] INDEX ON USING + // ()[ WHERE ()]`, with DESC spelled out and ASC implicit. + function canonicalizeIndexDef(def: string): IndexDescriptor { + const match = + /^CREATE (UNIQUE )?INDEX (\S+) ON \S+ USING (\S+) \((.*?)\)( WHERE .*)?$/.exec( + def, + ); + if (match === null) throw new Error(`unparsed index definition: ${def}`); + const [, unique, name, method, columnList, where] = match; + const columns = columnList! + .split(", ") + .map((column) => { + const desc = / DESC$/.test(column); + const bare = column.replace(/ (DESC|ASC)$/, "").replace(/ NULLS.*$/, ""); + return `${bare} ${desc ? "desc" : "asc"}`; + }) + .join(", "); + const flags = [ + unique !== undefined ? "unique" : null, + where !== undefined ? "partial" : null, + ].filter((flag) => flag !== null); + const suffix = flags.length > 0 ? ` [${flags.join(" ")}]` : ""; + return `${name} USING ${method}(${columns})${suffix}`; + } + + test("principal_mail: declares exactly the indexes the live table has, in the same column order", async () => { + await fromEmpty(async ({ db }) => { + await applyMailboxMigrations(db, "public"); + const rows = await db.execute<{ indexdef: string }>(sql` + SELECT indexdef FROM pg_indexes + WHERE schemaname = 'mailbox' + AND tablename = 'principal_mail' + AND indexname <> 'principal_mail_pkey' + `); + const live = rows.map((row) => canonicalizeIndexDef(row.indexdef)).sort(); + expect(declaredIndexes()).toEqual(live); + }); + }); +}); + +describe("expectedColumnTypes", () => { + test("is derived from the drizzle tables: the mail plane alone, since 0005 dropped the management table", () => { + const tables = new Set(expectedColumnTypes().map((e) => e.table)); + expect(tables).toEqual(new Set(["principal_mail"])); + }); + + test("expects zoneless timestamps and text ids on every relevant column", () => { + const byKey = new Map( + expectedColumnTypes().map((e) => [`${e.table}.${e.column}`, e.dataType]), + ); + expect(byKey.get("principal_mail.created_at")).toBe( + "timestamp without time zone", + ); + expect(byKey.get("principal_mail.id")).toBe("text"); + expect(byKey.get("principal_mail.raw")).toBe("bytea"); + expect(byKey.get("principal_mail.refs")).toBe("jsonb"); + }); +}); + +// `CREATE TABLE IF NOT EXISTS` matches on the table NAME only. A host that +// already owns a `principal_mail` would get a silent no-op, a ledger row, and +// every read decoding ITS columns through OUR codec. Each case plants such a +// table and asserts the boot is rejected without leaving a ledger row that +// would make the next boot skip the check. +describe("boot against a host table this package did not create", () => { + /** The ledger row count, or 0 when the ledger table itself does not exist. */ + async function ledgerRows(): Promise { + const [ledger] = await admin<{ exists: boolean }[]>` + SELECT to_regclass('"mailbox"."corbits_mailbox_migrations"') IS NOT NULL AS exists`; + if (!ledger!.exists) return 0; + const [row] = await admin<{ n: number }[]>` + SELECT count(*)::int AS n FROM "mailbox"."corbits_mailbox_migrations"`; + return row!.n; + } + + /** Throws if the boot SUCCEEDS, rather than yielding an `undefined`. */ + async function bootFailure(promise: Promise): Promise { + try { + await promise; + } catch (error) { + return error as Error; + } + throw new Error("expected the boot to be rejected, but it succeeded"); + } + + /** Plants a host `principal_mail` with the given column DDL. */ + async function plantPrincipalMail(columns: string): Promise { + await admin.unsafe(`CREATE SCHEMA "mailbox"`); + await admin.unsafe(`CREATE TABLE "mailbox"."principal_mail" (${columns})`); + } + + const BASE_COLUMNS = ` + "id" text PRIMARY KEY, + "tenant_id" text NOT NULL, + "principal_id" text NOT NULL, + "address" text NOT NULL, + "direction" text NOT NULL, + "raw" bytea NOT NULL, + "from_address" text, + "message_key" text, + "refs" jsonb`; + + test("rejects a pre-existing table whose column TYPE diverges", async () => { + await fromEmpty(async ({ db }) => { + // `created_at` still a `timestamptz`: invisible to every query until a + // non-UTC host serves the wrong page. + await plantPrincipalMail(`${BASE_COLUMNS}, + "subject" text, + "created_at" timestamptz NOT NULL DEFAULT now()`); + 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, " + + "expected timestamp without time zone", + ]); + // Rejected inside the migration's transaction, so the ledger row rolled + // back; otherwise the next boot would skip the check. + expect(await ledgerRows()).toBe(0); + }); + }); + + test("rejects a pre-existing table with a column MISSING outright", async () => { + await fromEmpty(async ({ db }) => { + // No `subject`, a column no index covers. `refs` stays present: it is + // GIN-indexed, so its absence would be rejected by the DDL, not this check. + await plantPrincipalMail(`${BASE_COLUMNS}, + "created_at" timestamp NOT NULL DEFAULT now()`); + const failure = await bootFailure(applyMailboxMigrations(db, "public")); + expect(failure).toBeInstanceOf(SchemaTypeMismatchError); + expect((failure as SchemaTypeMismatchError).mismatches).toEqual([ + "principal_mail.subject is missing (expected text)", + ]); + expect(await ledgerRows()).toBe(0); + }); + }); + + test("a rejected boot leaves the NEXT boot still rejecting", async () => { + await fromEmpty(async ({ db }) => { + await plantPrincipalMail(`${BASE_COLUMNS}, + "created_at" timestamptz NOT NULL DEFAULT now()`); + await expect(applyMailboxMigrations(db, "public")).rejects.toThrow( + SchemaTypeMismatchError, + ); + // A guard that only fires on the first boot is one a restart disables. + await expect(applyMailboxMigrations(db, "public")).rejects.toThrow( + SchemaTypeMismatchError, + ); + expect(await ledgerRows()).toBe(0); + }); + }); +}); diff --git a/src/native-store.ts b/src/native-store.ts index 3e803a2..128ab01 100644 --- a/src/native-store.ts +++ b/src/native-store.ts @@ -24,7 +24,7 @@ function pgTextArrayLiteral(items: readonly string[]): string { * * Deliberately reads via plain tagged `sql`, not the `schema.ts` drizzle table * objects: those objects are pinned by `schema-check.ts` and - * `schema-ddl-parity.test.ts` to the columns the OLD read/write paths depend + * `migrations.test.ts` to the columns the OLD read/write paths depend * on, and this slice must not widen what those assert. * * `MailboxStore`'s mutating methods (`append`/`addFlags`/`removeFlags`/ diff --git a/src/schema-check.test.ts b/src/schema-check.test.ts deleted file mode 100644 index affd0ef..0000000 --- a/src/schema-check.test.ts +++ /dev/null @@ -1,214 +0,0 @@ -// `CREATE TABLE IF NOT EXISTS` matches on the table NAME and nothing else. A -// host that already owns a table called `mailbox` or `principal_mail` gets a -// silent no-op, a ledger row saying the migration applied, and from then on -// every read in this package decoding ITS columns through OUR codec. Nothing -// errors; the data is just wrong. -// -// Each case here plants exactly that situation — a pre-existing table with a -// divergent column type, and one with a column missing outright — and asserts -// the boot is REJECTED, and that the rejection leaves no ledger row behind to -// make the next boot skip the check. -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, applyMailboxMigrations } from "./migrations.js"; -import { - expectedColumnTypes, - SchemaTypeMismatchError, -} from "./schema-check.js"; -import { createHostControlPlane, TEST_DATABASE_URL } from "./test-helpers.js"; - -const admin = postgres(TEST_DATABASE_URL, { onnotice: () => {} }); - -beforeAll(async () => { - await createHostControlPlane(drizzle(admin)); -}); - -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 applyMailboxMigrations(drizzle(admin), "public"); - await admin.end(); -}); - -function handle() { - const client = postgres(TEST_DATABASE_URL, { onnotice: () => {} }); - return { client, db: drizzle(client) }; -} - -// The tables are hard-qualified to the "mailbox" schema, so a pre-existing -// host table can only shadow ours from INSIDE that schema: each case drops it, -// recreates it empty, and plants the conflicting table there. -async function inFreshSchema( - _name: string, - fn: (h: ReturnType) => Promise, -): Promise { - await admin.unsafe(`DROP SCHEMA IF EXISTS "mailbox" CASCADE`); - await admin.unsafe(`CREATE SCHEMA "mailbox"`); - const h = handle(); - try { - await fn(h); - } finally { - await h.client.end(); - } -} - -async function countOf(query: string): Promise { - const rows = (await admin.unsafe(query)) as unknown as { n: number }[]; - return rows[0]!.n; -} - -/** The ledger row count, or 0 when the ledger table itself does not exist. */ -async function ledgerRows(schema: string): Promise { - const exists = await countOf( - `SELECT count(*)::int AS n FROM information_schema.tables - WHERE table_schema = '${schema}' - AND table_name = 'corbits_mailbox_migrations'`, - ); - if (exists === 0) return 0; - return countOf( - `SELECT count(*)::int AS n FROM "${schema}"."corbits_mailbox_migrations"`, - ); -} - -/** - * 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. - */ -async function bootFailure(promise: Promise): Promise { - try { - await promise; - } catch (error) { - return error as Error; - } - throw new Error("expected the boot to be rejected, but it succeeded"); -} - -describe("expectedColumnTypes", () => { - test("is derived from the drizzle tables: the mail plane alone, since 0005 dropped the management table", () => { - const expected = expectedColumnTypes(); - const tables = new Set(expected.map((e) => e.table)); - expect(tables).toEqual(new Set(["principal_mail"])); - }); - - test("expects zoneless timestamps and text ids on every relevant column", () => { - const byKey = new Map( - expectedColumnTypes().map((e) => [`${e.table}.${e.column}`, e.dataType]), - ); - expect(byKey.get("principal_mail.created_at")).toBe( - "timestamp without time zone", - ); - expect(byKey.get("principal_mail.id")).toBe("text"); - expect(byKey.get("principal_mail.raw")).toBe("bytea"); - expect(byKey.get("principal_mail.refs")).toBe("jsonb"); - }); -}); - -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 applyMailboxMigrations(db, "public"); - expect(await ledgerRows("mailbox")).toBe(MIGRATIONS.length); - }); - }); - - test("rejects a pre-existing table whose column TYPE diverges", async () => { - const schema = "mailbox"; - await inFreshSchema(schema, async ({ db }) => { - // The host's own `principal_mail`: same name, `created_at` still a - // `timestamptz` — the exact type this package just moved off, and the one - // whose difference is invisible to every query until a non-UTC host - // serves the wrong page. - await admin.unsafe(` - CREATE TABLE "${schema}"."principal_mail" ( - "id" text PRIMARY KEY, - "tenant_id" text NOT NULL, - "principal_id" text NOT NULL, - "address" text NOT NULL, - "direction" text NOT NULL, - "raw" bytea NOT NULL, - "subject" text, - "from_address" text, - "message_key" text, - "refs" jsonb, - "created_at" timestamptz NOT NULL DEFAULT now() - )`); - 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, " + - "expected timestamp without time zone", - ]); - }); - // THE POINT: the boot was rejected inside the migration's own transaction, - // so the ledger row rolled back with it. Had it been recorded, the next - // boot would skip the migration entirely and sail past the mismatch. - expect(await ledgerRows(schema)).toBe(0); - }); - - test("rejects a pre-existing table with a column MISSING outright", async () => { - const schema = "mailbox"; - await inFreshSchema(schema, async ({ db }) => { - // A host `principal_mail` carrying the scope and the frame but not - // `subject`. Nothing errors on such a schema: `subject` is read through - // the codec, so every message would just quietly fall back to the frame - // for it. - // - // The missing column is one NO index covers — see the case below for why - // that distinction matters. `refs` is present here for exactly that - // reason: it is GIN-indexed as of `0003_mail_references`, so its absence - // is now rejected by the DDL rather than by this check. - await admin.unsafe(` - CREATE TABLE "${schema}"."principal_mail" ( - "id" text PRIMARY KEY, - "tenant_id" text NOT NULL, - "principal_id" text NOT NULL, - "address" text NOT NULL, - "direction" text NOT NULL, - "raw" bytea NOT NULL, - "from_address" text, - "message_key" text, - "refs" jsonb, - "created_at" timestamp NOT NULL DEFAULT now() - )`); - const failure = await bootFailure(applyMailboxMigrations(db, "public")); - expect(failure).toBeInstanceOf(SchemaTypeMismatchError); - expect((failure as SchemaTypeMismatchError).mismatches).toEqual([ - "principal_mail.subject is missing (expected text)", - ]); - }); - expect(await ledgerRows(schema)).toBe(0); - }); - - test("a rejected boot leaves the NEXT boot still rejecting", async () => { - const schema = "mailbox"; - await inFreshSchema(schema, async ({ db }) => { - await admin.unsafe(` - CREATE TABLE "${schema}"."principal_mail" ( - "id" text PRIMARY KEY, - "tenant_id" text NOT NULL, - "principal_id" text NOT NULL, - "address" text NOT NULL, - "direction" text NOT NULL, - "raw" bytea NOT NULL, - "from_address" text, - "message_key" text, - "refs" jsonb, - "created_at" timestamptz NOT NULL DEFAULT now() - )`); - 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(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 deleted file mode 100644 index 5abd3d6..0000000 --- a/src/schema-ddl-parity.test.ts +++ /dev/null @@ -1,136 +0,0 @@ -// The drizzle table object is a public export, so a host can point -// `drizzle-kit push`/`generate` at it. If it declares an index the migrations -// do not create — or creates one with a different column order — that host's -// schema silently diverges from the one this package's queries were planned -// against. This suite diffs the two after a real migration run. -import { afterAll, beforeAll, describe, expect, it } from "bun:test"; -import postgres from "postgres"; -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 { 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 -// is exactly what the migration built there — running the (idempotent) -// migrations here keeps the suite independent of which file ran first. The -// control-plane stub tables must exist for the migration's FKs to land. -const SCHEMA = "mailbox"; -const client = postgres(TEST_DATABASE_URL, { onnotice: () => {} }); - -beforeAll(async () => { - const db = drizzle(client); - await createHostControlPlane(db); - await applyMailboxMigrations(db, "public"); -}); - -afterAll(async () => { - await client.end(); -}); - -/** `name USING method(col asc, col desc)` plus `unique`/`partial` markers. */ -type IndexDescriptor = string; - -// eslint-disable-next-line @typescript-eslint/no-explicit-any -- one canonicalizer -// for both tables; `getTableConfig` is invariant in its table generic. -function declaredIndexes(table: any): IndexDescriptor[] { - const { indexes } = getTableConfig(table); - return indexes - .map((index) => { - const config = index.config; - const columns = config.columns - .map((column) => { - // Every index here is over plain columns; an expression index would - // have no `.name` and must be added to this canonicalizer before it - // can be compared at all, rather than silently comparing as blank. - const name = (column as { name?: string }).name; - if (name === undefined) { - throw new Error( - `index ${config.name} uses an expression column this parity check cannot canonicalize`, - ); - } - const order = - (column as { indexConfig?: { order?: string } }).indexConfig - ?.order ?? "asc"; - return `${name} ${order}`; - }) - .join(", "); - const flags = [ - config.unique === true ? "unique" : null, - config.where !== undefined ? "partial" : null, - ].filter((flag) => flag !== null); - const suffix = flags.length > 0 ? ` [${flags.join(" ")}]` : ""; - // The access method is part of the descriptor, not decoration: a GIN - // index and a btree index over the same column serve different queries, - // and `refs @> …` is only servable by the former. - const method = (config as { method?: string }).method ?? "btree"; - return `${config.name} USING ${method}(${columns})${suffix}`; - }) - .sort(); -} - -// `pg_get_indexdef` renders `CREATE [UNIQUE] INDEX ON USING -// ()[ WHERE ()]`, with DESC spelled out and ASC left -// implicit. -function canonicalizeIndexDef(def: string): IndexDescriptor { - const match = - /^CREATE (UNIQUE )?INDEX (\S+) ON \S+ USING (\S+) \((.*?)\)( WHERE .*)?$/.exec( - def, - ); - if (match === null) throw new Error(`unparsed index definition: ${def}`); - const [, unique, name, method, columnList, where] = match; - const columns = columnList! - .split(", ") - .map((column) => { - const desc = / DESC$/.test(column); - const bare = column.replace(/ (DESC|ASC)$/, "").replace(/ NULLS.*$/, ""); - return `${bare} ${desc ? "desc" : "asc"}`; - }) - .join(", "); - const flags = [ - unique !== undefined ? "unique" : null, - where !== undefined ? "partial" : null, - ].filter((flag) => flag !== null); - const suffix = flags.length > 0 ? ` [${flags.join(" ")}]` : ""; - return `${name} USING ${method}(${columns})${suffix}`; -} - -async function liveIndexes(table: string): Promise { - const rows = await drizzle(client).execute<{ indexdef: string }>(sql` - SELECT indexdef FROM pg_indexes - WHERE schemaname = ${SCHEMA} - AND tablename = ${table} - AND indexname <> ${`${table}_pkey`} - `); - return rows.map((row) => canonicalizeIndexDef(row.indexdef)).sort(); -} - -// `mailbox.mailbox` (the pre-native management layer) was dropped in -// `0005_drop_pre_native_columns` — the mail plane is the only table left to -// hold to this parity. -const TABLES = [{ name: "principal_mail", declared: principalMail }] as const; - -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)); - }); - } - - it("keeps the keyset access path on the mail plane, where the split left it", async () => { - expect(await liveIndexes("principal_mail")).toContain( - "principal_mail_tenant_id_principal_id_created_at_id_idx USING btree(tenant_id asc, principal_id asc, created_at desc, id desc)", - ); - }); - - it("no longer has a live mailbox.mailbox table", async () => { - const rows = await drizzle(client).execute<{ exists: boolean }>(sql` - SELECT EXISTS ( - SELECT 1 FROM information_schema.tables - WHERE table_schema = ${SCHEMA} AND table_name = 'mailbox' - ) AS "exists" - `); - expect(rows[0]!.exists).toBe(false); - }); -}); diff --git a/src/schema.ts b/src/schema.ts index 12c25a7..bc4aa15 100644 --- a/src/schema.ts +++ b/src/schema.ts @@ -139,7 +139,7 @@ export const principalMail = mailboxPgSchema.table( // a host that points `drizzle-kit push`/`generate` at it recreates exactly // what is declared here — an index declared here but dropped by a migration // (or declared with a different column order) silently reintroduces itself - // into that host's schema. `schema-ddl-parity.test.ts` diffs the two, + // into that host's schema. `migrations.test.ts` diffs the two, // for BOTH tables. (t) => [ // (tenant_id, principal_id, created_at DESC, id DESC) — matches the list query's From 0af12c1ef8d0582b6fc0dca996bb9daeb06613d8 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 24 Sep 2026 23:50:21 -0700 Subject: [PATCH 2/3] test(sse): one SSE and bus-lifecycle suite Subscription, heartbeat and backpressure cases were spread over three files with overlapping setup and two duplicated cases. Closes CL-9272. --- src/bus-unsubscribe.test.ts | 70 ------- src/schema.ts | 3 +- src/sse-heartbeat.test.ts | 169 --------------- src/sse-stream.test.ts | 260 ----------------------- src/sse.test.ts | 397 ++++++++++++++++++++++++++++++++++++ 5 files changed, 398 insertions(+), 501 deletions(-) delete mode 100644 src/bus-unsubscribe.test.ts delete mode 100644 src/sse-heartbeat.test.ts delete mode 100644 src/sse-stream.test.ts create mode 100644 src/sse.test.ts diff --git a/src/bus-unsubscribe.test.ts b/src/bus-unsubscribe.test.ts deleted file mode 100644 index c558b8d..0000000 --- a/src/bus-unsubscribe.test.ts +++ /dev/null @@ -1,70 +0,0 @@ -import { describe, test, expect } from "bun:test"; -import { createInMemoryMailboxEventBus } from "./bus.js"; - -const P1 = { tenantId: "t1", principalId: "p1" }; - -describe("in-memory bus unsubscribe", () => { - test("a double-called unsubscribe must not evict a later subscriber", () => { - const bus = createInMemoryMailboxEventBus(); - const a: string[] = []; - const offA = bus.subscribe(P1, (e) => a.push(e.id)); - offA(); // tab A closes -> set for p1 becomes empty and is deleted - const b: string[] = []; - bus.subscribe(P1, (e) => b.push(e.id)); // tab B opens (new Set) - offA(); // mount.ts calls unsubscribe twice (onAbort + finally) - bus.publish(P1, { type: "mailbox", id: "evt" }); - expect(b).toEqual(["evt"]); // B is still subscribed - }); -}); - -describe("in-memory bus tenant isolation", () => { - test("the same principalId under two tenants is two mailboxes, not one", () => { - // A principal identifier is only unique within its tenant; a bus keyed on - // it alone would fan tenantB's events into tenantA's same-named principal. - const bus = createInMemoryMailboxEventBus(); - const forA: string[] = []; - const forB: string[] = []; - bus.subscribe({ tenantId: "tenantA", principalId: "alice" }, (e) => - forA.push(e.id), - ); - bus.subscribe({ tenantId: "tenantB", principalId: "alice" }, (e) => - forB.push(e.id), - ); - bus.publish( - { tenantId: "tenantB", principalId: "alice" }, - { type: "mailbox", id: "evt-b" }, - ); - expect(forA).toEqual([]); - expect(forB).toEqual(["evt-b"]); - }); -}); - -describe("in-memory bus listener isolation", () => { - test("a throwing listener does not starve later subscribers of the same event", () => { - // Fan-out is best-effort per connection. One bad listener must not turn - // publish into "first throw wins" and skip every open tab behind it. - const bus = createInMemoryMailboxEventBus(); - const seen: string[] = []; - bus.subscribe(P1, () => { - throw new Error("listener boom"); - }); - bus.subscribe(P1, (e) => seen.push(e.id)); - expect(() => - bus.publish(P1, { type: "mailbox", id: "evt-isolated" }), - ).not.toThrow(); - expect(seen).toEqual(["evt-isolated"]); - }); - - test("a throw mid-set still delivers to every remaining listener", () => { - const bus = createInMemoryMailboxEventBus(); - const order: string[] = []; - bus.subscribe(P1, () => order.push("a")); - bus.subscribe(P1, () => { - order.push("b"); - throw new Error("mid"); - }); - bus.subscribe(P1, () => order.push("c")); - bus.publish(P1, { type: "mailbox", id: "x" }); - expect(order).toEqual(["a", "b", "c"]); - }); -}); diff --git a/src/schema.ts b/src/schema.ts index bc4aa15..c029bfc 100644 --- a/src/schema.ts +++ b/src/schema.ts @@ -139,8 +139,7 @@ export const principalMail = mailboxPgSchema.table( // a host that points `drizzle-kit push`/`generate` at it recreates exactly // what is declared here — an index declared here but dropped by a migration // (or declared with a different column order) silently reintroduces itself - // into that host's schema. `migrations.test.ts` diffs the two, - // for BOTH tables. + // into that host's schema. `migrations.test.ts` diffs the two. (t) => [ // (tenant_id, principal_id, created_at DESC, id DESC) — matches the list query's // ORDER BY and its row-value cursor seek exactly, so paging is an Index diff --git a/src/sse-heartbeat.test.ts b/src/sse-heartbeat.test.ts deleted file mode 100644 index ec4d9db..0000000 --- a/src/sse-heartbeat.test.ts +++ /dev/null @@ -1,169 +0,0 @@ -// The SSE keep-alive. `describeRoute` promises "a heartbeat comment every 25s" -// and the README repeats it, but nothing asserted a heartbeat frame was ever -// emitted — the loop could have been dead and every suite would still pass, -// because an idle stream and a broken keep-alive look identical until a proxy -// drops the connection in production. -import { describe, expect, test } from "bun:test"; -import { createMailboxRoutes } from "./mount.js"; -import { createInMemoryMailboxEventBus } from "./bus.js"; -import { writeMailboxMessage } from "./write.js"; -import { allowAllGrants, mountAs, withTestDb, seedScope } from "./test-helpers.js"; -import type { MailboxDb } from "./db.js"; - -const SCOPE = { tenantId: "t1", principalId: "p1" }; - -function stream(db: MailboxDb, heartbeatIntervalMs: number) { - const bus = createInMemoryMailboxEventBus(); - const app = mountAs( - SCOPE, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@t1.example", - deliver: () => {}, - heartbeatIntervalMs, - }), - ); - return { app, bus }; -} - -/** Read frames until `done(text)` is satisfied, or give up after `timeoutMs`. */ -async function readUntil( - body: ReadableStream, - done: (text: string) => boolean, - timeoutMs = 5_000, -): Promise { - const reader = body.getReader(); - const decoder = new TextDecoder(); - let text = ""; - const deadline = Date.now() + timeoutMs; - try { - while (Date.now() < deadline) { - const chunk = await Promise.race([ - reader.read(), - new Promise<{ value: undefined; done: true }>((resolve) => - setTimeout(() => resolve({ value: undefined, done: true }), 250), - ), - ]); - if (chunk.value !== undefined) text += decoder.decode(chunk.value); - if (done(text)) return text; - } - return text; - } finally { - await reader.cancel(); - } -} - -describe("SSE heartbeat", () => { - test("emits a heartbeat comment on an otherwise idle stream", async () => { - const db = await withTestDb(); - const { app } = stream(db, 20); - const res = await app.request("/me/inbox/events"); - expect(res.status).toBe(200); - - // Nothing is ever published on this stream, so a heartbeat is the ONLY - // thing that can arrive. - const text = await readUntil(res.body!, (t) => t.includes(": heartbeat")); - expect(text).toContain(": heartbeat\n\n"); - }); - - test("keeps emitting heartbeats rather than sending exactly one", async () => { - const db = await withTestDb(); - const { app } = stream(db, 20); - const res = await app.request("/me/inbox/events"); - - const text = await readUntil( - res.body!, - (t) => t.split(": heartbeat").length - 1 >= 3, - ); - expect(text.split(": heartbeat").length - 1).toBeGreaterThanOrEqual(3); - }); - - test("heartbeats are comments, so they never look like mailbox events", async () => { - const db = await withTestDb(); - const { app } = stream(db, 20); - const res = await app.request("/me/inbox/events"); - - const text = await readUntil(res.body!, (t) => t.includes(": heartbeat")); - // A client parsing this must not see a nameless event or stray data — an - // SSE comment line starts with ':' and carries neither. - expect(text).not.toContain("event:"); - expect(text).not.toContain("data:"); - }); - - test("a real event still comes through while heartbeats are running", async () => { - const db = await withTestDb(); - await seedScope(db, SCOPE.tenantId, SCOPE.principalId); - const { app, bus } = stream(db, 20); - const res = await app.request("/me/inbox/events"); - - const body = res.body!; - // Registered on the next tick so the handler's subscription exists first. - setTimeout(() => { - void writeMailboxMessage( - db, - { - ...SCOPE, - address: "p@b.c", - fromAddress: "a@b.c", - subject: "hi", - body: "hello", - }, - bus, - ); - }, 60); - - const text = await readUntil(body, (t) => t.includes("event: mailbox")); - expect(text).toContain("event: mailbox"); - expect(text).toContain(": heartbeat"); - }); - - test("the default interval is the documented 25s, not the test override", async () => { - const db = await withTestDb(); - const bus = createInMemoryMailboxEventBus(); - const app = mountAs( - SCOPE, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@t1.example", - deliver: () => {}, - }), - ); - const res = await app.request("/me/inbox/events"); - - // With no override, nothing may arrive within a second — otherwise the - // override is leaking into the default and the 25s figure is fiction. - const text = await readUntil( - res.body!, - (t) => t.includes(": heartbeat"), - 1_000, - ); - expect(text).toBe(""); - }); - - test("mount refuses a non-positive or non-finite heartbeatIntervalMs", async () => { - // Zero/negative would spin a tight sleep/write loop per open connection; - // NaN/Infinity are the same class of host misconfiguration. Fail at mount, - // not on the first request, same as a bad vocabulary. - const db = await withTestDb(); - const bus = createInMemoryMailboxEventBus(); - for (const heartbeatIntervalMs of [0, -1, Number.NaN, Number.POSITIVE_INFINITY]) { - expect(() => - mountAs( - SCOPE, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@t1.example", - deliver: () => {}, - heartbeatIntervalMs, - }), - ), - ).toThrow(RangeError); - } - }); -}); diff --git a/src/sse-stream.test.ts b/src/sse-stream.test.ts deleted file mode 100644 index 92c23e4..0000000 --- a/src/sse-stream.test.ts +++ /dev/null @@ -1,260 +0,0 @@ -import { describe, test, expect, spyOn } from "bun:test"; -import { SSEStreamingApi } from "hono/streaming"; -import { createMailboxRoutes, MAX_PENDING_SSE_EVENTS } from "./mount.js"; -import { - createInMemoryMailboxEventBus, - type MailboxEventBus, - type MailboxEventScope, -} from "./bus.js"; -import { - allowAllGrants, - mountAs, - withTestDb, - seedScope, -} from "./test-helpers.js"; -import { writeMailboxMessage } from "./write.js"; - -describe("SSE stream", () => { - test("delivers an event to the subscribed principalId only", async () => { - const db = await withTestDb(); - await seedScope(db, "t1", "p1"); - const bus = createInMemoryMailboxEventBus(); - const app = mountAs( - { tenantId: "t1", principalId: "p1" }, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@t1.example", - deliver: () => {}, - }), - ); - const res = await app.request("/me/inbox/events"); - expect(res.status).toBe(200); - expect(res.headers.get("content-type")).toContain("text/event-stream"); - - const reader = res.body!.getReader(); - // give the handler a tick to register its subscription - await new Promise((r) => setTimeout(r, 50)); - await writeMailboxMessage( - db, - { - tenantId: "t1", - principalId: "p1", - address: "p@b.c", - fromAddress: "a@b.c", - subject: "hi", - body: "hello", - }, - bus, - ); - const chunk = await Promise.race([ - reader.read().then((r) => new TextDecoder().decode(r.value)), - new Promise((r) => setTimeout(() => r("__TIMEOUT__"), 3000)), - ]); - await reader.cancel(); - expect(chunk).not.toBe("__TIMEOUT__"); - expect(chunk).toContain("mailbox"); - }); - - test("unsubscribe isolation: one closed stream does not stop another", async () => { - const bus = createInMemoryMailboxEventBus(); - const scope = { tenantId: "t1", principalId: "p1" }; - const seenA: string[] = []; - const seenB: string[] = []; - const offA = bus.subscribe(scope, (e) => seenA.push(e.id)); - bus.subscribe(scope, (e) => seenB.push(e.id)); - offA(); - bus.publish(scope, { type: "mailbox", id: "x" }); - expect(seenA).toEqual([]); - expect(seenB).toEqual(["x"]); - }); - - test("tenant isolation end-to-end: same principalId, different tenant, no event", async () => { - const db = await withTestDb(); - await seedScope(db, "tenantA", "alice"); - await seedScope(db, "tenantB", "alice"); - const bus = createInMemoryMailboxEventBus(); - const app = mountAs( - { tenantId: "tenantA", principalId: "alice" }, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@tenantA.example", - deliver: () => {}, - }), - ); - const res = await app.request("/me/inbox/events"); - expect(res.status).toBe(200); - const reader = res.body!.getReader(); - // give the handler a tick to register its subscription - await new Promise((r) => setTimeout(r, 50)); - // ONE pending read for both races: a losing race branch would otherwise - // keep an orphaned read holding the next chunk. - const firstChunk = reader - .read() - .then((r) => new TextDecoder().decode(r.value)); - // tenantB's alice gets mail; tenantA's stream must stay silent. - await writeMailboxMessage( - db, - { - tenantId: "tenantB", - principalId: "alice", - address: "alice@b.example", - fromAddress: "a@b.example", - subject: "for the OTHER alice", - body: "hello", - }, - bus, - ); - const crossTenant = await Promise.race([ - firstChunk, - new Promise((r) => setTimeout(() => r("__TIMEOUT__"), 300)), - ]); - expect(crossTenant).toBe("__TIMEOUT__"); - // And the stream is still live for its OWN scope, so the silence above was - // isolation, not a dead connection. - await writeMailboxMessage( - db, - { - tenantId: "tenantA", - principalId: "alice", - address: "alice@a.example", - fromAddress: "a@a.example", - subject: "for this alice", - body: "hello", - }, - bus, - ); - const ownTenant = await Promise.race([ - firstChunk, - new Promise((r) => setTimeout(() => r("__TIMEOUT__"), 3000)), - ]); - await reader.cancel(); - expect(ownTenant).not.toBe("__TIMEOUT__"); - expect(ownTenant).toContain("mailbox"); - }); - - test("a consumer that stops reading is disconnected at the pending cap, not buffered for", async () => { - const db = await withTestDb(); - await seedScope(db, "t1", "p1"); - const bus = createInMemoryMailboxEventBus(); - const scope = { tenantId: "t1", principalId: "p1" }; - const app = mountAs( - scope, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@t1.example", - deliver: () => {}, - // Short heartbeat so the handler notices the overflow-close promptly. - heartbeatIntervalMs: 50, - }), - ); - const res = await app.request("/me/inbox/events"); - expect(res.status).toBe(200); - const reader = res.body!.getReader(); - // give the handler a tick to register its subscription - await new Promise((r) => setTimeout(r, 50)); - // The client never reads. Every event is a nudge the client would refetch - // from Postgres anyway, so past the cap the connection must close rather - // than park one pending write per event forever. - for (let i = 0; i <= MAX_PENDING_SSE_EVENTS + 5; i++) { - bus.publish(scope, { type: "mailbox", id: `evt-${i}` }); - } - const deadline = Date.now() + 3000; - let done = false; - while (!done && Date.now() < deadline) { - const result = await Promise.race([ - reader.read(), - new Promise<{ done: boolean }>((r) => - setTimeout(() => r({ done: false }), 200), - ), - ]); - done = result.done; - } - expect(done).toBe(true); - // The subscription is gone with the connection: publishing again reaches - // nobody and, more to the point, throws nothing. - bus.publish(scope, { type: "mailbox", id: "after-close" }); - }); - - test("a write failure while draining does not become an unhandled rejection", async () => { - // drain() is fired with `void`; a rejected writeSSE must be caught inside - // so it closes the stream cleanly instead of escaping as an unhandled - // rejection that the process (or a host) has to notice later. - // - // Hono's StreamingApi.write swallows writer errors, so cancelling the - // client never makes writeSSE reject in practice. Stub writeSSE to force - // the drain catch path and assert cleanup + unsubscribe. - const db = await withTestDb(); - await seedScope(db, "t1", "p1"); - const scope = { tenantId: "t1", principalId: "p1" }; - const realBus = createInMemoryMailboxEventBus(); - let activeSubs = 0; - const bus: MailboxEventBus = { - publish(s, event) { - realBus.publish(s, event); - }, - subscribe(s: MailboxEventScope, listener) { - activeSubs++; - const off = realBus.subscribe(s, listener); - let done = false; - return () => { - // mount unsubscribes from onAbort and finally; count once. - if (done) return; - done = true; - activeSubs--; - off(); - }; - }, - }; - const app = mountAs( - scope, - createMailboxRoutes({ - db, - requireGrant: allowAllGrants, - bus, - senderAddressFor: () => "sender@t1.example", - deliver: () => {}, - // Short heartbeat so the loop notices `closed` and runs finally promptly. - heartbeatIntervalMs: 50, - }), - ); - const writeSSE = spyOn( - SSEStreamingApi.prototype, - "writeSSE", - ).mockRejectedValue(new Error("simulated socket death")); - const rejections: unknown[] = []; - const onRejection = (reason: unknown) => { - rejections.push(reason); - }; - process.on("unhandledRejection", onRejection); - try { - const res = await app.request("/me/inbox/events"); - expect(res.status).toBe(200); - const reader = res.body!.getReader(); - // give the handler a tick to register its subscription - await new Promise((r) => setTimeout(r, 50)); - expect(activeSubs).toBe(1); - // One publish is enough: drain awaits writeSSE, which rejects. - bus.publish(scope, { type: "mailbox", id: "force-drain-reject" }); - // Wait past one heartbeat so the loop exits on `closed` and finally - // unsubscribes. - await new Promise((r) => setTimeout(r, 200)); - expect(writeSSE).toHaveBeenCalled(); - expect(rejections).toEqual([]); - expect(activeSubs).toBe(0); - // Publish after the failure path must still be a no-op, not a throw. - expect(() => - bus.publish(scope, { type: "mailbox", id: "still-safe" }), - ).not.toThrow(); - await reader.cancel(); - } finally { - process.off("unhandledRejection", onRejection); - writeSSE.mockRestore(); - } - }); -}); diff --git a/src/sse.test.ts b/src/sse.test.ts new file mode 100644 index 0000000..96c10d2 --- /dev/null +++ b/src/sse.test.ts @@ -0,0 +1,397 @@ +import { describe, expect, spyOn, test } from "bun:test"; +import { SSEStreamingApi } from "hono/streaming"; +import { createMailboxRoutes, MAX_PENDING_SSE_EVENTS } from "./mount.js"; +import { + createInMemoryMailboxEventBus, + type MailboxEventBus, + type MailboxEventScope, +} from "./bus.js"; +import { writeMailboxMessage } from "./write.js"; +import { allowAllGrants, mountAs, withTestDb, seedScope } from "./test-helpers.js"; +import type { MailboxDb } from "./db.js"; + +const SCOPE = { tenantId: "t1", principalId: "p1" }; + +function routes( + db: MailboxDb, + bus: MailboxEventBus, + scope: MailboxEventScope, + heartbeatIntervalMs?: number, +) { + return mountAs( + scope, + createMailboxRoutes({ + db, + requireGrant: allowAllGrants, + bus, + senderAddressFor: () => `sender@${scope.tenantId}.example`, + deliver: () => {}, + heartbeatIntervalMs, + }), + ); +} + +/** Read frames until `done(text)` is satisfied, or give up after `timeoutMs`. */ +async function readUntil( + body: ReadableStream, + done: (text: string) => boolean, + timeoutMs = 5_000, +): Promise { + const reader = body.getReader(); + const decoder = new TextDecoder(); + let text = ""; + const deadline = Date.now() + timeoutMs; + try { + while (Date.now() < deadline) { + const chunk = await Promise.race([ + reader.read(), + new Promise<{ value: undefined; done: true }>((resolve) => + setTimeout(() => resolve({ value: undefined, done: true }), 250), + ), + ]); + if (chunk.value !== undefined) text += decoder.decode(chunk.value); + if (done(text)) return text; + } + return text; + } finally { + await reader.cancel(); + } +} + +describe("in-memory bus", () => { + test("closing one subscriber leaves the others on the same scope", () => { + const bus = createInMemoryMailboxEventBus(); + const a: string[] = []; + const b: string[] = []; + const offA = bus.subscribe(SCOPE, (e) => a.push(e.id)); + bus.subscribe(SCOPE, (e) => b.push(e.id)); + offA(); + bus.publish(SCOPE, { type: "mailbox", id: "x" }); + expect(a).toEqual([]); + expect(b).toEqual(["x"]); + }); + + test("a double-called unsubscribe must not evict a later subscriber", () => { + const bus = createInMemoryMailboxEventBus(); + const a: string[] = []; + const offA = bus.subscribe(SCOPE, (e) => a.push(e.id)); + offA(); // tab A closes -> set for p1 becomes empty and is deleted + const b: string[] = []; + bus.subscribe(SCOPE, (e) => b.push(e.id)); // tab B opens (new Set) + offA(); // mount.ts calls unsubscribe twice (onAbort + finally) + bus.publish(SCOPE, { type: "mailbox", id: "evt" }); + expect(a).toEqual([]); + expect(b).toEqual(["evt"]); // B is still subscribed + }); + + test("a throw mid-set still delivers to every remaining listener", () => { + // Fan-out is best-effort per connection. One bad listener must not turn + // publish into "first throw wins" and skip every open tab behind it. + const bus = createInMemoryMailboxEventBus(); + const order: string[] = []; + bus.subscribe(SCOPE, () => order.push("a")); + bus.subscribe(SCOPE, () => { + order.push("b"); + throw new Error("mid"); + }); + bus.subscribe(SCOPE, () => order.push("c")); + expect(() => + bus.publish(SCOPE, { type: "mailbox", id: "x" }), + ).not.toThrow(); + expect(order).toEqual(["a", "b", "c"]); + }); +}); + +describe("SSE stream", () => { + test("delivers an event to the subscribed principalId only", async () => { + const db = await withTestDb(); + await seedScope(db, "t1", "p1"); + const bus = createInMemoryMailboxEventBus(); + const app = routes(db, bus, SCOPE); + const res = await app.request("/me/inbox/events"); + expect(res.status).toBe(200); + expect(res.headers.get("content-type")).toContain("text/event-stream"); + + const reader = res.body!.getReader(); + // give the handler a tick to register its subscription + await new Promise((r) => setTimeout(r, 50)); + await writeMailboxMessage( + db, + { + tenantId: "t1", + principalId: "p1", + address: "p@b.c", + fromAddress: "a@b.c", + subject: "hi", + body: "hello", + }, + bus, + ); + const chunk = await Promise.race([ + reader.read().then((r) => new TextDecoder().decode(r.value)), + new Promise((r) => setTimeout(() => r("__TIMEOUT__"), 3000)), + ]); + await reader.cancel(); + expect(chunk).not.toBe("__TIMEOUT__"); + expect(chunk).toContain("mailbox"); + }); + + test("tenant isolation end-to-end: same principalId, different tenant, no event", async () => { + const db = await withTestDb(); + await seedScope(db, "tenantA", "alice"); + await seedScope(db, "tenantB", "alice"); + const bus = createInMemoryMailboxEventBus(); + const app = routes(db, bus, { tenantId: "tenantA", principalId: "alice" }); + const res = await app.request("/me/inbox/events"); + expect(res.status).toBe(200); + const reader = res.body!.getReader(); + // give the handler a tick to register its subscription + await new Promise((r) => setTimeout(r, 50)); + // ONE pending read for both races: a losing race branch would otherwise + // keep an orphaned read holding the next chunk. + const firstChunk = reader + .read() + .then((r) => new TextDecoder().decode(r.value)); + // tenantB's alice gets mail; tenantA's stream must stay silent. + await writeMailboxMessage( + db, + { + tenantId: "tenantB", + principalId: "alice", + address: "alice@b.example", + fromAddress: "a@b.example", + subject: "for the OTHER alice", + body: "hello", + }, + bus, + ); + const crossTenant = await Promise.race([ + firstChunk, + new Promise((r) => setTimeout(() => r("__TIMEOUT__"), 300)), + ]); + expect(crossTenant).toBe("__TIMEOUT__"); + // And the stream is still live for its OWN scope, so the silence above was + // isolation, not a dead connection. + await writeMailboxMessage( + db, + { + tenantId: "tenantA", + principalId: "alice", + address: "alice@a.example", + fromAddress: "a@a.example", + subject: "for this alice", + body: "hello", + }, + bus, + ); + const ownTenant = await Promise.race([ + firstChunk, + new Promise((r) => setTimeout(() => r("__TIMEOUT__"), 3000)), + ]); + await reader.cancel(); + expect(ownTenant).not.toBe("__TIMEOUT__"); + expect(ownTenant).toContain("mailbox"); + }); + + test("a consumer that stops reading is disconnected at the pending cap, not buffered for", async () => { + const db = await withTestDb(); + await seedScope(db, "t1", "p1"); + const bus = createInMemoryMailboxEventBus(); + // Short heartbeat so the handler notices the overflow-close promptly. + const app = routes(db, bus, SCOPE, 50); + const res = await app.request("/me/inbox/events"); + expect(res.status).toBe(200); + const reader = res.body!.getReader(); + // give the handler a tick to register its subscription + await new Promise((r) => setTimeout(r, 50)); + // The client never reads. Every event is a nudge the client would refetch + // from Postgres anyway, so past the cap the connection must close rather + // than park one pending write per event forever. + for (let i = 0; i <= MAX_PENDING_SSE_EVENTS + 5; i++) { + bus.publish(SCOPE, { type: "mailbox", id: `evt-${i}` }); + } + const deadline = Date.now() + 3000; + let done = false; + while (!done && Date.now() < deadline) { + const result = await Promise.race([ + reader.read(), + new Promise<{ done: boolean }>((r) => + setTimeout(() => r({ done: false }), 200), + ), + ]); + done = result.done; + } + expect(done).toBe(true); + // The subscription is gone with the connection: publishing again reaches + // nobody and, more to the point, throws nothing. + bus.publish(SCOPE, { type: "mailbox", id: "after-close" }); + }); + + test("a write failure while draining does not become an unhandled rejection", async () => { + // drain() is fired with `void`; a rejected writeSSE must be caught inside + // so it closes the stream cleanly instead of escaping as an unhandled + // rejection that the process (or a host) has to notice later. + // + // Hono's StreamingApi.write swallows writer errors, so cancelling the + // client never makes writeSSE reject in practice. Stub writeSSE to force + // the drain catch path and assert cleanup + unsubscribe. + const db = await withTestDb(); + await seedScope(db, "t1", "p1"); + const realBus = createInMemoryMailboxEventBus(); + let activeSubs = 0; + const bus: MailboxEventBus = { + publish(s, event) { + realBus.publish(s, event); + }, + subscribe(s: MailboxEventScope, listener) { + activeSubs++; + const off = realBus.subscribe(s, listener); + let done = false; + return () => { + // mount unsubscribes from onAbort and finally; count once. + if (done) return; + done = true; + activeSubs--; + off(); + }; + }, + }; + // Short heartbeat so the loop notices `closed` and runs finally promptly. + const app = routes(db, bus, SCOPE, 50); + const writeSSE = spyOn( + SSEStreamingApi.prototype, + "writeSSE", + ).mockRejectedValue(new Error("simulated socket death")); + const rejections: unknown[] = []; + const onRejection = (reason: unknown) => { + rejections.push(reason); + }; + process.on("unhandledRejection", onRejection); + try { + const res = await app.request("/me/inbox/events"); + expect(res.status).toBe(200); + const reader = res.body!.getReader(); + // give the handler a tick to register its subscription + await new Promise((r) => setTimeout(r, 50)); + expect(activeSubs).toBe(1); + // One publish is enough: drain awaits writeSSE, which rejects. + bus.publish(SCOPE, { type: "mailbox", id: "force-drain-reject" }); + // Wait past one heartbeat so the loop exits on `closed` and finally + // unsubscribes. + await new Promise((r) => setTimeout(r, 200)); + expect(writeSSE).toHaveBeenCalled(); + expect(rejections).toEqual([]); + expect(activeSubs).toBe(0); + // Publish after the failure path must still be a no-op, not a throw. + expect(() => + bus.publish(SCOPE, { type: "mailbox", id: "still-safe" }), + ).not.toThrow(); + await reader.cancel(); + } finally { + process.off("unhandledRejection", onRejection); + writeSSE.mockRestore(); + } + }); +}); + +// `describeRoute` promises "a heartbeat comment every 25s". An idle stream and +// a dead keep-alive look identical until a proxy drops the connection, so the +// frames themselves are asserted. +describe("SSE heartbeat", () => { + test("emits a heartbeat comment on an otherwise idle stream", async () => { + const db = await withTestDb(); + const app = routes(db, createInMemoryMailboxEventBus(), SCOPE, 20); + const res = await app.request("/me/inbox/events"); + expect(res.status).toBe(200); + + // Nothing is ever published on this stream, so a heartbeat is the ONLY + // thing that can arrive. + const text = await readUntil(res.body!, (t) => t.includes(": heartbeat")); + expect(text).toContain(": heartbeat\n\n"); + }); + + test("keeps emitting heartbeats rather than sending exactly one", async () => { + const db = await withTestDb(); + const app = routes(db, createInMemoryMailboxEventBus(), SCOPE, 20); + const res = await app.request("/me/inbox/events"); + + const text = await readUntil( + res.body!, + (t) => t.split(": heartbeat").length - 1 >= 3, + ); + expect(text.split(": heartbeat").length - 1).toBeGreaterThanOrEqual(3); + }); + + test("heartbeats are comments, so they never look like mailbox events", async () => { + const db = await withTestDb(); + const app = routes(db, createInMemoryMailboxEventBus(), SCOPE, 20); + const res = await app.request("/me/inbox/events"); + + const text = await readUntil(res.body!, (t) => t.includes(": heartbeat")); + // A client parsing this must not see a nameless event or stray data — an + // SSE comment line starts with ':' and carries neither. + expect(text).not.toContain("event:"); + expect(text).not.toContain("data:"); + }); + + test("a real event still comes through while heartbeats are running", async () => { + const db = await withTestDb(); + await seedScope(db, SCOPE.tenantId, SCOPE.principalId); + const bus = createInMemoryMailboxEventBus(); + const app = routes(db, bus, SCOPE, 20); + const res = await app.request("/me/inbox/events"); + + const body = res.body!; + // Registered on the next tick so the handler's subscription exists first. + setTimeout(() => { + void writeMailboxMessage( + db, + { + ...SCOPE, + address: "p@b.c", + fromAddress: "a@b.c", + subject: "hi", + body: "hello", + }, + bus, + ); + }, 60); + + const text = await readUntil(body, (t) => t.includes("event: mailbox")); + expect(text).toContain("event: mailbox"); + expect(text).toContain(": heartbeat"); + }); + + test("the default interval is the documented 25s, not the test override", async () => { + const db = await withTestDb(); + const app = routes(db, createInMemoryMailboxEventBus(), SCOPE); + const res = await app.request("/me/inbox/events"); + + // With no override, nothing may arrive within a second — otherwise the + // override is leaking into the default and the 25s figure is fiction. + const text = await readUntil( + res.body!, + (t) => t.includes(": heartbeat"), + 1_000, + ); + expect(text).toBe(""); + }); + + test("mount refuses a non-positive or non-finite heartbeatIntervalMs", async () => { + // Zero/negative would spin a tight sleep/write loop per open connection; + // NaN/Infinity are the same class of host misconfiguration. Fail at mount, + // not on the first request, same as a bad vocabulary. + const db = await withTestDb(); + const bus = createInMemoryMailboxEventBus(); + for (const heartbeatIntervalMs of [ + 0, + -1, + Number.NaN, + Number.POSITIVE_INFINITY, + ]) { + expect(() => routes(db, bus, SCOPE, heartbeatIntervalMs)).toThrow( + RangeError, + ); + } + }); +}); From 2be4853e11c9c16f2b06e00794c38a263eb133b6 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 25 Sep 2026 07:48:18 -0700 Subject: [PATCH 3/3] docs(native-store): say why reads use raw sql today --- src/native-store.ts | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/src/native-store.ts b/src/native-store.ts index 128ab01..1dc5d2a 100644 --- a/src/native-store.ts +++ b/src/native-store.ts @@ -22,10 +22,9 @@ function pgTextArrayLiteral(items: readonly string[]): string { * functions already assume (`mailboxName` names the one mailbox * `store.messages` holds). * - * Deliberately reads via plain tagged `sql`, not the `schema.ts` drizzle table - * objects: those objects are pinned by `schema-check.ts` and - * `migrations.test.ts` to the columns the OLD read/write paths depend - * on, and this slice must not widen what those assert. + * Reads via plain tagged `sql`, not the `schema.ts` drizzle table objects: + * `schema-check.ts` asserts those objects against the live columns at boot, + * so they stay limited to what that check needs. * * `MailboxStore`'s mutating methods (`append`/`addFlags`/`removeFlags`/ * `remove`) are synchronous in `@intx/mailbox`'s interface — an in-memory backing