File
Blob: src/worker/git/receive/streamReceivePack.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | import type { RepoDurableObject } from "@/worker/do"; |
| 4 | import type { PackCatalogRow } from "@/worker/do/repo/db/schema"; |
| 5 | import type { ReceiveStatus } from "@/worker/git/operations/validation"; |
| 6 | |
| 7 | import { clientAbortedResponse, createLogger, getRepoStub } from "@/worker/common"; |
| 8 | import { |
| 9 | MAX_SIMULTANEOUS_CONNECTIONS, |
| 10 | SubrequestLimiter, |
| 11 | countSubrequest, |
| 12 | } from "@/worker/git/operations/limits"; |
| 13 | import { isValidRefName, validateReceiveCommands } from "@/worker/git/operations/validation"; |
| 14 | import { logOnce } from "@/worker/git/object-store/support"; |
| 15 | import { executeReceivePipeline, ReceivePipelineHttpError } from "./pipeline"; |
| 16 | import { readPktSectionStream } from "./pktSectionStream"; |
| 17 | import { |
| 18 | parseReceiveRequest, |
| 19 | type ParsedReceiveRequest, |
| 20 | type ReceiveCommandList, |
| 21 | type ReceiveNegotiatedCapabilities, |
| 22 | } from "./request"; |
| 23 | import { |
| 24 | buildReceiveResultResponse, |
| 25 | ReceiveSidebandWriter, |
| 26 | type ReceiveResponseMode, |
| 27 | } from "./response"; |
| 28 | import { |
| 29 | buildReceiveReportStatus, |
| 30 | buildReceiveUnpackFailureReport, |
| 31 | isReceiveAbort, |
| 32 | throwIfReceiveAborted, |
| 33 | } from "./support"; |
| 34 | |
| 35 | const RECEIVE_SUBREQUEST_BUDGET = 5_000; |
| 36 | |
| 37 | type RepoStub = DurableObjectStub<RepoDurableObject>; |
| 38 | type RepoStateChangeHandler = (change: { |
| 39 | changed: boolean; |
| 40 | empty: boolean; |
| 41 | }) => Promise<void> | void; |
| 42 | |
| 43 | function 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 | |
| 50 | function logReceiveEnd(log: Logger, status: number, extra?: Record<string, unknown>) { |
| 51 | log.info("receive:end", { status, ...extra }); |
| 52 | } |
| 53 | |
| 54 | function 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 | |
| 74 | function selectReceiveResponseMode( |
| 75 | capabilities: ReceiveNegotiatedCapabilities |
| 76 | ): ReceiveResponseMode { |
| 77 | return capabilities.sideBand64k ? "side-band-64k" : "plain"; |
| 78 | } |
| 79 | |
| 80 | function 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 | |
| 92 | function 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 | |
| 104 | function 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 | |
| 121 | function 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 | |
| 143 | function 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 | |
| 234 | export 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 | } |