Skip to content
File

Blob: test/fetch-streaming.worker.test.ts

typescript1146 lines
1import { it, expect, describe } from "vitest";
2import { env, exports as workerExports } from "cloudflare:workers";
3import { pktLine, delimPkt, flushPkt, concatChunks, decodePktLines } from "@/worker/git";
4import { handleFetchV2Streaming } from "@/worker/git/operations/uploadStream";
5import {
6 buildServeUploadPackPlan,
7 loadUploadPackSnapshot,
8 planUploadPack,
9} from "@/worker/git/operations/fetch/plan";
10import { uniqueRepoId, runDOWithRetry } from "./util/test-helpers";
11import { setupRepoForTests } from "./util/repoSeed";
12import { asBufferSource } from "@/worker/common";
13import { packRefsKey } from "@/worker/keys";
14import { runQueueMessage } from "./util/queue";
15import { createTestCacheContext, seedPackFirstRepo } from "./util/pack-first";
16import { makeTracingLimiter } from "./util/pack-indexer.helpers";
17import { buildAppendOnlyDelta, buildCopyPrefixDelta, buildPack } from "./util/git-pack";
18import { seedPackedRepoState } from "./util/packed-repo";
19import { computeOid, encodeGitObject } from "@/worker/git/core/objects";
20 
21function buildFetchBody({
22 wants,
23 haves,
24 done,
25}: {
26 wants: string[];
27 haves?: string[];
28 done?: boolean;
29}) {
30 const chunks: Uint8Array[] = [];
31 chunks.push(pktLine("command=fetch\n"));
32 chunks.push(delimPkt());
33 for (const w of wants) chunks.push(pktLine(`want ${w}\n`));
34 for (const h of haves || []) chunks.push(pktLine(`have ${h}\n`));
35 if (done) chunks.push(pktLine("done\n"));
36 chunks.push(flushPkt());
37 return concatChunks(chunks);
38}
39 
40/**
41 * Find the index of a byte sequence within a Uint8Array
42 */
43function findBytes(haystack: Uint8Array, needle: Uint8Array): number {
44 outer: for (let i = 0; i <= haystack.length - needle.length; i++) {
45 for (let j = 0; j < needle.length; j++) {
46 if (haystack[i + j] !== needle[j]) continue outer;
47 }
48 return i;
49 }
50 return -1;
51}
52 
53async function seedTwoCommitRepo(owner: string, repo: string) {
54 const repoId = `${owner}/${repo}`;
55 const doId = env.REPO_DO.idFromName(repoId);
56 const seeded = await seedPackFirstRepo(repoId);
57 
58 return {
59 repoId,
60 doId,
61 getStub: seeded.getStub,
62 firstCommit: seeded.baseCommit.oid,
63 secondCommit: seeded.nextCommit.oid,
64 };
65}
66 
67async function postFinalFetch(args: {
68 owner: string;
69 repo: string;
70 wants: string[];
71 haves: string[];
72}): Promise<Response> {
73 const body = buildFetchBody({
74 wants: args.wants,
75 haves: args.haves,
76 done: true,
77 });
78 return await workerExports.default.fetch(
79 `https://example.com/${args.owner}/${args.repo}/git-upload-pack`,
80 {
81 method: "POST",
82 headers: {
83 "Content-Type": "application/x-git-upload-pack-request",
84 "Git-Protocol": "version=2",
85 },
86 body: asBufferSource(body),
87 }
88 );
89}
90 
91async function expectRetryThenBackfillRepair(args: {
92 owner: string;
93 repo: string;
94 repoId: string;
95 doId: DurableObjectId;
96 packKey: string;
97 firstCommit: string;
98 secondCommit: string;
99}): Promise<void> {
100 const retryRes = await postFinalFetch({
101 owner: args.owner,
102 repo: args.repo,
103 wants: [args.secondCommit],
104 haves: [args.firstCommit],
105 });
106 expect(retryRes.status).toBe(503);
107 expect(retryRes.headers.get("Retry-After")).toBe("10");
108 expect(await retryRes.text()).not.toContain("packfile");
109 
110 const queueResult = await runQueueMessage({
111 kind: "pack-ref-backfill",
112 doId: args.doId.toString(),
113 repoId: args.repoId,
114 packKey: args.packKey,
115 });
116 expect(queueResult).toEqual({ acked: true, retried: false });
117 await expect(env.REPO_BUCKET.head(packRefsKey(args.packKey))).resolves.toBeTruthy();
118 
119 const okRes = await postFinalFetch({
120 owner: args.owner,
121 repo: args.repo,
122 wants: [args.secondCommit],
123 haves: [args.firstCommit],
124 });
125 expect(okRes.status).toBe(200);
126 const okBytes = new Uint8Array(await okRes.arrayBuffer());
127 expect(new TextDecoder().decode(okBytes.subarray(4, 13))).toBe("packfile\n");
128}
129 
130describe("git fetch streaming (default)", () => {
131 it("handles fetch with streaming by default", async () => {
132 const owner = "o";
133 const repo = uniqueRepoId("streaming");
134 await setupRepoForTests(env, owner, repo);
135 const repoId = `${owner}/${repo}`;
136 
137 // Seed a repository with some commits
138 const id = env.REPO_DO.idFromName(repoId);
139 const { commitOid } = await runDOWithRetry(
140 () => env.REPO_DO.get(id),
141 async (instance) => await instance.seedMinimalRepo()
142 );
143 
144 const body = buildFetchBody({ wants: [commitOid], done: true });
145 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
146 
147 const res = await workerExports.default.fetch(url, {
148 method: "POST",
149 headers: {
150 "Content-Type": "application/x-git-upload-pack-request",
151 "Git-Protocol": "version=2",
152 },
153 body,
154 } as any);
155 
156 expect(res.status).toBe(200);
157 expect(res.headers.get("Content-Type")).toContain("git-upload-pack-result");
158 
159 const bytes = new Uint8Array(await res.arrayBuffer());
160 
161 const lines = decodePktLines(bytes);
162 let hasAcknowledgments = false;
163 let hasPackfile = false;
164 let inPackfile = false;
165 const packData: Uint8Array[] = [];
166 let hasSideband = false;
167 
168 for (const line of lines) {
169 if (line.type === "line" && line.text === "acknowledgments\n") {
170 hasAcknowledgments = true;
171 }
172 if (line.type === "line" && line.text === "packfile\n") {
173 hasPackfile = true;
174 inPackfile = true;
175 }
176 }
177 
178 expect(
179 hasAcknowledgments,
180 "Response should NOT contain 'acknowledgments\\n' when done=true"
181 ).toBe(false);
182 expect(hasPackfile, "Response should contain 'packfile\\n' pkt-line").toBe(true);
183 
184 for (const line of lines) {
185 if (line.type === "line" && line.text === "packfile\n") {
186 inPackfile = true;
187 } else if (inPackfile && line.type === "line" && line.raw) {
188 // Check if this is sideband data (first byte is 0x01, 0x02, or 0x03)
189 if (
190 line.raw.length > 0 &&
191 (line.raw[0] === 0x01 || line.raw[0] === 0x02 || line.raw[0] === 0x03)
192 ) {
193 hasSideband = true;
194 if (line.raw[0] === 0x01) {
195 // Channel 1: pack data
196 packData.push(line.raw.subarray(1));
197 }
198 }
199 }
200 }
201 
202 expect(hasSideband).toBe(true);
203 const pack = concatChunks(packData);
204 
205 expect(pack.length).toBeGreaterThan(0);
206 
207 // Verify pack signature
208 const packSig = new TextDecoder().decode(pack.subarray(0, 4));
209 expect(packSig).toBe("PACK");
210 
211 // Verify pack header
212 const dv = new DataView(pack.buffer, pack.byteOffset, pack.byteLength);
213 const version = dv.getUint32(4);
214 const objCount = dv.getUint32(8);
215 expect(version).toBe(2);
216 expect(objCount).toBeGreaterThan(0);
217 
218 // Verify SHA-1 trailer (last 20 bytes)
219 expect(pack.length).toBeGreaterThanOrEqual(32); // At least header + SHA-1
220 const packBody = pack.subarray(0, pack.length - 20);
221 const expectedSha = pack.subarray(pack.length - 20);
222 const actualSha = new Uint8Array(await crypto.subtle.digest("SHA-1", asBufferSource(packBody)));
223 expect(Array.from(actualSha)).toEqual(Array.from(expectedSha));
224 });
225 
226 it("handles incremental fetch with haves", async () => {
227 const owner = "o";
228 const repo = uniqueRepoId("incremental");
229 await setupRepoForTests(env, owner, repo);
230 const repoId = `${owner}/${repo}`;
231 
232 // Seed repository and get multiple commits
233 const id = env.REPO_DO.idFromName(repoId);
234 const { commitOid, parentOid } = await runDOWithRetry(
235 () => env.REPO_DO.get(id),
236 async (instance) => {
237 const firstResult = await instance.seedMinimalRepo();
238 const secondResult = await instance.seedMinimalRepo();
239 return {
240 commitOid: secondResult.commitOid,
241 parentOid: firstResult.commitOid,
242 };
243 }
244 );
245 
246 // First, do negotiation without done
247 const negotiateBody = buildFetchBody({
248 wants: [commitOid],
249 haves: [parentOid],
250 done: false,
251 });
252 
253 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
254 const negotiateRes = await workerExports.default.fetch(url, {
255 method: "POST",
256 headers: {
257 "Content-Type": "application/x-git-upload-pack-request",
258 "Git-Protocol": "version=2",
259 // Streaming is now default, no header needed
260 },
261 body: negotiateBody,
262 } as any);
263 
264 expect(negotiateRes.status).toBe(200);
265 const negotiateBytes = new Uint8Array(await negotiateRes.arrayBuffer());
266 const negotiateLines = decodePktLines(negotiateBytes);
267 let hasAcknowledgments = false;
268 let hasPackfile = false;
269 let hasParentAck = false;
270 
271 for (const line of negotiateLines) {
272 if (line.type === "line") {
273 if (line.text === "acknowledgments\n") hasAcknowledgments = true;
274 if (line.text === "packfile\n") hasPackfile = true;
275 if (line.text && line.text.includes(`ACK ${parentOid}`)) hasParentAck = true;
276 }
277 }
278 
279 // Should only have acknowledgments, no packfile
280 expect(
281 hasAcknowledgments,
282 "Negotiation response should contain 'acknowledgments\\n' pkt-line"
283 ).toBe(true);
284 expect(hasPackfile, "Negotiation response should NOT contain 'packfile\\n' pkt-line").toBe(
285 false
286 );
287 expect(hasParentAck, `Negotiation should ACK parent ${parentOid}`).toBe(true);
288 
289 // Now fetch with done
290 const fetchBody = buildFetchBody({
291 wants: [commitOid],
292 haves: [parentOid],
293 done: true,
294 });
295 
296 const fetchRes = await workerExports.default.fetch(url, {
297 method: "POST",
298 headers: {
299 "Content-Type": "application/x-git-upload-pack-request",
300 "Git-Protocol": "version=2",
301 // Streaming is now default, no header needed
302 },
303 body: fetchBody,
304 } as any);
305 
306 expect(fetchRes.status).toBe(200);
307 const fetchBytes = new Uint8Array(await fetchRes.arrayBuffer());
308 const fetchLines = decodePktLines(fetchBytes);
309 let hasFetchAcknowledgments = false;
310 let hasFetchPackfile = false;
311 
312 for (const line of fetchLines) {
313 if (line.type === "line") {
314 if (line.text === "acknowledgments\n") hasFetchAcknowledgments = true;
315 if (line.text === "packfile\n") hasFetchPackfile = true;
316 }
317 }
318 
319 // Should go straight to packfile when done=true
320 expect(
321 hasFetchAcknowledgments,
322 "Final fetch with done=true should NOT contain 'acknowledgments\\n'"
323 ).toBe(false);
324 expect(hasFetchPackfile, "Final fetch should contain 'packfile\\n' pkt-line").toBe(true);
325 
326 // Parse and verify pack data
327 let inPackfile = false;
328 const packData: Uint8Array[] = [];
329 
330 for (const line of fetchLines) {
331 if (line.type === "line" && line.text === "packfile\n") {
332 inPackfile = true;
333 } else if (inPackfile && line.type === "line" && line.raw?.[0] === 0x01) {
334 packData.push(line.raw.subarray(1));
335 }
336 }
337 
338 const pack = concatChunks(packData);
339 expect(pack.length).toBeGreaterThan(0);
340 
341 // Verify it's a valid pack
342 const packSig = new TextDecoder().decode(pack.subarray(0, 4));
343 expect(packSig).toBe("PACK");
344 });
345 
346 it("returns retry before packfile when an active pack ref sidecar is missing and backfill repairs it", async () => {
347 const owner = "o";
348 const repo = uniqueRepoId("missing-ref-sidecar");
349 await setupRepoForTests(env, owner, repo);
350 const repoId = `${owner}/${repo}`;
351 const id = env.REPO_DO.idFromName(repoId);
352 const getStub = () => env.REPO_DO.get(id);
353 const { commitOid: firstCommit } = await runDOWithRetry(
354 getStub,
355 async (instance) => await instance.seedMinimalRepo()
356 );
357 const { commitOid: secondCommit } = await runDOWithRetry(
358 getStub,
359 async (instance) => await instance.seedMinimalRepo()
360 );
361 
362 const activeCatalog = await getStub().getActivePackCatalog();
363 expect(activeCatalog.length).toBeGreaterThan(0);
364 const targetPack = activeCatalog[0]!;
365 await env.REPO_BUCKET.delete(packRefsKey(targetPack.packKey));
366 
367 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
368 const fetchBody = buildFetchBody({
369 wants: [secondCommit],
370 haves: [firstCommit],
371 done: true,
372 });
373 
374 const retryRes = await workerExports.default.fetch(url, {
375 method: "POST",
376 headers: {
377 "Content-Type": "application/x-git-upload-pack-request",
378 "Git-Protocol": "version=2",
379 },
380 body: fetchBody,
381 } as any);
382 expect(retryRes.status).toBe(503);
383 expect(retryRes.headers.get("Retry-After")).toBe("10");
384 expect(await retryRes.text()).not.toContain("packfile");
385 
386 const queueResult = await runQueueMessage({
387 kind: "pack-ref-backfill",
388 doId: id.toString(),
389 repoId,
390 packKey: targetPack.packKey,
391 });
392 expect(queueResult).toEqual({ acked: true, retried: false });
393 await expect(env.REPO_BUCKET.head(packRefsKey(targetPack.packKey))).resolves.toBeTruthy();
394 
395 const okRes = await workerExports.default.fetch(url, {
396 method: "POST",
397 headers: {
398 "Content-Type": "application/x-git-upload-pack-request",
399 "Git-Protocol": "version=2",
400 },
401 body: fetchBody,
402 } as any);
403 expect(okRes.status).toBe(200);
404 const okBytes = new Uint8Array(await okRes.arrayBuffer());
405 expect(new TextDecoder().decode(okBytes.subarray(4, 13))).toBe("packfile\n");
406 });
407 
408 it("backfills refs for an already-active same-pack REF_DELTA pack", async () => {
409 const owner = "o";
410 const repo = uniqueRepoId("missing-ref-sidecar-ref-delta");
411 await setupRepoForTests(env, owner, repo);
412 const repoId = `${owner}/${repo}`;
413 const id = env.REPO_DO.idFromName(repoId);
414 const getStub = () => env.REPO_DO.get(id);
415 
416 const baseBlobPayload = new TextEncoder().encode("base\n");
417 const midSuffix = new TextEncoder().encode("mid\n");
418 const finalSuffix = new TextEncoder().encode("final\n");
419 const baseBlob = await encodeGitObject("blob", baseBlobPayload);
420 
421 const midPayload = new Uint8Array(baseBlobPayload.length + midSuffix.length);
422 midPayload.set(baseBlobPayload, 0);
423 midPayload.set(midSuffix, baseBlobPayload.length);
424 const midOid = await computeOid("blob", midPayload);
425 
426 const finalPayload = new Uint8Array(midPayload.length + finalSuffix.length);
427 finalPayload.set(midPayload, 0);
428 finalPayload.set(finalSuffix, midPayload.length);
429 
430 const packBytes = await buildPack([
431 {
432 type: "ref-delta",
433 baseOid: midOid,
434 delta: buildAppendOnlyDelta(midPayload, finalSuffix),
435 },
436 {
437 type: "ref-delta",
438 baseOid: baseBlob.oid,
439 delta: buildAppendOnlyDelta(baseBlobPayload, midSuffix),
440 },
441 { type: "blob", payload: baseBlobPayload },
442 ]);
443 
444 await seedPackedRepoState({
445 env,
446 repoId,
447 getStub,
448 packs: [{ name: "pack-ref-delta-chain.pack", packBytes }],
449 });
450 
451 const activeCatalog = await getStub().getActivePackCatalog();
452 expect(activeCatalog).toHaveLength(1);
453 const targetPack = activeCatalog[0]!;
454 await env.REPO_BUCKET.delete(packRefsKey(targetPack.packKey));
455 
456 const queueResult = await runQueueMessage({
457 kind: "pack-ref-backfill",
458 doId: id.toString(),
459 repoId,
460 packKey: targetPack.packKey,
461 });
462 
463 expect(queueResult).toEqual({ acked: true, retried: false });
464 await expect(env.REPO_BUCKET.head(packRefsKey(targetPack.packKey))).resolves.toBeTruthy();
465 });
466 
467 it("backfills refs when the newest external duplicate base points back to the target pack", async () => {
468 const owner = "o";
469 const repo = uniqueRepoId("missing-ref-sidecar-external-duplicate-cycle");
470 await setupRepoForTests(env, owner, repo);
471 const repoId = `${owner}/${repo}`;
472 const id = env.REPO_DO.idFromName(repoId);
473 const getStub = () => env.REPO_DO.get(id);
474 
475 const targetPayload = new TextEncoder().encode("target prefix\n");
476 const duplicateSuffix = new TextEncoder().encode("duplicate suffix\n");
477 const duplicatePayload = new Uint8Array(targetPayload.length + duplicateSuffix.length);
478 duplicatePayload.set(targetPayload, 0);
479 duplicatePayload.set(duplicateSuffix, targetPayload.length);
480 
481 const targetOid = await computeOid("blob", targetPayload);
482 const duplicateOid = await computeOid("blob", duplicatePayload);
483 
484 const olderPack = await buildPack([{ type: "blob", payload: duplicatePayload }]);
485 const targetPack = await buildPack([
486 {
487 type: "ref-delta",
488 baseOid: duplicateOid,
489 delta: buildCopyPrefixDelta(duplicatePayload, targetPayload.length),
490 },
491 ]);
492 const newerPack = await buildPack([
493 {
494 type: "ref-delta",
495 baseOid: targetOid,
496 delta: buildAppendOnlyDelta(targetPayload, duplicateSuffix),
497 },
498 ]);
499 
500 const seeded = await seedPackedRepoState({
501 env,
502 repoId,
503 getStub,
504 // The seeder indexes the reversed list, so this caller order creates
505 // older -> target -> newer catalog history while active reads remain
506 // newest-first.
507 packs: [
508 { name: "pack-newer-duplicate.pack", packBytes: newerPack },
509 { name: "pack-target-cycle.pack", packBytes: targetPack },
510 { name: "pack-older-base.pack", packBytes: olderPack },
511 ],
512 });
513 
514 const targetPackKey = seeded.packKeys[1]!;
515 const activeCatalog = await getStub().getActivePackCatalog();
516 expect(activeCatalog.map((row) => row.packKey)).toContain(targetPackKey);
517 await env.REPO_BUCKET.delete(packRefsKey(targetPackKey));
518 
519 const queueResult = await runQueueMessage({
520 kind: "pack-ref-backfill",
521 doId: id.toString(),
522 repoId,
523 packKey: targetPackKey,
524 });
525 
526 expect(queueResult).toEqual({ acked: true, retried: false });
527 await expect(env.REPO_BUCKET.head(packRefsKey(targetPackKey))).resolves.toBeTruthy();
528 });
529 
530 it("plans final fetch from valid sidecars without closure-time range reads", async () => {
531 const owner = "o";
532 const repo = uniqueRepoId("valid-ref-sidecar-no-range");
533 await setupRepoForTests(env, owner, repo);
534 const { repoId, firstCommit, secondCommit } = await seedTwoCommitRepo(owner, repo);
535 const cacheCtx = createTestCacheContext(`https://example.com/${repoId}/git-upload-pack`);
536 const labels: string[] = [];
537 cacheCtx.memo = {
538 ...(cacheCtx.memo || {}),
539 limiter: makeTracingLimiter(labels),
540 };
541 
542 const snapshotLoad = await loadUploadPackSnapshot(env, repoId, cacheCtx);
543 expect(snapshotLoad.type).toBe("Ready");
544 if (snapshotLoad.type !== "Ready") return;
545 
546 labels.length = 0;
547 const plan = await buildServeUploadPackPlan(
548 env,
549 repoId,
550 snapshotLoad.snapshot,
551 [secondCommit],
552 [firstCommit],
553 undefined,
554 cacheCtx
555 );
556 
557 expect(plan.type).toBe("Serve");
558 expect(plan.neededOids.length).toBeGreaterThan(0);
559 expect(labels).toContain("r2:get-pack-refs");
560 expect(labels).not.toContain("r2:get-range");
561 });
562 
563 it("returns retry before packfile when an active pack ref sidecar is corrupt", async () => {
564 const owner = "o";
565 const repo = uniqueRepoId("corrupt-ref-sidecar");
566 await setupRepoForTests(env, owner, repo);
567 const { repoId, doId, getStub, firstCommit, secondCommit } = await seedTwoCommitRepo(
568 owner,
569 repo
570 );
571 const activeCatalog = await getStub().getActivePackCatalog();
572 expect(activeCatalog.length).toBeGreaterThan(0);
573 const targetPack = activeCatalog[0]!;
574 const refsKey = packRefsKey(targetPack.packKey);
575 const refsObject = await env.REPO_BUCKET.get(refsKey);
576 if (!refsObject) throw new Error("missing ref sidecar");
577 const refsBytes = new Uint8Array(await refsObject.arrayBuffer());
578 refsBytes[0] ^= 0xff;
579 await env.REPO_BUCKET.put(refsKey, asBufferSource(refsBytes));
580 
581 await expectRetryThenBackfillRepair({
582 owner,
583 repo,
584 repoId,
585 doId,
586 packKey: targetPack.packKey,
587 firstCommit,
588 secondCommit,
589 });
590 });
591 
592 it("returns retry before packfile when an active pack ref sidecar is stale", async () => {
593 const owner = "o";
594 const repo = uniqueRepoId("stale-ref-sidecar");
595 await setupRepoForTests(env, owner, repo);
596 const { repoId, doId, getStub, firstCommit, secondCommit } = await seedTwoCommitRepo(
597 owner,
598 repo
599 );
600 const activeCatalog = await getStub().getActivePackCatalog();
601 expect(activeCatalog.length).toBeGreaterThan(0);
602 const targetPack = activeCatalog[0]!;
603 const refsKey = packRefsKey(targetPack.packKey);
604 const refsObject = await env.REPO_BUCKET.get(refsKey);
605 if (!refsObject) throw new Error("missing ref sidecar");
606 const refsBytes = new Uint8Array(await refsObject.arrayBuffer());
607 refsBytes[40] ^= 0xff;
608 await env.REPO_BUCKET.put(refsKey, asBufferSource(refsBytes));
609 
610 await expectRetryThenBackfillRepair({
611 owner,
612 repo,
613 repoId,
614 doId,
615 packKey: targetPack.packKey,
616 firstCommit,
617 secondCommit,
618 });
619 });
620 
621 it("does not require pack ref sidecars for negotiation-only upload-pack planning", async () => {
622 const owner = "o";
623 const repo = uniqueRepoId("negotiation-ref-sidecar");
624 await setupRepoForTests(env, owner, repo);
625 const repoId = `${owner}/${repo}`;
626 const id = env.REPO_DO.idFromName(repoId);
627 const getStub = () => env.REPO_DO.get(id);
628 const { commitOid: firstCommit } = await runDOWithRetry(
629 getStub,
630 async (instance) => await instance.seedMinimalRepo()
631 );
632 const { commitOid: secondCommit } = await runDOWithRetry(
633 getStub,
634 async (instance) => await instance.seedMinimalRepo()
635 );
636 
637 const activeCatalog = await getStub().getActivePackCatalog();
638 expect(activeCatalog.length).toBeGreaterThan(0);
639 await env.REPO_BUCKET.delete(packRefsKey(activeCatalog[0]!.packKey));
640 
641 const plan = await planUploadPack(env, repoId, [secondCommit], [firstCommit], false);
642 
643 expect(plan.type).toBe("Serve");
644 if (plan.type !== "Serve") return;
645 expect(plan.neededOids).toEqual([]);
646 expect(plan.ackOids).toContain(firstCommit);
647 });
648 
649 it("handles initial clone (no haves) with streaming", async () => {
650 const owner = "o";
651 const repo = uniqueRepoId("clone");
652 await setupRepoForTests(env, owner, repo);
653 const repoId = `${owner}/${repo}`;
654 
655 // Seed a repository
656 const id = env.REPO_DO.idFromName(repoId);
657 const { commitOid } = await runDOWithRetry(
658 () => env.REPO_DO.get(id),
659 async (instance) => instance.seedMinimalRepo()
660 );
661 
662 // Clone with no haves
663 const body = buildFetchBody({ wants: [commitOid], done: true });
664 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
665 
666 const res = await workerExports.default.fetch(url, {
667 method: "POST",
668 headers: {
669 "Content-Type": "application/x-git-upload-pack-request",
670 "Git-Protocol": "version=2",
671 },
672 body,
673 } as any);
674 
675 expect(res.status).toBe(200);
676 
677 const bytes = new Uint8Array(await res.arrayBuffer());
678 const lines = decodePktLines(bytes);
679 
680 // Should include NAK since there are no common haves
681 let hasNak = false;
682 let hasPackfile = false;
683 let progressMessages = 0;
684 const packData: Uint8Array[] = [];
685 
686 for (const line of lines) {
687 if (line.type === "line") {
688 if (line.text === "NAK\n") hasNak = true;
689 if (line.text === "packfile\n") hasPackfile = true;
690 if (hasPackfile && line.raw?.[0] === 0x01) {
691 packData.push(line.raw.subarray(1));
692 } else if (hasPackfile && line.raw?.[0] === 0x02) {
693 progressMessages++;
694 }
695 }
696 }
697 
698 // With done=true, there are no acknowledgments (no NAK)
699 expect(hasNak, "Clone response with done=true should NOT contain NAK").toBe(false);
700 expect(hasPackfile, "Clone response should contain packfile").toBe(true);
701 
702 // Verify pack contains all objects (tree + commit at minimum)
703 const pack = concatChunks(packData);
704 const dv = new DataView(pack.buffer, pack.byteOffset, pack.byteLength);
705 const objCount = dv.getUint32(8);
706 expect(objCount).toBeGreaterThanOrEqual(2); // At least tree + commit
707 });
708 
709 it("handles repositories with packs created by default", async () => {
710 const owner = "o";
711 const repo = uniqueRepoId("with-pack");
712 await setupRepoForTests(env, owner, repo);
713 const repoId = `${owner}/${repo}`;
714 
715 // Seed repository with packed objects (default behavior)
716 const id = env.REPO_DO.idFromName(repoId);
717 const { commitOid } = await runDOWithRetry(
718 () => env.REPO_DO.get(id),
719 async (instance) => instance.seedMinimalRepo() // Default: withPack=true
720 );
721 
722 const body = buildFetchBody({ wants: [commitOid], done: true });
723 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
724 
725 // Streaming is now the default
726 const res = await workerExports.default.fetch(url, {
727 method: "POST",
728 headers: {
729 "Content-Type": "application/x-git-upload-pack-request",
730 "Git-Protocol": "version=2",
731 },
732 body,
733 } as any);
734 
735 // Should succeed with streaming
736 expect(res.status).toBe(200);
737 expect(res.headers.get("Content-Type")).toContain("git-upload-pack-result");
738 
739 const bytes = new Uint8Array(await res.arrayBuffer());
740 const lines = decodePktLines(bytes);
741 let hasAcknowledgments = false;
742 let hasPackfile = false;
743 
744 for (const line of lines) {
745 if (line.type === "line") {
746 if (line.text === "acknowledgments\n") hasAcknowledgments = true;
747 if (line.text === "packfile\n") hasPackfile = true;
748 }
749 }
750 
751 // Verify basic structure
752 // When done=true, response goes straight to packfile
753 expect(
754 hasAcknowledgments,
755 "Response should NOT contain 'acknowledgments\\n' when done=true"
756 ).toBe(false);
757 expect(hasPackfile, "Response should contain 'packfile\\n' pkt-line").toBe(true);
758 
759 // Find and verify pack data
760 const packStart = findBytes(bytes, new TextEncoder().encode("PACK"));
761 expect(packStart).toBeGreaterThan(-1);
762 
763 const pack = bytes.subarray(packStart);
764 const dv = new DataView(pack.buffer, pack.byteOffset, pack.byteLength);
765 expect(dv.getUint32(4)).toBe(2); // version
766 expect(dv.getUint32(8)).toBeGreaterThan(0); // object count
767 });
768 
769 it("returns 503 when pack assembly fails", async () => {
770 const owner = "o";
771 const repo = uniqueRepoId("fail");
772 await setupRepoForTests(env, owner, repo);
773 
774 // Request fetch for non-existent objects
775 const body = buildFetchBody({
776 wants: ["deadbeefdeadbeefdeadbeefdeadbeefdeadbeef"],
777 done: true,
778 });
779 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
780 
781 const res = await workerExports.default.fetch(url, {
782 method: "POST",
783 headers: {
784 "Content-Type": "application/x-git-upload-pack-request",
785 "Git-Protocol": "version=2",
786 },
787 body,
788 } as any);
789 
790 // Should return 503 since objects don't exist
791 expect(res.status).toBe(503);
792 expect(res.headers.get("Retry-After")).toBeDefined();
793 });
794 
795 it("handles request abort mid-stream gracefully", async () => {
796 const owner = "o";
797 const repo = uniqueRepoId("abort");
798 await setupRepoForTests(env, owner, repo);
799 const repoId = `${owner}/${repo}`;
800 
801 // Seed repository with multiple objects to ensure streaming takes some time
802 const id = env.REPO_DO.idFromName(repoId);
803 const commits: string[] = [];
804 
805 await runDOWithRetry(
806 () => env.REPO_DO.get(id),
807 async (instance) => {
808 // Create multiple commits to ensure pack has content
809 for (let i = 0; i < 5; i++) {
810 const result = await instance.seedMinimalRepo();
811 commits.push(result.commitOid);
812 }
813 return commits;
814 }
815 );
816 
817 const body = buildFetchBody({ wants: commits, done: true });
818 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
819 const abortController = new AbortController();
820 
821 const fetchPromise = workerExports.default.fetch(url, {
822 method: "POST",
823 headers: {
824 "Content-Type": "application/x-git-upload-pack-request",
825 "Git-Protocol": "version=2",
826 },
827 body,
828 signal: abortController.signal,
829 } as any);
830 
831 // Abort after a short delay to interrupt the stream
832 setTimeout(() => abortController.abort(), 10);
833 
834 // The fetch should be aborted
835 try {
836 const res = await fetchPromise;
837 // If we get a response, check if it's the expected abort response
838 // Some implementations might return 499 Client Closed Request
839 if (res.status === 499) {
840 expect(res.status).toBe(499);
841 } else {
842 // Otherwise the stream might have completed before abort
843 expect(res.status).toBe(200);
844 }
845 } catch (e: any) {
846 // AbortError is expected
847 expect(e.name).toBe("AbortError");
848 }
849 });
850 
851 it("emits band-3 fatal message on mid-stream error", async () => {
852 const owner = "o";
853 const repo = uniqueRepoId("fatal");
854 await setupRepoForTests(env, owner, repo);
855 const repoId = `${owner}/${repo}`;
856 
857 // This test is tricky to implement without mocking R2 failures
858 // We'll create a scenario where pack assembly could fail mid-stream
859 // by requesting objects that exist in DO but might fail during assembly
860 
861 // Seed repository
862 const id = env.REPO_DO.idFromName(repoId);
863 const { commitOid } = await runDOWithRetry(
864 () => env.REPO_DO.get(id),
865 async (instance) => instance.seedMinimalRepo()
866 );
867 
868 // To truly test band-3 fatal, we'd need to inject an R2 failure
869 // Since we can't easily mock R2 in the test environment,
870 // we'll test that the protocol structure is correct for error cases
871 
872 // Create a malformed request that might trigger an error during processing
873 const body = buildFetchBody({
874 wants: [commitOid],
875 done: true,
876 });
877 
878 // Corrupt the body slightly to potentially trigger an error
879 const corruptedBody = new Uint8Array(body.length + 10);
880 corruptedBody.set(body, 0);
881 // Add some garbage that might confuse the parser after valid data
882 corruptedBody.set(new Uint8Array([0xff, 0xff, 0xff, 0xff]), body.length);
883 
884 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
885 
886 const res = await workerExports.default.fetch(url, {
887 method: "POST",
888 headers: {
889 "Content-Type": "application/x-git-upload-pack-request",
890 "Git-Protocol": "version=2",
891 },
892 body: corruptedBody,
893 } as any);
894 
895 // The server should handle the corruption gracefully
896 // Either by returning an error status or completing with valid data
897 expect([200, 400, 500, 503].includes(res.status)).toBe(true);
898 
899 if (res.status === 200) {
900 // If it succeeded, check for valid response structure
901 const bytes = new Uint8Array(await res.arrayBuffer());
902 const lines = decodePktLines(bytes);
903 
904 // Check if there's a band-3 fatal message
905 for (const line of lines) {
906 if (line.type === "line" && line.raw?.[0] === 0x03) {
907 const fatalMsg = new TextDecoder().decode(line.raw.subarray(1));
908 expect(fatalMsg).toContain("fatal:");
909 }
910 }
911 
912 // Note: hasFatal might be false if the server recovered from the corruption
913 }
914 });
915 
916 it("handles abort signal during negotiation phase", async () => {
917 const owner = "o";
918 const repo = uniqueRepoId("abort-negotiation");
919 await setupRepoForTests(env, owner, repo);
920 const repoId = `${owner}/${repo}`;
921 
922 // Seed repository with commits
923 const id = env.REPO_DO.idFromName(repoId);
924 const { commitOid, parentOid } = await runDOWithRetry(
925 () => env.REPO_DO.get(id),
926 async (instance) => {
927 const first = await instance.seedMinimalRepo();
928 const second = await instance.seedMinimalRepo();
929 return {
930 commitOid: second.commitOid,
931 parentOid: first.commitOid,
932 };
933 }
934 );
935 
936 // Test abort during negotiation (done=false)
937 const body = buildFetchBody({
938 wants: [commitOid],
939 haves: [parentOid],
940 done: false,
941 });
942 
943 const abortController = new AbortController();
944 
945 // Abort immediately
946 abortController.abort();
947 
948 const res = await handleFetchV2Streaming(env as Env, repoId, body, abortController.signal);
949 expect(res.status).toBe(499);
950 });
951 
952 it("verifies streaming response includes progress messages", async () => {
953 const owner = "o";
954 const repo = uniqueRepoId("progress");
955 await setupRepoForTests(env, owner, repo);
956 const repoId = `${owner}/${repo}`;
957 
958 // Create a repository with enough content to trigger progress messages
959 const id = env.REPO_DO.idFromName(repoId);
960 const commits: string[] = [];
961 
962 await runDOWithRetry(
963 () => env.REPO_DO.get(id),
964 async (instance) => {
965 // Create multiple commits
966 for (let i = 0; i < 3; i++) {
967 const result = await instance.seedMinimalRepo();
968 commits.push(result.commitOid);
969 }
970 // Try to trigger packing if possible
971 // Note: this might fall back to loose objects
972 return commits;
973 }
974 );
975 
976 const body = buildFetchBody({ wants: commits, done: true });
977 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
978 
979 const res = await workerExports.default.fetch(url, {
980 method: "POST",
981 headers: {
982 "Content-Type": "application/x-git-upload-pack-request",
983 "Git-Protocol": "version=2",
984 },
985 body,
986 } as any);
987 
988 expect(res.status).toBe(200);
989 
990 const bytes = new Uint8Array(await res.arrayBuffer());
991 const lines = decodePktLines(bytes);
992 
993 // Look for band-2 progress messages
994 const progressMessages: string[] = [];
995 let inPackfile = false;
996 
997 for (const line of lines) {
998 if (line.type === "line" && line.text === "packfile\n") {
999 inPackfile = true;
1000 }
1001 if (inPackfile && line.type === "line" && line.raw?.[0] === 0x02) {
1002 // Band 2: progress message
1003 const msg = new TextDecoder().decode(line.raw.subarray(1));
1004 progressMessages.push(msg);
1005 }
1006 }
1007 
1008 expect(progressMessages.length).toBeGreaterThan(0);
1009 const packStart = findBytes(bytes, new TextEncoder().encode("PACK"));
1010 expect(packStart).toBeGreaterThan(-1);
1011 });
1012 
1013 it("emits pack preparation progress before pack data for initial fetches", async () => {
1014 const owner = "o";
1015 const repo = uniqueRepoId("early-progress");
1016 await setupRepoForTests(env, owner, repo);
1017 const repoId = `${owner}/${repo}`;
1018 
1019 // Seed repository with some commits and ensure they're packed
1020 const id = env.REPO_DO.idFromName(repoId);
1021 const { commitOid } = await runDOWithRetry(
1022 () => env.REPO_DO.get(id),
1023 async (instance) => instance.seedMinimalRepo()
1024 );
1025 
1026 const body = buildFetchBody({ wants: [commitOid], done: true });
1027 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
1028 
1029 const res = await workerExports.default.fetch(url, {
1030 method: "POST",
1031 headers: {
1032 "Content-Type": "application/x-git-upload-pack-request",
1033 "Git-Protocol": "version=2",
1034 },
1035 body,
1036 } as any);
1037 
1038 expect(res.status).toBe(200);
1039 
1040 const bytes = new Uint8Array(await res.arrayBuffer());
1041 const lines = decodePktLines(bytes);
1042 
1043 // Collect all progress messages and pack data in order
1044 const orderedOutput: { type: "progress" | "data"; content: string | number }[] = [];
1045 let inPackfile = false;
1046 
1047 for (const line of lines) {
1048 if (line.type === "line" && line.text === "packfile\n") {
1049 inPackfile = true;
1050 } else if (inPackfile && line.type === "line" && line.raw) {
1051 if (line.raw[0] === 0x02) {
1052 // Band 2: progress message
1053 const msg = new TextDecoder().decode(line.raw.subarray(1));
1054 orderedOutput.push({ type: "progress", content: msg });
1055 } else if (line.raw[0] === 0x01) {
1056 // Band 1: pack data - just record the first byte to prove data arrived
1057 orderedOutput.push({ type: "data", content: line.raw[1] });
1058 }
1059 }
1060 }
1061 
1062 // Verify we got progress messages
1063 const progressMessages = orderedOutput.filter((o) => o.type === "progress");
1064 expect(progressMessages.length).toBeGreaterThan(0);
1065 
1066 expect(progressMessages[0]?.content).toBe("Preparing pack...\n");
1067 
1068 // Verify progress comes before data
1069 const firstProgressIdx = orderedOutput.findIndex((o) => o.type === "progress");
1070 const firstDataIdx = orderedOutput.findIndex((o) => o.type === "data");
1071 
1072 expect(firstProgressIdx).toBeGreaterThanOrEqual(0);
1073 expect(firstDataIdx).toBeGreaterThanOrEqual(0);
1074 expect(firstProgressIdx).toBeLessThan(firstDataIdx);
1075 });
1076 
1077 it("emits pack preparation progress before pack data when haves are present", async () => {
1078 const owner = "o";
1079 const repo = uniqueRepoId("have-progress");
1080 await setupRepoForTests(env, owner, repo);
1081 const repoId = `${owner}/${repo}`;
1082 
1083 const id = env.REPO_DO.idFromName(repoId);
1084 const { commitOid, parentOid } = await runDOWithRetry(
1085 () => env.REPO_DO.get(id),
1086 async (instance) => {
1087 const first = await instance.seedMinimalRepo();
1088 const second = await instance.seedMinimalRepo();
1089 return {
1090 commitOid: second.commitOid,
1091 parentOid: first.commitOid,
1092 };
1093 }
1094 );
1095 
1096 const body = buildFetchBody({
1097 wants: [commitOid],
1098 haves: [parentOid],
1099 done: true,
1100 });
1101 const url = `https://example.com/${owner}/${repo}/git-upload-pack`;
1102 
1103 const res = await workerExports.default.fetch(url, {
1104 method: "POST",
1105 headers: {
1106 "Content-Type": "application/x-git-upload-pack-request",
1107 "Git-Protocol": "version=2",
1108 },
1109 body,
1110 } as any);
1111 
1112 expect(res.status).toBe(200);
1113 
1114 const bytes = new Uint8Array(await res.arrayBuffer());
1115 const lines = decodePktLines(bytes);
1116 
1117 const orderedOutput: { type: "progress" | "data"; content: string | number }[] = [];
1118 let inPackfile = false;
1119 
1120 for (const line of lines) {
1121 if (line.type === "line" && line.text === "packfile\n") {
1122 inPackfile = true;
1123 } else if (inPackfile && line.type === "line" && line.raw) {
1124 if (line.raw[0] === 0x02) {
1125 const msg = new TextDecoder().decode(line.raw.subarray(1));
1126 orderedOutput.push({ type: "progress", content: msg });
1127 } else if (line.raw[0] === 0x01) {
1128 orderedOutput.push({ type: "data", content: line.raw[1] });
1129 }
1130 }
1131 }
1132 
1133 const progressMessages = orderedOutput.filter((o) => o.type === "progress");
1134 expect(progressMessages[0]?.content).toBe("Preparing pack...\n");
1135 
1136 const prepareIdx = orderedOutput.findIndex(
1137 (o) => o.type === "progress" && o.content === "Preparing pack...\n"
1138 );
1139 const firstDataIdx = orderedOutput.findIndex((o) => o.type === "data");
1140 
1141 expect(prepareIdx).toBeGreaterThanOrEqual(0);
1142 expect(firstDataIdx).toBeGreaterThanOrEqual(0);
1143 expect(prepareIdx).toBeLessThan(firstDataIdx);
1144 });
1145});