File
Blob: src/worker/git/receive/r2Upload.ts
| 1 | import { bytesToHex, createDigestStream } from "@/worker/common"; |
| 2 | import { packIndexKey, packRefsKey } from "@/worker/keys"; |
| 3 | import { SubrequestLimiter } from "../operations/limits"; |
| 4 | import { appendBytes, cloneBytes } from "./bytes"; |
| 5 | |
| 6 | const MULTIPART_PART_BYTES = 8 * 1024 * 1024; |
| 7 | const PACK_HEADER_BYTES = 12; |
| 8 | const PACK_TRAILER_BYTES = 20; |
| 9 | const UPLOAD_PROGRESS_STEPS = 20; |
| 10 | |
| 11 | export type StagedPackUpload = { |
| 12 | packKey: string; |
| 13 | packBytes: number; |
| 14 | cleanup(): Promise<void>; |
| 15 | }; |
| 16 | |
| 17 | function formatProgressBytes(bytes: number): string { |
| 18 | if (bytes >= 1024 * 1024) { |
| 19 | return `${(bytes / 1024 / 1024).toFixed(1)} MiB`; |
| 20 | } |
| 21 | if (bytes >= 1024) { |
| 22 | return `${(bytes / 1024).toFixed(1)} KiB`; |
| 23 | } |
| 24 | return `${bytes} B`; |
| 25 | } |
| 26 | |
| 27 | function emitKnownLengthUploadProgress(args: { |
| 28 | onProgress?: (message: string) => void; |
| 29 | uploadedBytes: number; |
| 30 | totalBytes: number; |
| 31 | progressInterval: number; |
| 32 | lastReportedStep: number; |
| 33 | }): number { |
| 34 | if (!args.onProgress || args.totalBytes <= 0) return args.lastReportedStep; |
| 35 | |
| 36 | const percent = Math.floor((args.uploadedBytes / args.totalBytes) * 100); |
| 37 | const nextStep = Math.min( |
| 38 | UPLOAD_PROGRESS_STEPS, |
| 39 | Math.floor(args.uploadedBytes / args.progressInterval) |
| 40 | ); |
| 41 | if (nextStep <= args.lastReportedStep && args.uploadedBytes < args.totalBytes) { |
| 42 | return args.lastReportedStep; |
| 43 | } |
| 44 | |
| 45 | if (args.uploadedBytes >= args.totalBytes) { |
| 46 | args.onProgress( |
| 47 | `Uploading pack to object storage: 100% (${formatProgressBytes(args.totalBytes)}/${formatProgressBytes(args.totalBytes)}), done.\n` |
| 48 | ); |
| 49 | return UPLOAD_PROGRESS_STEPS; |
| 50 | } |
| 51 | |
| 52 | args.onProgress( |
| 53 | `Uploading pack to object storage: ${percent}% (${formatProgressBytes(args.uploadedBytes)}/${formatProgressBytes(args.totalBytes)})\r` |
| 54 | ); |
| 55 | return nextStep; |
| 56 | } |
| 57 | |
| 58 | function emitStreamingUploadProgress(args: { |
| 59 | onProgress?: (message: string) => void; |
| 60 | uploadedBytes: number; |
| 61 | reportEveryBytes: number; |
| 62 | lastReportedBytes: number; |
| 63 | }): number { |
| 64 | if (!args.onProgress) return args.lastReportedBytes; |
| 65 | if (args.uploadedBytes < args.reportEveryBytes) return args.lastReportedBytes; |
| 66 | if (args.uploadedBytes - args.lastReportedBytes < args.reportEveryBytes) { |
| 67 | return args.lastReportedBytes; |
| 68 | } |
| 69 | args.onProgress( |
| 70 | `Uploading pack to object storage: ${formatProgressBytes(args.uploadedBytes)} uploaded\r` |
| 71 | ); |
| 72 | return args.uploadedBytes; |
| 73 | } |
| 74 | |
| 75 | function parseContentLength(request: Request): number | undefined { |
| 76 | const raw = request.headers.get("Content-Length"); |
| 77 | if (!raw) return undefined; |
| 78 | const parsed = Number.parseInt(raw, 10); |
| 79 | if (!Number.isFinite(parsed) || parsed < 0) return undefined; |
| 80 | return parsed; |
| 81 | } |
| 82 | |
| 83 | function getRemainingBodyLength(request: Request, bytesConsumed: number): number | undefined { |
| 84 | const contentLength = parseContentLength(request); |
| 85 | if (contentLength === undefined || contentLength < bytesConsumed) return undefined; |
| 86 | return contentLength - bytesConsumed; |
| 87 | } |
| 88 | |
| 89 | function validatePackHeader(prefix: Uint8Array<ArrayBufferLike>): void { |
| 90 | if (prefix.byteLength < PACK_HEADER_BYTES) return; |
| 91 | if (prefix[0] !== 0x50 || prefix[1] !== 0x41 || prefix[2] !== 0x43 || prefix[3] !== 0x4b) { |
| 92 | throw new Error("receive-pack body did not begin with a valid PACK header."); |
| 93 | } |
| 94 | const version = (prefix[4] << 24) | (prefix[5] << 16) | (prefix[6] << 8) | prefix[7]; |
| 95 | if (version !== 2) { |
| 96 | throw new Error(`Unsupported pack version ${version}.`); |
| 97 | } |
| 98 | } |
| 99 | |
| 100 | async function updateTrailerWindow( |
| 101 | digestWriter: WritableStreamDefaultWriter<Uint8Array>, |
| 102 | previousTail: Uint8Array<ArrayBufferLike>, |
| 103 | nextChunk: Uint8Array<ArrayBufferLike> |
| 104 | ): Promise<Uint8Array<ArrayBufferLike>> { |
| 105 | const combined = new Uint8Array(previousTail.byteLength + nextChunk.byteLength); |
| 106 | combined.set(previousTail, 0); |
| 107 | combined.set(nextChunk, previousTail.byteLength); |
| 108 | |
| 109 | if (combined.byteLength <= PACK_TRAILER_BYTES) { |
| 110 | return combined; |
| 111 | } |
| 112 | |
| 113 | const digestBytes = combined.subarray(0, combined.byteLength - PACK_TRAILER_BYTES); |
| 114 | const nextTail = combined.subarray(combined.byteLength - PACK_TRAILER_BYTES); |
| 115 | await digestWriter.write(digestBytes); |
| 116 | return cloneBytes(nextTail); |
| 117 | } |
| 118 | |
| 119 | async function stageKnownLengthPack(args: { |
| 120 | env: Env; |
| 121 | packKey: string; |
| 122 | expectedLength: number; |
| 123 | packStream: ReadableStream<Uint8Array>; |
| 124 | limiter: SubrequestLimiter; |
| 125 | countSubrequest(op: string, n?: number): void; |
| 126 | onProgress?: (message: string) => void; |
| 127 | }): Promise<StagedPackUpload> { |
| 128 | if (args.expectedLength <= 0) { |
| 129 | throw new Error("Streaming receive expected a non-empty pack body."); |
| 130 | } |
| 131 | |
| 132 | const fixedLengthStream = new FixedLengthStream(args.expectedLength); |
| 133 | const uploadWriter = fixedLengthStream.writable.getWriter(); |
| 134 | const digestStream = createDigestStream("SHA-1"); |
| 135 | const digestWriter = digestStream.getWriter(); |
| 136 | const reader = args.packStream.getReader(); |
| 137 | |
| 138 | args.countSubrequest("r2:put-pack"); |
| 139 | const putPromise = args.limiter.run("r2:put-pack", async () => { |
| 140 | return await args.env.REPO_BUCKET.put(args.packKey, fixedLengthStream.readable); |
| 141 | }); |
| 142 | let totalBytes = 0; |
| 143 | let headerPrefix: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 144 | let tail: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 145 | const progressInterval = Math.max(1, Math.floor(args.expectedLength / UPLOAD_PROGRESS_STEPS)); |
| 146 | let lastReportedStep = 0; |
| 147 | |
| 148 | args.onProgress?.( |
| 149 | `Uploading pack to object storage: 0% (0 B/${formatProgressBytes(args.expectedLength)})\r` |
| 150 | ); |
| 151 | |
| 152 | try { |
| 153 | while (true) { |
| 154 | const next = await reader.read(); |
| 155 | if (next.done) break; |
| 156 | const chunk = cloneBytes(next.value); |
| 157 | |
| 158 | totalBytes += chunk.byteLength; |
| 159 | if (headerPrefix.byteLength < PACK_HEADER_BYTES) { |
| 160 | const headerNeeded = PACK_HEADER_BYTES - headerPrefix.byteLength; |
| 161 | headerPrefix = appendBytes(headerPrefix, chunk.subarray(0, headerNeeded)); |
| 162 | validatePackHeader(headerPrefix); |
| 163 | } |
| 164 | |
| 165 | tail = await updateTrailerWindow(digestWriter, tail, chunk); |
| 166 | await uploadWriter.write(chunk); |
| 167 | lastReportedStep = emitKnownLengthUploadProgress({ |
| 168 | onProgress: args.onProgress, |
| 169 | uploadedBytes: totalBytes, |
| 170 | totalBytes: args.expectedLength, |
| 171 | progressInterval, |
| 172 | lastReportedStep, |
| 173 | }); |
| 174 | } |
| 175 | |
| 176 | if (totalBytes !== args.expectedLength) { |
| 177 | throw new Error("Received pack length did not match Content-Length."); |
| 178 | } |
| 179 | if (totalBytes < PACK_HEADER_BYTES + PACK_TRAILER_BYTES) { |
| 180 | throw new Error("Received pack body was too short."); |
| 181 | } |
| 182 | |
| 183 | await digestWriter.close(); |
| 184 | const computedDigest = new Uint8Array(await digestStream.digest); |
| 185 | if (bytesToHex(computedDigest) !== bytesToHex(tail)) { |
| 186 | throw new Error("Received pack trailer SHA-1 did not match the streamed body."); |
| 187 | } |
| 188 | |
| 189 | await uploadWriter.close(); |
| 190 | await putPromise; |
| 191 | return { |
| 192 | packKey: args.packKey, |
| 193 | packBytes: totalBytes, |
| 194 | async cleanup() { |
| 195 | await deleteStagedPackArtifacts({ |
| 196 | env: args.env, |
| 197 | packKey: args.packKey, |
| 198 | limiter: args.limiter, |
| 199 | countSubrequest: args.countSubrequest, |
| 200 | }); |
| 201 | }, |
| 202 | }; |
| 203 | } catch (error) { |
| 204 | try { |
| 205 | await uploadWriter.abort(error); |
| 206 | } catch {} |
| 207 | try { |
| 208 | await reader.cancel(error); |
| 209 | } catch {} |
| 210 | try { |
| 211 | await deletePackArtifact({ |
| 212 | env: args.env, |
| 213 | key: args.packKey, |
| 214 | limiter: args.limiter, |
| 215 | countSubrequest: args.countSubrequest, |
| 216 | op: "r2:delete-staged-pack", |
| 217 | }); |
| 218 | } catch {} |
| 219 | throw error; |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | async function uploadMultipartPart(args: { |
| 224 | upload: R2MultipartUpload; |
| 225 | partNumber: number; |
| 226 | bytes: Uint8Array; |
| 227 | limiter: SubrequestLimiter; |
| 228 | countSubrequest(op: string, n?: number): void; |
| 229 | }): Promise<R2UploadedPart> { |
| 230 | args.countSubrequest("r2:upload-pack-part"); |
| 231 | return await args.limiter.run("r2:upload-pack-part", async () => { |
| 232 | return await args.upload.uploadPart(args.partNumber, args.bytes); |
| 233 | }); |
| 234 | } |
| 235 | |
| 236 | async function stageMultipartPack(args: { |
| 237 | env: Env; |
| 238 | packKey: string; |
| 239 | packStream: ReadableStream<Uint8Array>; |
| 240 | limiter: SubrequestLimiter; |
| 241 | countSubrequest(op: string, n?: number): void; |
| 242 | onProgress?: (message: string) => void; |
| 243 | }): Promise<StagedPackUpload> { |
| 244 | args.countSubrequest("r2:create-pack-multipart"); |
| 245 | const upload = await args.limiter.run("r2:create-pack-multipart", async () => { |
| 246 | return await args.env.REPO_BUCKET.createMultipartUpload(args.packKey); |
| 247 | }); |
| 248 | const digestStream = createDigestStream("SHA-1"); |
| 249 | const digestWriter = digestStream.getWriter(); |
| 250 | const reader = args.packStream.getReader(); |
| 251 | |
| 252 | const uploadedParts: R2UploadedPart[] = []; |
| 253 | let buffered: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 254 | let totalBytes = 0; |
| 255 | let headerPrefix: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 256 | let tail: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 257 | let partNumber = 1; |
| 258 | let lastReportedBytes = 0; |
| 259 | |
| 260 | args.onProgress?.("Uploading pack to object storage: streaming upload started\n"); |
| 261 | |
| 262 | try { |
| 263 | while (true) { |
| 264 | const next = await reader.read(); |
| 265 | if (next.done) break; |
| 266 | const chunk = cloneBytes(next.value); |
| 267 | |
| 268 | totalBytes += chunk.byteLength; |
| 269 | if (headerPrefix.byteLength < PACK_HEADER_BYTES) { |
| 270 | const headerNeeded = PACK_HEADER_BYTES - headerPrefix.byteLength; |
| 271 | headerPrefix = appendBytes(headerPrefix, chunk.subarray(0, headerNeeded)); |
| 272 | validatePackHeader(headerPrefix); |
| 273 | } |
| 274 | |
| 275 | tail = await updateTrailerWindow(digestWriter, tail, chunk); |
| 276 | buffered = appendBytes(buffered, chunk); |
| 277 | lastReportedBytes = emitStreamingUploadProgress({ |
| 278 | onProgress: args.onProgress, |
| 279 | uploadedBytes: totalBytes, |
| 280 | reportEveryBytes: MULTIPART_PART_BYTES, |
| 281 | lastReportedBytes, |
| 282 | }); |
| 283 | |
| 284 | while (buffered.byteLength >= MULTIPART_PART_BYTES) { |
| 285 | const partBytes = buffered.slice(0, MULTIPART_PART_BYTES); |
| 286 | buffered = buffered.slice(MULTIPART_PART_BYTES); |
| 287 | uploadedParts.push( |
| 288 | await uploadMultipartPart({ |
| 289 | upload, |
| 290 | partNumber, |
| 291 | bytes: partBytes, |
| 292 | limiter: args.limiter, |
| 293 | countSubrequest: args.countSubrequest, |
| 294 | }) |
| 295 | ); |
| 296 | partNumber++; |
| 297 | } |
| 298 | } |
| 299 | |
| 300 | if (totalBytes < PACK_HEADER_BYTES + PACK_TRAILER_BYTES) { |
| 301 | throw new Error("Received pack body was too short."); |
| 302 | } |
| 303 | |
| 304 | if (buffered.byteLength === 0 && uploadedParts.length === 0) { |
| 305 | throw new Error("Streaming receive expected a non-empty pack body."); |
| 306 | } |
| 307 | |
| 308 | uploadedParts.push( |
| 309 | await uploadMultipartPart({ |
| 310 | upload, |
| 311 | partNumber, |
| 312 | bytes: buffered, |
| 313 | limiter: args.limiter, |
| 314 | countSubrequest: args.countSubrequest, |
| 315 | }) |
| 316 | ); |
| 317 | |
| 318 | await digestWriter.close(); |
| 319 | const computedDigest = new Uint8Array(await digestStream.digest); |
| 320 | if (bytesToHex(computedDigest) !== bytesToHex(tail)) { |
| 321 | throw new Error("Received pack trailer SHA-1 did not match the streamed body."); |
| 322 | } |
| 323 | |
| 324 | args.countSubrequest("r2:complete-pack-multipart"); |
| 325 | await args.limiter.run("r2:complete-pack-multipart", async () => { |
| 326 | await upload.complete(uploadedParts); |
| 327 | }); |
| 328 | args.onProgress?.( |
| 329 | `Uploading pack to object storage: done (${formatProgressBytes(totalBytes)})\n` |
| 330 | ); |
| 331 | |
| 332 | return { |
| 333 | packKey: args.packKey, |
| 334 | packBytes: totalBytes, |
| 335 | async cleanup() { |
| 336 | await deleteStagedPackArtifacts({ |
| 337 | env: args.env, |
| 338 | packKey: args.packKey, |
| 339 | limiter: args.limiter, |
| 340 | countSubrequest: args.countSubrequest, |
| 341 | }); |
| 342 | }, |
| 343 | }; |
| 344 | } catch (error) { |
| 345 | try { |
| 346 | await reader.cancel(error); |
| 347 | } catch {} |
| 348 | try { |
| 349 | await upload.abort(); |
| 350 | } catch {} |
| 351 | try { |
| 352 | await deletePackArtifact({ |
| 353 | env: args.env, |
| 354 | key: args.packKey, |
| 355 | limiter: args.limiter, |
| 356 | countSubrequest: args.countSubrequest, |
| 357 | op: "r2:delete-staged-pack", |
| 358 | }); |
| 359 | } catch {} |
| 360 | throw error; |
| 361 | } |
| 362 | } |
| 363 | |
| 364 | async function deletePackArtifact(args: { |
| 365 | env: Env; |
| 366 | key: string; |
| 367 | limiter: SubrequestLimiter; |
| 368 | countSubrequest(op: string, n?: number): void; |
| 369 | op: string; |
| 370 | }): Promise<void> { |
| 371 | args.countSubrequest(args.op); |
| 372 | await args.limiter.run(args.op, async () => { |
| 373 | await args.env.REPO_BUCKET.delete(args.key); |
| 374 | }); |
| 375 | } |
| 376 | |
| 377 | async function deleteStagedPackArtifacts(args: { |
| 378 | env: Env; |
| 379 | packKey: string; |
| 380 | limiter: SubrequestLimiter; |
| 381 | countSubrequest(op: string, n?: number): void; |
| 382 | }): Promise<void> { |
| 383 | // Cleanup is intentionally split per artifact so logs and request accounting |
| 384 | // can identify which platform call hit the budget or failed. |
| 385 | const failures: string[] = []; |
| 386 | const artifacts = [ |
| 387 | { key: args.packKey, op: "r2:delete-staged-pack" }, |
| 388 | { key: packIndexKey(args.packKey), op: "r2:delete-pack-idx" }, |
| 389 | { key: packRefsKey(args.packKey), op: "r2:delete-pack-refs" }, |
| 390 | ]; |
| 391 | |
| 392 | for (const artifact of artifacts) { |
| 393 | try { |
| 394 | await deletePackArtifact({ |
| 395 | env: args.env, |
| 396 | key: artifact.key, |
| 397 | limiter: args.limiter, |
| 398 | countSubrequest: args.countSubrequest, |
| 399 | op: artifact.op, |
| 400 | }); |
| 401 | } catch (error) { |
| 402 | failures.push(`${artifact.op}:${String(error)}`); |
| 403 | } |
| 404 | } |
| 405 | |
| 406 | if (failures.length > 0) { |
| 407 | throw new Error(`staged pack cleanup failed for ${failures.join(", ")}`); |
| 408 | } |
| 409 | } |
| 410 | |
| 411 | export async function stagePackToR2(args: { |
| 412 | env: Env; |
| 413 | request: Request; |
| 414 | packStream: ReadableStream<Uint8Array>; |
| 415 | packKey: string; |
| 416 | bytesConsumed: number; |
| 417 | limiter: SubrequestLimiter; |
| 418 | countSubrequest(op: string, n?: number): void; |
| 419 | onProgress?: (message: string) => void; |
| 420 | }): Promise<StagedPackUpload> { |
| 421 | const remainingLength = getRemainingBodyLength(args.request, args.bytesConsumed); |
| 422 | if (remainingLength !== undefined) { |
| 423 | return await stageKnownLengthPack({ |
| 424 | env: args.env, |
| 425 | packKey: args.packKey, |
| 426 | expectedLength: remainingLength, |
| 427 | packStream: args.packStream, |
| 428 | limiter: args.limiter, |
| 429 | countSubrequest: args.countSubrequest, |
| 430 | onProgress: args.onProgress, |
| 431 | }); |
| 432 | } |
| 433 | |
| 434 | return await stageMultipartPack({ |
| 435 | env: args.env, |
| 436 | packKey: args.packKey, |
| 437 | packStream: args.packStream, |
| 438 | limiter: args.limiter, |
| 439 | countSubrequest: args.countSubrequest, |
| 440 | onProgress: args.onProgress, |
| 441 | }); |
| 442 | } |
| 443 | |
| 444 | export async function deleteStagedPack(upload: StagedPackUpload | undefined): Promise<void> { |
| 445 | if (!upload) return; |
| 446 | await upload.cleanup(); |
| 447 | } |