Skip to content
File

Blob: src/worker/git/operations/fetch/plan.ts

typescript261 lines
1import type { CacheContext } from "@/worker/cache";
2import type { Logger } from "@/worker/common/logger";
3import type { SnapshotLoadResult } from "@/worker/git/pack/snapshot";
4import type { OrderedPackSnapshot, ServeUploadPackPlan, UploadPackPlan } from "./types";
5import type { PackRefSnapshotEntry, PackRefSnapshotLoadResult } from "@/worker/git/pack/refIndex";
6 
7import { createLogger } from "@/worker/common";
8import { buildInitialCloneNeeded, loadOrderedPackSnapshot } from "@/worker/git/pack/snapshot";
9import { getDoIdFromPath } from "@/worker/keys";
10import { findCommonHaves } from "../closure";
11import { computeNeededFromPackRefs } from "./refClosure";
12import { loadPackRefView } from "@/worker/git/pack/refIndex";
13 
14export class FetchPlanRetryError extends Error {
15 readonly reason: "missing-ref-index" | "closure-budget-exceeded";
16 readonly retryAfterSeconds: number;
17 
18 constructor(reason: "missing-ref-index" | "closure-budget-exceeded") {
19 super(reason);
20 this.name = "FetchPlanRetryError";
21 this.reason = reason;
22 this.retryAfterSeconds = 10;
23 }
24}
25 
26export async function loadUploadPackSnapshot(
27 env: Env,
28 repoId: string,
29 cacheCtx?: CacheContext
30): Promise<SnapshotLoadResult> {
31 // Snapshot readiness stays outside the streaming response so callers can
32 // still convert "not ready" into an HTTP retry signal before headers commit.
33 const log = createLogger(env.LOG_LEVEL, { service: "StreamPlan", repoId });
34 const snapshotLoad = await loadOrderedPackSnapshot(env, repoId, cacheCtx, log);
35 if (snapshotLoad.type === "RepositoryNotReady") {
36 log.warn("stream:plan:repository-not-ready", { reason: snapshotLoad.reason });
37 }
38 return snapshotLoad;
39}
40 
41function schedulePackRefBackfill(args: {
42 env: Env;
43 repoId: string;
44 packKey: string;
45 cacheCtx?: CacheContext;
46 log: Logger;
47 reason: string;
48}): void {
49 const doId = getDoIdFromPath(args.packKey);
50 if (!doId) {
51 args.log.warn("stream:fetch:ref-index-backfill-skipped", {
52 packKey: args.packKey,
53 reason: "missing-do-id",
54 });
55 return;
56 }
57 
58 const send = args.env.REPO_TASKS_QUEUE.send({
59 kind: "pack-ref-backfill",
60 doId,
61 repoId: args.repoId,
62 packKey: args.packKey,
63 })
64 .then(() => {
65 args.log.info("stream:fetch:ref-index-backfill-queued", {
66 packKey: args.packKey,
67 reason: args.reason,
68 });
69 })
70 .catch((error) => {
71 args.log.warn("stream:fetch:ref-index-backfill-enqueue-failed", {
72 packKey: args.packKey,
73 reason: args.reason,
74 error: String(error),
75 });
76 });
77 
78 if (args.cacheCtx) {
79 args.cacheCtx.ctx.waitUntil(send);
80 } else {
81 send.catch(() => {});
82 }
83}
84 
85export async function loadPackRefSnapshot(
86 env: Env,
87 repoId: string,
88 snapshot: OrderedPackSnapshot,
89 cacheCtx?: CacheContext
90): Promise<PackRefSnapshotLoadResult> {
91 const log = createLogger(env.LOG_LEVEL, { service: "StreamPlan", repoId });
92 const packs: PackRefSnapshotEntry[] = [];
93 const missing: Array<{
94 packKey: string;
95 packBytes: number;
96 reason: "missing" | "corrupt" | "stale";
97 detail?: string;
98 }> = [];
99 
100 for (const pack of snapshot.packs) {
101 const load = await loadPackRefView(env, pack.packKey, pack.idx, cacheCtx);
102 if (load.type === "Ready") {
103 packs.push({
104 packKey: pack.packKey,
105 packBytes: pack.packBytes,
106 idx: pack.idx,
107 refs: load.view,
108 });
109 continue;
110 }
111 
112 const reason = load.type === "Missing" ? "missing" : load.kind;
113 const detail = load.type === "Invalid" ? load.reason : undefined;
114 missing.push({
115 packKey: pack.packKey,
116 packBytes: pack.packBytes,
117 reason,
118 detail,
119 });
120 log.warn("stream:fetch:ref-index-missing", {
121 packKey: pack.packKey,
122 reason,
123 detail,
124 });
125 schedulePackRefBackfill({
126 env,
127 repoId,
128 packKey: pack.packKey,
129 cacheCtx,
130 log,
131 reason,
132 });
133 }
134 
135 log.info("stream:plan:ref-snapshot", {
136 packs: snapshot.packs.length,
137 loaded: packs.length,
138 missing: missing.length,
139 });
140 
141 if (missing.length > 0) {
142 return { type: "Missing", packs: missing };
143 }
144 
145 return { type: "Ready", packs };
146}
147 
148export async function buildServeUploadPackPlan(
149 env: Env,
150 repoId: string,
151 snapshot: OrderedPackSnapshot,
152 wants: string[],
153 haves: string[],
154 signal?: AbortSignal,
155 cacheCtx?: CacheContext,
156 onProgress?: (message: string) => void
157): Promise<ServeUploadPackPlan> {
158 const log = createLogger(env.LOG_LEVEL, { service: "StreamPlan", repoId });
159 
160 if (haves.length === 0) {
161 onProgress?.("Selecting objects to send...\n");
162 const neededOids = buildInitialCloneNeeded(snapshot);
163 log.info("stream:plan:init-clone", {
164 packs: snapshot.packs.length,
165 needed: neededOids.length,
166 });
167 return {
168 type: "Serve",
169 repoId,
170 snapshot,
171 neededOids,
172 ackOids: [],
173 signal,
174 cacheCtx,
175 };
176 }
177 
178 const refSnapshot = await loadPackRefSnapshot(env, repoId, snapshot, cacheCtx);
179 if (refSnapshot.type === "Missing") {
180 throw new FetchPlanRetryError("missing-ref-index");
181 }
182 
183 const closure = await computeNeededFromPackRefs({
184 logLevel: env.LOG_LEVEL,
185 repoId,
186 packs: refSnapshot.packs,
187 wants,
188 haves,
189 onProgress,
190 });
191 if (closure.type === "BudgetExceeded") {
192 log.warn("stream:plan:closure-budget-exceeded", {
193 reason: closure.reason,
194 needed: closure.neededOids.length,
195 seen: closure.stats.seen,
196 queued: closure.stats.queued,
197 missing: closure.stats.missing,
198 edgeVisits: closure.stats.edgeVisits,
199 duplicateQueueSkips: closure.stats.duplicateQueueSkips,
200 });
201 throw new FetchPlanRetryError("closure-budget-exceeded");
202 }
203 const neededOids = closure.neededOids;
204 
205 log.info("stream:plan:serve", {
206 packs: snapshot.packs.length,
207 needed: neededOids.length,
208 ackOids: 0,
209 });
210 
211 return {
212 type: "Serve",
213 repoId,
214 snapshot,
215 neededOids,
216 ackOids: [],
217 signal,
218 cacheCtx,
219 };
220}
221 
222export async function planUploadPack(
223 env: Env,
224 repoId: string,
225 wants: string[],
226 haves: string[],
227 done: boolean,
228 signal?: AbortSignal,
229 cacheCtx?: CacheContext
230): Promise<UploadPackPlan> {
231 const snapshotLoad = await loadUploadPackSnapshot(env, repoId, cacheCtx);
232 if (snapshotLoad.type === "RepositoryNotReady") {
233 return { type: "RepositoryNotReady" };
234 }
235 
236 if (!done) {
237 const ackOids = haves.length > 0 ? await findCommonHaves(env, repoId, haves, cacheCtx) : [];
238 return {
239 type: "Serve",
240 repoId,
241 snapshot: snapshotLoad.snapshot,
242 neededOids: [],
243 ackOids,
244 signal,
245 cacheCtx,
246 };
247 }
248 
249 const servePlan = await buildServeUploadPackPlan(
250 env,
251 repoId,
252 snapshotLoad.snapshot,
253 wants,
254 haves,
255 signal,
256 cacheCtx
257 );
258 
259 return servePlan;
260}