Skip to content
File

Blob: src/worker/git/pack/rewrite/selectionResolve.ts

typescript602 lines
1import type { OrderedPackSnapshot } from "@/worker/git/operations/fetch/types";
2import type { IdxView, PackedObjectResult } from "@/worker/git/object-store/types";
3import type { Logger } from "@/worker/common/logger";
4import type { PackHeaderEx } from "../packMeta";
5 
6import { bytesToHex, deflate, hexToBytes } from "@/worker/common";
7import { encodeObjHeader, objTypeCode } from "@/worker/git/core/objects";
8import {
9 collectPackedObjectCandidates,
10 findOffsetIndex,
11 findOidRunInIdx,
12} from "@/worker/git/object-store";
13import { materializePackedObjectCandidate } from "@/worker/git/object-store/materialize";
14import {
15 claimCanonicalOwner,
16 clonePackHeader,
17 type DuplicateHeaderCache,
18 type SelectedOidLookup,
19 type SelectionStats,
20} from "./ownership";
21import {
22 ensurePackReadState,
23 growSelectionTable,
24 readSelectedHeader,
25 recordRewriteFailure,
26 selectionDependsOn,
27 selectionKey,
28 setSelectionEntryIdentity,
29 storeSelectionHeader,
30 type PackReadState,
31 type RewriteOptions,
32 type SelectionTable,
33} from "./shared";
34 
35const SYNTHETIC_OBJECT_MAX_BYTES = 8 * 1024 * 1024;
36const SYNTHETIC_PAYLOAD_TOTAL_MAX_BYTES = 32 * 1024 * 1024;
37 
38/**
39 * Result of reading a single entry header and resolving its delta base.
40 *
41 * When `redirectTo` is set, `sel` was redirected to another already-selected
42 * owner for the same OID. When `supersedeSel` is set, `sel` took ownership
43 * back from an older duplicate and that previous owner should be compacted out
44 * after all phases complete.
45 */
46export type HeaderResolveResult =
47 | { ok: true; ofsBaseCanonicalized: boolean; redirectTo?: number; supersedeSel?: number }
48 | { ok: false };
49 
50type RefDeltaBaseChoice = {
51 baseSel: number;
52 candidateCount: number;
53};
54 
55/**
56 * Add a (packSlot, entryIndex) pair to the selection table if not already
57 * present. Returns the selection slot (new or existing).
58 */
59export function addEntry(
60 table: SelectionTable,
61 dedupMap: Map<number, number>,
62 packSlot: number,
63 entryIndex: number,
64 idx: IdxView
65): number {
66 const key = selectionKey(packSlot, entryIndex);
67 const existing = dedupMap.get(key);
68 if (existing !== undefined) return existing;
69 
70 if (table.count >= table.capacity) growSelectionTable(table);
71 
72 const sel = table.count++;
73 setSelectionEntryIdentity(table, sel, packSlot, entryIndex, idx);
74 table.baseSlots[sel] = -1;
75 
76 dedupMap.set(key, sel);
77 return sel;
78}
79 
80/**
81 * Read the header for a selection slot, store all fields into the table,
82 * and immediately resolve any delta base. New bases are pushed to
83 * `secondaryQueue` for chase in the next iteration.
84 */
85export async function readHeaderAndResolveBase(
86 table: SelectionTable,
87 sel: number,
88 snapshot: OrderedPackSnapshot,
89 readerStates: Map<number, PackReadState>,
90 dedupMap: Map<number, number>,
91 oidOwners: SelectedOidLookup,
92 duplicateHeaderCache: DuplicateHeaderCache,
93 stats: SelectionStats,
94 secondaryQueue: number[],
95 env: Env,
96 log: Logger,
97 warnedFlags: Set<string>,
98 options?: RewriteOptions
99): Promise<HeaderResolveResult> {
100 // Skip if already read (entry added as both needed and as a base)
101 table.queuedForHeader[sel] = 0;
102 if (table.typeCodes[sel] !== 0) return { ok: true, ofsBaseCanonicalized: false };
103 
104 const packSlot = table.packSlots[sel];
105 const pack = snapshot.packs[packSlot];
106 const readState = await ensurePackReadState(
107 env,
108 pack,
109 packSlot,
110 readerStates,
111 log,
112 warnedFlags,
113 options
114 );
115 
116 const offset = table.offsets[sel];
117 const header = await readSelectedHeader(readState, offset);
118 if (!header) {
119 log.warn("rewrite:header-read-failed", { packKey: pack.packKey, offset });
120 return { ok: false };
121 }
122 duplicateHeaderCache.set(
123 selectionKey(packSlot, table.entryIndices[sel]),
124 clonePackHeader(header)
125 );
126 
127 if (!storeSelectionHeader(table, sel, offset, table.nextOffsets[sel], header)) {
128 log.warn("rewrite:invalid-payload-length", { packKey: pack.packKey, offset });
129 return { ok: false };
130 }
131 const resolvedHeader = header;
132 
133 let canonicalized = false;
134 const ownership = await claimCanonicalOwner(
135 table,
136 sel,
137 snapshot,
138 readerStates,
139 dedupMap,
140 oidOwners,
141 duplicateHeaderCache,
142 stats,
143 env,
144 log,
145 warnedFlags,
146 options
147 );
148 if (ownership.kind === "error") {
149 return { ok: false };
150 }
151 if (ownership.kind === "redirect") {
152 if (ownership.upgradedOwner) {
153 log.debug("rewrite:duplicate-owner-upgraded", {
154 fromPackKey: pack.packKey,
155 offset,
156 targetSel: ownership.targetSel,
157 });
158 } else {
159 log.debug("rewrite:duplicate-owner-redirect", {
160 fromPackKey: pack.packKey,
161 offset,
162 targetSel: ownership.targetSel,
163 });
164 }
165 return {
166 ok: true,
167 ofsBaseCanonicalized: ownership.canonicalized,
168 redirectTo: ownership.targetSel,
169 };
170 }
171 if (ownership.kind === "takeover") {
172 log.debug("rewrite:duplicate-owner-taken-over-by-ofs-base", {
173 fromPackKey: pack.packKey,
174 offset,
175 previousOwnerSel: ownership.previousOwnerSel,
176 });
177 canonicalized = ownership.canonicalized;
178 // The current row stays live after reclaiming ownership, so its delta base
179 // still needs to be resolved below before the row can be streamed.
180 return await resolveDeltaBaseAndFinish(ownership.previousOwnerSel);
181 }
182 if (ownership.kind === "swapped") {
183 log.debug("rewrite:delta-canonicalized-to-full", {
184 fromPackKey: pack.packKey,
185 offset,
186 });
187 return { ok: true, ofsBaseCanonicalized: true };
188 }
189 canonicalized = ownership.canonicalized;
190 
191 return await resolveDeltaBaseAndFinish();
192 
193 async function resolveDeltaBaseAndFinish(supersedeSel?: number): Promise<HeaderResolveResult> {
194 if (
195 !(await resolveDeltaBaseFromHeader(
196 table,
197 sel,
198 snapshot,
199 dedupMap,
200 secondaryQueue,
201 env,
202 log,
203 resolvedHeader,
204 options
205 ))
206 ) {
207 return { ok: false };
208 }
209 
210 return {
211 ok: true,
212 ofsBaseCanonicalized: canonicalized,
213 supersedeSel,
214 };
215 }
216}
217 
218export async function resolveDeltaBaseFromHeader(
219 table: SelectionTable,
220 sel: number,
221 snapshot: OrderedPackSnapshot,
222 dedupMap: Map<number, number>,
223 secondaryQueue: number[],
224 env: Env,
225 log: Logger,
226 header: PackHeaderEx,
227 options?: RewriteOptions
228): Promise<boolean> {
229 const packSlot = table.packSlots[sel];
230 const pack = snapshot.packs[packSlot];
231 const offset = table.offsets[sel];
232 
233 // Shared helper for both the first header-read pass and the retained-
234 // redirect repair pass. Keeping this logic in one place avoids a second
235 // implementation drifting on subtle rules like OFS pinning or REF_DELTA
236 // base OID storage.
237 
238 // --- OFS_DELTA: base is at (same pack, offset - baseRel) ------------------
239 if (header.type === 6 && header.baseRel !== undefined) {
240 const baseOffset = offset - header.baseRel;
241 const baseIndex = findOffsetIndex(pack.idx, baseOffset);
242 if (baseIndex === undefined) {
243 log.warn("rewrite:missing-ofs-base", { packKey: pack.packKey, offset, baseOffset });
244 recordRewriteFailure(options, {
245 reason: "missing-ofs-base",
246 retryable: false,
247 details: { packKey: pack.packKey, offset, baseOffset },
248 });
249 return false;
250 }
251 
252 // OFS_DELTA bases stay pack-local. The source pack's offset ordering is
253 // already acyclic; cross-pack canonicalization here can manufacture cycles
254 // when a newer duplicate of the same OID is itself stored as a delta.
255 const baseSel = addEntry(table, dedupMap, packSlot, baseIndex, pack.idx);
256 table.ofsPinned[baseSel] = 1;
257 
258 if (baseSel === sel) {
259 // Self-referential OFS_DELTA (e.g. baseRel is zero or points back to
260 // own offset). This indicates pack corruption โ€” abort the selection so
261 // the caller can retry or investigate.
262 log.warn("rewrite:self-referential-delta", {
263 packKey: pack.packKey,
264 offset,
265 deltaType: "ofs",
266 baseRel: header.baseRel,
267 });
268 recordRewriteFailure(options, {
269 reason: "self-referential-ofs-delta",
270 retryable: false,
271 details: { packKey: pack.packKey, offset, baseRel: header.baseRel },
272 });
273 return false;
274 }
275 table.baseSlots[sel] = baseSel;
276 if (table.typeCodes[baseSel] === 0 && !table.queuedForHeader[baseSel]) {
277 table.queuedForHeader[baseSel] = 1;
278 secondaryQueue.push(baseSel);
279 }
280 }
281 
282 // --- REF_DELTA: base identified by OID, may be in any pack ----------------
283 if (header.type === 7 && header.baseOid) {
284 // Store raw base OID bytes for streaming (avoids hex round-trip later)
285 const rawBytes = hexToBytes(header.baseOid);
286 if (!table.baseOidRaw) {
287 table.baseOidRaw = new Uint8Array(table.capacity * 20);
288 }
289 table.baseOidRaw.set(rawBytes, sel * 20);
290 
291 const choice = chooseRefDeltaBase(
292 table,
293 sel,
294 snapshot,
295 dedupMap,
296 log,
297 rawBytes,
298 header.baseOid
299 );
300 if (choice.baseSel < 0) {
301 const reason = choice.candidateCount === 0 ? "missing-ref-base" : "ref-base-cycle-unresolved";
302 log.warn(`rewrite:${reason}`, {
303 sel,
304 packKey: pack.packKey,
305 offset,
306 baseOid: header.baseOid,
307 candidateCount: choice.candidateCount,
308 });
309 return await materializeSelectionAsFullObject({
310 table,
311 sel,
312 snapshot,
313 env,
314 log,
315 options,
316 reason,
317 baseOid: header.baseOid,
318 candidateCount: choice.candidateCount,
319 });
320 }
321 
322 const baseSel = choice.baseSel;
323 if (baseSel === sel) {
324 // Self-referential REF_DELTA with no full-object duplicate available.
325 // This pack cannot be rewritten into a topologically valid output.
326 log.warn("rewrite:self-referential-delta", {
327 packKey: pack.packKey,
328 offset,
329 deltaType: "ref",
330 baseOid: header.baseOid,
331 });
332 return await materializeSelectionAsFullObject({
333 table,
334 sel,
335 snapshot,
336 env,
337 log,
338 options,
339 reason: "self-referential-ref-delta",
340 baseOid: header.baseOid,
341 candidateCount: choice.candidateCount,
342 });
343 }
344 
345 table.baseSlots[sel] = baseSel;
346 if (table.typeCodes[baseSel] === 0 && !table.queuedForHeader[baseSel]) {
347 table.queuedForHeader[baseSel] = 1;
348 secondaryQueue.push(baseSel);
349 }
350 }
351 
352 return true;
353}
354 
355type MaterializeSelectionArgs = {
356 table: SelectionTable;
357 sel: number;
358 snapshot: OrderedPackSnapshot;
359 env: Env;
360 log: Logger;
361 options?: RewriteOptions;
362 reason: string;
363 baseOid: string;
364 candidateCount: number;
365};
366 
367type MaterializeOidArgs = {
368 snapshot: OrderedPackSnapshot;
369 env: Env;
370 log: Logger;
371 options?: RewriteOptions;
372 oid: string | Uint8Array;
373 visited: Set<string>;
374};
375 
376function selectedOidHex(table: SelectionTable, sel: number): string {
377 return bytesToHex(table.oidsRaw.subarray(sel * 20, sel * 20 + 20));
378}
379 
380function syntheticPayloadTotal(table: SelectionTable): number {
381 let total = 0;
382 for (let sel = 0; sel < table.syntheticPayloads.length; sel++) {
383 total += table.syntheticPayloads[sel]?.byteLength ?? 0;
384 }
385 return total;
386}
387 
388async function materializeOidFromSnapshot(
389 args: MaterializeOidArgs
390): Promise<PackedObjectResult | undefined> {
391 const limiter = args.options?.limiter;
392 const countSubrequest = args.options?.countSubrequest;
393 if (!limiter || !countSubrequest) return undefined;
394 
395 const candidates = collectPackedObjectCandidates(args.snapshot.packs, args.oid);
396 for (const candidate of candidates) {
397 if (args.options?.signal?.aborted) return undefined;
398 
399 const object = await materializePackedObjectCandidate({
400 env: args.env,
401 candidate,
402 limiter,
403 countSubrequest,
404 log: args.log,
405 cyclePolicy: "miss",
406 resolveRefBase: async (baseOid, nextVisited) => {
407 return await materializeOidFromSnapshot({
408 ...args,
409 oid: baseOid,
410 visited: nextVisited,
411 });
412 },
413 visited: args.visited,
414 signal: args.options?.signal,
415 });
416 if (object) return object;
417 }
418 
419 return undefined;
420}
421 
422async function materializeSelectionAsFullObject(args: MaterializeSelectionArgs): Promise<boolean> {
423 const oid = selectedOidHex(args.table, args.sel);
424 const object = await materializeOidFromSnapshot({
425 snapshot: args.snapshot,
426 env: args.env,
427 log: args.log,
428 options: args.options,
429 oid,
430 visited: new Set<string>(),
431 });
432 if (!object) {
433 args.log.warn("rewrite:cycle-breaker-materialize-miss", {
434 sel: args.sel,
435 oid,
436 reason: args.reason,
437 baseOid: args.baseOid,
438 candidateCount: args.candidateCount,
439 });
440 recordRewriteFailure(args.options, {
441 reason: "cycle-breaker-materialize-miss",
442 retryable: false,
443 details: {
444 sel: args.sel,
445 oid,
446 baseOid: args.baseOid,
447 candidateCount: args.candidateCount,
448 },
449 });
450 return false;
451 }
452 
453 if (object.oid !== oid) {
454 args.log.warn("rewrite:cycle-breaker-oid-mismatch", {
455 sel: args.sel,
456 oid,
457 materializedOid: object.oid,
458 });
459 recordRewriteFailure(args.options, {
460 reason: "cycle-breaker-oid-mismatch",
461 retryable: false,
462 details: { sel: args.sel, oid, materializedOid: object.oid },
463 });
464 return false;
465 }
466 
467 if (object.payload.byteLength > SYNTHETIC_OBJECT_MAX_BYTES) {
468 args.log.warn("rewrite:cycle-breaker-object-too-large", {
469 sel: args.sel,
470 oid,
471 payloadBytes: object.payload.byteLength,
472 maxBytes: SYNTHETIC_OBJECT_MAX_BYTES,
473 });
474 recordRewriteFailure(args.options, {
475 reason: "synthetic-object-too-large",
476 retryable: false,
477 details: {
478 sel: args.sel,
479 oid,
480 payloadBytes: object.payload.byteLength,
481 maxBytes: SYNTHETIC_OBJECT_MAX_BYTES,
482 },
483 });
484 return false;
485 }
486 
487 const compressedPayload = await deflate(object.payload);
488 const existingPayloadBytes = args.table.syntheticPayloads[args.sel]?.byteLength ?? 0;
489 const nextSyntheticTotal =
490 syntheticPayloadTotal(args.table) - existingPayloadBytes + compressedPayload.byteLength;
491 if (nextSyntheticTotal > SYNTHETIC_PAYLOAD_TOTAL_MAX_BYTES) {
492 args.log.warn("rewrite:cycle-breaker-total-too-large", {
493 sel: args.sel,
494 oid,
495 compressedBytes: compressedPayload.byteLength,
496 totalBytes: nextSyntheticTotal,
497 maxBytes: SYNTHETIC_PAYLOAD_TOTAL_MAX_BYTES,
498 });
499 recordRewriteFailure(args.options, {
500 reason: "synthetic-payload-budget-exceeded",
501 retryable: false,
502 details: {
503 sel: args.sel,
504 oid,
505 compressedBytes: compressedPayload.byteLength,
506 totalBytes: nextSyntheticTotal,
507 maxBytes: SYNTHETIC_PAYLOAD_TOTAL_MAX_BYTES,
508 },
509 });
510 return false;
511 }
512 
513 const typeCode = objTypeCode(object.type);
514 const headerBytes = encodeObjHeader(typeCode, object.payload.byteLength);
515 if (headerBytes.byteLength > 5) {
516 recordRewriteFailure(args.options, {
517 reason: "synthetic-header-too-large",
518 retryable: false,
519 details: { sel: args.sel, oid, headerBytes: headerBytes.byteLength },
520 });
521 return false;
522 }
523 
524 args.table.typeCodes[args.sel] = typeCode;
525 args.table.headerLens[args.sel] = headerBytes.byteLength;
526 args.table.payloadLens[args.sel] = compressedPayload.byteLength;
527 args.table.sizeVarBuf.set(headerBytes, args.sel * 5);
528 args.table.sizeVarLens[args.sel] = headerBytes.byteLength;
529 args.table.baseSlots[args.sel] = -1;
530 args.table.queuedForHeader[args.sel] = 0;
531 args.table.syntheticPayloads[args.sel] = compressedPayload;
532 
533 args.log.info("rewrite:cycle-breaker-materialized", {
534 sel: args.sel,
535 oid,
536 type: object.type,
537 reason: args.reason,
538 baseOid: args.baseOid,
539 candidateCount: args.candidateCount,
540 payloadBytes: object.payload.byteLength,
541 compressedBytes: compressedPayload.byteLength,
542 syntheticTotalBytes: nextSyntheticTotal,
543 });
544 return true;
545}
546 
547/**
548 * Choose a REF_DELTA base by scanning duplicate OID runs in snapshot order.
549 *
550 * The hot path only uses the already-loaded idx views and the partially wired
551 * selection table. If a duplicate candidate is already selected, adding
552 * `sel -> candidateSel` must not make the selected base chain cycle back into
553 * `sel`; otherwise the chooser falls through to the next duplicate candidate.
554 */
555function chooseRefDeltaBase(
556 table: SelectionTable,
557 sel: number,
558 snapshot: OrderedPackSnapshot,
559 dedupMap: Map<number, number>,
560 log: Logger,
561 rawBytes: Uint8Array,
562 baseOid: string
563): RefDeltaBaseChoice {
564 const currentPackSlot = table.packSlots[sel];
565 const currentEntryIndex = table.entryIndices[sel];
566 let candidateCount = 0;
567 
568 for (let candidatePackSlot = 0; candidatePackSlot < snapshot.packs.length; candidatePackSlot++) {
569 const candidatePack = snapshot.packs[candidatePackSlot]!;
570 const run = findOidRunInIdx(candidatePack.idx, rawBytes);
571 if (!run) continue;
572 
573 for (let entryIndex = run.startIndex; entryIndex <= run.endIndex; entryIndex++) {
574 candidateCount++;
575 if (candidatePackSlot === currentPackSlot && entryIndex === currentEntryIndex) {
576 continue;
577 }
578 
579 const candidateKey = selectionKey(candidatePackSlot, entryIndex);
580 const candidateSel = dedupMap.get(candidateKey);
581 if (candidateSel !== undefined) {
582 if (candidateSel === sel || selectionDependsOn(table, candidateSel, sel)) {
583 log.debug("rewrite:ref-base-candidate-cycle-skipped", {
584 sel,
585 baseOid,
586 candidatePackSlot,
587 candidateEntryIndex: entryIndex,
588 selectedCandidateSel: candidateSel,
589 });
590 continue;
591 }
592 return { baseSel: candidateSel, candidateCount };
593 }
594 
595 const baseSel = addEntry(table, dedupMap, candidatePackSlot, entryIndex, candidatePack.idx);
596 return { baseSel, candidateCount };
597 }
598 }
599 
600 return { baseSel: -1, candidateCount };
601}