File
Blob: src/worker/git/pack/rewrite/selection.ts
| 1 | import type { OrderedPackSnapshot } from "@/worker/git/operations/fetch/types"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | |
| 4 | import { hexToBytes } from "@/worker/common"; |
| 5 | import { |
| 6 | createSelectedOidLookup, |
| 7 | type DuplicateHeaderCache, |
| 8 | type SelectionStats, |
| 9 | } from "./ownership"; |
| 10 | import { collapseUnsafeRedirectOwners, compactDeadSlots } from "./selectionCompact"; |
| 11 | import { |
| 12 | collectRetainedRedirectsNeedingBaseResolution, |
| 13 | resolveRetainedRedirectBase, |
| 14 | } from "./selectionRetained"; |
| 15 | import { addEntry, readHeaderAndResolveBase, type HeaderResolveResult } from "./selectionResolve"; |
| 16 | import { |
| 17 | allocateSelectionTable, |
| 18 | resolveOrderedEntryByOid, |
| 19 | sortSelectionSlots, |
| 20 | type PackReadState, |
| 21 | type RewriteOptions, |
| 22 | type SelectionTable, |
| 23 | } from "./shared"; |
| 24 | |
| 25 | export type BuildSelectionResult = { |
| 26 | table: SelectionTable; |
| 27 | readerStates: Map<number, PackReadState>; |
| 28 | addedDeltaBases: number; |
| 29 | }; |
| 30 | |
| 31 | export async function buildSelection( |
| 32 | env: Env, |
| 33 | snapshot: OrderedPackSnapshot, |
| 34 | neededOids: string[], |
| 35 | log: Logger, |
| 36 | warnedFlags: Set<string>, |
| 37 | options?: RewriteOptions |
| 38 | ): Promise<BuildSelectionResult | undefined> { |
| 39 | const table = allocateSelectionTable(Math.max(neededOids.length, 16)); |
| 40 | const readerStates = new Map<number, PackReadState>(); |
| 41 | /** Maps selectionKey(packSlot, entryIndex) → selection slot for exact pack-position dedup. */ |
| 42 | const dedupMap = new Map<number, number>(); |
| 43 | const oidOwners = createSelectedOidLookup(table.capacity); |
| 44 | const duplicateHeaderCache: DuplicateHeaderCache = new Map(); |
| 45 | const stats: SelectionStats = { |
| 46 | duplicateRedirects: 0, |
| 47 | duplicateOwnerUpgrades: 0, |
| 48 | duplicateOfsOwnerTakeovers: 0, |
| 49 | duplicateHeaderProbes: 0, |
| 50 | }; |
| 51 | |
| 52 | const resolveStart = Date.now(); |
| 53 | |
| 54 | // Resolve all needed OIDs to their selected pack slot and entry index. |
| 55 | for (const oid of neededOids) { |
| 56 | if (options?.signal?.aborted) return undefined; |
| 57 | |
| 58 | const oidBytes = hexToBytes(oid); |
| 59 | const location = resolveOrderedEntryByOid(snapshot, oidBytes); |
| 60 | if (!location) { |
| 61 | log.warn("rewrite:missing-needed-object", { oid }); |
| 62 | return undefined; |
| 63 | } |
| 64 | |
| 65 | addEntry(table, dedupMap, location.packSlot, location.entryIndex, location.pack.idx); |
| 66 | } |
| 67 | |
| 68 | const initialResolvedCount = table.count; |
| 69 | const resolveMs = Date.now() - resolveStart; |
| 70 | |
| 71 | // --- Early passthrough: skip header scan + base chase when every object |
| 72 | // in a single pack is already selected. The passthrough stream path |
| 73 | // reads raw pack bytes and needs none of the header data. ---------- |
| 74 | if (snapshot.packs.length === 1 && table.count === snapshot.packs[0].idx.count) { |
| 75 | log.info("rewrite:selection", { |
| 76 | requestedOids: neededOids.length, |
| 77 | selectedEntries: table.count, |
| 78 | addedDeltaBases: 0, |
| 79 | resolveMs, |
| 80 | headerReadMs: 0, |
| 81 | baseChaseIterations: 0, |
| 82 | }); |
| 83 | return { table, readerStates, addedDeltaBases: 0 }; |
| 84 | } |
| 85 | |
| 86 | const headerStart = Date.now(); |
| 87 | |
| 88 | // Counter for selected delta entries replaced with a full-object duplicate |
| 89 | // or redirected to an already-selected full owner for the same OID. |
| 90 | let duplicateCanonicalizations = 0; |
| 91 | |
| 92 | // Duplicate-OID rows that will be compacted out after selection completes. |
| 93 | // Most entries redirect current sel → owner sel; OFS-pinned takeovers invert |
| 94 | // that and retire the previous owner instead. |
| 95 | const deadSlots = new Map<number, number>(); |
| 96 | |
| 97 | /** Collect the result of readHeaderAndResolveBase. Returns false on failure. */ |
| 98 | function collectResult(sel: number, result: HeaderResolveResult): boolean { |
| 99 | if (!result.ok) return false; |
| 100 | if (result.ofsBaseCanonicalized) duplicateCanonicalizations++; |
| 101 | if (result.redirectTo !== undefined) { |
| 102 | deadSlots.set(sel, result.redirectTo); |
| 103 | stats.duplicateRedirects++; |
| 104 | } |
| 105 | if (result.supersedeSel !== undefined) { |
| 106 | deadSlots.delete(sel); |
| 107 | deadSlots.set(result.supersedeSel, sel); |
| 108 | stats.duplicateRedirects++; |
| 109 | } |
| 110 | return true; |
| 111 | } |
| 112 | |
| 113 | async function processHeaderBatch(batch: Uint32Array | number[]): Promise<boolean> { |
| 114 | for (const sel of batch) { |
| 115 | if (options?.signal?.aborted) return false; |
| 116 | const result = await readHeaderAndResolveBase( |
| 117 | table, |
| 118 | sel, |
| 119 | snapshot, |
| 120 | readerStates, |
| 121 | dedupMap, |
| 122 | oidOwners, |
| 123 | duplicateHeaderCache, |
| 124 | stats, |
| 125 | secondaryQueue, |
| 126 | env, |
| 127 | log, |
| 128 | warnedFlags, |
| 129 | options |
| 130 | ); |
| 131 | if (!collectResult(sel, result)) return false; |
| 132 | } |
| 133 | return true; |
| 134 | } |
| 135 | |
| 136 | async function drainSecondaryQueue(): Promise<boolean> { |
| 137 | while (secondaryQueue.length > 0) { |
| 138 | baseChaseIterations++; |
| 139 | if (options?.signal?.aborted) return false; |
| 140 | |
| 141 | const batch = secondaryQueue; |
| 142 | secondaryQueue = []; |
| 143 | |
| 144 | // Keep each pass mostly forward-only within a pack so the reader can |
| 145 | // stay on a small sliding window instead of bouncing around R2. |
| 146 | sortSelectionSlots(table, batch); |
| 147 | if (!(await processHeaderBatch(batch))) return false; |
| 148 | } |
| 149 | return true; |
| 150 | } |
| 151 | |
| 152 | // Sort once, then read object headers in offset order. |
| 153 | // Sorting by (packSlot, offset) maximizes SequentialReader locality. |
| 154 | const sortedSels = buildSortedIndex(table.count); |
| 155 | let secondaryQueue: number[] = []; |
| 156 | |
| 157 | if (!(await processHeaderBatch(sortedSels))) return undefined; |
| 158 | |
| 159 | // Chase delta bases until no new bases are discovered. |
| 160 | let baseChaseIterations = 0; |
| 161 | if (!(await drainSecondaryQueue())) return undefined; |
| 162 | |
| 163 | // Dead-slot pruning may keep some redirected duplicate deltas live when the |
| 164 | // selected owner still depends on their exact row. Those redirected rows |
| 165 | // returned early before resolving their own base chains, so finish that work |
| 166 | // now before compaction decides which duplicates remain in the output. |
| 167 | // |
| 168 | // Footgun: "redirected" does not mean "safe to ignore". A redirected row is |
| 169 | // only dead if compaction really removes it. Once dead-slot pruning decides |
| 170 | // the row must stay live, streaming will visit it like any other row and |
| 171 | // therefore requires `baseSlots[sel]` to be fully wired first. |
| 172 | let retainedRedirectResolutions = 0; |
| 173 | while (deadSlots.size > 0) { |
| 174 | const retainedRedirects = collectRetainedRedirectsNeedingBaseResolution(table, deadSlots); |
| 175 | if (retainedRedirects.length === 0) break; |
| 176 | |
| 177 | sortSelectionSlots(table, retainedRedirects); |
| 178 | |
| 179 | for (const sel of retainedRedirects) { |
| 180 | if (options?.signal?.aborted) return undefined; |
| 181 | const resolved = await resolveRetainedRedirectBase( |
| 182 | table, |
| 183 | sel, |
| 184 | snapshot, |
| 185 | readerStates, |
| 186 | dedupMap, |
| 187 | duplicateHeaderCache, |
| 188 | secondaryQueue, |
| 189 | env, |
| 190 | log, |
| 191 | warnedFlags, |
| 192 | options |
| 193 | ); |
| 194 | if (!resolved) return undefined; |
| 195 | retainedRedirectResolutions++; |
| 196 | } |
| 197 | |
| 198 | if (!(await drainSecondaryQueue())) return undefined; |
| 199 | } |
| 200 | |
| 201 | // Some redirected duplicates only stayed live because the chosen owner still |
| 202 | // depended on their exact row. Keeping both rows in the final pack makes the |
| 203 | // rewrite fetchable again, but Git still rejects the output because the same |
| 204 | // object OID appears twice. Collapse those cases back to one live row by |
| 205 | // rewriting the owner slot to stream the retained duplicate's encoding. |
| 206 | // |
| 207 | // Footgun: the owner slot index must stay stable here because children may |
| 208 | // already point at it. Only the row's source pack position and header/base |
| 209 | // metadata change; the selection slot itself remains the canonical owner. |
| 210 | let collapsedUnsafeRedirectOwners = 0; |
| 211 | if (deadSlots.size > 0) { |
| 212 | collapsedUnsafeRedirectOwners = collapseUnsafeRedirectOwners(table, deadSlots, log); |
| 213 | } |
| 214 | |
| 215 | // --- Compact dead duplicate-OID slots out of the table. Most redirected |
| 216 | // duplicates can collapse onto their surviving owner. The retained- |
| 217 | // redirect repair above already rewrote the few topology-sensitive cases |
| 218 | // back to one live row before this final compaction pass runs. |
| 219 | if (deadSlots.size > 0) { |
| 220 | compactDeadSlots(table, deadSlots, log); |
| 221 | } |
| 222 | |
| 223 | const headerReadMs = Date.now() - headerStart; |
| 224 | const addedDeltaBases = table.count - initialResolvedCount; |
| 225 | |
| 226 | log.info("rewrite:selection", { |
| 227 | requestedOids: neededOids.length, |
| 228 | selectedEntries: table.count, |
| 229 | addedDeltaBases, |
| 230 | duplicateCanonicalizations, |
| 231 | ownerLookupEntries: oidOwners.count, |
| 232 | duplicateRedirects: stats.duplicateRedirects, |
| 233 | duplicateOwnerUpgrades: stats.duplicateOwnerUpgrades, |
| 234 | duplicateOfsOwnerTakeovers: stats.duplicateOfsOwnerTakeovers, |
| 235 | duplicateHeaderProbes: stats.duplicateHeaderProbes, |
| 236 | retainedRedirectResolutions, |
| 237 | collapsedUnsafeRedirectOwners, |
| 238 | resolveMs, |
| 239 | headerReadMs, |
| 240 | baseChaseIterations, |
| 241 | }); |
| 242 | |
| 243 | return { table, readerStates, addedDeltaBases }; |
| 244 | |
| 245 | function buildSortedIndex(count: number): Uint32Array { |
| 246 | const indices = new Uint32Array(count); |
| 247 | for (let i = 0; i < count; i++) indices[i] = i; |
| 248 | return sortSelectionSlots(table, indices) as Uint32Array; |
| 249 | } |
| 250 | } |