Skip to content
File

Blob: src/worker/git/pack/indexer/scan.ts

typescript504 lines
1/**
2 * Streaming pack scanner (Pass 1).
3 *
4 * Reads a .pack file from R2 in sequential chunks, parses every entry header,
5 * inflates compressed data to determine span boundaries, computes OIDs for
6 * non-delta objects, and records metadata for delta objects. The entire pack is
7 * never buffered in memory; at most one object's inflated payload is held at a
8 * time before being discarded.
9 */
10 
11import { bytesToHex } from "@/worker/common/hex";
12import { createDigestStream } from "@/worker/common/webtypes";
13import { computeOidBytes } from "@/worker/git/core/objects";
14import { readPackRange } from "@/worker/git/pack/packMeta";
15import { typeCodeToObjectType } from "@/worker/git/object-store/support";
16import { PackRefsBuilder } from "@/worker/git/pack/refIndex";
17 
18import { InflateCursor, CRC32_INIT, crc32Update, crc32Finish } from "./inflateCursor";
19import { allocateEntryTable } from "./types";
20import type { IndexerOptions, ScanResult, RefBaseOids } from "./types";
21 
22const DEFAULT_CHUNK_SIZE = 1_048_576; // 1 MiB
23const DELTA_HEADER_CAPTURE_LIMIT = 16;
24const PACK_HEADER_BYTES = 12;
25const PACK_TRAILER_BYTES = 20;
26const MIN_PACK_BYTES = PACK_HEADER_BYTES + PACK_TRAILER_BYTES;
27// Smallest possible packed entry:
28// - 1 byte Git pack object header
29// - 8 bytes for the smallest valid zlib-wrapped empty payload
30//
31// Delta entries are always larger than this, so using the constant as a
32// lower bound can only reject impossible object counts, never valid packs.
33const MIN_PACKED_ENTRY_BYTES = 9;
34// The scanner stores one typed-array row per object. Keep a hard ceiling so a
35// malformed pack header cannot force unbounded metadata allocation in the
36// shared worker isolate before any object bytes are validated.
37const MAX_INDEXABLE_OBJECT_COUNT = 250_000;
38const MAX_PACK_OFFSET = 0xffffffff;
39const MAX_PACK_OBJECT_SIZE = 0xffffffff;
40const SCAN_PROGRESS_STEPS = 20;
41 
42function emitScanProgress(
43 onProgress: IndexerOptions["onProgress"],
44 processed: number,
45 total: number
46): void {
47 if (!onProgress || total <= 0) return;
48 const percent = Math.round((processed / total) * 100);
49 if (processed >= total) {
50 onProgress(`Scanning pack objects: 100% (${total}/${total}), done.\n`);
51 return;
52 }
53 onProgress(`Scanning pack objects: ${percent}% (${processed}/${total})\r`);
54}
55 
56function isReservedPackType(type: number): boolean {
57 return type === 0 || type === 5;
58}
59 
60function readDeltaSizeVarint(
61 data: Uint8Array,
62 pos: number,
63 fieldName: string
64): { value: number; nextPos: number } {
65 let value = 0;
66 let factor = 1;
67 let cursor = pos;
68 
69 while (cursor < data.length) {
70 const b = data[cursor++];
71 value += (b & 0x7f) * factor;
72 if (value > MAX_PACK_OBJECT_SIZE) {
73 throw new Error(`scan: ${fieldName} exceeds supported 32-bit size range`);
74 }
75 if (!(b & 0x80)) {
76 return { value, nextPos: cursor };
77 }
78 if (factor > Number.MAX_SAFE_INTEGER / 128) {
79 throw new Error(`scan: ${fieldName} is too large to decode safely`);
80 }
81 factor *= 128;
82 }
83 
84 throw new Error(`scan: truncated ${fieldName} header`);
85}
86 
87// ---------------------------------------------------------------------------
88// Pack header varint parsing helpers (inline to avoid extra R2 reads)
89// ---------------------------------------------------------------------------
90 
91/**
92 * Parse a pack entry header from an in-memory buffer at the given position.
93 * Returns the type code, decompressed size, header length, and delta metadata.
94 *
95 * The logic mirrors `readPackHeaderExFromBuf` in packMeta.ts but returns the
96 * decoded size directly and operates on a position-based cursor so the scanner
97 * can work from its sliding buffer without copying.
98 */
99function parseEntryHeader(
100 buf: Uint8Array,
101 pos: number
102): {
103 type: number;
104 size: number;
105 headerLen: number;
106 baseOidBytes?: Uint8Array;
107 baseRel?: number;
108} | null {
109 if (pos >= buf.length) return null;
110 
111 const start = pos;
112 let c = buf[pos++];
113 const type = (c >> 4) & 0x07;
114 if (isReservedPackType(type)) {
115 throw new Error(`scan: invalid reserved pack type ${type} at offset ${start}`);
116 }
117 let size = c & 0x0f;
118 let factor = 16;
119 
120 while (c & 0x80) {
121 if (pos >= buf.length) return null;
122 c = buf[pos++];
123 size += (c & 0x7f) * factor;
124 if (size > MAX_PACK_OBJECT_SIZE) {
125 throw new Error(`scan: entry size exceeds supported 32-bit range at offset ${start}`);
126 }
127 if (c & 0x80) {
128 // The indexer stores sizes in Uint32Array-backed tables. Reject absurdly
129 // long varints here instead of letting arithmetic drift past safe integer
130 // precision and then mislabeling the failure later in the scan.
131 if (factor > Number.MAX_SAFE_INTEGER / 128) {
132 throw new Error(`scan: entry size is too large to decode safely at offset ${start}`);
133 }
134 factor *= 128;
135 }
136 }
137 
138 const sizeVarLen = pos - start;
139 
140 if (type === 7) {
141 // REF_DELTA: 20-byte base OID follows the size varint
142 if (pos + 20 > buf.length) return null;
143 // Borrow a zero-copy view here and copy it into the flat ref-base table
144 // immediately after header parsing. The backing scan buffer is not stable
145 // once the reader advances to the next range window.
146 const baseOidBytes = buf.subarray(pos, pos + 20);
147 return { type, size, headerLen: sizeVarLen + 20, baseOidBytes };
148 }
149 
150 if (type === 6) {
151 // OFS_DELTA: variable-length negative offset follows the size varint
152 if (pos >= buf.length) return null;
153 let b = buf[pos++];
154 let x = b & 0x7f;
155 while (b & 0x80) {
156 if (pos >= buf.length) return null;
157 b = buf[pos++];
158 x = (x + 1) * 128 + (b & 0x7f);
159 if (x > MAX_PACK_OFFSET) {
160 throw new Error(`scan: OFS_DELTA base distance exceeds 32-bit range at offset ${start}`);
161 }
162 }
163 return { type, size, headerLen: pos - start, baseRel: x };
164 }
165 
166 return { type, size, headerLen: sizeVarLen };
167}
168 
169/**
170 * Read the base_size and result_size varints from the start of a delta
171 * instruction stream. These are the first two varints in the inflated delta
172 * payload (before any copy/insert opcodes).
173 */
174function readDeltaResultSize(data: Uint8Array): number {
175 const baseSize = readDeltaSizeVarint(data, 0, "delta base-size");
176 const resultSize = readDeltaSizeVarint(data, baseSize.nextPos, "delta result-size");
177 return resultSize.value;
178}
179 
180function validatePackObjectCount(packSize: number, objectCount: number): void {
181 if (objectCount > MAX_INDEXABLE_OBJECT_COUNT) {
182 throw new Error(
183 `scan: object count ${objectCount} exceeds safe isolate limit ${MAX_INDEXABLE_OBJECT_COUNT}`
184 );
185 }
186 
187 const bytesAvailableForEntries = packSize - PACK_HEADER_BYTES - PACK_TRAILER_BYTES;
188 const maxPossibleObjects = Math.floor(bytesAvailableForEntries / MIN_PACKED_ENTRY_BYTES);
189 if (objectCount > maxPossibleObjects) {
190 throw new Error(
191 `scan: object count ${objectCount} cannot fit in ${packSize} bytes of pack data`
192 );
193 }
194}
195 
196// ---------------------------------------------------------------------------
197// Buffered reader – manages a sliding window over sequential R2 range reads
198// ---------------------------------------------------------------------------
199 
200class BufferedPackReader {
201 private env: Env;
202 private packKey: string;
203 private packSize: number;
204 private chunkSize: number;
205 private limiter: IndexerOptions["limiter"];
206 private countSub: IndexerOptions["countSubrequest"];
207 private signal?: AbortSignal;
208 
209 /** Current in-memory buffer. */
210 buf: Uint8Array<ArrayBufferLike> = new Uint8Array(0);
211 /** Current read position within `buf`. */
212 pos = 0;
213 /** Absolute pack offset corresponding to buf[0]. */
214 bufAbsStart = 0;
215 
216 constructor(opts: IndexerOptions) {
217 this.env = opts.env;
218 this.packKey = opts.packKey;
219 this.packSize = opts.packSize;
220 this.chunkSize = opts.chunkSize ?? DEFAULT_CHUNK_SIZE;
221 this.limiter = opts.limiter;
222 this.countSub = opts.countSubrequest;
223 this.signal = opts.signal;
224 }
225 
226 /** Absolute offset of the current read position in the pack. */
227 get absPos(): number {
228 return this.bufAbsStart + this.pos;
229 }
230 
231 /** Number of unread bytes remaining in the current buffer. */
232 get remaining(): number {
233 return this.buf.length - this.pos;
234 }
235 
236 /**
237 * Ensure the buffer has at least `minBytes` available from the current
238 * position. Reads a new chunk from R2 if necessary, concatenating with
239 * any leftover bytes from the current buffer.
240 */
241 async ensure(minBytes: number): Promise<void> {
242 while (this.remaining < minBytes) {
243 // `buf` always represents one contiguous window: unread bytes from the
244 // previous chunk followed by freshly fetched bytes. Reading from
245 // `bufAbsStart + buf.length` therefore continues exactly where the
246 // current window ends.
247 const nextAbsOffset = this.bufAbsStart + this.buf.length;
248 const bytesLeft = this.packSize - nextAbsOffset;
249 if (bytesLeft <= 0) return;
250 const readLen = Math.min(this.chunkSize, bytesLeft);
251 
252 const chunk = await readPackRange(this.env, this.packKey, nextAbsOffset, readLen, {
253 limiter: this.limiter,
254 countSubrequest: this.countSub,
255 signal: this.signal,
256 });
257 if (!chunk) throw new Error("scan: unexpected R2 read failure");
258 
259 const leftover = this.buf.subarray(this.pos);
260 if (leftover.length === 0) {
261 this.bufAbsStart = nextAbsOffset;
262 this.buf = chunk;
263 this.pos = 0;
264 continue;
265 }
266 
267 const newBuf = new Uint8Array(leftover.length + chunk.length);
268 newBuf.set(leftover, 0);
269 newBuf.set(chunk, leftover.length);
270 
271 this.bufAbsStart += this.pos;
272 this.buf = newBuf;
273 this.pos = 0;
274 }
275 }
276 
277 /** Consume `n` bytes from the buffer and advance the cursor. */
278 consume(n: number): Uint8Array {
279 const slice = this.buf.subarray(this.pos, this.pos + n);
280 this.pos += n;
281 return slice;
282 }
283}
284 
285// ---------------------------------------------------------------------------
286// scanPack
287// ---------------------------------------------------------------------------
288 
289export async function scanPack(opts: IndexerOptions): Promise<ScanResult> {
290 const { env, packKey, packSize, log } = opts;
291 if (!Number.isSafeInteger(packSize) || packSize < MIN_PACK_BYTES) {
292 throw new Error(`scan: pack size ${packSize} is smaller than the minimum valid pack size`);
293 }
294 const reader = new BufferedPackReader(opts);
295 if (packSize - PACK_TRAILER_BYTES > MAX_PACK_OFFSET) {
296 throw new Error(
297 `scan: pack offsets above ${MAX_PACK_OFFSET} bytes are not supported by this indexer yet`
298 );
299 }
300 
301 // ---- 1. Read and validate the 12-byte pack header ----
302 await reader.ensure(PACK_HEADER_BYTES);
303 const headerBuf = reader.consume(PACK_HEADER_BYTES);
304 
305 const magic =
306 String.fromCharCode(headerBuf[0]) +
307 String.fromCharCode(headerBuf[1]) +
308 String.fromCharCode(headerBuf[2]) +
309 String.fromCharCode(headerBuf[3]);
310 if (magic !== "PACK") throw new Error("scan: invalid pack magic");
311 
312 const hdv = new DataView(headerBuf.buffer, headerBuf.byteOffset, 12);
313 const version = hdv.getUint32(4, false);
314 if (version !== 2) throw new Error(`scan: unsupported pack version ${version}`);
315 
316 const objectCount = hdv.getUint32(8, false);
317 validatePackObjectCount(packSize, objectCount);
318 log.info("scan:start", { packKey, packSize, objectCount });
319 opts.onProgress?.(`Scanning pack objects: 0% (0/${objectCount})\r`);
320 
321 // ---- 2. Allocate entry table ----
322 const table = allocateEntryTable(objectCount);
323 const refsBuilder = new PackRefsBuilder(objectCount);
324 const refBaseOids: RefBaseOids = new Uint8Array(objectCount * 20);
325 let refDeltaCount = 0;
326 let resolvedCount = 0;
327 
328 // ---- 3. Streaming SHA-1 digest (covers everything except the trailing 20 bytes) ----
329 const digestStream = createDigestStream("SHA-1");
330 const digestWriter = digestStream.getWriter();
331 // Feed the 12-byte header to the digest.
332 await digestWriter.write(headerBuf);
333 
334 const inflator = new InflateCursor();
335 const progressInterval = Math.max(1, Math.floor(objectCount / SCAN_PROGRESS_STEPS));
336 
337 // ---- 4. Sequential scan of all entries ----
338 for (let i = 0; i < objectCount; i++) {
339 const entryStart = reader.absPos;
340 
341 // Make sure we have enough bytes to parse the header (up to ~30 bytes for
342 // type varint + OFS_DELTA distance or REF_DELTA OID).
343 await reader.ensure(Math.min(64, packSize - entryStart));
344 
345 const header = parseEntryHeader(reader.buf, reader.pos);
346 if (!header) throw new Error(`scan: failed to parse header at offset ${entryStart}`);
347 
348 // Record basic metadata.
349 table.offsets[i] = entryStart;
350 table.types[i] = header.type;
351 table.headerLens[i] = header.headerLen;
352 table.decompressedSizes[i] = header.size; // overwritten for deltas below
353 
354 if (header.type === 6 && header.baseRel !== undefined) {
355 if (header.baseRel > entryStart) {
356 throw new Error(
357 `scan: OFS_DELTA at offset ${entryStart} points before the start of the pack`
358 );
359 }
360 table.ofsBaseOffsets[i] = entryStart - header.baseRel;
361 }
362 if (header.type === 7 && header.baseOidBytes) {
363 refBaseOids.set(header.baseOidBytes, i * 20);
364 refDeltaCount++;
365 }
366 
367 // CRC-32 accumulation starts with the header bytes.
368 const headerBytes = reader.buf.subarray(reader.pos, reader.pos + header.headerLen);
369 let crc = crc32Update(CRC32_INIT, headerBytes, 0, headerBytes.length);
370 
371 // Feed header bytes to the SHA-1 digest.
372 await digestWriter.write(headerBytes);
373 
374 // Advance past the header.
375 reader.pos += header.headerLen;
376 
377 // ---- Inflate the compressed payload ----
378 const captureLimit = typeCodeToObjectType(header.type)
379 ? header.size
380 : DELTA_HEADER_CAPTURE_LIMIT;
381 inflator.reset({ captureLimit });
382 let firstInflatePush = true;
383 while (!inflator.finished) {
384 // Ensure we have data to feed.
385 // The first push must include the full 2-byte zlib header. A 1-byte
386 // first chunk used to mis-parse valid packs at range boundaries.
387 const minBytes = firstInflatePush ? 2 : 1;
388 if (reader.remaining < minBytes) {
389 await reader.ensure(minBytes);
390 }
391 if (reader.remaining < minBytes) {
392 throw new Error(`scan: unexpected EOF while inflating entry at offset ${entryStart}`);
393 }
394 const available = reader.buf.subarray(reader.pos, reader.pos + reader.remaining);
395 inflator.push(available);
396 firstInflatePush = false;
397 
398 const consumed = inflator.consumedInputBytes;
399 if (consumed <= 0 && !inflator.finished) {
400 throw new Error(`scan: inflate stalled at offset ${reader.absPos}`);
401 }
402 
403 // Feed consumed compressed bytes to CRC and digest.
404 const compressedSlice = reader.buf.subarray(reader.pos, reader.pos + consumed);
405 crc = crc32Update(crc, compressedSlice, 0, compressedSlice.length);
406 await digestWriter.write(compressedSlice);
407 
408 reader.pos += consumed;
409 }
410 
411 // Record span end and CRC.
412 table.spanEnds[i] = reader.absPos;
413 table.crc32s[i] = crc32Finish(crc);
414 
415 // ---- Compute OID for non-delta objects ----
416 const baseType = typeCodeToObjectType(header.type);
417 if (baseType) {
418 const inflated = inflator.output;
419 if (inflated.length !== header.size) {
420 throw new Error(
421 `scan: inflated ${baseType} size mismatch at offset ${entryStart} (expected ${header.size}, got ${inflated.length})`
422 );
423 }
424 // Non-delta: hash "<type> <size>\0<payload>" to get the OID.
425 table.oids.set(await computeOidBytes(baseType, inflated), i * 20);
426 table.objectTypes[i] = header.type;
427 refsBuilder.recordObject(i, baseType, inflated);
428 table.resolved[i] = 1;
429 table.decompressedSizes[i] = inflated.length;
430 resolvedCount++;
431 } else {
432 if (inflator.outputLength !== header.size) {
433 throw new Error(
434 `scan: inflated delta size mismatch at offset ${entryStart} (expected ${header.size}, got ${inflator.outputLength})`
435 );
436 }
437 // Pack headers store the size of the inflated delta *program*, not the
438 // final post-apply object size. The apply step needs the latter, so we
439 // read the delta's declared result size here and stash it for resolve().
440 const resultSize = readDeltaResultSize(inflator.capturedOutput);
441 table.decompressedSizes[i] = resultSize;
442 table.resolved[i] = 0;
443 }
444 
445 // Log progress periodically.
446 if ((i + 1) % 10000 === 0 || i + 1 === objectCount) {
447 log.debug("scan:progress", { processed: i + 1, total: objectCount });
448 }
449 if ((i + 1) % progressInterval === 0 || i + 1 === objectCount) {
450 emitScanProgress(opts.onProgress, i + 1, objectCount);
451 }
452 }
453 
454 // ---- 5. Validate trailing SHA-1 checksum ----
455 // A valid pack has no slack bytes between the final entry and the trailing
456 // 20-byte checksum. The digest covers exactly bytes [0, packSize - 20), so
457 // accepting extra bytes here would let a malformed pack smuggle undeclared
458 // data past the scanner while still reusing the original trailer hash.
459 const trailingOffset = packSize - PACK_TRAILER_BYTES;
460 if (reader.absPos !== trailingOffset) {
461 throw new Error(
462 `scan: expected indexed entries to end at ${trailingOffset}, got ${reader.absPos}`
463 );
464 }
465 
466 await digestWriter.close();
467 const computedHash = new Uint8Array(await digestStream.digest);
468 
469 // Read the trailing 20 bytes.
470 const trailer = await readPackRange(env, packKey, trailingOffset, PACK_TRAILER_BYTES, {
471 limiter: opts.limiter,
472 countSubrequest: opts.countSubrequest,
473 signal: opts.signal,
474 });
475 if (!trailer || trailer.length !== PACK_TRAILER_BYTES) {
476 throw new Error("scan: failed to read pack trailer checksum");
477 }
478 
479 // Compare.
480 for (let i = 0; i < 20; i++) {
481 if (computedHash[i] !== trailer[i]) {
482 throw new Error(
483 `scan: pack checksum mismatch (computed ${bytesToHex(computedHash)} != trailing ${bytesToHex(trailer)})`
484 );
485 }
486 }
487 
488 log.info("scan:done", {
489 objectCount,
490 resolved: resolvedCount,
491 deltas: objectCount - resolvedCount,
492 });
493 
494 return {
495 table,
496 refBaseOids,
497 refDeltaCount,
498 resolvedCount,
499 objectCount,
500 packChecksum: trailer,
501 refsBuilder,
502 };
503}