File
Blob: src/worker/git/operations/fetch/sideband.ts
| 1 | import { pktLine, flushPkt } from "@/worker/git/core"; |
| 2 | import type { Logger } from "@/worker/common/logger"; |
| 3 | |
| 4 | export const SIDEBAND_PAYLOAD_MAX_BYTES = 65_515; |
| 5 | |
| 6 | type SidebandEnqueueController = { |
| 7 | enqueue(chunk: Uint8Array): void; |
| 8 | }; |
| 9 | |
| 10 | export 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 | |
| 31 | export 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 | |
| 42 | export 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 | |
| 97 | export 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 | |
| 113 | export function emitProgress( |
| 114 | controller: ReadableStreamDefaultController<Uint8Array>, |
| 115 | message: string |
| 116 | ) { |
| 117 | enqueueSidebandPayload(controller, 2, new TextEncoder().encode(message)); |
| 118 | } |
| 119 | |
| 120 | export function emitFatal( |
| 121 | controller: ReadableStreamDefaultController<Uint8Array>, |
| 122 | message: string |
| 123 | ) { |
| 124 | enqueueSidebandPayload(controller, 3, new TextEncoder().encode(`fatal: ${message}\n`)); |
| 125 | } |
| 126 | |
| 127 | export 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 | } |