Skip to content
File

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

typescript121 lines
1import { DELIM, FLUSH, RESPONSE_END } from "@/worker/git/core/pktline";
2import { appendBytes, cloneBytes } from "./bytes";
3 
4const MAX_COMMAND_SECTION_BYTES = 256 * 1024;
5 
6type 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 
20function 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 
65export 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}