Apply PR #24553: fix(session): harden shell cancellation

This commit is contained in:
opencode-agent[bot] 2026-04-27 21:01:39 +00:00
commit 591da2faba
5 changed files with 190 additions and 151 deletions

View file

@ -1,10 +1,10 @@
import { Cause, Deferred, Effect, Exit, Fiber, Schema, Scope, SynchronizedRef } from "effect" import { Cause, Deferred, Effect, Exit, Fiber, Latch, Schema, Scope, SynchronizedRef } from "effect"
export interface Runner<A, E = never> { export interface Runner<A, E = never> {
readonly state: State<A, E> readonly state: State<A, E>
readonly busy: boolean readonly busy: boolean
readonly ensureRunning: (work: Effect.Effect<A, E>) => Effect.Effect<A, E> readonly ensureRunning: (work: Effect.Effect<A, E>) => Effect.Effect<A, E>
readonly startShell: (work: Effect.Effect<A, E>) => Effect.Effect<A, E> readonly startShell: (work: Effect.Effect<A, E>, ready?: Latch.Latch) => Effect.Effect<A, E>
readonly cancel: Effect.Effect<void> readonly cancel: Effect.Effect<void>
} }
@ -18,6 +18,8 @@ interface RunHandle<A, E> {
interface ShellHandle<A, E> { interface ShellHandle<A, E> {
id: number id: number
cancelled: Deferred.Deferred<void>
ready?: Latch.Latch
fiber: Fiber.Fiber<A, E> fiber: Fiber.Fiber<A, E>
} }
@ -59,6 +61,9 @@ export const make = <A, E = never>(
? Deferred.fail(done, new Cancelled()).pipe(Effect.asVoid) ? Deferred.fail(done, new Cancelled()).pipe(Effect.asVoid)
: Deferred.done(done, exit).pipe(Effect.asVoid) : Deferred.done(done, exit).pipe(Effect.asVoid)
const awaitDone = (done: Deferred.Deferred<A, E | Cancelled>) =>
Deferred.await(done).pipe(Effect.catchTag("RunnerCancelled", (e) => onInterrupt ?? Effect.die(e)))
const idleIfCurrent = () => const idleIfCurrent = () =>
SynchronizedRef.modify(ref, (st) => [st._tag === "Idle" ? idle : Effect.void, st] as const).pipe(Effect.flatten) SynchronizedRef.modify(ref, (st) => [st._tag === "Idle" ? idle : Effect.void, st] as const).pipe(Effect.flatten)
@ -89,7 +94,9 @@ export const make = <A, E = never>(
SynchronizedRef.modifyEffect( SynchronizedRef.modifyEffect(
ref, ref,
Effect.fnUntraced(function* (st) { Effect.fnUntraced(function* (st) {
if (st._tag === "Shell" && st.shell.id === id) return [idle, { _tag: "Idle" }] as const if (st._tag === "Shell" && st.shell.id === id) {
return [idle, { _tag: "Idle" }] as const
}
if (st._tag === "ShellThenRun" && st.shell.id === id) { if (st._tag === "ShellThenRun" && st.shell.id === id) {
const run = yield* startRun(st.run.work, st.run.done) const run = yield* startRun(st.run.work, st.run.done)
return [Effect.void, { _tag: "Running", run }] as const return [Effect.void, { _tag: "Running", run }] as const
@ -98,7 +105,12 @@ export const make = <A, E = never>(
}), }),
).pipe(Effect.flatten) ).pipe(Effect.flatten)
const stopShell = (shell: ShellHandle<A, E>) => Fiber.interrupt(shell.fiber) const stopShell = (shell: ShellHandle<A, E>) =>
Effect.gen(function* () {
if (shell.ready) yield* shell.ready.await.pipe(Effect.exit, Effect.asVoid)
yield* Deferred.succeed(shell.cancelled, undefined).pipe(Effect.asVoid)
yield* Fiber.interrupt(shell.fiber)
})
const ensureRunning = (work: Effect.Effect<A, E>) => const ensureRunning = (work: Effect.Effect<A, E>) =>
SynchronizedRef.modifyEffect( SynchronizedRef.modifyEffect(
@ -107,30 +119,25 @@ export const make = <A, E = never>(
switch (st._tag) { switch (st._tag) {
case "Running": case "Running":
case "ShellThenRun": case "ShellThenRun":
return [Deferred.await(st.run.done), st] as const return [awaitDone(st.run.done), st] as const
case "Shell": { case "Shell": {
const run = { const run = {
id: next(), id: next(),
done: yield* Deferred.make<A, E | Cancelled>(), done: yield* Deferred.make<A, E | Cancelled>(),
work, work,
} satisfies PendingHandle<A, E> } satisfies PendingHandle<A, E>
return [Deferred.await(run.done), { _tag: "ShellThenRun", shell: st.shell, run }] as const return [awaitDone(run.done), { _tag: "ShellThenRun", shell: st.shell, run }] as const
} }
case "Idle": { case "Idle": {
const done = yield* Deferred.make<A, E | Cancelled>() const done = yield* Deferred.make<A, E | Cancelled>()
const run = yield* startRun(work, done) const run = yield* startRun(work, done)
return [Deferred.await(done), { _tag: "Running", run }] as const return [awaitDone(done), { _tag: "Running", run }] as const
} }
} }
}), }),
).pipe( ).pipe(Effect.flatten)
Effect.flatten,
Effect.catch(
(e): Effect.Effect<A, E> => (e instanceof Cancelled ? (onInterrupt ?? Effect.die(e)) : Effect.fail(e as E)),
),
)
const startShell = (work: Effect.Effect<A, E>) => const startShell = (work: Effect.Effect<A, E>, ready?: Latch.Latch) =>
SynchronizedRef.modifyEffect( SynchronizedRef.modifyEffect(
ref, ref,
Effect.fnUntraced(function* (st) { Effect.fnUntraced(function* (st) {
@ -145,13 +152,20 @@ export const make = <A, E = never>(
} }
yield* busy yield* busy
const id = next() const id = next()
const cancelled = yield* Deferred.make<void>()
const fiber = yield* work.pipe(Effect.ensuring(finishShell(id)), Effect.forkChild) const fiber = yield* work.pipe(Effect.ensuring(finishShell(id)), Effect.forkChild)
const shell = { id, fiber } satisfies ShellHandle<A, E> const shell = { id, cancelled, ready, fiber } satisfies ShellHandle<A, E>
return [ return [
Effect.gen(function* () { Effect.gen(function* () {
const exit = yield* Fiber.await(fiber) const exit = yield* Fiber.await(fiber)
if (Exit.isSuccess(exit)) return exit.value if (Exit.isSuccess(exit)) return exit.value
if (Cause.hasInterruptsOnly(exit.cause) && onInterrupt) return yield* onInterrupt if (
Cause.hasInterruptsOnly(exit.cause) ||
((yield* Deferred.isDone(cancelled)) && Cause.hasInterrupts(exit.cause) && !Cause.hasDies(exit.cause))
) {
if (onInterrupt) return yield* onInterrupt
return yield* Effect.die(new Cancelled())
}
return yield* Effect.failCause(exit.cause) return yield* Effect.failCause(exit.cause)
}), }),
{ _tag: "Shell", shell }, { _tag: "Shell", shell },
@ -183,8 +197,8 @@ export const make = <A, E = never>(
case "ShellThenRun": case "ShellThenRun":
return [ return [
Effect.gen(function* () { Effect.gen(function* () {
yield* Deferred.fail(st.run.done, new Cancelled()).pipe(Effect.asVoid)
yield* stopShell(st.shell) yield* stopShell(st.shell)
yield* Deferred.fail(st.run.done, new Cancelled()).pipe(Effect.asVoid)
yield* idleIfCurrent() yield* idleIfCurrent()
}), }),
{ _tag: "Idle" } as const, { _tag: "Idle" } as const,

View file

@ -46,7 +46,7 @@ import { AppFileSystem } from "@opencode-ai/core/filesystem"
import { Truncate } from "@/tool/truncate" import { Truncate } from "@/tool/truncate"
import { decodeDataUrl } from "@/util/data-url" import { decodeDataUrl } from "@/util/data-url"
import { Process } from "@/util/process" import { Process } from "@/util/process"
import { Cause, Effect, Exit, Layer, Option, Scope, Context, Schema } from "effect" import { Cause, Effect, Exit, Latch, Layer, Option, Scope, Context, Schema } from "effect"
import { zod } from "@/util/effect-zod" import { zod } from "@/util/effect-zod"
import { withStatics } from "@/util/schema" import { withStatics } from "@/util/schema"
import * as EffectLogger from "@opencode-ai/core/effect/logger" import * as EffectLogger from "@opencode-ai/core/effect/logger"
@ -725,9 +725,12 @@ NOTE: At any point in time through this workflow you should feel free to ask the
} satisfies MessageV2.TextPart) } satisfies MessageV2.TextPart)
}) })
const shellImpl = Effect.fn("SessionPrompt.shellImpl")(function* (input: ShellInput) { const shellImpl = Effect.fn("SessionPrompt.shellImpl")(function* (input: ShellInput, ready?: Latch.Latch) {
return yield* Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const markReady = ready ? ready.open.pipe(Effect.asVoid) : Effect.void
const { msg, part, cwd } = yield* Effect.gen(function* () {
const ctx = yield* InstanceState.context const ctx = yield* InstanceState.context
const run = yield* runner()
const session = yield* sessions.get(input.sessionID) const session = yield* sessions.get(input.sessionID)
if (session.revert) { if (session.revert) {
yield* revert.cleanup(session) yield* revert.cleanup(session)
@ -789,27 +792,15 @@ NOTE: At any point in time through this workflow you should feel free to ask the
}, },
} }
yield* sessions.updatePart(part) yield* sessions.updatePart(part)
return { msg, part, cwd: ctx.directory }
}).pipe(Effect.ensuring(markReady))
const cfg = yield* config.get() const cfg = yield* config.get()
const sh = Shell.preferred(cfg.shell) const sh = Shell.preferred(cfg.shell)
const cwd = ctx.directory
const args = Shell.args(sh, input.command, cwd) const args = Shell.args(sh, input.command, cwd)
const shellEnv = yield* plugin.trigger(
"shell.env",
{ cwd, sessionID: input.sessionID, callID: part.callID },
{ env: {} },
)
const cmd = ChildProcess.make(sh, args, {
cwd,
extendEnv: true,
env: { ...shellEnv.env, TERM: "dumb" },
stdin: "ignore",
forceKillAfter: "3 seconds",
})
let output = "" let output = ""
let aborted = false let aborted = false
const finish = Effect.uninterruptible( const finish = Effect.uninterruptible(
Effect.gen(function* () { Effect.gen(function* () {
if (aborted) { if (aborted) {
@ -833,35 +824,46 @@ NOTE: At any point in time through this workflow you should feel free to ask the
}), }),
) )
const exit = yield* Effect.gen(function* () { const exit = yield* restore(
Effect.gen(function* () {
const shellEnv = yield* plugin.trigger(
"shell.env",
{ cwd, sessionID: input.sessionID, callID: part.callID },
{ env: {} },
)
const cmd = ChildProcess.make(sh, args, {
cwd,
extendEnv: true,
env: { ...shellEnv.env, TERM: "dumb" },
stdin: "ignore",
forceKillAfter: "3 seconds",
})
const handle = yield* spawner.spawn(cmd) const handle = yield* spawner.spawn(cmd)
yield* Stream.runForEach(Stream.decodeText(handle.all), (chunk) => yield* Stream.runForEach(Stream.decodeText(handle.all), (chunk) =>
Effect.sync(() => { Effect.gen(function* () {
output += chunk output += chunk
if (part.state.status === "running") { if (part.state.status === "running") {
part.state.metadata = { output, description: "" } part.state.metadata = { output, description: "" }
void run.fork(sessions.updatePart(part)) yield* sessions.updatePart(part)
} }
}), }),
) )
yield* handle.exitCode yield* handle.exitCode
}).pipe( }).pipe(Effect.scoped, Effect.orDie),
Effect.scoped, ).pipe(Effect.exit)
Effect.onInterrupt(() =>
Effect.sync(() => {
aborted = true
}),
),
Effect.orDie,
Effect.ensuring(finish),
Effect.exit,
)
if (Exit.isFailure(exit) && !Cause.hasInterruptsOnly(exit.cause)) { if (Exit.isFailure(exit) && Cause.hasInterrupts(exit.cause) && !Cause.hasDies(exit.cause)) {
aborted = true
}
yield* finish
if (Exit.isFailure(exit) && !aborted && !Cause.hasInterruptsOnly(exit.cause)) {
return yield* Effect.failCause(exit.cause) return yield* Effect.failCause(exit.cause)
} }
return { info: msg, parts: [part] } return { info: msg, parts: [part] }
}),
)
}) })
const getModel = Effect.fn("SessionPrompt.getModel")(function* ( const getModel = Effect.fn("SessionPrompt.getModel")(function* (
@ -1512,7 +1514,8 @@ NOTE: At any point in time through this workflow you should feel free to ask the
const shell: (input: ShellInput) => Effect.Effect<MessageV2.WithParts> = Effect.fn("SessionPrompt.shell")( const shell: (input: ShellInput) => Effect.Effect<MessageV2.WithParts> = Effect.fn("SessionPrompt.shell")(
function* (input: ShellInput) { function* (input: ShellInput) {
return yield* state.startShell(input.sessionID, lastAssistant(input.sessionID), shellImpl(input)) const ready = yield* Latch.make()
return yield* state.startShell(input.sessionID, lastAssistant(input.sessionID), shellImpl(input, ready), ready)
}, },
) )

View file

@ -1,6 +1,6 @@
import { InstanceState } from "@/effect/instance-state" import { InstanceState } from "@/effect/instance-state"
import { Runner } from "@/effect/runner" import { Runner } from "@/effect/runner"
import { Effect, Layer, Scope, Context } from "effect" import { Effect, Latch, Layer, Scope, Context } from "effect"
import * as Session from "./session" import * as Session from "./session"
import { MessageV2 } from "./message-v2" import { MessageV2 } from "./message-v2"
import { SessionID } from "./schema" import { SessionID } from "./schema"
@ -18,6 +18,7 @@ export interface Interface {
sessionID: SessionID, sessionID: SessionID,
onInterrupt: Effect.Effect<MessageV2.WithParts>, onInterrupt: Effect.Effect<MessageV2.WithParts>,
work: Effect.Effect<MessageV2.WithParts>, work: Effect.Effect<MessageV2.WithParts>,
ready?: Latch.Latch,
) => Effect.Effect<MessageV2.WithParts> ) => Effect.Effect<MessageV2.WithParts>
} }
@ -98,8 +99,9 @@ export const layer = Layer.effect(
sessionID: SessionID, sessionID: SessionID,
onInterrupt: Effect.Effect<MessageV2.WithParts>, onInterrupt: Effect.Effect<MessageV2.WithParts>,
work: Effect.Effect<MessageV2.WithParts>, work: Effect.Effect<MessageV2.WithParts>,
ready?: Latch.Latch,
) { ) {
return yield* (yield* runner(sessionID, onInterrupt)).startShell(work) return yield* (yield* runner(sessionID, onInterrupt)).startShell(work, ready)
}) })
return Service.of({ assertNotBusy, cancel, ensureRunning, startShell }) return Service.of({ assertNotBusy, cancel, ensureRunning, startShell })

View file

@ -334,6 +334,22 @@ describe("Runner", () => {
}), }),
) )
it.live(
"cancel does not mask shell defects",
Effect.gen(function* () {
const s = yield* Scope.Scope
const runner = Runner.make<string>(s, { onInterrupt: Effect.succeed("interrupted") })
const sh = yield* runner
.startShell(Effect.never.pipe(Effect.ensuring(Effect.die("boom")), Effect.as("ignored")))
.pipe(Effect.forkChild)
yield* Effect.sleep("10 millis")
yield* runner.cancel
expect(Exit.isFailure(yield* Fiber.await(sh))).toBe(true)
}),
)
// --- shell→run handoff --- // --- shell→run handoff ---
it.live( it.live(

View file

@ -1472,6 +1472,10 @@ unix(
const exit = yield* Fiber.await(loop) const exit = yield* Fiber.await(loop)
expect(Exit.isSuccess(exit)).toBe(true) expect(Exit.isSuccess(exit)).toBe(true)
if (Exit.isSuccess(exit)) {
const tool = completedTool(exit.value.parts)
expect(tool?.state.output).toContain("User aborted the command")
}
yield* Fiber.await(sh) yield* Fiber.await(sh)
}), }),