File
Blob: src/worker/dispatch/shared/run-steps/command-stream.ts
| 1 | import { |
| 2 | parseSSEStream, |
| 3 | type ExecEvent, |
| 4 | type ExecutionSession, |
| 5 | type Process, |
| 6 | type StreamOptions, |
| 7 | } from "@cloudflare/sandbox"; |
| 8 | |
| 9 | import type { LogStream } from "@/contracts"; |
| 10 | |
| 11 | import type { LogBatcher } from "./logging"; |
| 12 | |
| 13 | export interface CommandStreamResult { |
| 14 | stdout: string; |
| 15 | stderr: string; |
| 16 | exitCode: number | null; |
| 17 | terminalEvent: "complete" | "error" | "interrupted"; |
| 18 | errorMessage: string | null; |
| 19 | } |
| 20 | |
| 21 | interface CommandStreamHooks { |
| 22 | onStart?: (event: ExecEvent) => Promise<void> | void; |
| 23 | } |
| 24 | |
| 25 | const appendChunk = ( |
| 26 | result: CommandStreamResult, |
| 27 | batcher: LogBatcher, |
| 28 | stream: Exclude<LogStream, "system">, |
| 29 | chunk: string | undefined, |
| 30 | ): void => { |
| 31 | if (!chunk) { |
| 32 | return; |
| 33 | } |
| 34 | |
| 35 | batcher.push(stream, chunk); |
| 36 | result[stream] += chunk; |
| 37 | }; |
| 38 | |
| 39 | export const executeSessionCommandStream = async ( |
| 40 | session: Pick<ExecutionSession, "execStream">, |
| 41 | command: string, |
| 42 | options: StreamOptions | undefined, |
| 43 | batcher: LogBatcher, |
| 44 | hooks?: CommandStreamHooks, |
| 45 | ): Promise<CommandStreamResult> => { |
| 46 | const result: CommandStreamResult = { |
| 47 | stdout: "", |
| 48 | stderr: "", |
| 49 | exitCode: null, |
| 50 | terminalEvent: "interrupted", |
| 51 | errorMessage: null, |
| 52 | }; |
| 53 | |
| 54 | try { |
| 55 | const stream = await session.execStream(command, options); |
| 56 | for await (const event of parseSSEStream<ExecEvent>(stream)) { |
| 57 | switch (event.type) { |
| 58 | case "start": |
| 59 | await hooks?.onStart?.(event); |
| 60 | break; |
| 61 | case "stdout": |
| 62 | appendChunk(result, batcher, "stdout", event.data); |
| 63 | break; |
| 64 | case "stderr": |
| 65 | appendChunk(result, batcher, "stderr", event.data); |
| 66 | break; |
| 67 | case "complete": |
| 68 | result.terminalEvent = "complete"; |
| 69 | result.exitCode = event.exitCode ?? event.result?.exitCode ?? 1; |
| 70 | break; |
| 71 | case "error": |
| 72 | result.terminalEvent = "error"; |
| 73 | result.exitCode = 1; |
| 74 | result.errorMessage = event.error ?? event.data ?? "Command stream emitted an error event."; |
| 75 | break; |
| 76 | } |
| 77 | |
| 78 | if (result.terminalEvent !== "interrupted") { |
| 79 | break; |
| 80 | } |
| 81 | } |
| 82 | } catch (error) { |
| 83 | result.errorMessage = error instanceof Error ? error.message : String(error); |
| 84 | } finally { |
| 85 | await batcher.flush(); |
| 86 | } |
| 87 | |
| 88 | if (result.terminalEvent === "interrupted" && result.errorMessage === null) { |
| 89 | result.errorMessage = "Command stream ended without a terminal event."; |
| 90 | } |
| 91 | |
| 92 | return result; |
| 93 | }; |
| 94 | |
| 95 | const isLiveProcess = (process: Process): boolean => process.status === "starting" || process.status === "running"; |
| 96 | |
| 97 | export const resolveExecStreamProcess = async ( |
| 98 | session: Pick<ExecutionSession, "listProcesses">, |
| 99 | command: string, |
| 100 | pid?: number, |
| 101 | ): Promise<Process | null> => { |
| 102 | const liveProcesses = (await session.listProcesses()).filter(isLiveProcess); |
| 103 | if (pid !== undefined) { |
| 104 | const pidMatch = liveProcesses.find((process) => process.pid === pid); |
| 105 | if (pidMatch) { |
| 106 | return pidMatch; |
| 107 | } |
| 108 | } |
| 109 | |
| 110 | const commandMatches = liveProcesses |
| 111 | .filter((process) => process.command === command) |
| 112 | .sort((left, right) => right.startTime.getTime() - left.startTime.getTime()); |
| 113 | if (commandMatches.length > 0) { |
| 114 | return commandMatches[0] ?? null; |
| 115 | } |
| 116 | |
| 117 | const latestLiveProcess = [...liveProcesses].sort( |
| 118 | (left, right) => right.startTime.getTime() - left.startTime.getTime(), |
| 119 | ); |
| 120 | return latestLiveProcess[0] ?? null; |
| 121 | }; |