Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 43 additions & 13 deletions docs/contributing/architecture/data-storage.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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)

Expand Down
Original file line number Diff line number Diff line change
@@ -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;
4 changes: 2 additions & 2 deletions packages/worker/src/admin/mailbox-maintenance.node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
17 changes: 15 additions & 2 deletions packages/worker/src/admin/mailbox-maintenance.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand Down
25 changes: 11 additions & 14 deletions packages/worker/src/app/retention.node.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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 () => {
Expand Down
3 changes: 2 additions & 1 deletion packages/worker/src/app/retention.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down
Loading
Loading