Skip to content
File

Blob: src/worker/git/operations/fetch/refClosure.ts

typescript524 lines
1import type { PackRefSnapshotEntry } from "@/worker/git/pack/refIndex";
2 
3import { bytesToHex, createLogger, hexToBytes, isValidOid } from "@/worker/common";
4import { findOidIndexFromBytes } from "@/worker/git/object-store";
5import {
6 getPackRefRawRefAt,
7 getPackRefTypeCode,
8 visitPackRefRawRefsAt,
9} from "@/worker/git/pack/refIndex";
10 
11const HAVE_CAP = 128;
12const MAINLINE_ENRICHMENT_BUDGET = 20;
13const CLOSURE_TIMEOUT_MS = 49_000;
14const MISSING_REF_CAP = 1024;
15const OID_BYTES = 20;
16 
17export type RefClosureStats = {
18 indexedObjects: number;
19 queued: number;
20 seen: number;
21 needed: number;
22 missing: number;
23 edgeVisits: number;
24 duplicateQueueSkips: number;
25};
26 
27export type RefClosureResult =
28 | {
29 type: "Ready";
30 neededOids: string[];
31 ackOids: string[];
32 stats: RefClosureStats;
33 }
34 | {
35 type: "BudgetExceeded";
36 neededOids: string[];
37 ackOids: string[];
38 reason: "timeout" | "missing-ref-budget";
39 stats: RefClosureStats;
40 };
41 
42type ClosureIndex = {
43 packBaseOrdinals: Uint32Array;
44 objectCount: number;
45};
46 
47type LocatedObject = {
48 packSlot: number;
49 oidIndex: number;
50 ordinal: number;
51};
52 
53type CommonHave = {
54 oid: string;
55 located: LocatedObject;
56};
57 
58type LocatedObjectQueue = {
59 packSlots: Uint32Array;
60 oidIndices: Uint32Array;
61 cursor: number;
62 count: number;
63};
64 
65// Mainline enrichment is intentionally tiny, so raw OIDs keep that side walk
66// simple without affecting the bounded final closure queue below.
67type RawOidQueue = {
68 rawOids: Uint8Array;
69 cursor: number;
70 count: number;
71};
72 
73function buildClosureIndex(packs: PackRefSnapshotEntry[]): ClosureIndex {
74 const packBaseOrdinals = new Uint32Array(packs.length + 1);
75 let objectCount = 0;
76 
77 for (let packSlot = 0; packSlot < packs.length; packSlot++) {
78 packBaseOrdinals[packSlot] = objectCount;
79 objectCount += packs[packSlot]!.idx.count;
80 }
81 packBaseOrdinals[packs.length] = objectCount;
82 
83 return { packBaseOrdinals, objectCount };
84}
85 
86function locateObject(
87 packs: PackRefSnapshotEntry[],
88 closureIndex: ClosureIndex,
89 rawOid: Uint8Array,
90 rawOidStart: number
91): LocatedObject | undefined {
92 for (let packSlot = 0; packSlot < packs.length; packSlot++) {
93 const oidIndex = findOidIndexFromBytes(packs[packSlot]!.idx, rawOid, rawOidStart);
94 if (oidIndex < 0) continue;
95 return {
96 packSlot,
97 oidIndex,
98 ordinal: closureIndex.packBaseOrdinals[packSlot]! + oidIndex,
99 };
100 }
101 return undefined;
102}
103 
104function createLocatedObjectQueue(initialEntries: number): LocatedObjectQueue {
105 const initialCapacity = Math.max(initialEntries, 16);
106 return {
107 packSlots: new Uint32Array(initialCapacity),
108 oidIndices: new Uint32Array(initialCapacity),
109 cursor: 0,
110 count: 0,
111 };
112}
113 
114function ensureLocatedObjectQueueCapacity(queue: LocatedObjectQueue, nextCount: number): void {
115 if (nextCount <= queue.packSlots.length) return;
116 
117 let nextCapacity = queue.packSlots.length;
118 while (nextCapacity < nextCount) nextCapacity *= 2;
119 
120 const nextPackSlots = new Uint32Array(nextCapacity);
121 const nextOidIndices = new Uint32Array(nextCapacity);
122 nextPackSlots.set(queue.packSlots);
123 nextOidIndices.set(queue.oidIndices);
124 queue.packSlots = nextPackSlots;
125 queue.oidIndices = nextOidIndices;
126}
127 
128function pushLocatedObject(queue: LocatedObjectQueue, located: LocatedObject): void {
129 ensureLocatedObjectQueueCapacity(queue, queue.count + 1);
130 queue.packSlots[queue.count] = located.packSlot;
131 queue.oidIndices[queue.count] = located.oidIndex;
132 queue.count++;
133}
134 
135function createRawOidQueue(initialEntries: number): RawOidQueue {
136 return {
137 rawOids: new Uint8Array(Math.max(initialEntries, 16) * OID_BYTES),
138 cursor: 0,
139 count: 0,
140 };
141}
142 
143function ensureRawOidQueueCapacity(queue: RawOidQueue, nextCount: number): void {
144 if (nextCount * OID_BYTES <= queue.rawOids.byteLength) return;
145 
146 let nextCapacity = queue.rawOids.byteLength / OID_BYTES;
147 while (nextCapacity < nextCount) nextCapacity *= 2;
148 
149 const nextRawOids = new Uint8Array(nextCapacity * OID_BYTES);
150 nextRawOids.set(queue.rawOids);
151 queue.rawOids = nextRawOids;
152}
153 
154function pushRawOid(queue: RawOidQueue, rawOid: Uint8Array, rawOidStart: number): void {
155 ensureRawOidQueueCapacity(queue, queue.count + 1);
156 queue.rawOids.set(rawOid.subarray(rawOidStart, rawOidStart + OID_BYTES), queue.count * OID_BYTES);
157 queue.count++;
158}
159 
160function findCommonHavesInSnapshot(
161 packs: PackRefSnapshotEntry[],
162 closureIndex: ClosureIndex,
163 haves: string[]
164): CommonHave[] {
165 const cappedHaves = haves.slice(0, HAVE_CAP);
166 const found: CommonHave[] = [];
167 const seenFlags = new Uint8Array(closureIndex.objectCount);
168 
169 for (const have of cappedHaves) {
170 const oid = have.toLowerCase();
171 if (!isValidOid(oid)) continue;
172 
173 const rawOid = hexToBytes(oid);
174 const located = locateObject(packs, closureIndex, rawOid, 0);
175 if (!located) continue;
176 if (seenFlags[located.ordinal]) continue;
177 
178 seenFlags[located.ordinal] = 1;
179 found.push({ oid, located });
180 }
181 
182 return found;
183}
184 
185function enqueueLocatedObject(
186 queue: LocatedObjectQueue,
187 queuedFlags: Uint8Array,
188 located: LocatedObject,
189 duplicateQueueSkips: { value: number }
190): void {
191 if (queuedFlags[located.ordinal]) {
192 duplicateQueueSkips.value++;
193 return;
194 }
195 
196 queuedFlags[located.ordinal] = 1;
197 pushLocatedObject(queue, located);
198}
199 
200function recordMissingOid(args: {
201 oid: string;
202 missingSeen: Set<string>;
203 missingNeeded: Set<string>;
204 includeNeeded: boolean;
205}): "recorded" | "duplicate" | "budget-exceeded" {
206 if (args.missingSeen.has(args.oid)) return "duplicate";
207 if (args.missingSeen.size >= MISSING_REF_CAP) return "budget-exceeded";
208 
209 args.missingSeen.add(args.oid);
210 if (args.includeNeeded) {
211 args.missingNeeded.add(args.oid);
212 }
213 return "recorded";
214}
215 
216function addNeededObject(
217 located: LocatedObject,
218 neededFlags: Uint8Array,
219 neededPackSlots: { values: Uint32Array },
220 neededOidIndices: { values: Uint32Array },
221 neededCountRef: { value: number }
222): void {
223 if (neededFlags[located.ordinal]) return;
224 neededFlags[located.ordinal] = 1;
225 
226 if (neededCountRef.value >= neededPackSlots.values.length) {
227 const nextLength = Math.max(neededPackSlots.values.length * 2, 16);
228 const nextPackSlots = new Uint32Array(nextLength);
229 const nextOidIndices = new Uint32Array(nextLength);
230 nextPackSlots.set(neededPackSlots.values);
231 nextOidIndices.set(neededOidIndices.values);
232 neededPackSlots.values = nextPackSlots;
233 neededOidIndices.values = nextOidIndices;
234 }
235 
236 neededPackSlots.values[neededCountRef.value] = located.packSlot;
237 neededOidIndices.values[neededCountRef.value] = located.oidIndex;
238 neededCountRef.value++;
239}
240 
241function buildNeededOids(args: {
242 packs: PackRefSnapshotEntry[];
243 neededPackSlots: Uint32Array;
244 neededOidIndices: Uint32Array;
245 neededCount: number;
246 missingNeeded: Set<string>;
247}): string[] {
248 const neededOids: string[] = [];
249 for (let index = 0; index < args.neededCount; index++) {
250 const packSlot = args.neededPackSlots[index]!;
251 const oidIndex = args.neededOidIndices[index]!;
252 const rawNames = args.packs[packSlot]!.idx.rawNames;
253 const oidStart = oidIndex * OID_BYTES;
254 neededOids.push(bytesToHex(rawNames.subarray(oidStart, oidStart + OID_BYTES)));
255 }
256 
257 for (const oid of args.missingNeeded) {
258 neededOids.push(oid);
259 }
260 return neededOids;
261}
262 
263function buildRefClosureStats(args: {
264 closureIndex: ClosureIndex;
265 queue: LocatedObjectQueue;
266 seenCount: number;
267 neededCount: number;
268 missingSeen: Set<string>;
269 missingNeeded: Set<string>;
270 edgeVisits: number;
271 duplicateQueueSkips: number;
272}): RefClosureStats {
273 return {
274 indexedObjects: args.closureIndex.objectCount,
275 queued: args.queue.count,
276 seen: args.seenCount,
277 needed: args.neededCount + args.missingNeeded.size,
278 missing: args.missingSeen.size,
279 edgeVisits: args.edgeVisits,
280 duplicateQueueSkips: args.duplicateQueueSkips,
281 };
282}
283 
284export async function computeNeededFromPackRefs(args: {
285 logLevel?: string;
286 repoId: string;
287 packs: PackRefSnapshotEntry[];
288 wants: string[];
289 haves: string[];
290 onProgress?: (message: string) => void;
291}): Promise<RefClosureResult> {
292 const log = createLogger(args.logLevel, { service: "RefClosure", repoId: args.repoId });
293 const startTime = Date.now();
294 const closureIndex = buildClosureIndex(args.packs);
295 const stopFlags = new Uint8Array(closureIndex.objectCount);
296 const missingStop = new Set<string>();
297 let stopCount = 0;
298 
299 args.onProgress?.("Finding common commits...\n");
300 const commonHaves = findCommonHavesInSnapshot(args.packs, closureIndex, args.haves);
301 const ackOids = commonHaves.map((have) => have.oid);
302 for (const have of commonHaves) {
303 if (stopFlags[have.located.ordinal]) continue;
304 stopFlags[have.located.ordinal] = 1;
305 stopCount++;
306 }
307 
308 if (commonHaves.length > 0 && commonHaves.length < 10) {
309 const mainlineQueue = createRawOidQueue(commonHaves.length + MAINLINE_ENRICHMENT_BUDGET);
310 for (const have of commonHaves) {
311 const rawNames = args.packs[have.located.packSlot]!.idx.rawNames;
312 pushRawOid(mainlineQueue, rawNames, have.located.oidIndex * OID_BYTES);
313 }
314 
315 let walked = 0;
316 while (mainlineQueue.cursor < mainlineQueue.count && walked < MAINLINE_ENRICHMENT_BUDGET) {
317 if (Date.now() - startTime > 2_000) break;
318 
319 const oidStart = mainlineQueue.cursor * OID_BYTES;
320 mainlineQueue.cursor++;
321 const located = locateObject(args.packs, closureIndex, mainlineQueue.rawOids, oidStart);
322 if (!located) continue;
323 
324 // Commit sidecars store refs as [tree, first-parent, ...remaining-parents].
325 if (getPackRefTypeCode(args.packs[located.packSlot]!.refs, located.oidIndex) !== 1) {
326 continue;
327 }
328 
329 const firstParent = getPackRefRawRefAt(
330 args.packs[located.packSlot]!.refs,
331 located.oidIndex,
332 1
333 );
334 if (!firstParent) continue;
335 
336 const parentLocated = locateObject(args.packs, closureIndex, firstParent, 0);
337 if (parentLocated) {
338 if (stopFlags[parentLocated.ordinal]) continue;
339 stopFlags[parentLocated.ordinal] = 1;
340 stopCount++;
341 pushRawOid(mainlineQueue, firstParent, 0);
342 walked++;
343 continue;
344 }
345 
346 const parentOid = bytesToHex(firstParent);
347 if (missingStop.has(parentOid)) continue;
348 missingStop.add(parentOid);
349 stopCount++;
350 walked++;
351 }
352 
353 log.debug("stream:plan:mainline-enriched", { stopSize: stopCount, walked });
354 }
355 
356 args.onProgress?.("Selecting objects to send...\n");
357 
358 const seenFlags = new Uint8Array(closureIndex.objectCount);
359 const queuedFlags = new Uint8Array(closureIndex.objectCount);
360 const neededFlags = new Uint8Array(closureIndex.objectCount);
361 const missingSeen = new Set<string>();
362 const missingNeeded = new Set<string>();
363 const queue = createLocatedObjectQueue(args.wants.length);
364 const neededPackSlots = { values: new Uint32Array(Math.max(args.wants.length, 16)) };
365 const neededOidIndices = { values: new Uint32Array(Math.max(args.wants.length, 16)) };
366 const neededCount = { value: 0 };
367 let seenCount = 0;
368 let edgeVisits = 0;
369 const duplicateQueueSkips = { value: 0 };
370 
371 const buildStats = () =>
372 buildRefClosureStats({
373 closureIndex,
374 queue,
375 seenCount,
376 neededCount: neededCount.value,
377 missingSeen,
378 missingNeeded,
379 edgeVisits,
380 duplicateQueueSkips: duplicateQueueSkips.value,
381 });
382 
383 const buildBudgetExceededResult = (
384 reason: "timeout" | "missing-ref-budget"
385 ): RefClosureResult => ({
386 type: "BudgetExceeded",
387 reason,
388 neededOids: buildNeededOids({
389 packs: args.packs,
390 neededPackSlots: neededPackSlots.values,
391 neededOidIndices: neededOidIndices.values,
392 neededCount: neededCount.value,
393 missingNeeded,
394 }),
395 ackOids,
396 stats: buildStats(),
397 });
398 
399 for (const want of args.wants) {
400 const normalized = want.toLowerCase();
401 if (!isValidOid(normalized)) {
402 const result = recordMissingOid({
403 oid: normalized,
404 missingSeen,
405 missingNeeded,
406 includeNeeded: true,
407 });
408 if (result === "budget-exceeded") return buildBudgetExceededResult("missing-ref-budget");
409 continue;
410 }
411 
412 const rawOid = hexToBytes(normalized);
413 const located = locateObject(args.packs, closureIndex, rawOid, 0);
414 if (!located) {
415 const result = recordMissingOid({
416 oid: normalized,
417 missingSeen,
418 missingNeeded,
419 includeNeeded: true,
420 });
421 if (result === "budget-exceeded") return buildBudgetExceededResult("missing-ref-budget");
422 continue;
423 }
424 
425 enqueueLocatedObject(queue, queuedFlags, located, duplicateQueueSkips);
426 }
427 
428 log.info("stream:plan:closure-start", {
429 wants: args.wants.length,
430 haves: args.haves.length,
431 ackOids: ackOids.length,
432 indexedObjects: closureIndex.objectCount,
433 stopSet: stopCount,
434 queued: queue.count,
435 });
436 
437 while (queue.cursor < queue.count) {
438 if (Date.now() - startTime > CLOSURE_TIMEOUT_MS) {
439 return buildBudgetExceededResult("timeout");
440 }
441 
442 const packSlot = queue.packSlots[queue.cursor]!;
443 const oidIndex = queue.oidIndices[queue.cursor]!;
444 queue.cursor++;
445 const ordinal = closureIndex.packBaseOrdinals[packSlot]! + oidIndex;
446 const located: LocatedObject = { packSlot, oidIndex, ordinal };
447 
448 if (seenFlags[located.ordinal]) continue;
449 seenFlags[located.ordinal] = 1;
450 seenCount++;
451 
452 if (stopFlags[located.ordinal]) {
453 log.debug("stream:plan:hit-stop", {
454 packKey: args.packs[located.packSlot]!.packKey,
455 oidIndex: located.oidIndex,
456 });
457 continue;
458 }
459 
460 addNeededObject(located, neededFlags, neededPackSlots, neededOidIndices, neededCount);
461 let budgetExceeded = false;
462 visitPackRefRawRefsAt(
463 args.packs[located.packSlot]!.refs,
464 located.oidIndex,
465 (rawRefs, start) => {
466 edgeVisits++;
467 if (budgetExceeded) return;
468 
469 const refLocated = locateObject(args.packs, closureIndex, rawRefs, start);
470 if (refLocated) {
471 enqueueLocatedObject(queue, queuedFlags, refLocated, duplicateQueueSkips);
472 return;
473 }
474 
475 const oid = bytesToHex(rawRefs.subarray(start, start + OID_BYTES));
476 const includeNeeded = !missingStop.has(oid);
477 const result = recordMissingOid({
478 oid,
479 missingSeen,
480 missingNeeded,
481 includeNeeded,
482 });
483 if (result === "budget-exceeded") {
484 budgetExceeded = true;
485 return;
486 }
487 
488 if (result === "recorded" && !includeNeeded) {
489 log.debug("stream:plan:hit-stop", { oid });
490 }
491 }
492 );
493 if (budgetExceeded) {
494 return buildBudgetExceededResult("missing-ref-budget");
495 }
496 }
497 
498 const stats = buildStats();
499 log.info("stream:plan:closure-complete", {
500 needed: stats.needed,
501 seen: stats.seen,
502 queued: stats.queued,
503 missing: stats.missing,
504 edgeVisits: stats.edgeVisits,
505 duplicateQueueSkips: stats.duplicateQueueSkips,
506 stopSet: stopCount,
507 ackOids: ackOids.length,
508 timeMs: Date.now() - startTime,
509 });
510 
511 return {
512 type: "Ready",
513 neededOids: buildNeededOids({
514 packs: args.packs,
515 neededPackSlots: neededPackSlots.values,
516 neededOidIndices: neededOidIndices.values,
517 neededCount: neededCount.value,
518 missingNeeded,
519 }),
520 ackOids,
521 stats,
522 };
523}