import * as Effect from "effect/Effect"; import * as Stream from "effect/Stream"; import { parseSSEStream, type ExecEvent } from "@cloudflare/sandbox"; import { SHELL_OUTPUT_BYTES, type ShellResult } from "../shared/shell"; import { AgentFailure, agentCall } from "./agent-io"; import type { ShellHost } from "./shell-tool"; export function shellResults( host: () => ShellHost, command: string, caller?: AbortSignal, ): Stream.Stream { return Stream.unwrap( Effect.sync(() => { const owner = host(); const controller = new AbortController(); let lease: { id: string; expiresAt: number } | undefined; let timer: ReturnType | undefined; let closing: Promise | undefined; const result: ShellResult = { status: "running", stdout: "", stderr: "", outputBytes: 0, exitCode: null, truncated: false, cleanup: "pending", workspaceLifetime: "invocation", }; const close = () => Effect.suspend(() => { if (!lease) return Effect.succeed(false); const id = lease.id; closing ??= Effect.runPromise( agentCall(() => owner.close(id)).pipe( Effect.timeoutOrElse({ duration: 5_000, orElse: () => Effect.succeed(false), }), Effect.catchTag("AgentFailure", () => Effect.succeed(false)), ), ); return Effect.promise(() => closing!); }); const stop = (reason: "cancelled" | "timeout" | "output_limit") => { if (!controller.signal.aborted) controller.abort(reason); void Effect.runPromise(close()); }; const cancel = () => stop("cancelled"); const release = Effect.gen(function* () { if (timer) clearTimeout(timer); caller?.removeEventListener("abort", cancel); if (lease && Date.now() >= lease.expiresAt && !caller?.aborted) result.error = "timeout"; if (controller.signal.aborted) result.error = controller.signal.reason as ShellResult["error"]; if (result.status === "running" || result.error) result.status = "failed"; result.error ??= result.status === "failed" ? "unavailable" : undefined; result.cleanup = (yield* close()) ? "closed" : "pending"; }); const execute = Stream.unwrap( Effect.gen(function* () { yield* Effect.acquireRelease( Effect.sync(() => { caller?.addEventListener("abort", cancel, { once: true }); if (caller?.aborted) cancel(); }), () => release, ); if (controller.signal.aborted) return yield* Effect.fail( new AgentFailure({ message: "Shell interrupted" }), ); // Native reservation cannot be interrupted. Hand its result to this // scope before allowing iterator.return() to run the release finalizer. const reservation = yield* agentCall(() => owner.reserve()).pipe( Effect.tap((reserved) => Effect.sync(() => { lease = reserved; }), ), Effect.uninterruptible, ); if (controller.signal.aborted) return yield* Effect.fail( new AgentFailure({ message: "Shell interrupted" }), ); const id = reservation.id; timer = setTimeout( () => stop("timeout"), Math.max(0, reservation.expiresAt - Date.now()), ); let lastYield = 0; const events = Stream.unwrap( agentCall(() => owner.launch(id, command)).pipe( Effect.map((stream) => Stream.fromAsyncIterable( parseSSEStream(stream, controller.signal), () => new AgentFailure({ message: "Shell unavailable" }), ), ), ), ).pipe( Stream.map((event) => { let emit = false, done = false; if (event.type === "stdout" || event.type === "stderr") { const bytes = new TextEncoder().encode(event.data ?? ""), remaining = SHELL_OUTPUT_BYTES - result.outputBytes; result[event.type] += new TextDecoder().decode( bytes.subarray(0, remaining), { stream: bytes.byteLength > remaining }, ); result.outputBytes += Math.min(bytes.byteLength, remaining); if (bytes.byteLength >= remaining) { result.truncated = true; result.error = "output_limit"; stop("output_limit"); done = true; } else if (Date.now() - lastYield >= 250) { lastYield = Date.now(); emit = true; } } else if (event.type === "complete") { result.exitCode = event.exitCode ?? event.result?.exitCode ?? null; result.status = result.exitCode === 0 ? "succeeded" : "failed"; if (result.status === "failed") result.error = "command_failed"; done = true; } else if (event.type === "error") { result.error = "unavailable"; done = true; } return { done, value: emit ? { ...result } : undefined }; }), Stream.takeUntil((event) => event.done), Stream.map((event) => event.value), Stream.filter((value): value is ShellResult => value !== undefined), ); return Stream.concat(Stream.make({ ...result }), events); }), ).pipe( Stream.catchCause(() => { result.error ??= "unavailable"; return Stream.empty; }), ); // Close before emitting the final result. An early iterator return also // closes this scope, without starting another command or consuming more SSE. return Stream.concat( Stream.scoped(execute), Stream.fromEffect(Effect.sync(() => ({ ...result }))), ); }), ); }