Skip to content
File

Blob: src/worker/dispatch/shared/run-steps/command-stream.ts

typescript122 lines
1import {
2 parseSSEStream,
3 type ExecEvent,
4 type ExecutionSession,
5 type Process,
6 type StreamOptions,
7} from "@cloudflare/sandbox";
8 
9import type { LogStream } from "@/contracts";
10 
11import type { LogBatcher } from "./logging";
12 
13export interface CommandStreamResult {
14 stdout: string;
15 stderr: string;
16 exitCode: number | null;
17 terminalEvent: "complete" | "error" | "interrupted";
18 errorMessage: string | null;
19}
20 
21interface CommandStreamHooks {
22 onStart?: (event: ExecEvent) => Promise<void> | void;
23}
24 
25const 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 
39export 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 
95const isLiveProcess = (process: Process): boolean => process.status === "starting" || process.status === "running";
96 
97export 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};