Skip to content
File

Blob: src/worker/r2/blobs.ts

typescript163 lines
1const MULTIPART_PART_BYTES = 8 * 1024 * 1024;
2 
3export class BodyTooLargeError extends Error {
4 constructor(readonly maxBytes: number) {
5 super("Request body exceeds configured maximum size");
6 }
7}
8 
9function chunkBytes(chunk: Uint8Array | ArrayBuffer): Uint8Array<ArrayBuffer> {
10 const source = chunk instanceof Uint8Array ? chunk : new Uint8Array(chunk);
11 const copy = new Uint8Array(source.byteLength);
12 copy.set(source);
13 return copy;
14}
15 
16function appendBytes(left: Uint8Array, right: Uint8Array): Uint8Array<ArrayBuffer> {
17 if (left.byteLength === 0) return chunkBytes(right);
18 const combined = new Uint8Array(left.byteLength + right.byteLength);
19 combined.set(left, 0);
20 combined.set(right, left.byteLength);
21 return combined;
22}
23 
24async function putKnownLength(
25 bucket: R2Bucket,
26 input: {
27 key: string;
28 body: ReadableStream | null;
29 contentType: string | null;
30 contentLength: number;
31 maxBytes: number;
32 },
33): Promise<{ size: number; etag: string | null }> {
34 if (input.contentLength > input.maxBytes) throw new BodyTooLargeError(input.maxBytes);
35 if (!input.body) {
36 const object = await bucket.put(input.key, new Uint8Array(), {
37 httpMetadata: { contentType: input.contentType ?? "application/octet-stream" },
38 });
39 return { size: object.size ?? 0, etag: object.etag ?? null };
40 }
41 
42 const fixedLengthStream = new FixedLengthStream(input.contentLength);
43 const writer = fixedLengthStream.writable.getWriter();
44 const reader = input.body.getReader();
45 const putPromise = bucket.put(input.key, fixedLengthStream.readable, {
46 httpMetadata: { contentType: input.contentType ?? "application/octet-stream" },
47 });
48 let size = 0;
49 
50 try {
51 while (true) {
52 const next = await reader.read();
53 if (next.done) break;
54 const bytes = chunkBytes(next.value);
55 size += bytes.byteLength;
56 if (size > input.maxBytes) throw new BodyTooLargeError(input.maxBytes);
57 await writer.write(bytes);
58 }
59 if (size !== input.contentLength) {
60 throw new Error("Request body length did not match Content-Length");
61 }
62 await writer.close();
63 const object = await putPromise;
64 return { size: object.size ?? size, etag: object.etag ?? null };
65 } catch (cause) {
66 try {
67 await writer.abort(cause);
68 } catch {}
69 try {
70 await reader.cancel(cause);
71 } catch {}
72 try {
73 await bucket.delete(input.key);
74 } catch {}
75 throw cause;
76 }
77}
78 
79async function putMultipart(
80 bucket: R2Bucket,
81 input: {
82 key: string;
83 body: ReadableStream | null;
84 contentType: string | null;
85 maxBytes: number;
86 },
87): Promise<{ size: number; etag: string | null }> {
88 if (!input.body) {
89 const object = await bucket.put(input.key, new Uint8Array(), {
90 httpMetadata: { contentType: input.contentType ?? "application/octet-stream" },
91 });
92 return { size: object.size ?? 0, etag: object.etag ?? null };
93 }
94 
95 const upload = await bucket.createMultipartUpload(input.key, {
96 httpMetadata: { contentType: input.contentType ?? "application/octet-stream" },
97 });
98 const reader = input.body.getReader();
99 const parts: R2UploadedPart[] = [];
100 let buffered: Uint8Array<ArrayBuffer> = new Uint8Array();
101 let partNumber = 1;
102 let size = 0;
103 
104 try {
105 while (true) {
106 const next = await reader.read();
107 if (next.done) break;
108 const bytes = chunkBytes(next.value);
109 size += bytes.byteLength;
110 if (size > input.maxBytes) throw new BodyTooLargeError(input.maxBytes);
111 buffered = appendBytes(buffered, bytes);
112 
113 while (buffered.byteLength >= MULTIPART_PART_BYTES) {
114 const partBytes = buffered.slice(0, MULTIPART_PART_BYTES);
115 buffered = buffered.slice(MULTIPART_PART_BYTES);
116 parts.push(await upload.uploadPart(partNumber, partBytes));
117 partNumber += 1;
118 }
119 }
120 
121 if (buffered.byteLength > 0 || parts.length === 0) {
122 parts.push(await upload.uploadPart(partNumber, buffered));
123 }
124 const object = await upload.complete(parts);
125 return { size: object.size ?? size, etag: object.etag ?? null };
126 } catch (cause) {
127 try {
128 await reader.cancel(cause);
129 } catch {}
130 try {
131 await upload.abort();
132 } catch {}
133 try {
134 await bucket.delete(input.key);
135 } catch {}
136 throw cause;
137 }
138}
139 
140export async function putFileBlob(
141 bucket: R2Bucket,
142 input: {
143 key: string;
144 body: ReadableStream | null;
145 contentType: string | null;
146 contentLength: number | null;
147 maxBytes: number;
148 },
149): Promise<{ size: number; etag: string | null }> {
150 if (input.contentLength !== null) {
151 return await putKnownLength(bucket, { ...input, contentLength: input.contentLength });
152 }
153 return await putMultipart(bucket, input);
154}
155 
156export async function getFileBlob(bucket: R2Bucket, key: string): Promise<R2ObjectBody | null> {
157 return await bucket.get(key);
158}
159 
160export async function deleteFileBlob(bucket: R2Bucket, key: string): Promise<void> {
161 await bucket.delete(key);
162}