Skip to content
File

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

typescript464 lines
1import type {
2 OrderedPackSnapshot,
3 OrderedPackSnapshotEntry,
4} from "@/worker/git/operations/fetch/types";
5import type { Logger } from "@/worker/common/logger";
6import type { Limiter } from "@/worker/git/operations/limits";
7import type { IdxView } from "@/worker/git/object-store/types";
8import type { PackHeaderEx } from "../packMeta";
9 
10import { createLogger } from "@/worker/common";
11import { findFirstPackedObjectCandidate, getNextOffsetByIndex } from "@/worker/git/object-store";
12import { SequentialReader } from "@/worker/git/pack/indexer/resolve/reader";
13import { readPackHeaderExFromBuf, readPackRange } from "../packMeta";
14 
15export const HEADER_READ_BYTES = 128;
16export const DEFAULT_CHUNK_SIZE = 4_194_304;
17export const WHOLE_PACK_MAX_BYTES = 8 * 1024 * 1024;
18export const WHOLE_PACK_TOTAL_BUDGET = 32 * 1024 * 1024;
19export const HEADER_STABILITY_CAP = 16;
20 
21export type RewriteFailure = {
22 reason: string;
23 retryable: boolean;
24 details?: Record<string, unknown>;
25};
26 
27export type RewriteFailureRecorder = {
28 value?: RewriteFailure;
29};
30 
31export type RewriteOptions = {
32 signal?: AbortSignal;
33 limiter?: Limiter;
34 countSubrequest?: (n?: number) => boolean | void;
35 onProgress?: (msg: string) => void;
36 failure?: RewriteFailureRecorder;
37};
38 
39/**
40 * Flat typed-array representation of selected pack entries, modeled after
41 * the `PackEntryTable` pattern in `src/worker/git/pack/indexer/types.ts`.
42 *
43 * Every field is indexed by a selection slot (`sel`). The table grows on
44 * demand when delta-base chasing discovers bases beyond the initial
45 * `neededOids` set.
46 */
47export interface SelectionTable {
48 count: number;
49 capacity: number;
50 
51 /* identity — set during resolve */
52 packSlots: Uint8Array;
53 entryIndices: Uint32Array;
54 offsets: Float64Array;
55 nextOffsets: Float64Array;
56 oidsRaw: Uint8Array; // capacity * 20, selected object OIDs
57 
58 /* header data — set during header-read pass */
59 typeCodes: Uint8Array;
60 headerLens: Uint16Array;
61 payloadLens: Float64Array;
62 sizeVarBuf: Uint8Array; // capacity * 5, concatenated varint bytes
63 sizeVarLens: Uint8Array; // 1–5 per entry
64 
65 /* delta relationships */
66 baseSlots: Int32Array; // sel of base, -1 = non-delta
67 baseOidRaw: Uint8Array | null; // lazy; capacity * 20, for REF_DELTA
68 queuedForHeader: Uint8Array; // 1 = already queued for header read
69 ofsPinned: Uint8Array; // 1 = exact pack position is required by an OFS_DELTA child
70 syntheticPayloads: Array<Uint8Array | undefined>; // compressed full-object payloads by sel
71 
72 /* output layout — set during convergence / topology sort */
73 outputOffsets: Float64Array;
74 outputHeaderLens: Uint16Array;
75 outputOrder: Uint32Array; // topology-sorted selection indices
76}
77 
78export function allocateSelectionTable(capacity: number): SelectionTable {
79 return {
80 count: 0,
81 capacity,
82 packSlots: new Uint8Array(capacity),
83 entryIndices: new Uint32Array(capacity),
84 offsets: new Float64Array(capacity),
85 nextOffsets: new Float64Array(capacity),
86 oidsRaw: new Uint8Array(capacity * 20),
87 typeCodes: new Uint8Array(capacity),
88 headerLens: new Uint16Array(capacity),
89 payloadLens: new Float64Array(capacity),
90 sizeVarBuf: new Uint8Array(capacity * 5),
91 sizeVarLens: new Uint8Array(capacity),
92 baseSlots: new Int32Array(capacity).fill(-1),
93 baseOidRaw: null, // allocated lazily on first REF_DELTA
94 queuedForHeader: new Uint8Array(capacity),
95 ofsPinned: new Uint8Array(capacity),
96 syntheticPayloads: new Array<Uint8Array | undefined>(capacity),
97 outputOffsets: new Float64Array(capacity),
98 outputHeaderLens: new Uint16Array(capacity),
99 outputOrder: new Uint32Array(capacity),
100 };
101}
102 
103/** Double the table capacity, preserving existing data. */
104export function growSelectionTable(table: SelectionTable): void {
105 const next = Math.max(table.capacity * 2, 64);
106 
107 function grow<T extends ArrayLike<number> & { set(src: T): void }>(
108 old: T,
109 ctor: new (len: number) => T,
110 len: number
111 ): T {
112 const arr = new ctor(len);
113 arr.set(old);
114 return arr;
115 }
116 
117 table.packSlots = grow(table.packSlots, Uint8Array, next);
118 table.entryIndices = grow(table.entryIndices, Uint32Array, next);
119 table.offsets = grow(table.offsets, Float64Array, next);
120 table.nextOffsets = grow(table.nextOffsets, Float64Array, next);
121 table.typeCodes = grow(table.typeCodes, Uint8Array, next);
122 table.headerLens = grow(table.headerLens, Uint16Array, next);
123 table.payloadLens = grow(table.payloadLens, Float64Array, next);
124 table.sizeVarLens = grow(table.sizeVarLens, Uint8Array, next);
125 table.outputOffsets = grow(table.outputOffsets, Float64Array, next);
126 table.outputHeaderLens = grow(table.outputHeaderLens, Uint16Array, next);
127 table.outputOrder = grow(table.outputOrder, Uint32Array, next);
128 
129 table.queuedForHeader = grow(table.queuedForHeader, Uint8Array, next);
130 table.ofsPinned = grow(table.ofsPinned, Uint8Array, next);
131 table.syntheticPayloads.length = next;
132 
133 // sizeVarBuf is capacity * 5
134 const oldSvBuf = table.sizeVarBuf;
135 table.sizeVarBuf = new Uint8Array(next * 5);
136 table.sizeVarBuf.set(oldSvBuf);
137 
138 const oldOidsRaw = table.oidsRaw;
139 table.oidsRaw = new Uint8Array(next * 20);
140 table.oidsRaw.set(oldOidsRaw);
141 
142 // baseSlots: new slots default to -1
143 const oldBaseSlots = table.baseSlots;
144 table.baseSlots = new Int32Array(next).fill(-1);
145 table.baseSlots.set(oldBaseSlots);
146 
147 // baseOidRaw: grow only if already allocated
148 if (table.baseOidRaw) {
149 const oldRaw = table.baseOidRaw;
150 table.baseOidRaw = new Uint8Array(next * 20);
151 table.baseOidRaw.set(oldRaw);
152 }
153 
154 table.capacity = next;
155}
156 
157/** Pack (packSlot, entryIndex) into a single number for Map<number, number>. */
158export function selectionKey(packSlot: number, entryIndex: number): number {
159 return packSlot * 0x1_0000_0000 + entryIndex;
160}
161 
162/**
163 * Copy the pack-position identity for a row.
164 *
165 * This helper keeps the per-row identity fields together so add/replace flows
166 * cannot forget to carry the raw OID bytes along with the new pack position.
167 */
168export function setSelectionEntryIdentity(
169 table: SelectionTable,
170 sel: number,
171 packSlot: number,
172 entryIndex: number,
173 idx: IdxView
174): void {
175 table.packSlots[sel] = packSlot;
176 table.entryIndices[sel] = entryIndex;
177 table.offsets[sel] = idx.offsets[entryIndex];
178 
179 const nextOffset = getNextOffsetByIndex(idx, entryIndex);
180 if (nextOffset === undefined) {
181 throw new Error(`rewrite: missing next offset for pack#${packSlot} entry#${entryIndex}`);
182 }
183 table.nextOffsets[sel] = nextOffset;
184 table.oidsRaw.set(idx.rawNames.subarray(entryIndex * 20, entryIndex * 20 + 20), sel * 20);
185 table.syntheticPayloads[sel] = undefined;
186}
187 
188/** Store the parsed pack header into the selection row. */
189export function storeSelectionHeader(
190 table: SelectionTable,
191 sel: number,
192 offset: number,
193 nextOffset: number,
194 header: PackHeaderEx
195): boolean {
196 table.typeCodes[sel] = header.type;
197 table.headerLens[sel] = header.headerLen;
198 
199 const payloadLength = nextOffset - offset - header.headerLen;
200 if (payloadLength < 0) {
201 return false;
202 }
203 table.payloadLens[sel] = payloadLength;
204 
205 const svStart = sel * 5;
206 table.sizeVarBuf.set(header.sizeVarBytes, svStart);
207 table.sizeVarLens[sel] = header.sizeVarBytes.length;
208 return true;
209}
210 
211/**
212 * Compare two selection slots by source-pack traversal order.
213 *
214 * Rewrite header reads and payload streaming both favor `(packSlot, offset)`
215 * ordering so the `SequentialReader` stays on a mostly forward path.
216 */
217export function compareSelectionSlots(
218 table: SelectionTable,
219 leftSel: number,
220 rightSel: number
221): number {
222 const packDiff = table.packSlots[leftSel] - table.packSlots[rightSel];
223 if (packDiff !== 0) return packDiff;
224 return table.offsets[leftSel] - table.offsets[rightSel];
225}
226 
227/** Sort selection slots in-place using source-pack traversal order. */
228export function sortSelectionSlots(
229 table: SelectionTable,
230 slots: Uint32Array | number[]
231): Uint32Array | number[] {
232 slots.sort((leftSel, rightSel) => compareSelectionSlots(table, leftSel, rightSel));
233 return slots;
234}
235 
236/**
237 * Selection dependencies are a single base chain per row, so checking whether
238 * `startSel` depends on `targetSel` is a bounded linked-list walk. The helper
239 * intentionally follows only `baseSlots`; callers use it in hot paths before
240 * any output graph has been allocated.
241 */
242export function selectionDependsOn(
243 table: SelectionTable,
244 startSel: number,
245 targetSel: number
246): boolean {
247 let cur = startSel;
248 for (let depth = 0; depth < table.count; depth++) {
249 const baseSel = table.baseSlots[cur];
250 if (baseSel < 0) return false;
251 if (baseSel === targetSel) return true;
252 cur = baseSel;
253 }
254 return false;
255}
256 
257type CopySelectionRowOptions = {
258 preserveTargetOfsPinned?: boolean;
259};
260 
261/**
262 * Copy all planner-phase row fields from one selection slot to another.
263 *
264 * This intentionally excludes layout outputs (`outputOffsets`,
265 * `outputHeaderLens`, `outputOrder`) because every row rewrite happens before
266 * topology and output sizing run.
267 *
268 * `baseSlots` and `baseOidRaw` move with the row because they are part of the
269 * row's resolved delta wiring, not derived output state. Callers that need
270 * extra semantics such as queue clearing or OFS pin merging still apply those
271 * adjustments explicitly after the copy.
272 */
273export function copySelectionRow(
274 table: SelectionTable,
275 targetSel: number,
276 sourceSel: number,
277 options?: CopySelectionRowOptions
278): void {
279 const targetPinned = table.ofsPinned[targetSel];
280 
281 table.packSlots[targetSel] = table.packSlots[sourceSel];
282 table.entryIndices[targetSel] = table.entryIndices[sourceSel];
283 table.offsets[targetSel] = table.offsets[sourceSel];
284 table.nextOffsets[targetSel] = table.nextOffsets[sourceSel];
285 table.oidsRaw.set(table.oidsRaw.subarray(sourceSel * 20, sourceSel * 20 + 20), targetSel * 20);
286 table.typeCodes[targetSel] = table.typeCodes[sourceSel];
287 table.headerLens[targetSel] = table.headerLens[sourceSel];
288 table.payloadLens[targetSel] = table.payloadLens[sourceSel];
289 table.sizeVarLens[targetSel] = table.sizeVarLens[sourceSel];
290 table.baseSlots[targetSel] = table.baseSlots[sourceSel];
291 table.sizeVarBuf.set(table.sizeVarBuf.subarray(sourceSel * 5, sourceSel * 5 + 5), targetSel * 5);
292 table.queuedForHeader[targetSel] = table.queuedForHeader[sourceSel];
293 table.ofsPinned[targetSel] = options?.preserveTargetOfsPinned
294 ? targetPinned
295 : table.ofsPinned[sourceSel];
296 table.syntheticPayloads[targetSel] = table.syntheticPayloads[sourceSel];
297 
298 if (table.baseOidRaw) {
299 table.baseOidRaw.set(
300 table.baseOidRaw.subarray(sourceSel * 20, sourceSel * 20 + 20),
301 targetSel * 20
302 );
303 }
304}
305 
306export function recordRewriteFailure(
307 options: RewriteOptions | undefined,
308 failure: RewriteFailure
309): void {
310 if (!options?.failure || options.failure.value) return;
311 options.failure.value = failure;
312}
313 
314export type PackReadState = {
315 pack: OrderedPackSnapshotEntry;
316 reader: SequentialReader;
317 wholePack?: Uint8Array;
318};
319 
320export function buildPackHeader(objectCount: number): Uint8Array {
321 const header = new Uint8Array(12);
322 header.set(new TextEncoder().encode("PACK"), 0);
323 const view = new DataView(header.buffer);
324 view.setUint32(4, 2);
325 view.setUint32(8, objectCount);
326 return header;
327}
328 
329export function countRewriteSubrequest(
330 log: Logger,
331 warnedFlags: Set<string>,
332 options: RewriteOptions | undefined,
333 flag: string,
334 details: Record<string, unknown>,
335 n?: number
336): boolean | void {
337 const withinBudget = options?.countSubrequest?.(n);
338 if (withinBudget === false && !warnedFlags.has(flag)) {
339 warnedFlags.add(flag);
340 log.warn("soft-budget-exhausted", details);
341 }
342 return withinBudget;
343}
344 
345function getRequiredLimiter(options?: RewriteOptions): Limiter {
346 if (!options?.limiter) {
347 throw new Error("rewrite: limiter required");
348 }
349 return options.limiter;
350}
351 
352async function loadWholePack(
353 env: Env,
354 pack: OrderedPackSnapshotEntry,
355 log: Logger,
356 warnedFlags: Set<string>,
357 options?: RewriteOptions
358): Promise<Uint8Array | undefined> {
359 return await readPackRange(env, pack.packKey, 0, pack.packBytes, {
360 limiter: getRequiredLimiter(options),
361 signal: options?.signal,
362 countSubrequest: (n?: number) =>
363 countRewriteSubrequest(
364 log,
365 warnedFlags,
366 options,
367 `rewrite-whole-pack:${pack.packKey}`,
368 { op: "r2:get-range", packKey: pack.packKey },
369 n
370 ),
371 });
372}
373 
374function createPackReadState(
375 env: Env,
376 pack: OrderedPackSnapshotEntry,
377 log: Logger,
378 warnedFlags: Set<string>,
379 options?: RewriteOptions
380): PackReadState {
381 const readerLog = createLogger(env.LOG_LEVEL, { service: "RewritePackReader" });
382 return {
383 pack,
384 reader: new SequentialReader(
385 env,
386 pack.packKey,
387 pack.packBytes,
388 DEFAULT_CHUNK_SIZE,
389 getRequiredLimiter(options),
390 (n?: number) =>
391 countRewriteSubrequest(
392 log,
393 warnedFlags,
394 options,
395 `rewrite-range:${pack.packKey}`,
396 { op: "r2:get-range", packKey: pack.packKey },
397 n
398 ),
399 readerLog,
400 options?.signal
401 ),
402 };
403}
404 
405export async function ensurePackReadState(
406 env: Env,
407 pack: OrderedPackSnapshotEntry,
408 packSlot: number,
409 readerStates: Map<number, PackReadState>,
410 log: Logger,
411 warnedFlags: Set<string>,
412 options?: RewriteOptions
413): Promise<PackReadState> {
414 const existing = readerStates.get(packSlot);
415 if (existing) return existing;
416 
417 const state = createPackReadState(env, pack, log, warnedFlags, options);
418 if (pack.packBytes <= WHOLE_PACK_MAX_BYTES) {
419 let loaded = 0;
420 for (const s of readerStates.values()) {
421 if (s.wholePack) loaded += s.wholePack.length;
422 }
423 if (loaded + pack.packBytes <= WHOLE_PACK_TOTAL_BUDGET) {
424 state.wholePack = await loadWholePack(env, pack, log, warnedFlags, options);
425 }
426 }
427 readerStates.set(packSlot, state);
428 return state;
429}
430 
431/**
432 * Read a pack entry header at the given byte offset.
433 * For the wholePack path the subarray is stable; for the SequentialReader
434 * path the returned `sizeVarBytes` references the reader buffer and must
435 * be copied before the next preload.
436 */
437export async function readSelectedHeader(
438 state: PackReadState,
439 offset: number
440): Promise<PackHeaderEx | undefined> {
441 if (state.wholePack) {
442 return readPackHeaderExFromBuf(state.wholePack, offset);
443 }
444 
445 const bytesLeft = Math.max(0, state.pack.packBytes - offset);
446 const headerBytes = await state.reader.readRange(offset, Math.min(HEADER_READ_BYTES, bytesLeft));
447 return readPackHeaderExFromBuf(headerBytes, 0);
448}
449 
450/** Search packs in snapshot order; first match wins (duplicate selection). */
451export function resolveOrderedEntryByOid(
452 snapshot: OrderedPackSnapshot,
453 oid: string | Uint8Array
454): { packSlot: number; pack: OrderedPackSnapshotEntry; entryIndex: number } | undefined {
455 const candidate = findFirstPackedObjectCandidate(snapshot.packs, oid);
456 if (!candidate) return undefined;
457 
458 return {
459 packSlot: candidate.packSlot,
460 pack: snapshot.packs[candidate.packSlot]!,
461 entryIndex: candidate.objectIndex,
462 };
463}