Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158import * 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<ShellResult> { return Stream.unwrap( Effect.sync(() => { const owner = host(); const controller = new AbortController(); let lease: { id: string; expiresAt: number } | undefined; let timer: ReturnType<typeof setTimeout> | undefined; let closing: Promise<boolean> | 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<ExecEvent>(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 }))), ); }), );}