Skip to content
File

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

typescript137 lines
1import type { OrderedPackSnapshot } from "@/worker/git/operations/fetch/types";
2 
3import { createLogger } from "@/worker/common";
4import {
5 buildSelection,
6 buildOutputOrder,
7 canPassthroughSinglePack,
8 computeHeaderLengths,
9} from "./rewrite/plan";
10import {
11 ensurePackReadState,
12 type RewriteFailure,
13 type RewriteFailureRecorder,
14 type RewriteOptions,
15} from "./rewrite/shared";
16import { createPassthroughStream, createRewriteStream } from "./rewrite/stream";
17 
18export type PackRewriteResult =
19 | { status: "ok"; stream: ReadableStream<Uint8Array> }
20 | { status: "failed"; failure: RewriteFailure };
21 
22export async function rewritePackResult(
23 env: Env,
24 snapshot: OrderedPackSnapshot,
25 neededOids: string[],
26 options?: RewriteOptions
27): Promise<PackRewriteResult> {
28 const log = createLogger(env.LOG_LEVEL, { service: "PackRewrite" });
29 const startedAt = Date.now();
30 const warnedFlags = new Set<string>();
31 const failure: RewriteFailureRecorder = options?.failure || {};
32 const rewriteOptions: RewriteOptions = { ...options, failure };
33 
34 function failed(reason: string, retryable: boolean, details?: Record<string, unknown>) {
35 return {
36 status: "failed" as const,
37 failure: failure.value || { reason, retryable, details },
38 };
39 }
40 
41 if (rewriteOptions.signal?.aborted) {
42 return failed("aborted", true);
43 }
44 if (!rewriteOptions.limiter) {
45 throw new Error("rewrite: limiter required");
46 }
47 if (!rewriteOptions.countSubrequest) {
48 throw new Error("rewrite: countSubrequest required");
49 }
50 
51 const selection = await buildSelection(
52 env,
53 snapshot,
54 neededOids,
55 log,
56 warnedFlags,
57 rewriteOptions
58 );
59 if (!selection) {
60 return failed("selection-failed", true, { needed: neededOids.length });
61 }
62 
63 const { table, readerStates } = selection;
64 
65 if (canPassthroughSinglePack(snapshot, table)) {
66 const readState = await ensurePackReadState(
67 env,
68 snapshot.packs[0]!,
69 0,
70 readerStates,
71 log,
72 warnedFlags,
73 rewriteOptions
74 );
75 
76 log.info("rewrite:passthrough", {
77 packKey: snapshot.packs[0]?.packKey,
78 objects: table.count,
79 });
80 
81 return {
82 status: "ok",
83 stream: createPassthroughStream({
84 env,
85 snapshotPack: snapshot.packs[0]!,
86 readState,
87 log,
88 warnedFlags,
89 options: rewriteOptions,
90 onComplete: () => {
91 log.info("rewrite:stream-complete", {
92 passthrough: true,
93 wholePackLoads: countWholePackLoads(readerStates),
94 timeMs: Date.now() - startedAt,
95 });
96 },
97 }),
98 };
99 }
100 
101 if (!buildOutputOrder(table, log)) {
102 return failed("topology-incomplete", false, { selected: table.count });
103 }
104 if (!computeHeaderLengths(table, log)) {
105 return failed("header-lengths-did-not-converge", false, { selected: table.count });
106 }
107 
108 return {
109 status: "ok",
110 stream: createRewriteStream(table, snapshot, readerStates, log, rewriteOptions, () => {
111 log.info("rewrite:stream-complete", {
112 passthrough: false,
113 wholePackLoads: countWholePackLoads(readerStates),
114 timeMs: Date.now() - startedAt,
115 });
116 }),
117 };
118}
119 
120export async function rewritePack(
121 env: Env,
122 snapshot: OrderedPackSnapshot,
123 neededOids: string[],
124 options?: RewriteOptions
125): Promise<ReadableStream<Uint8Array> | undefined> {
126 const result = await rewritePackResult(env, snapshot, neededOids, options);
127 return result.status === "ok" ? result.stream : undefined;
128}
129 
130function countWholePackLoads(readerStates: Map<number, { wholePack?: Uint8Array }>): number {
131 let count = 0;
132 for (const state of readerStates.values()) {
133 if (state.wholePack) count++;
134 }
135 return count;
136}