Skip to content
File

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

typescript677 lines
1import { describe, expect, it } from "vitest";
2import { createExecutionContext } from "cloudflare:test";
3import { env, exports as workerExports } from "cloudflare:workers";
4import { concatChunks, flushPkt, pktLine } from "@/worker/git/core";
5import { computeOid, encodeGitObject } from "@/worker/git/core/objects";
6import { handleStreamingReceivePackPOST } from "@/worker/git/receive/streamReceivePack";
7import { buildFetchBody } from "./util/fetch-protocol";
8import { buildAppendOnlyDelta, buildPack, zero40 } from "./util/git-pack";
9import { buildTreePayload } from "./util/packed-repo";
10import {
11 callStubWithRetry,
12 deleteLooseObjectCopies,
13 toRequestBody,
14 uniqueRepoId,
15} from "./util/test-helpers";
16import { lookupPushAuth, setupRepoForTests } from "./util/repoSeed";
17import { seedPackFirstRepo } from "./util/pack-first";
18import { doPrefix, packRefsKey, r2PackDirPrefix } from "@/worker/keys";
19import {
20 buildStreamingReceiveBody,
21 decodeReceiveSideband,
22 decodeReportStatus,
23 promoteToStreaming,
24 pushStreamingUpdate,
25} from "./util/streaming-helpers";
26 
27function streamBody(bytes: Uint8Array, chunkSize = 1024): ReadableStream<Uint8Array> {
28 return new ReadableStream<Uint8Array>({
29 start(controller) {
30 for (let offset = 0; offset < bytes.byteLength; offset += chunkSize) {
31 controller.enqueue(bytes.subarray(offset, offset + chunkSize));
32 }
33 controller.close();
34 },
35 });
36}
37 
38function abortingStreamBody(
39 bytes: Uint8Array,
40 abortController: AbortController,
41 options?: {
42 chunkSize?: number;
43 abortAfterChunks?: number;
44 }
45): ReadableStream<Uint8Array> {
46 const chunkSize = options?.chunkSize ?? 256;
47 const abortAfterChunks = options?.abortAfterChunks ?? 1;
48 let offset = 0;
49 let emittedChunks = 0;
50 
51 return new ReadableStream<Uint8Array>({
52 pull(controller) {
53 if (abortController.signal.aborted) {
54 const error = new Error("client aborted");
55 error.name = "AbortError";
56 controller.error(error);
57 return;
58 }
59 if (offset >= bytes.byteLength) {
60 controller.close();
61 return;
62 }
63 
64 controller.enqueue(bytes.subarray(offset, offset + chunkSize));
65 offset += chunkSize;
66 emittedChunks++;
67 
68 if (emittedChunks >= abortAfterChunks && !abortController.signal.aborted) {
69 abortController.abort();
70 }
71 },
72 });
73}
74 
75async function listStagedReceivePacks(repoId: string): Promise<string[]> {
76 const doId = env.REPO_DO.idFromName(repoId);
77 const prefix = r2PackDirPrefix(doPrefix(doId.toString()));
78 const listed = await env.REPO_BUCKET.list({ prefix });
79 return listed.objects.map((object) => object.key).filter((key) => key.includes("/pack-rx-"));
80}
81 
82function pushAuthFromUrl(url: string): string | undefined {
83 const match = /https?:\/\/[^/]+\/([^/]+)\/([^/]+)\/git-receive-pack/.exec(url);
84 return match ? lookupPushAuth(match[1]!, match[2]!) : undefined;
85}
86 
87async function pushBody(
88 url: string,
89 body: Uint8Array,
90 options?: {
91 stream?: boolean;
92 authHeader?: string;
93 }
94): Promise<Response> {
95 const headers: Record<string, string> = {
96 "Content-Type": "application/x-git-receive-pack-request",
97 };
98 const auth = options?.authHeader ?? pushAuthFromUrl(url);
99 if (auth) headers.Authorization = auth;
100 return await workerExports.default.fetch(url, {
101 method: "POST",
102 headers,
103 body: options?.stream ? streamBody(body) : body,
104 } as any);
105}
106 
107describe("streaming receive-pack", () => {
108 it("returns 499 when the request is already aborted before receive work starts", async () => {
109 const owner = "o";
110 const repo = uniqueRepoId("stream-receive-aborted-start");
111 await setupRepoForTests(env, owner, repo);
112 const repoId = `${owner}/${repo}`;
113 const seeded = await seedPackFirstRepo(repoId);
114 await promoteToStreaming(owner, repo);
115 
116 const abortController = new AbortController();
117 abortController.abort();
118 
119 const body = concatChunks([
120 pktLine(
121 `${seeded.nextCommit.oid} ${seeded.nextCommit.oid} refs/heads/main\0 report-status ofs-delta agent=test\n`
122 ),
123 flushPkt(),
124 ]);
125 const request = new Request(`https://example.com/${owner}/${repo}/git-receive-pack`, {
126 method: "POST",
127 headers: { "Content-Type": "application/x-git-receive-pack-request" },
128 body: toRequestBody(body),
129 signal: abortController.signal,
130 });
131 
132 const response = await handleStreamingReceivePackPOST(
133 env,
134 repoId,
135 request,
136 createExecutionContext()
137 );
138 expect(response.status).toBe(499);
139 
140 const activity = await callStubWithRetry(seeded.getStub, (stub) => stub.getRepoActivity());
141 expect(activity).toBeNull();
142 expect(await listStagedReceivePacks(repoId)).toEqual([]);
143 });
144 
145 it("streams a create push and fetch still works after deleting all loose copies", async () => {
146 const owner = "o";
147 const repo = uniqueRepoId("stream-receive-create");
148 await setupRepoForTests(env, owner, repo);
149 const repoId = `${owner}/${repo}`;
150 const seeded = await seedPackFirstRepo(repoId);
151 await promoteToStreaming(owner, repo);
152 
153 const author = "You <you@example.com> 0 +0000";
154 const blobPayload = new TextEncoder().encode("version three\n");
155 const blob = await encodeGitObject("blob", blobPayload);
156 const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blob.oid }]);
157 const tree = await encodeGitObject("tree", treePayload);
158 const commitPayload = new TextEncoder().encode(
159 `tree ${tree.oid}\n` +
160 `parent ${seeded.nextCommit.oid}\n` +
161 `author ${author}\n` +
162 `committer ${author}\n\n` +
163 `third commit\n`
164 );
165 const commit = await encodeGitObject("commit", commitPayload);
166 const pack = await buildPack([
167 { type: "blob", payload: blobPayload },
168 { type: "tree", payload: treePayload },
169 { type: "commit", payload: commitPayload },
170 ]);
171 const body = concatChunks([
172 pktLine(
173 `${seeded.nextCommit.oid} ${commit.oid} refs/heads/main\0 report-status ofs-delta agent=test\n`
174 ),
175 flushPkt(),
176 pack,
177 ]);
178 
179 const response = await pushBody(`https://example.com/${owner}/${repo}/git-receive-pack`, body, {
180 stream: true,
181 });
182 expect(response.status).toBe(200);
183 expect(decodeReportStatus(new Uint8Array(await response.arrayBuffer()))).toContain(
184 "ok refs/heads/main"
185 );
186 
187 await deleteLooseObjectCopies(env, seeded.getStub, seeded.objectOids);
188 
189 const rawResponse = await workerExports.default.fetch(
190 `https://example.com/${owner}/${repo}/raw?oid=${blob.oid}&name=README.md`
191 );
192 expect(rawResponse.status).toBe(200);
193 expect(await rawResponse.text()).toBe("version three\n");
194 
195 const fetchResponse = await workerExports.default.fetch(
196 `https://example.com/${owner}/${repo}/git-upload-pack`,
197 {
198 method: "POST",
199 headers: {
200 "Content-Type": "application/x-git-upload-pack-request",
201 "Git-Protocol": "version=2",
202 },
203 body: toRequestBody(
204 buildFetchBody({
205 wants: [commit.oid],
206 haves: [seeded.nextCommit.oid],
207 done: true,
208 })
209 ),
210 }
211 );
212 expect(fetchResponse.status).toBe(200);
213 const fetchBytes = new Uint8Array(await fetchResponse.arrayBuffer());
214 expect(new TextDecoder().decode(fetchBytes.subarray(4, 13))).toBe("packfile\n");
215 });
216 
217 it("reports upload, scan, resolve, and final status over side-band-64k", async () => {
218 const owner = "o";
219 const repo = uniqueRepoId("stream-receive-sideband");
220 await setupRepoForTests(env, owner, repo);
221 const repoId = `${owner}/${repo}`;
222 const seeded = await seedPackFirstRepo(repoId);
223 await promoteToStreaming(owner, repo);
224 
225 const author = "You <you@example.com> 0 +0000";
226 const basePayload = new TextEncoder().encode("version two\n");
227 const suffix = new TextEncoder().encode("sideband progress\n");
228 const delta = buildAppendOnlyDelta(basePayload, suffix);
229 const blobPayload = new Uint8Array(basePayload.byteLength + suffix.byteLength);
230 blobPayload.set(basePayload, 0);
231 blobPayload.set(suffix, basePayload.byteLength);
232 const blobOid = await computeOid("blob", blobPayload);
233 const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blobOid }]);
234 const tree = await encodeGitObject("tree", treePayload);
235 const commitPayload = new TextEncoder().encode(
236 `tree ${tree.oid}\n` +
237 `parent ${seeded.nextCommit.oid}\n` +
238 `author ${author}\n` +
239 `committer ${author}\n\n` +
240 `sideband progress\n`
241 );
242 const commit = await encodeGitObject("commit", commitPayload);
243 const pack = await buildPack([
244 { type: "ref-delta", baseOid: seeded.nextBlob.oid, delta },
245 { type: "tree", payload: treePayload },
246 { type: "commit", payload: commitPayload },
247 ]);
248 
249 const response = await pushBody(
250 `https://example.com/${owner}/${repo}/git-receive-pack`,
251 concatChunks([
252 pktLine(
253 `${seeded.nextCommit.oid} ${commit.oid} refs/heads/main\0 report-status side-band-64k ofs-delta agent=test\n`
254 ),
255 flushPkt(),
256 pack,
257 ]),
258 { stream: true }
259 );
260 expect(response.status).toBe(200);
261 
262 const decoded = decodeReceiveSideband(new Uint8Array(await response.arrayBuffer()));
263 expect(decoded.progress.some((line) => line.includes("Uploading pack to object storage"))).toBe(
264 true
265 );
266 expect(decoded.progress.some((line) => line.includes("Scanning pack objects"))).toBe(true);
267 expect(decoded.progress.some((line) => line.includes("Resolving deltas"))).toBe(true);
268 expect(decoded.progress.some((line) => line.includes("Writing pack index"))).toBe(true);
269 expect(decoded.progress.some((line) => line.includes("Writing pack reference index"))).toBe(
270 true
271 );
272 expect(decoded.reportStatus).toContain("ok refs/heads/main");
273 expect(decoded.fatal).toEqual([]);
274 
275 const catalog = await callStubWithRetry(seeded.getStub, (stub) => stub.getActivePackCatalog());
276 const receivedPack = catalog.find((row) => row.packKey.includes("/pack-rx-"));
277 expect(receivedPack).toBeTruthy();
278 await expect(env.REPO_BUCKET.head(packRefsKey(receivedPack!.packKey))).resolves.toBeTruthy();
279 });
280 
281 it("suppresses side-band progress when the client requests quiet", async () => {
282 const owner = "o";
283 const repo = uniqueRepoId("stream-receive-sideband-quiet");
284 await setupRepoForTests(env, owner, repo);
285 const repoId = `${owner}/${repo}`;
286 const seeded = await seedPackFirstRepo(repoId);
287 await promoteToStreaming(owner, repo);
288 
289 const push = await buildStreamingReceiveBody({
290 parentOid: seeded.nextCommit.oid,
291 nextText: "quiet progress\n",
292 commitMessage: "quiet progress",
293 capabilities: "report-status side-band-64k quiet ofs-delta agent=test",
294 });
295 
296 const response = await pushBody(
297 `https://example.com/${owner}/${repo}/git-receive-pack`,
298 push.body,
299 {
300 stream: true,
301 }
302 );
303 expect(response.status).toBe(200);
304 
305 const decoded = decodeReceiveSideband(new Uint8Array(await response.arrayBuffer()));
306 expect(decoded.progress).toEqual([]);
307 expect(decoded.reportStatus).toContain("ok refs/heads/main");
308 expect(decoded.fatal).toEqual([]);
309 });
310 
311 it("handles delete-only pushes in streaming mode", async () => {
312 const owner = "o";
313 const repo = uniqueRepoId("stream-receive-delete");
314 const seededRepo = await setupRepoForTests(env, owner, repo);
315 const repoId = `${owner}/${repo}`;
316 const seeded = await seedPackFirstRepo(repoId);
317 await promoteToStreaming(owner, repo);
318 
319 const author = "You <you@example.com> 0 +0000";
320 const blobPayload = new TextEncoder().encode("feature branch\n");
321 const blob = await encodeGitObject("blob", blobPayload);
322 const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blob.oid }]);
323 const tree = await encodeGitObject("tree", treePayload);
324 const commitPayload = new TextEncoder().encode(
325 `tree ${tree.oid}\n` +
326 `parent ${seeded.nextCommit.oid}\n` +
327 `author ${author}\n` +
328 `committer ${author}\n\n` +
329 `feature commit\n`
330 );
331 const commit = await encodeGitObject("commit", commitPayload);
332 const createPack = await buildPack([
333 { type: "blob", payload: blobPayload },
334 { type: "tree", payload: treePayload },
335 { type: "commit", payload: commitPayload },
336 ]);
337 
338 const createResponse = await pushBody(
339 `https://example.com/${owner}/${repo}/git-receive-pack`,
340 concatChunks([
341 pktLine(
342 `${zero40()} ${commit.oid} refs/heads/feature\0 report-status ofs-delta agent=test\n`
343 ),
344 flushPkt(),
345 createPack,
346 ])
347 );
348 expect(createResponse.status).toBe(200);
349 
350 const deleteResponse = await pushBody(
351 `https://example.com/${owner}/${repo}/git-receive-pack`,
352 concatChunks([
353 pktLine(`${commit.oid} ${zero40()} refs/heads/feature\0 report-status\n`),
354 flushPkt(),
355 ])
356 );
357 expect(deleteResponse.status).toBe(200);
358 expect(decodeReportStatus(new Uint8Array(await deleteResponse.arrayBuffer()))).toContain(
359 "ok refs/heads/feature"
360 );
361 
362 const refsResponse = await workerExports.default.fetch(
363 `https://example.com/${owner}/${repo}/admin/refs`,
364 { headers: { Cookie: seededRepo.cookieHeader } }
365 );
366 const refs = (await refsResponse.json()) as Array<{ name: string; oid: string }>;
367 expect(refs.find((ref) => ref.name === "refs/heads/feature")).toBeUndefined();
368 });
369 
370 it("rejects stale old-oids and leaves no staged receive packs behind", async () => {
371 const owner = "o";
372 const repo = uniqueRepoId("stream-receive-stale");
373 await setupRepoForTests(env, owner, repo);
374 const repoId = `${owner}/${repo}`;
375 const seeded = await seedPackFirstRepo(repoId);
376 await promoteToStreaming(owner, repo);
377 
378 const author = "You <you@example.com> 0 +0000";
379 const blobPayload = new TextEncoder().encode("stale branch\n");
380 const blob = await encodeGitObject("blob", blobPayload);
381 const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blob.oid }]);
382 const tree = await encodeGitObject("tree", treePayload);
383 const commitPayload = new TextEncoder().encode(
384 `tree ${tree.oid}\n` +
385 `parent ${seeded.nextCommit.oid}\n` +
386 `author ${author}\n` +
387 `committer ${author}\n\n` +
388 `stale commit\n`
389 );
390 const commit = await encodeGitObject("commit", commitPayload);
391 const pack = await buildPack([
392 { type: "blob", payload: blobPayload },
393 { type: "tree", payload: treePayload },
394 { type: "commit", payload: commitPayload },
395 ]);
396 
397 const response = await pushBody(
398 `https://example.com/${owner}/${repo}/git-receive-pack`,
399 concatChunks([
400 pktLine(`${zero40()} ${commit.oid} refs/heads/main\0 report-status ofs-delta agent=test\n`),
401 flushPkt(),
402 pack,
403 ]),
404 { stream: true }
405 );
406 expect(response.status).toBe(200);
407 const lines = decodeReportStatus(new Uint8Array(await response.arrayBuffer()));
408 expect(lines.some((line) => line.startsWith("ng refs/heads/main stale old-oid"))).toBe(true);
409 expect(await listStagedReceivePacks(repoId)).toEqual([]);
410 });
411 
412 it("accepts thin packs with active external bases, rejects missing ones, and clears the receive lease after failure", async () => {
413 const owner = "o";
414 const repo = uniqueRepoId("stream-receive-thin");
415 await setupRepoForTests(env, owner, repo);
416 const repoId = `${owner}/${repo}`;
417 const seeded = await seedPackFirstRepo(repoId);
418 await promoteToStreaming(owner, repo);
419 
420 const author = "You <you@example.com> 0 +0000";
421 const basePayload = new TextEncoder().encode("version two\n");
422 const suffix = new TextEncoder().encode("delta tail\n");
423 const delta = buildAppendOnlyDelta(basePayload, suffix);
424 const blobPayload = new Uint8Array(basePayload.byteLength + suffix.byteLength);
425 blobPayload.set(basePayload, 0);
426 blobPayload.set(suffix, basePayload.byteLength);
427 const blobOid = await computeOid("blob", blobPayload);
428 const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blobOid }]);
429 const tree = await encodeGitObject("tree", treePayload);
430 const commitPayload = new TextEncoder().encode(
431 `tree ${tree.oid}\n` +
432 `parent ${seeded.nextCommit.oid}\n` +
433 `author ${author}\n` +
434 `committer ${author}\n\n` +
435 `thin commit\n`
436 );
437 const commit = await encodeGitObject("commit", commitPayload);
438 const goodPack = await buildPack([
439 { type: "ref-delta", baseOid: seeded.nextBlob.oid, delta },
440 { type: "tree", payload: treePayload },
441 { type: "commit", payload: commitPayload },
442 ]);
443 
444 const goodResponse = await pushBody(
445 `https://example.com/${owner}/${repo}/git-receive-pack`,
446 concatChunks([
447 pktLine(
448 `${seeded.nextCommit.oid} ${commit.oid} refs/heads/main\0 report-status ofs-delta agent=test\n`
449 ),
450 flushPkt(),
451 goodPack,
452 ])
453 );
454 expect(goodResponse.status).toBe(200);
455 expect(decodeReportStatus(new Uint8Array(await goodResponse.arrayBuffer()))).toContain(
456 "ok refs/heads/main"
457 );
458 const packKeysBeforeBadPush = await listStagedReceivePacks(repoId);
459 
460 const badTreePayload = buildTreePayload([
461 { mode: "100644", name: "README.md", oid: "ab".repeat(20) },
462 ]);
463 const badTree = await encodeGitObject("tree", badTreePayload);
464 const badCommitPayload = new TextEncoder().encode(
465 `tree ${badTree.oid}\n` +
466 `parent ${commit.oid}\n` +
467 `author ${author}\n` +
468 `committer ${author}\n\n` +
469 `bad thin commit\n`
470 );
471 const badCommit = await encodeGitObject("commit", badCommitPayload);
472 const missingBasePack = await buildPack([
473 {
474 type: "ref-delta",
475 baseOid: "cd".repeat(20),
476 delta: buildAppendOnlyDelta(
477 new TextEncoder().encode("base\n"),
478 new TextEncoder().encode("missing\n")
479 ),
480 },
481 { type: "tree", payload: badTreePayload },
482 { type: "commit", payload: badCommitPayload },
483 ]);
484 
485 const badResponse = await pushBody(
486 `https://example.com/${owner}/${repo}/git-receive-pack`,
487 concatChunks([
488 pktLine(
489 `${commit.oid} ${badCommit.oid} refs/heads/main\0 report-status ofs-delta agent=test\n`
490 ),
491 flushPkt(),
492 missingBasePack,
493 ]),
494 { stream: true }
495 );
496 expect(badResponse.status).toBe(400);
497 expect(await listStagedReceivePacks(repoId)).toEqual(packKeysBeforeBadPush);
498 
499 const activityAfterBadPush = await callStubWithRetry(seeded.getStub, (stub) =>
500 stub.getRepoActivity()
501 );
502 expect(activityAfterBadPush).toBeNull();
503 
504 const retryPush = await pushStreamingUpdate(owner, repo, commit.oid, "cleanup retry\n");
505 expect(retryPush.commitOid).not.toBe(commit.oid);
506 });
507 
508 it("returns unpack error over side-band-64k when resolve fails after streaming has started", async () => {
509 const owner = "o";
510 const repo = uniqueRepoId("stream-receive-sideband-failure");
511 await setupRepoForTests(env, owner, repo);
512 const repoId = `${owner}/${repo}`;
513 const seeded = await seedPackFirstRepo(repoId);
514 await promoteToStreaming(owner, repo);
515 
516 const author = "You <you@example.com> 0 +0000";
517 const badTreePayload = buildTreePayload([
518 { mode: "100644", name: "README.md", oid: "ab".repeat(20) },
519 ]);
520 const badTree = await encodeGitObject("tree", badTreePayload);
521 const badCommitPayload = new TextEncoder().encode(
522 `tree ${badTree.oid}\n` +
523 `parent ${seeded.nextCommit.oid}\n` +
524 `author ${author}\n` +
525 `committer ${author}\n\n` +
526 `bad sideband commit\n`
527 );
528 const badCommit = await encodeGitObject("commit", badCommitPayload);
529 const missingBasePack = await buildPack([
530 {
531 type: "ref-delta",
532 baseOid: "cd".repeat(20),
533 delta: buildAppendOnlyDelta(
534 new TextEncoder().encode("base\n"),
535 new TextEncoder().encode("missing\n")
536 ),
537 },
538 { type: "tree", payload: badTreePayload },
539 { type: "commit", payload: badCommitPayload },
540 ]);
541 
542 const response = await pushBody(
543 `https://example.com/${owner}/${repo}/git-receive-pack`,
544 concatChunks([
545 pktLine(
546 `${seeded.nextCommit.oid} ${badCommit.oid} refs/heads/main\0 report-status side-band-64k ofs-delta agent=test\n`
547 ),
548 flushPkt(),
549 missingBasePack,
550 ]),
551 { stream: true }
552 );
553 expect(response.status).toBe(200);
554 
555 const decoded = decodeReceiveSideband(new Uint8Array(await response.arrayBuffer()));
556 expect(decoded.progress.some((line) => line.includes("Uploading pack to object storage"))).toBe(
557 true
558 );
559 expect(decoded.reportStatus.some((line) => line.startsWith("unpack error"))).toBe(true);
560 expect(decoded.reportStatus.some((line) => line.startsWith("ng refs/heads/main"))).toBe(true);
561 expect(await listStagedReceivePacks(repoId)).toEqual([]);
562 
563 const activity = await callStubWithRetry(seeded.getStub, (stub) => stub.getRepoActivity());
564 expect(activity).toBeNull();
565 });
566 
567 it("returns 499 and cleans up when the request aborts during the streaming upload", async () => {
568 const owner = "o";
569 const repo = uniqueRepoId("stream-receive-abort-upload");
570 await setupRepoForTests(env, owner, repo);
571 const repoId = `${owner}/${repo}`;
572 const seeded = await seedPackFirstRepo(repoId);
573 await promoteToStreaming(owner, repo);
574 
575 const author = "You <you@example.com> 0 +0000";
576 const blobPayload = new TextEncoder().encode("aborted upload\n");
577 const blob = await encodeGitObject("blob", blobPayload);
578 const treePayload = buildTreePayload([{ mode: "100644", name: "README.md", oid: blob.oid }]);
579 const tree = await encodeGitObject("tree", treePayload);
580 const commitPayload = new TextEncoder().encode(
581 `tree ${tree.oid}\n` +
582 `parent ${seeded.nextCommit.oid}\n` +
583 `author ${author}\n` +
584 `committer ${author}\n\n` +
585 `aborted upload\n`
586 );
587 const commit = await encodeGitObject("commit", commitPayload);
588 const pack = await buildPack([
589 { type: "blob", payload: blobPayload },
590 { type: "tree", payload: treePayload },
591 { type: "commit", payload: commitPayload },
592 ]);
593 const body = concatChunks([
594 pktLine(
595 `${seeded.nextCommit.oid} ${commit.oid} refs/heads/main\0 report-status ofs-delta agent=test\n`
596 ),
597 flushPkt(),
598 pack,
599 ]);
600 
601 const abortController = new AbortController();
602 const request = new Request(`https://example.com/${owner}/${repo}/git-receive-pack`, {
603 method: "POST",
604 headers: { "Content-Type": "application/x-git-receive-pack-request" },
605 body: abortingStreamBody(body, abortController),
606 signal: abortController.signal,
607 });
608 
609 const response = await handleStreamingReceivePackPOST(
610 env,
611 repoId,
612 request,
613 createExecutionContext()
614 );
615 expect(response.status).toBe(499);
616 expect(await listStagedReceivePacks(repoId)).toEqual([]);
617 
618 const activity = await callStubWithRetry(seeded.getStub, (stub) => stub.getRepoActivity());
619 expect(activity).toBeNull();
620 
621 const retryPush = await pushStreamingUpdate(
622 owner,
623 repo,
624 seeded.nextCommit.oid,
625 "after abort\n"
626 );
627 expect(retryPush.commitOid).not.toBe(seeded.nextCommit.oid);
628 });
629 
630 it("returns 503 when a streaming receive lease is already active", async () => {
631 const owner = "o";
632 const repo = uniqueRepoId("stream-receive-busy");
633 await setupRepoForTests(env, owner, repo);
634 const repoId = `${owner}/${repo}`;
635 const seeded = await seedPackFirstRepo(repoId);
636 await promoteToStreaming(owner, repo);
637 
638 const begin = await callStubWithRetry<any>(seeded.getStub, (stub) => stub.beginReceive());
639 if (!begin.ok) {
640 throw new Error("expected test receive lease to be granted");
641 }
642 
643 const response = await pushBody(
644 `https://example.com/${owner}/${repo}/git-receive-pack`,
645 concatChunks([
646 pktLine(`${zero40()} ${seeded.nextCommit.oid} refs/heads/main\0 report-status\n`),
647 flushPkt(),
648 ])
649 );
650 expect(response.status).toBe(503);
651 expect(response.headers.get("Retry-After")).toBe("10");
652 
653 await callStubWithRetry(seeded.getStub, (stub) => stub.abortReceive(begin.lease.token));
654 });
655 
656 it("rejects invalid refs without leaving staged receive packs behind", async () => {
657 const owner = "o";
658 const repo = uniqueRepoId("stream-receive-invalid-ref");
659 await setupRepoForTests(env, owner, repo);
660 const repoId = `${owner}/${repo}`;
661 await seedPackFirstRepo(repoId);
662 await promoteToStreaming(owner, repo);
663 
664 const response = await pushBody(
665 `https://example.com/${owner}/${repo}/git-receive-pack`,
666 concatChunks([
667 pktLine(`${zero40()} ${"a".repeat(40)} HEAD\0 report-status ofs-delta agent=test\n`),
668 flushPkt(),
669 ])
670 );
671 expect(response.status).toBe(200);
672 const lines = decodeReportStatus(new Uint8Array(await response.arrayBuffer()));
673 expect(lines.some((line) => line.startsWith("unpack error invalid-ref"))).toBe(true);
674 expect(await listStagedReceivePacks(repoId)).toEqual([]);
675 });
676});