File
Blob: test/pack-indexer.resolve.ofs.worker.test.ts
| 1 | import { describe, expect, it } from "vitest"; |
| 2 | import { env } from "cloudflare:workers"; |
| 3 | import { buildPack, buildAppendOnlyDelta } from "./util/git-pack"; |
| 4 | import { uniqueRepoId } from "./util/test-helpers"; |
| 5 | import { |
| 6 | makeCountSubrequest, |
| 7 | makeLimiter, |
| 8 | makeTracingLimiter, |
| 9 | packIndexerLog as log, |
| 10 | } from "./util/pack-indexer.helpers"; |
| 11 | |
| 12 | import { |
| 13 | allocateEntryTable, |
| 14 | isResolveAbortedError, |
| 15 | scanPack, |
| 16 | resolveDeltasAndWriteIdx, |
| 17 | } from "@/worker/git/pack/indexer"; |
| 18 | import { computeOid } from "@/worker/git/core/objects"; |
| 19 | import { bytesToHex } from "@/worker/common/hex"; |
| 20 | import { DEFAULT_SUBREQUEST_BUDGET } from "@/worker/git/operations/limits"; |
| 21 | import { getOidHexAt, parseIdxView } from "@/worker/git/object-store"; |
| 22 | import { packIndexKey } from "@/worker/keys"; |
| 23 | import { getBasePayload } from "@/worker/git/pack/indexer/resolve/materialize"; |
| 24 | import { PayloadLRU } from "@/worker/git/pack/indexer/resolve/payloadCache"; |
| 25 | import { SequentialReader } from "@/worker/git/pack/indexer/resolve/reader"; |
| 26 | |
| 27 | async 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 | |
| 37 | describe("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 | }); |