diff --git a/docs/contributing/architecture/data-storage.md b/docs/contributing/architecture/data-storage.md index 3232aaca64..ad39281474 100644 --- a/docs/contributing/architecture/data-storage.md +++ b/docs/contributing/architecture/data-storage.md @@ -1353,11 +1353,34 @@ 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 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 + 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:** 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, @@ -1381,16 +1404,22 @@ 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 + 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 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 @@ -1399,11 +1428,12 @@ below. `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. + 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. Contextless provider-id reverse lookups must not - enumerate per-user Mailbox objects or resolve a system compatibility mirror. + 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..5a180fadeb --- /dev/null +++ b/packages/worker/migrations/0132-email-outbound-provider-index-detach.sql @@ -0,0 +1,60 @@ +-- 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, + 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) +); + +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); + +-- 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/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.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/app/retention.ts b/packages/worker/src/app/retention.ts index 048a1122ed..737eae1543 100644 --- a/packages/worker/src/app/retention.ts +++ b/packages/worker/src/app/retention.ts @@ -711,7 +711,8 @@ export async function pruneUserEmailMessagesForRetention(input: { idColumn: 'message_id', ids: messageIds, }) - // Provider-index rows cascade via FK ON DELETE CASCADE from email_messages. + // 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/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..7b1f1115d2 --- /dev/null +++ b/packages/worker/src/email/outbound-provider-index-detach-migration.node.test.ts @@ -0,0 +1,234 @@ +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() +} + +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) + 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 }) + expect(compatibilityTrigger(sqlite)).toBeUndefined() + + 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(compatibilityTrigger(sqlite)).toEqual({ + sql: expect.stringContaining('DELETE FROM email_outbound_provider_index'), + }) + 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) + + 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', + }) + + // 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-2') + expect(cleanup.changes).toBe(1) + expect( + sqlite + .prepare( + `SELECT provider_message_id + FROM email_outbound_provider_index + ORDER BY provider_message_id`, + ) + .all(), + ).toEqual([]) + 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..fb68537b89 100644 --- a/packages/worker/src/email/repo.ts +++ b/packages/worker/src/email/repo.ts @@ -1305,8 +1305,16 @@ 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. + // 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( + `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), @@ -1353,7 +1361,15 @@ 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( + `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..1e8f53d8a0 100644 --- a/packages/worker/src/email/test-schema.ts +++ b/packages/worker/src/email/test-schema.ts @@ -364,18 +364,23 @@ WHERE direction = 'outbound' `CREATE TABLE IF NOT EXISTS 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 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/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/tools/migration-ledger.json b/tools/migration-ledger.json index 007d135c77..c4383fc5f3 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": "d3bb5b668057e29c3ad42cc1962c312ac4a7764e4f183268be3d5ddc29901dff" } ] }