Skip to content
File

Blob: src/worker/git/pack/rewrite/stream.ts

typescript369 lines
1import type {
2 OrderedPackSnapshot,
3 OrderedPackSnapshotEntry,
4} from "@/worker/git/operations/fetch/types";
5import type { Logger } from "@/worker/common/logger";
6 
7import { createDigestStream } from "@/worker/common";
8import { isResolveAbortedError } from "@/worker/git/pack/indexer/resolve/errors";
9import { encodeOfsDeltaDistance } from "../packMeta";
10import {
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 
23function 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 
57async 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 
111export 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 
204export 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 
253export 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}