File
Blob: test/fetch-streaming.worker.test.ts
| 1 | import { it, expect, describe } from "vitest"; |
| 2 | import { env, exports as workerExports } from "cloudflare:workers"; |
| 3 | import { pktLine, delimPkt, flushPkt, concatChunks, decodePktLines } from "@/worker/git"; |
| 4 | import { handleFetchV2Streaming } from "@/worker/git/operations/uploadStream"; |
| 5 | import { |
| 6 | buildServeUploadPackPlan, |
| 7 | loadUploadPackSnapshot, |
| 8 | planUploadPack, |
| 9 | } from "@/worker/git/operations/fetch/plan"; |
| 10 | import { uniqueRepoId, runDOWithRetry } from "./util/test-helpers"; |
| 11 | import { setupRepoForTests } from "./util/repoSeed"; |
| 12 | import { asBufferSource } from "@/worker/common"; |
| 13 | import { packRefsKey } from "@/worker/keys"; |
| 14 | import { runQueueMessage } from "./util/queue"; |
| 15 | import { createTestCacheContext, seedPackFirstRepo } from "./util/pack-first"; |
| 16 | import { makeTracingLimiter } from "./util/pack-indexer.helpers"; |
| 17 | import { buildAppendOnlyDelta, buildCopyPrefixDelta, buildPack } from "./util/git-pack"; |
| 18 | import { seedPackedRepoState } from "./util/packed-repo"; |
| 19 | import { computeOid, encodeGitObject } from "@/worker/git/core/objects"; |
| 20 | |
| 21 | function 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 | */ |
| 43 | function 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 | |
| 53 | async 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 | |
| 67 | async 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 | |
| 91 | async 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 | |
| 130 | describe("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 | }); |