Skip to content
File

Blob: src/worker/git/receive/r2Upload.ts

typescript448 lines
1import { bytesToHex, createDigestStream } from "@/worker/common";
2import { packIndexKey, packRefsKey } from "@/worker/keys";
3import { SubrequestLimiter } from "../operations/limits";
4import { appendBytes, cloneBytes } from "./bytes";
5 
6const MULTIPART_PART_BYTES = 8 * 1024 * 1024;
7const PACK_HEADER_BYTES = 12;
8const PACK_TRAILER_BYTES = 20;
9const UPLOAD_PROGRESS_STEPS = 20;
10 
11export type StagedPackUpload = {
12 packKey: string;
13 packBytes: number;
14 cleanup(): Promise<void>;
15};
16 
17function 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 
27function 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 
58function 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 
75function 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 
83function 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 
89function 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 
100async 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 
119async 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 
223async 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 
236async 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 
364async 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 
377async 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 
411export 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 
444export async function deleteStagedPack(upload: StagedPackUpload | undefined): Promise<void> {
445 if (!upload) return;
446 await upload.cleanup();
447}