File
Blob: src/worker/r2/blobs.ts
| 1 | const MULTIPART_PART_BYTES = 8 * 1024 * 1024; |
| 2 | |
| 3 | export class BodyTooLargeError extends Error { |
| 4 | constructor(readonly maxBytes: number) { |
| 5 | super("Request body exceeds configured maximum size"); |
| 6 | } |
| 7 | } |
| 8 | |
| 9 | function 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 | |
| 16 | function 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 | |
| 24 | async 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 | |
| 79 | async 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 | |
| 140 | export 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 | |
| 156 | export async function getFileBlob(bucket: R2Bucket, key: string): Promise<R2ObjectBody | null> { |
| 157 | return await bucket.get(key); |
| 158 | } |
| 159 | |
| 160 | export async function deleteFileBlob(bucket: R2Bucket, key: string): Promise<void> { |
| 161 | await bucket.delete(key); |
| 162 | } |