Skip to content
File

Blob: src/worker/git/operations/fetch/sideband.ts

typescript169 lines
1import { pktLine, flushPkt } from "@/worker/git/core";
2import type { Logger } from "@/worker/common/logger";
3 
4export const SIDEBAND_PAYLOAD_MAX_BYTES = 65_515;
5 
6type SidebandEnqueueController = {
7 enqueue(chunk: Uint8Array): void;
8};
9 
10export function createSidebandPacketChunks(
11 band: 1 | 2 | 3,
12 payload: Uint8Array,
13 maxChunk: number = SIDEBAND_PAYLOAD_MAX_BYTES
14): Uint8Array[] {
15 const chunks: Uint8Array[] = [];
16 for (let off = 0; off < payload.byteLength; off += maxChunk) {
17 const slice = payload.subarray(off, Math.min(off + maxChunk, payload.byteLength));
18 const banded = new Uint8Array(1 + slice.byteLength);
19 banded[0] = band;
20 banded.set(slice, 1);
21 chunks.push(pktLine(banded));
22 }
23 if (payload.byteLength === 0) {
24 const banded = new Uint8Array(1);
25 banded[0] = band;
26 chunks.push(pktLine(banded));
27 }
28 return chunks;
29}
30 
31export function enqueueSidebandPayload(
32 controller: SidebandEnqueueController,
33 band: 1 | 2 | 3,
34 payload: Uint8Array,
35 maxChunk?: number
36): void {
37 for (const chunk of createSidebandPacketChunks(band, payload, maxChunk)) {
38 controller.enqueue(chunk);
39 }
40}
41 
42export class SidebandProgressMux {
43 private progressMessages: string[] = [];
44 private progressIdx = 0;
45 private lastProgressTime = 0;
46 private inProgress = false;
47 private resolveFirstProgress?: () => void;
48 private firstProgressPromise: Promise<void>;
49 private readonly intervalMs: number;
50 
51 constructor(intervalMs = 100) {
52 this.intervalMs = intervalMs;
53 this.firstProgressPromise = new Promise<void>((resolve) => {
54 this.resolveFirstProgress = resolve;
55 });
56 }
57 
58 push(msg: string): void {
59 this.progressMessages.push(msg);
60 if (this.resolveFirstProgress) {
61 this.resolveFirstProgress();
62 this.resolveFirstProgress = undefined;
63 }
64 }
65 
66 async waitForFirst(timeoutMs = 20): Promise<void> {
67 await Promise.race([this.firstProgressPromise, new Promise((r) => setTimeout(r, timeoutMs))]);
68 }
69 
70 shouldSendProgress(): boolean {
71 const now = Date.now();
72 return (
73 now - this.lastProgressTime >= this.intervalMs &&
74 !this.inProgress &&
75 this.progressIdx < this.progressMessages.length
76 );
77 }
78 
79 async sendPending(emitFn: (msg: string) => void): Promise<void> {
80 if (this.shouldSendProgress()) {
81 this.inProgress = true;
82 while (this.progressIdx < this.progressMessages.length) {
83 emitFn(this.progressMessages[this.progressIdx++]);
84 }
85 this.lastProgressTime = Date.now();
86 this.inProgress = false;
87 }
88 }
89 
90 sendRemaining(emitFn: (msg: string) => void): void {
91 while (this.progressIdx < this.progressMessages.length) {
92 emitFn(this.progressMessages[this.progressIdx++]);
93 }
94 }
95}
96 
97export function createSidebandTransform(options?: {
98 onProgress?: (msg: string) => void;
99 signal?: AbortSignal;
100}): TransformStream<Uint8Array, Uint8Array> {
101 return new TransformStream<Uint8Array, Uint8Array>({
102 async transform(chunk, controller) {
103 if (options?.signal?.aborted) {
104 controller.terminate();
105 return;
106 }
107 
108 enqueueSidebandPayload(controller, 1, chunk);
109 },
110 });
111}
112 
113export function emitProgress(
114 controller: ReadableStreamDefaultController<Uint8Array>,
115 message: string
116) {
117 enqueueSidebandPayload(controller, 2, new TextEncoder().encode(message));
118}
119 
120export function emitFatal(
121 controller: ReadableStreamDefaultController<Uint8Array>,
122 message: string
123) {
124 enqueueSidebandPayload(controller, 3, new TextEncoder().encode(`fatal: ${message}\n`));
125}
126 
127export async function pipePackWithSideband(
128 packStream: ReadableStream<Uint8Array>,
129 controller: ReadableStreamDefaultController<Uint8Array>,
130 options: {
131 signal?: AbortSignal;
132 progressMux: SidebandProgressMux;
133 log: Logger;
134 }
135): Promise<void> {
136 const { signal, progressMux, log } = options;
137 
138 try {
139 const sidebandTransform = createSidebandTransform({ signal });
140 const reader = packStream.pipeThrough(sidebandTransform).getReader();
141 
142 await progressMux.waitForFirst();
143 progressMux.sendRemaining((msg) => emitProgress(controller, msg));
144 
145 while (true) {
146 if (signal?.aborted) {
147 log.debug("pipe:aborted");
148 reader.cancel();
149 break;
150 }
151 
152 const { done, value } = await reader.read();
153 if (done) break;
154 
155 await progressMux.sendPending((msg) => emitProgress(controller, msg));
156 controller.enqueue(value);
157 }
158 
159 progressMux.sendRemaining((msg) => emitProgress(controller, msg));
160 controller.enqueue(flushPkt());
161 } catch (error) {
162 log.error("pipe:error", { error: String(error) });
163 try {
164 emitFatal(controller, String(error));
165 } catch {}
166 throw error;
167 }
168}