Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
136 changes: 76 additions & 60 deletions packages/opencode/src/cli/cmd/tui.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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()
}
Expand All @@ -303,7 +319,7 @@ export const TuiThreadCommand = cmd({
unguard?.()
} catch {}
}
process.exit(0)
process.exit(process.exitCode ?? 0)
},
})
// scratch
97 changes: 87 additions & 10 deletions packages/opencode/src/util/rpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }))
})
}
}
}
Expand All @@ -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<T extends Definition>(target: {
postMessage: (data: string) => void | null
onmessage: ((this: Worker, ev: MessageEvent<any>) => any) | null
addEventListener?: (type: "error" | "close" | "messageerror", listener: (event: WorkerLifecycleEvent) => void) => void
}) {
const pending = new Map<number, (result: any) => void>()
const pending = new Map<number, { resolve: (result: any) => void; reject: (error: unknown) => void }>()
const listeners = new Map<string, Set<(data: any) => 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 extends keyof T>(method: Method, input: Parameters<T[Method]>[0]): Promise<ReturnType<T[Method]>> {
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<Data>(event: string, handler: (data: Data) => void) {
Expand All @@ -60,6 +124,19 @@ export function client<T extends Definition>(target: {
handlers!.delete(handler)
}
},
onDisconnect(handler: (error: Error) => void) {
if (dead) {
handler(dead)
return () => {}
}
disconnectHandlers.add(handler)
return () => {
disconnectHandlers.delete(handler)
}
},
expectDisconnect() {
disconnectExpected = true
},
}
}

Expand Down
1 change: 1 addition & 0 deletions packages/opencode/test/fixture/crashing-worker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
throw new Error("Worker crashed on purpose")
Loading
Loading