Skip to content
File

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

typescript972 lines
1import { describe, expect, it, vi } from "vitest";
2import { env, exports as workerExports } from "cloudflare:workers";
3import { getRepoStub } from "@/worker/common";
4import { bytesToHex } from "@/worker/common/hex";
5import { encodeGitObject } from "@/worker/git/core/objects";
6import { concatChunks, decodePktLines } from "@/worker/git";
7import { packIndexKey, packRefsKey } from "@/worker/keys";
8import { buildFetchBody } from "./util/fetch-protocol";
9import {
10 deleteLooseObjectCopies,
11 uniqueRepoId,
12 runDOWithRetry,
13 seedPackedRepoState,
14 buildTreePayload,
15 buildPack,
16 buildAppendOnlyDelta,
17} from "./util/test-helpers";
18import { setupRepoForTests } from "./util/repoSeed";
19import { seedPackFirstRepo } from "./util/pack-first";
20import { indexTestPack } from "./util/test-indexer";
21import { decodeReportStatus, promoteToStreaming } from "./util/streaming-helpers";
22import { asTypedStorage, type RepoStateSchema } from "@/worker/do/repo/repoState";
23import {
24 getDb,
25 getPackCatalogRow,
26 listActivePackCatalog,
27 upsertPackCatalogRow,
28} from "@/worker/do/repo/db";
29import { AUTO_COMPACTION_MAX_SOURCE_OBJECTS } from "@/worker/do/repo/catalog/compaction/plan";
30import { createQueueSendResponse } from "./util/queue";
31import {
32 compactOnce,
33 deleteSupersededOnce,
34 collectPackObjects,
35 pushOverflowingStreamingHistory,
36} from "./util/compaction-helpers";
37 
38type DebugState = {
39 activePacks?: Array<{ key: string; tier: number; kind: string }>;
40 supersededPacks?: Array<{ key: string; tier: number; kind: string }>;
41 compaction?: { queued?: boolean };
42};
43 
44async function getDebugState(
45 owner: string,
46 repo: string,
47 cookieHeader: string
48): Promise<DebugState> {
49 const response = await workerExports.default.fetch(
50 `https://example.com/${owner}/${repo}/admin/debug-state`,
51 { headers: { Cookie: cookieHeader } }
52 );
53 expect(response.status).toBe(200);
54 return (await response.json()) as DebugState;
55}
56 
57describe("streaming compaction", () => {
58 it("previews and requests real compaction work only after streaming overflow exists", async () => {
59 const owner = "o";
60 const repo = uniqueRepoId("stream-compaction-admin");
61 const seededRepo = await setupRepoForTests(env, owner, repo);
62 const repoId = `${owner}/${repo}`;
63 const seeded = await seedPackFirstRepo(repoId);
64 await promoteToStreaming(owner, repo);
65 
66 const sendSpy = vi
67 .spyOn(env.REPO_TASKS_QUEUE, "send")
68 .mockImplementation(async () => createQueueSendResponse());
69 try {
70 await pushOverflowingStreamingHistory({
71 owner,
72 repo,
73 repoId,
74 startingCommitOid: seeded.nextCommit.oid,
75 updates: 4,
76 });
77 sendSpy.mockClear();
78 
79 const previewResponse = await workerExports.default.fetch(
80 `https://example.com/${owner}/${repo}/admin/compact`,
81 {
82 method: "POST",
83 headers: {
84 Cookie: seededRepo.cookieHeader,
85 "Content-Type": "application/json",
86 Origin: "https://example.com",
87 },
88 body: JSON.stringify({}),
89 }
90 );
91 expect(previewResponse.status).toBe(200);
92 const previewJson = (await previewResponse.json()) as {
93 action?: string;
94 status?: string;
95 plan?: {
96 sourcePacks?: Array<{ packKey?: string }>;
97 sourceTier?: number;
98 targetTier?: number;
99 };
100 };
101 expect(previewJson.action).toBe("preview");
102 expect(previewJson.status).toBe("ok");
103 expect(previewJson.plan?.sourcePacks?.length).toBe(4);
104 expect(previewJson.plan?.sourceTier).toBe(0);
105 expect(previewJson.plan?.targetTier).toBe(1);
106 
107 const requestResponse = await workerExports.default.fetch(
108 `https://example.com/${owner}/${repo}/admin/compact`,
109 {
110 method: "POST",
111 headers: {
112 Cookie: seededRepo.cookieHeader,
113 "Content-Type": "application/json",
114 Origin: "https://example.com",
115 },
116 body: JSON.stringify({ dryRun: false }),
117 }
118 );
119 expect(requestResponse.status).toBe(202);
120 const requestJson = (await requestResponse.json()) as {
121 action?: string;
122 status?: string;
123 shouldEnqueue?: boolean;
124 };
125 expect(requestJson.action).toBe("request");
126 expect(requestJson.status).toBe("queued");
127 expect(requestJson.shouldEnqueue).toBe(true);
128 expect(sendSpy).toHaveBeenCalledWith({
129 kind: "compaction",
130 doId: env.REPO_DO.idFromName(repoId).toString(),
131 repoId,
132 });
133 } finally {
134 sendSpy.mockRestore();
135 }
136 });
137 
138 it("keeps the admin request queued when queue enqueue fails after DO state is recorded", async () => {
139 const owner = "o";
140 const repo = uniqueRepoId("stream-compaction-admin-enqueue-failure");
141 const seededRepo = await setupRepoForTests(env, owner, repo);
142 const repoId = `${owner}/${repo}`;
143 const seeded = await seedPackFirstRepo(repoId);
144 await promoteToStreaming(owner, repo);
145 
146 const sendSpy = vi
147 .spyOn(env.REPO_TASKS_QUEUE, "send")
148 .mockImplementation(async () => createQueueSendResponse());
149 try {
150 await pushOverflowingStreamingHistory({
151 owner,
152 repo,
153 repoId,
154 startingCommitOid: seeded.nextCommit.oid,
155 updates: 4,
156 });
157 sendSpy.mockClear();
158 sendSpy.mockRejectedValue(new Error("queue unavailable"));
159 
160 const requestResponse = await workerExports.default.fetch(
161 `https://example.com/${owner}/${repo}/admin/compact`,
162 {
163 method: "POST",
164 headers: {
165 Cookie: seededRepo.cookieHeader,
166 "Content-Type": "application/json",
167 Origin: "https://example.com",
168 },
169 body: JSON.stringify({ dryRun: false }),
170 }
171 );
172 expect(requestResponse.status).toBe(202);
173 const requestJson = (await requestResponse.json()) as {
174 action?: string;
175 status?: string;
176 shouldEnqueue?: boolean;
177 };
178 expect(requestJson.action).toBe("request");
179 expect(requestJson.status).toBe("queued");
180 expect(requestJson.shouldEnqueue).toBe(true);
181 
182 const stateAfterRequest = await getDebugState(owner, repo, seededRepo.cookieHeader);
183 expect(stateAfterRequest.compaction?.queued).toBe(true);
184 expect(sendSpy).toHaveBeenCalledTimes(1);
185 } finally {
186 sendSpy.mockRestore();
187 }
188 });
189 
190 it("previews compaction plan even after clearing compactionWantedAt", async () => {
191 const owner = "o";
192 const repo = uniqueRepoId("stream-compaction-preview-cleared");
193 const seededRepo = await setupRepoForTests(env, owner, repo);
194 const repoId = `${owner}/${repo}`;
195 const seeded = await seedPackFirstRepo(repoId);
196 await promoteToStreaming(owner, repo);
197 
198 const sendSpy = vi
199 .spyOn(env.REPO_TASKS_QUEUE, "send")
200 .mockImplementation(async () => createQueueSendResponse());
201 try {
202 await pushOverflowingStreamingHistory({
203 owner,
204 repo,
205 repoId,
206 startingCommitOid: seeded.nextCommit.oid,
207 updates: 4,
208 });
209 sendSpy.mockClear();
210 
211 // Request compaction so compactionWantedAt is set.
212 const requestResponse = await workerExports.default.fetch(
213 `https://example.com/${owner}/${repo}/admin/compact`,
214 {
215 method: "POST",
216 headers: {
217 Cookie: seededRepo.cookieHeader,
218 "Content-Type": "application/json",
219 Origin: "https://example.com",
220 },
221 body: JSON.stringify({ dryRun: false }),
222 }
223 );
224 expect(requestResponse.status).toBe(202);
225 
226 // Clear the recorded request.
227 const clearResponse = await workerExports.default.fetch(
228 `https://example.com/${owner}/${repo}/admin/compact`,
229 {
230 method: "DELETE",
231 headers: { Cookie: seededRepo.cookieHeader, Origin: "https://example.com" },
232 }
233 );
234 expect(clearResponse.status).toBe(200);
235 const clearJson = (await clearResponse.json()) as { cleared?: boolean };
236 expect(clearJson.cleared).toBe(true);
237 
238 // Preview should still show the plan with queued: false.
239 const previewResponse = await workerExports.default.fetch(
240 `https://example.com/${owner}/${repo}/admin/compact`,
241 {
242 method: "POST",
243 headers: {
244 Cookie: seededRepo.cookieHeader,
245 "Content-Type": "application/json",
246 Origin: "https://example.com",
247 },
248 body: JSON.stringify({}),
249 }
250 );
251 expect(previewResponse.status).toBe(200);
252 const previewJson = (await previewResponse.json()) as {
253 action?: string;
254 status?: string;
255 queued?: boolean;
256 plan?: { sourcePacks?: unknown[] };
257 };
258 expect(previewJson.action).toBe("preview");
259 expect(previewJson.status).toBe("ok");
260 expect(previewJson.queued).toBe(false);
261 expect(previewJson.plan?.sourcePacks?.length).toBe(4);
262 } finally {
263 sendSpy.mockRestore();
264 }
265 });
266 
267 it("compacts superseded packs and keeps fetch and raw reads correct without loose objects", async () => {
268 const owner = "o";
269 const repo = uniqueRepoId("stream-compaction-run");
270 const seededRepo = await setupRepoForTests(env, owner, repo);
271 const repoId = `${owner}/${repo}`;
272 const seeded = await seedPackFirstRepo(repoId);
273 const stub = getRepoStub(env, repoId);
274 await promoteToStreaming(owner, repo);
275 
276 const pushed = await pushOverflowingStreamingHistory({
277 owner,
278 repo,
279 repoId,
280 startingCommitOid: seeded.nextCommit.oid,
281 updates: 4,
282 });
283 
284 await deleteLooseObjectCopies(env, seeded.getStub, [
285 ...seeded.objectOids,
286 ...pushed.objectOids,
287 ]);
288 
289 const compacted = await compactOnce(repoId);
290 expect(compacted.acked).toBe(true);
291 expect(compacted.retried).toBe(false);
292 
293 const stateAfterCompaction = await getDebugState(owner, repo, seededRepo.cookieHeader);
294 expect(stateAfterCompaction.compaction?.queued).toBe(false);
295 expect(stateAfterCompaction.activePacks?.some((pack) => pack.kind === "compact")).toBe(true);
296 expect(stateAfterCompaction.supersededPacks?.length).toBe(4);
297 
298 const supersededPackKeys = (stateAfterCompaction.supersededPacks || []).map((pack) => pack.key);
299 const beforeDelete = await collectPackObjects(supersededPackKeys);
300 expect(beforeDelete.every((entry) => entry.exists && entry.idxExists && entry.refsExists)).toBe(
301 true
302 );
303 
304 const rawResponse = await workerExports.default.fetch(
305 `https://example.com/${owner}/${repo}/raw?oid=${pushed.objectOids.at(-3)}&name=README.md`
306 );
307 expect(rawResponse.status).toBe(200);
308 expect(await rawResponse.text()).toBe("streaming update 3\n");
309 
310 const fetchResponse = await workerExports.default.fetch(
311 `https://example.com/${owner}/${repo}/git-upload-pack`,
312 {
313 method: "POST",
314 headers: {
315 "Content-Type": "application/x-git-upload-pack-request",
316 "Git-Protocol": "version=2",
317 },
318 body: buildFetchBody({
319 wants: [pushed.currentCommitOid],
320 haves: [seeded.nextCommit.oid],
321 done: true,
322 }),
323 } as any
324 );
325 expect(fetchResponse.status).toBe(200);
326 expect(
327 decodeReportStatus(new Uint8Array(await fetchResponse.arrayBuffer())).length
328 ).toBeGreaterThan(0);
329 
330 const deleted = await deleteSupersededOnce(repoId, supersededPackKeys);
331 expect(deleted.acked).toBe(true);
332 expect(deleted.retried).toBe(false);
333 
334 const afterDelete = await collectPackObjects(supersededPackKeys);
335 expect(
336 afterDelete.every((entry) => !entry.exists && !entry.idxExists && !entry.refsExists)
337 ).toBe(true);
338 
339 // The superseded catalog rows remain visible for admin/debug until explicit cleanup.
340 const finalState = await getDebugState(owner, repo, seededRepo.cookieHeader);
341 expect(finalState.supersededPacks?.length).toBe(4);
342 
343 void stub;
344 });
345 
346 it("skips an oversized base pack when compacting receive-pack overflow", async () => {
347 const owner = "o";
348 const repo = uniqueRepoId("stream-compaction-large-base");
349 const seededRepo = await setupRepoForTests(env, owner, repo);
350 const repoId = `${owner}/${repo}`;
351 const seeded = await seedPackFirstRepo(repoId);
352 await promoteToStreaming(owner, repo);
353 
354 const pushed = await pushOverflowingStreamingHistory({
355 owner,
356 repo,
357 repoId,
358 startingCommitOid: seeded.nextCommit.oid,
359 updates: 4,
360 });
361 
362 const basePackKey = seeded.packKeys[0];
363 if (!basePackKey) throw new Error("missing seeded base pack");
364 await runDOWithRetry(seeded.getStub, async (_instance, state) => {
365 const db = getDb(state.storage);
366 const baseRow = await getPackCatalogRow(db, basePackKey);
367 if (!baseRow) throw new Error("missing seeded base catalog row");
368 await upsertPackCatalogRow(db, {
369 ...baseRow,
370 objectCount: AUTO_COMPACTION_MAX_SOURCE_OBJECTS + 1,
371 });
372 });
373 
374 const stub = getRepoStub(env, repoId);
375 const preview = await stub.previewCompaction();
376 expect(preview.status).toBe("ok");
377 expect(preview.plan?.sourcePacks.map((pack) => pack.seqLo)).toEqual([2, 3, 4, 5]);
378 
379 const compacted = await compactOnce(repoId);
380 expect(compacted.acked).toBe(true);
381 expect(compacted.retried).toBe(false);
382 
383 const stateAfterCompaction = await getDebugState(owner, repo, seededRepo.cookieHeader);
384 expect(stateAfterCompaction.compaction?.queued).toBe(false);
385 expect(stateAfterCompaction.activePacks?.some((pack) => pack.key === basePackKey)).toBe(true);
386 expect(stateAfterCompaction.supersededPacks?.some((pack) => pack.key === basePackKey)).toBe(
387 false
388 );
389 expect(stateAfterCompaction.supersededPacks?.length).toBe(4);
390 
391 const rawResponse = await workerExports.default.fetch(
392 `https://example.com/${owner}/${repo}/raw?oid=${pushed.objectOids.at(-3)}&name=README.md`
393 );
394 expect(rawResponse.status).toBe(200);
395 expect(await rawResponse.text()).toBe("streaming update 3\n");
396 });
397 
398 it("acks and clears queued compaction when every safe window is over budget", async () => {
399 const owner = "o";
400 const repo = uniqueRepoId("stream-compaction-blocked-budget");
401 const seededRepo = await setupRepoForTests(env, owner, repo);
402 const repoId = `${owner}/${repo}`;
403 const seeded = await seedPackFirstRepo(repoId);
404 await promoteToStreaming(owner, repo);
405 
406 await pushOverflowingStreamingHistory({
407 owner,
408 repo,
409 repoId,
410 startingCommitOid: seeded.nextCommit.oid,
411 updates: 4,
412 });
413 
414 await runDOWithRetry(seeded.getStub, async (_instance, state) => {
415 const db = getDb(state.storage);
416 const activeRows = await listActivePackCatalog(db);
417 const overBudgetPerPack = Math.floor(AUTO_COMPACTION_MAX_SOURCE_OBJECTS / 4) + 1;
418 for (const row of activeRows) {
419 await upsertPackCatalogRow(db, {
420 ...row,
421 objectCount: overBudgetPerPack,
422 });
423 }
424 });
425 
426 const compacted = await compactOnce(repoId);
427 expect(compacted.acked).toBe(true);
428 expect(compacted.retried).toBe(false);
429 
430 const stateAfterCompaction = await getDebugState(owner, repo, seededRepo.cookieHeader);
431 expect(stateAfterCompaction.compaction?.queued).toBe(false);
432 expect(stateAfterCompaction.activePacks?.length).toBe(5);
433 expect(stateAfterCompaction.supersededPacks?.length ?? 0).toBe(0);
434 });
435 
436 it("admin remove deletes pack, idx, and ref sidecar for a superseded pack", async () => {
437 const owner = "o";
438 const repo = uniqueRepoId("stream-compaction-admin-remove");
439 const seededRepo = await setupRepoForTests(env, owner, repo);
440 const repoId = `${owner}/${repo}`;
441 const seeded = await seedPackFirstRepo(repoId);
442 await promoteToStreaming(owner, repo);
443 
444 await pushOverflowingStreamingHistory({
445 owner,
446 repo,
447 repoId,
448 startingCommitOid: seeded.nextCommit.oid,
449 updates: 4,
450 });
451 
452 const compacted = await compactOnce(repoId);
453 expect(compacted.acked).toBe(true);
454 expect(compacted.retried).toBe(false);
455 
456 const stateAfterCompaction = await getDebugState(owner, repo, seededRepo.cookieHeader);
457 const supersededPackKey = stateAfterCompaction.supersededPacks?.[0]?.key;
458 if (!supersededPackKey) throw new Error("missing superseded pack");
459 
460 await expect(env.REPO_BUCKET.head(supersededPackKey)).resolves.toBeTruthy();
461 await expect(env.REPO_BUCKET.head(packIndexKey(supersededPackKey))).resolves.toBeTruthy();
462 await expect(env.REPO_BUCKET.head(packRefsKey(supersededPackKey))).resolves.toBeTruthy();
463 
464 const packName = supersededPackKey.split("/").pop();
465 if (!packName) throw new Error("missing pack name");
466 
467 const deleteResponse = await workerExports.default.fetch(
468 `https://example.com/${owner}/${repo}/admin/pack/${encodeURIComponent(packName)}`,
469 {
470 method: "DELETE",
471 headers: { Cookie: seededRepo.cookieHeader, Origin: "https://example.com" },
472 }
473 );
474 expect(deleteResponse.status).toBe(200);
475 const deleteJson = (await deleteResponse.json()) as {
476 ok?: boolean;
477 deletedPack?: boolean;
478 deletedIndex?: boolean;
479 deletedRefs?: boolean;
480 deletedMetadata?: boolean;
481 packState?: string;
482 };
483 expect(deleteJson.ok).toBe(true);
484 expect(deleteJson.packState).toBe("superseded");
485 expect(deleteJson.deletedPack).toBe(true);
486 expect(deleteJson.deletedIndex).toBe(true);
487 expect(deleteJson.deletedRefs).toBe(true);
488 expect(deleteJson.deletedMetadata).toBe(true);
489 
490 await expect(env.REPO_BUCKET.head(supersededPackKey)).resolves.toBeNull();
491 await expect(env.REPO_BUCKET.head(packIndexKey(supersededPackKey))).resolves.toBeNull();
492 await expect(env.REPO_BUCKET.head(packRefsKey(supersededPackKey))).resolves.toBeNull();
493 });
494 
495 it("returns retry when a receive lease appears before compaction commit", async () => {
496 const owner = "o";
497 const repo = uniqueRepoId("stream-compaction-receive-priority");
498 const seededRepo = await setupRepoForTests(env, owner, repo);
499 const repoId = `${owner}/${repo}`;
500 const seeded = await seedPackFirstRepo(repoId);
501 const getStub = () => env.REPO_DO.get(env.REPO_DO.idFromName(repoId));
502 await promoteToStreaming(owner, repo);
503 
504 await pushOverflowingStreamingHistory({
505 owner,
506 repo,
507 repoId,
508 startingCommitOid: seeded.nextCommit.oid,
509 updates: 4,
510 });
511 
512 const stub = getStub();
513 const begin = await stub.beginCompaction();
514 expect(begin.ok).toBe(true);
515 if (!begin.ok) {
516 throw new Error("expected compaction to begin");
517 }
518 
519 await runDOWithRetry(getStub, async (_instance, state) => {
520 const store = asTypedStorage<RepoStateSchema>(state.storage);
521 const now = Date.now();
522 await store.put("receiveLease", {
523 token: "receive-priority",
524 createdAt: now,
525 expiresAt: now + 60_000,
526 });
527 });
528 
529 const result = await stub.commitCompaction({
530 token: begin.lease.token,
531 sourcePacks: begin.sourcePacks,
532 targetTier: begin.targetTier,
533 packsetVersion: begin.packsetVersion,
534 stagedPack: {
535 packKey: `${begin.sourcePacks[0]!.packKey}.fake-compaction`,
536 packBytes: begin.sourcePacks[0]!.packBytes,
537 idxBytes: begin.sourcePacks[0]!.idxBytes,
538 objectCount: begin.sourcePacks[0]!.objectCount,
539 },
540 });
541 expect(result.status).toBe("retry");
542 if (result.status === "retry") {
543 expect(result.reason).toBe("receive-active");
544 }
545 
546 const state = await getDebugState(owner, repo, seededRepo.cookieHeader);
547 expect(state.activePacks?.every((pack) => pack.kind !== "compact")).toBe(true);
548 });
549 
550 it("returns retry when packsetVersion changes before compaction commit", async () => {
551 const owner = "o";
552 const repo = uniqueRepoId("stream-compaction-packset-changed");
553 const seededRepo = await setupRepoForTests(env, owner, repo);
554 const repoId = `${owner}/${repo}`;
555 const seeded = await seedPackFirstRepo(repoId);
556 const getStub = () => env.REPO_DO.get(env.REPO_DO.idFromName(repoId));
557 await promoteToStreaming(owner, repo);
558 
559 await pushOverflowingStreamingHistory({
560 owner,
561 repo,
562 repoId,
563 startingCommitOid: seeded.nextCommit.oid,
564 updates: 4,
565 });
566 
567 const stub = getStub();
568 const begin = await stub.beginCompaction();
569 expect(begin.ok).toBe(true);
570 if (!begin.ok) {
571 throw new Error("expected compaction to begin");
572 }
573 
574 // Bump the packset version behind the compaction lease's back.
575 await runDOWithRetry(getStub, async (_instance, state) => {
576 const store = asTypedStorage<RepoStateSchema>(state.storage);
577 const current = (await store.get("packsetVersion")) || 0;
578 await store.put("packsetVersion", current + 1);
579 });
580 
581 const result = await stub.commitCompaction({
582 token: begin.lease.token,
583 sourcePacks: begin.sourcePacks,
584 targetTier: begin.targetTier,
585 packsetVersion: begin.packsetVersion,
586 stagedPack: {
587 packKey: `${begin.sourcePacks[0]!.packKey}.fake-compaction`,
588 packBytes: begin.sourcePacks[0]!.packBytes,
589 idxBytes: begin.sourcePacks[0]!.idxBytes,
590 objectCount: begin.sourcePacks[0]!.objectCount,
591 },
592 });
593 expect(result.status).toBe("retry");
594 if (result.status === "retry") {
595 expect(result.reason).toBe("packset-changed");
596 }
597 
598 const state = await getDebugState(owner, repo, seededRepo.cookieHeader);
599 expect(state.activePacks?.every((pack) => pack.kind !== "compact")).toBe(true);
600 });
601 
602 it("compacts successfully when a non-source pack contains a duplicate identity REF_DELTA", async () => {
603 // Regression test for the compaction self-referential delta bug.
604 //
605 // When the active catalog snapshot is newest-first and a newer non-source
606 // pack contains a REF_DELTA whose resolved OID equals its baseOid (an
607 // identity delta), resolveOrderedEntryByOid picks that entry for a needed
608 // OID and the base chase loops back to the same entry, creating a cycle
609 // that the topology sort cannot order.
610 //
611 // The fix reorders the compaction snapshot so source packs are searched
612 // first, ensuring the authoritative full-object entry is selected.
613 
614 const owner = "o";
615 const repo = uniqueRepoId("stream-compaction-self-ref");
616 const seededRepo = await setupRepoForTests(env, owner, repo);
617 const repoId = `${owner}/${repo}`;
618 const id = env.REPO_DO.idFromName(repoId);
619 const getStub = () => env.REPO_DO.get(id);
620 
621 const author = "You <you@example.com> 0 +0000";
622 
623 // Shared blob that will appear in both a source pack and the newer
624 // non-source pack as an identity REF_DELTA.
625 const sharedBlobPayload = new TextEncoder().encode("shared content\n");
626 const sharedBlob = await encodeGitObject("blob", sharedBlobPayload);
627 
628 // Build 4 source packs (oldest) with distinct commits, one containing the
629 // shared blob. Each pack needs at least one unique object so the OID sets
630 // differ and compaction has real work to do.
631 const sourcePacks: Array<{ name: string; packBytes: Uint8Array }> = [];
632 let parentOid: string | undefined;
633 const allObjectOids: string[] = [];
634 
635 for (let i = 0; i < 4; i++) {
636 const blobPayload = new TextEncoder().encode(`source content ${i}\n`);
637 const blob = await encodeGitObject("blob", blobPayload);
638 
639 // Include the shared blob in the first source pack so it is a needed OID.
640 const treeEntries = [{ mode: "100644" as const, name: `file-${i}.txt`, oid: blob.oid }];
641 if (i === 0) {
642 treeEntries.push({ mode: "100644" as const, name: "shared.txt", oid: sharedBlob.oid });
643 }
644 const treePayload = buildTreePayload(treeEntries);
645 const tree = await encodeGitObject("tree", treePayload);
646 
647 const commitText =
648 `tree ${tree.oid}\n` +
649 (parentOid ? `parent ${parentOid}\n` : "") +
650 `author ${author}\ncommitter ${author}\n\nsource ${i}\n`;
651 const commit = await encodeGitObject("commit", new TextEncoder().encode(commitText));
652 parentOid = commit.oid;
653 
654 const objects: Array<{ type: "blob" | "tree" | "commit"; payload: Uint8Array }> = [
655 { type: "blob", payload: blobPayload },
656 { type: "tree", payload: treePayload },
657 { type: "commit", payload: new TextEncoder().encode(commitText) },
658 ];
659 if (i === 0) {
660 objects.unshift({ type: "blob", payload: sharedBlobPayload });
661 }
662 
663 sourcePacks.push({ name: `pack-source-${i}.pack`, packBytes: await buildPack(objects) });
664 allObjectOids.push(blob.oid, tree.oid, commit.oid);
665 if (i === 0) allObjectOids.push(sharedBlob.oid);
666 }
667 
668 // Build the newest pack with an identity REF_DELTA for the shared blob.
669 // The delta copies the entire base content, so the resolved OID equals
670 // the base OID — exactly the scenario that triggers the self-loop.
671 const identityDelta = buildAppendOnlyDelta(sharedBlobPayload, new Uint8Array(0));
672 const newestPackBytes = await buildPack([
673 { type: "ref-delta", baseOid: sharedBlob.oid, delta: identityDelta },
674 ]);
675 const newestPack = { name: "pack-newest-dup.pack", packBytes: newestPackBytes };
676 
677 // seedPackedRepoState expects packs newest-first; it indexes oldest-first
678 // internally so the REF_DELTA base in the source packs is available.
679 const lastCommitOid = parentOid!;
680 await seedPackedRepoState({
681 env,
682 repoId,
683 getStub,
684 packs: [newestPack, ...sourcePacks],
685 refs: [{ name: "refs/heads/main", oid: lastCommitOid }],
686 head: { target: "refs/heads/main", oid: lastCommitOid },
687 });
688 
689 // Verify compaction plan selects the 4 source packs (oldest tier-0).
690 const preState = await getDebugState(owner, repo, seededRepo.cookieHeader);
691 expect(preState.activePacks?.length).toBe(5);
692 expect(preState.activePacks?.filter((p) => p.tier === 0).length).toBe(5);
693 
694 // Request and run compaction.
695 const stub = getStub();
696 const request = await stub.requestCompaction();
697 expect(request.status).toBe("queued");
698 
699 const result = await compactOnce(repoId);
700 expect(result.acked).toBe(true);
701 expect(result.retried).toBe(false);
702 
703 // Verify post-compaction state: one compacted pack, 4 source packs superseded.
704 const postState = await getDebugState(owner, repo, seededRepo.cookieHeader);
705 expect(postState.activePacks?.some((p) => p.kind === "compact")).toBe(true);
706 expect(postState.supersededPacks?.length).toBe(4);
707 expect(postState.compaction?.queued).toBe(false);
708 });
709 
710 it("fetch after compaction does not produce duplicate REF_DELTA entries", async () => {
711 // Regression test for broken `git pull` after compaction.
712 //
713 // After compaction, the active snapshot is newest-first:
714 // [compacted pack (newest seqHi), non-source pack (older)]
715 // If the non-source pack has an OFS_DELTA whose in-pack base chain
716 // reaches an entry with the same OID as a full object in the compacted
717 // pack, the OFS_DELTA base chase adds that entry by pack-local offset —
718 // bypassing OID-level dedup. The output pack then has two entries for
719 // the same OID, causing git's index-pack to reject the pack with
720 // "REF_DELTA at offset X already resolved (duplicate base Y)".
721 //
722 // The fix canonicalizes OFS_DELTA bases via resolveOrderedEntryByOid so
723 // the position-based dedup in addEntry catches cross-pack duplicates.
724 
725 const owner = "o";
726 const repo = uniqueRepoId("stream-compaction-fetch-dup");
727 const seededRepo = await setupRepoForTests(env, owner, repo);
728 const repoId = `${owner}/${repo}`;
729 const id = env.REPO_DO.idFromName(repoId);
730 const getStub = () => env.REPO_DO.get(id);
731 
732 const author = "You <you@example.com> 0 +0000";
733 
734 // Shared blob: appears as a full object in a source pack, and as an
735 // identity REF_DELTA in the non-source pack. After compaction, both
736 // the compacted pack and the non-source pack contain this OID.
737 const sharedBlobPayload = new TextEncoder().encode("shared content for fetch test\n");
738 const sharedBlob = await encodeGitObject("blob", sharedBlobPayload);
739 
740 // Child blob: stored as OFS_DELTA based on the shared blob in the
741 // non-source pack. Its base chase is the path that triggers the
742 // cross-pack duplicate. Must be referenced by a tree to be needed.
743 const childBlobFullPayload = new Uint8Array([
744 ...sharedBlobPayload,
745 ...new TextEncoder().encode("extra\n"),
746 ]);
747 const childBlob = await encodeGitObject("blob", childBlobFullPayload);
748 
749 // Build 4 source packs with distinct commits. Source pack 0 includes
750 // the shared blob as a full object.
751 const sourcePacks: Array<{ name: string; packBytes: Uint8Array }> = [];
752 let parentOid: string | undefined;
753 
754 for (let i = 0; i < 4; i++) {
755 const blobPayload = new TextEncoder().encode(`source content fetch ${i}\n`);
756 const blob = await encodeGitObject("blob", blobPayload);
757 
758 const treeEntries = [{ mode: "100644" as const, name: `file-${i}.txt`, oid: blob.oid }];
759 if (i === 0) {
760 treeEntries.push({ mode: "100644" as const, name: "shared.txt", oid: sharedBlob.oid });
761 }
762 const treePayload = buildTreePayload(treeEntries);
763 const tree = await encodeGitObject("tree", treePayload);
764 
765 const commitText =
766 `tree ${tree.oid}\n` +
767 (parentOid ? `parent ${parentOid}\n` : "") +
768 `author ${author}\ncommitter ${author}\n\nsource fetch ${i}\n`;
769 const commit = await encodeGitObject("commit", new TextEncoder().encode(commitText));
770 parentOid = commit.oid;
771 
772 const objects: Array<{ type: "blob" | "tree" | "commit"; payload: Uint8Array }> = [
773 { type: "blob", payload: blobPayload },
774 { type: "tree", payload: treePayload },
775 { type: "commit", payload: new TextEncoder().encode(commitText) },
776 ];
777 if (i === 0) {
778 objects.unshift({ type: "blob", payload: sharedBlobPayload });
779 }
780 
781 sourcePacks.push({ name: `pack-src-${i}.pack`, packBytes: await buildPack(objects) });
782 }
783 
784 // Build a newest pack that contributes a real commit to the graph.
785 // Layout (by entry index):
786 // 0: identity REF_DELTA for sharedBlob (resolves to sharedBlob.oid)
787 // 1: OFS_DELTA child blob based on entry 0 (resolves to childBlob.oid)
788 // 2: tree referencing the child blob
789 // 3: commit (child of last source commit) referencing tree 2
790 //
791 // During fetch, the child blob (entry 1) is needed because the tree
792 // references it. Its OFS_DELTA base chase reaches entry 0 (shared blob
793 // OID), which was already selected from the compacted pack — the
794 // cross-pack duplicate that this test guards against.
795 const identityDelta = buildAppendOnlyDelta(sharedBlobPayload, new Uint8Array(0));
796 const childDelta = buildAppendOnlyDelta(sharedBlobPayload, new TextEncoder().encode("extra\n"));
797 
798 const newestTreePayload = buildTreePayload([
799 { mode: "100644" as const, name: "child.txt", oid: childBlob.oid },
800 { mode: "100644" as const, name: "shared.txt", oid: sharedBlob.oid },
801 ]);
802 const newestTree = await encodeGitObject("tree", newestTreePayload);
803 
804 const newestCommitText =
805 `tree ${newestTree.oid}\n` +
806 `parent ${parentOid}\n` +
807 `author ${author}\ncommitter ${author}\n\nnewest with child blob\n`;
808 const newestCommit = await encodeGitObject(
809 "commit",
810 new TextEncoder().encode(newestCommitText)
811 );
812 
813 const newestPackBytes = await buildPack([
814 { type: "ref-delta", baseOid: sharedBlob.oid, delta: identityDelta },
815 { type: "ofs-delta", baseIndex: 0, delta: childDelta },
816 { type: "tree", payload: newestTreePayload },
817 { type: "commit", payload: new TextEncoder().encode(newestCommitText) },
818 ]);
819 const newestPack = { name: "pack-newest-fetch-dup.pack", packBytes: newestPackBytes };
820 
821 // HEAD = newest commit so the fetch graph reaches objects in every pack.
822 await seedPackedRepoState({
823 env,
824 repoId,
825 getStub,
826 packs: [newestPack, ...sourcePacks],
827 refs: [{ name: "refs/heads/main", oid: newestCommit.oid }],
828 head: { target: "refs/heads/main", oid: newestCommit.oid },
829 });
830 
831 // Run compaction — merges the 4 source packs into one compacted pack.
832 // The non-source pack (newest) stays active.
833 const stub = getStub();
834 await stub.requestCompaction();
835 const compactResult = await compactOnce(repoId);
836 expect(compactResult.acked).toBe(true);
837 expect(compactResult.retried).toBe(false);
838 
839 const postState = await getDebugState(owner, repo, seededRepo.cookieHeader);
840 expect(postState.activePacks?.some((p) => p.kind === "compact")).toBe(true);
841 
842 // Fetch all objects (clone scenario) — this is the path that was broken.
843 const fetchResponse = await workerExports.default.fetch(
844 `https://example.com/${owner}/${repo}/git-upload-pack`,
845 {
846 method: "POST",
847 headers: {
848 "Content-Type": "application/x-git-upload-pack-request",
849 "Git-Protocol": "version=2",
850 },
851 body: buildFetchBody({ wants: [newestCommit.oid], done: true }),
852 } as any
853 );
854 expect(fetchResponse.status).toBe(200);
855 
856 // Extract sideband-encoded pack bytes from the response.
857 const bytes = new Uint8Array(await fetchResponse.arrayBuffer());
858 const lines = decodePktLines(bytes);
859 const packChunks: Uint8Array[] = [];
860 let inPackfile = false;
861 for (const line of lines) {
862 if (line.type === "line" && line.text === "packfile\n") {
863 inPackfile = true;
864 continue;
865 }
866 if (inPackfile && line.type === "line" && line.raw && line.raw[0] === 0x01) {
867 packChunks.push(line.raw.subarray(1));
868 }
869 }
870 const packOut = concatChunks(packChunks);
871 
872 // Basic pack header sanity check.
873 expect(new TextDecoder().decode(packOut.subarray(0, 4))).toBe("PACK");
874 
875 // Index the returned pack and verify no duplicate OIDs.
876 const verifyKey = `verify/compaction-fetch-dup-${Date.now()}.pack`;
877 await env.REPO_BUCKET.put(verifyKey, packOut);
878 const verifyResult = await indexTestPack(env, verifyKey, packOut.byteLength);
879 
880 const oidSet = new Set<string>();
881 for (let i = 0; i < verifyResult.idxView.count; i++) {
882 const oidBytes = verifyResult.idxView.rawNames.subarray(i * 20, (i + 1) * 20);
883 oidSet.add(bytesToHex(oidBytes));
884 }
885 // Intentionally stricter than git index-pack (which tolerates duplicate
886 // full objects and OFS identity deltas, but rejects duplicate REF_DELTAs).
887 // Our rewrite should never produce ANY duplicate OIDs in the output pack.
888 expect(oidSet.size).toBe(verifyResult.idxView.count);
889 
890 // Also verify after superseded pack deletion — compacted pack is sole
891 // source for the merged objects, no duplicates possible.
892 const supersededKeys =
893 postState.supersededPacks?.map((p) => p.key).filter((k): k is string => !!k) ?? [];
894 if (supersededKeys.length > 0) {
895 await deleteSupersededOnce(repoId, supersededKeys);
896 }
897 
898 const fetchResponse2 = await workerExports.default.fetch(
899 `https://example.com/${owner}/${repo}/git-upload-pack`,
900 {
901 method: "POST",
902 headers: {
903 "Content-Type": "application/x-git-upload-pack-request",
904 "Git-Protocol": "version=2",
905 },
906 body: buildFetchBody({ wants: [newestCommit.oid], done: true }),
907 } as any
908 );
909 expect(fetchResponse2.status).toBe(200);
910 
911 const bytes2 = new Uint8Array(await fetchResponse2.arrayBuffer());
912 const lines2 = decodePktLines(bytes2);
913 const packChunks2: Uint8Array[] = [];
914 let inPackfile2 = false;
915 for (const line of lines2) {
916 if (line.type === "line" && line.text === "packfile\n") {
917 inPackfile2 = true;
918 continue;
919 }
920 if (inPackfile2 && line.type === "line" && line.raw && line.raw[0] === 0x01) {
921 packChunks2.push(line.raw.subarray(1));
922 }
923 }
924 const packOut2 = concatChunks(packChunks2);
925 expect(new TextDecoder().decode(packOut2.subarray(0, 4))).toBe("PACK");
926 
927 const verifyKey2 = `verify/compaction-fetch-post-delete-${Date.now()}.pack`;
928 await env.REPO_BUCKET.put(verifyKey2, packOut2);
929 const verifyResult2 = await indexTestPack(env, verifyKey2, packOut2.byteLength);
930 
931 const oidSet2 = new Set<string>();
932 for (let i = 0; i < verifyResult2.idxView.count; i++) {
933 const oidBytes = verifyResult2.idxView.rawNames.subarray(i * 20, (i + 1) * 20);
934 oidSet2.add(bytesToHex(oidBytes));
935 }
936 expect(oidSet2.size).toBe(verifyResult2.idxView.count);
937 });
938 
939 it("keeps active pack counts bounded after repeated pushes and compaction drains", async () => {
940 const owner = "o";
941 const repo = uniqueRepoId("stream-compaction-bounded");
942 const seededRepo = await setupRepoForTests(env, owner, repo);
943 const repoId = `${owner}/${repo}`;
944 const seeded = await seedPackFirstRepo(repoId);
945 await promoteToStreaming(owner, repo);
946 
947 await pushOverflowingStreamingHistory({
948 owner,
949 repo,
950 repoId,
951 startingCommitOid: seeded.nextCommit.oid,
952 updates: 8,
953 });
954 
955 for (let attempt = 0; attempt < 6; attempt++) {
956 const queuedState = await getDebugState(owner, repo, seededRepo.cookieHeader);
957 if (!queuedState.compaction?.queued) break;
958 const result = await compactOnce(repoId);
959 expect(result.acked || result.retried).toBe(true);
960 }
961 
962 const finalState = await getDebugState(owner, repo, seededRepo.cookieHeader);
963 const counts = new Map<number, number>();
964 for (const pack of finalState.activePacks || []) {
965 counts.set(pack.tier, (counts.get(pack.tier) || 0) + 1);
966 }
967 for (const count of counts.values()) {
968 expect(count).toBeLessThanOrEqual(4);
969 }
970 });
971});