Skip to content
File

Blob: src/worker/git/receive/pipeline.ts

typescript388 lines
1import type { CacheContext } from "@/worker/cache";
2import type { Logger } from "@/worker/common/logger";
3import type { RepoDurableObject } from "@/worker/do";
4import type { PackCatalogRow } from "@/worker/do/repo/db/schema";
5import type { ReceiveCommand, ReceiveStatus } from "@/worker/git/operations/validation";
6 
7import { SubrequestLimiter } from "@/worker/git/operations/limits";
8import {
9 resolveDeltasAndWriteIdx,
10 runPackConnectivityCheck,
11 scanPack,
12} from "@/worker/git/pack/indexer";
13import { doPrefix, r2PackKey } from "@/worker/keys";
14import { deleteStagedPack, stagePackToR2, type StagedPackUpload } from "./r2Upload";
15import { buildReceiveReportStatus, isReceiveAbort, throwIfReceiveAborted } from "./support";
16 
17type RepoStub = DurableObjectStub<RepoDurableObject>;
18 
19export type ReceivePipelineResult = {
20 reportStatusBody: Uint8Array;
21 changed: boolean;
22 empty: boolean;
23 packKey?: string;
24 packBytes?: number;
25};
26 
27export class ReceivePipelineHttpError extends Error {
28 readonly status: number;
29 readonly reason: string;
30 
31 constructor(status: number, reason: string, message: string) {
32 super(message);
33 this.name = "ReceivePipelineHttpError";
34 this.status = status;
35 this.reason = reason;
36 }
37}
38 
39type ReceiveCleanupAttempt = "inline" | "retry";
40 
41async function abortReceiveLease(args: {
42 stub: RepoStub;
43 leaseToken: string;
44 log: Logger;
45 reason: string;
46 attempt: ReceiveCleanupAttempt;
47}): Promise<boolean> {
48 try {
49 const cleared = await args.stub.abortReceive(args.leaseToken);
50 if (!cleared) {
51 args.log.warn("receive:abort-missed", {
52 reason: args.reason,
53 attempt: args.attempt,
54 leaseToken: args.leaseToken,
55 });
56 }
57 return cleared;
58 } catch (error) {
59 args.log.warn("receive:abort-failed", {
60 reason: args.reason,
61 attempt: args.attempt,
62 leaseToken: args.leaseToken,
63 error: String(error),
64 });
65 return false;
66 }
67}
68 
69async function cleanupStagedPack(args: {
70 stagedUpload: StagedPackUpload | undefined;
71 log: Logger;
72 reason: string;
73 attempt: ReceiveCleanupAttempt;
74}): Promise<boolean> {
75 if (!args.stagedUpload) return true;
76 
77 try {
78 await deleteStagedPack(args.stagedUpload);
79 return true;
80 } catch (error) {
81 args.log.warn("receive:staged-pack-cleanup-failed", {
82 reason: args.reason,
83 attempt: args.attempt,
84 packKey: args.stagedUpload.packKey,
85 error: String(error),
86 });
87 return false;
88 }
89}
90 
91async function cleanupFailedReceive(args: {
92 ctx: ExecutionContext;
93 stub: RepoStub;
94 leaseToken: string;
95 stagedUpload: StagedPackUpload | undefined;
96 log: Logger;
97 reason: string;
98}): Promise<void> {
99 const leaseCleared = await abortReceiveLease({
100 stub: args.stub,
101 leaseToken: args.leaseToken,
102 log: args.log,
103 reason: args.reason,
104 attempt: "inline",
105 });
106 const stagedPackDeleted = await cleanupStagedPack({
107 stagedUpload: args.stagedUpload,
108 log: args.log,
109 reason: args.reason,
110 attempt: "inline",
111 });
112 
113 if (leaseCleared && stagedPackDeleted) return;
114 
115 args.log.warn("receive:cleanup-retry-scheduled", {
116 reason: args.reason,
117 leaseToken: args.leaseToken,
118 packKey: args.stagedUpload?.packKey,
119 });
120 args.ctx.waitUntil(
121 (async () => {
122 const retryLeaseCleared =
123 leaseCleared ||
124 (await abortReceiveLease({
125 stub: args.stub,
126 leaseToken: args.leaseToken,
127 log: args.log,
128 reason: args.reason,
129 attempt: "retry",
130 }));
131 const retryStagedPackDeleted =
132 stagedPackDeleted ||
133 (await cleanupStagedPack({
134 stagedUpload: args.stagedUpload,
135 log: args.log,
136 reason: args.reason,
137 attempt: "retry",
138 }));
139 
140 if (!retryLeaseCleared || !retryStagedPackDeleted) {
141 args.log.error("receive:cleanup-retry-incomplete", {
142 reason: args.reason,
143 leaseCleared: retryLeaseCleared,
144 stagedPackDeleted: retryStagedPackDeleted,
145 leaseToken: args.leaseToken,
146 packKey: args.stagedUpload?.packKey,
147 });
148 }
149 })()
150 );
151}
152 
153type ExecuteReceivePipelineArgs = {
154 env: Env;
155 repoId: string;
156 request: Request;
157 ctx: ExecutionContext;
158 packStream: ReadableStream<Uint8Array>;
159 bytesConsumed: number;
160 stub: RepoStub;
161 leaseToken: string;
162 activeCatalog: PackCatalogRow[];
163 commands: ReceiveCommand[];
164 log: Logger;
165 cacheCtx: CacheContext;
166 limiter: SubrequestLimiter;
167 countSubrequest(op: string, n?: number): void;
168 onProgress?: (message: string) => void;
169};
170 
171function buildReceiveResult(args: {
172 unpackOk: boolean;
173 unpackMessage?: string;
174 commands: ReceiveCommand[];
175 statuses: ReceiveStatus[];
176 changed: boolean;
177 empty: boolean;
178 packKey?: string;
179 packBytes?: number;
180}): ReceivePipelineResult {
181 return {
182 reportStatusBody: buildReceiveReportStatus({
183 unpackOk: args.unpackOk,
184 unpackMessage: args.unpackMessage,
185 commands: args.commands,
186 statuses: args.statuses,
187 }),
188 changed: args.changed,
189 empty: args.empty,
190 packKey: args.packKey,
191 packBytes: args.packBytes,
192 };
193}
194 
195export async function executeReceivePipeline(
196 args: ExecuteReceivePipelineArgs
197): Promise<ReceivePipelineResult> {
198 let stagedUpload: StagedPackUpload | undefined;
199 
200 try {
201 const hasNonDelete = args.commands.some((command) => !/^0{40}$/i.test(command.newOid));
202 let stagedPack:
203 | {
204 packKey: string;
205 packBytes: number;
206 idxBytes: number;
207 objectCount: number;
208 }
209 | undefined;
210 
211 if (hasNonDelete) {
212 const packKey = r2PackKey(
213 doPrefix(args.stub.id.toString()),
214 `pack-rx-${args.leaseToken}.pack`
215 );
216 stagedUpload = await stagePackToR2({
217 env: args.env,
218 request: args.request,
219 packStream: args.packStream,
220 packKey,
221 bytesConsumed: args.bytesConsumed,
222 limiter: args.limiter,
223 countSubrequest: args.countSubrequest,
224 onProgress: args.onProgress,
225 });
226 throwIfReceiveAborted(args.request, args.log, "stage-pack");
227 
228 const scanResult = await scanPack({
229 env: args.env,
230 packKey: stagedUpload.packKey,
231 packSize: stagedUpload.packBytes,
232 limiter: args.limiter,
233 countSubrequest: (n = 1) => args.countSubrequest("r2:scan-pack", n),
234 log: args.log,
235 signal: args.request.signal,
236 onProgress: args.onProgress,
237 });
238 throwIfReceiveAborted(args.request, args.log, "scan-pack");
239 
240 const resolveResult = await resolveDeltasAndWriteIdx({
241 env: args.env,
242 packKey: stagedUpload.packKey,
243 packSize: stagedUpload.packBytes,
244 limiter: args.limiter,
245 countSubrequest: (n = 1) => args.countSubrequest("r2:resolve-pack", n),
246 log: args.log,
247 scanResult,
248 activeCatalog: args.activeCatalog,
249 cacheCtx: args.cacheCtx,
250 repoId: args.repoId,
251 signal: args.request.signal,
252 onProgress: args.onProgress,
253 });
254 throwIfReceiveAborted(args.request, args.log, "resolve-pack");
255 
256 const connectivityStatuses = args.commands.map((command) => ({
257 ref: command.ref,
258 ok: true,
259 }));
260 args.onProgress?.("Checking received object connectivity\n");
261 await runPackConnectivityCheck({
262 env: args.env,
263 repoId: args.repoId,
264 newPackKey: stagedUpload.packKey,
265 newIdxView: resolveResult.idxView,
266 newPackSize: stagedUpload.packBytes,
267 activeCatalog: args.activeCatalog,
268 commands: args.commands,
269 statuses: connectivityStatuses,
270 log: args.log,
271 cacheCtx: args.cacheCtx,
272 });
273 throwIfReceiveAborted(args.request, args.log, "connectivity-check");
274 
275 if (!connectivityStatuses.every((status) => status.ok)) {
276 args.countSubrequest("do:abort-receive");
277 await cleanupFailedReceive({
278 ctx: args.ctx,
279 stub: args.stub,
280 leaseToken: args.leaseToken,
281 stagedUpload,
282 log: args.log,
283 reason: "connectivity-rejected",
284 });
285 args.log.warn("receive:connectivity-rejected", {
286 conflictCount: connectivityStatuses.filter((status) => !status.ok).length,
287 });
288 return buildReceiveResult({
289 unpackOk: true,
290 commands: args.commands,
291 statuses: connectivityStatuses,
292 changed: false,
293 empty: false,
294 });
295 }
296 
297 stagedPack = {
298 packKey: stagedUpload.packKey,
299 packBytes: stagedUpload.packBytes,
300 idxBytes: resolveResult.idxBytes,
301 objectCount: resolveResult.objectCount,
302 };
303 }
304 
305 args.countSubrequest("do:finalize-receive");
306 throwIfReceiveAborted(args.request, args.log, "finalize-receive");
307 args.onProgress?.("Updating refs\n");
308 const finalize = await args.stub.finalizeReceive({
309 token: args.leaseToken,
310 commands: args.commands,
311 stagedPack,
312 });
313 
314 if (finalize.status === "lease_mismatch") {
315 await cleanupStagedPack({
316 stagedUpload,
317 log: args.log,
318 reason: "finalize-lease-mismatch",
319 attempt: "inline",
320 });
321 args.log.warn("receive:lease-mismatch", { leaseToken: args.leaseToken });
322 throw new ReceivePipelineHttpError(
323 503,
324 "lease-mismatch",
325 "Repository receive lease expired before commit."
326 );
327 }
328 
329 if (finalize.status === "ref_conflict") {
330 await cleanupStagedPack({
331 stagedUpload,
332 log: args.log,
333 reason: "finalize-ref-conflict",
334 attempt: "inline",
335 });
336 args.log.warn("receive:ref-conflict", {
337 conflictCount: finalize.statuses.filter((status) => !status.ok).length,
338 stage: "finalize",
339 });
340 return buildReceiveResult({
341 unpackOk: true,
342 commands: args.commands,
343 statuses: finalize.statuses,
344 changed: false,
345 empty: false,
346 });
347 }
348 
349 if (finalize.shouldQueueCompaction) {
350 args.log.info("receive:compaction-requested", { repoId: args.repoId });
351 args.ctx.waitUntil(
352 args.env.REPO_TASKS_QUEUE.send({
353 kind: "compaction",
354 doId: args.stub.id.toString(),
355 repoId: args.repoId,
356 }).catch((error) => {
357 args.log.warn("receive:compaction-enqueue-failed", {
358 repoId: args.repoId,
359 error: String(error),
360 });
361 })
362 );
363 }
364 
365 return buildReceiveResult({
366 unpackOk: true,
367 commands: args.commands,
368 statuses: finalize.statuses,
369 changed: finalize.changed,
370 empty: finalize.empty,
371 packKey: stagedPack?.packKey,
372 packBytes: stagedPack?.packBytes,
373 });
374 } catch (error) {
375 args.countSubrequest("do:abort-receive");
376 const aborted = isReceiveAbort(args.request, error);
377 await cleanupFailedReceive({
378 ctx: args.ctx,
379 stub: args.stub,
380 leaseToken: args.leaseToken,
381 stagedUpload,
382 log: args.log,
383 reason: aborted ? "receive-aborted" : "receive-error",
384 });
385 throw error;
386 }
387}