File
Blob: src/worker/email-handler.ts
| 1 | import { nanoid } from "nanoid"; |
| 2 | import PostalMime from "postal-mime"; |
| 3 | import { and, count, eq } from "drizzle-orm"; |
| 4 | import { createDb } from "@/worker/db"; |
| 5 | import { attachments, domains, emails, inboxes } from "@/worker/db/schema"; |
| 6 | import { createLogger, errorContext } from "@/worker/logger"; |
| 7 | import { toNewEmailNotification } from "@/worker/serializers/email"; |
| 8 | import { |
| 9 | getAttachmentStorageKey, |
| 10 | getBodyStorageKey, |
| 11 | storeAttachment, |
| 12 | storeEmailBody, |
| 13 | storeRawEmail, |
| 14 | deleteStorageForEmails, |
| 15 | } from "@/worker/services/storage"; |
| 16 | |
| 17 | const logger = createLogger("email-handler"); |
| 18 | const MAX_INBOUND_EMAIL_SIZE_BYTES = 10 * 1024 * 1024; |
| 19 | const MAX_INBOUND_ATTACHMENTS = 10; |
| 20 | const MAX_EMAILS_PER_INBOX = 100; |
| 21 | |
| 22 | function parseRecipientAddress(address: string) { |
| 23 | const raw = address.trim(); |
| 24 | const normalized = raw.toLowerCase(); |
| 25 | const atIndex = normalized.indexOf("@"); |
| 26 | |
| 27 | if (atIndex <= 0 || atIndex === normalized.length - 1 || normalized.indexOf("@", atIndex + 1) !== -1) { |
| 28 | return null; |
| 29 | } |
| 30 | |
| 31 | const localPart = normalized.slice(0, atIndex); |
| 32 | const domain = normalized.slice(atIndex + 1); |
| 33 | const plusIndex = localPart.indexOf("+"); |
| 34 | const canonicalLocalPart = plusIndex === -1 ? localPart : localPart.slice(0, plusIndex); |
| 35 | |
| 36 | if (!canonicalLocalPart) { |
| 37 | return null; |
| 38 | } |
| 39 | |
| 40 | return { |
| 41 | raw, |
| 42 | domain, |
| 43 | canonicalLocalPart, |
| 44 | }; |
| 45 | } |
| 46 | |
| 47 | function toArrayBuffer(value: string | ArrayBuffer | Uint8Array): ArrayBuffer { |
| 48 | if (value instanceof ArrayBuffer) { |
| 49 | return value; |
| 50 | } |
| 51 | |
| 52 | if (value instanceof Uint8Array) { |
| 53 | const buffer = new ArrayBuffer(value.byteLength); |
| 54 | new Uint8Array(buffer).set(value); |
| 55 | return buffer; |
| 56 | } |
| 57 | |
| 58 | const encoded = new TextEncoder().encode(value); |
| 59 | const buffer = new ArrayBuffer(encoded.byteLength); |
| 60 | new Uint8Array(buffer).set(encoded); |
| 61 | return buffer; |
| 62 | } |
| 63 | |
| 64 | export async function handleIncomingEmail(message: ForwardableEmailMessage, env: Env, ctx: ExecutionContext) { |
| 65 | const recipient = parseRecipientAddress(message.to); |
| 66 | |
| 67 | if (!recipient) { |
| 68 | logger.warn("email_rejected", "Rejected inbound email with invalid address", { |
| 69 | address: message.to.trim(), |
| 70 | }); |
| 71 | message.setReject("Invalid address"); |
| 72 | return; |
| 73 | } |
| 74 | |
| 75 | const db = createDb(env.DB); |
| 76 | const [domainRecord, inbox] = await Promise.all([ |
| 77 | db.query.domains.findFirst({ |
| 78 | where: eq(domains.domain, recipient.domain), |
| 79 | }), |
| 80 | db.query.inboxes.findFirst({ |
| 81 | where: and(eq(inboxes.localPart, recipient.canonicalLocalPart), eq(inboxes.domain, recipient.domain)), |
| 82 | }), |
| 83 | ]); |
| 84 | |
| 85 | if (!domainRecord?.isActive || !inbox) { |
| 86 | logger.warn("email_rejected", "Rejected inbound email for missing or disabled inbox", { |
| 87 | address: recipient.raw, |
| 88 | canonicalAddress: `${recipient.canonicalLocalPart}@${recipient.domain}`, |
| 89 | domain: recipient.domain, |
| 90 | }); |
| 91 | message.setReject("Address not found"); |
| 92 | return; |
| 93 | } |
| 94 | |
| 95 | if (!inbox.isPermanent && inbox.expiresAt && inbox.expiresAt.getTime() < Date.now()) { |
| 96 | logger.warn("email_rejected", "Rejected inbound email for expired inbox", { |
| 97 | address: recipient.raw, |
| 98 | canonicalAddress: inbox.fullAddress, |
| 99 | }); |
| 100 | message.setReject("Inbox expired"); |
| 101 | return; |
| 102 | } |
| 103 | |
| 104 | if (message.rawSize > MAX_INBOUND_EMAIL_SIZE_BYTES) { |
| 105 | logger.warn("email_rejected", "Rejected inbound email that exceeded size limit", { |
| 106 | address: recipient.raw, |
| 107 | canonicalAddress: inbox.fullAddress, |
| 108 | sizeBytes: message.rawSize, |
| 109 | maxSizeBytes: MAX_INBOUND_EMAIL_SIZE_BYTES, |
| 110 | }); |
| 111 | message.setReject("Message too large"); |
| 112 | return; |
| 113 | } |
| 114 | |
| 115 | const quotaRows = await db.select({ total: count() }).from(emails).where(eq(emails.inboxId, inbox.id)); |
| 116 | const emailCount = quotaRows[0]?.total ?? 0; |
| 117 | |
| 118 | if (emailCount >= MAX_EMAILS_PER_INBOX) { |
| 119 | logger.warn("email_rejected", "Rejected inbound email for inbox quota", { |
| 120 | address: recipient.raw, |
| 121 | canonicalAddress: inbox.fullAddress, |
| 122 | emailCount, |
| 123 | maxEmailsPerInbox: MAX_EMAILS_PER_INBOX, |
| 124 | }); |
| 125 | message.setReject("Inbox is full"); |
| 126 | return; |
| 127 | } |
| 128 | |
| 129 | try { |
| 130 | const rawEmail = await new Response(message.raw).arrayBuffer(); |
| 131 | const parser = new PostalMime(); |
| 132 | const parsed = await parser.parse(rawEmail); |
| 133 | const parsedAttachments = parsed.attachments ?? []; |
| 134 | |
| 135 | if (parsedAttachments.length > MAX_INBOUND_ATTACHMENTS) { |
| 136 | logger.warn("email_rejected", "Rejected inbound email with too many attachments", { |
| 137 | address: recipient.raw, |
| 138 | canonicalAddress: inbox.fullAddress, |
| 139 | attachmentCount: parsedAttachments.length, |
| 140 | maxAttachments: MAX_INBOUND_ATTACHMENTS, |
| 141 | }); |
| 142 | message.setReject("Too many attachments"); |
| 143 | return; |
| 144 | } |
| 145 | |
| 146 | const emailId = nanoid(); |
| 147 | const receivedAt = new Date(); |
| 148 | const bodyKey = getBodyStorageKey(emailId); |
| 149 | const attachmentRecords: Array<{ |
| 150 | row: typeof attachments.$inferInsert; |
| 151 | content: ArrayBuffer; |
| 152 | }> = []; |
| 153 | |
| 154 | for (const attachment of parsedAttachments) { |
| 155 | const attachmentId = nanoid(); |
| 156 | const content = toArrayBuffer(attachment.content); |
| 157 | attachmentRecords.push({ |
| 158 | content, |
| 159 | row: { |
| 160 | id: attachmentId, |
| 161 | emailId, |
| 162 | filename: attachment.filename, |
| 163 | contentType: attachment.mimeType, |
| 164 | sizeBytes: content.byteLength, |
| 165 | storageKey: getAttachmentStorageKey(emailId, attachmentId, attachment.filename), |
| 166 | }, |
| 167 | }); |
| 168 | } |
| 169 | |
| 170 | const attachmentRows = attachmentRecords.map((attachment) => attachment.row); |
| 171 | |
| 172 | const batchOps = [ |
| 173 | db.insert(emails).values({ |
| 174 | id: emailId, |
| 175 | inboxId: inbox.id, |
| 176 | recipientAddress: recipient.raw, |
| 177 | fromAddress: parsed.from?.address ?? message.from, |
| 178 | fromName: parsed.from?.name ?? null, |
| 179 | subject: parsed.subject ?? "(no subject)", |
| 180 | receivedAt, |
| 181 | sizeBytes: message.rawSize, |
| 182 | hasAttachments: attachmentRecords.length > 0, |
| 183 | bodyKey, |
| 184 | }), |
| 185 | ...(attachmentRows.length > 0 ? [db.insert(attachments).values(attachmentRows)] : []), |
| 186 | ] as const; |
| 187 | |
| 188 | await db.batch(batchOps as any); |
| 189 | |
| 190 | try { |
| 191 | await storeRawEmail(env.STORAGE, emailId, rawEmail); |
| 192 | await storeEmailBody(env.STORAGE, emailId, { |
| 193 | text: parsed.text, |
| 194 | html: typeof parsed.html === "string" ? parsed.html : null, |
| 195 | }); |
| 196 | |
| 197 | for (const attachment of attachmentRecords) { |
| 198 | await storeAttachment(env.STORAGE, emailId, attachment.row.id, { |
| 199 | content: attachment.content, |
| 200 | filename: attachment.row.filename, |
| 201 | contentType: attachment.row.contentType, |
| 202 | }); |
| 203 | } |
| 204 | } catch (storageError) { |
| 205 | const [cleanupResult, deleteResult] = await Promise.allSettled([ |
| 206 | deleteStorageForEmails(env.STORAGE, [emailId]), |
| 207 | db.delete(emails).where(eq(emails.id, emailId)), |
| 208 | ]); |
| 209 | |
| 210 | logger.error("email_storage_failed", "Failed to persist inbound email content", { |
| 211 | address: recipient.raw, |
| 212 | canonicalAddress: inbox.fullAddress, |
| 213 | emailId, |
| 214 | cleanupFailed: cleanupResult.status === "rejected", |
| 215 | deleteFailed: deleteResult.status === "rejected", |
| 216 | ...errorContext(storageError), |
| 217 | }); |
| 218 | message.setReject("Could not process email"); |
| 219 | return; |
| 220 | } |
| 221 | |
| 222 | const notification = toNewEmailNotification({ |
| 223 | id: emailId, |
| 224 | recipientAddress: recipient.raw, |
| 225 | fromAddress: parsed.from?.address ?? message.from, |
| 226 | fromName: parsed.from?.name ?? null, |
| 227 | subject: parsed.subject ?? "(no subject)", |
| 228 | receivedAt, |
| 229 | isRead: false, |
| 230 | hasAttachments: attachmentRecords.length > 0, |
| 231 | sizeBytes: message.rawSize, |
| 232 | }); |
| 233 | |
| 234 | const stub = env.INBOX_WS.getByName(inbox.fullAddress); |
| 235 | ctx.waitUntil( |
| 236 | (async () => { |
| 237 | try { |
| 238 | await stub.notifyNewEmail(notification); |
| 239 | } catch (error) { |
| 240 | logger.warn("email_notification_failed", "Stored inbound email but websocket notification failed", { |
| 241 | address: inbox.fullAddress, |
| 242 | emailId, |
| 243 | ...errorContext(error), |
| 244 | }); |
| 245 | } |
| 246 | })(), |
| 247 | ); |
| 248 | |
| 249 | logger.info("email_stored", "Stored inbound email", { |
| 250 | address: inbox.fullAddress, |
| 251 | recipientAddress: recipient.raw, |
| 252 | emailId, |
| 253 | attachmentCount: attachmentRecords.length, |
| 254 | subject: parsed.subject ?? "(no subject)", |
| 255 | }); |
| 256 | } catch (error) { |
| 257 | logger.error("email_processing_failed", "Failed to process inbound email", { |
| 258 | address: recipient.raw, |
| 259 | canonicalAddress: inbox.fullAddress, |
| 260 | ...errorContext(error), |
| 261 | }); |
| 262 | message.setReject("Could not process email"); |
| 263 | } |
| 264 | } |