From 43881706ad4a86c7614af56bc3b3181acdc526fd Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 3 Aug 2026 04:51:59 +0000 Subject: [PATCH 1/3] feat(email): detach provider index from legacy graph Co-authored-by: Kent C. Dodds --- .../contributing/architecture/data-storage.md | 69 ++++-- ...2-email-outbound-provider-index-detach.sql | 51 +++++ .../admin/mailbox-maintenance.node.test.ts | 4 +- .../worker/src/admin/mailbox-maintenance.ts | 17 +- packages/worker/src/app/retention.ts | 8 +- .../email-user-graph-authority.node.test.ts | 120 ++++++++++ .../src/email/email-user-graph-authority.ts | 51 +++++ ...ovider-index-detach-migration.node.test.ts | 210 ++++++++++++++++++ .../src/email/outbound-provider-index.ts | 14 ++ .../outbound-provider-index.workers.test.ts | 38 ++-- packages/worker/src/email/repo.ts | 14 +- .../system-email-authority.workers.test.ts | 18 +- packages/worker/src/email/test-schema.ts | 17 +- .../admin-mailbox-maintenance.node.test.ts | 1 + .../admin/admin-mailbox-maintenance.ts | 5 + .../email-user-graph-statement-classifier.ts | 79 +++++++ tools/migration-ledger.json | 4 + 17 files changed, 658 insertions(+), 62 deletions(-) create mode 100644 packages/worker/migrations/0132-email-outbound-provider-index-detach.sql create mode 100644 packages/worker/src/email/email-user-graph-authority.node.test.ts create mode 100644 packages/worker/src/email/email-user-graph-authority.ts create mode 100644 packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts create mode 100644 packages/worker/src/test-support/email-user-graph-statement-classifier.ts diff --git a/docs/contributing/architecture/data-storage.md b/docs/contributing/architecture/data-storage.md index 3232aaca64..1346a67884 100644 --- a/docs/contributing/architecture/data-storage.md +++ b/docs/contributing/architecture/data-storage.md @@ -1353,11 +1353,30 @@ cover the high-risk live surface today. `mailbox_parity_mismatch_count`, account not marked for deletion) **and** the default-off `mailbox-read-cutover` flag is enabled per user. Live gate evaluations record flag exposures (session cache or cutover memo chokepoint). -4. **D1 write-off / event retirement** — stop writing moved user-mail metadata - to D1; retire dual-write and event/mirror machinery used only for the - migration. -5. **Later contract migrations** — drop retired D1 user-mail tables/columns only - after verification. No premature schema deletion. +4. **D1 write-off / event retirement** — retire dual-write and event/mirror + machinery used only for the migration as part of the ordered step 5 cutover + below. +5. **USER graph contract (5a then 5b)** — deploy the 5a prerequisite first: + migration `0132-email-outbound-provider-index-detach.sql` atomically rebuilds + the retained global provider index without its cross-store-invalid + `email_messages` foreign key, and the authority guard contract lands without + changing live paths. Verify production reports + `status.outboundProviderIndex.foreignKeyDetached: true` before proceeding. A + later 5a cutover then moves **all** USER inbound, outbound, provider-index + synchronization, classification, explicit delete, and retention graph + mutations to Mailbox-only authority and freezes the four shared USER graph + tables (`email_threads`, `email_messages`, `email_attachments`, and + `email_delivery_events`). After the frozen-table verification window, take + and verify a fresh production backup; only then may 5b drop the retired USER + graph tables and migration-only machinery. No USER row or graph table is + destructively removed by the prerequisite deployment. + + **Rollback caveat:** after Mailbox-only writes begin, the frozen D1 USER + graph is stale and must never be re-enabled as authority by a code rollback. + A rollback must remain Mailbox-authoritative (or explicitly rebuild D1 from + Mailbox before restoring old code). After 5b drops the tables, rollback also + requires the verified fresh backup/schema restore; reverting application code + alone is not safe. The every-5-minute `mailbox_parity` scheduled lane (queue-isolated sibling in `scheduled-lanes.ts`) owns backfill of all owner messages and delivery events, @@ -1381,29 +1400,33 @@ below. - **Provider-message reverse lookup** — outbound Cloudflare sending webhooks resolve owner/message through the derived D1 table `email_outbound_provider_index` (migration - `0128-email-outbound-provider-index.sql`), keyed by + `0128-email-outbound-provider-index.sql`, detached by prerequisite migration + `0132-email-outbound-provider-index-detach.sql`), keyed by `(provider, provider_message_id)` with `user_id`, `message_id`, `inbox_id`, and created/updated timestamps (indexes on `user_id` and unique `message_id`). `email_messages.provider_message_id` remains authoritative: outbound inserts with a provider id and `updateEmailMessageDelivery` commit the message row plus index sync in one `db.batch` (index owner/inbox fields come from the - authoritative message row, never caller input). `message_id` references - `email_messages(id)` with `ON DELETE CASCADE`, so message deletes clear index - rows; account deletion still inventories/deletes by `user_id` for coverage. - Outbound send separates provider acceptance from terminal D1/index - persistence: once the provider returns a `providerMessageId`, persistence uses - bounded D1 retries and must not mark the message `failed`, clear the id, or - resend. Account export treats the table as derived global lookup - (`includeInExport: false` / `derivedData.email_outbound_provider_index`) - because authoritative outbound message rows are already exported. - `recordProviderEmailDeliveryEvent` resolves index-first, then loads the - owner-scoped message by `user_id`/`message_id` (no full-table provider scan). - System outbound is unsupported, and the verified `no-system-provider-links` - disposition means `system:email` rows are never added to this legacy-FK index. - Aggregate parity (`loadOutboundProviderIndexParityReport`; counts only) is - surfaced on `admin_mailbox_maintenance` `status.outboundProviderIndex` for - production verification. Contextless provider-id reverse lookups must not - enumerate per-user Mailbox objects or resolve a system compatibility mirror. + authoritative message row, never caller input). `message_id` is an opaque + owner-scoped key with no foreign key to `email_messages`: legacy message + deletion cannot cascade across the upcoming Mailbox/D1 store boundary. + Explicit message and account deletion therefore remove index rows by + `message_id`/`user_id`. Outbound send separates provider acceptance from + terminal D1/index persistence: once the provider returns a + `providerMessageId`, persistence uses bounded D1 retries and must not mark the + message `failed`, clear the id, or resend. Account export treats the table as + derived global lookup (`includeInExport: false` / + `derivedData.email_outbound_provider_index`) because authoritative outbound + message rows are already exported. `recordProviderEmailDeliveryEvent` resolves + index-first, then loads the owner-scoped message by `user_id`/`message_id` (no + full-table provider scan). System outbound is unsupported, and the verified + `no-system-provider-links` disposition means `system:email` rows are never + added to this global index. Aggregate parity + (`loadOutboundProviderIndexParityReport`; counts only) is surfaced on + `admin_mailbox_maintenance` `status.outboundProviderIndex` for production + verification, alongside the live schema check `foreignKeyDetached`. + Contextless provider-id reverse lookups must not enumerate per-user Mailbox + objects or resolve a system compatibility mirror. ### Inbound durability boundary (USER Mailbox authority) diff --git a/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql b/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql new file mode 100644 index 0000000000..e201d6d2b2 --- /dev/null +++ b/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql @@ -0,0 +1,51 @@ +-- Step 5a prerequisite: keep the global provider reverse lookup in D1 while +-- detaching it from the legacy USER email_messages graph. Wrangler applies this +-- migration file atomically, so a failed copy restores the original table. +-- +-- system:email has dedicated D1 graph authority and does not support outbound +-- provider links. USER message authority moves to Mailbox in the next slice; +-- message_id therefore remains an opaque owner-scoped lookup key, not a +-- cross-store foreign key. + +CREATE TABLE email_outbound_provider_index_next ( + provider TEXT NOT NULL CHECK (length(provider) > 0), + provider_message_id TEXT NOT NULL CHECK (length(provider_message_id) > 0), + user_id TEXT NOT NULL CHECK ( + length(user_id) > 0 AND user_id <> 'system:email' + ), + message_id TEXT NOT NULL CHECK (length(message_id) > 0), + inbox_id TEXT CHECK (inbox_id IS NULL OR length(inbox_id) > 0), + created_at TEXT NOT NULL CHECK (length(created_at) > 0), + updated_at TEXT NOT NULL CHECK (length(updated_at) > 0), + PRIMARY KEY (provider, provider_message_id) +); + +INSERT INTO email_outbound_provider_index_next ( + provider, + provider_message_id, + user_id, + message_id, + inbox_id, + created_at, + updated_at +) +SELECT + provider, + provider_message_id, + user_id, + message_id, + inbox_id, + created_at, + updated_at +FROM email_outbound_provider_index; + +DROP TABLE email_outbound_provider_index; + +ALTER TABLE email_outbound_provider_index_next + RENAME TO email_outbound_provider_index; + +CREATE INDEX idx_email_outbound_provider_index_user_id + ON email_outbound_provider_index(user_id); + +CREATE UNIQUE INDEX idx_email_outbound_provider_index_message_id + ON email_outbound_provider_index(message_id); diff --git a/packages/worker/src/admin/mailbox-maintenance.node.test.ts b/packages/worker/src/admin/mailbox-maintenance.node.test.ts index 712e62d218..98ea99f33c 100644 --- a/packages/worker/src/admin/mailbox-maintenance.node.test.ts +++ b/packages/worker/src/admin/mailbox-maintenance.node.test.ts @@ -195,8 +195,7 @@ function createMaintenanceDb() { inbox_id TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, - PRIMARY KEY (provider, provider_message_id), - FOREIGN KEY (message_id) REFERENCES email_messages(id) ON DELETE CASCADE + PRIMARY KEY (provider, provider_message_id) ); CREATE TABLE email_delivery_events ( id TEXT PRIMARY KEY, @@ -543,6 +542,7 @@ test('loadAdminMailboxMaintenanceStatus aggregates buckets without owner ids', a incomplete: 1, eligible: 1, outboundProviderIndex: { + foreignKeyDetached: true, linkedMessageCount: 1, indexCount: 1, missingFromIndexCount: 0, diff --git a/packages/worker/src/admin/mailbox-maintenance.ts b/packages/worker/src/admin/mailbox-maintenance.ts index e7c4f87440..d59bd2ad78 100644 --- a/packages/worker/src/admin/mailbox-maintenance.ts +++ b/packages/worker/src/admin/mailbox-maintenance.ts @@ -21,6 +21,7 @@ import { type MailboxCountResult, } from '#worker/email/mailbox-types.ts' import { + isOutboundProviderIndexForeignKeyDetached, loadOutboundProviderIndexParityReport, type OutboundProviderIndexParityReport, } from '#worker/email/outbound-provider-index.ts' @@ -81,7 +82,10 @@ export type AdminMailboxMaintenanceStatus = { * Fleet-wide aggregate D1 outbound provider reverse-index parity. Counts * only — no owner ids, message ids, or email content. */ - outboundProviderIndex: OutboundProviderIndexParityReport + outboundProviderIndex: OutboundProviderIndexParityReport & { + /** True only when the deployed schema has no legacy message_id FK. */ + foreignKeyDetached: boolean + } /** * Dedicated system-email authority versus its legacy rollback mirror. * Aggregate counts only; no email content is exposed. @@ -468,10 +472,19 @@ export async function loadAdminMailboxMaintenanceStatus(input: { } } - const [outboundProviderIndex, systemEmailGraph] = await Promise.all([ + const [ + outboundProviderIndexParity, + outboundProviderIndexForeignKeyDetached, + systemEmailGraph, + ] = await Promise.all([ loadOutboundProviderIndexParityReport({ db: input.db }), + isOutboundProviderIndexForeignKeyDetached(input.db), loadSystemEmailGraphParityReport({ db: input.db }), ]) + const outboundProviderIndex = { + ...outboundProviderIndexParity, + foreignKeyDetached: outboundProviderIndexForeignKeyDetached, + } return { generatedAt, diff --git a/packages/worker/src/app/retention.ts b/packages/worker/src/app/retention.ts index 048a1122ed..a61e66cd15 100644 --- a/packages/worker/src/app/retention.ts +++ b/packages/worker/src/app/retention.ts @@ -711,7 +711,13 @@ export async function pruneUserEmailMessagesForRetention(input: { idColumn: 'message_id', ids: messageIds, }) - // Provider-index rows cascade via FK ON DELETE CASCADE from email_messages. + // The retained global provider index has no cross-store message FK. + await deleteByIds({ + db: input.db, + table: 'email_outbound_provider_index', + idColumn: 'message_id', + ids: messageIds, + }) result.deletedMessages = await deleteByIds({ db: input.db, table: 'email_messages', diff --git a/packages/worker/src/email/email-user-graph-authority.node.test.ts b/packages/worker/src/email/email-user-graph-authority.node.test.ts new file mode 100644 index 0000000000..9da222dfc7 --- /dev/null +++ b/packages/worker/src/email/email-user-graph-authority.node.test.ts @@ -0,0 +1,120 @@ +import { expect, test } from 'vitest' +import { + assertUserEmailGraphD1StatementsAllowedAfterCutover, + classifyEmailGraphD1Statement, +} from '#worker/test-support/email-user-graph-statement-classifier.ts' +import { + assertEmailGraphAuthority, + assertSystemEmailGraphOwner, + assertUserEmailGraphOwner, + emailGraphAuthorityForOwner, + emailUserGraphAuthority, + systemEmailGraphAuthority, +} from './email-user-graph-authority.ts' +import { systemEmailOwnerId } from './email-owner.ts' + +test('email graph authority reserves Mailbox for USER and dedicated D1 for system', () => { + expect(emailGraphAuthorityForOwner('user-1')).toBe(emailUserGraphAuthority) + expect(emailGraphAuthorityForOwner(systemEmailOwnerId)).toBe( + systemEmailGraphAuthority, + ) + expect(() => + assertEmailGraphAuthority({ + ownerId: 'user-1', + authority: emailUserGraphAuthority, + }), + ).not.toThrow() + expect(() => + assertEmailGraphAuthority({ + ownerId: systemEmailOwnerId, + authority: systemEmailGraphAuthority, + }), + ).not.toThrow() + expect(() => + assertEmailGraphAuthority({ + ownerId: 'user-1', + authority: systemEmailGraphAuthority, + }), + ).toThrow(/requires mailbox authority/iu) + expect(() => + assertEmailGraphAuthority({ + ownerId: systemEmailOwnerId, + authority: emailUserGraphAuthority, + }), + ).toThrow(/requires dedicated-system-d1 authority/iu) + expect(() => assertUserEmailGraphOwner(systemEmailOwnerId)).toThrow( + /requires dedicated D1 graph authority/iu, + ) + expect(() => assertSystemEmailGraphOwner('user-1')).toThrow( + /require Mailbox authority/iu, + ) +}) + +test('test classifier rejects captured USER graph mutations but permits reads and provider lookup writes', () => { + expect( + classifyEmailGraphD1Statement(` + WITH candidate AS (SELECT id FROM email_messages WHERE user_id = ?) + DELETE FROM email_messages + WHERE id IN (SELECT id FROM candidate) + `), + ).toEqual({ + sharedGraphWrites: ['email_messages'], + dedicatedSystemGraphWrites: [], + }) + expect( + classifyEmailGraphD1Statement(` + INSERT INTO system_email_delivery_events (id) VALUES (?) + `), + ).toEqual({ + sharedGraphWrites: [], + dedicatedSystemGraphWrites: ['system_email_delivery_events'], + }) + expect( + classifyEmailGraphD1Statement(` + -- UPDATE email_threads is commentary, not a statement. + SELECT message.id + FROM email_messages message + WHERE message.user_id = ? + `), + ).toEqual({ + sharedGraphWrites: [], + dedicatedSystemGraphWrites: [], + }) + + expect(() => + assertUserEmailGraphD1StatementsAllowedAfterCutover({ + ownerId: 'user-1', + statements: [ + `SELECT * FROM email_messages WHERE user_id = ?`, + `INSERT INTO email_outbound_provider_index ( + provider, provider_message_id, user_id, message_id, + created_at, updated_at + ) VALUES (?, ?, ?, ?, ?, ?)`, + `UPDATE users SET updated_at = ? WHERE stable_user_id = ?`, + ], + }), + ).not.toThrow() + expect(() => + assertUserEmailGraphD1StatementsAllowedAfterCutover({ + ownerId: 'user-1', + statements: [ + `SELECT * FROM email_messages WHERE user_id = ?`, + `UPDATE email_threads SET updated_at = ? WHERE id = ?`, + ], + }), + ).toThrow(/D1 graph write to email_threads/iu) + expect(() => + assertUserEmailGraphD1StatementsAllowedAfterCutover({ + ownerId: 'user-1', + statements: [ + `REPLACE INTO "email_attachments" (id, message_id) VALUES (?, ?)`, + ], + }), + ).toThrow(/D1 graph write to email_attachments/iu) + expect(() => + assertUserEmailGraphD1StatementsAllowedAfterCutover({ + ownerId: systemEmailOwnerId, + statements: [], + }), + ).toThrow(/requires dedicated D1 graph authority/iu) +}) diff --git a/packages/worker/src/email/email-user-graph-authority.ts b/packages/worker/src/email/email-user-graph-authority.ts new file mode 100644 index 0000000000..4633e03c2b --- /dev/null +++ b/packages/worker/src/email/email-user-graph-authority.ts @@ -0,0 +1,51 @@ +import { isSystemEmailOwner } from './email-owner.ts' + +/** + * Step 5a prerequisite contract only. + * + * No production write path consumes these assertions in this slice, so this + * module does not yet protect live traffic. The step 5a write-cutover must route + * every USER graph mutation through this contract when it removes shared-D1 + * writes. + */ + +export const emailUserGraphAuthority = 'mailbox' as const +export const systemEmailGraphAuthority = 'dedicated-system-d1' as const + +export type EmailGraphAuthority = + | typeof emailUserGraphAuthority + | typeof systemEmailGraphAuthority + +export function emailGraphAuthorityForOwner( + ownerId: string, +): EmailGraphAuthority { + return isSystemEmailOwner(ownerId) + ? systemEmailGraphAuthority + : emailUserGraphAuthority +} + +export function assertEmailGraphAuthority(input: { + ownerId: string + authority: EmailGraphAuthority +}): void { + const expected = emailGraphAuthorityForOwner(input.ownerId) + if (input.authority !== expected) { + throw new Error( + `Email graph owner ${input.ownerId} requires ${expected} authority, not ${input.authority}.`, + ) + } +} + +export function assertUserEmailGraphOwner(ownerId: string): void { + if (isSystemEmailOwner(ownerId)) { + throw new Error( + 'The reserved system email owner requires dedicated D1 graph authority.', + ) + } +} + +export function assertSystemEmailGraphOwner(ownerId: string): void { + if (!isSystemEmailOwner(ownerId)) { + throw new Error('USER email graph owners require Mailbox authority.') + } +} diff --git a/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts b/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts new file mode 100644 index 0000000000..7272f187d0 --- /dev/null +++ b/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts @@ -0,0 +1,210 @@ +import { DatabaseSync } from 'node:sqlite' +import { expect, test } from 'vitest' +import { + applyMigrationLikeD1, + applyMigrationsBefore, +} from '#worker/test-support/system-email-graph-migration.ts' +import { createD1FromSqlite } from '#worker/test-support/create-d1-from-sqlite.ts' +import { systemEmailOwnerId } from './email-owner.ts' +import { + emailOutboundProviderCloudflare, + getOutboundProviderIndexRow, +} from './outbound-provider-index.ts' + +const providerIndexDetachMigration = + '0132-email-outbound-provider-index-detach.sql' + +function tableColumns(db: DatabaseSync) { + return db.prepare(`PRAGMA table_info(email_outbound_provider_index)`).all() +} + +function namedIndexes(db: DatabaseSync) { + return db + .prepare( + `SELECT name, "unique" + FROM pragma_index_list('email_outbound_provider_index') + WHERE origin != 'pk' + ORDER BY name`, + ) + .all() +} + +test('0132 atomically detaches provider lookup from the legacy USER graph', async () => { + using sqlite = new DatabaseSync(':memory:') + sqlite.exec('PRAGMA foreign_keys = ON') + applyMigrationsBefore(sqlite, providerIndexDetachMigration) + sqlite.exec(` + INSERT INTO email_messages ( + id, direction, user_id, from_address, processing_status, + provider_message_id, created_at, updated_at + ) VALUES + ( + 'user-message-1', 'outbound', 'user-1', 'sender@example.com', + 'sent', 'provider-user-1', '2026-08-03T00:00:00.000Z', + '2026-08-03T00:01:00.000Z' + ), + ( + 'user-message-2', 'outbound', 'user-2', 'sender@example.com', + 'sent', 'provider-user-2', '2026-08-03T00:02:00.000Z', + '2026-08-03T00:03:00.000Z' + ), + ( + 'system-message', 'inbound', '${systemEmailOwnerId}', + 'sender@example.net', 'stored', NULL, + '2026-08-03T00:04:00.000Z', '2026-08-03T00:05:00.000Z' + ); + + INSERT INTO email_outbound_provider_index ( + provider, provider_message_id, user_id, message_id, inbox_id, + created_at, updated_at + ) VALUES + ( + '${emailOutboundProviderCloudflare}', 'provider-user-1', 'user-1', + 'user-message-1', 'inbox-1', '2026-08-03T00:00:00.000Z', + '2026-08-03T00:01:00.000Z' + ), + ( + '${emailOutboundProviderCloudflare}', 'provider-user-2', 'user-2', + 'user-message-2', NULL, '2026-08-03T00:02:00.000Z', + '2026-08-03T00:03:00.000Z' + ), + ( + '${emailOutboundProviderCloudflare}', 'invalid-system-provider', + '${systemEmailOwnerId}', 'system-message', NULL, + '2026-08-03T00:04:00.000Z', '2026-08-03T00:05:00.000Z' + ); + `) + const columnsBefore = tableColumns(sqlite) + const indexesBefore = namedIndexes(sqlite) + + expect(() => + applyMigrationLikeD1(sqlite, providerIndexDetachMigration), + ).toThrow(/CHECK constraint failed/u) + expect(tableColumns(sqlite)).toEqual(columnsBefore) + expect(namedIndexes(sqlite)).toEqual(indexesBefore) + expect( + sqlite + .prepare(`PRAGMA foreign_key_list(email_outbound_provider_index)`) + .all(), + ).toContainEqual( + expect.objectContaining({ + table: 'email_messages', + from: 'message_id', + to: 'id', + on_delete: 'CASCADE', + }), + ) + expect( + sqlite + .prepare( + `SELECT COUNT(*) AS count + FROM email_outbound_provider_index`, + ) + .get(), + ).toEqual({ count: 3 }) + + sqlite + .prepare( + `DELETE FROM email_outbound_provider_index + WHERE user_id = ?`, + ) + .run(systemEmailOwnerId) + applyMigrationLikeD1(sqlite, providerIndexDetachMigration) + + expect(tableColumns(sqlite)).toEqual(columnsBefore) + expect(namedIndexes(sqlite)).toEqual(indexesBefore) + expect( + sqlite + .prepare(`PRAGMA foreign_key_list(email_outbound_provider_index)`) + .all(), + ).toEqual([]) + expect( + sqlite + .prepare( + `SELECT + provider, provider_message_id, user_id, message_id, inbox_id, + created_at, updated_at + FROM email_outbound_provider_index + ORDER BY provider_message_id`, + ) + .all(), + ).toEqual([ + { + provider: emailOutboundProviderCloudflare, + provider_message_id: 'provider-user-1', + user_id: 'user-1', + message_id: 'user-message-1', + inbox_id: 'inbox-1', + created_at: '2026-08-03T00:00:00.000Z', + updated_at: '2026-08-03T00:01:00.000Z', + }, + { + provider: emailOutboundProviderCloudflare, + provider_message_id: 'provider-user-2', + user_id: 'user-2', + message_id: 'user-message-2', + inbox_id: null, + created_at: '2026-08-03T00:02:00.000Z', + updated_at: '2026-08-03T00:03:00.000Z', + }, + ]) + expect( + sqlite + .prepare( + `SELECT COUNT(*) AS count + FROM email_outbound_provider_index + WHERE user_id = ?`, + ) + .get(systemEmailOwnerId), + ).toEqual({ count: 0 }) + expect(() => + sqlite + .prepare( + `INSERT INTO email_outbound_provider_index ( + provider, provider_message_id, user_id, message_id, + created_at, updated_at + ) VALUES (?, ?, ?, ?, ?, ?)`, + ) + .run( + emailOutboundProviderCloudflare, + 'provider-system-after-cutover', + systemEmailOwnerId, + 'system-message', + '2026-08-03T00:06:00.000Z', + '2026-08-03T00:06:00.000Z', + ), + ).toThrow(/CHECK constraint failed/u) + + sqlite + .prepare(`DELETE FROM email_messages WHERE id = ?`) + .run('user-message-1') + const d1 = createD1FromSqlite(sqlite) + await expect( + getOutboundProviderIndexRow({ + db: d1, + providerMessageId: 'provider-user-1', + }), + ).resolves.toMatchObject({ + userId: 'user-1', + messageId: 'user-message-1', + inboxId: 'inbox-1', + }) + + const cleanup = sqlite + .prepare( + `DELETE FROM email_outbound_provider_index + WHERE user_id = ?`, + ) + .run('user-1') + expect(cleanup.changes).toBe(1) + expect( + sqlite + .prepare( + `SELECT provider_message_id + FROM email_outbound_provider_index + ORDER BY provider_message_id`, + ) + .all(), + ).toEqual([{ provider_message_id: 'provider-user-2' }]) + expect(sqlite.prepare('PRAGMA foreign_key_check').all()).toEqual([]) +}) diff --git a/packages/worker/src/email/outbound-provider-index.ts b/packages/worker/src/email/outbound-provider-index.ts index df39301f2d..d7f60d770f 100644 --- a/packages/worker/src/email/outbound-provider-index.ts +++ b/packages/worker/src/email/outbound-provider-index.ts @@ -26,6 +26,20 @@ export type OutboundProviderIndexParityReport = { parity: boolean } +export async function isOutboundProviderIndexForeignKeyDetached( + db: D1Database, +): Promise { + const result = await db + .prepare(`PRAGMA foreign_key_list(email_outbound_provider_index)`) + .all<{ table: string; from: string; to: string }>() + return !result.results.some( + (foreignKey) => + foreignKey.table === 'email_messages' && + foreignKey.from === 'message_id' && + foreignKey.to === 'id', + ) +} + export function classifyOutboundProviderIndexParity( counts: Omit, ): OutboundProviderIndexParityReport { diff --git a/packages/worker/src/email/outbound-provider-index.workers.test.ts b/packages/worker/src/email/outbound-provider-index.workers.test.ts index 27d0ead910..e539250d50 100644 --- a/packages/worker/src/email/outbound-provider-index.workers.test.ts +++ b/packages/worker/src/email/outbound-provider-index.workers.test.ts @@ -7,6 +7,7 @@ import { deleteOutboundProviderIndexByMessageId, emailOutboundProviderCloudflare, getOutboundProviderIndexRow, + isOutboundProviderIndexForeignKeyDetached, loadOutboundProviderIndexParityReport, } from './outbound-provider-index.ts' import { OutboundEmailPersistenceError, sendOutboundEmail } from './outbound.ts' @@ -56,6 +57,9 @@ async function seedVerifiedAccount(email: string) { test('insert and delivery update dual-write the outbound provider index', async () => { await ensureEmailTestSchema(env.APP_DB) + await expect( + isOutboundProviderIndexForeignKeyDetached(env.APP_DB), + ).resolves.toBe(true) const userId = `idx-user-${crypto.randomUUID()}` const providerMessageId = `provider-${crypto.randomUUID()}` const message = await insertEmailMessage({ @@ -252,22 +256,24 @@ test('provider delivery lifecycle resolves through the index for user owners', a ).toBe('unmatched') }) -test('provider reverse lookup refuses legacy system owner links', async () => { +test('provider index blocks system owner links', async () => { await ensureEmailTestSchema(env.APP_DB) const providerMessageId = `provider-system-${crypto.randomUUID()}` - await insertEmailMessage({ - db: env.APP_DB, - message: { - direction: 'outbound', - userId: 'system:email', - fromAddress: 'support@example.com', - toAddresses: ['recipient@example.net'], - subject: 'Unsupported system link', - processingStatus: 'sent', - providerMessageId, - sentAt: '2026-08-01T13:00:00.000Z', - }, - }) + await expect( + insertEmailMessage({ + db: env.APP_DB, + message: { + direction: 'outbound', + userId: 'system:email', + fromAddress: 'support@example.com', + toAddresses: ['recipient@example.net'], + subject: 'Unsupported system link', + processingStatus: 'sent', + providerMessageId, + sentAt: '2026-08-01T13:00:00.000Z', + }, + }), + ).rejects.toThrow(/CHECK constraint failed/u) expect( await getOutboundEmailMessageByProviderMessageId({ @@ -281,8 +287,8 @@ test('provider reverse lookup refuses legacy system owner links', async () => { userId: 'system:email', }), ).toMatchObject({ - linkedMessageCount: 1, - indexCount: 1, + linkedMessageCount: 0, + indexCount: 0, parity: true, }) }) diff --git a/packages/worker/src/email/repo.ts b/packages/worker/src/email/repo.ts index cc7de93570..ae359c855e 100644 --- a/packages/worker/src/email/repo.ts +++ b/packages/worker/src/email/repo.ts @@ -1305,8 +1305,14 @@ export async function deleteEmailMessageById(input: { : d1CapturedBlobs // Atomic batch: a partial delete (attachments gone, message left) would // lose the storage_key values needed to delete the R2 blobs on retry. - // Provider-index rows cascade via FK ON DELETE CASCADE from email_messages. + // The global provider index intentionally has no cross-store message FK. await input.db.batch([ + input.db + .prepare( + `DELETE FROM email_outbound_provider_index + WHERE message_id = ?`, + ) + .bind(input.messageId), input.db .prepare(`DELETE FROM email_attachments WHERE message_id = ?`) .bind(input.messageId), @@ -1354,6 +1360,12 @@ export async function deleteEmailMessageProjectionById(input: { return { messageDeleted: false } } await input.db.batch([ + input.db + .prepare( + `DELETE FROM email_outbound_provider_index + WHERE message_id = ?`, + ) + .bind(input.messageId), input.db .prepare(`DELETE FROM email_attachments WHERE message_id = ?`) .bind(input.messageId), diff --git a/packages/worker/src/email/system-email-authority.workers.test.ts b/packages/worker/src/email/system-email-authority.workers.test.ts index 1162d2531b..f2d0f985b9 100644 --- a/packages/worker/src/email/system-email-authority.workers.test.ts +++ b/packages/worker/src/email/system-email-authority.workers.test.ts @@ -382,14 +382,16 @@ test('dedicated authority operations fail closed if provider links appear after ) .bind(systemEmailOwnerId, now, now) .run() - await env.APP_DB.prepare( - `INSERT INTO email_outbound_provider_index ( - provider, provider_message_id, user_id, message_id, created_at, - updated_at - ) VALUES ('resend', 'provider-message-after-cutover', ?, ?, ?, ?)`, - ) - .bind(systemEmailOwnerId, 'post-cutover-provider-link', now, now) - .run() + await expect( + env.APP_DB.prepare( + `INSERT INTO email_outbound_provider_index ( + provider, provider_message_id, user_id, message_id, created_at, + updated_at + ) VALUES ('resend', 'provider-message-after-cutover', ?, ?, ?, ?)`, + ) + .bind(systemEmailOwnerId, 'post-cutover-provider-link', now, now) + .run(), + ).rejects.toThrow(/CHECK constraint failed/u) await expect( listSystemEmailMessages({ db: env.APP_DB, limit: 10 }), diff --git a/packages/worker/src/email/test-schema.ts b/packages/worker/src/email/test-schema.ts index c231f1eaca..70b99eaf9f 100644 --- a/packages/worker/src/email/test-schema.ts +++ b/packages/worker/src/email/test-schema.ts @@ -362,15 +362,14 @@ ON email_messages(provider_message_id) WHERE direction = 'outbound' AND provider_message_id IS NOT NULL;`, `CREATE TABLE IF NOT EXISTS email_outbound_provider_index ( - provider TEXT NOT NULL, - provider_message_id TEXT NOT NULL, - user_id TEXT NOT NULL, - message_id TEXT NOT NULL, - inbox_id TEXT, - created_at TEXT NOT NULL, - updated_at TEXT NOT NULL, - PRIMARY KEY (provider, provider_message_id), - FOREIGN KEY (message_id) REFERENCES email_messages(id) ON DELETE CASCADE + provider TEXT NOT NULL CHECK (length(provider) > 0), + provider_message_id TEXT NOT NULL CHECK (length(provider_message_id) > 0), + user_id TEXT NOT NULL CHECK (length(user_id) > 0 AND user_id <> 'system:email'), + message_id TEXT NOT NULL CHECK (length(message_id) > 0), + inbox_id TEXT CHECK (inbox_id IS NULL OR length(inbox_id) > 0), + created_at TEXT NOT NULL CHECK (length(created_at) > 0), + updated_at TEXT NOT NULL CHECK (length(updated_at) > 0), + PRIMARY KEY (provider, provider_message_id) );`, `CREATE INDEX IF NOT EXISTS idx_email_outbound_provider_index_user_id ON email_outbound_provider_index(user_id);`, diff --git a/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.node.test.ts b/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.node.test.ts index 67038e2ed5..015280e8b2 100644 --- a/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.node.test.ts +++ b/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.node.test.ts @@ -63,6 +63,7 @@ const emptyStatus = { newestCheckedAt: '2026-08-01T11:00:00.000Z', earliestCutoverAt: '2026-07-31T00:00:00.000Z', outboundProviderIndex: { + foreignKeyDetached: true, linkedMessageCount: 0, indexCount: 0, missingFromIndexCount: 0, diff --git a/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.ts b/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.ts index 30009aa4a6..4afdab2995 100644 --- a/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.ts +++ b/packages/worker/src/mcp/capabilities/admin/admin-mailbox-maintenance.ts @@ -46,6 +46,11 @@ const mailboxCountSchema = z.object({ }) const outboundProviderIndexParitySchema = z.object({ + foreignKeyDetached: z + .boolean() + .describe( + 'True when the deployed provider-index schema has no message_id foreign key to legacy email_messages.', + ), linkedMessageCount: z .number() .int() diff --git a/packages/worker/src/test-support/email-user-graph-statement-classifier.ts b/packages/worker/src/test-support/email-user-graph-statement-classifier.ts new file mode 100644 index 0000000000..00f6fa1f2b --- /dev/null +++ b/packages/worker/src/test-support/email-user-graph-statement-classifier.ts @@ -0,0 +1,79 @@ +import { assertUserEmailGraphOwner } from '#worker/email/email-user-graph-authority.ts' + +const sharedEmailGraphTables = [ + 'email_threads', + 'email_messages', + 'email_attachments', + 'email_delivery_events', +] as const + +const dedicatedSystemEmailGraphTables = [ + 'system_email_threads', + 'system_email_messages', + 'system_email_attachments', + 'system_email_delivery_events', +] as const + +export type EmailGraphD1StatementClassification = { + sharedGraphWrites: Array<(typeof sharedEmailGraphTables)[number]> + dedicatedSystemGraphWrites: Array< + (typeof dedicatedSystemEmailGraphTables)[number] + > +} + +const mutationTargetPattern = + /\b(?:insert(?:\s+or\s+\w+)?\s+into|replace(?:\s+or\s+\w+)?\s+into|update(?:\s+or\s+\w+)?|delete\s+from)\s+["`[]?([a-z_][a-z0-9_]*)/giu + +function withoutSqlComments(sql: string): string { + return sql.replace(/\/\*[\s\S]*?\*\//gu, ' ').replace(/--[^\r\n]*/gu, ' ') +} + +/** + * Test-only classifier for SQL captured from exercised live code paths. + * + * It identifies mutation targets, including a mutation after a WITH clause. + * Reads and the global email_outbound_provider_index are intentionally not + * graph writes. This helper has no production call site and is not a runtime + * cutover guard. + */ +export function classifyEmailGraphD1Statement( + sql: string, +): EmailGraphD1StatementClassification { + const targets = new Set() + for (const match of withoutSqlComments(sql).matchAll(mutationTargetPattern)) { + const target = match[1]?.toLowerCase() + if (target != null) targets.add(target) + } + return { + sharedGraphWrites: sharedEmailGraphTables.filter((table) => + targets.has(table), + ), + dedicatedSystemGraphWrites: dedicatedSystemEmailGraphTables.filter( + (table) => targets.has(table), + ), + } +} + +/** + * Apply this to every D1 statement captured while exercising a USER live path + * after cutover. It fails on either legacy shared-graph writes or accidental + * writes into the dedicated system graph; non-graph D1 writes remain allowed. + */ +export function assertUserEmailGraphD1StatementsAllowedAfterCutover(input: { + ownerId: string + statements: ReadonlyArray +}): void { + assertUserEmailGraphOwner(input.ownerId) + for (const sql of input.statements) { + const classified = classifyEmailGraphD1Statement(sql) + const forbidden = [ + ...classified.sharedGraphWrites, + ...classified.dedicatedSystemGraphWrites, + ] + if (forbidden.length > 0) { + throw new Error( + `Post-cutover USER email path attempted a D1 graph write to ${forbidden.join(', ')}.`, + ) + } + } +} diff --git a/tools/migration-ledger.json b/tools/migration-ledger.json index 007d135c77..ef32a92b4f 100644 --- a/tools/migration-ledger.json +++ b/tools/migration-ledger.json @@ -539,6 +539,10 @@ { "filename": "0131-system-email-graph-authority.sql", "sha256": "e2b8cecc90501c54c6a7391da9e2d0ab334c62c6ff749bbb128167a7db2480e8" + }, + { + "filename": "0132-email-outbound-provider-index-detach.sql", + "sha256": "c84f2ca81d5155b35cf8f85f99bc2ea2ef52e138318c458ad503c5ad4bb1296d" } ] } From 8178a876c3bad85872a5e2e4913448d6e74ac46f Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 3 Aug 2026 05:09:46 +0000 Subject: [PATCH 2/3] fix(email): preserve provider cleanup on rollback Co-authored-by: Kent C. Dodds --- .../contributing/architecture/data-storage.md | 55 ++++---- ...2-email-outbound-provider-index-detach.sql | 27 ++-- packages/worker/src/app/retention.ts | 9 +- .../email-user-graph-authority.node.test.ts | 120 ------------------ .../src/email/email-user-graph-authority.ts | 51 -------- ...ovider-index-detach-migration.node.test.ts | 36 +++++- packages/worker/src/email/test-schema.ts | 20 ++- .../email-user-graph-statement-classifier.ts | 79 ------------ tools/migration-ledger.json | 2 +- 9 files changed, 95 insertions(+), 304 deletions(-) delete mode 100644 packages/worker/src/email/email-user-graph-authority.node.test.ts delete mode 100644 packages/worker/src/email/email-user-graph-authority.ts delete mode 100644 packages/worker/src/test-support/email-user-graph-statement-classifier.ts diff --git a/docs/contributing/architecture/data-storage.md b/docs/contributing/architecture/data-storage.md index 1346a67884..ad39281474 100644 --- a/docs/contributing/architecture/data-storage.md +++ b/docs/contributing/architecture/data-storage.md @@ -1359,8 +1359,9 @@ cover the high-risk live surface today. 5. **USER graph contract (5a then 5b)** — deploy the 5a prerequisite first: migration `0132-email-outbound-provider-index-detach.sql` atomically rebuilds the retained global provider index without its cross-store-invalid - `email_messages` foreign key, and the authority guard contract lands without - changing live paths. Verify production reports + `email_messages` foreign key and installs a compatibility trigger that + preserves the old cascade behavior for legacy message deletes. No live USER + authority path flips in this prerequisite. Verify production reports `status.outboundProviderIndex.foreignKeyDetached: true` before proceeding. A later 5a cutover then moves **all** USER inbound, outbound, provider-index synchronization, classification, explicit delete, and retention graph @@ -1371,12 +1372,15 @@ cover the high-risk live surface today. graph tables and migration-only machinery. No USER row or graph table is destructively removed by the prerequisite deployment. - **Rollback caveat:** after Mailbox-only writes begin, the frozen D1 USER - graph is stale and must never be re-enabled as authority by a code rollback. - A rollback must remain Mailbox-authoritative (or explicitly rebuild D1 from - Mailbox before restoring old code). After 5b drops the tables, rollback also - requires the verified fresh backup/schema restore; reverting application code - alone is not safe. + **Rollback caveat:** before the Mailbox-only write flip, old code that issues + a single `DELETE FROM email_messages` remains safe because the schema trigger + atomically removes the matching provider-index row. After Mailbox-only writes + begin, the frozen D1 USER graph is stale and must never be re-enabled as + authority by a code rollback. A rollback must remain Mailbox-authoritative + (or explicitly rebuild D1 from Mailbox before restoring old code). Step 5b + removes the compatibility trigger with `email_messages`; after that drop, + rollback also requires the verified fresh backup/schema restore, and + reverting application code alone is not safe. The every-5-minute `mailbox_parity` scheduled lane (queue-isolated sibling in `scheduled-lanes.ts`) owns backfill of all owner messages and delivery events, @@ -1409,22 +1413,25 @@ below. plus index sync in one `db.batch` (index owner/inbox fields come from the authoritative message row, never caller input). `message_id` is an opaque owner-scoped key with no foreign key to `email_messages`: legacy message - deletion cannot cascade across the upcoming Mailbox/D1 store boundary. - Explicit message and account deletion therefore remove index rows by - `message_id`/`user_id`. Outbound send separates provider acceptance from - terminal D1/index persistence: once the provider returns a - `providerMessageId`, persistence uses bounded D1 retries and must not mark the - message `failed`, clear the id, or resend. Account export treats the table as - derived global lookup (`includeInExport: false` / - `derivedData.email_outbound_provider_index`) because authoritative outbound - message rows are already exported. `recordProviderEmailDeliveryEvent` resolves - index-first, then loads the owner-scoped message by `user_id`/`message_id` (no - full-table provider scan). System outbound is unsupported, and the verified - `no-system-provider-links` disposition means `system:email` rows are never - added to this global index. Aggregate parity - (`loadOutboundProviderIndexParityReport`; counts only) is surfaced on - `admin_mailbox_maintenance` `status.outboundProviderIndex` for production - verification, alongside the live schema check `foreignKeyDetached`. + deletion cannot cascade through a cross-store relationship. While the legacy + table exists, trigger `email_messages_delete_outbound_provider_index` + atomically deletes index rows by `OLD.id`, preserving the old FK behavior for + rollback code; step 5b removes the trigger with `email_messages`. Current + explicit message deletion also cleans the index inside the same `db.batch` for + cross-version compatibility, and account deletion inventories/deletes by + `user_id`. Outbound send separates provider acceptance from terminal D1/index + persistence: once the provider returns a `providerMessageId`, persistence uses + bounded D1 retries and must not mark the message `failed`, clear the id, or + resend. Account export treats the table as derived global lookup + (`includeInExport: false` / `derivedData.email_outbound_provider_index`) + because authoritative outbound message rows are already exported. + `recordProviderEmailDeliveryEvent` resolves index-first, then loads the + owner-scoped message by `user_id`/`message_id` (no full-table provider scan). + System outbound is unsupported, and the verified `no-system-provider-links` + disposition means `system:email` rows are never added to this global index. + Aggregate parity (`loadOutboundProviderIndexParityReport`; counts only) is + surfaced on `admin_mailbox_maintenance` `status.outboundProviderIndex` for + production verification, alongside the live schema check `foreignKeyDetached`. Contextless provider-id reverse lookups must not enumerate per-user Mailbox objects or resolve a system compatibility mirror. diff --git a/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql b/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql index e201d6d2b2..5a180fadeb 100644 --- a/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql +++ b/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql @@ -8,15 +8,13 @@ -- cross-store foreign key. CREATE TABLE email_outbound_provider_index_next ( - provider TEXT NOT NULL CHECK (length(provider) > 0), - provider_message_id TEXT NOT NULL CHECK (length(provider_message_id) > 0), - user_id TEXT NOT NULL CHECK ( - length(user_id) > 0 AND user_id <> 'system:email' - ), - message_id TEXT NOT NULL CHECK (length(message_id) > 0), - inbox_id TEXT CHECK (inbox_id IS NULL OR length(inbox_id) > 0), - created_at TEXT NOT NULL CHECK (length(created_at) > 0), - updated_at TEXT NOT NULL CHECK (length(updated_at) > 0), + provider TEXT NOT NULL, + provider_message_id TEXT NOT NULL, + user_id TEXT NOT NULL CHECK (user_id <> 'system:email'), + message_id TEXT NOT NULL, + inbox_id TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, PRIMARY KEY (provider, provider_message_id) ); @@ -49,3 +47,14 @@ CREATE INDEX idx_email_outbound_provider_index_user_id CREATE UNIQUE INDEX idx_email_outbound_provider_index_message_id ON email_outbound_provider_index(message_id); + +-- Preserve the old FK's rollback behavior while legacy email_messages remains. +-- SQLite runs trigger statements in the same transaction as the DELETE, so a +-- failure rolls back both the message and index cleanup. Step 5b removes this +-- compatibility trigger with the legacy table. +CREATE TRIGGER email_messages_delete_outbound_provider_index +AFTER DELETE ON email_messages +BEGIN + DELETE FROM email_outbound_provider_index + WHERE message_id = OLD.id; +END; diff --git a/packages/worker/src/app/retention.ts b/packages/worker/src/app/retention.ts index a61e66cd15..737eae1543 100644 --- a/packages/worker/src/app/retention.ts +++ b/packages/worker/src/app/retention.ts @@ -711,13 +711,8 @@ export async function pruneUserEmailMessagesForRetention(input: { idColumn: 'message_id', ids: messageIds, }) - // The retained global provider index has no cross-store message FK. - await deleteByIds({ - db: input.db, - table: 'email_outbound_provider_index', - idColumn: 'message_id', - ids: messageIds, - }) + // The schema compatibility trigger atomically deletes provider-index rows + // with each legacy message DELETE. result.deletedMessages = await deleteByIds({ db: input.db, table: 'email_messages', diff --git a/packages/worker/src/email/email-user-graph-authority.node.test.ts b/packages/worker/src/email/email-user-graph-authority.node.test.ts deleted file mode 100644 index 9da222dfc7..0000000000 --- a/packages/worker/src/email/email-user-graph-authority.node.test.ts +++ /dev/null @@ -1,120 +0,0 @@ -import { expect, test } from 'vitest' -import { - assertUserEmailGraphD1StatementsAllowedAfterCutover, - classifyEmailGraphD1Statement, -} from '#worker/test-support/email-user-graph-statement-classifier.ts' -import { - assertEmailGraphAuthority, - assertSystemEmailGraphOwner, - assertUserEmailGraphOwner, - emailGraphAuthorityForOwner, - emailUserGraphAuthority, - systemEmailGraphAuthority, -} from './email-user-graph-authority.ts' -import { systemEmailOwnerId } from './email-owner.ts' - -test('email graph authority reserves Mailbox for USER and dedicated D1 for system', () => { - expect(emailGraphAuthorityForOwner('user-1')).toBe(emailUserGraphAuthority) - expect(emailGraphAuthorityForOwner(systemEmailOwnerId)).toBe( - systemEmailGraphAuthority, - ) - expect(() => - assertEmailGraphAuthority({ - ownerId: 'user-1', - authority: emailUserGraphAuthority, - }), - ).not.toThrow() - expect(() => - assertEmailGraphAuthority({ - ownerId: systemEmailOwnerId, - authority: systemEmailGraphAuthority, - }), - ).not.toThrow() - expect(() => - assertEmailGraphAuthority({ - ownerId: 'user-1', - authority: systemEmailGraphAuthority, - }), - ).toThrow(/requires mailbox authority/iu) - expect(() => - assertEmailGraphAuthority({ - ownerId: systemEmailOwnerId, - authority: emailUserGraphAuthority, - }), - ).toThrow(/requires dedicated-system-d1 authority/iu) - expect(() => assertUserEmailGraphOwner(systemEmailOwnerId)).toThrow( - /requires dedicated D1 graph authority/iu, - ) - expect(() => assertSystemEmailGraphOwner('user-1')).toThrow( - /require Mailbox authority/iu, - ) -}) - -test('test classifier rejects captured USER graph mutations but permits reads and provider lookup writes', () => { - expect( - classifyEmailGraphD1Statement(` - WITH candidate AS (SELECT id FROM email_messages WHERE user_id = ?) - DELETE FROM email_messages - WHERE id IN (SELECT id FROM candidate) - `), - ).toEqual({ - sharedGraphWrites: ['email_messages'], - dedicatedSystemGraphWrites: [], - }) - expect( - classifyEmailGraphD1Statement(` - INSERT INTO system_email_delivery_events (id) VALUES (?) - `), - ).toEqual({ - sharedGraphWrites: [], - dedicatedSystemGraphWrites: ['system_email_delivery_events'], - }) - expect( - classifyEmailGraphD1Statement(` - -- UPDATE email_threads is commentary, not a statement. - SELECT message.id - FROM email_messages message - WHERE message.user_id = ? - `), - ).toEqual({ - sharedGraphWrites: [], - dedicatedSystemGraphWrites: [], - }) - - expect(() => - assertUserEmailGraphD1StatementsAllowedAfterCutover({ - ownerId: 'user-1', - statements: [ - `SELECT * FROM email_messages WHERE user_id = ?`, - `INSERT INTO email_outbound_provider_index ( - provider, provider_message_id, user_id, message_id, - created_at, updated_at - ) VALUES (?, ?, ?, ?, ?, ?)`, - `UPDATE users SET updated_at = ? WHERE stable_user_id = ?`, - ], - }), - ).not.toThrow() - expect(() => - assertUserEmailGraphD1StatementsAllowedAfterCutover({ - ownerId: 'user-1', - statements: [ - `SELECT * FROM email_messages WHERE user_id = ?`, - `UPDATE email_threads SET updated_at = ? WHERE id = ?`, - ], - }), - ).toThrow(/D1 graph write to email_threads/iu) - expect(() => - assertUserEmailGraphD1StatementsAllowedAfterCutover({ - ownerId: 'user-1', - statements: [ - `REPLACE INTO "email_attachments" (id, message_id) VALUES (?, ?)`, - ], - }), - ).toThrow(/D1 graph write to email_attachments/iu) - expect(() => - assertUserEmailGraphD1StatementsAllowedAfterCutover({ - ownerId: systemEmailOwnerId, - statements: [], - }), - ).toThrow(/requires dedicated D1 graph authority/iu) -}) diff --git a/packages/worker/src/email/email-user-graph-authority.ts b/packages/worker/src/email/email-user-graph-authority.ts deleted file mode 100644 index 4633e03c2b..0000000000 --- a/packages/worker/src/email/email-user-graph-authority.ts +++ /dev/null @@ -1,51 +0,0 @@ -import { isSystemEmailOwner } from './email-owner.ts' - -/** - * Step 5a prerequisite contract only. - * - * No production write path consumes these assertions in this slice, so this - * module does not yet protect live traffic. The step 5a write-cutover must route - * every USER graph mutation through this contract when it removes shared-D1 - * writes. - */ - -export const emailUserGraphAuthority = 'mailbox' as const -export const systemEmailGraphAuthority = 'dedicated-system-d1' as const - -export type EmailGraphAuthority = - | typeof emailUserGraphAuthority - | typeof systemEmailGraphAuthority - -export function emailGraphAuthorityForOwner( - ownerId: string, -): EmailGraphAuthority { - return isSystemEmailOwner(ownerId) - ? systemEmailGraphAuthority - : emailUserGraphAuthority -} - -export function assertEmailGraphAuthority(input: { - ownerId: string - authority: EmailGraphAuthority -}): void { - const expected = emailGraphAuthorityForOwner(input.ownerId) - if (input.authority !== expected) { - throw new Error( - `Email graph owner ${input.ownerId} requires ${expected} authority, not ${input.authority}.`, - ) - } -} - -export function assertUserEmailGraphOwner(ownerId: string): void { - if (isSystemEmailOwner(ownerId)) { - throw new Error( - 'The reserved system email owner requires dedicated D1 graph authority.', - ) - } -} - -export function assertSystemEmailGraphOwner(ownerId: string): void { - if (!isSystemEmailOwner(ownerId)) { - throw new Error('USER email graph owners require Mailbox authority.') - } -} diff --git a/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts b/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts index 7272f187d0..7b1f1115d2 100644 --- a/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts +++ b/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts @@ -29,7 +29,18 @@ function namedIndexes(db: DatabaseSync) { .all() } -test('0132 atomically detaches provider lookup from the legacy USER graph', async () => { +function compatibilityTrigger(db: DatabaseSync) { + return db + .prepare( + `SELECT sql + FROM sqlite_schema + WHERE type = 'trigger' + AND name = 'email_messages_delete_outbound_provider_index'`, + ) + .get() +} + +test('0132 detaches the FK while preserving rollback message-delete behavior', async () => { using sqlite = new DatabaseSync(':memory:') sqlite.exec('PRAGMA foreign_keys = ON') applyMigrationsBefore(sqlite, providerIndexDetachMigration) @@ -102,6 +113,7 @@ test('0132 atomically detaches provider lookup from the legacy USER graph', asyn ) .get(), ).toEqual({ count: 3 }) + expect(compatibilityTrigger(sqlite)).toBeUndefined() sqlite .prepare( @@ -118,6 +130,9 @@ test('0132 atomically detaches provider lookup from the legacy USER graph', asyn .prepare(`PRAGMA foreign_key_list(email_outbound_provider_index)`) .all(), ).toEqual([]) + expect(compatibilityTrigger(sqlite)).toEqual({ + sql: expect.stringContaining('DELETE FROM email_outbound_provider_index'), + }) expect( sqlite .prepare( @@ -175,9 +190,6 @@ test('0132 atomically detaches provider lookup from the legacy USER graph', asyn ), ).toThrow(/CHECK constraint failed/u) - sqlite - .prepare(`DELETE FROM email_messages WHERE id = ?`) - .run('user-message-1') const d1 = createD1FromSqlite(sqlite) await expect( getOutboundProviderIndexRow({ @@ -190,12 +202,24 @@ test('0132 atomically detaches provider lookup from the legacy USER graph', asyn inboxId: 'inbox-1', }) + // This is the old rollback path: one legacy message DELETE, with no app + // provider-index statement. The schema trigger preserves FK-cascade behavior. + sqlite + .prepare(`DELETE FROM email_messages WHERE id = ?`) + .run('user-message-1') + await expect( + getOutboundProviderIndexRow({ + db: d1, + providerMessageId: 'provider-user-1', + }), + ).resolves.toBeNull() + const cleanup = sqlite .prepare( `DELETE FROM email_outbound_provider_index WHERE user_id = ?`, ) - .run('user-1') + .run('user-2') expect(cleanup.changes).toBe(1) expect( sqlite @@ -205,6 +229,6 @@ test('0132 atomically detaches provider lookup from the legacy USER graph', asyn ORDER BY provider_message_id`, ) .all(), - ).toEqual([{ provider_message_id: 'provider-user-2' }]) + ).toEqual([]) expect(sqlite.prepare('PRAGMA foreign_key_check').all()).toEqual([]) }) diff --git a/packages/worker/src/email/test-schema.ts b/packages/worker/src/email/test-schema.ts index 70b99eaf9f..1e8f53d8a0 100644 --- a/packages/worker/src/email/test-schema.ts +++ b/packages/worker/src/email/test-schema.ts @@ -362,19 +362,25 @@ ON email_messages(provider_message_id) WHERE direction = 'outbound' AND provider_message_id IS NOT NULL;`, `CREATE TABLE IF NOT EXISTS email_outbound_provider_index ( - provider TEXT NOT NULL CHECK (length(provider) > 0), - provider_message_id TEXT NOT NULL CHECK (length(provider_message_id) > 0), - user_id TEXT NOT NULL CHECK (length(user_id) > 0 AND user_id <> 'system:email'), - message_id TEXT NOT NULL CHECK (length(message_id) > 0), - inbox_id TEXT CHECK (inbox_id IS NULL OR length(inbox_id) > 0), - created_at TEXT NOT NULL CHECK (length(created_at) > 0), - updated_at TEXT NOT NULL CHECK (length(updated_at) > 0), + provider TEXT NOT NULL, + provider_message_id TEXT NOT NULL, + user_id TEXT NOT NULL CHECK (user_id <> 'system:email'), + message_id TEXT NOT NULL, + inbox_id TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, PRIMARY KEY (provider, provider_message_id) );`, `CREATE INDEX IF NOT EXISTS idx_email_outbound_provider_index_user_id ON email_outbound_provider_index(user_id);`, `CREATE UNIQUE INDEX IF NOT EXISTS idx_email_outbound_provider_index_message_id ON email_outbound_provider_index(message_id);`, + `CREATE TRIGGER IF NOT EXISTS email_messages_delete_outbound_provider_index +AFTER DELETE ON email_messages +BEGIN + DELETE FROM email_outbound_provider_index + WHERE message_id = OLD.id; +END;`, `CREATE UNIQUE INDEX IF NOT EXISTS idx_email_delivery_events_provider_event_id ON email_delivery_events(provider_event_id) WHERE provider_event_id IS NOT NULL;`, diff --git a/packages/worker/src/test-support/email-user-graph-statement-classifier.ts b/packages/worker/src/test-support/email-user-graph-statement-classifier.ts deleted file mode 100644 index 00f6fa1f2b..0000000000 --- a/packages/worker/src/test-support/email-user-graph-statement-classifier.ts +++ /dev/null @@ -1,79 +0,0 @@ -import { assertUserEmailGraphOwner } from '#worker/email/email-user-graph-authority.ts' - -const sharedEmailGraphTables = [ - 'email_threads', - 'email_messages', - 'email_attachments', - 'email_delivery_events', -] as const - -const dedicatedSystemEmailGraphTables = [ - 'system_email_threads', - 'system_email_messages', - 'system_email_attachments', - 'system_email_delivery_events', -] as const - -export type EmailGraphD1StatementClassification = { - sharedGraphWrites: Array<(typeof sharedEmailGraphTables)[number]> - dedicatedSystemGraphWrites: Array< - (typeof dedicatedSystemEmailGraphTables)[number] - > -} - -const mutationTargetPattern = - /\b(?:insert(?:\s+or\s+\w+)?\s+into|replace(?:\s+or\s+\w+)?\s+into|update(?:\s+or\s+\w+)?|delete\s+from)\s+["`[]?([a-z_][a-z0-9_]*)/giu - -function withoutSqlComments(sql: string): string { - return sql.replace(/\/\*[\s\S]*?\*\//gu, ' ').replace(/--[^\r\n]*/gu, ' ') -} - -/** - * Test-only classifier for SQL captured from exercised live code paths. - * - * It identifies mutation targets, including a mutation after a WITH clause. - * Reads and the global email_outbound_provider_index are intentionally not - * graph writes. This helper has no production call site and is not a runtime - * cutover guard. - */ -export function classifyEmailGraphD1Statement( - sql: string, -): EmailGraphD1StatementClassification { - const targets = new Set() - for (const match of withoutSqlComments(sql).matchAll(mutationTargetPattern)) { - const target = match[1]?.toLowerCase() - if (target != null) targets.add(target) - } - return { - sharedGraphWrites: sharedEmailGraphTables.filter((table) => - targets.has(table), - ), - dedicatedSystemGraphWrites: dedicatedSystemEmailGraphTables.filter( - (table) => targets.has(table), - ), - } -} - -/** - * Apply this to every D1 statement captured while exercising a USER live path - * after cutover. It fails on either legacy shared-graph writes or accidental - * writes into the dedicated system graph; non-graph D1 writes remain allowed. - */ -export function assertUserEmailGraphD1StatementsAllowedAfterCutover(input: { - ownerId: string - statements: ReadonlyArray -}): void { - assertUserEmailGraphOwner(input.ownerId) - for (const sql of input.statements) { - const classified = classifyEmailGraphD1Statement(sql) - const forbidden = [ - ...classified.sharedGraphWrites, - ...classified.dedicatedSystemGraphWrites, - ] - if (forbidden.length > 0) { - throw new Error( - `Post-cutover USER email path attempted a D1 graph write to ${forbidden.join(', ')}.`, - ) - } - } -} diff --git a/tools/migration-ledger.json b/tools/migration-ledger.json index ef32a92b4f..c4383fc5f3 100644 --- a/tools/migration-ledger.json +++ b/tools/migration-ledger.json @@ -542,7 +542,7 @@ }, { "filename": "0132-email-outbound-provider-index-detach.sql", - "sha256": "c84f2ca81d5155b35cf8f85f99bc2ea2ef52e138318c458ad503c5ad4bb1296d" + "sha256": "d3bb5b668057e29c3ad42cc1962c312ac4a7764e4f183268be3d5ddc29901dff" } ] } From dd8e7e59004a77d0cc56b972648e763f13bb468d Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 3 Aug 2026 05:09:52 +0000 Subject: [PATCH 3/3] test(email): cover rollback provider cleanup Co-authored-by: Kent C. Dodds --- .../worker/src/app/retention.node.test.ts | 25 ++++++++----------- packages/worker/src/email/repo.ts | 6 ++++- 2 files changed, 16 insertions(+), 15 deletions(-) diff --git a/packages/worker/src/app/retention.node.test.ts b/packages/worker/src/app/retention.node.test.ts index ed911e7561..b6116596ae 100644 --- a/packages/worker/src/app/retention.node.test.ts +++ b/packages/worker/src/app/retention.node.test.ts @@ -193,14 +193,19 @@ function createRetentionDb() { CREATE TABLE email_outbound_provider_index ( provider TEXT NOT NULL, provider_message_id TEXT NOT NULL, - user_id TEXT NOT NULL, + user_id TEXT NOT NULL CHECK (user_id <> 'system:email'), message_id TEXT NOT NULL, inbox_id TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, - PRIMARY KEY (provider, provider_message_id), - FOREIGN KEY (message_id) REFERENCES email_messages(id) ON DELETE CASCADE + PRIMARY KEY (provider, provider_message_id) ); + CREATE TRIGGER email_messages_delete_outbound_provider_index + AFTER DELETE ON email_messages + BEGIN + DELETE FROM email_outbound_provider_index + WHERE message_id = OLD.id; + END; CREATE TABLE email_attachments ( id TEXT PRIMARY KEY, message_id TEXT NOT NULL, @@ -615,17 +620,9 @@ test('email message retention deletes old user rows, R2 blobs, attachments, and created_at, updated_at ) VALUES ('cloudflare-email', 'prov-old-blob', 'user-1', 'msg-old-blob', NULL, ?, ?), - ('cloudflare-email', 'prov-fresh', 'user-1', 'msg-fresh', NULL, ?, ?), - ('cloudflare-email', 'prov-system', 'system:email', 'msg-old-system', NULL, ?, ?)`, - ) - .run( - indexCreatedAt, - indexCreatedAt, - daysAgo(1), - daysAgo(1), - indexCreatedAt, - indexCreatedAt, + ('cloudflare-email', 'prov-fresh', 'user-1', 'msg-fresh', NULL, ?, ?)`, ) + .run(indexCreatedAt, indexCreatedAt, daysAgo(1), daysAgo(1)) const blobDelete = vi.fn(async () => undefined) const result = await pruneUserEmailMessagesForRetention({ @@ -666,7 +663,7 @@ test('email message retention deletes old user rows, R2 blobs, attachments, and ) .all() as Array<{ provider_message_id: string }> ).map((row) => row.provider_message_id), - ).toEqual(['prov-fresh', 'prov-system']) + ).toEqual(['prov-fresh']) }) test('email message retention never deletes rows whose blob cannot be deleted first', async () => { diff --git a/packages/worker/src/email/repo.ts b/packages/worker/src/email/repo.ts index ae359c855e..fb68537b89 100644 --- a/packages/worker/src/email/repo.ts +++ b/packages/worker/src/email/repo.ts @@ -1305,7 +1305,9 @@ export async function deleteEmailMessageById(input: { : d1CapturedBlobs // Atomic batch: a partial delete (attachments gone, message left) would // lose the storage_key values needed to delete the R2 blobs on retry. - // The global provider index intentionally has no cross-store message FK. + // Keep explicit provider-index cleanup in this same atomic batch for + // mixed-version schemas. The old FK and new compatibility trigger both + // tolerate the row already being absent. await input.db.batch([ input.db .prepare( @@ -1359,6 +1361,8 @@ export async function deleteEmailMessageProjectionById(input: { if (row.user_id !== input.expectedUserId) { return { messageDeleted: false } } + // Same-batch explicit cleanup keeps this projection delete safe across the + // legacy-FK and compatibility-trigger schemas. await input.db.batch([ input.db .prepare(