const MULTIPART_PART_BYTES = 8 * 1024 * 1024; export class BodyTooLargeError extends Error { constructor(readonly maxBytes: number) { super("Request body exceeds configured maximum size"); } } function chunkBytes(chunk: Uint8Array | ArrayBuffer): Uint8Array { const source = chunk instanceof Uint8Array ? chunk : new Uint8Array(chunk); const copy = new Uint8Array(source.byteLength); copy.set(source); return copy; } function appendBytes(left: Uint8Array, right: Uint8Array): Uint8Array { if (left.byteLength === 0) return chunkBytes(right); const combined = new Uint8Array(left.byteLength + right.byteLength); combined.set(left, 0); combined.set(right, left.byteLength); return combined; } async function putKnownLength( bucket: R2Bucket, input: { key: string; body: ReadableStream | null; contentType: string | null; contentLength: number; maxBytes: number; }, ): Promise<{ size: number; etag: string | null }> { if (input.contentLength > input.maxBytes) throw new BodyTooLargeError(input.maxBytes); if (!input.body) { const object = await bucket.put(input.key, new Uint8Array(), { httpMetadata: { contentType: input.contentType ?? "application/octet-stream" }, }); return { size: object.size ?? 0, etag: object.etag ?? null }; } const fixedLengthStream = new FixedLengthStream(input.contentLength); const writer = fixedLengthStream.writable.getWriter(); const reader = input.body.getReader(); const putPromise = bucket.put(input.key, fixedLengthStream.readable, { httpMetadata: { contentType: input.contentType ?? "application/octet-stream" }, }); let size = 0; try { while (true) { const next = await reader.read(); if (next.done) break; const bytes = chunkBytes(next.value); size += bytes.byteLength; if (size > input.maxBytes) throw new BodyTooLargeError(input.maxBytes); await writer.write(bytes); } if (size !== input.contentLength) { throw new Error("Request body length did not match Content-Length"); } await writer.close(); const object = await putPromise; return { size: object.size ?? size, etag: object.etag ?? null }; } catch (cause) { try { await writer.abort(cause); } catch {} try { await reader.cancel(cause); } catch {} try { await bucket.delete(input.key); } catch {} throw cause; } } async function putMultipart( bucket: R2Bucket, input: { key: string; body: ReadableStream | null; contentType: string | null; maxBytes: number; }, ): Promise<{ size: number; etag: string | null }> { if (!input.body) { const object = await bucket.put(input.key, new Uint8Array(), { httpMetadata: { contentType: input.contentType ?? "application/octet-stream" }, }); return { size: object.size ?? 0, etag: object.etag ?? null }; } const upload = await bucket.createMultipartUpload(input.key, { httpMetadata: { contentType: input.contentType ?? "application/octet-stream" }, }); const reader = input.body.getReader(); const parts: R2UploadedPart[] = []; let buffered: Uint8Array = new Uint8Array(); let partNumber = 1; let size = 0; try { while (true) { const next = await reader.read(); if (next.done) break; const bytes = chunkBytes(next.value); size += bytes.byteLength; if (size > input.maxBytes) throw new BodyTooLargeError(input.maxBytes); buffered = appendBytes(buffered, bytes); while (buffered.byteLength >= MULTIPART_PART_BYTES) { const partBytes = buffered.slice(0, MULTIPART_PART_BYTES); buffered = buffered.slice(MULTIPART_PART_BYTES); parts.push(await upload.uploadPart(partNumber, partBytes)); partNumber += 1; } } if (buffered.byteLength > 0 || parts.length === 0) { parts.push(await upload.uploadPart(partNumber, buffered)); } const object = await upload.complete(parts); return { size: object.size ?? size, etag: object.etag ?? null }; } catch (cause) { try { await reader.cancel(cause); } catch {} try { await upload.abort(); } catch {} try { await bucket.delete(input.key); } catch {} throw cause; } } export async function putFileBlob( bucket: R2Bucket, input: { key: string; body: ReadableStream | null; contentType: string | null; contentLength: number | null; maxBytes: number; }, ): Promise<{ size: number; etag: string | null }> { if (input.contentLength !== null) { return await putKnownLength(bucket, { ...input, contentLength: input.contentLength }); } return await putMultipart(bucket, input); } export async function getFileBlob(bucket: R2Bucket, key: string): Promise { return await bucket.get(key); } export async function deleteFileBlob(bucket: R2Bucket, key: string): Promise { await bucket.delete(key); }