Skip to content
File

Blob: src/worker/git/operations/uploadStream/index.ts

typescript157 lines
1import type { CacheContext } from "@/worker/cache";
2import type { ServeUploadPackPlan } from "../fetch/types";
3 
4import { pktLine } from "@/worker/git/core";
5import { responseCacheControl } from "@/worker/cache/policy";
6import { createLogger } from "@/worker/common";
7import { getLimiter, countSubrequest } from "../limits";
8import { parseFetchArgs } from "../args";
9import { findCommonHaves } from "../closure";
10import { buildAckOnlyResponse } from "../fetch/protocol";
11import { repositoryNotReadyResponse } from "../fetch/responses";
12import {
13 buildServeUploadPackPlan,
14 FetchPlanRetryError,
15 loadUploadPackSnapshot,
16} from "../fetch/plan";
17import { resolvePackStreamResult } from "../fetch/execute";
18import {
19 SidebandProgressMux,
20 emitProgress,
21 emitFatal,
22 pipePackWithSideband,
23} from "../fetch/sideband";
24 
25export * from "../fetch/types";
26 
27function fetchPlanRetryResponse(error: FetchPlanRetryError): Response {
28 return new Response("Repository fetch planning is not ready, please retry in a few moments.\n", {
29 status: 503,
30 headers: {
31 "Retry-After": String(error.retryAfterSeconds),
32 "Content-Type": "text/plain; charset=utf-8",
33 "X-Git-Error": error.reason,
34 },
35 });
36}
37 
38export async function handleFetchV2Streaming(
39 env: Env,
40 repoId: string,
41 body: Uint8Array,
42 signal?: AbortSignal,
43 cacheCtx?: CacheContext
44): Promise<Response> {
45 const { wants, haves, done } = parseFetchArgs(body);
46 const log = createLogger(env.LOG_LEVEL, { service: "StreamFetchV2", repoId });
47 
48 if (signal?.aborted) {
49 return new Response("client aborted\n", { status: 499 });
50 }
51 
52 if (wants.length === 0) {
53 return buildAckOnlyResponse([], cacheCtx);
54 }
55 
56 if (!done) {
57 let ackOids: string[] = [];
58 if (haves.length > 0) {
59 ackOids = await findCommonHaves(env, repoId, haves, cacheCtx);
60 log.debug("stream:fetch:negotiation", { haves: haves.length, acks: ackOids.length });
61 }
62 return buildAckOnlyResponse(ackOids, cacheCtx);
63 }
64 
65 // Keep fetch readiness ahead of the response so clients still receive the
66 // current 503 + Retry-After signal while the repository is not yet fetchable.
67 const snapshotStart = Date.now();
68 const snapshotLoad = await loadUploadPackSnapshot(env, repoId, cacheCtx);
69 if (snapshotLoad.type === "RepositoryNotReady") {
70 log.warn("stream:fetch:repository-not-ready", { reason: snapshotLoad.reason });
71 return repositoryNotReadyResponse();
72 }
73 const snapshot = snapshotLoad.snapshot;
74 
75 log.info("stream:fetch:snapshot-ready", {
76 wants: wants.length,
77 haves: haves.length,
78 packs: snapshot.packs.length,
79 timeMs: Date.now() - snapshotStart,
80 });
81 
82 const planStart = Date.now();
83 log.info("stream:fetch:planning-start", {
84 wants: wants.length,
85 haves: haves.length,
86 });
87 let plan: ServeUploadPackPlan;
88 try {
89 plan = await buildServeUploadPackPlan(env, repoId, snapshot, wants, haves, signal, cacheCtx);
90 } catch (error) {
91 if (error instanceof FetchPlanRetryError) {
92 log.warn("stream:fetch:planning-retry", { reason: error.reason });
93 return fetchPlanRetryResponse(error);
94 }
95 throw error;
96 }
97 log.info("stream:fetch:planning-complete", {
98 needed: plan.neededOids.length,
99 timeMs: Date.now() - planStart,
100 });
101 
102 const responseStream = new ReadableStream<Uint8Array>({
103 async start(controller) {
104 const streamLog = createLogger(env.LOG_LEVEL, { service: "StreamFetchV2", repoId });
105 try {
106 controller.enqueue(pktLine("packfile\n"));
107 // Once the response body has started, later failures must travel over
108 // Git sideband because the HTTP status line is already committed.
109 emitProgress(controller, "Preparing pack...\n");
110 
111 const progressMux = new SidebandProgressMux();
112 const limiter = getLimiter(plan.cacheCtx);
113 const packResult = await resolvePackStreamResult(env, plan, {
114 signal: plan.signal,
115 limiter,
116 countSubrequest: (n?: number) => countSubrequest(plan.cacheCtx, n),
117 onProgress: (msg) => progressMux.push(msg),
118 });
119 
120 if (packResult.status !== "ok") {
121 streamLog.warn("stream:fetch:assemble-unavailable", {
122 needed: plan.neededOids.length,
123 reason: packResult.failure.reason,
124 retryable: packResult.failure.retryable,
125 details: packResult.failure.details,
126 });
127 emitFatal(controller, `Unable to assemble pack: ${packResult.failure.reason}`);
128 controller.close();
129 return;
130 }
131 
132 await pipePackWithSideband(packResult.stream, controller, {
133 signal: plan.signal,
134 progressMux,
135 log: streamLog,
136 });
137 
138 controller.close();
139 } catch (error) {
140 streamLog.error("stream:response:error", { error: String(error) });
141 try {
142 emitFatal(controller, String(error));
143 } catch {}
144 controller.error(error);
145 }
146 },
147 });
148 
149 return new Response(responseStream, {
150 status: 200,
151 headers: {
152 "Content-Type": "application/x-git-upload-pack-result",
153 "Cache-Control": responseCacheControl(cacheCtx),
154 },
155 });
156}