( source: ReadableStream<Uint8Array>, )
| 50 | } |
| 51 | |
| 52 | export function decodeRuntimeEvents( |
| 53 | source: ReadableStream<Uint8Array>, |
| 54 | ): ReadableStream<WorkspaceRuntimeEvent<ExecEncoding>> { |
| 55 | const decoder = new TextDecoder(); |
| 56 | let buffered = ""; |
| 57 | return source.pipeThrough( |
| 58 | new TransformStream<Uint8Array, WorkspaceRuntimeEvent<ExecEncoding>>({ |
| 59 | transform(chunk, controller) { |
| 60 | buffered += decoder.decode(chunk, { stream: true }); |
| 61 | const lines = buffered.split("\n"); |
| 62 | buffered = lines.pop() ?? ""; |
| 63 | for (const line of lines) { |
| 64 | if (line) controller.enqueue(decodeFrame(JSON.parse(line) as RuntimeFrame)); |
| 65 | } |
| 66 | }, |
| 67 | flush(controller) { |
| 68 | buffered += decoder.decode(); |
| 69 | if (buffered.trim()) controller.enqueue(decodeFrame(JSON.parse(buffered) as RuntimeFrame)); |
| 70 | }, |
| 71 | }), |
| 72 | ); |
| 73 | } |
no test coverage detected