File
Blob: src/worker/git/receive/pktSectionStream.ts
| 1 | import { DELIM, FLUSH, RESPONSE_END } from "@/worker/git/core/pktline"; |
| 2 | import { appendBytes, cloneBytes } from "./bytes"; |
| 3 | |
| 4 | const MAX_COMMAND_SECTION_BYTES = 256 * 1024; |
| 5 | |
| 6 | type ParsedPktSection = |
| 7 | | { |
| 8 | status: "ok"; |
| 9 | lines: string[]; |
| 10 | offset: number; |
| 11 | } |
| 12 | | { |
| 13 | status: "incomplete"; |
| 14 | } |
| 15 | | { |
| 16 | status: "invalid"; |
| 17 | message: string; |
| 18 | }; |
| 19 | |
| 20 | function parsePktSectionPrefix(bytes: Uint8Array): ParsedPktSection { |
| 21 | const decoder = new TextDecoder(); |
| 22 | const lines: string[] = []; |
| 23 | let offset = 0; |
| 24 | |
| 25 | while (offset + 4 <= bytes.byteLength) { |
| 26 | const header = decoder.decode(bytes.subarray(offset, offset + 4)); |
| 27 | offset += 4; |
| 28 | |
| 29 | if (header === FLUSH) { |
| 30 | return { |
| 31 | status: "ok", |
| 32 | lines, |
| 33 | offset, |
| 34 | }; |
| 35 | } |
| 36 | |
| 37 | if (header === DELIM || header === RESPONSE_END) { |
| 38 | return { |
| 39 | status: "invalid", |
| 40 | message: "Unexpected protocol delimiter in receive-pack command section.", |
| 41 | }; |
| 42 | } |
| 43 | |
| 44 | const length = Number.parseInt(header, 16); |
| 45 | if (!Number.isFinite(length) || length < 4) { |
| 46 | return { |
| 47 | status: "invalid", |
| 48 | message: "Malformed pkt-line header in receive-pack command section.", |
| 49 | }; |
| 50 | } |
| 51 | |
| 52 | const payloadLength = length - 4; |
| 53 | if (offset + payloadLength > bytes.byteLength) { |
| 54 | return { status: "incomplete" }; |
| 55 | } |
| 56 | |
| 57 | const payload = bytes.subarray(offset, offset + payloadLength); |
| 58 | offset += payloadLength; |
| 59 | lines.push(decoder.decode(payload).replace(/\r?\n$/, "")); |
| 60 | } |
| 61 | |
| 62 | return { status: "incomplete" }; |
| 63 | } |
| 64 | |
| 65 | export async function readPktSectionStream(body: ReadableStream<Uint8Array>): Promise<{ |
| 66 | lines: string[]; |
| 67 | bytesConsumed: number; |
| 68 | packStream: ReadableStream<Uint8Array>; |
| 69 | }> { |
| 70 | const reader = body.getReader(); |
| 71 | let buffered: Uint8Array<ArrayBufferLike> = new Uint8Array(0); |
| 72 | |
| 73 | while (true) { |
| 74 | const parsed = parsePktSectionPrefix(buffered); |
| 75 | if (parsed.status === "ok") { |
| 76 | const prefixRemainder = buffered.subarray(parsed.offset); |
| 77 | const packStream = new ReadableStream<Uint8Array>({ |
| 78 | start(controller) { |
| 79 | if (prefixRemainder.byteLength > 0) { |
| 80 | controller.enqueue(prefixRemainder); |
| 81 | } |
| 82 | }, |
| 83 | async pull(controller) { |
| 84 | const next = await reader.read(); |
| 85 | if (next.done) { |
| 86 | controller.close(); |
| 87 | return; |
| 88 | } |
| 89 | controller.enqueue(next.value); |
| 90 | }, |
| 91 | async cancel(reason) { |
| 92 | await reader.cancel(reason); |
| 93 | }, |
| 94 | }); |
| 95 | |
| 96 | return { |
| 97 | lines: parsed.lines, |
| 98 | bytesConsumed: parsed.offset, |
| 99 | packStream, |
| 100 | }; |
| 101 | } |
| 102 | |
| 103 | if (parsed.status === "invalid") { |
| 104 | await reader.cancel(parsed.message); |
| 105 | throw new Error(parsed.message); |
| 106 | } |
| 107 | |
| 108 | const next = await reader.read(); |
| 109 | if (next.done) { |
| 110 | await reader.cancel("receive-pack command section ended before flush"); |
| 111 | throw new Error("Receive-pack command section ended before the flush packet."); |
| 112 | } |
| 113 | |
| 114 | buffered = appendBytes(buffered, cloneBytes(next.value)); |
| 115 | if (buffered.byteLength > MAX_COMMAND_SECTION_BYTES) { |
| 116 | await reader.cancel("receive-pack command section exceeded the supported size"); |
| 117 | throw new Error("Receive-pack command section exceeded the supported size."); |
| 118 | } |
| 119 | } |
| 120 | } |