Skip to content
File

Blob: test/pack-indexer.resolve.ofs.worker.test.ts

typescript441 lines
1import { describe, expect, it } from "vitest";
2import { env } from "cloudflare:workers";
3import { buildPack, buildAppendOnlyDelta } from "./util/git-pack";
4import { uniqueRepoId } from "./util/test-helpers";
5import {
6 makeCountSubrequest,
7 makeLimiter,
8 makeTracingLimiter,
9 packIndexerLog as log,
10} from "./util/pack-indexer.helpers";
11 
12import {
13 allocateEntryTable,
14 isResolveAbortedError,
15 scanPack,
16 resolveDeltasAndWriteIdx,
17} from "@/worker/git/pack/indexer";
18import { computeOid } from "@/worker/git/core/objects";
19import { bytesToHex } from "@/worker/common/hex";
20import { DEFAULT_SUBREQUEST_BUDGET } from "@/worker/git/operations/limits";
21import { getOidHexAt, parseIdxView } from "@/worker/git/object-store";
22import { packIndexKey } from "@/worker/keys";
23import { getBasePayload } from "@/worker/git/pack/indexer/resolve/materialize";
24import { PayloadLRU } from "@/worker/git/pack/indexer/resolve/payloadCache";
25import { SequentialReader } from "@/worker/git/pack/indexer/resolve/reader";
26 
27async function expectResolveAborted(promise: Promise<unknown>): Promise<void> {
28 try {
29 await promise;
30 } catch (error) {
31 expect(isResolveAbortedError(error)).toBe(true);
32 return;
33 }
34 throw new Error("expected resolve to abort");
35}
36 
37describe("resolveDeltasAndWriteIdx OFS_DELTA", () => {
38 it("resolves OFS_DELTA and writes valid idx", async () => {
39 const baseBlobPayload = new TextEncoder().encode("base content here\n");
40 const suffix = new TextEncoder().encode("appended text\n");
41 const delta = buildAppendOnlyDelta(baseBlobPayload, suffix);
42 
43 const expectedPayload = new Uint8Array(baseBlobPayload.length + suffix.length);
44 expectedPayload.set(baseBlobPayload, 0);
45 expectedPayload.set(suffix, baseBlobPayload.length);
46 const expectedOid = await computeOid("blob", expectedPayload);
47 
48 const packBytes = await buildPack([
49 { type: "blob", payload: baseBlobPayload },
50 { type: "ofs-delta", baseIndex: 0, delta },
51 ]);
52 
53 const packKey = "test/resolve-ofs.pack";
54 await env.REPO_BUCKET.put(packKey, packBytes);
55 const head = await env.REPO_BUCKET.head(packKey);
56 
57 const scanResult = await scanPack({
58 env,
59 packKey,
60 packSize: head!.size,
61 limiter: makeLimiter(),
62 countSubrequest: () => {},
63 log,
64 });
65 
66 expect(scanResult.objectCount).toBe(2);
67 expect(scanResult.table.resolved[0]).toBe(1);
68 expect(scanResult.table.resolved[1]).toBe(0);
69 
70 const repoId = uniqueRepoId();
71 const resolveResult = await resolveDeltasAndWriteIdx({
72 env,
73 packKey,
74 packSize: head!.size,
75 limiter: makeLimiter(),
76 countSubrequest: () => {},
77 log,
78 scanResult,
79 repoId,
80 });
81 
82 expect(scanResult.table.resolved[1]).toBe(1);
83 expect(bytesToHex(scanResult.table.oids.subarray(20, 40))).toBe(expectedOid);
84 
85 const idxObj = await env.REPO_BUCKET.get(packIndexKey(packKey));
86 expect(idxObj).not.toBeNull();
87 
88 const idxBuf = new Uint8Array(await idxObj!.arrayBuffer());
89 const idxView = parseIdxView(packKey, idxBuf, head!.size);
90 expect(idxView).not.toBeUndefined();
91 expect(idxView!.count).toBe(2);
92 expect([getOidHexAt(idxView!, 0), getOidHexAt(idxView!, 1)]).toContain(expectedOid);
93 
94 expect(resolveResult.idxView.count).toBe(2);
95 expect(resolveResult.objectCount).toBe(2);
96 });
97 
98 it("rejects an already-aborted resolve before writing idx", async () => {
99 const baseBlobPayload = new TextEncoder().encode("abort base\n");
100 const suffix = new TextEncoder().encode("abort tail\n");
101 const packKey = "test/resolve-aborted-before-start.pack";
102 const packBytes = await buildPack([
103 { type: "blob", payload: baseBlobPayload },
104 { type: "ofs-delta", baseIndex: 0, delta: buildAppendOnlyDelta(baseBlobPayload, suffix) },
105 ]);
106 await env.REPO_BUCKET.put(packKey, packBytes);
107 const head = await env.REPO_BUCKET.head(packKey);
108 
109 const scanResult = await scanPack({
110 env,
111 packKey,
112 packSize: head!.size,
113 limiter: makeLimiter(),
114 countSubrequest: () => {},
115 log,
116 });
117 
118 const abortController = new AbortController();
119 abortController.abort();
120 
121 await expectResolveAborted(
122 resolveDeltasAndWriteIdx({
123 env,
124 packKey,
125 packSize: head!.size,
126 limiter: makeLimiter(),
127 countSubrequest: () => {},
128 log,
129 scanResult,
130 repoId: uniqueRepoId(),
131 signal: abortController.signal,
132 })
133 );
134 
135 expect(await env.REPO_BUCKET.get(packIndexKey(packKey))).toBeNull();
136 });
137 
138 it("re-materializes evicted bases when the LRU budget is tiny", async () => {
139 const baseBlobPayload = new TextEncoder().encode("base\n");
140 const midSuffix = new TextEncoder().encode("mid\n");
141 const finalSuffix = new TextEncoder().encode("final\n");
142 const midDelta = buildAppendOnlyDelta(baseBlobPayload, midSuffix);
143 
144 const midPayload = new Uint8Array(baseBlobPayload.length + midSuffix.length);
145 midPayload.set(baseBlobPayload, 0);
146 midPayload.set(midSuffix, baseBlobPayload.length);
147 const finalDelta = buildAppendOnlyDelta(midPayload, finalSuffix);
148 
149 const finalPayload = new Uint8Array(midPayload.length + finalSuffix.length);
150 finalPayload.set(midPayload, 0);
151 finalPayload.set(finalSuffix, midPayload.length);
152 const finalOid = await computeOid("blob", finalPayload);
153 
154 const packBytes = await buildPack([
155 { type: "blob", payload: baseBlobPayload },
156 { type: "ofs-delta", baseIndex: 0, delta: midDelta },
157 { type: "ofs-delta", baseIndex: 1, delta: finalDelta },
158 ]);
159 
160 const packKey = "test/resolve-lru-rematerialize.pack";
161 await env.REPO_BUCKET.put(packKey, packBytes);
162 const head = await env.REPO_BUCKET.head(packKey);
163 
164 const scanResult = await scanPack({
165 env,
166 packKey,
167 packSize: head!.size,
168 limiter: makeLimiter(),
169 countSubrequest: () => {},
170 log,
171 });
172 
173 const repoId = uniqueRepoId();
174 await resolveDeltasAndWriteIdx({
175 env,
176 packKey,
177 packSize: head!.size,
178 limiter: makeLimiter(),
179 countSubrequest: () => {},
180 log,
181 scanResult,
182 repoId,
183 lruBudget: 1,
184 });
185 
186 expect(bytesToHex(scanResult.table.oids.subarray(40, 60))).toBe(finalOid);
187 });
188 
189 it("routes idx writes through the shared limiter and subrequest counter", async () => {
190 const baseBlobPayload = new TextEncoder().encode("budget test\n");
191 const packBytes = await buildPack([{ type: "blob", payload: baseBlobPayload }]);
192 
193 const packKey = "test/subreq-budget.pack";
194 await env.REPO_BUCKET.put(packKey, packBytes);
195 const head = await env.REPO_BUCKET.head(packKey);
196 
197 const labels: string[] = [];
198 const limiter = makeTracingLimiter(labels);
199 const counter = { count: 0 };
200 const scanResult = await scanPack({
201 env,
202 packKey,
203 packSize: head!.size,
204 limiter,
205 countSubrequest: makeCountSubrequest(counter),
206 log,
207 });
208 
209 const repoId = uniqueRepoId();
210 await resolveDeltasAndWriteIdx({
211 env,
212 packKey,
213 packSize: head!.size,
214 limiter,
215 countSubrequest: makeCountSubrequest(counter),
216 log,
217 scanResult,
218 repoId,
219 });
220 
221 expect(labels).toContain("r2:get-range");
222 expect(labels).toContain("r2:put-pack-idx");
223 expect(counter.count).toBeGreaterThan(0);
224 expect(counter.count).toBeLessThan(DEFAULT_SUBREQUEST_BUDGET);
225 });
226 
227 it("streams pass-2 inflates in multiple range reads when the resolve chunk size is tiny", async () => {
228 const baseBlobPayload = new Uint8Array(8 * 1024);
229 for (let i = 0; i < baseBlobPayload.length; i++) {
230 baseBlobPayload[i] = (i * 31) & 0xff;
231 }
232 
233 const suffix = new Uint8Array(4 * 1024);
234 for (let i = 0; i < suffix.length; i++) {
235 suffix[i] = (255 - i * 17) & 0xff;
236 }
237 
238 const expectedPayload = new Uint8Array(baseBlobPayload.length + suffix.length);
239 expectedPayload.set(baseBlobPayload, 0);
240 expectedPayload.set(suffix, baseBlobPayload.length);
241 const expectedOid = await computeOid("blob", expectedPayload);
242 
243 const packKey = "test/resolve-streamed-pass2.pack";
244 const packBytes = await buildPack([
245 { type: "blob", payload: baseBlobPayload },
246 { type: "ofs-delta", baseIndex: 0, delta: buildAppendOnlyDelta(baseBlobPayload, suffix) },
247 ]);
248 await env.REPO_BUCKET.put(packKey, packBytes);
249 const head = await env.REPO_BUCKET.head(packKey);
250 
251 const scanResult = await scanPack({
252 env,
253 packKey,
254 packSize: head!.size,
255 limiter: makeLimiter(),
256 countSubrequest: () => {},
257 log,
258 });
259 
260 const labels: string[] = [];
261 const limiter = makeTracingLimiter(labels);
262 const counter = { count: 0 };
263 
264 await resolveDeltasAndWriteIdx({
265 env,
266 packKey,
267 packSize: head!.size,
268 chunkSize: 32,
269 limiter,
270 countSubrequest: makeCountSubrequest(counter),
271 log,
272 scanResult,
273 repoId: uniqueRepoId(),
274 });
275 
276 expect(bytesToHex(scanResult.table.oids.subarray(20, 40))).toBe(expectedOid);
277 expect(labels.filter((label) => label === "r2:get-range").length).toBeGreaterThan(1);
278 expect(counter.count).toBeLessThan(DEFAULT_SUBREQUEST_BUDGET);
279 });
280 
281 it("aborts mid-resolve without writing an idx", async () => {
282 const entries: (
283 | { type: "blob"; payload: Uint8Array }
284 | { type: "ofs-delta"; baseIndex: number; delta: Uint8Array }
285 )[] = [];
286 let currentPayload = new TextEncoder().encode("base\n");
287 entries.push({ type: "blob", payload: currentPayload });
288 
289 for (let i = 0; i < 96; i++) {
290 const suffix = new TextEncoder().encode(String.fromCharCode(97 + (i % 26)));
291 entries.push({
292 type: "ofs-delta",
293 baseIndex: i,
294 delta: buildAppendOnlyDelta(currentPayload, suffix),
295 });
296 
297 const nextPayload = new Uint8Array(currentPayload.length + suffix.length);
298 nextPayload.set(currentPayload, 0);
299 nextPayload.set(suffix, currentPayload.length);
300 currentPayload = nextPayload;
301 }
302 
303 const packKey = "test/resolve-mid-abort.pack";
304 const packBytes = await buildPack(entries);
305 await env.REPO_BUCKET.put(packKey, packBytes);
306 const head = await env.REPO_BUCKET.head(packKey);
307 
308 const scanResult = await scanPack({
309 env,
310 packKey,
311 packSize: head!.size,
312 limiter: makeLimiter(),
313 countSubrequest: () => {},
314 log,
315 });
316 
317 const abortController = new AbortController();
318 const counter = { count: 0 };
319 await expectResolveAborted(
320 resolveDeltasAndWriteIdx({
321 env,
322 packKey,
323 packSize: head!.size,
324 chunkSize: 16,
325 limiter: makeLimiter(),
326 countSubrequest: (n = 1) => {
327 counter.count += n;
328 if (counter.count >= 8 && !abortController.signal.aborted) {
329 abortController.abort();
330 }
331 },
332 log,
333 scanResult,
334 repoId: uniqueRepoId(),
335 lruBudget: 1,
336 signal: abortController.signal,
337 })
338 );
339 
340 expect(counter.count).toBeGreaterThanOrEqual(8);
341 expect(await env.REPO_BUCKET.get(packIndexKey(packKey))).toBeNull();
342 });
343 
344 it("re-materializes a longer OFS chain when the LRU budget is tiny", async () => {
345 const entries: (
346 | { type: "blob"; payload: Uint8Array }
347 | { type: "ofs-delta"; baseIndex: number; delta: Uint8Array }
348 )[] = [];
349 let currentPayload = new TextEncoder().encode("base\n");
350 entries.push({ type: "blob", payload: currentPayload });
351 
352 for (let i = 0; i < 128; i++) {
353 const suffix = new TextEncoder().encode(String.fromCharCode(97 + (i % 26)));
354 entries.push({
355 type: "ofs-delta",
356 baseIndex: i,
357 delta: buildAppendOnlyDelta(currentPayload, suffix),
358 });
359 
360 const nextPayload = new Uint8Array(currentPayload.length + suffix.length);
361 nextPayload.set(currentPayload, 0);
362 nextPayload.set(suffix, currentPayload.length);
363 currentPayload = nextPayload;
364 }
365 
366 const packKey = "test/resolve-long-ofs-rematerialize.pack";
367 const packBytes = await buildPack(entries);
368 await env.REPO_BUCKET.put(packKey, packBytes);
369 const head = await env.REPO_BUCKET.head(packKey);
370 
371 const scanResult = await scanPack({
372 env,
373 packKey,
374 packSize: head!.size,
375 limiter: makeLimiter(),
376 countSubrequest: () => {},
377 log,
378 });
379 
380 await resolveDeltasAndWriteIdx({
381 env,
382 packKey,
383 packSize: head!.size,
384 limiter: makeLimiter(),
385 countSubrequest: () => {},
386 log,
387 scanResult,
388 repoId: uniqueRepoId(),
389 lruBudget: 1,
390 });
391 
392 const expectedOid = await computeOid("blob", currentPayload);
393 const lastStart = scanResult.table.oids.length - 20;
394 expect(bytesToHex(scanResult.table.oids.subarray(lastStart, lastStart + 20))).toBe(expectedOid);
395 });
396 
397 it("fails fast when rematerialization encounters a cyclic base chain", async () => {
398 const table = allocateEntryTable(2);
399 table.types[0] = 6;
400 table.types[1] = 6;
401 
402 const reader = new SequentialReader(
403 env,
404 "test/materialize-cycle.pack",
405 0,
406 1,
407 makeLimiter(),
408 () => {},
409 log
410 );
411 const baseIndex = new Int32Array([1, 0]);
412 
413 await expect(
414 getBasePayload(
415 {
416 env,
417 packKey: "test/materialize-cycle.pack",
418 packSize: 0,
419 limiter: makeLimiter(),
420 countSubrequest: () => {},
421 log,
422 scanResult: {
423 table,
424 refBaseOids: new Uint8Array(40),
425 refDeltaCount: 0,
426 resolvedCount: 2,
427 objectCount: 2,
428 packChecksum: new Uint8Array(20),
429 },
430 repoId: uniqueRepoId(),
431 },
432 0,
433 new PayloadLRU(1, 2),
434 reader,
435 table,
436 baseIndex
437 )
438 ).rejects.toThrow(/cycle or runaway traversal/);
439 });
440});