Skip to content
File

Blob: src/worker/services/storage.ts

typescript136 lines
1const BODY_PREFIX = "bodies";
2const ATTACHMENT_PREFIX = "attachments";
3const RAW_PREFIX = "raw";
4const STORAGE_DELETE_BATCH_SIZE = 1000;
5const STORAGE_EMAIL_BATCH_SIZE = 25;
6 
7function chunk<T>(items: T[], size: number) {
8 const batches: T[][] = [];
9 
10 for (let index = 0; index < items.length; index += size) {
11 batches.push(items.slice(index, index + size));
12 }
13 
14 return batches;
15}
16 
17function safeFilename(filename: string | null | undefined) {
18 return (filename ?? "attachment.bin").replace(/[^a-zA-Z0-9._-]/g, "_");
19}
20 
21export function getBodyStorageKey(emailId: string) {
22 return `${BODY_PREFIX}/${emailId}.json`;
23}
24 
25export function getRawStorageKey(emailId: string) {
26 return `${RAW_PREFIX}/${emailId}.eml`;
27}
28 
29export function getAttachmentStorageKey(emailId: string, attachmentId: string, filename?: string | null) {
30 return `${ATTACHMENT_PREFIX}/${emailId}/${attachmentId}/${safeFilename(filename)}`;
31}
32 
33export async function storeRawEmail(bucket: R2Bucket, emailId: string, rawEmail: ArrayBuffer) {
34 const key = getRawStorageKey(emailId);
35 await bucket.put(key, rawEmail, {
36 httpMetadata: {
37 contentType: "message/rfc822",
38 },
39 });
40 return key;
41}
42 
43export async function storeEmailBody(
44 bucket: R2Bucket,
45 emailId: string,
46 payload: { text?: string | null; html?: string | null },
47) {
48 const key = getBodyStorageKey(emailId);
49 await bucket.put(key, JSON.stringify({ text: payload.text ?? null, html: payload.html ?? null }), {
50 httpMetadata: {
51 contentType: "application/json; charset=utf-8",
52 },
53 });
54 return key;
55}
56 
57export async function readEmailBody(bucket: R2Bucket, key: string | null) {
58 if (!key) {
59 return { text: null, html: null };
60 }
61 
62 const object = await bucket.get(key);
63 if (!object) {
64 return { text: null, html: null };
65 }
66 
67 const body = (await object.json()) as { text?: string | null; html?: string | null };
68 return {
69 text: body.text ?? null,
70 html: body.html ?? null,
71 };
72}
73 
74export async function storeAttachment(
75 bucket: R2Bucket,
76 emailId: string,
77 attachmentId: string,
78 attachment: {
79 content: ArrayBuffer;
80 filename?: string | null;
81 contentType?: string | null;
82 },
83) {
84 const key = getAttachmentStorageKey(emailId, attachmentId, attachment.filename);
85 await bucket.put(key, attachment.content, {
86 httpMetadata: {
87 contentType: attachment.contentType ?? "application/octet-stream",
88 },
89 });
90 return key;
91}
92 
93export async function deleteStorageKeys(bucket: R2Bucket, keys: Array<string | null | undefined>) {
94 const filtered = [...new Set(keys.filter((key): key is string => Boolean(key)))];
95 if (filtered.length === 0) {
96 return;
97 }
98 
99 for (const batch of chunk(filtered, STORAGE_DELETE_BATCH_SIZE)) {
100 await bucket.delete(batch);
101 }
102}
103 
104async function listKeysByPrefix(bucket: R2Bucket, prefix: string) {
105 const keys: string[] = [];
106 let cursor: string | undefined;
107 
108 do {
109 const result = await bucket.list({ prefix, cursor });
110 for (const object of result.objects) {
111 keys.push(object.key);
112 }
113 cursor = result.truncated ? result.cursor : undefined;
114 } while (cursor);
115 
116 return keys;
117}
118 
119export async function deleteStorageForEmails(bucket: R2Bucket, emailIds: string[]) {
120 if (emailIds.length === 0) {
121 return;
122 }
123 
124 for (const emailBatch of chunk(emailIds, STORAGE_EMAIL_BATCH_SIZE)) {
125 const keyGroups = await Promise.all(
126 emailBatch.flatMap((emailId) => [
127 listKeysByPrefix(bucket, `${BODY_PREFIX}/${emailId}`),
128 listKeysByPrefix(bucket, `${ATTACHMENT_PREFIX}/${emailId}/`),
129 listKeysByPrefix(bucket, `${RAW_PREFIX}/${emailId}`),
130 ]),
131 );
132 
133 await deleteStorageKeys(bucket, keyGroups.flat());
134 }
135}