File
Blob: src/worker/git/pack/rewrite/shared.ts
| 1 | import type { |
| 2 | OrderedPackSnapshot, |
| 3 | OrderedPackSnapshotEntry, |
| 4 | } from "@/worker/git/operations/fetch/types"; |
| 5 | import type { Logger } from "@/worker/common/logger"; |
| 6 | import type { Limiter } from "@/worker/git/operations/limits"; |
| 7 | import type { IdxView } from "@/worker/git/object-store/types"; |
| 8 | import type { PackHeaderEx } from "../packMeta"; |
| 9 | |
| 10 | import { createLogger } from "@/worker/common"; |
| 11 | import { findFirstPackedObjectCandidate, getNextOffsetByIndex } from "@/worker/git/object-store"; |
| 12 | import { SequentialReader } from "@/worker/git/pack/indexer/resolve/reader"; |
| 13 | import { readPackHeaderExFromBuf, readPackRange } from "../packMeta"; |
| 14 | |
| 15 | export const HEADER_READ_BYTES = 128; |
| 16 | export const DEFAULT_CHUNK_SIZE = 4_194_304; |
| 17 | export const WHOLE_PACK_MAX_BYTES = 8 * 1024 * 1024; |
| 18 | export const WHOLE_PACK_TOTAL_BUDGET = 32 * 1024 * 1024; |
| 19 | export const HEADER_STABILITY_CAP = 16; |
| 20 | |
| 21 | export type RewriteFailure = { |
| 22 | reason: string; |
| 23 | retryable: boolean; |
| 24 | details?: Record<string, unknown>; |
| 25 | }; |
| 26 | |
| 27 | export type RewriteFailureRecorder = { |
| 28 | value?: RewriteFailure; |
| 29 | }; |
| 30 | |
| 31 | export 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 | */ |
| 47 | export 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 | |
| 78 | export 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. */ |
| 104 | export 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>. */ |
| 158 | export 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 | */ |
| 168 | export 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. */ |
| 189 | export 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 | */ |
| 217 | export 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. */ |
| 228 | export 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 | */ |
| 242 | export 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 | |
| 257 | type 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 | */ |
| 273 | export 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 | |
| 306 | export 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 | |
| 314 | export type PackReadState = { |
| 315 | pack: OrderedPackSnapshotEntry; |
| 316 | reader: SequentialReader; |
| 317 | wholePack?: Uint8Array; |
| 318 | }; |
| 319 | |
| 320 | export 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 | |
| 329 | export 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 | |
| 345 | function getRequiredLimiter(options?: RewriteOptions): Limiter { |
| 346 | if (!options?.limiter) { |
| 347 | throw new Error("rewrite: limiter required"); |
| 348 | } |
| 349 | return options.limiter; |
| 350 | } |
| 351 | |
| 352 | async 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 | |
| 374 | function 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 | |
| 405 | export 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 | */ |
| 437 | export 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). */ |
| 451 | export 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 | } |