opencode/packages/core/test/session-shell.test.ts
Kit Langton acaca2cc01 fix(util): separate process exit from capture completion
Report exit and running state from the child exit signal, independently of buffered output. On scope release, discard abandoned capture and retain the existing pipe-close or capture-deadline wait before applying process-group cleanup policy.

Cover unread output after confirmed process exit, successful descendant survival, and the approved policy that a parent exiting successfully before its invocation timeout retains success while the bounded capture grace finishes.
2026-08-31 15:12:01 -04:00

416 lines
18 KiB
TypeScript

import { describe, expect } from "bun:test"
import fs from "fs/promises"
import path from "path"
import { Cause, Context, Deferred, Effect, Exit, Fiber, Layer, Option, Schedule, Stream } from "effect"
import { TestClock } from "effect/testing"
import { Bus } from "@opencode-ai/core/bus"
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
import { Location } from "@opencode-ai/core/location"
import { LocationServiceMap } from "@opencode-ai/core/location-service-map"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { Session } from "@opencode-ai/core/session"
import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionExecution } from "@opencode-ai/core/session/execution"
import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator"
import { Shell } from "@opencode-ai/core/shell"
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
import { location } from "./fixture/location"
import { tmpdirScoped } from "./fixture/tmpdir"
import { testEffect } from "./lib/effect"
class ExecutionControl extends Context.Service<
ExecutionControl,
{
started: Deferred.Deferred<void>
release: Deferred.Deferred<void>
wakes: Session.ID[]
}
>()("test/SessionShell/ExecutionControl") {}
const controlLayer = Layer.effect(
ExecutionControl,
Effect.gen(function* () {
return { started: yield* Deferred.make<void>(), release: yield* Deferred.make<void>(), wakes: [] }
}),
)
const executionLayer = Layer.effect(
SessionExecution.Service,
Effect.gen(function* () {
const control = yield* ExecutionControl
const coordinator = yield* SessionRunCoordinator.make<Session.ID, never>({
drain: () => Deferred.succeed(control.started, undefined).pipe(Effect.andThen(Deferred.await(control.release))),
})
return SessionExecution.Service.of({
active: coordinator.active,
isActive: coordinator.isActive,
resume: coordinator.run,
interrupt: (sessionID) => coordinator.interrupt(sessionID),
awaitIdle: coordinator.awaitIdle,
wake: (sessionID) =>
Effect.sync(() => control.wakes.push(sessionID)).pipe(Effect.andThen(coordinator.wake(sessionID))),
})
}),
)
const it = testEffect(
AppNodeBuilder.build(LayerNode.group([Bus.node, Session.node, SessionExecution.node, LocationServiceMap.node]), [
Bus.node.replace(Bus.configured({ persist: true })),
SessionExecution.node.replace(executionLayer.pipe(Layer.provide(controlLayer))),
]).pipe(Layer.provideMerge(controlLayer)),
)
const setup = Effect.gen(function* () {
const tmp = yield* tmpdirScoped()
const session = yield* Session.Service
const created = yield* session.create({ location: Location.Ref.make({ directory: AbsolutePath.make(tmp.path) }) })
const locations = yield* LocationServiceMap.Service
const shell = yield* Shell.Service.pipe(Effect.provide(locations.get(created.location)))
const control = yield* ExecutionControl
const execution = yield* SessionExecution.Service
return { tmp, session, created, shell, control, execution }
})
const log = (session: Session.Interface, sessionID: Session.ID, follow = false) =>
session.log({ sessionID, follow }).pipe(Stream.filter((event) => !Bus.isSynced(event)))
// File gates prove the child process is running without relying on elapsed-time assertions.
const launch = Effect.fn(function* (fixture: Effect.Success<typeof setup>, name: string) {
const release = Effect.promise(() => Bun.write(path.join(fixture.tmp.path, `${name}.release`), "")).pipe(
Effect.asVoid,
)
yield* Effect.addFinalizer(() => release)
const command =
process.platform === "win32"
? `Write-Output '${name} started'; New-Item '${name}.started' -ItemType File | Out-Null; while (!(Test-Path '${name}.release') -and (Test-Path $PWD)) { Start-Sleep -Milliseconds 10 }; Write-Output '${name} finished'`
: `printf '${name} started\\n'; touch ${name}.started; while [ ! -f ${name}.release ] && [ -d "$PWD" ]; do sleep 0.01; done; printf '${name} finished\\n'`
const caller = yield* fixture.session
.shell({ sessionID: fixture.created.id, command })
.pipe(Effect.scoped, Effect.forkScoped)
yield* Effect.promise(() => Bun.file(path.join(fixture.tmp.path, `${name}.started`)).exists()).pipe(
Effect.repeat({ until: (exists) => exists, schedule: Schedule.spaced("10 millis") }),
Effect.timeout("5 seconds"),
)
const event = yield* log(fixture.session, fixture.created.id, true).pipe(
Stream.filter((event): event is SessionEvent.Shell.Started => event.type === "session.shell.started"),
Stream.filter((event) => event.data.shell.command === command),
Stream.runHead,
Effect.map(Option.getOrThrow),
Effect.timeout("5 seconds"),
)
return { caller, shellID: event.data.shell.id, release }
})
describe("Session.shell", () => {
it.live("runs shells concurrently with an active model and waits for each shell's own completion", () =>
Effect.gen(function* () {
const fixture = yield* setup
const model = yield* fixture.execution.resume(fixture.created.id).pipe(Effect.forkScoped)
yield* Deferred.await(fixture.control.started).pipe(Effect.timeout("5 seconds"))
const first = yield* launch(fixture, "first")
const second = yield* launch(fixture, "second")
expect(yield* fixture.execution.active).toContain(fixture.created.id)
expect(first.caller.pollUnsafe()).toBeUndefined()
expect(second.caller.pollUnsafe()).toBeUndefined()
expect(yield* fixture.session.messages({ sessionID: fixture.created.id, order: "asc" })).toMatchObject([
{ type: "shell", shellID: first.shellID, status: "running", metadata: { background: true } },
{ type: "shell", shellID: second.shellID, status: "running", metadata: { background: true } },
])
yield* second.release
yield* Fiber.join(second.caller).pipe(Effect.timeout("5 seconds"))
expect(first.caller.pollUnsafe()).toBeUndefined()
expect(yield* fixture.session.inbox(fixture.created.id)).toMatchObject([
{
type: "synthetic",
payload: {
metadata: { source: "shell", shellID: second.shellID, state: "completed" },
text: expect.stringContaining("second finished"),
},
},
])
yield* first.release
yield* Fiber.join(first.caller).pipe(Effect.timeout("5 seconds"))
expect(
(yield* fixture.session.inbox(fixture.created.id))
.filter((item) => item.type === "synthetic")
.map((item) => item.payload.metadata?.shellID),
).toEqual([second.shellID, first.shellID])
expect(yield* fixture.execution.active).toContain(fixture.created.id)
expect(fixture.control.wakes).toEqual([])
yield* Deferred.succeed(fixture.control.release, undefined)
yield* Fiber.join(model).pipe(Effect.timeout("5 seconds"))
}),
)
it.live("routes shell completion to the moved Session while retaining output from its original Location", () =>
Effect.gen(function* () {
const fixture = yield* setup
const destination = Location.Ref.make({ directory: AbsolutePath.make((yield* tmpdirScoped()).path) })
const command = yield* launch(fixture, "moved")
const bus = yield* Bus.Service
const Done = Bus.ephemeral({ type: "test.shell.move.done", schema: {} })
const readers = yield* Effect.forEach([fixture.created.location, destination], (ref) =>
bus.subscribe([SessionEvent.Shell.Ended, SessionEvent.InboxEnqueued, Done]).pipe(
Stream.takeUntil((event) => event.type === Done.type),
Stream.runCollect,
Effect.provideService(Location.Service, location(ref)),
Effect.forkScoped({ startImmediately: true }),
),
)
// This fixture's runner does not drain move requests; project the actual move event instead.
yield* bus.publish(SessionEvent.Moved, {
sessionID: fixture.created.id,
location: destination,
projectID: fixture.created.projectID,
})
expect((yield* fixture.session.get(fixture.created.id)).location).toEqual(destination)
yield* command.release
yield* Fiber.join(command.caller).pipe(Effect.timeout("5 seconds"))
yield* bus.publish(Done, {})
const events = yield* Effect.forEach(readers, (reader) => Fiber.join(reader).pipe(Effect.timeout("5 seconds")))
expect(events.map((items) => items.map((event) => event.type))).toEqual([
[Done.type],
["session.shell.ended", "session.inbox.enqueued", Done.type],
])
expect(events[1][0]).toMatchObject({
data: { output: { output: expect.stringContaining("moved finished") } },
})
const inbox = yield* fixture.session.inbox(fixture.created.id)
expect(inbox).toHaveLength(1)
expect(inbox[0]).toMatchObject({
type: "synthetic",
payload: {
text: expect.stringContaining("moved finished"),
metadata: { shellID: command.shellID, state: "completed" },
},
})
}),
)
it.live("does not suppress a normal prompt wake while a shell is running or wake again when it completes", () =>
Effect.gen(function* () {
const fixture = yield* setup
const command = yield* launch(fixture, "prompt")
const prompt = yield* fixture.session.prompt({ sessionID: fixture.created.id, text: "Continue while this runs" })
yield* Deferred.await(fixture.control.started).pipe(Effect.timeout("5 seconds"))
expect(fixture.control.wakes).toEqual([fixture.created.id])
expect(command.caller.pollUnsafe()).toBeUndefined()
yield* Deferred.succeed(fixture.control.release, undefined)
yield* fixture.execution.awaitIdle(fixture.created.id)
yield* command.release
yield* Fiber.join(command.caller).pipe(Effect.timeout("5 seconds"))
expect(fixture.control.wakes).toEqual([fixture.created.id])
expect(yield* fixture.execution.active).toEqual(new Set())
expect(yield* fixture.session.inbox(fixture.created.id)).toMatchObject([
{ id: prompt.id, type: "user" },
{ type: "synthetic", payload: { metadata: { source: "shell" } } },
])
}),
)
for (const exit of [0, 7]) {
it.live(`records output and admits one completion without waking the model after exit ${exit}`, () =>
Effect.gen(function* () {
const fixture = yield* setup
const command =
process.platform === "win32"
? `Write-Output 'user output'; [Console]::Error.WriteLine('user error'); exit ${exit}`
: `printf 'user output\\n'; printf 'user error\\n' >&2; exit ${exit}`
yield* fixture.session.shell({ sessionID: fixture.created.id, command })
const events = yield* log(fixture.session, fixture.created.id).pipe(Stream.runCollect)
expect(events.map((event) => event.type)).toEqual([
"session.created",
"session.shell.started",
"session.shell.ended",
"session.inbox.enqueued",
])
expect(events[1]).toMatchObject({ data: { shell: { metadata: { background: true } } } })
const messages = yield* fixture.session.messages({ sessionID: fixture.created.id })
expect(messages).toHaveLength(1)
expect(messages[0]).toMatchObject({
type: "shell",
command,
status: "exited",
exit,
metadata: { background: true },
})
const message = messages[0]
if (message?.type !== "shell") return yield* Effect.die("Missing shell projection")
expect(message.time.completed).toBeDefined()
expect(message.output?.output).toContain("user output")
expect(message.output?.output).toContain("user error")
expect(message.output?.truncated).toBe(false)
const inbox = yield* fixture.session.inbox(fixture.created.id)
expect(inbox).toHaveLength(1)
expect(inbox[0]).toMatchObject({
type: "synthetic",
delivery: "steer",
})
const completion = inbox[0]
if (completion?.type !== "synthetic") return yield* Effect.die("Missing shell completion")
expect(completion.payload.metadata).toEqual({
source: "shell",
shellID: message.shellID,
state: "completed",
exit,
truncated: false,
})
expect(completion.payload.description).toBeUndefined()
expect(completion.payload.text).toContain(command)
expect(completion.payload.text).toContain("user output")
expect(completion.payload.text).toContain("user error")
expect(completion.payload.text).toContain(`exited with code ${exit}`)
expect(fixture.control.wakes).toEqual([])
}),
)
}
it.effect("keeps success when the invocation timeout expires during post-exit capture", () =>
Effect.gen(function* () {
const fixture = yield* setup
const pidFile = path.join(fixture.tmp.path, "child.pid")
yield* Effect.addFinalizer(() =>
Effect.tryPromise(async () => process.kill(Number(await fs.readFile(pidFile, "utf8")), "SIGKILL")).pipe(
Effect.ignore,
),
)
const info = yield* fixture.shell.create({
command: `node "${path.join(import.meta.dir, "fixture/held-stdio.cjs")}" exit "${pidFile}"`,
timeout: 500,
})
const completion = yield* fixture.shell.wait(info.id).pipe(Effect.forkScoped)
// Wait for the real process without advancing its invocation timeout or capture deadline.
yield* fixture.shell
.get(info.id)
.pipe(
Effect.repeat({ until: (info) => info.status === "exited", schedule: Schedule.spaced("10 millis") }),
Effect.timeout("3 seconds"),
TestClock.withLive,
)
yield* TestClock.adjust("500 millis")
expect(yield* fixture.shell.get(info.id)).toMatchObject({ status: "exited", exit: 0 })
expect(completion.pollUnsafe()).toBeUndefined()
yield* TestClock.adjust("500 millis")
expect(yield* Fiber.join(completion).pipe(Effect.timeout("3 seconds"), TestClock.withLive)).toMatchObject({
status: "exited",
exit: 0,
})
const result = yield* fixture.shell.result(info)
expect(result.capture?.output).toContain("foreground-out")
expect(result.capture?.output).toContain("foreground-err")
const pid = Number(yield* Effect.promise(() => fs.readFile(pidFile, "utf8")))
expect(() => process.kill(pid, 0)).not.toThrow()
}),
)
for (const outcome of [
{
status: "killed",
state: "cancelled",
text: "Command cancelled",
output: "Shell command output is no longer available.",
},
{ status: "timeout", state: "completed", text: "Command timed out", output: "timeout started" },
]) {
it.live(`records a ${outcome.status} shell and admits its completion without waking the model`, () =>
Effect.gen(function* () {
const fixture = yield* setup
const command = yield* launch(fixture, outcome.status)
yield* outcome.status === "killed"
? fixture.shell.remove(command.shellID)
: fixture.shell.timeout(command.shellID, 1)
yield* Fiber.join(command.caller).pipe(Effect.timeout("5 seconds"))
expect(yield* fixture.session.messages({ sessionID: fixture.created.id })).toMatchObject([
{
type: "shell",
shellID: command.shellID,
status: outcome.status,
time: { completed: expect.anything() },
output: { output: expect.stringContaining(outcome.output) },
},
])
const inbox = yield* fixture.session.inbox(fixture.created.id)
expect(inbox).toHaveLength(1)
expect(inbox).toMatchObject([
{
type: "synthetic",
payload: {
text: expect.stringContaining(outcome.text),
metadata: { source: "shell", shellID: command.shellID, state: outcome.state },
},
},
])
expect(inbox[0]).not.toHaveProperty("payload.metadata.exit")
expect(fixture.control.wakes).toEqual([])
}),
)
}
it.live("admits a spawn failure without waking the model before failing the caller", () =>
Effect.gen(function* () {
const fixture = yield* setup
// Keep the Location initialized, but invalidate the real child's working directory.
yield* Effect.promise(() => fs.rm(fixture.tmp.path, { recursive: true, force: true }))
const command = "echo should-not-start"
const exit = yield* fixture.session.shell({ sessionID: fixture.created.id, command }).pipe(Effect.exit)
expect(Exit.isFailure(exit)).toBe(true)
expect(yield* fixture.session.messages({ sessionID: fixture.created.id })).toEqual([])
const inbox = yield* fixture.session.inbox(fixture.created.id)
expect(inbox).toHaveLength(1)
expect(inbox[0]).toMatchObject({
type: "synthetic",
payload: { text: expect.stringContaining(command), metadata: { source: "shell", state: "error" } },
})
expect(fixture.control.wakes).toEqual([])
}),
)
it.live("keeps the shell and completion alive when the waiting caller is interrupted", () =>
Effect.gen(function* () {
const fixture = yield* setup
const command = yield* launch(fixture, "detached")
yield* Fiber.interrupt(command.caller).pipe(Effect.timeout("5 seconds"))
const exit = yield* Fiber.await(command.caller)
expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true)
expect(yield* fixture.shell.get(command.shellID)).toMatchObject({ status: "running" })
expect(yield* fixture.session.inbox(fixture.created.id)).toEqual([])
yield* command.release
yield* log(fixture.session, fixture.created.id, true).pipe(
Stream.filter((event) => event.type === "session.inbox.enqueued"),
Stream.runHead,
Effect.timeout("5 seconds"),
)
expect(yield* fixture.session.messages({ sessionID: fixture.created.id })).toMatchObject([
{
type: "shell",
shellID: command.shellID,
status: "exited",
exit: 0,
output: { output: expect.stringContaining("detached finished") },
},
])
expect(yield* fixture.session.inbox(fixture.created.id)).toMatchObject([
{
type: "synthetic",
payload: {
text: expect.stringContaining("detached finished"),
metadata: { shellID: command.shellID, state: "completed" },
},
},
])
expect(fixture.control.wakes).toEqual([])
}),
)
})