diff --git a/CHANGELOG.md b/CHANGELOG.md index c057ab62..a1c3a3d9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -48,6 +48,14 @@ parallel copies under `docs/` or `scripts/notes/`. At cut time: rename owns idle rebuild; delivery generation owns session identity, so interrupt, /clear, and /new abort the outstanding overlay, skip minting a grant, and notify the operator instead of delivering into a rebuilt agent. +- TUI quit, crash, and process signals await once-only runtime shutdown so live + shell-guard children are reaped. Teardown failure after a completed session + exits 1; SIGINT, SIGTERM, and SIGHUP still exit 128+n. +- Persist close_agent surfaces leftover-child dispose failure so a worker + that survives reap is not reported as a successful shutdown. +- Leftover exec dispose is reported as a failed run (stderr + status failed), and + parent toolset dispose finishes remaining workers and posix teardown before + surfacing leftover-child failure. ### Changed diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index ce5e45b2..35460091 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -369,7 +369,7 @@ tool call - **Secret Guard** (`secret-guard-plugin.ts`) — Hard-denies path-keyed tool calls (`read_file`, `write_file`, …) that would put a sensitive file into (or write it from) the model context. Runs before the permission plugin, so the path-arg deny holds even under `--dangerously-skip-permissions`. Shell commands that _reference_ a sensitive path (tokenized so `cat .env`, `bun --env-file=.env run …`, and quote/env-assignment forms are detected) are not hard-denied here: they require operator approval via the permission gate, and auto mode forces an ask through the auto-shell policy (`sensitive-path` rule). Once the operator approves, the command runs. Shell detection is best-effort: token matching defeats quoting and env-assignment/redirection forms but not dynamic path construction (variable indirection, `printf` assembly). Tool-result secret scrub still redacts credential-shaped output. - **Authorization** (`run-shell-authz.ts`, wired by `authz-plugin.ts`) — Denies catastrophic shell command patterns by regex, and hard-blocks shell `find`, head-position `rg`, and recursive `grep -r` (they can walk huge trees and OOM the host). Bounded `grep`/`search_files` tools remain practical alternatives (timeout + output caps); the patterns match those three command shapes only — an `ls -R`, `fd`, or scripted `os.walk` is just as unbounded and is not caught, so the block message tells the model not to substitute one. The permission gate’s shell auto-allow path consults the same policy so it never pre-approves a command authz would reject. - **Permission** (`permission-plugin.ts`) — Delegates consequential calls to the permission gate. -- **Shell Guard** (`shell-guard-plugin.ts`) — Corbits Code-only replacement for stock `run_shell` (interchange stays unpatched): no built-in default timeout (optional per-call or `settings.shell.timeoutMs`; `maxTimeoutMs` clamps only a resolved timeout), 512KB display cap with head+tail retention (the process keeps running when the cap is hit), process-group kill on timeout/abort only, and `background: true` — the call returns a `shell_id` at once (registry in `src/shell/background-shell.ts`), the process group keeps running past the turn, completion is delivered on a later turn via `buildShellBackgroundMessage`, and `shell_collect` retrieves or cancels (schema advertised by `advertiseShellGuardTimeout`; evaluated by the permission chain at start time like any shell call). Also applies a 10s wall-clock budget to `grep`/`search_files`. +- **Shell Guard** (`shell-guard-plugin.ts`) — Corbits Code-only replacement for stock `run_shell` (interchange stays unpatched): no built-in default timeout (optional per-call or `settings.shell.timeoutMs`; `maxTimeoutMs` clamps only a resolved timeout), 512KB display cap with head+tail retention (the process keeps running when the cap is hit), process-group kill on timeout, abort, and plugin dispose (live children tracked in the plugin and reaped by `posixTools.dispose`), and `background: true` — the call returns a `shell_id` at once (registry in `src/shell/background-shell.ts`), the process group keeps running past the turn, completion is delivered on a later turn via `buildShellBackgroundMessage`, and `shell_collect` retrieves or cancels (schema advertised by `advertiseShellGuardTimeout`; evaluated by the permission chain at start time like any shell call). Also applies a 10s wall-clock budget to `grep`/`search_files`. Ripgrep detached spawns are not tracked. - **Read File Guard** (`read-file-guard-plugin.ts`) — Corbits Code-only short-circuit for `read_file` on real filesystem paths and configured `tool-output://` URIs (interchange stays unpatched): streaming reads that never decode the whole file in one pass, caps model-facing output at 50KB, defaults to 2000 lines, truncates long lines with recovery hints, samples the first chunk to reject binary, and stops at an 8MB scan ceiling. Emits `offset` continuation notices so the model can page without losing file or spill content on disk. - **Verify** (`verify-plugin.ts`) — Re-reads after `write_file` / `edit_file` and errors on mismatch. Per-path serialization (`file-mutation-lock.ts`) prevents parallel edits on one file from tripping verification. - **Edit file line range** (`edit-file-line-range-plugin.ts`) — Corbits Code-only short-circuit for `edit_file` mode B (`start_line`/`end_line`/`new_string`), same pattern as shell-guard; schema advertised via `advertiseEditFileLineRange`. Modes are mutually exclusive: a call supplying both `old_string` and `start_line`/`end_line` is rejected with a recoverable error naming which fields to omit (no file-content disambiguation). diff --git a/src/agent/fleet-verbs-mount.test.ts b/src/agent/fleet-verbs-mount.test.ts index b822cf01..11ce7112 100644 --- a/src/agent/fleet-verbs-mount.test.ts +++ b/src/agent/fleet-verbs-mount.test.ts @@ -53,6 +53,164 @@ describe("primary fleet verb mount", () => { await toolset.dispose(); }); + test("createAgentToolset dispose rejects when a fleet closeOne throws leftover children", async () => { + const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-")); + const { createAgentToolset } = await import("./tools.js"); + const permissionGate = { + check: async () => ({ allowed: true }), + getSkipPermissions: () => false, + } as never; + const sessions = createSubAgentSessionStore(); + const worker = sessions.start({ description: "d", agentId: "a", brief: "b" }); + sessions.markRunning(worker.id); + sessions.registerClose(worker.id, async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }); + + const toolset = await createAgentToolset({ + cwd, + permissionGate, + onOperatorGate: async () => ({ kind: "option", index: 0 }), + subAgent: { + provider: { + providerName: "test", + baseURL: "http://127.0.0.1:0", + model: "test-model", + }, + getWorkdirBase: () => cwd, + sessions, + }, + }); + + await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/); + }); + + test("createAgentToolset dispose rejects when a retained completed persist worker leaves children", async () => { + const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-")); + const { createAgentToolset } = await import("./tools.js"); + const permissionGate = { + check: async () => ({ allowed: true }), + getSkipPermissions: () => false, + } as never; + const sessions = createSubAgentSessionStore(); + const worker = sessions.start({ + description: "d", + agentId: "a", + brief: "b", + retained: true, + }); + sessions.registerClose(worker.id, async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }); + sessions.complete(worker.id, "done", { agentRetained: true }); + + const toolset = await createAgentToolset({ + cwd, + permissionGate, + onOperatorGate: async () => ({ kind: "option", index: 0 }), + subAgent: { + provider: { + providerName: "test", + baseURL: "http://127.0.0.1:0", + model: "test-model", + }, + getWorkdirBase: () => cwd, + sessions, + }, + }); + + await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/); + }); + + test("createAgentToolset dispose closes remaining retained workers after the first leftover", async () => { + const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-")); + const { createAgentToolset } = await import("./tools.js"); + const permissionGate = { + check: async () => ({ allowed: true }), + getSkipPermissions: () => false, + } as never; + const sessions = createSubAgentSessionStore(); + const first = sessions.start({ + description: "d1", + agentId: "a", + brief: "b", + retained: true, + }); + const second = sessions.start({ + description: "d2", + agentId: "a", + brief: "b", + retained: true, + }); + let firstCloseCalls = 0; + let secondCloseCalls = 0; + sessions.registerClose(first.id, async () => { + firstCloseCalls += 1; + throw new Error("1 shell child process still live after 2000ms reap"); + }); + sessions.registerClose(second.id, async () => { + secondCloseCalls += 1; + }); + sessions.complete(first.id, "done", { agentRetained: true }); + sessions.complete(second.id, "done", { agentRetained: true }); + + const toolset = await createAgentToolset({ + cwd, + permissionGate, + onOperatorGate: async () => ({ kind: "option", index: 0 }), + subAgent: { + provider: { + providerName: "test", + baseURL: "http://127.0.0.1:0", + model: "test-model", + }, + getWorkdirBase: () => cwd, + sessions, + }, + }); + + await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/); + expect(firstCloseCalls).toBe(1); + expect(secondCloseCalls).toBe(1); + }); + + test("createAgentToolset dispose rejects when a retained running persist worker leaves children", async () => { + const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-")); + const { createAgentToolset } = await import("./tools.js"); + const permissionGate = { + check: async () => ({ allowed: true }), + getSkipPermissions: () => false, + } as never; + const sessions = createSubAgentSessionStore(); + const worker = sessions.start({ + description: "d", + agentId: "a", + brief: "b", + retained: true, + }); + sessions.markRunning(worker.id); + sessions.registerClose(worker.id, async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }); + + const toolset = await createAgentToolset({ + cwd, + permissionGate, + onOperatorGate: async () => ({ kind: "option", index: 0 }), + subAgent: { + provider: { + providerName: "test", + baseURL: "http://127.0.0.1:0", + model: "test-model", + }, + getWorkdirBase: () => cwd, + sessions, + }, + }); + + await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/); + }); + test("createAgentToolset omits fleet verbs when subAgent is not set", async () => { const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-")); const { createAgentToolset } = await import("./tools.js"); diff --git a/src/agent/tools.ts b/src/agent/tools.ts index de2c0787..3e1ccb5b 100644 --- a/src/agent/tools.ts +++ b/src/agent/tools.ts @@ -105,6 +105,13 @@ export const ASK_OPERATOR_OPTION_MAX_CHARS = 48; /** Cap on the ask_operator question (UTF-16 code units). */ export const ASK_OPERATOR_QUESTION_MAX_CHARS = 160; +function rethrowToolsetDisposeFailures(failures: unknown[]): void { + const first = failures[0]; + if (first === undefined) return; + if (failures.length === 1) throw first; + throw new AggregateError(failures, "toolset leftover dispose failed"); +} + const SubmitOutputArgs = type({ "summary?": "string", "step?": "string", @@ -982,11 +989,28 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise { + const failures: unknown[] = []; + // Kill every live background process group before the posix teardown so + // /clear, interrupt, and reload cannot leave orphans behind. + backgroundShells.disposeAll("session closed"); + try { + await posixTools.dispose(); + } catch (err: unknown) { + failures.push(err); + } const fleetSessions = fleetSessionsForDispose; if (fleetSessions !== undefined) { - fleetSessions.cancelAll("parent session closed"); + try { + await fleetSessions.cancelAll("parent session closed"); + } catch (err: unknown) { + failures.push(err); + } for (const session of [...fleetSessions.list()].reverse()) { - await fleetSessions.closeOne(session.id, DEFAULT_CLOSE_DEADLINE_MS); + try { + await fleetSessions.closeOne(session.id, DEFAULT_CLOSE_DEADLINE_MS); + } catch (err: unknown) { + failures.push(err); + } } } await Promise.allSettled([...inFlightConnections.values()]); @@ -997,11 +1021,8 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise client.close().catch(() => undefined)), ); connectedClients.clear(); - // Kill every live background process group before the posix teardown so - // /clear, interrupt, and reload cannot leave orphans behind. - backgroundShells.disposeAll("session closed"); - await posixTools.dispose(); await disposeWebSearchClients(); + rethrowToolsetDisposeFailures(failures); })(); return disposal; }; diff --git a/src/exec/runner.ts b/src/exec/runner.ts index 2ec6d9a5..5a8f55ae 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -22,6 +22,7 @@ import { type SubAgentSessionStore, } from "../subagent/index.js"; import { getProcessAdmissionQueue } from "../subagent/admission.js"; +import { awaitCloseWithoutHidingLeftover } from "../subagent/dispose.js"; import type { ContextStore, InferenceSource, InboundMessage } from "@intx/types/runtime"; import { OPERATOR_ORIGINATED_FLAG } from "../agent/message-provenance.js"; import { loadAgentProfiles } from "../agent/profiles.js"; @@ -44,6 +45,7 @@ import { sessionDir, } from "../session/index.js"; import { setActiveRun } from "../session/active-run.js"; +import { setActiveDisposeHost, clearActiveDisposeHost } from "../session/active-host.js"; import { finalizeRunState, saveState, type ConnectedMcpServer } from "../session/state.js"; import { resolveExecRunStatus, type RunSink } from "../session/run-sink.js"; import { createRunSummary } from "../session/hooks.js"; @@ -126,30 +128,69 @@ export function execUserFailureMessage( } /** - * Headless analogue of TUI `runtime-shutdown`: abort live workers, then close - * the primary agent and dispose the toolset. `cancelAll` is fire-and-forget — - * it does not serialize `closeOne`. + * Headless analogue of TUI `runtime-shutdown`: dispose the toolset (posix + * process-group reap) before waiting on agent.close so a hung close cannot + * skip killing detached run_shell children. `cancelAll` is awaited so a + * leftover-child throw is visible. Once-only per runtime object so the send + * path, `finally`, and signal host cannot double-dispose. */ -export async function disposeExecRuntime(args: { +const execDisposeInFlight = new WeakMap>(); + +function rethrowExecDisposeFailures(failures: unknown[]): void { + const first = failures[0]; + if (first === undefined) return; + if (failures.length === 1) throw first; + throw new AggregateError(failures, "exec runtime dispose failed"); +} + +export function disposeExecRuntime(args: { agent: { close: () => Promise } | null; toolset: { dispose: () => Promise } | null; subAgentSessions: Pick | null; }): Promise { - args.subAgentSessions?.cancelAll("Session closed"); - if (args.agent !== null) { - await args.agent.close().catch((err: unknown) => { - logger.debug("agent.close during exec finally failed: {error}", { - error: formatCaughtError(err), - }); - }); + const key = args.toolset ?? args.agent ?? args.subAgentSessions; + if (key !== null) { + const existing = execDisposeInFlight.get(key); + if (existing !== undefined) return existing; } + + const run = runExecDispose(args); + if (key !== null) execDisposeInFlight.set(key, run); + return run; +} + +async function runExecDispose(args: { + agent: { close: () => Promise } | null; + toolset: { dispose: () => Promise } | null; + subAgentSessions: Pick | null; +}): Promise { + const failures: unknown[] = []; if (args.toolset !== null) { - await args.toolset.dispose().catch((err: unknown) => { + try { + await args.toolset.dispose(); + } catch (err: unknown) { logger.debug("toolset.dispose during exec finally failed: {error}", { error: formatCaughtError(err), }); - }); + failures.push(err); + } } + try { + await args.subAgentSessions?.cancelAll("Session closed"); + } catch (err) { + failures.push(err); + } + if (args.agent !== null) { + try { + await awaitCloseWithoutHidingLeftover(args.agent.close(), failures[0]); + } catch (err: unknown) { + logger.debug("agent.close during exec finally failed: {error}", { + error: formatCaughtError(err), + }); + failures.push(err); + } + } + rethrowExecDisposeFailures(failures); } /** @@ -288,6 +329,7 @@ export async function runExec(config: Config): Promise { let runSink: RunSink | null = null; let providerFailureObserved = false; let providerError: InferenceErrorLike | undefined; + let result: ExecResult | undefined; const persist = async ( status: "running" | "done" | "failed" | "cancelled", @@ -490,6 +532,7 @@ export async function runExec(config: Config): Promise { ...(extraToolPlugins.length > 0 ? { extraToolPlugins } : {}), }); toolset = agentToolset; + setActiveDisposeHost(() => disposeExecRuntime({ agent, toolset, subAgentSessions })); const systemPrompt = overlay.systemPrompt ?? @@ -797,7 +840,7 @@ export async function runExec(config: Config): Promise { stderr.write(`Error: ${userMessage}\n`); const persistStatus = summaryStatus === "cancelled" ? "cancelled" : "failed"; await persist(persistStatus, { error: diagnosticMessage }); - return { + result = { exitCode: 1, sessionId, text: textOut, @@ -810,10 +853,11 @@ export async function runExec(config: Config): Promise { provider: config.providerName, model: config.model, }; + return result; } await persist("done"); - return { + result = { exitCode: 0, sessionId, text: textOut, @@ -825,13 +869,14 @@ export async function runExec(config: Config): Promise { provider: config.providerName, model: config.model, }; + return result; } catch (err) { const diagnosticMessage = formatCaughtError(err); logger.error("exec failed: {error}", { error: diagnosticMessage }); const userMessage = execUserFailureMessage(config, err, providerFailureObserved, providerError); stderr.write(`Error: ${userMessage}\n`); await persist("failed", { error: diagnosticMessage }); - return { + result = { exitCode: 1, sessionId, text: textOut, @@ -844,8 +889,21 @@ export async function runExec(config: Config): Promise { provider: config.providerName, model: config.model, }; + return result; } finally { - await disposeExecRuntime({ agent, toolset, subAgentSessions }); + try { + await disposeExecRuntime({ agent, toolset, subAgentSessions }); + } catch (err: unknown) { + const message = formatCaughtError(err); + logger.error("runtime dispose failed: {error}", { error: message }); + stderr.write(`Error: runtime dispose failed: ${message}\n`); + if (result !== undefined) { + result.exitCode = 1; + result.status = "failed"; + result.error = `runtime dispose failed: ${message}`; + } + } + clearActiveDisposeHost(); } } diff --git a/src/index.ts b/src/index.ts index 310f9e21..a609a636 100644 --- a/src/index.ts +++ b/src/index.ts @@ -118,6 +118,31 @@ export async function main(argv: readonly string[]): Promise { // the process down) can't re-enter either path a second time. let terminating = false; +export const RUNTIME_TEARDOWN_DEADLINE_MS = 2_000; + +async function awaitActiveDisposeHost(context: string): Promise { + const dispose = getActiveDisposeHost(); + if (dispose === null) return; + let timer: ReturnType | undefined; + try { + await Promise.race([ + Promise.resolve(dispose()), + new Promise((_, reject) => { + timer = setTimeout(() => { + reject(new Error(`runtime teardown exceeded ${RUNTIME_TEARDOWN_DEADLINE_MS}ms`)); + }, RUNTIME_TEARDOWN_DEADLINE_MS); + if (typeof timer.unref === "function") timer.unref(); + }), + ]); + } catch (disposeErr: unknown) { + process.stderr.write( + `host dispose failed ${context}: ${disposeErr instanceof Error ? disposeErr.message : String(disposeErr)}\n`, + ); + } finally { + if (timer !== undefined) clearTimeout(timer); + } +} + // Exported so an integration test can register these process-level handlers // and inject a crash without spawning the full TUI stack. export async function handleFatal(kind: CrashKind, error: unknown): Promise { @@ -129,20 +154,11 @@ export async function handleFatal(kind: CrashKind, error: unknown): Promise { if (terminating) return; terminating = true; - try { - getActiveDisposeHost()?.(); - } catch (disposeErr: unknown) { - process.stderr.write( - `host dispose failed handling ${signal}: ${disposeErr instanceof Error ? disposeErr.message : String(disposeErr)}\n`, - ); - } - // Same fence as handleFatal: any snapshot still queued in writeChains must - // see isCrashed and step aside before saveCrashState renames run.json. - markCrashed(); - void finalizeActiveRunOnSignal(signal).finally(() => { + void (async () => { + // Same fence as handleFatal: start teardown, then flip isCrashed + // before awaiting so queued snapshot writes cannot clobber the + // terminal write. See markCrashed's doc comment for the residual + // window this cannot close. + const teardown = awaitActiveDisposeHost(`handling ${signal}`); + markCrashed(); + await teardown; + await finalizeActiveRunOnSignal(signal); process.exit(128 + SIGNAL_EXIT_NUMBER[signal]); - }); + })(); }); } } diff --git a/src/plugins/shell-guard-plugin.test.ts b/src/plugins/shell-guard-plugin.test.ts index 34814519..c41b2c8e 100644 --- a/src/plugins/shell-guard-plugin.test.ts +++ b/src/plugins/shell-guard-plugin.test.ts @@ -4,8 +4,8 @@ import { join } from "node:path"; import { tmpdir } from "node:os"; import { realpathSync } from "node:fs"; import type { ToolCall, ToolResult } from "@intx/types/runtime"; - -import { spawnSync } from "node:child_process"; +import { EventEmitter } from "node:events"; +import { spawnSync, type ChildProcess } from "node:child_process"; import { randomUUID } from "node:crypto"; import { createBackgroundShellRegistry } from "../shell/background-shell.js"; @@ -14,6 +14,7 @@ import { MAX_SHELL_OUTPUT_BYTES, advertiseShellGuardTimeout, resolveShellTimeoutMs, + reapLiveChildren, runGuardedShell, shellGuardPlugin, } from "./shell-guard-plugin.js"; @@ -747,4 +748,148 @@ describe("shellGuardPlugin", () => { expect(result.isError).toBe(true); expect(result.content).toMatch(/aborted/); }); + + test("dispose without abort kills tagged grandchildren and is idempotent", async () => { + if (process.platform === "win32") return; + const plugin = shellGuardPlugin(process.cwd()); + const handler = plugin.middleware!(fallback); + const token = `ic_guard_dispose_${randomUUID()}`; + const cmd = `bash -c 'IC_GUARD_TAG=${token} sleep 600 & IC_GUARD_TAG=${token} exec sleep 600'`; + const running = handler( + { id: "dispose-live", name: "run_shell", arguments: { command: cmd } }, + neverAbort(), + ); + try { + const started = Date.now(); + while (Date.now() - started < 5_000) { + const probe = spawnSync("pgrep", ["-f", token], { encoding: "utf8" }); + if ((probe.stdout?.trim() ?? "").length > 0) break; + await new Promise((r) => setTimeout(r, 50)); + } + expect(spawnSync("pgrep", ["-f", token], { encoding: "utf8" }).stdout?.trim() ?? "").not.toBe( + "", + ); + expect(plugin.dispose).toBeDefined(); + await plugin.dispose!(); + await new Promise((r) => setTimeout(r, 300)); + const after = spawnSync("pgrep", ["-f", token], { encoding: "utf8" }); + expect(after.stdout?.trim() ?? "").toBe(""); + expect(after.status).not.toBe(0); + await plugin.dispose!(); + await running; + } finally { + spawnSync("pkill", ["-9", "-f", token]); + } + }); + + test("dispose refuses a queued run_shell so it cannot stay running after reap", async () => { + if (process.platform === "win32") return; + const plugin = shellGuardPlugin(process.cwd()); + const handler = plugin.middleware!(fallback); + const token1 = `ic_guard_queued1_${randomUUID()}`; + const token2 = `ic_guard_queued2_${randomUUID()}`; + const first = handler( + { + id: "q-live", + name: "run_shell", + arguments: { command: `IC_GUARD_TAG=${token1} sleep 600` }, + }, + neverAbort(), + ); + try { + const started = Date.now(); + while (Date.now() - started < 5_000) { + const probe = spawnSync("pgrep", ["-f", token1], { encoding: "utf8" }); + if ((probe.stdout?.trim() ?? "").length > 0) break; + await new Promise((r) => setTimeout(r, 50)); + } + expect( + spawnSync("pgrep", ["-f", token1], { encoding: "utf8" }).stdout?.trim() ?? "", + ).not.toBe(""); + + const queued = handler( + { + id: "q-wait", + name: "run_shell", + arguments: { command: `IC_GUARD_TAG=${token2} sleep 600` }, + }, + neverAbort(), + ); + + let disposeError: unknown; + try { + await plugin.dispose!(); + } catch (err) { + disposeError = err; + } + + await first; + await Promise.race([queued, new Promise((r) => setTimeout(r, 400))]); + await new Promise((r) => setTimeout(r, 200)); + + const leftover1 = + spawnSync("pgrep", ["-f", token1], { encoding: "utf8" }).stdout?.trim() ?? ""; + const leftover2 = + spawnSync("pgrep", ["-f", token2], { encoding: "utf8" }).stdout?.trim() ?? ""; + expect(leftover1).toBe(""); + expect(leftover2).toBe(""); + if (leftover2.length > 0) { + expect(disposeError).toBeDefined(); + } else { + const queuedResult = await queued; + expect(queuedResult.isError).toBe(true); + expect(String(queuedResult.content)).toMatch(/disposed/); + expect(disposeError).toBeUndefined(); + } + spawnSync("pkill", ["-9", "-f", token2]); + await Promise.race([queued, new Promise((r) => setTimeout(r, 1_000))]); + } finally { + spawnSync("pkill", ["-9", "-f", token1]); + spawnSync("pkill", ["-9", "-f", token2]); + } + }, 15_000); + + test("overlapping dispose joins the in-flight reap", async () => { + if (process.platform === "win32") return; + const plugin = shellGuardPlugin(process.cwd()); + const handler = plugin.middleware!(fallback); + const token = `ic_guard_join_${randomUUID()}`; + const running = handler( + { + id: "join-live", + name: "run_shell", + arguments: { command: `IC_GUARD_TAG=${token} sleep 600` }, + }, + neverAbort(), + ); + try { + const started = Date.now(); + while (Date.now() - started < 5_000) { + const probe = spawnSync("pgrep", ["-f", token], { encoding: "utf8" }); + if ((probe.stdout?.trim() ?? "").length > 0) break; + await new Promise((r) => setTimeout(r, 50)); + } + expect(plugin.dispose).toBeDefined(); + const first = plugin.dispose!(); + const second = plugin.dispose!(); + expect(second).toBe(first); + await Promise.all([first, second]); + await running; + const after = spawnSync("pgrep", ["-f", token], { encoding: "utf8" }); + expect(after.stdout?.trim() ?? "").toBe(""); + } finally { + spawnSync("pkill", ["-9", "-f", token]); + } + }); + + test("dispose fails when a child survives the reap window", async () => { + const child = Object.assign(new EventEmitter(), { + exitCode: null, + signalCode: null, + kill: () => true, + }) as ChildProcess; + await expect(reapLiveChildren(new Set([child]))).rejects.toThrow( + /still live after 2000ms reap/, + ); + }, 10_000); }); diff --git a/src/plugins/shell-guard-plugin.ts b/src/plugins/shell-guard-plugin.ts index db1309f5..95a7da70 100644 --- a/src/plugins/shell-guard-plugin.ts +++ b/src/plugins/shell-guard-plugin.ts @@ -1,4 +1,4 @@ -import { spawn } from "node:child_process"; +import { spawn, type ChildProcess } from "node:child_process"; import { realpathSync } from "node:fs"; import type { ToolPlugin } from "@intx/tools-posix"; import { killProcessTree, type BackgroundShellRegistry } from "../shell/background-shell.js"; @@ -207,11 +207,55 @@ export class BoundedShellOutput { } } +function waitChildClose(child: ChildProcess): Promise { + if (child.exitCode !== null || child.signalCode !== null) return Promise.resolve(); + return new Promise((resolve) => { + child.once("close", () => resolve()); + child.once("error", () => resolve()); + }); +} + +const SHELL_GUARD_DISPOSE_REAP_MS = 2_000; + +function childStillLive(child: ChildProcess): boolean { + return child.exitCode === null && child.signalCode === null; +} + +// Abort SIGKILLs the process group immediately (runGuardedShell onAbort). This +// window is only a backstop for children still tracked at dispose. Leftovers +// after it must fail teardown; do not stretch the process-exit 2s deadline. +export async function reapLiveChildren(liveChildren: Set): Promise { + const remaining = [...liveChildren]; + for (const child of remaining) killProcessTree(child); + if (remaining.length === 0) return; + const closed = Promise.all(remaining.map(waitChildClose)); + let timer: ReturnType | undefined; + const timedOut = new Promise<"timeout">((resolve) => { + timer = setTimeout(() => resolve("timeout"), SHELL_GUARD_DISPOSE_REAP_MS); + }); + try { + const winner = await Promise.race([closed.then(() => "closed" as const), timedOut]); + if (winner === "closed") return; + const stillLive = remaining.filter(childStillLive); + if (stillLive.length > 0) { + throw new Error( + `${stillLive.length} shell child process${stillLive.length === 1 ? "" : "es"} still live after ${SHELL_GUARD_DISPOSE_REAP_MS}ms reap`, + ); + } + } finally { + if (timer !== undefined) clearTimeout(timer); + } +} export async function runGuardedShell( args: RunShellArgs, signal: AbortSignal, + liveChildren?: Set, + isDisposed?: () => boolean, ): Promise { signal.throwIfAborted(); + if (isDisposed?.()) { + throw new Error("run_shell refused: shell guard disposed"); + } // Arm setTimeout only when a positive timeout was resolved. No built-in default. const timeoutMs = args.timeout !== undefined && args.timeout > 0 ? args.timeout : undefined; @@ -230,8 +274,13 @@ export async function runGuardedShell( // settings.env overrides layered on top. env: args.env !== undefined ? { ...process.env, ...args.env } : undefined, }); + liveChildren?.add(child); + const dropLive = () => { + liveChildren?.delete(child); + }; if (child.stdout === null || child.stderr === null) { + dropLive(); reject(new Error("child process streams are null; stdio misconfigured")); return; } @@ -300,10 +349,12 @@ export async function runGuardedShell( }; child.on("error", (err) => { + dropLive(); settle(new Error(`failed to spawn command: ${args.command}`, { cause: err })); }); child.on("close", (code, sig) => { + dropLive(); if (settled) return; const exitCode = code ?? (sig !== null ? 128 : 1); finishOutput(exitCode, false); @@ -347,11 +398,20 @@ export function shellGuardPlugin( const maxOutputBytes = timeoutConfig?.maxOutputBytes ?? MAX_SHELL_OUTPUT_BYTES; const sessionRoot = realpathSync(cwd); let retainedShellCwd = sessionRoot; + const liveChildren = new Set(); + let disposed = false; + let disposal: Promise | undefined; // Serialize run_shell so concurrent tools cannot race retained cwd updates // (last-writer-wins or a non-cd call finishing after a cd and resetting cwd). let shellChain: Promise = Promise.resolve(); const enqueueShell = (fn: () => Promise): Promise => { - const run = shellChain.then(fn, fn); + const runUnlessDisposed = (): Promise => { + if (disposed) { + return Promise.reject(new Error("run_shell refused: shell guard disposed")); + } + return fn(); + }; + const run = shellChain.then(runUnlessDisposed, runUnlessDisposed); shellChain = run.then( () => undefined, () => undefined, @@ -442,6 +502,8 @@ export function shellGuardPlugin( ...(env !== undefined ? { env } : {}), }, signal, + liveChildren, + () => disposed, ); const parsed = parsePwdProbeOutput(output); if (perCallCwdRaw === undefined && parsed.finalCwd !== undefined) { @@ -473,7 +535,11 @@ export function shellGuardPlugin( isError: true, }; } - }); + }).catch((err: unknown) => ({ + callId: call.id, + content: err instanceof Error ? err.message : String(err), + isError: true, + })); } if (SEARCH_TOOLS.has(call.name)) { @@ -527,5 +593,15 @@ export function shellGuardPlugin( return next(call, signal); }, + dispose: () => { + if (disposal !== undefined) return disposal; + disposed = true; + disposal = (async () => { + await reapLiveChildren(liveChildren); + await shellChain; + await reapLiveChildren(liveChildren); + })(); + return disposal; + }, }; } diff --git a/src/session/active-host.test.ts b/src/session/active-host.test.ts index 889fa136..9948672f 100644 --- a/src/session/active-host.test.ts +++ b/src/session/active-host.test.ts @@ -33,4 +33,16 @@ describe("active-host", () => { setActiveDisposeHost(second); expect(getActiveDisposeHost()).toBe(second); }); + + test("accepts an async dispose handle", async () => { + let ran = false; + const handle = async () => { + ran = true; + }; + setActiveDisposeHost(handle); + const active = getActiveDisposeHost(); + expect(active).toBe(handle); + await active?.(); + expect(ran).toBe(true); + }); }); diff --git a/src/session/active-host.ts b/src/session/active-host.ts index 046d7f63..0a2e4da1 100644 --- a/src/session/active-host.ts +++ b/src/session/active-host.ts @@ -5,9 +5,11 @@ // has mounted. Cleared the moment runTUI itself finalizes (normally or via // its own crash path) so a signal arriving after teardown has nothing left // to call. -let activeDisposeHost: (() => void) | null = null; +export type ActiveDisposeHost = () => void | Promise; -export function setActiveDisposeHost(disposeHost: () => void): void { +let activeDisposeHost: ActiveDisposeHost | null = null; + +export function setActiveDisposeHost(disposeHost: ActiveDisposeHost): void { activeDisposeHost = disposeHost; } @@ -15,6 +17,6 @@ export function clearActiveDisposeHost(): void { activeDisposeHost = null; } -export function getActiveDisposeHost(): (() => void) | null { +export function getActiveDisposeHost(): ActiveDisposeHost | null { return activeDisposeHost; } diff --git a/src/subagent/dispose.ts b/src/subagent/dispose.ts index 20dd5931..75f4e7eb 100644 --- a/src/subagent/dispose.ts +++ b/src/subagent/dispose.ts @@ -39,14 +39,62 @@ export const DEFAULT_CLOSE_DEADLINE_MS = 30_000; /** * Honest limits for plugin-spawn teardown (for operator docs and output notes). - * Corbits Code can dispose posix tools and LSP sidecars per sub-agent session; OS - * children spawned inside shell-guard and ripgrep middleware are aborted via the - * tool AbortSignal on cancel/close but are not centrally registered without - * upstream spawn hooks on those plugins. + * Corbits Code disposes posix tools and LSP sidecars per sub-agent session. + * Shell-guard tracks live `run_shell` children and kills the process group on + * plugin dispose (`posixTools.dispose`). Ripgrep detached spawns are not + * tracked in a global registry. */ export const SUBAGENT_PLUGIN_SPAWN_TEARDOWN_LIMITS = - "Per sub-agent session Corbits Code runs agent.close(), drains in-flight tool middleware (best-effort), then posixTools.dispose() (LSP and plugin dispose callbacks). " + - "run_shell and ripgrep spawns honor AbortSignal process-group kill but are not tracked in a global registry until shell-guard/ripgrep expose spawn hooks."; + "Per sub-agent session Corbits Code runs posixTools.dispose() (LSP and plugin dispose callbacks, including in-flight tool drain), then agent.close() and stream drain. " + + "run_shell children are tracked in the shell-guard plugin and killed on posixTools.dispose; ripgrep detached spawns are not tracked in a global registry."; + +/** Fail a hung close instead of resolving as successful teardown. */ +export async function awaitBoundedTeardown( + teardown: Promise, + deadlineMs: number, +): Promise { + let teardownError: unknown; + let timedOut = false; + let timer: ReturnType | undefined; + try { + await Promise.race([ + teardown.then( + () => undefined, + (err: unknown) => { + teardownError = err; + }, + ), + new Promise((resolve) => { + timer = setTimeout(() => { + timedOut = true; + resolve(); + }, deadlineMs); + }), + ]); + } finally { + if (timer !== undefined) clearTimeout(timer); + } + if (teardownError !== undefined) throw teardownError; + if (timedOut) throw new Error(`session close exceeded ${deadlineMs}ms`); +} + +/** + * Always start `close`. Await it only when posix/toolset dispose already + * succeeded; a leftover throw must not wait unbounded on a hung close. + */ +export async function awaitCloseWithoutHidingLeftover( + close: Promise, + leftover: unknown, +): Promise { + if (leftover === undefined) { + await close; + return; + } + void close.then( + () => undefined, + () => undefined, + ); +} export interface SubAgentSpawnSnapshot { inFlightToolCalls: number; @@ -110,19 +158,21 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput) if (input.signal !== undefined && input.closeOnAbort !== undefined) { input.signal.removeEventListener("abort", input.closeOnAbort); } + let posixError: unknown; try { - await input.agent?.close(); - } catch { - // ignore + await input.posixTools.dispose(); + } catch (err: unknown) { + posixError = err; } try { - await input.streamPromise; + await awaitCloseWithoutHidingLeftover(input.agent?.close() ?? Promise.resolve(), posixError); } catch { // ignore } try { - await input.posixTools.dispose(); + await awaitCloseWithoutHidingLeftover(input.streamPromise ?? Promise.resolve(), posixError); } catch { - // LSP shutdown can fail when several sub-agents exit together. + // ignore } + if (posixError !== undefined) throw posixError; } diff --git a/src/subagent/index.test.ts b/src/subagent/index.test.ts index 5aafcab2..8c4603e1 100644 --- a/src/subagent/index.test.ts +++ b/src/subagent/index.test.ts @@ -72,6 +72,80 @@ describe("sub-agent teardown", () => { expect(disposeCount).toBe(2); }); + test("disposeSubAgentSession reaps posix tools before waiting on agent.close", async () => { + const order: string[] = []; + let releaseClose!: () => void; + const closeGate = new Promise((resolve) => { + releaseClose = resolve; + }); + const pending = disposeSubAgentSession({ + agent: { + close: async () => { + order.push("close-start"); + await closeGate; + order.push("close-end"); + }, + }, + posixTools: { + dispose: async () => { + order.push("posix"); + }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 20)); + expect(order).toEqual(["posix", "close-start"]); + releaseClose(); + await pending; + expect(order).toEqual(["posix", "close-start", "close-end"]); + }); + + test("disposeSubAgentSession does not treat a throwing posix dispose as success", async () => { + const posixTools = { + dispose: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }; + + await expect( + disposeSubAgentSession({ + agent: { close: async () => undefined }, + posixTools, + }), + ).rejects.toThrow(/still live after 2000ms reap/); + }); + + test("disposeSubAgentSession surfaces leftover posix dispose when agent.close hangs", async () => { + const posixTools = { + dispose: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }; + let closeStarted = false; + const pending = disposeSubAgentSession({ + agent: { + close: () => { + closeStarted = true; + return new Promise(() => {}); + }, + }, + posixTools, + }); + const result = await Promise.race([ + pending.then( + () => ({ kind: "resolved" as const }), + (err: unknown) => ({ kind: "rejected" as const, err }), + ), + new Promise<{ kind: "timeout" }>((resolve) => { + setTimeout(() => resolve({ kind: "timeout" }), 200); + }), + ]); + expect(closeStarted).toBe(true); + expect(result.kind).toBe("rejected"); + if (result.kind !== "rejected") throw new Error("expected leftover reject"); + expect(result.err).toBeInstanceOf(Error); + expect((result.err as Error).message).toMatch(/still live after 2000ms reap/); + }); + test("spawn registry tracks in-flight plugin tool calls", async () => { const { plugin, snapshot } = createSubAgentSpawnRegistryPlugin(); expect(plugin.middleware).toBeDefined(); @@ -96,9 +170,10 @@ describe("sub-agent teardown", () => { expect(snapshot().inFlightToolCalls).toBe(0); }); - test("teardown limits document missing global spawn registry", () => { + test("teardown limits document shell-guard dispose reaping", () => { expect(SUBAGENT_PLUGIN_SPAWN_TEARDOWN_LIMITS).toContain("posixTools.dispose"); - expect(SUBAGENT_PLUGIN_SPAWN_TEARDOWN_LIMITS).toContain("spawn hooks"); + expect(SUBAGENT_PLUGIN_SPAWN_TEARDOWN_LIMITS).toContain("shell-guard"); + expect(SUBAGENT_PLUGIN_SPAWN_TEARDOWN_LIMITS).toContain("ripgrep"); }); }); diff --git a/src/subagent/lifecycle-tools.test.ts b/src/subagent/lifecycle-tools.test.ts index 0731bbb4..57f2d4a8 100644 --- a/src/subagent/lifecycle-tools.test.ts +++ b/src/subagent/lifecycle-tools.test.ts @@ -86,9 +86,49 @@ describe("close_agent", () => { // Exercise the store directly with a short deadline (the tool itself // uses the real ~30s bound, which would make this test slow). const started = Date.now(); - const childStatus = await sessions.closeOne(wedgedChild.id, 25); + await expect(sessions.closeOne(wedgedChild.id, 25)).rejects.toThrow( + /session close exceeded 25ms/, + ); expect(Date.now() - started).toBeLessThan(500); - expect(childStatus).toBe("shutdown"); + }); + + test("closes remaining siblings after a leftover-child throw, then fails", async () => { + const sessions = createSubAgentSessionStore(); + const parent = sessions.start({ description: "parent", agentId: "a", brief: "b" }); + const leftover = sessions.start({ + description: "leftover", + agentId: "a", + brief: "b", + parentSessionId: parent.id, + }); + const sibling = sessions.start({ + description: "sibling", + agentId: "a", + brief: "b", + parentSessionId: parent.id, + }); + const closedOrder: string[] = []; + sessions.registerClose(leftover.id, async () => { + closedOrder.push(leftover.id); + throw new Error("1 shell child process still live after 2000ms reap"); + }); + sessions.registerClose(sibling.id, async () => { + closedOrder.push(sibling.id); + }); + sessions.registerClose(parent.id, async () => { + closedOrder.push(parent.id); + }); + + const closeAgent = createCloseAgentTool({ + sessions, + fleetRecords: createFleetMailbox(sessions), + }); + await expect(callTool(closeAgent, { target: parent.id })).rejects.toThrow( + /still live after 2000ms reap/, + ); + expect(closedOrder).toContain(leftover.id); + expect(closedOrder).toContain(sibling.id); + expect(closedOrder).toContain(parent.id); }); }); diff --git a/src/subagent/lifecycle-tools.ts b/src/subagent/lifecycle-tools.ts index 4eb6b753..9c92c569 100644 --- a/src/subagent/lifecycle-tools.ts +++ b/src/subagent/lifecycle-tools.ts @@ -41,8 +41,8 @@ export const closeAgentToolDefinition: ToolDefinition = { description: "Permanently close a worker session by agent_id, closing its descendants first. Bounded " + `by a ~${Math.round(DEFAULT_CLOSE_DEADLINE_MS / 1000)}s cleanup deadline per session so a wedged worker cannot hang ` + - "this call — a session that misses the deadline is still marked shutdown; its teardown just " + - "keeps running in the background. Unblocks any in-flight wait_agents on these ids immediately with " + + "this call — a session that misses the deadline is still marked shutdown and the call fails " + + "instead of reporting success while children may still be live. Unblocks any in-flight wait_agents on these ids immediately with " + "status 'interrupted'. Closing is permanent: a closed session cannot be resumed.", inputSchema: { type: "object", @@ -184,14 +184,28 @@ export function createCloseAgentTool(deps: CloseAgentToolDeps): AgentTool { .map((s) => ({ id: s.id, parentSessionId: s.parentSessionId })); const order = descendantsClosingOrder(nodes, target); const closed: { agent_id: string; status: AgentLifecycleStatus }[] = []; + const failures: unknown[] = []; for (const id of order) { // Terminalize the wait mailbox before teardown. closeOne flips strip // status to "cancelled", which kills the soft-interrupt fallback that // still requires status === "running" — without this, in-flight // wait_agents hangs until timeout. deps.fleetRecords.interrupt(id); - const status = await deps.sessions.closeOne(id, DEFAULT_CLOSE_DEADLINE_MS); - closed.push({ agent_id: id, status }); + try { + const status = await deps.sessions.closeOne(id, DEFAULT_CLOSE_DEADLINE_MS); + closed.push({ agent_id: id, status }); + } catch (err: unknown) { + failures.push(err); + const after = deps.sessions.get(id); + closed.push({ + agent_id: id, + status: after === undefined ? "not_found" : after.lifecycleStatus, + }); + } + } + if (failures.length === 1) throw failures[0]; + if (failures.length > 1) { + throw new AggregateError(failures, "close_agent leftover dispose failed"); } const own = closed.find((c) => c.agent_id === target); return lifecycleResult( diff --git a/src/subagent/retain-salvage.test.ts b/src/subagent/retain-salvage.test.ts index a960ee8c..fc541779 100644 --- a/src/subagent/retain-salvage.test.ts +++ b/src/subagent/retain-salvage.test.ts @@ -17,14 +17,11 @@ describe("retained session lifecycle", () => { // agentRetained:false, exactly as its real call site does whenever // result.agentRetained isn't true. store.complete(s.id, "Stopped: deadline\n\nPartial work...", { agentRetained: false }); - const after = store.get(s.id); - console.log("lifecycleStatus:", after?.lifecycleStatus, "retained:", after?.retained); const outcome = store.resumeOne(s.id, "more"); - console.log("resumeOne outcome:", JSON.stringify(outcome)); expect(outcome.ok).toBe(false); }); - test("cancelAll does not close retained completed sessions", () => { + test("cancelAll closes salvaged retained completed sessions", async () => { const store = createSubAgentSessionStore({ maxCompleted: 5 }); const s = store.start({ description: "worker", @@ -37,8 +34,7 @@ describe("retained session lifecycle", () => { closed = true; }); store.complete(s.id, "done"); - const cancelled = store.cancelAll("parent stop"); - console.log("cancelAll returned:", cancelled, "| close invoked:", closed); + await store.cancelAll("parent stop"); expect(closed).toBe(true); }); @@ -63,11 +59,10 @@ describe("retained session lifecycle", () => { store.registerClose(s.id, async () => {}); store.complete(s.id, "done"); } - console.log("sessions retained despite maxRetained=3:", store.list().length); expect(store.list().length).toBeLessThanOrEqual(3); }); - test("a genuinely retained clean completion IS resumable, and cancelAll releases it", () => { + test("a genuinely retained clean completion IS resumable, and cancelAll releases it", async () => { const store = createSubAgentSessionStore({ maxCompleted: 5 }); const s = store.start({ description: "worker", @@ -86,7 +81,7 @@ describe("retained session lifecycle", () => { store.registerFollowup(s.id, async () => "next"); expect(store.resumeOne(s.id, "more").ok).toBe(true); expect(closed).toBe(false); - store.cancelAll("parent stop"); + await store.cancelAll("parent stop"); expect(closed).toBe(true); }); diff --git a/src/subagent/run-persist-close.test.ts b/src/subagent/run-persist-close.test.ts new file mode 100644 index 00000000..8816fbf9 --- /dev/null +++ b/src/subagent/run-persist-close.test.ts @@ -0,0 +1,210 @@ +/** + * Persist close_agent must surface a leftover-child posix dispose, not treat + * it as a successful bounded close. + */ +import { describe, expect, test } from "bun:test"; +import { mkdtemp } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import type { ReactorEmittedEvent } from "@intx/inference"; + +import { withMockedModuleDuring } from "../../tests/helpers/mock-module.js"; +import { createPermissionGate } from "../permission/gate.js"; +import type { RunSubAgentParams } from "./types.js"; + +const permissionGate = createPermissionGate({ + approvals: [], + interactive: false, + skipPermissions: true, + reactorGated: false, +}); + +function stubAgent() { + return { + send: async () => { + await new Promise((resolve) => setTimeout(resolve, 20)); + return { + type: "reply" as const, + reply: "done", + turn: { role: "assistant", content: [] }, + }; + }, + stream: () => (async function* (): AsyncGenerator {})(), + deliver: () => {}, + close: async () => {}, + setSource: () => {}, + setSources: () => {}, + history: async () => [], + checkpoints: async () => [], + readAt: async () => [], + blobReader: {}, + }; +} + +describe("persist close_agent leftover dispose", () => { + test("onAgentReady close rejects when posix dispose reports leftover children", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-persist-close-")); + + await withMockedModuleDuring( + import.meta.resolve("@intx/tools-posix"), + (real: typeof import("@intx/tools-posix")) => ({ + ...real, + createPosixTools: (opts: Parameters[0]) => + Object.assign(real.createPosixTools(opts), { + dispose: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }), + }), + async () => + withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => + stubAgent() as unknown as Awaited< + ReturnType + >, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + let handles: + | { + close: (deadlineMs?: number) => Promise; + } + | undefined; + const params: RunSubAgentParams = { + cwd, + workdirBase: join(cwd, ".ctx"), + permissionGate, + provider: { providerName: "test", baseURL: "http://localhost", model: "test-model" }, + description: "persist close leftover probe", + prompt: "finish the first turn", + persist: true, + onAgentReady: (h) => { + handles = h; + }, + }; + const result = await runSubAgent(params); + expect(result.agentRetained).toBe(true); + if (handles === undefined) throw new Error("onAgentReady never fired"); + await expect(handles.close(1000)).rejects.toThrow(/still live after 2000ms reap/); + }, + ), + ); + }); + + test("onAgentReady close reaps posix tools before a hung agent.close and fails the deadline", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-persist-close-hung-")); + let posixDisposed = false; + + await withMockedModuleDuring( + import.meta.resolve("@intx/tools-posix"), + (real: typeof import("@intx/tools-posix")) => ({ + ...real, + createPosixTools: (opts: Parameters[0]) => + Object.assign(real.createPosixTools(opts), { + dispose: async () => { + posixDisposed = true; + }, + }), + }), + async () => + withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => + ({ + ...stubAgent(), + close: () => new Promise(() => {}), + }) as unknown as Awaited>, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + let handles: + | { + close: (deadlineMs?: number) => Promise; + } + | undefined; + const params: RunSubAgentParams = { + cwd, + workdirBase: join(cwd, ".ctx"), + permissionGate, + provider: { providerName: "test", baseURL: "http://localhost", model: "test-model" }, + description: "persist close hung close probe", + prompt: "finish the first turn", + persist: true, + onAgentReady: (h) => { + handles = h; + }, + }; + const result = await runSubAgent(params); + expect(result.agentRetained).toBe(true); + if (handles === undefined) throw new Error("onAgentReady never fired"); + await expect(handles.close(50)).rejects.toThrow(/session close exceeded 50ms/); + expect(posixDisposed).toBe(true); + }, + ), + ); + }); + + test("onAgentReady close surfaces leftover posix dispose when agent.close hangs", async () => { + const cwd = await mkdtemp(join(tmpdir(), "corbits-persist-close-leftover-hang-")); + let closeStarted = false; + + await withMockedModuleDuring( + import.meta.resolve("@intx/tools-posix"), + (real: typeof import("@intx/tools-posix")) => ({ + ...real, + createPosixTools: (opts: Parameters[0]) => + Object.assign(real.createPosixTools(opts), { + dispose: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }), + }), + async () => + withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => + ({ + ...stubAgent(), + close: () => { + closeStarted = true; + return new Promise(() => {}); + }, + }) as unknown as Awaited>, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + let handles: + | { + close: (deadlineMs?: number) => Promise; + } + | undefined; + const params: RunSubAgentParams = { + cwd, + workdirBase: join(cwd, ".ctx"), + permissionGate, + provider: { providerName: "test", baseURL: "http://localhost", model: "test-model" }, + description: "persist close leftover hung close probe", + prompt: "finish the first turn", + persist: true, + onAgentReady: (h) => { + handles = h; + }, + }; + const result = await runSubAgent(params); + expect(result.agentRetained).toBe(true); + if (handles === undefined) throw new Error("onAgentReady never fired"); + await expect(handles.close(200)).rejects.toThrow(/still live after 2000ms reap/); + expect(closeStarted).toBe(true); + }, + ), + ); + }); +}); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index 2794b0ba..1a41a773 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -130,6 +130,7 @@ import { disposeSubAgentSession, isSubAgentCancelError, DEFAULT_CLOSE_DEADLINE_MS, + awaitBoundedTeardown, } from "./dispose.js"; import { createFleetMailbox, @@ -1106,12 +1107,6 @@ async function runSubAgentInner( } catch { // close is idempotent; ignore races with disposeSubAgentSession. } - try { - backgroundShells.disposeAll("sub-agent closed"); - await posixTools.dispose(); - } catch { - // ignore - } })(); }; if (runController.signal.aborted) { @@ -1123,33 +1118,33 @@ async function runSubAgentInner( // Hand the caller a bounded, idempotent close it can call at any time // (close_agent) — independent of whether this run ends up retained. // Aborting first stops a still-running turn before tearing down; on an - // already-finished turn the abort is a no-op. The timeout races teardown - // itself so a wedged descendant cannot hang the caller — see dispose.ts - // for the close()-ordering issue that can stall it. + // already-finished turn the abort is a no-op. posix dispose/reap runs + // before waiting on agent.close so a wedged close cannot skip killing + // detached run_shell children. The deadline abandons a hung close and + // fails rather than reporting success while children may still be live. if (params.onAgentReady !== undefined) { const boundedClose = async (deadlineMs = DEFAULT_CLOSE_DEADLINE_MS): Promise => { if (!runController.signal.aborted) runController.abort(new Error("closed by close_agent")); - const teardown = disposeSubAgentSession({ - signal: runController.signal, - ...(closeOnAbort !== undefined ? { closeOnAbort } : {}), - agent, - ...(streamPromise !== undefined ? { streamPromise } : {}), - posixTools, - }).catch(() => { - // Best-effort: a wedged descendant must not reject the caller. - }); - await Promise.race([ - teardown, - new Promise((resolve) => setTimeout(resolve, deadlineMs)), - ]); - // The finally block kept the parent-abort forwarding listener alive - // for a persisted session (see runController.dispose's doc); now that - // this session is actually closing, tear it down for real. - runController.dispose(); + try { + await awaitBoundedTeardown( + disposeSubAgentSession({ + signal: runController.signal, + ...(closeOnAbort !== undefined ? { closeOnAbort } : {}), + agent, + ...(streamPromise !== undefined ? { streamPromise } : {}), + posixTools, + }), + deadlineMs, + ); + } finally { + // The finally block kept the parent-abort forwarding listener alive + // for a persisted session (see runController.dispose's doc); now that + // this session is actually closing, tear it down for real. + runController.dispose(); + } }; // Interrupt only fires interruptController — never runController/ - // close, so it cannot hit the close()-ordering wedge documented in - // dispose.ts. + // close, so it cannot hang teardown on a wedged agent.close. const interrupt = (): void => { if (!interruptController.signal.aborted) { interruptController.abort(new Error("interrupted by interrupt_agent")); diff --git a/src/subagent/session-store.test.ts b/src/subagent/session-store.test.ts index 91830468..9facbbf1 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -574,19 +574,66 @@ describe("CL-6943 reusable worker sessions", () => { expect(store.resumeOne("missing", "more")).toEqual({ ok: false, status: "not_found" }); }); - test("closeOne is bounded by its deadline when the registered close hangs forever", async () => { + test("cancelAll then closeOne does not swallow a leftover-child throw as shutdown success", async () => { + const store = createSubAgentSessionStore(); + const session = store.start({ + description: "d", + agentId: "a", + brief: "b", + retained: true, + }); + store.registerClose(session.id, async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }); + store.complete(session.id, "done", { agentRetained: true }); + + const leftover = /still live after 2000ms reap/; + let cancelThrew = false; + try { + await store.cancelAll("parent stop"); + } catch (err) { + expect(err).toBeInstanceOf(Error); + expect((err as Error).message).toMatch(leftover); + cancelThrew = true; + } + let closeThrew = false; + let closeStatus: string | undefined; + try { + closeStatus = await store.closeOne(session.id, 1000); + } catch (err) { + expect(err).toBeInstanceOf(Error); + expect((err as Error).message).toMatch(leftover); + closeThrew = true; + } + expect(cancelThrew || closeThrew).toBe(true); + expect(cancelThrew === false && closeThrew === false && closeStatus === "shutdown").toBe(false); + }); + + test("closeOne fails a hung close instead of reporting shutdown success", async () => { const store = createSubAgentSessionStore(); const session = store.start({ description: "d", agentId: "a", brief: "b", retained: true }); store.registerClose(session.id, () => new Promise(() => {})); // never resolves const started = Date.now(); - const status = await store.closeOne(session.id, 25); + await expect(store.closeOne(session.id, 25)).rejects.toThrow(/session close exceeded 25ms/); expect(Date.now() - started).toBeLessThan(500); - expect(status).toBe("shutdown"); expect(store.get(session.id)?.lifecycleStatus).toBe("shutdown"); expect(store.get(session.id)?.retained).toBe(false); }); + test("closeOne rejects when the registered close throws a leftover child after reap", async () => { + const store = createSubAgentSessionStore(); + const session = store.start({ description: "d", agentId: "a", brief: "b", retained: true }); + store.registerClose(session.id, async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }); + + await expect(store.closeOne(session.id, 1000)).rejects.toThrow(/still live after 2000ms reap/); + expect(store.get(session.id)?.lifecycleStatus).toBe("shutdown"); + expect(store.get(session.id)?.retained).toBe(false); + expect(await store.closeOne(session.id, 1000)).toBe("shutdown"); + }); + test("closeOne is idempotent and returns not_found for an unknown id", async () => { const store = createSubAgentSessionStore(); const session = store.start({ description: "d", agentId: "a", brief: "b", retained: true }); diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index 33c02114..7d1dda95 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -7,7 +7,7 @@ import type { ReactorEmittedEvent } from "@intx/inference"; import { getLogger } from "@intx/log"; import { LOG_NAMESPACE_ROOT } from "../branding.js"; -import { DEFAULT_CLOSE_DEADLINE_MS } from "./dispose.js"; +import { awaitBoundedTeardown, DEFAULT_CLOSE_DEADLINE_MS } from "./dispose.js"; import { isAlreadyClosed, isLiveStrip, @@ -22,6 +22,20 @@ import type { AdmissionQueue, AdmissionStatus } from "./admission.js"; const log = getLogger([LOG_NAMESPACE_ROOT, "subagent", "session-store"]); +async function invokeCloseBounded( + close: (deadlineMs?: number) => Promise, + deadlineMs: number, +): Promise { + try { + await awaitBoundedTeardown(close(deadlineMs), deadlineMs); + } catch (err: unknown) { + log.warn("session close raced deadline: {error}", { + error: err instanceof Error ? err.message : String(err), + }); + throw err; + } +} + export type SubAgentSessionStatus = "running" | "done" | "failed" | "cancelled"; /** @@ -194,8 +208,11 @@ export interface SubAgentSessionStore { // Abort a running session and mark it cancelled. Returns true when a running // session was cancelled; false if missing or already terminal. cancel(id: string, reason?: string): boolean; - // Cancel every running session. Returns the ids that transitioned. - cancelAll(reason?: string): string[]; + // Cancel every running session. Closes retained workers with the same + // deadline race as closeOne: leftover-child throws and hung closes reject + // instead of reporting success while children may still be live. Returns + // the ids that transitioned to cancelled. + cancelAll(reason?: string): Promise; // CL-6943: flips a "pending_init" session to "running" once its agent // object actually exists. No-op on an unknown id or one already past init. markRunning(id: string): void; @@ -1226,14 +1243,12 @@ export function createSubAgentSessionStore( // Bounded here too, defense-in-depth against a caller-registered // close that does not honor its own deadline argument — a wedged // descendant must not hang the whole close_agent call. - await Promise.race([ - close(deadlineMs).catch((err: unknown) => { - log.warn("session close raced deadline: {error}", { - error: err instanceof Error ? err.message : String(err), - }); - }), - new Promise((resolve) => setTimeout(resolve, deadlineMs)), - ]); + let closeError: unknown; + try { + await invokeCloseBounded(close, deadlineMs); + } catch (err: unknown) { + closeError = err; + } if (keepFailed) { // fail() already stamped failed; invoke leftover teardown without // rewriting that to shutdown. @@ -1243,6 +1258,7 @@ export function createSubAgentSessionStore( deliverHandles.delete(id); runInFlight.delete(id); pruneCompleted(); + if (closeError !== undefined) throw closeError; const after = sessions.get(id); return after === undefined ? "not_found" : projectLifecycleStatus(after.lifecycle); } @@ -1266,6 +1282,7 @@ export function createSubAgentSessionStore( deliverHandles.delete(id); runInFlight.delete(id); pruneCompleted(); + if (closeError !== undefined) throw closeError; return "shutdown"; }, @@ -1481,7 +1498,7 @@ export function createSubAgentSessionStore( return cancelSession(id, reason); }, - cancelAll(reason = DEFAULT_CANCEL_REASON): string[] { + async cancelAll(reason = DEFAULT_CANCEL_REASON): Promise { // Snapshot before cancelSession: markCancelled clears retained, and a // resumed retained worker is strip-live so the first loop would otherwise // skip the close-handle pass (CL-7001). @@ -1493,10 +1510,17 @@ export function createSubAgentSessionStore( for (const session of running) { if (cancelSession(session.id, reason)) cancelled.push(session.id); } + const pendingCloses: Promise[] = []; for (const id of retainedIds) { const session = sessions.get(id); if (session === undefined || session.lifecycle.state === "shutdown") continue; - releaseHandles(id); + const close = closeHandles.get(id); + cancelAskInternal(id, "session handles released"); + closeHandles.delete(id); + cancelHandles.delete(id); + interruptHandles.delete(id); + followupHandles.delete(id); + deliverHandles.delete(id); mutate(id, (s) => { s.lifecycle = { state: "shutdown", @@ -1505,6 +1529,15 @@ export function createSubAgentSessionStore( }; s.retained = false; }); + if (close !== undefined) { + pendingCloses.push(invokeCloseBounded(close, DEFAULT_CLOSE_DEADLINE_MS)); + } + } + const results = await Promise.allSettled(pendingCloses); + const failures = results.flatMap((r) => (r.status === "rejected" ? [r.reason] : [])); + if (failures.length === 1) throw failures[0]; + if (failures.length > 1) { + throw new AggregateError(failures, "session cancelAll close failed"); } return cancelled; }, diff --git a/src/subagent/types.ts b/src/subagent/types.ts index 11769ef1..efe7e393 100644 --- a/src/subagent/types.ts +++ b/src/subagent/types.ts @@ -190,9 +190,8 @@ export type RunSubAgentParams = { * * Always fired regardless of `persist`, so a caller can act on a * still-running session too, not only a retained one. The deadline - * argument to `close` bounds how long teardown may take; a wedged close is - * abandoned (not awaited further) once it elapses rather than hanging the - * caller. + * argument to `close` bounds how long teardown may take; a wedged close + * fails rather than reporting success while children may still be live. */ onAgentReady?: (handles: { close: (deadlineMs?: number) => Promise; diff --git a/src/tui/runner-exit-code.test.ts b/src/tui/runner-exit-code.test.ts index 8d4860f0..3d4fc90f 100644 --- a/src/tui/runner-exit-code.test.ts +++ b/src/tui/runner-exit-code.test.ts @@ -65,6 +65,16 @@ describe("resolveExitCode", () => { }); expect(code).toBe(0); }); + + test("returns 1 when teardown failed even if status is done", () => { + const code = resolveExitCode({ + runError: undefined, + sinkError: undefined, + status: "done", + teardownFailed: true, + }); + expect(code).toBe(1); + }); }); describe("resolveLocalSettingsPath", () => { diff --git a/src/tui/runner/exit.test.ts b/src/tui/runner/exit.test.ts new file mode 100644 index 00000000..e71e7612 --- /dev/null +++ b/src/tui/runner/exit.test.ts @@ -0,0 +1,111 @@ +import { describe, expect, spyOn, test } from "bun:test"; +import { getLogger } from "@intx/log"; + +import { LOG_NAMESPACE_ROOT } from "../../branding.js"; +import { finalizeTUIRun } from "./exit.js"; +import type { RunnerServices, RunnerState } from "./state.js"; + +function stubQuit(args: { awaitTail: () => Promise; shutdownRuntime: () => Promise }): { + state: RunnerState; + services: RunnerServices; +} { + const state = { + host: { + waitUntilExit: async () => undefined, + }, + shutdownRuntime: args.shutdownRuntime, + runError: undefined, + streamPromise: Promise.resolve(), + config: { cwd: "/tmp", task: "t" }, + sessionId: "s", + startedAt: 1, + runTaskTitle: "t", + connectedMcpServers: [], + liveSource: { id: "p", model: "m" }, + } as unknown as RunnerState; + const services = { + sessionOps: { + enqueue: async () => undefined, + awaitTail: args.awaitTail, + }, + cycleRecorder: { dispose: async () => "" }, + mcpConnectController: new AbortController(), + runSink: { + getTurnCollector: () => null, + getRunError: () => undefined, + getStatus: () => "done", + getTurnCount: () => 0, + getTokenUsage: () => ({ + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + thinking: 0, + }), + getToolCallCount: () => 0, + }, + crashGuard: { markFinalized: () => undefined, isFinalized: () => false }, + activeRunHandle: { task: "", startedAt: 0, model: "" }, + hookManager: { dispatchPostRun: async () => undefined }, + liveSessionMode: "orchestrator", + } as unknown as RunnerServices; + return { state, services }; +} + +describe("finalizeTUIRun quit order", () => { + test("starts runtime shutdown without waiting on a hung session-op tail", async () => { + const order: string[] = []; + let settleTail!: (err: Error) => void; + const hungTail = new Promise((_, reject) => { + settleTail = reject; + }); + const { state, services } = stubQuit({ + awaitTail: async () => { + order.push("tail"); + await hungTail; + }, + shutdownRuntime: async () => { + order.push("shutdown"); + }, + }); + + const pending = finalizeTUIRun(state, services); + try { + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(order[0]).toBe("shutdown"); + } finally { + settleTail(new Error("stop")); + } + await expect(pending).rejects.toThrow("stop"); + }); + + test("logs a runtime shutdown failure instead of swallowing it", async () => { + const logger = getLogger([LOG_NAMESPACE_ROOT, "tui"]); + const errorSpy = spyOn(logger, "error"); + let settleTail!: (err: Error) => void; + const hungTail = new Promise((_, reject) => { + settleTail = reject; + }); + const { state, services } = stubQuit({ + awaitTail: () => hungTail, + shutdownRuntime: async () => { + throw new Error("plugin dispose failed"); + }, + }); + + const pending = finalizeTUIRun(state, services); + try { + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(errorSpy).toHaveBeenCalled(); + const logged = errorSpy.mock.calls as unknown as readonly (readonly unknown[])[]; + const first = logged[0]; + expect(first).toBeDefined(); + expect(String(first?.[0])).toMatch(/shutdown/i); + expect(first?.[1]).toEqual({ error: "plugin dispose failed" }); + } finally { + errorSpy.mockRestore(); + settleTail(new Error("stop")); + } + await expect(pending).rejects.toThrow("stop"); + }); +}); diff --git a/src/tui/runner/exit.ts b/src/tui/runner/exit.ts index bdca53ad..e72e080a 100644 --- a/src/tui/runner/exit.ts +++ b/src/tui/runner/exit.ts @@ -43,29 +43,37 @@ const tuiLogger = getLogger([LOG_NAMESPACE_ROOT, "tui"]); export function resetSessionForRotation( state: Pick, services: Pick, -): void { +): Promise { + let cancelledWorkers: Promise = Promise.resolve([]); const reset = (): void => { services.deliveryGeneration.bump(); cancelFeedbackCapture(); services.emitter.emit("session.clear"); - services.subAgentSessions.cancelAll("Session cleared"); + cancelledWorkers = services.subAgentSessions.cancelAll("Session cleared"); }; if (state.withFleetPublicationSuspended === undefined) { reset(); } else { state.withFleetPublicationSuspended(reset); } + return cancelledWorkers; } export interface ResolveExitCodeArgs { runError: string | undefined; sinkError: string | undefined; status: RunSummary["status"]; + teardownFailed?: boolean; } export function resolveExitCode(args: ResolveExitCodeArgs): number { - const { runError, sinkError, status } = args; - if (runError !== undefined || sinkError !== undefined || status !== "done") { + const { runError, sinkError, status, teardownFailed } = args; + if ( + teardownFailed === true || + runError !== undefined || + sinkError !== undefined || + status !== "done" + ) { return 1; } return 0; @@ -465,12 +473,13 @@ export async function createRunLifecycle( // abort handles → child agent.close) before clearing the session store so // /clear does not leave orphaned child reactors burning tokens. const newSession = (): void => { - resetSessionForRotation(state, services); + const cancelledWorkers = resetSessionForRotation(state, services); // Backend rotation is always enqueued regardless of contention; the queue // serialises it behind any in-progress op. Sub-agents nest under the new // session automatically because getWorkdirBase reads the live sessionId. void enqueueOp(async () => { try { + await cancelledWorkers; // Tear the old agent down and dispose the recorder before workdir is // repointed: the pump can deliver stray deltas until the stream // settles, and a dead cycle's partial must land in the session that @@ -549,9 +558,21 @@ export async function finalizeTUIRun( services: RunnerServices, ): Promise { await hostOf(state).waitUntilExit(); - // Stop inference and every worker before persistence, hooks, or telemetry can - // delay process exit. Closing the terminal is a process-lifetime boundary. - await state.shutdownRuntime?.(); + // Stop workers before awaiting the session-op tail so a hung enqueue cannot + // delay abort/reap. Persistence, hooks, and telemetry stay after stop. + // Toolset dispose lives inside shutdownRuntime so quit, crash, and signals + // share one owner. + let teardownFailed = false; + try { + await state.shutdownRuntime?.(); + } catch (err) { + teardownFailed = true; + tuiLogger.error("runtime shutdown failed: {error}", { + error: err instanceof Error ? err.message : String(err), + }); + } + await services.sessionOps.awaitTail(); + state.stopFleetReporting?.(); // Quitting mid-stream is an abnormal end for the in-flight cycle: nothing // downstream delivers its terminal event once the app is gone. @@ -613,17 +634,16 @@ export async function finalizeTUIRun( // PerfTrace OTEL export runs once at process exit in main (flushPerfToOtel). await getTelemetry().flush(); - await services.sessionOps.awaitTail(); try { await state.streamPromise; } catch { // ignore } - await services.toolset.dispose(); return resolveExitCode({ runError: state.runError, sinkError, status: services.runSink.getStatus(), + teardownFailed, }); } diff --git a/src/tui/runner/host.ts b/src/tui/runner/host.ts index 24f61828..e7d9f1e7 100644 --- a/src/tui/runner/host.ts +++ b/src/tui/runner/host.ts @@ -355,7 +355,10 @@ export async function mountRunnerHost(deps: RunnerHostDeps): Promise // every operator already knows across two keys, and Ctrl+D stays the // prompt's delete-character-under-cursor. + let disposed = false; const dispose = (): void => { + if (disposed) return; + disposed = true; stopBranchWatch(); deps.eventEmitter.off("event", onCostEvent); deps.eventEmitter.off("session.clear", onSessionClear); diff --git a/src/tui/runner/index.ts b/src/tui/runner/index.ts index 3c9e9f44..27a42b1d 100644 --- a/src/tui/runner/index.ts +++ b/src/tui/runner/index.ts @@ -199,7 +199,7 @@ export async function runTUI(initialConfig: Config): Promise { // short-circuits once the clean path has marked the run finalized, and a // throw after that point still has to give the terminal back. try { - start.crashGuard.invokeDisposeHost(); + await start.crashGuard.invokeDisposeHost(); } catch (disposeErr: unknown) { tuiLogger.warn("crash finalize: host dispose failed: {error}", { error: disposeErr instanceof Error ? disposeErr.message : String(disposeErr), diff --git a/src/tui/runner/shutdown.ts b/src/tui/runner/shutdown.ts index e6ea1bd1..eb2c9b0e 100644 --- a/src/tui/runner/shutdown.ts +++ b/src/tui/runner/shutdown.ts @@ -1,33 +1,50 @@ +import { awaitCloseWithoutHidingLeftover } from "../../subagent/dispose.js"; + export interface RuntimeShutdownDeps { disposeHost: () => void; - cancelWorkers: () => void; + cancelWorkers: () => void | Promise; closeAgent: () => Promise; + disposeToolset: () => Promise; +} + +function rethrowShutdownFailures(failures: unknown[]): void { + const first = failures[0]; + if (first === undefined) return; + if (failures.length === 1) throw first; + throw new AggregateError(failures, "runtime shutdown failed"); } /** Start every process-owned teardown path once, even when exit races a signal. */ export function createRuntimeShutdown(deps: RuntimeShutdownDeps): () => Promise { - let started = false; - let completion = Promise.resolve(); + let completion: Promise | undefined; return (): Promise => { - if (started) return completion; - started = true; + if (completion !== undefined) return completion; - try { - deps.disposeHost(); - } catch { - // Every teardown leg is best-effort; one failure must not strand the rest. - } - try { - deps.cancelWorkers(); - } catch { - // The primary agent still needs its abort even if a worker hook misbehaves. - } - try { - completion = deps.closeAgent().catch(() => undefined); - } catch { - completion = Promise.resolve(); - } + completion = (async () => { + const failures: unknown[] = []; + try { + deps.disposeHost(); + } catch (err) { + failures.push(err); + } + try { + await deps.disposeToolset(); + } catch (err) { + failures.push(err); + } + try { + await deps.cancelWorkers(); + } catch (err) { + failures.push(err); + } + try { + await awaitCloseWithoutHidingLeftover(deps.closeAgent(), failures[0]); + } catch (err) { + failures.push(err); + } + rethrowShutdownFailures(failures); + })(); return completion; }; } diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index 75213566..dc7f48d8 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -128,15 +128,14 @@ export function wirePostStartup( const shutdownRuntime = createRuntimeShutdown({ disposeHost: hostOf(state).dispose, - cancelWorkers: () => { - services.subAgentSessions.cancelAll("Session closed"); + cancelWorkers: async () => { + await services.subAgentSessions.cancelAll("Session closed"); }, closeAgent: () => liveAgent(state).close(), + disposeToolset: () => services.toolset.dispose(), }); state.shutdownRuntime = shutdownRuntime; - services.crashGuard.setDisposeHost(() => { - void shutdownRuntime(); - }); + services.crashGuard.setDisposeHost(() => shutdownRuntime()); setActiveDisposeHost(() => services.crashGuard.invokeDisposeHost()); // Harness inference.error events omit providerId; stamp the live catalog id diff --git a/src/tui/runtime-shutdown.test.ts b/src/tui/runtime-shutdown.test.ts index 208c3681..c5545333 100644 --- a/src/tui/runtime-shutdown.test.ts +++ b/src/tui/runtime-shutdown.test.ts @@ -3,51 +3,190 @@ import { describe, expect, test } from "bun:test"; import { createRuntimeShutdown } from "./runner/shutdown.js"; describe("runtime shutdown", () => { - test("restores the terminal, cancels workers, and closes the primary agent", async () => { + test("restores the terminal, disposes the toolset, then closes the primary agent", async () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, + disposeToolset: async () => { + calls.push("toolset"); + }, }); await shutdown(); - expect(calls).toEqual(["host", "workers", "agent"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); test("runs teardown only once when exit and a signal race", async () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, + disposeToolset: async () => { + calls.push("toolset"); + }, }); await Promise.all([shutdown(), shutdown()]); - expect(calls).toEqual(["host", "workers", "agent"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); - test("still aborts workers and the primary agent when host disposal fails", async () => { + test("still runs remaining legs and rejects when host disposal fails", async () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => { calls.push("host"); throw new Error("renderer failure"); }, - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, + disposeToolset: async () => { + calls.push("toolset"); + }, }); - await shutdown(); + await expect(shutdown()).rejects.toThrow("renderer failure"); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); + }); + + test("rejects when toolset dispose throws after other legs ran", async () => { + const calls: string[] = []; + const shutdown = createRuntimeShutdown({ + disposeHost: () => calls.push("host"), + cancelWorkers: () => { + calls.push("workers"); + }, + closeAgent: async () => { + calls.push("agent"); + }, + disposeToolset: async () => { + calls.push("toolset"); + throw new Error("plugin dispose failed"); + }, + }); + + await expect(shutdown()).rejects.toThrow("plugin dispose failed"); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); + }); + + test("rejects when async cancelWorkers throws leftover children", async () => { + const calls: string[] = []; + const shutdown = createRuntimeShutdown({ + disposeHost: () => calls.push("host"), + cancelWorkers: async () => { + calls.push("workers"); + throw new Error("1 shell child process still live after 2000ms reap"); + }, + closeAgent: async () => { + calls.push("agent"); + }, + disposeToolset: async () => { + calls.push("toolset"); + }, + }); + + await expect(shutdown()).rejects.toThrow(/still live after 2000ms reap/); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); + }); + + test("awaits an async toolset dispose before resolving", async () => { + const calls: string[] = []; + let resolveToolset!: () => void; + const toolsetGate = new Promise((resolve) => { + resolveToolset = resolve; + }); + const shutdown = createRuntimeShutdown({ + disposeHost: () => calls.push("host"), + cancelWorkers: () => { + calls.push("workers"); + }, + closeAgent: async () => { + calls.push("agent"); + }, + disposeToolset: async () => { + await toolsetGate; + calls.push("toolset"); + }, + }); - expect(calls).toEqual(["host", "workers", "agent"]); + const pending = shutdown(); + await Promise.resolve(); + expect(calls).toEqual(["host"]); + resolveToolset(); + await pending; + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); + }); + + test("reaps the toolset before waiting on a hung agent close", async () => { + const calls: string[] = []; + let releaseClose!: () => void; + const closeGate = new Promise((resolve) => { + releaseClose = resolve; + }); + const shutdown = createRuntimeShutdown({ + disposeHost: () => calls.push("host"), + cancelWorkers: () => { + calls.push("workers"); + }, + closeAgent: async () => { + await closeGate; + calls.push("agent"); + }, + disposeToolset: async () => { + calls.push("toolset"); + }, + }); + + const pending = shutdown(); + await new Promise((resolve) => setTimeout(resolve, 20)); + expect(calls).toEqual(["host", "toolset", "workers"]); + releaseClose(); + await pending; + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); + }); + + test("surfaces leftover toolset dispose when agent.close hangs", async () => { + let closeStarted = false; + const shutdown = createRuntimeShutdown({ + disposeHost: () => undefined, + cancelWorkers: () => undefined, + closeAgent: () => { + closeStarted = true; + return new Promise(() => {}); + }, + disposeToolset: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }); + const result = await Promise.race([ + shutdown().then( + () => ({ kind: "resolved" as const }), + (err: unknown) => ({ kind: "rejected" as const, err }), + ), + new Promise<{ kind: "timeout" }>((resolve) => { + setTimeout(() => resolve({ kind: "timeout" }), 200); + }), + ]); + expect(closeStarted).toBe(true); + expect(result.kind).toBe("rejected"); + if (result.kind !== "rejected") throw new Error("expected leftover reject"); + expect(result.err).toBeInstanceOf(Error); + expect((result.err as Error).message).toMatch(/still live after 2000ms reap/); }); }); diff --git a/src/tui/session-start.test.ts b/src/tui/session-start.test.ts index 30f7f03a..7122a9e7 100644 --- a/src/tui/session-start.test.ts +++ b/src/tui/session-start.test.ts @@ -63,4 +63,23 @@ describe("createTUICrashGuard", () => { }, ); }); + + test("invokeDisposeHost returns an async dispose handle", async () => { + const { createTUICrashGuard } = await import("./session-start.js"); + const guard = createTUICrashGuard(() => ({ + cwd: "/cwd", + sessionId: "session", + startedAt: 1, + runTaskTitle: "task", + providerName: "provider", + model: "model", + })); + let ran = false; + guard.setDisposeHost(async () => { + await Promise.resolve(); + ran = true; + }); + await guard.invokeDisposeHost(); + expect(ran).toBe(true); + }); }); diff --git a/src/tui/session-start.ts b/src/tui/session-start.ts index 40aeb713..ccea7285 100644 --- a/src/tui/session-start.ts +++ b/src/tui/session-start.ts @@ -69,8 +69,8 @@ export interface TUICrashGuard { isFinalized: () => boolean; markFinalized: () => void; setPartialFlush: (flush: () => Promise) => void; - invokeDisposeHost: () => void; - setDisposeHost: (dispose: () => void) => void; + invokeDisposeHost: () => void | Promise; + setDisposeHost: (dispose: () => void | Promise) => void; bindLiveSession: (get: () => TUILiveSession) => void; finalizeOnCrash: (err: unknown) => Promise; } @@ -87,7 +87,7 @@ export function createTUICrashGuard(getLiveSession: () => TUILiveSession): TUICr // Bound once the host is mounted. Without this the crash path leaves the // renderer alive, so the alternate screen, mouse reporting and raw mode are // never disabled and the operator's terminal is left wedged. - let disposeHost: () => void = () => {}; + let disposeHost: () => void | Promise = () => {}; let getSession = getLiveSession; const finalizeOnCrash = async (err: unknown): Promise => { @@ -146,9 +146,7 @@ export function createTUICrashGuard(getLiveSession: () => TUILiveSession): TUICr setPartialFlush: (flush) => { flushPartialOnCrash = flush; }, - invokeDisposeHost: () => { - disposeHost(); - }, + invokeDisposeHost: () => disposeHost(), setDisposeHost: (dispose) => { disposeHost = dispose; }, diff --git a/tests/fixtures/exec-shutdown-reap/simulate-reap.ts b/tests/fixtures/exec-shutdown-reap/simulate-reap.ts new file mode 100644 index 00000000..559c3c5a --- /dev/null +++ b/tests/fixtures/exec-shutdown-reap/simulate-reap.ts @@ -0,0 +1,104 @@ +// Spawned by tests/integration/exec-shutdown-reap.test.ts. Starts a tagged +// sleep through the real shell-guard plugin, registers exec dispose as the +// process dispose host, then takes the requested exit path so the parent can +// assert the child was reaped. +import { spawnSync } from "node:child_process"; +import { writeFileSync } from "node:fs"; +import type { ToolCall, ToolResult } from "@intx/types/runtime"; + +import { disposeExecRuntime } from "../../../src/exec/runner.js"; +import { shellGuardPlugin } from "../../../src/plugins/shell-guard-plugin.js"; +import { setActiveDisposeHost } from "../../../src/session/active-host.js"; + +const token = process.env["REAP_TOKEN"]; +const path = process.env["REAP_PATH"]; +const countPath = process.env["REAP_COUNT_PATH"]; +if (token === undefined || path === undefined || countPath === undefined) { + throw new Error("REAP_TOKEN, REAP_PATH, and REAP_COUNT_PATH must be set"); +} + +const reapToken = token; +const exitPath = path; +const disposeCountPath = countPath; + +// Handlers must be installed before READY. Importing src/index.js is slow, and +// the parent sends the signal as soon as it sees READY. +if (exitPath === "crash") { + const { installCrashHandlers } = await import("../../../src/index.js"); + installCrashHandlers(); +} else if (exitPath === "signal") { + const { installSignalHandlers } = await import("../../../src/index.js"); + installSignalHandlers(); +} + +const fallback = async (call: ToolCall): Promise => ({ + callId: call.id, + content: "FALLBACK", +}); + +const plugin = shellGuardPlugin(process.cwd()); +if (plugin.middleware === undefined || plugin.dispose === undefined) { + throw new Error("shell-guard plugin missing middleware or dispose"); +} + +let disposeCalls = 0; +const pluginDispose = plugin.dispose.bind(plugin); +const countedDispose = async (): Promise => { + disposeCalls += 1; + writeFileSync(disposeCountPath, String(disposeCalls)); + await pluginDispose(); +}; + +const toolset = { dispose: countedDispose }; +const hungClose = + exitPath === "crash" || exitPath === "signal" + ? { close: () => new Promise(() => undefined) } + : null; +const host = (): Promise => + disposeExecRuntime({ + agent: hungClose, + toolset, + subAgentSessions: null, + }); + +const handler = plugin.middleware(fallback); +const cmd = `bash -c 'IC_GUARD_TAG=${reapToken} sleep 600 & IC_GUARD_TAG=${reapToken} exec sleep 600'`; +void handler( + { id: "reap-live", name: "run_shell", arguments: { command: cmd } }, + new AbortController().signal, +); + +async function waitForTag(): Promise { + const started = Date.now(); + while (Date.now() - started < 5_000) { + const probe = spawnSync("pgrep", ["-f", reapToken], { encoding: "utf8" }); + if ((probe.stdout?.trim() ?? "").length > 0) return; + await new Promise((resolve) => setTimeout(resolve, 50)); + } + throw new Error(`tagged child did not appear: ${reapToken}`); +} + +await waitForTag(); +setActiveDisposeHost(host); +writeFileSync(disposeCountPath, "0"); +process.stdout.write("READY\n"); + +if (exitPath === "quit") { + await host(); + await host(); + writeFileSync(disposeCountPath, String(disposeCalls)); + process.exit(0); +} + +if (exitPath === "crash") { + setImmediate(() => { + throw new Error("simulated reap crash"); + }); + await new Promise(() => undefined); +} + +if (exitPath === "signal") { + await new Promise(() => undefined); +} + +throw new Error(`unknown REAP_PATH: ${exitPath}`); diff --git a/tests/integration/exec-shutdown-reap.test.ts b/tests/integration/exec-shutdown-reap.test.ts new file mode 100644 index 00000000..80016631 --- /dev/null +++ b/tests/integration/exec-shutdown-reap.test.ts @@ -0,0 +1,91 @@ +import { spawnSync } from "node:child_process"; +import { mkdtempSync, readFileSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { randomUUID } from "node:crypto"; + +import { describe, expect, test } from "bun:test"; + +const FIXTURE = join(import.meta.dirname, "../fixtures/exec-shutdown-reap/simulate-reap.ts"); + +async function readLine(stream: ReadableStream): Promise { + const reader = stream.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + while (!buffer.includes("\n")) { + const { value, done } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }); + } + reader.releaseLock(); + return buffer; +} + +async function waitUntilGone(token: string): Promise { + const started = Date.now(); + let stdout = ""; + while (Date.now() - started < 5_000) { + const probe = spawnSync("pgrep", ["-f", token], { encoding: "utf8" }); + stdout = probe.stdout?.trim() ?? ""; + if (stdout.length === 0) return stdout; + await new Promise((resolve) => setTimeout(resolve, 50)); + } + return stdout; +} + +describe.skipIf(process.platform === "win32")( + "integration — exec shutdown reaps shell-guard children", + () => { + test.each([ + ["quit", "", 0], + ["crash", "", 1], + ["signal", "SIGINT", 130], + ["signal", "SIGTERM", 143], + ["signal", "SIGHUP", 129], + ] as const)( + "%s %s reaps the tagged child and disposes once", + async (path, signal, expectedExitCode) => { + const home = mkdtempSync(join(tmpdir(), "corbits-exec-reap-home-")); + const token = `ic_reap_${randomUUID()}`; + const countPath = join(home, "dispose-count"); + + try { + const proc = Bun.spawn(["bun", "run", FIXTURE], { + cwd: home, + env: { + ...process.env, + HOME: home, + REAP_TOKEN: token, + REAP_PATH: path, + REAP_COUNT_PATH: countPath, + }, + stdout: "pipe", + stderr: "pipe", + }); + + const ready = await readLine(proc.stdout); + if (!ready.includes("READY")) { + const errText = await new Response(proc.stderr).text(); + throw new Error( + `fixture did not report READY: ${JSON.stringify(ready)} stderr=${errText}`, + ); + } + + if (signal === "SIGINT" || signal === "SIGTERM" || signal === "SIGHUP") { + proc.kill(signal); + } + const exitCode = await proc.exited; + expect(exitCode).toBe(expectedExitCode); + + const leftover = await waitUntilGone(token); + expect(leftover).toBe(""); + expect(readFileSync(countPath, "utf8").trim()).toBe("1"); + } finally { + spawnSync("pkill", ["-9", "-f", token]); + rmSync(home, { recursive: true, force: true }); + } + }, + 15_000, + ); + }, +); diff --git a/tests/unit/exec/runner.test.ts b/tests/unit/exec/runner.test.ts index a856fbdc..f24400af 100644 --- a/tests/unit/exec/runner.test.ts +++ b/tests/unit/exec/runner.test.ts @@ -15,6 +15,7 @@ import { } from "../../../src/exec/runner.js"; import { BUILD_TOOLS, SKYWALKER_TOOLS } from "../../../src/agent/directors/tool-sets.js"; import { clearActiveRun, getActiveRun, setActiveRun } from "../../../src/session/active-run.js"; +import { getActiveDisposeHost } from "../../../src/session/active-host.js"; import { loadState, type RunState } from "../../../src/session/state.js"; import type { AgentToolset } from "../../../src/agent/tools.js"; import { createSubAgentSessionStore } from "../../../src/subagent/session-store.js"; @@ -288,6 +289,89 @@ describe("runExec", () => { rmSync(home, { recursive: true, force: true }); } }); + + test("dispose failure after toolset exists is once-only and forces a nonzero exit", async () => { + const previous = getActiveRun(); + clearActiveRun(); + const cwd = mkdtempSync(join(tmpdir(), "corbits-exec-dispose-cwd-")); + const home = mkdtempSync(join(tmpdir(), "corbits-exec-dispose-home-")); + const sessionId = "exec-dispose-fail"; + let disposeCalls = 0; + const dummySource = { id: "test", provider: "test", model: "test" } as InferenceSource; + const stderrChunks: string[] = []; + const origWrite = process.stderr.write.bind(process.stderr); + process.stderr.write = ((chunk: string | Uint8Array, ...rest: unknown[]) => { + stderrChunks.push(typeof chunk === "string" ? chunk : Buffer.from(chunk).toString("utf8")); + return origWrite(chunk as never, ...(rest as never[])); + }) as typeof process.stderr.write; + try { + await withMockedModuleDuring( + import.meta.resolve("node:os"), + (real: typeof import("node:os")) => ({ ...real, homedir: () => home }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../../../src/agent/tools.js"), + (real: typeof import("../../../src/agent/tools.js")) => ({ + ...real, + createAgentToolset: async (): Promise => + ({ + dispose: () => { + disposeCalls += 1; + return Promise.reject(new Error("plugin dispose failed")); + }, + }) as AgentToolset, + }), + async () => { + await withMockedModuleDuring( + import.meta.resolve("../../../src/session/assemble-runtime.js"), + (real: typeof import("../../../src/session/assemble-runtime.js")) => ({ + ...real, + resolveLiveSessionSources: () => ({ + sources: [dummySource], + defaultSource: dummySource.id, + selected: dummySource, + }), + assembleChatAgent: () => ({ + directorHolder: {}, + buildAgent: async () => { + throw new Error("buildAgent should not run"); + }, + }), + assembleSessionLifecycle: async () => { + expect(getActiveDisposeHost()).not.toBeNull(); + throw new Error("stop-after-toolset"); + }, + }), + async () => { + const { runExec: runExecUnderMock } = await import("../../../src/exec/runner.js"); + const result = await runExecUnderMock({ + ...bareConfig("do the thing"), + cwd, + sessionId, + director: "builder", + globalSettingsPath: join(home, "settings.json"), + providers: [], + }); + expect(result.exitCode).toBe(1); + expect(result.status).toBe("failed"); + expect(result.error).toMatch(/plugin dispose failed|runtime dispose failed/i); + expect(stderrChunks.join("")).toMatch(/runtime dispose failed/i); + expect(disposeCalls).toBe(1); + expect(getActiveDisposeHost()).toBeNull(); + }, + ); + }, + ); + }, + ); + } finally { + process.stderr.write = origWrite; + if (previous !== null) setActiveRun(previous); + else clearActiveRun(); + rmSync(cwd, { recursive: true, force: true }); + rmSync(home, { recursive: true, force: true }); + } + }); }); describe("disposeExecRuntime", () => { @@ -316,7 +400,116 @@ describe("disposeExecRuntime", () => { expect(aborted).toBe(1); expect(store.get(worker.id)?.status).toBe("cancelled"); - expect(calls).toEqual(["agent", "toolset"]); + expect(calls).toEqual(["toolset", "agent"]); + }); + + test("runs teardown only once when called concurrently", async () => { + const calls: string[] = []; + const toolset = { + dispose: async () => { + calls.push("toolset"); + }, + }; + const args = { + agent: { + close: async () => { + calls.push("agent"); + }, + }, + toolset, + subAgentSessions: null, + }; + + await Promise.all([disposeExecRuntime(args), disposeExecRuntime(args)]); + + expect(calls).toEqual(["toolset", "agent"]); + }); + + test("reaps the toolset before waiting on a hung agent close", async () => { + const calls: string[] = []; + let releaseClose!: () => void; + const closeGate = new Promise((resolve) => { + releaseClose = resolve; + }); + const pending = disposeExecRuntime({ + agent: { + close: async () => { + await closeGate; + calls.push("agent"); + }, + }, + toolset: { + dispose: async () => { + calls.push("toolset"); + }, + }, + subAgentSessions: null, + }); + await new Promise((resolve) => setTimeout(resolve, 20)); + expect(calls).toEqual(["toolset"]); + releaseClose(); + await pending; + expect(calls).toEqual(["toolset", "agent"]); + }); + + test("rejects leftover-child dispose from the toolset", async () => { + await expect( + disposeExecRuntime({ + agent: { close: async () => undefined }, + toolset: { + dispose: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }, + subAgentSessions: null, + }), + ).rejects.toThrow(/still live after 2000ms reap/); + }); + + test("surfaces leftover toolset dispose when agent.close hangs", async () => { + let closeStarted = false; + const pending = disposeExecRuntime({ + agent: { + close: () => { + closeStarted = true; + return new Promise(() => {}); + }, + }, + toolset: { + dispose: async () => { + throw new Error("1 shell child process still live after 2000ms reap"); + }, + }, + subAgentSessions: null, + }); + const result = await Promise.race([ + pending.then( + () => ({ kind: "resolved" as const }), + (err: unknown) => ({ kind: "rejected" as const, err }), + ), + new Promise<{ kind: "timeout" }>((resolve) => { + setTimeout(() => resolve({ kind: "timeout" }), 200); + }), + ]); + expect(closeStarted).toBe(true); + expect(result.kind).toBe("rejected"); + if (result.kind !== "rejected") throw new Error("expected leftover reject"); + expect(result.err).toBeInstanceOf(Error); + expect((result.err as Error).message).toMatch(/still live after 2000ms reap/); + }); + + test("rejects when toolset dispose fails", async () => { + await expect( + disposeExecRuntime({ + agent: { close: async () => undefined }, + toolset: { + dispose: async () => { + throw new Error("plugin dispose failed"); + }, + }, + subAgentSessions: null, + }), + ).rejects.toThrow("plugin dispose failed"); }); }); diff --git a/tests/unit/subagent-session-store.test.ts b/tests/unit/subagent-session-store.test.ts index b91e09a9..7f0bf592 100644 --- a/tests/unit/subagent-session-store.test.ts +++ b/tests/unit/subagent-session-store.test.ts @@ -225,7 +225,7 @@ describe("createSubAgentSessionStore", () => { expect(aborted).toBe(1); }); - test("cancelAll aborts every running session", () => { + test("cancelAll aborts every running session", async () => { let n = 0; const store = createSubAgentSessionStore({ createId: () => `s-${++n}`, @@ -238,7 +238,7 @@ describe("createSubAgentSessionStore", () => { store.registerCancel("s-1", () => aborted.push("s-1")); store.registerCancel("s-2", () => aborted.push("s-2")); store.registerCancel("s-3", () => aborted.push("s-3")); // already done — ignored - const cancelled = store.cancelAll("Parent stop"); + const cancelled = await store.cancelAll("Parent stop"); expect(cancelled.sort()).toEqual(["s-1", "s-2"]); expect(aborted.sort()).toEqual(["s-1", "s-2"]); expect(store.get("s-1")?.status).toBe("cancelled");