File
Blob: src/worker/git/operations/fetch/refClosure.ts
| 1 | import type { PackRefSnapshotEntry } from "@/worker/git/pack/refIndex"; |
| 2 | |
| 3 | import { bytesToHex, createLogger, hexToBytes, isValidOid } from "@/worker/common"; |
| 4 | import { findOidIndexFromBytes } from "@/worker/git/object-store"; |
| 5 | import { |
| 6 | getPackRefRawRefAt, |
| 7 | getPackRefTypeCode, |
| 8 | visitPackRefRawRefsAt, |
| 9 | } from "@/worker/git/pack/refIndex"; |
| 10 | |
| 11 | const HAVE_CAP = 128; |
| 12 | const MAINLINE_ENRICHMENT_BUDGET = 20; |
| 13 | const CLOSURE_TIMEOUT_MS = 49_000; |
| 14 | const MISSING_REF_CAP = 1024; |
| 15 | const OID_BYTES = 20; |
| 16 | |
| 17 | export 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 | |
| 27 | export 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 | |
| 42 | type ClosureIndex = { |
| 43 | packBaseOrdinals: Uint32Array; |
| 44 | objectCount: number; |
| 45 | }; |
| 46 | |
| 47 | type LocatedObject = { |
| 48 | packSlot: number; |
| 49 | oidIndex: number; |
| 50 | ordinal: number; |
| 51 | }; |
| 52 | |
| 53 | type CommonHave = { |
| 54 | oid: string; |
| 55 | located: LocatedObject; |
| 56 | }; |
| 57 | |
| 58 | type 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. |
| 67 | type RawOidQueue = { |
| 68 | rawOids: Uint8Array; |
| 69 | cursor: number; |
| 70 | count: number; |
| 71 | }; |
| 72 | |
| 73 | function 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 | |
| 86 | function 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 | |
| 104 | function 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 | |
| 114 | function 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 | |
| 128 | function 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 | |
| 135 | function 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 | |
| 143 | function 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 | |
| 154 | function 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 | |
| 160 | function 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 | |
| 185 | function 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 | |
| 200 | function 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 | |
| 216 | function 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 | |
| 241 | function 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 | |
| 263 | function 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 | |
| 284 | export 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 | } |