diff --git a/packages/opencode/src/cli/cmd/tui.ts b/packages/opencode/src/cli/cmd/tui.ts index 95ffac7ea51d..8a328c3a3f5d 100644 --- a/packages/opencode/src/cli/cmd/tui.ts +++ b/packages/opencode/src/cli/cmd/tui.ts @@ -218,83 +218,99 @@ export const TuiThreadCommand = cmd({ } process.on("SIGUSR2", reload) + let triggerFatal: ((error: unknown) => void) | undefined + let stopped = false const stop = async () => { if (stopped) return stopped = true process.off("SIGUSR2", reload) await withTimeout(client.call("shutdown", undefined), 5000).catch(() => {}) + // Keep disconnects observable until immediately before our own termination. + client.expectDisconnect() worker.terminate() } - const prompt = await input(args.prompt) - const config = await TuiConfig.get() + try { + const prompt = await input(args.prompt) + const config = await TuiConfig.get() - const network = resolveNetworkOptionsNoConfig(args) - const external = hasArg("--port") || hasArg("--hostname") || network.mdns === true + const network = resolveNetworkOptionsNoConfig(args) + const external = hasArg("--port") || hasArg("--hostname") || network.mdns === true - const headers = external ? ServerAuth.headers() : undefined + const headers = external ? ServerAuth.headers() : undefined - const transport = external - ? { - url: (await client.call("server", network)).url, - fetch: undefined, - events: undefined, - headers, - } - : { - url: "http://opencode.internal", - fetch: createWorkerFetch(client), - events: createEventSource(client), - } + const transport = external + ? { + url: (await client.call("server", network)).url, + fetch: undefined, + events: undefined, + headers, + } + : { + url: "http://opencode.internal", + fetch: createWorkerFetch(client), + events: createEventSource(client), + } - try { - await validateSession({ - url: transport.url, - sessionID: args.session, - directory: cwd, - fetch: transport.fetch, - headers, - }) - } catch (error) { - UI.error(errorMessage(error)) - process.exitCode = 1 - return - } + try { + await validateSession({ + url: transport.url, + sessionID: args.session, + directory: cwd, + fetch: transport.fetch, + headers, + }) + } catch (error) { + UI.error(errorMessage(error)) + process.exitCode = 1 + return + } - setTimeout(() => { - client.call("checkUpgrade", { directory: cwd }).catch(() => {}) - }, 1000).unref?.() + setTimeout(() => { + client.call("checkUpgrade", { directory: cwd }).catch(() => {}) + }, 1000).unref?.() - try { const { Effect } = await import("effect") const { run } = await import("../tui/layer") const { createLegacyTuiPluginHost } = await import("@/plugin/tui/runtime") - await Effect.runPromise( - run({ - url: transport.url, - async onSnapshot() { - const tui = writeHeapSnapshot("tui.heapsnapshot") - const server = await client.call("snapshot", undefined) - return [tui, server] - }, - config, - pluginHost: createLegacyTuiPluginHost(), - directory: cwd, - fetch: transport.fetch, - headers: transport.headers, - events: transport.events, - args: { - continue: args.continue, - sessionID: args.session, - agent: args.agent, - model: args.model, - prompt, - fork: args.fork, - auto: args.auto || args.yolo || args["dangerously-skip-permissions"], - }, - }), - ) + try { + await Effect.runPromise( + run({ + url: transport.url, + async onSnapshot() { + const tui = writeHeapSnapshot("tui.heapsnapshot") + const server = await client.call("snapshot", undefined) + return [tui, server] + }, + config, + pluginHost: createLegacyTuiPluginHost(), + directory: cwd, + fetch: transport.fetch, + headers: transport.headers, + events: transport.events, + onReady: (controls) => { + triggerFatal = controls.triggerFatal + client.onDisconnect((error) => { + if (triggerFatal) return triggerFatal(error) + UI.error(errorMessage(error)) + process.exitCode = 1 + }) + }, + args: { + continue: args.continue, + sessionID: args.session, + agent: args.agent, + model: args.model, + prompt, + fork: args.fork, + auto: args.auto || args.yolo || args["dangerously-skip-permissions"], + }, + }), + ) + } finally { + triggerFatal = undefined + } } finally { await stop() } @@ -303,7 +319,7 @@ export const TuiThreadCommand = cmd({ unguard?.() } catch {} } - process.exit(0) + process.exit(process.exitCode ?? 0) }, }) // scratch diff --git a/packages/opencode/src/util/rpc.ts b/packages/opencode/src/util/rpc.ts index 02586ebcfc60..6e2cb27323f2 100644 --- a/packages/opencode/src/util/rpc.ts +++ b/packages/opencode/src/util/rpc.ts @@ -6,8 +6,18 @@ export function listen(rpc: Definition) { onmessage = async (evt) => { const parsed = JSON.parse(evt.data) if (parsed.type === "rpc.request") { - const result = await rpc[parsed.method](parsed.input) - postMessage(JSON.stringify({ type: "rpc.result", result, id: parsed.id })) + // Promise-wrapped so a synchronous throw also lands in .catch() below. + await new Promise((resolve) => resolve(rpc[parsed.method](parsed.input))) + .then((result) => postMessage(JSON.stringify({ type: "rpc.result", result, id: parsed.id }))) + .catch((error) => { + let message: string + try { + message = error instanceof Error ? error.message : String(error) + } catch { + message = "Unknown error" + } + postMessage(JSON.stringify({ type: "rpc.result", error: message, id: parsed.id })) + }) } } } @@ -16,37 +26,91 @@ export function emit(event: string, data: unknown) { postMessage(JSON.stringify({ type: "rpc.event", event, data })) } +type WorkerLifecycleEvent = Event & { error?: unknown; message?: string; code?: number } + export function client(target: { postMessage: (data: string) => void | null onmessage: ((this: Worker, ev: MessageEvent) => any) | null + addEventListener?: (type: "error" | "close" | "messageerror", listener: (event: WorkerLifecycleEvent) => void) => void }) { - const pending = new Map void>() + const pending = new Map void; reject: (error: unknown) => void }>() const listeners = new Map void>>() + const disconnectHandlers = new Set<(error: Error) => void>() let id = 0 + let dead: Error | undefined + let disconnectExpected = false + const rejectAllPending = (error: Error) => { + if (dead || disconnectExpected) return + dead = error + for (const entry of pending.values()) entry.reject(error) + pending.clear() + for (const handler of disconnectHandlers) handler(error) + } + target.addEventListener?.("error", (event) => { + rejectAllPending(event?.error instanceof Error ? event.error : new Error(event?.message || "Worker error")) + }) + target.addEventListener?.("close", (event) => { + // Code 0 is ambiguous, so only expectDisconnect() marks an intentional shutdown. + rejectAllPending( + new Error(`Worker exited unexpectedly${typeof event?.code === "number" ? ` (code ${event.code})` : ""}`), + ) + }) + target.addEventListener?.("messageerror", () => { + rejectAllPending(new Error("Worker sent a message that could not be deserialized")) + }) target.onmessage = async (evt) => { - const parsed = JSON.parse(evt.data) + const parsed = (() => { + try { + return JSON.parse(evt.data) + } catch (cause) { + rejectAllPending(new Error("Worker sent invalid RPC JSON", { cause })) + } + })() + if (parsed === undefined) return + if (typeof parsed !== "object" || parsed === null || !("type" in parsed)) { + rejectAllPending(new Error("Worker sent an invalid RPC message")) + return + } if (parsed.type === "rpc.result") { - const resolve = pending.get(parsed.id) - if (resolve) { - resolve(parsed.result) + if (typeof parsed.id !== "number") { + rejectAllPending(new Error("Worker sent an invalid RPC result")) + return + } + const entry = pending.get(parsed.id) + if (entry) { + if ("error" in parsed) entry.reject(new Error(parsed.error)) + else entry.resolve(parsed.result) pending.delete(parsed.id) } + return } if (parsed.type === "rpc.event") { + if (typeof parsed.event !== "string") { + rejectAllPending(new Error("Worker sent an invalid RPC event")) + return + } const handlers = listeners.get(parsed.event) if (handlers) { for (const handler of handlers) { handler(parsed.data) } } + return } + rejectAllPending(new Error("Worker sent an unknown RPC message")) } return { call(method: Method, input: Parameters[0]): Promise> { const requestId = id++ - return new Promise((resolve) => { - pending.set(requestId, resolve) - target.postMessage(JSON.stringify({ type: "rpc.request", method, input, id: requestId })) + return new Promise((resolve, reject) => { + if (dead) return reject(dead) + pending.set(requestId, { resolve, reject }) + try { + target.postMessage(JSON.stringify({ type: "rpc.request", method, input, id: requestId })) + } catch (error) { + pending.delete(requestId) + reject(error) + } }) }, on(event: string, handler: (data: Data) => void) { @@ -60,6 +124,19 @@ export function client(target: { handlers!.delete(handler) } }, + onDisconnect(handler: (error: Error) => void) { + if (dead) { + handler(dead) + return () => {} + } + disconnectHandlers.add(handler) + return () => { + disconnectHandlers.delete(handler) + } + }, + expectDisconnect() { + disconnectExpected = true + }, } } diff --git a/packages/opencode/test/fixture/crashing-worker.ts b/packages/opencode/test/fixture/crashing-worker.ts new file mode 100644 index 000000000000..b289dcb18201 --- /dev/null +++ b/packages/opencode/test/fixture/crashing-worker.ts @@ -0,0 +1 @@ +throw new Error("Worker crashed on purpose") diff --git a/packages/opencode/test/util/rpc.test.ts b/packages/opencode/test/util/rpc.test.ts new file mode 100644 index 000000000000..35f377e4c129 --- /dev/null +++ b/packages/opencode/test/util/rpc.test.ts @@ -0,0 +1,247 @@ +import { describe, expect, test, afterEach } from "bun:test" +import { Rpc } from "../../src/util/rpc" + +describe("util.rpc", () => { + const originalOnmessage = (globalThis as any).onmessage + const originalPostMessage = (globalThis as any).postMessage + + afterEach(() => { + ;(globalThis as any).onmessage = originalOnmessage + ;(globalThis as any).postMessage = originalPostMessage + }) + + test("listen replies with an error instead of hanging when a handler throws", async () => { + const sent: any[] = [] + ;(globalThis as any).postMessage = (data: string) => sent.push(JSON.parse(data)) + + Rpc.listen({ + boom: async () => { + throw new Error("no such column: name") + }, + }) + + await (globalThis as any).onmessage({ + data: JSON.stringify({ type: "rpc.request", method: "boom", input: undefined, id: 1 }), + }) + + expect(sent).toHaveLength(1) + expect(sent[0]).toMatchObject({ type: "rpc.result", id: 1, error: "no such column: name" }) + }) + + test("listen falls back to a safe message when the thrown value cannot be stringified", async () => { + const sent: any[] = [] + ;(globalThis as any).postMessage = (data: string) => sent.push(JSON.parse(data)) + const unstringifiable = new Proxy( + {}, + { + get() { + throw new Error("nope") + }, + }, + ) + + Rpc.listen({ + boom: async () => { + throw unstringifiable + }, + }) + + await (globalThis as any).onmessage({ + data: JSON.stringify({ type: "rpc.request", method: "boom", input: undefined, id: 2 }), + }) + + expect(sent).toHaveLength(1) + expect(sent[0]).toMatchObject({ type: "rpc.result", id: 2, error: "Unknown error" }) + }) + + test("client.call rejects instead of hanging forever when the reply carries an error", async () => { + const target = { postMessage: (_data: string) => {}, onmessage: null as any } + const client = Rpc.client<{ boom: (input: undefined) => void }>(target) + + queueMicrotask(() => { + target.onmessage!({ + data: JSON.stringify({ type: "rpc.result", id: 0, error: "no such column: name" }), + }) + }) + + await expect(client.call("boom", undefined)).rejects.toThrow("no such column: name") + }) + + test("client.call still resolves normally for a successful reply", async () => { + const target = { postMessage: (_data: string) => {}, onmessage: null as any } + const client = Rpc.client<{ ping: (input: undefined) => string }>(target) + + queueMicrotask(() => { + target.onmessage!({ + data: JSON.stringify({ type: "rpc.result", id: 0, result: "pong" }), + }) + }) + + await expect(client.call("ping", undefined)).resolves.toBe("pong") + }) + + test("client.call rejects instead of hanging when postMessage throws synchronously", async () => { + const target = { + postMessage: (_data: string) => { + throw new Error("channel closed") + }, + onmessage: null as any, + } + const client = Rpc.client<{ boom: (input: undefined) => void }>(target) + + await expect(client.call("boom", undefined)).rejects.toThrow("channel closed") + }) + + function createFakeWorker() { + const target = new EventTarget() as EventTarget & { + postMessage: (data: string) => void + onmessage: ((ev: MessageEvent) => any) | null + } + target.postMessage = () => {} + target.onmessage = null + return target + } + + test("an in-flight call rejects instead of hanging forever when the worker crashes", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + worker.dispatchEvent(new ErrorEvent("error", { error: new Error("segfault") })) + + await expect(pending).rejects.toThrow("segfault") + }) + + test("an in-flight call rejects instead of hanging forever when a real Bun Worker crashes", async () => { + const file = new URL("../fixture/crashing-worker.ts", import.meta.url) + const worker = new Worker(file) + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + try { + await expect(pending).rejects.toThrow("Worker crashed on purpose") + } finally { + worker.terminate() + } + }) + + test("an in-flight call rejects instead of hanging forever when the worker exits unexpectedly", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + worker.dispatchEvent(new CloseEvent("close", { code: 1 })) + + await expect(pending).rejects.toThrow(/exited unexpectedly/) + }) + + test("an in-flight call rejects instead of hanging forever when the worker sends an undeserializable message", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + worker.dispatchEvent(new MessageEvent("messageerror")) + + await expect(pending).rejects.toThrow(/could not be deserialized/) + }) + + test("an in-flight call rejects when the worker sends invalid RPC JSON", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + worker.onmessage?.(new MessageEvent("message", { data: "not json" })) + + await expect(pending).rejects.toThrow("invalid RPC JSON") + }) + + test("an in-flight call rejects when the worker sends an invalid RPC envelope", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + worker.onmessage?.(new MessageEvent("message", { data: "null" })) + + await expect(pending).rejects.toThrow("invalid RPC message") + }) + + test("an in-flight call rejects when the worker sends an unknown RPC message", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const pending = client.call("fetch", undefined) + worker.onmessage?.(new MessageEvent("message", { data: JSON.stringify({ type: "unknown" }) })) + + await expect(pending).rejects.toThrow("unknown RPC message") + }) + + test("a call made after the worker has already died rejects immediately instead of hanging", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + worker.dispatchEvent(new CloseEvent("close", { code: 1 })) + + await expect(client.call("fetch", undefined)).rejects.toThrow(/exited unexpectedly/) + }) + + test("onDisconnect fires even when the worker dies with nothing in flight", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const seen: Error[] = [] + client.onDisconnect((error) => seen.push(error)) + + worker.dispatchEvent(new ErrorEvent("error", { error: new Error("segfault") })) + + expect(seen).toHaveLength(1) + expect(seen[0].message).toBe("segfault") + }) + + test("onDisconnect does not fire once expectDisconnect() has been called", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const seen: Error[] = [] + client.onDisconnect((error) => seen.push(error)) + + client.expectDisconnect() + worker.dispatchEvent(new CloseEvent("close", { code: 0 })) + + expect(seen).toHaveLength(0) + }) + + test("an unexpected exit with code 0 is fatal without expectDisconnect()", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + const seen: Error[] = [] + client.onDisconnect((error) => seen.push(error)) + + worker.dispatchEvent(new CloseEvent("close", { code: 0 })) + + expect(seen).toHaveLength(1) + }) + + test("a crash during an in-flight call is not suppressed just because expectDisconnect() follows shortly after", async () => { + const worker = createFakeWorker() + const client = Rpc.client<{ shutdown: (input: undefined) => void }>(worker) + + const pending = client.call("shutdown", undefined) + worker.dispatchEvent(new CloseEvent("close", { code: 1 })) + await expect(pending).rejects.toThrow(/exited unexpectedly/) + + expect(() => client.expectDisconnect()).not.toThrow() + }) + + test("onDisconnect invokes a handler registered after the worker already died", () => { + const worker = createFakeWorker() + const client = Rpc.client<{ fetch: (input: undefined) => void }>(worker) + + worker.dispatchEvent(new CloseEvent("close", { code: 1 })) + + const seen: Error[] = [] + client.onDisconnect((error) => seen.push(error)) + + expect(seen).toHaveLength(1) + }) +}) diff --git a/packages/tui/src/app.tsx b/packages/tui/src/app.tsx index 57f372ef709a..07165b48d3fd 100644 --- a/packages/tui/src/app.tsx +++ b/packages/tui/src/app.tsx @@ -149,6 +149,8 @@ export type TuiInput = { headers?: RequestInit["headers"] events?: EventSource pluginHost: TuiPluginHost + // Exposes fatal shutdown after renderer setup so idle transport failures restore the terminal. + onReady?: (controls: { triggerFatal: (error: unknown) => void }) => void } function errorMessage(error: unknown) { @@ -234,6 +236,28 @@ export const run = Effect.fn("Tui.run")(function* (input: TuiInput) { () => Effect.sync(() => process.off("SIGHUP", onSighup)), ) renderer.once("destroy", () => Deferred.doneUnsafe(shutdown, Effect.void)) + // Route process and transport failures through renderer teardown and error reporting. + yield* Effect.acquireRelease( + Effect.sync(() => { + let handled = false + const onFatal = (error: unknown) => { + if (handled) return + handled = true + // undefined is the clean-exit sentinel; normalize reasonless failures. + exit.reason = error ?? new Error("Unhandled fatal error in TUI process") + if (!renderer.isDestroyed) destroyRenderer(renderer) + } + process.on("unhandledRejection", onFatal) + process.on("uncaughtException", onFatal) + input.onReady?.({ triggerFatal: onFatal }) + return onFatal + }), + (onFatal) => + Effect.sync(() => { + process.off("unhandledRejection", onFatal) + process.off("uncaughtException", onFatal) + }), + ) const pluginRuntime = createPluginRuntime() yield* Effect.tryPromise(async () => { @@ -351,13 +375,20 @@ export const run = Effect.fn("Tui.run")(function* (input: TuiInput) { }, renderer) }) yield* Deferred.await(shutdown) - return { epilogue: exit.epilogue, reason: exit.reason } + // Captured here, before scope finalizers unmount the render tree (which + // resets the epilogue via Session's onCleanup). exit.reason is read below + // instead, since a transport failure (triggerFatal via onDisconnect) can + // still land while those finalizers are running and must not be missed. + return { epilogue: exit.epilogue } }), ) yield* Effect.sync(() => { win32FlushInputBuffer() - if (result.reason !== undefined) - process.stderr.write((cliErrorMessage(result.reason) ?? errorFormat(result.reason)) + "\n") + if (exit.reason !== undefined) { + process.stderr.write((cliErrorMessage(exit.reason) ?? errorFormat(exit.reason)) + "\n") + // cliErrorMessage may already set process.exitCode for a tagged CliError; don't override it. + if (process.exitCode === undefined) process.exitCode = 1 + } if (result.epilogue) process.stdout.write(result.epilogue + "\n") }) }) diff --git a/packages/tui/test/app-lifecycle.test.tsx b/packages/tui/test/app-lifecycle.test.tsx index 570663424718..3d21054483e3 100644 --- a/packages/tui/test/app-lifecycle.test.tsx +++ b/packages/tui/test/app-lifecycle.test.tsx @@ -126,3 +126,61 @@ test("app.exit prints the session epilogue after scoped cleanup", async () => { mock.restore() } }) + +test("fatal errors raised during cleanup are reported after finalizers", async () => { + const setup = await createTestRenderer({ width: 80, height: 24, useThread: false }) + const core = await import("@opentui/core") + mock.module("@opentui/core", () => ({ ...core, createCliRenderer: async () => setup.renderer })) + const events = createEventSource() + const calls = createFetch() + const originalWrite = process.stderr.write.bind(process.stderr) + const originalExitCode = process.exitCode + let stderr = "" + let triggerFatal!: (error: unknown) => void + let ready!: () => void + const mounted = new Promise((resolve) => { + ready = resolve + }) + + process.stderr.write = ((chunk: string | Uint8Array) => { + stderr += String(chunk) + return true + }) as typeof process.stderr.write + process.exitCode = undefined + + try { + const { run } = await import("../src/app") + const task = Effect.runPromise( + run({ + url: "http://test", + directory, + config: createTuiResolvedConfig({ plugin_enabled: {} }), + fetch: calls.fetch, + events: events.source, + args: {}, + onReady(controls) { + triggerFatal = controls.triggerFatal + ready() + }, + pluginHost: { + async start() {}, + async dispose() { + triggerFatal(new Error("worker disconnected during cleanup")) + }, + }, + }).pipe(Effect.provide(AppNodeBuilder.build(Global.node))), + ) + + await mounted + process.emit("SIGHUP") + await task + + expect(stderr).toContain("worker disconnected during cleanup") + expect(Number(process.exitCode)).toBe(1) + } finally { + process.stderr.write = originalWrite + process.exitCode = originalExitCode + if (!setup.renderer.isDestroyed) setup.renderer.destroy() + mock.restore() + } +})