Skip to content
File

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

typescript400 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 { ReceiveStatus } from "@/worker/git/operations/validation";
6 
7import { clientAbortedResponse, createLogger, getRepoStub } from "@/worker/common";
8import {
9 MAX_SIMULTANEOUS_CONNECTIONS,
10 SubrequestLimiter,
11 countSubrequest,
12} from "@/worker/git/operations/limits";
13import { isValidRefName, validateReceiveCommands } from "@/worker/git/operations/validation";
14import { logOnce } from "@/worker/git/object-store/support";
15import { executeReceivePipeline, ReceivePipelineHttpError } from "./pipeline";
16import { readPktSectionStream } from "./pktSectionStream";
17import {
18 parseReceiveRequest,
19 type ParsedReceiveRequest,
20 type ReceiveCommandList,
21 type ReceiveNegotiatedCapabilities,
22} from "./request";
23import {
24 buildReceiveResultResponse,
25 ReceiveSidebandWriter,
26 type ReceiveResponseMode,
27} from "./response";
28import {
29 buildReceiveReportStatus,
30 buildReceiveUnpackFailureReport,
31 isReceiveAbort,
32 throwIfReceiveAborted,
33} from "./support";
34 
35const RECEIVE_SUBREQUEST_BUDGET = 5_000;
36 
37type RepoStub = DurableObjectStub<RepoDurableObject>;
38type RepoStateChangeHandler = (change: {
39 changed: boolean;
40 empty: boolean;
41}) => Promise<void> | void;
42 
43function countReceiveSubrequest(cacheCtx: CacheContext, log: Logger, op: string, n = 1) {
44 if (countSubrequest(cacheCtx, n)) return;
45 logOnce(cacheCtx, `receive-soft-budget:${op}`, () => {
46 log.warn("soft-budget-exhausted", { op });
47 });
48}
49 
50function logReceiveEnd(log: Logger, status: number, extra?: Record<string, unknown>) {
51 log.info("receive:end", { status, ...extra });
52}
53 
54function buildReceiveCacheContext(
55 request: Request,
56 ctx: ExecutionContext,
57 repoId: string
58): { cacheCtx: CacheContext; limiter: SubrequestLimiter } {
59 const limiter = new SubrequestLimiter(MAX_SIMULTANEOUS_CONNECTIONS);
60 return {
61 cacheCtx: {
62 req: request,
63 ctx,
64 memo: {
65 repoId,
66 limiter,
67 subreqBudget: RECEIVE_SUBREQUEST_BUDGET,
68 },
69 },
70 limiter,
71 };
72}
73 
74function selectReceiveResponseMode(
75 capabilities: ReceiveNegotiatedCapabilities
76): ReceiveResponseMode {
77 return capabilities.sideBand64k ? "side-band-64k" : "plain";
78}
79 
80function scheduleRepoStateChange(
81 ctx: ExecutionContext,
82 onRepoStateChanged: RepoStateChangeHandler | undefined,
83 change: {
84 changed: boolean;
85 empty: boolean;
86 }
87): void {
88 if (!onRepoStateChanged || !change.changed) return;
89 ctx.waitUntil(Promise.resolve().then(() => onRepoStateChanged(change)));
90}
91 
92function buildInvalidRefResponse(args: {
93 mode: ReceiveResponseMode;
94 commands: ReceiveCommandList;
95}): Response {
96 return buildReceiveResultResponse({
97 mode: args.mode,
98 reportStatusBody: buildReceiveUnpackFailureReport(args.commands, "invalid-ref", "invalid"),
99 changed: false,
100 empty: false,
101 });
102}
103 
104function buildPreflightConflictResponse(args: {
105 mode: ReceiveResponseMode;
106 commands: ReceiveCommandList;
107 statuses: ReceiveStatus[];
108}): Response {
109 return buildReceiveResultResponse({
110 mode: args.mode,
111 reportStatusBody: buildReceiveReportStatus({
112 unpackOk: true,
113 commands: args.commands,
114 statuses: args.statuses,
115 }),
116 changed: false,
117 empty: false,
118 });
119}
120 
121function getErrorStatus(error: unknown): number {
122 if (error instanceof ReceivePipelineHttpError) {
123 return error.status;
124 }
125 
126 const message = String(error);
127 const lower = message.toLowerCase();
128 if (lower.includes("unsupported pack version") || lower.includes("pack header")) {
129 return 415;
130 }
131 if (
132 lower.includes("malformed") ||
133 lower.includes("missing") ||
134 lower.includes("ended before") ||
135 lower.includes("could not be resolved") ||
136 lower.includes("delta")
137 ) {
138 return 400;
139 }
140 return 500;
141}
142 
143function createSidebandReceiveResponse(args: {
144 env: Env;
145 repoId: string;
146 request: Request;
147 ctx: ExecutionContext;
148 stub: RepoStub;
149 log: Logger;
150 cacheCtx: CacheContext;
151 limiter: SubrequestLimiter;
152 leaseToken: string;
153 activeCatalog: PackCatalogRow[];
154 commands: ParsedReceiveRequest["commands"];
155 capabilities: ReceiveNegotiatedCapabilities;
156 packStream: ReadableStream<Uint8Array>;
157 bytesConsumed: number;
158 onRepoStateChanged?: RepoStateChangeHandler | undefined;
159}): Response {
160 const responseStream = new ReadableStream<Uint8Array>({
161 async start(controller) {
162 const writer = new ReceiveSidebandWriter(controller);
163 const onProgress = args.capabilities.quiet
164 ? undefined
165 : (message: string) => writer.progress(message);
166 let closed = false;
167 const close = () => {
168 if (closed) return;
169 closed = true;
170 controller.close();
171 };
172 
173 try {
174 const result = await executeReceivePipeline({
175 env: args.env,
176 repoId: args.repoId,
177 request: args.request,
178 ctx: args.ctx,
179 packStream: args.packStream,
180 bytesConsumed: args.bytesConsumed,
181 stub: args.stub,
182 leaseToken: args.leaseToken,
183 activeCatalog: args.activeCatalog,
184 commands: args.commands,
185 log: args.log,
186 cacheCtx: args.cacheCtx,
187 limiter: args.limiter,
188 countSubrequest: (op, n = 1) => countReceiveSubrequest(args.cacheCtx, args.log, op, n),
189 onProgress,
190 });
191 
192 scheduleRepoStateChange(args.ctx, args.onRepoStateChanged, {
193 changed: result.changed,
194 empty: result.empty,
195 });
196 writer.reportStatus(result.reportStatusBody);
197 logReceiveEnd(args.log, 200, {
198 changed: result.changed,
199 empty: result.empty,
200 packKey: result.packKey,
201 packBytes: result.packBytes,
202 });
203 } catch (error) {
204 if (isReceiveAbort(args.request, error)) {
205 logReceiveEnd(args.log, 499, { reason: "client-aborted" });
206 close();
207 return;
208 }
209 
210 args.log.error("receive:error", { error: String(error) });
211 writer.reportStatus(
212 buildReceiveUnpackFailureReport(
213 args.commands,
214 error instanceof ReceivePipelineHttpError ? error.message : String(error)
215 )
216 );
217 logReceiveEnd(args.log, 200, { reason: "sideband-unpack-error" });
218 } finally {
219 close();
220 }
221 },
222 });
223 
224 return new Response(responseStream, {
225 status: 200,
226 headers: {
227 "Content-Type": "application/x-git-receive-pack-result",
228 // Credentialed mutating path; never share in caches.
229 "Cache-Control": "no-store",
230 },
231 });
232}
233 
234export async function handleStreamingReceivePackPOST(
235 env: Env,
236 repoId: string,
237 request: Request,
238 ctx: ExecutionContext,
239 options?: {
240 onRepoStateChanged?: RepoStateChangeHandler | undefined;
241 }
242): Promise<Response> {
243 const stub = getRepoStub(env, repoId);
244 const log = createLogger(env.LOG_LEVEL, {
245 service: "StreamingReceivePack",
246 repoId,
247 });
248 log.info("receive:start", { mode: "streaming" });
249 
250 if (!request.body) {
251 logReceiveEnd(log, 400, { reason: "missing-body" });
252 return new Response("Missing receive-pack request body\n", { status: 400 });
253 }
254 if (request.signal.aborted) {
255 logReceiveEnd(log, 499, { reason: "client-aborted" });
256 return clientAbortedResponse();
257 }
258 
259 const { cacheCtx, limiter } = buildReceiveCacheContext(request, ctx, repoId);
260 
261 countReceiveSubrequest(cacheCtx, log, "do:begin-receive");
262 const begin = await stub.beginReceive();
263 if (!begin.ok) {
264 log.warn("receive:block-busy", { retryAfter: begin.retryAfter, mode: "streaming" });
265 logReceiveEnd(log, 503, { reason: "receive-lease-active" });
266 return new Response("Repository is busy receiving; please retry shortly.\n", {
267 status: 503,
268 headers: {
269 "Retry-After": String(begin.retryAfter),
270 "Content-Type": "text/plain; charset=utf-8",
271 },
272 });
273 }
274 
275 let pipelineStarted = false;
276 try {
277 const { lines, bytesConsumed, packStream } = await readPktSectionStream(request.body);
278 throwIfReceiveAborted(request, log, "read-command-section");
279 
280 const parsedRequest = parseReceiveRequest(lines);
281 const responseMode = selectReceiveResponseMode(parsedRequest.capabilities);
282 
283 const invalidCommand = parsedRequest.commands.find((command) => !isValidRefName(command.ref));
284 if (invalidCommand) {
285 countReceiveSubrequest(cacheCtx, log, "do:abort-receive");
286 await stub.abortReceive(begin.lease.token).catch(() => {});
287 log.warn("receive:invalid-ref", { ref: invalidCommand.ref });
288 const response = buildInvalidRefResponse({
289 mode: responseMode,
290 commands: parsedRequest.commands,
291 });
292 logReceiveEnd(log, response.status, { reason: "invalid-ref", changed: false, empty: false });
293 return response;
294 }
295 
296 const preflightStatuses = validateReceiveCommands(begin.refs, parsedRequest.commands);
297 if (!preflightStatuses.every((status) => status.ok)) {
298 countReceiveSubrequest(cacheCtx, log, "do:abort-receive");
299 await stub.abortReceive(begin.lease.token).catch(() => {});
300 log.warn("receive:ref-conflict", {
301 conflictCount: preflightStatuses.filter((status) => !status.ok).length,
302 stage: "preflight",
303 });
304 const response = buildPreflightConflictResponse({
305 mode: responseMode,
306 commands: parsedRequest.commands,
307 statuses: preflightStatuses,
308 });
309 logReceiveEnd(log, response.status, { reason: "preflight-ref-conflict", changed: false });
310 return response;
311 }
312 
313 if (responseMode === "side-band-64k") {
314 return createSidebandReceiveResponse({
315 env,
316 repoId,
317 request,
318 ctx,
319 stub,
320 log,
321 cacheCtx,
322 limiter,
323 leaseToken: begin.lease.token,
324 activeCatalog: begin.activeCatalog,
325 commands: parsedRequest.commands,
326 capabilities: parsedRequest.capabilities,
327 packStream,
328 bytesConsumed,
329 onRepoStateChanged: options?.onRepoStateChanged,
330 });
331 }
332 
333 pipelineStarted = true;
334 const result = await executeReceivePipeline({
335 env,
336 repoId,
337 request,
338 ctx,
339 packStream,
340 bytesConsumed,
341 stub,
342 leaseToken: begin.lease.token,
343 activeCatalog: begin.activeCatalog,
344 commands: parsedRequest.commands,
345 log,
346 cacheCtx,
347 limiter,
348 countSubrequest: (op, n = 1) => countReceiveSubrequest(cacheCtx, log, op, n),
349 });
350 
351 scheduleRepoStateChange(ctx, options?.onRepoStateChanged, {
352 changed: result.changed,
353 empty: result.empty,
354 });
355 
356 const response = buildReceiveResultResponse({
357 mode: "plain",
358 reportStatusBody: result.reportStatusBody,
359 changed: result.changed,
360 empty: result.empty,
361 });
362 logReceiveEnd(log, response.status, {
363 changed: result.changed,
364 empty: result.empty,
365 packKey: result.packKey,
366 packBytes: result.packBytes,
367 });
368 return response;
369 } catch (error) {
370 if (!pipelineStarted) {
371 countReceiveSubrequest(cacheCtx, log, "do:abort-receive");
372 await stub.abortReceive(begin.lease.token).catch(() => {});
373 }
374 
375 if (isReceiveAbort(request, error)) {
376 log.info("receive:aborted", { error: String(error) });
377 logReceiveEnd(log, 499, { reason: "client-aborted" });
378 return clientAbortedResponse();
379 }
380 
381 log.error("receive:error", { error: String(error) });
382 
383 if (error instanceof ReceivePipelineHttpError) {
384 const response = new Response(`${error.message}\n`, {
385 status: error.status,
386 headers: { "Content-Type": "text/plain; charset=utf-8" },
387 });
388 logReceiveEnd(log, response.status, { reason: error.reason });
389 return response;
390 }
391 
392 const response = new Response(`${String(error)}\n`, {
393 status: getErrorStatus(error),
394 headers: { "Content-Type": "text/plain; charset=utf-8" },
395 });
396 logReceiveEnd(log, response.status, { reason: "error" });
397 return response;
398 }
399}