File
Blob: src/worker/git/operations/fetch/plan.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | import type { SnapshotLoadResult } from "@/worker/git/pack/snapshot"; |
| 4 | import type { OrderedPackSnapshot, ServeUploadPackPlan, UploadPackPlan } from "./types"; |
| 5 | import type { PackRefSnapshotEntry, PackRefSnapshotLoadResult } from "@/worker/git/pack/refIndex"; |
| 6 | |
| 7 | import { createLogger } from "@/worker/common"; |
| 8 | import { buildInitialCloneNeeded, loadOrderedPackSnapshot } from "@/worker/git/pack/snapshot"; |
| 9 | import { getDoIdFromPath } from "@/worker/keys"; |
| 10 | import { findCommonHaves } from "../closure"; |
| 11 | import { computeNeededFromPackRefs } from "./refClosure"; |
| 12 | import { loadPackRefView } from "@/worker/git/pack/refIndex"; |
| 13 | |
| 14 | export 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 | |
| 26 | export 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 | |
| 41 | function 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 | |
| 85 | export 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 | |
| 148 | export 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 | |
| 222 | export 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 | } |