File
Blob: src/worker/git/pack/rewrite/stream.ts
| 1 | import type { |
| 2 | OrderedPackSnapshot, |
| 3 | OrderedPackSnapshotEntry, |
| 4 | } from "@/worker/git/operations/fetch/types"; |
| 5 | import type { Logger } from "@/worker/common/logger"; |
| 6 | |
| 7 | import { createDigestStream } from "@/worker/common"; |
| 8 | import { isResolveAbortedError } from "@/worker/git/pack/indexer/resolve/errors"; |
| 9 | import { encodeOfsDeltaDistance } from "../packMeta"; |
| 10 | import { |
| 11 | WHOLE_PACK_MAX_BYTES, |
| 12 | buildPackHeader, |
| 13 | countRewriteSubrequest, |
| 14 | type PackReadState, |
| 15 | type RewriteOptions, |
| 16 | type SelectionTable, |
| 17 | } from "./shared"; |
| 18 | |
| 19 | // --------------------------------------------------------------------------- |
| 20 | // Entry header construction |
| 21 | // --------------------------------------------------------------------------- |
| 22 | |
| 23 | function buildEntryHeaderBytes(table: SelectionTable, sel: number): Uint8Array | undefined { |
| 24 | const type = table.typeCodes[sel]; |
| 25 | const svStart = sel * 5; |
| 26 | const svLen = table.sizeVarLens[sel]; |
| 27 | if (svLen === 0) return undefined; |
| 28 | |
| 29 | if (type === 6) { |
| 30 | const base = table.baseSlots[sel]; |
| 31 | if (base < 0) return undefined; |
| 32 | const distance = table.outputOffsets[sel] - table.outputOffsets[base]; |
| 33 | const distBytes = encodeOfsDeltaDistance(distance); |
| 34 | const out = new Uint8Array(svLen + distBytes.length); |
| 35 | out.set(table.sizeVarBuf.subarray(svStart, svStart + svLen), 0); |
| 36 | out.set(distBytes, svLen); |
| 37 | return out; |
| 38 | } |
| 39 | |
| 40 | if (type === 7) { |
| 41 | if (!table.baseOidRaw) return undefined; |
| 42 | const baseOidBytes = table.baseOidRaw.subarray(sel * 20, sel * 20 + 20); |
| 43 | const out = new Uint8Array(svLen + 20); |
| 44 | out.set(table.sizeVarBuf.subarray(svStart, svStart + svLen), 0); |
| 45 | out.set(baseOidBytes, svLen); |
| 46 | return out; |
| 47 | } |
| 48 | |
| 49 | // Non-delta: just the size varint (subarray is safe — table is immutable during streaming) |
| 50 | return table.sizeVarBuf.subarray(svStart, svStart + svLen); |
| 51 | } |
| 52 | |
| 53 | // --------------------------------------------------------------------------- |
| 54 | // Payload emission |
| 55 | // --------------------------------------------------------------------------- |
| 56 | |
| 57 | async function emitPackPayload( |
| 58 | controller: ReadableStreamDefaultController<Uint8Array>, |
| 59 | writer: WritableStreamDefaultWriter<Uint8Array>, |
| 60 | table: SelectionTable, |
| 61 | sel: number, |
| 62 | state: PackReadState | undefined |
| 63 | ): Promise<void> { |
| 64 | const syntheticPayload = table.syntheticPayloads[sel]; |
| 65 | if (syntheticPayload) { |
| 66 | await writer.write(syntheticPayload); |
| 67 | controller.enqueue(syntheticPayload); |
| 68 | return; |
| 69 | } |
| 70 | |
| 71 | if (!state) { |
| 72 | throw new Error( |
| 73 | `rewrite: missing read state for pack#${table.packSlots[sel]} entry#${table.entryIndices[sel]}` |
| 74 | ); |
| 75 | } |
| 76 | |
| 77 | const payloadStart = table.offsets[sel] + table.headerLens[sel]; |
| 78 | let bytesLeft = table.payloadLens[sel]; |
| 79 | if (bytesLeft <= 0) return; |
| 80 | |
| 81 | if (state.wholePack) { |
| 82 | const payload = state.wholePack.subarray(payloadStart, payloadStart + bytesLeft); |
| 83 | await writer.write(payload); |
| 84 | controller.enqueue(payload); |
| 85 | return; |
| 86 | } |
| 87 | |
| 88 | let currentOffset = payloadStart; |
| 89 | while (bytesLeft > 0) { |
| 90 | let window = await state.reader.readWindow(currentOffset, bytesLeft); |
| 91 | if (window.length === 0) { |
| 92 | window = await state.reader.readRange(currentOffset, Math.min(bytesLeft, 1)); |
| 93 | } |
| 94 | if (window.length === 0) { |
| 95 | throw new Error( |
| 96 | `rewrite: unexpected EOF while streaming pack#${table.packSlots[sel]} entry#${table.entryIndices[sel]}` |
| 97 | ); |
| 98 | } |
| 99 | |
| 100 | await writer.write(window); |
| 101 | controller.enqueue(window); |
| 102 | currentOffset += window.length; |
| 103 | bytesLeft -= window.length; |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | // --------------------------------------------------------------------------- |
| 108 | // Passthrough stream (single pack, all objects selected) |
| 109 | // --------------------------------------------------------------------------- |
| 110 | |
| 111 | export async function passthroughSinglePack( |
| 112 | env: Env, |
| 113 | snapshotPack: OrderedPackSnapshotEntry, |
| 114 | readState: PackReadState, |
| 115 | controller: ReadableStreamDefaultController<Uint8Array>, |
| 116 | log: Logger, |
| 117 | warnedFlags: Set<string>, |
| 118 | options?: RewriteOptions |
| 119 | ): Promise<"completed" | "aborted"> { |
| 120 | const digestStream = createDigestStream("SHA-1"); |
| 121 | const writer = digestStream.getWriter(); |
| 122 | |
| 123 | if (options?.signal?.aborted) { |
| 124 | await writer.abort(); |
| 125 | return "aborted"; |
| 126 | } |
| 127 | |
| 128 | const emit = async (chunk: Uint8Array) => { |
| 129 | await writer.write(chunk); |
| 130 | controller.enqueue(chunk); |
| 131 | }; |
| 132 | |
| 133 | options?.onProgress?.(`Enumerating objects: ${snapshotPack.idx.count}, from 1 packs\n`); |
| 134 | |
| 135 | if (readState.wholePack) { |
| 136 | if (readState.wholePack.length < 20) { |
| 137 | throw new Error("rewrite: passthrough pack read failed"); |
| 138 | } |
| 139 | if (options?.signal?.aborted) { |
| 140 | await writer.abort(); |
| 141 | return "aborted"; |
| 142 | } |
| 143 | await emit(readState.wholePack.subarray(0, readState.wholePack.length - 20)); |
| 144 | } else if (snapshotPack.packBytes <= WHOLE_PACK_MAX_BYTES) { |
| 145 | throw new Error("rewrite: missing whole-pack preload for passthrough"); |
| 146 | } else { |
| 147 | countRewriteSubrequest( |
| 148 | log, |
| 149 | warnedFlags, |
| 150 | options, |
| 151 | `rewrite-passthrough:${snapshotPack.packKey}`, |
| 152 | { op: "r2:get-pack", packKey: snapshotPack.packKey } |
| 153 | ); |
| 154 | await options!.limiter!.run("r2:get-pack", async () => { |
| 155 | const packObject = await env.REPO_BUCKET.get(snapshotPack.packKey); |
| 156 | if (!packObject?.body) { |
| 157 | throw new Error("rewrite: passthrough pack stream unavailable"); |
| 158 | } |
| 159 | |
| 160 | const reader = packObject.body.getReader(); |
| 161 | let trailing = new Uint8Array(0); |
| 162 | while (true) { |
| 163 | if (options?.signal?.aborted) { |
| 164 | await reader.cancel(); |
| 165 | await writer.abort(); |
| 166 | throw new Error("rewrite: passthrough aborted"); |
| 167 | } |
| 168 | const { done, value } = await reader.read(); |
| 169 | if (done) break; |
| 170 | if (!value) continue; |
| 171 | |
| 172 | const chunk = new Uint8Array(trailing.length + value.length); |
| 173 | chunk.set(trailing, 0); |
| 174 | chunk.set(value, trailing.length); |
| 175 | if (chunk.length <= 20) { |
| 176 | trailing = chunk; |
| 177 | continue; |
| 178 | } |
| 179 | |
| 180 | const bodyChunk = chunk.subarray(0, chunk.length - 20); |
| 181 | trailing = chunk.slice(chunk.length - 20); |
| 182 | if (options?.signal?.aborted) { |
| 183 | await reader.cancel(); |
| 184 | await writer.abort(); |
| 185 | throw new Error("rewrite: passthrough aborted"); |
| 186 | } |
| 187 | await emit(bodyChunk); |
| 188 | } |
| 189 | |
| 190 | if (trailing.length < 20) { |
| 191 | throw new Error("rewrite: truncated passthrough pack"); |
| 192 | } |
| 193 | }); |
| 194 | } |
| 195 | |
| 196 | await writer.close(); |
| 197 | options?.onProgress?.( |
| 198 | `Counting objects: 100% (${snapshotPack.idx.count}/${snapshotPack.idx.count}), done.\n` |
| 199 | ); |
| 200 | controller.enqueue(new Uint8Array(await digestStream.digest)); |
| 201 | return "completed"; |
| 202 | } |
| 203 | |
| 204 | export function createPassthroughStream(args: { |
| 205 | env: Env; |
| 206 | snapshotPack: OrderedPackSnapshotEntry; |
| 207 | readState: PackReadState; |
| 208 | log: Logger; |
| 209 | warnedFlags: Set<string>; |
| 210 | options?: RewriteOptions; |
| 211 | onComplete?: () => void; |
| 212 | }): ReadableStream<Uint8Array> { |
| 213 | return new ReadableStream<Uint8Array>({ |
| 214 | async start(controller) { |
| 215 | try { |
| 216 | const status = await passthroughSinglePack( |
| 217 | args.env, |
| 218 | args.snapshotPack, |
| 219 | args.readState, |
| 220 | controller, |
| 221 | args.log, |
| 222 | args.warnedFlags, |
| 223 | args.options |
| 224 | ); |
| 225 | if (status === "aborted") { |
| 226 | args.log.debug("rewrite:passthrough-aborted"); |
| 227 | controller.close(); |
| 228 | return; |
| 229 | } |
| 230 | args.onComplete?.(); |
| 231 | controller.close(); |
| 232 | } catch (error) { |
| 233 | if ( |
| 234 | isResolveAbortedError(error) || |
| 235 | args.options?.signal?.aborted || |
| 236 | (error instanceof Error && error.message === "rewrite: passthrough aborted") |
| 237 | ) { |
| 238 | args.log.debug("rewrite:passthrough-aborted"); |
| 239 | controller.close(); |
| 240 | return; |
| 241 | } |
| 242 | args.log.error("rewrite:passthrough-error", { error: String(error) }); |
| 243 | controller.error(error); |
| 244 | } |
| 245 | }, |
| 246 | }); |
| 247 | } |
| 248 | |
| 249 | // --------------------------------------------------------------------------- |
| 250 | // Rewrite stream (multi-pack or partial selection) |
| 251 | // --------------------------------------------------------------------------- |
| 252 | |
| 253 | export function createRewriteStream( |
| 254 | table: SelectionTable, |
| 255 | snapshot: OrderedPackSnapshot, |
| 256 | readStates: Map<number, PackReadState>, |
| 257 | log: Logger, |
| 258 | options?: RewriteOptions, |
| 259 | onComplete?: () => void |
| 260 | ): ReadableStream<Uint8Array> { |
| 261 | return new ReadableStream<Uint8Array>({ |
| 262 | async start(controller) { |
| 263 | let writer: WritableStreamDefaultWriter<Uint8Array> | undefined; |
| 264 | try { |
| 265 | const digestStream = createDigestStream("SHA-1"); |
| 266 | const digestWriter = digestStream.getWriter(); |
| 267 | writer = digestWriter; |
| 268 | |
| 269 | const emit = async (chunk: Uint8Array) => { |
| 270 | await digestWriter.write(chunk); |
| 271 | controller.enqueue(chunk); |
| 272 | }; |
| 273 | |
| 274 | await emit(buildPackHeader(table.count)); |
| 275 | options?.onProgress?.( |
| 276 | `Enumerating objects: ${table.count}, from ${readStates.size} packs\n` |
| 277 | ); |
| 278 | |
| 279 | const progressInterval = Math.max(1, Math.floor(table.count / 10)); |
| 280 | let streamed = 0; |
| 281 | |
| 282 | for (let i = 0; i < table.count; i++) { |
| 283 | if (options?.signal?.aborted) { |
| 284 | log.debug("rewrite:stream-aborted"); |
| 285 | await digestWriter.abort(); |
| 286 | controller.close(); |
| 287 | return; |
| 288 | } |
| 289 | |
| 290 | const sel = table.outputOrder[i]; |
| 291 | const packSlot = table.packSlots[sel]; |
| 292 | const readState = readStates.get(packSlot); |
| 293 | const syntheticPayload = table.syntheticPayloads[sel]; |
| 294 | const headerBytes = buildEntryHeaderBytes(table, sel); |
| 295 | if (!readState && !syntheticPayload) { |
| 296 | const pack = snapshot.packs[packSlot]; |
| 297 | log.error("rewrite:missing-read-state", { |
| 298 | sel, |
| 299 | packSlot, |
| 300 | entryIndex: table.entryIndices[sel], |
| 301 | packKey: pack?.packKey, |
| 302 | typeCode: table.typeCodes[sel], |
| 303 | baseSel: table.baseSlots[sel], |
| 304 | }); |
| 305 | throw new Error( |
| 306 | `rewrite: missing read state for ${pack?.packKey}#${table.entryIndices[sel]}` |
| 307 | ); |
| 308 | } |
| 309 | |
| 310 | if (!headerBytes) { |
| 311 | const pack = snapshot.packs[packSlot]; |
| 312 | const svLen = table.sizeVarLens[sel]; |
| 313 | const typeCode = table.typeCodes[sel]; |
| 314 | const baseSel = table.baseSlots[sel]; |
| 315 | const basePackSlot = baseSel >= 0 ? table.packSlots[baseSel] : undefined; |
| 316 | const baseEntryIndex = baseSel >= 0 ? table.entryIndices[baseSel] : undefined; |
| 317 | log.error("rewrite:invalid-header-state", { |
| 318 | sel, |
| 319 | packSlot, |
| 320 | entryIndex: table.entryIndices[sel], |
| 321 | packKey: pack?.packKey, |
| 322 | typeCode, |
| 323 | sizeVarLen: svLen, |
| 324 | hasBaseOidRaw: !!table.baseOidRaw, |
| 325 | baseSel, |
| 326 | basePackSlot, |
| 327 | baseEntryIndex, |
| 328 | }); |
| 329 | throw new Error( |
| 330 | `rewrite: invalid header state for ${pack?.packKey}#${table.entryIndices[sel]}` |
| 331 | ); |
| 332 | } |
| 333 | |
| 334 | await emit(headerBytes); |
| 335 | await emitPackPayload(controller, writer, table, sel, readState); |
| 336 | |
| 337 | streamed++; |
| 338 | if (streamed % progressInterval === 0 || streamed === table.count) { |
| 339 | const percent = Math.round((streamed / table.count) * 100); |
| 340 | if (streamed === table.count) { |
| 341 | options?.onProgress?.( |
| 342 | `Counting objects: 100% (${table.count}/${table.count}), done.\n` |
| 343 | ); |
| 344 | } else { |
| 345 | options?.onProgress?.(`Counting objects: ${percent}% (${streamed}/${table.count})\r`); |
| 346 | } |
| 347 | } |
| 348 | } |
| 349 | |
| 350 | await digestWriter.close(); |
| 351 | controller.enqueue(new Uint8Array(await digestStream.digest)); |
| 352 | onComplete?.(); |
| 353 | controller.close(); |
| 354 | } catch (error) { |
| 355 | if (isResolveAbortedError(error) || options?.signal?.aborted) { |
| 356 | log.debug("rewrite:stream-aborted"); |
| 357 | try { |
| 358 | await writer?.abort(); |
| 359 | } catch {} |
| 360 | controller.close(); |
| 361 | return; |
| 362 | } |
| 363 | log.error("rewrite:stream-error", { error: String(error) }); |
| 364 | controller.error(error); |
| 365 | } |
| 366 | }, |
| 367 | }); |
| 368 | } |