Skip to content
File

Blob: tests/worker/email-handler.test.ts

typescript444 lines
1import { createExecutionContext, env } from "cloudflare:test";
2import { beforeEach, describe, expect, it, vi } from "vitest";
3import type { NewEmailNotification } from "@/shared/contracts";
4import { createDb } from "@/worker/db";
5import { attachments, emails } from "@/worker/db/schema";
6import { handleIncomingEmail } from "@/worker/email-handler";
7import { getRawStorageKey, readEmailBody } from "@/worker/services/storage";
8import { resetWorkerState, seedDomain, seedEmail, seedInbox } from "./api/helpers";
9 
10const encoder = new TextEncoder();
11 
12function getDb() {
13 return createDb(env.DB.withSession("first-primary"));
14}
15 
16interface FakeMessageOptions {
17 from?: string;
18 raw?: string;
19 rawSize?: number;
20 to: string;
21}
22 
23function createExecutionContextStub() {
24 const waitUntilPromises: Promise<unknown>[] = [];
25 const ctx = createExecutionContext();
26 const originalWaitUntil = ctx.waitUntil.bind(ctx);
27 
28 ctx.waitUntil = (promise) => {
29 waitUntilPromises.push(promise);
30 originalWaitUntil(promise);
31 };
32 
33 return {
34 ctx,
35 waitUntilPromises,
36 };
37}
38 
39function createMessage(options: FakeMessageOptions) {
40 const rejectReasons: string[] = [];
41 const raw =
42 options.raw ??
43 buildMimeMessage({
44 from: options.from ?? "sender@example.com",
45 to: options.to,
46 subject: "Test subject",
47 text: "Plain text body",
48 });
49 const rawBytes = encoder.encode(raw);
50 const message: ForwardableEmailMessage = {
51 from: options.from ?? "sender@example.com",
52 forward: vi.fn<(rcptTo: string, headers?: Headers) => Promise<EmailSendResult>>().mockResolvedValue({
53 messageId: "forwarded",
54 }),
55 headers: new Headers(),
56 raw: new ReadableStream<Uint8Array>({
57 start(controller) {
58 controller.enqueue(rawBytes);
59 controller.close();
60 },
61 }),
62 rawSize: options.rawSize ?? rawBytes.byteLength,
63 reply: vi.fn<(message: EmailMessage) => Promise<EmailSendResult>>().mockResolvedValue({
64 messageId: "reply",
65 }),
66 setReject(reason: string) {
67 rejectReasons.push(reason);
68 },
69 to: options.to,
70 };
71 
72 return {
73 message,
74 rejectReasons,
75 };
76}
77 
78function mockInboxWebSocketNotification(
79 address: string,
80 implementation?: (payload: NewEmailNotification) => Promise<void>,
81) {
82 const stub = env.INBOX_WS.getByName(address);
83 const notifyNewEmail = vi.spyOn(stub, "notifyNewEmail");
84 
85 if (implementation) {
86 notifyNewEmail.mockImplementation(implementation);
87 } else {
88 notifyNewEmail.mockResolvedValue(undefined);
89 }
90 
91 vi.spyOn(env.INBOX_WS, "getByName").mockReturnValue(stub);
92 
93 return notifyNewEmail;
94}
95 
96function buildMimeMessage(options: {
97 attachments?: Array<{
98 content: string;
99 contentType: string;
100 filename: string;
101 }>;
102 from: string;
103 html?: string;
104 subject: string;
105 text?: string;
106 to: string;
107}) {
108 const attachments = options.attachments ?? [];
109 const hasText = typeof options.text === "string" && options.text.length > 0;
110 const hasHtml = typeof options.html === "string" && options.html.length > 0;
111 
112 if (attachments.length === 0 && hasText && !hasHtml) {
113 return [
114 `From: <${options.from}>`,
115 `To: <${options.to}>`,
116 `Subject: ${options.subject}`,
117 "MIME-Version: 1.0",
118 "Content-Type: text/plain; charset=utf-8",
119 "",
120 options.text,
121 "",
122 ].join("\r\n");
123 }
124 
125 const mixedBoundary = "mixed_boundary";
126 const alternativeBoundary = "alt_boundary";
127 const parts: string[] = [];
128 
129 if (hasText || hasHtml) {
130 if (hasText && hasHtml) {
131 parts.push(
132 [
133 `--${mixedBoundary}`,
134 `Content-Type: multipart/alternative; boundary=\"${alternativeBoundary}\"`,
135 "",
136 `--${alternativeBoundary}`,
137 "Content-Type: text/plain; charset=utf-8",
138 "Content-Transfer-Encoding: 7bit",
139 "",
140 options.text,
141 `--${alternativeBoundary}`,
142 "Content-Type: text/html; charset=utf-8",
143 "Content-Transfer-Encoding: 7bit",
144 "",
145 options.html,
146 `--${alternativeBoundary}--`,
147 ].join("\r\n"),
148 );
149 } else if (hasText) {
150 parts.push(
151 [
152 `--${mixedBoundary}`,
153 "Content-Type: text/plain; charset=utf-8",
154 "Content-Transfer-Encoding: 7bit",
155 "",
156 options.text,
157 ].join("\r\n"),
158 );
159 } else if (hasHtml) {
160 parts.push(
161 [
162 `--${mixedBoundary}`,
163 "Content-Type: text/html; charset=utf-8",
164 "Content-Transfer-Encoding: 7bit",
165 "",
166 options.html,
167 ].join("\r\n"),
168 );
169 }
170 }
171 
172 for (const attachment of attachments) {
173 parts.push(
174 [
175 `--${mixedBoundary}`,
176 `Content-Type: ${attachment.contentType}; name=\"${attachment.filename}\"`,
177 `Content-Disposition: attachment; filename=\"${attachment.filename}\"`,
178 "Content-Transfer-Encoding: base64",
179 "",
180 Buffer.from(attachment.content, "utf8").toString("base64"),
181 ].join("\r\n"),
182 );
183 }
184 
185 parts.push(`--${mixedBoundary}--`);
186 
187 return [
188 `From: <${options.from}>`,
189 `To: <${options.to}>`,
190 `Subject: ${options.subject}`,
191 "MIME-Version: 1.0",
192 `Content-Type: multipart/mixed; boundary=\"${mixedBoundary}\"`,
193 "",
194 parts.join("\r\n"),
195 "",
196 ].join("\r\n");
197}
198 
199describe("worker email handler", () => {
200 beforeEach(async () => {
201 vi.restoreAllMocks();
202 await resetWorkerState();
203 });
204 
205 it("rejects invalid recipient addresses", async () => {
206 const { ctx } = createExecutionContextStub();
207 const { message, rejectReasons } = createMessage({
208 to: "not-an-email",
209 });
210 
211 await handleIncomingEmail(message, env, ctx);
212 
213 expect(rejectReasons).toEqual(["Invalid address"]);
214 });
215 
216 it("rejects messages for missing or inactive domains", async () => {
217 await seedDomain("inactive.test", false);
218 const { ctx: missingCtx } = createExecutionContextStub();
219 const missing = createMessage({ to: "reader@missing.test" });
220 const { ctx: inactiveCtx } = createExecutionContextStub();
221 const inactive = createMessage({ to: "reader@inactive.test" });
222 
223 await handleIncomingEmail(missing.message, env, missingCtx);
224 await handleIncomingEmail(inactive.message, env, inactiveCtx);
225 
226 expect(missing.rejectReasons).toEqual(["Address not found"]);
227 expect(inactive.rejectReasons).toEqual(["Address not found"]);
228 });
229 
230 it("rejects messages when the inbox does not exist", async () => {
231 await seedDomain("mail.test", true);
232 const { ctx } = createExecutionContextStub();
233 const { message, rejectReasons } = createMessage({
234 to: "reader@mail.test",
235 });
236 
237 await handleIncomingEmail(message, env, ctx);
238 
239 expect(rejectReasons).toEqual(["Address not found"]);
240 });
241 
242 it("rejects messages for expired temporary inboxes", async () => {
243 await seedDomain("mail.test", true);
244 await seedInbox({
245 address: "reader@mail.test",
246 createdAt: new Date("2026-01-01T00:00:00.000Z"),
247 expiresAt: new Date("2026-01-02T00:00:00.000Z"),
248 });
249 const { ctx } = createExecutionContextStub();
250 const { message, rejectReasons } = createMessage({
251 to: "reader@mail.test",
252 });
253 
254 await handleIncomingEmail(message, env, ctx);
255 
256 expect(rejectReasons).toEqual(["Inbox expired"]);
257 });
258 
259 it("rejects oversize messages", async () => {
260 await seedDomain("mail.test", true);
261 await seedInbox({
262 address: "reader@mail.test",
263 });
264 const { ctx } = createExecutionContextStub();
265 const { message, rejectReasons } = createMessage({
266 to: "reader@mail.test",
267 rawSize: 10 * 1024 * 1024 + 1,
268 });
269 
270 await handleIncomingEmail(message, env, ctx);
271 
272 expect(rejectReasons).toEqual(["Message too large"]);
273 });
274 
275 it("rejects inboxes that are over quota", async () => {
276 await seedDomain("mail.test", true);
277 const inbox = await seedInbox({
278 address: "reader@mail.test",
279 });
280 
281 for (let index = 0; index < 100; index += 1) {
282 await seedEmail({
283 id: `email-${index}`,
284 address: inbox.fullAddress,
285 inboxId: inbox.id,
286 });
287 }
288 
289 const { ctx } = createExecutionContextStub();
290 const { message, rejectReasons } = createMessage({
291 to: "reader@mail.test",
292 });
293 
294 await handleIncomingEmail(message, env, ctx);
295 
296 expect(rejectReasons).toEqual(["Inbox is full"]);
297 });
298 
299 it("rejects messages with too many attachments", async () => {
300 await seedDomain("mail.test", true);
301 await seedInbox({
302 address: "reader@mail.test",
303 });
304 const attachments = Array.from({ length: 11 }, (_, index) => ({
305 filename: `attachment-${index}.txt`,
306 contentType: "text/plain",
307 content: `attachment-${index}`,
308 }));
309 const { ctx } = createExecutionContextStub();
310 const { message, rejectReasons } = createMessage({
311 to: "reader@mail.test",
312 raw: buildMimeMessage({
313 from: "sender@example.com",
314 to: "reader@mail.test",
315 subject: "Too many attachments",
316 text: "Hello",
317 attachments,
318 }),
319 });
320 
321 await handleIncomingEmail(message, env, ctx);
322 
323 expect(rejectReasons).toEqual(["Too many attachments"]);
324 });
325 
326 it("stores raw email, parsed body, attachment metadata, and R2 objects", async () => {
327 await seedDomain("mail.test", true);
328 const inbox = await seedInbox({
329 address: "reader@mail.test",
330 });
331 const notifyNewEmail = mockInboxWebSocketNotification(inbox.fullAddress);
332 const { ctx, waitUntilPromises } = createExecutionContextStub();
333 const { message, rejectReasons } = createMessage({
334 to: "reader@mail.test",
335 raw: buildMimeMessage({
336 from: "sender@example.com",
337 to: "reader@mail.test",
338 subject: "Stored message",
339 text: "Plain text body",
340 html: "<p>HTML body</p>",
341 attachments: [
342 {
343 filename: "invoice.txt",
344 contentType: "text/plain",
345 content: "attachment body",
346 },
347 ],
348 }),
349 });
350 
351 await handleIncomingEmail(message, env, ctx);
352 await Promise.allSettled(waitUntilPromises);
353 
354 expect(rejectReasons).toEqual([]);
355 
356 const storedEmail = await getDb().query.emails.findFirst({
357 where: (table, { eq }) => eq(table.inboxId, inbox.id),
358 });
359 const storedAttachments = await getDb().query.attachments.findMany({
360 where: (table, { eq }) => eq(table.emailId, storedEmail?.id ?? "missing"),
361 });
362 const storedBody = await readEmailBody(env.STORAGE, storedEmail?.bodyKey ?? null);
363 
364 expect(storedEmail?.recipientAddress).toBe("reader@mail.test");
365 expect(storedEmail?.subject).toBe("Stored message");
366 expect(storedEmail?.hasAttachments).toBe(true);
367 expect(storedBody).toEqual({
368 text: "Plain text body\n",
369 html: "<p>HTML body</p>\n",
370 });
371 expect(storedAttachments).toHaveLength(1);
372 expect(await env.STORAGE.get(getRawStorageKey(storedEmail?.id ?? "missing"))).not.toBeNull();
373 expect(await env.STORAGE.get(storedAttachments[0]?.storageKey ?? "missing")).not.toBeNull();
374 expect(notifyNewEmail).toHaveBeenCalledTimes(1);
375 });
376 
377 it("routes plus aliases to the base inbox while preserving the delivered recipient address", async () => {
378 await seedDomain("mail.test", true);
379 const inbox = await seedInbox({
380 address: "reader@mail.test",
381 });
382 mockInboxWebSocketNotification(inbox.fullAddress);
383 const { ctx, waitUntilPromises } = createExecutionContextStub();
384 const { message, rejectReasons } = createMessage({
385 to: "Reader+tag@mail.test",
386 });
387 
388 await handleIncomingEmail(message, env, ctx);
389 await Promise.allSettled(waitUntilPromises);
390 
391 const storedEmail = await getDb().query.emails.findFirst({
392 where: (table, { eq }) => eq(table.inboxId, inbox.id),
393 });
394 
395 expect(rejectReasons).toEqual([]);
396 expect(storedEmail?.recipientAddress).toBe("Reader+tag@mail.test");
397 });
398 
399 it("cleans up and rejects the message when storage persistence fails", async () => {
400 await seedDomain("mail.test", true);
401 const inbox = await seedInbox({
402 address: "reader@mail.test",
403 });
404 vi.spyOn(env.STORAGE, "put").mockRejectedValueOnce(new Error("storage failed"));
405 const { ctx } = createExecutionContextStub();
406 const { message, rejectReasons } = createMessage({
407 to: "reader@mail.test",
408 });
409 
410 await handleIncomingEmail(message, env, ctx);
411 
412 const storedEmail = await getDb().query.emails.findFirst({
413 where: (table, { eq }) => eq(table.inboxId, inbox.id),
414 });
415 const storedAttachmentRows = await getDb().select().from(attachments);
416 
417 expect(rejectReasons).toEqual(["Could not process email"]);
418 expect(storedEmail).toBeUndefined();
419 expect(storedAttachmentRows).toHaveLength(0);
420 });
421 
422 it("does not roll back successful storage when websocket notification fails", async () => {
423 await seedDomain("mail.test", true);
424 const inbox = await seedInbox({
425 address: "reader@mail.test",
426 });
427 mockInboxWebSocketNotification(inbox.fullAddress, async () => {
428 throw new Error("notify failed");
429 });
430 const { ctx, waitUntilPromises } = createExecutionContextStub();
431 const { message, rejectReasons } = createMessage({
432 to: "reader@mail.test",
433 });
434 
435 await handleIncomingEmail(message, env, ctx);
436 await Promise.allSettled(waitUntilPromises);
437 
438 const storedEmailRows = await getDb().select().from(emails);
439 
440 expect(rejectReasons).toEqual([]);
441 expect(storedEmailRows).toHaveLength(1);
442 });
443});