From 859b7dfbb1a4dff3b1e55e873a8332367b68165b Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 20:53:50 -0700 Subject: [PATCH 01/14] Kill live shell-guard children on plugin dispose --- docs/ARCHITECTURE.md | 2 +- src/plugins/shell-guard-plugin.test.ts | 31 ++++++++++++++++++++ src/plugins/shell-guard-plugin.ts | 40 +++++++++++++++++++++++++- src/subagent/dispose.ts | 10 +++---- src/subagent/index.test.ts | 5 ++-- 5 files changed, 79 insertions(+), 9 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index ce5e45b27..354600910 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/plugins/shell-guard-plugin.test.ts b/src/plugins/shell-guard-plugin.test.ts index 348145195..b6d4fdc9b 100644 --- a/src/plugins/shell-guard-plugin.test.ts +++ b/src/plugins/shell-guard-plugin.test.ts @@ -747,4 +747,35 @@ 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]); + } + }); }); diff --git a/src/plugins/shell-guard-plugin.ts b/src/plugins/shell-guard-plugin.ts index db1309f5b..2c52c477d 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,9 +207,32 @@ 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; + +async function reapLiveChildren(liveChildren: Set): Promise { + const remaining = [...liveChildren]; + for (const child of remaining) killProcessTree(child); + if (remaining.length === 0) return; + await Promise.race([ + Promise.all(remaining.map(waitChildClose)), + new Promise((resolve) => { + const timer = setTimeout(resolve, SHELL_GUARD_DISPOSE_REAP_MS); + if (typeof timer.unref === "function") timer.unref(); + }), + ]); +} export async function runGuardedShell( args: RunShellArgs, signal: AbortSignal, + liveChildren?: Set, ): Promise { signal.throwIfAborted(); @@ -230,8 +253,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 +328,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,6 +377,8 @@ export function shellGuardPlugin( const maxOutputBytes = timeoutConfig?.maxOutputBytes ?? MAX_SHELL_OUTPUT_BYTES; const sessionRoot = realpathSync(cwd); let retainedShellCwd = sessionRoot; + const liveChildren = new Set(); + 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(); @@ -442,6 +474,7 @@ export function shellGuardPlugin( ...(env !== undefined ? { env } : {}), }, signal, + liveChildren, ); const parsed = parsePwdProbeOutput(output); if (perCallCwdRaw === undefined && parsed.finalCwd !== undefined) { @@ -527,5 +560,10 @@ export function shellGuardPlugin( return next(call, signal); }, + dispose: () => { + if (disposal !== undefined) return disposal; + disposal = reapLiveChildren(liveChildren); + return disposal; + }, }; } diff --git a/src/subagent/dispose.ts b/src/subagent/dispose.ts index 20dd5931e..fcf9c928e 100644 --- a/src/subagent/dispose.ts +++ b/src/subagent/dispose.ts @@ -39,14 +39,14 @@ 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."; + "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."; export interface SubAgentSpawnSnapshot { inFlightToolCalls: number; diff --git a/src/subagent/index.test.ts b/src/subagent/index.test.ts index 5aafcab21..2f82b3872 100644 --- a/src/subagent/index.test.ts +++ b/src/subagent/index.test.ts @@ -96,9 +96,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"); }); }); From bafb37663500ae28f38173dfbfbd19882e1f1980 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 20:54:54 -0700 Subject: [PATCH 02/14] Dispose the toolset inside once-only runtime shutdown --- src/tui/runner/shutdown.ts | 37 ++++++++++++++----- src/tui/runner/wiring.ts | 1 + src/tui/runtime-shutdown.test.ts | 62 +++++++++++++++++++++++++++++--- 3 files changed, 86 insertions(+), 14 deletions(-) diff --git a/src/tui/runner/shutdown.ts b/src/tui/runner/shutdown.ts index e6ea1bd10..e74eaeb42 100644 --- a/src/tui/runner/shutdown.ts +++ b/src/tui/runner/shutdown.ts @@ -2,6 +2,14 @@ export interface RuntimeShutdownDeps { disposeHost: () => void; cancelWorkers: () => void; 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. */ @@ -13,21 +21,32 @@ export function createRuntimeShutdown(deps: RuntimeShutdownDeps): () => Promise< if (started) return completion; started = true; + const failures: unknown[] = []; + try { deps.disposeHost(); - } catch { - // Every teardown leg is best-effort; one failure must not strand the rest. + } catch (err) { + failures.push(err); } 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(); + } catch (err) { + failures.push(err); } + + completion = (async () => { + try { + await deps.closeAgent(); + } catch (err) { + failures.push(err); + } + try { + await deps.disposeToolset(); + } catch (err) { + failures.push(err); + } + rethrowShutdownFailures(failures); + })(); return completion; }; } diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index 752135666..b8d5e8efe 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -132,6 +132,7 @@ export function wirePostStartup( services.subAgentSessions.cancelAll("Session closed"); }, closeAgent: () => liveAgent(state).close(), + disposeToolset: () => services.toolset.dispose(), }); state.shutdownRuntime = shutdownRuntime; services.crashGuard.setDisposeHost(() => { diff --git a/src/tui/runtime-shutdown.test.ts b/src/tui/runtime-shutdown.test.ts index 208c36811..9f56ceb9e 100644 --- a/src/tui/runtime-shutdown.test.ts +++ b/src/tui/runtime-shutdown.test.ts @@ -3,7 +3,7 @@ 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, cancels workers, closes the primary agent, and disposes the toolset", async () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), @@ -11,11 +11,14 @@ describe("runtime shutdown", () => { closeAgent: async () => { calls.push("agent"); }, + disposeToolset: async () => { + calls.push("toolset"); + }, }); await shutdown(); - expect(calls).toEqual(["host", "workers", "agent"]); + expect(calls).toEqual(["host", "workers", "agent", "toolset"]); }); test("runs teardown only once when exit and a signal race", async () => { @@ -26,14 +29,17 @@ describe("runtime shutdown", () => { 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", "workers", "agent", "toolset"]); }); - 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: () => { @@ -44,10 +50,56 @@ describe("runtime shutdown", () => { closeAgent: async () => { calls.push("agent"); }, + disposeToolset: async () => { + calls.push("toolset"); + }, }); - await shutdown(); + await expect(shutdown()).rejects.toThrow("renderer failure"); + expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + }); + + 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", "workers", "agent", "toolset"]); + }); + + 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"); + }, + }); + const pending = shutdown(); + await Promise.resolve(); expect(calls).toEqual(["host", "workers", "agent"]); + resolveToolset(); + await pending; + expect(calls).toEqual(["host", "workers", "agent", "toolset"]); }); }); From 3d8f5ef81f8f622b89966ddb18ef897dc36b811d Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 21:24:26 -0700 Subject: [PATCH 03/14] Await runtime shutdown on crash and signal paths --- CHANGELOG.md | 3 + src/exec/runner.ts | 76 +++++++++++++++++---- src/index.ts | 61 ++++++++++------- src/session/active-host.test.ts | 12 ++++ src/session/active-host.ts | 8 ++- src/tui/runner-exit-code.test.ts | 10 +++ src/tui/runner/exit.ts | 23 +++++-- src/tui/runner/host.ts | 3 + src/tui/runner/index.ts | 2 +- src/tui/runner/wiring.ts | 4 +- src/tui/session-start.test.ts | 19 ++++++ src/tui/session-start.ts | 10 ++- tests/unit/exec/runner.test.ts | 110 +++++++++++++++++++++++++++++++ 13 files changed, 287 insertions(+), 54 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c057ab62b..0a686ea30 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -48,6 +48,9 @@ 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. ### Changed diff --git a/src/exec/runner.ts b/src/exec/runner.ts index 2ec6d9a52..87cce5897 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -44,6 +44,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"; @@ -128,28 +129,66 @@ 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`. + * it does not serialize `closeOne`. 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 { + 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 { - args.subAgentSessions?.cancelAll("Session closed"); + const failures: unknown[] = []; + try { + args.subAgentSessions?.cancelAll("Session closed"); + } catch (err) { + failures.push(err); + } if (args.agent !== null) { - await args.agent.close().catch((err: unknown) => { + try { + await args.agent.close(); + } catch (err: unknown) { logger.debug("agent.close during exec finally failed: {error}", { error: formatCaughtError(err), }); - }); + failures.push(err); + } } 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); + } } + rethrowExecDisposeFailures(failures); } /** @@ -288,6 +327,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 +530,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 +838,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 +851,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 +867,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 +887,19 @@ 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) { + logger.debug("disposeExecRuntime failed: {error}", { + error: formatCaughtError(err), + }); + if (result !== undefined && result.exitCode === 0) { + result.exitCode = 1; + } + } + clearActiveDisposeHost(); } } diff --git a/src/index.ts b/src/index.ts index 310f9e21b..a1a44eaa5 100644 --- a/src/index.ts +++ b/src/index.ts @@ -118,6 +118,28 @@ 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; + try { + await Promise.race([ + Promise.resolve(dispose()), + new Promise((_, reject) => { + const 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`, + ); + } +} + // 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 +151,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/session/active-host.test.ts b/src/session/active-host.test.ts index 889fa136f..9948672f7 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 046d7f63d..0a2e4da12 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/tui/runner-exit-code.test.ts b/src/tui/runner-exit-code.test.ts index 8d4860f0c..3d4fc90f1 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.ts b/src/tui/runner/exit.ts index bdca53ad4..4e1091d17 100644 --- a/src/tui/runner/exit.ts +++ b/src/tui/runner/exit.ts @@ -61,11 +61,17 @@ 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; @@ -549,9 +555,17 @@ export async function finalizeTUIRun( services: RunnerServices, ): Promise { await hostOf(state).waitUntilExit(); + await services.sessionOps.awaitTail(); // 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?.(); + // Toolset dispose lives inside shutdownRuntime so quit, crash, and signals + // share one owner. + let teardownFailed = false; + try { + await state.shutdownRuntime?.(); + } catch { + teardownFailed = true; + } 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 +627,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 24f618289..e7d9f1e76 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 3c9e9f443..27a42b1d5 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/wiring.ts b/src/tui/runner/wiring.ts index b8d5e8efe..c60851a39 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -135,9 +135,7 @@ export function wirePostStartup( 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/session-start.test.ts b/src/tui/session-start.test.ts index 30f7f03a2..7122a9e75 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 40aeb713a..ccea72854 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/unit/exec/runner.test.ts b/tests/unit/exec/runner.test.ts index a856fbdcc..a2d2d0776 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,79 @@ 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; + 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(disposeCalls).toBe(1); + expect(getActiveDisposeHost()).toBeNull(); + }, + ); + }, + ); + }, + ); + } finally { + if (previous !== null) setActiveRun(previous); + else clearActiveRun(); + rmSync(cwd, { recursive: true, force: true }); + rmSync(home, { recursive: true, force: true }); + } + }); }); describe("disposeExecRuntime", () => { @@ -318,6 +392,42 @@ describe("disposeExecRuntime", () => { expect(store.get(worker.id)?.status).toBe("cancelled"); expect(calls).toEqual(["agent", "toolset"]); }); + + 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(["agent", "toolset"]); + }); + + 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"); + }); }); describe("resolveExecDirectorOverlay", () => { From e0b04a82c019e50fd6b6c582563cff045c5846b4 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 21:24:36 -0700 Subject: [PATCH 04/14] Reap exec children in shutdown integration tests --- .../exec-shutdown-reap/simulate-reap.ts | 100 ++++++++++++++++++ tests/integration/exec-shutdown-reap.test.ts | 91 ++++++++++++++++ 2 files changed, 191 insertions(+) create mode 100644 tests/fixtures/exec-shutdown-reap/simulate-reap.ts create mode 100644 tests/integration/exec-shutdown-reap.test.ts 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 000000000..a2b06664e --- /dev/null +++ b/tests/fixtures/exec-shutdown-reap/simulate-reap.ts @@ -0,0 +1,100 @@ +// 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 host = (): Promise => + disposeExecRuntime({ + agent: null, + 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 000000000..80016631a --- /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, + ); + }, +); From 3fb80f00c871edaa1c593c96e708404778f49e8c Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 21:25:51 -0700 Subject: [PATCH 05/14] Format shell-guard dispose tests --- src/plugins/shell-guard-plugin.test.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/plugins/shell-guard-plugin.test.ts b/src/plugins/shell-guard-plugin.test.ts index b6d4fdc9b..a98d8ca8c 100644 --- a/src/plugins/shell-guard-plugin.test.ts +++ b/src/plugins/shell-guard-plugin.test.ts @@ -765,7 +765,9 @@ describe("shellGuardPlugin", () => { 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(spawnSync("pgrep", ["-f", token], { encoding: "utf8" }).stdout?.trim() ?? "").not.toBe( + "", + ); expect(plugin.dispose).toBeDefined(); await plugin.dispose!(); await new Promise((r) => setTimeout(r, 300)); From 2cb5d943192fed29a9b5c382e991cea9a4a88200 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 22:03:18 -0700 Subject: [PATCH 06/14] Abort TUI quit before awaiting the session-op tail Stop workers first so a hung session-op cannot delay abort and reap. Log shutdown failures at error while still mapping teardown failure to exit 1. --- src/tui/runner/exit.test.ts | 111 ++++++++++++++++++++++++++++++++++++ src/tui/runner/exit.ts | 12 ++-- 2 files changed, 119 insertions(+), 4 deletions(-) create mode 100644 src/tui/runner/exit.test.ts diff --git a/src/tui/runner/exit.test.ts b/src/tui/runner/exit.test.ts new file mode 100644 index 000000000..e71e7612b --- /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 4e1091d17..2a7fd8d04 100644 --- a/src/tui/runner/exit.ts +++ b/src/tui/runner/exit.ts @@ -555,17 +555,21 @@ export async function finalizeTUIRun( services: RunnerServices, ): Promise { await hostOf(state).waitUntilExit(); - await services.sessionOps.awaitTail(); - // Stop inference and every worker before persistence, hooks, or telemetry can - // delay process exit. Closing the terminal is a process-lifetime boundary. + // 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 { + } 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. From 52a2f25a9a9c04f46ddfd1943c4d7600077ced44 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Tue, 8 Sep 2026 22:03:24 -0700 Subject: [PATCH 07/14] Fail teardown when shell children survive reap A leftover after the two-second backstop must reject dispose so the exit 1 path can fire. Abort already SIGKILLs the process group at abort start. --- src/plugins/shell-guard-plugin.test.ts | 16 +++++++++++-- src/plugins/shell-guard-plugin.ts | 33 +++++++++++++++++++------- 2 files changed, 39 insertions(+), 10 deletions(-) diff --git a/src/plugins/shell-guard-plugin.test.ts b/src/plugins/shell-guard-plugin.test.ts index a98d8ca8c..f28278cdb 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"; @@ -780,4 +781,15 @@ describe("shellGuardPlugin", () => { 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 2c52c477d..1a5fd6f1a 100644 --- a/src/plugins/shell-guard-plugin.ts +++ b/src/plugins/shell-guard-plugin.ts @@ -217,17 +217,34 @@ function waitChildClose(child: ChildProcess): Promise { const SHELL_GUARD_DISPOSE_REAP_MS = 2_000; -async function reapLiveChildren(liveChildren: Set): Promise { +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; - await Promise.race([ - Promise.all(remaining.map(waitChildClose)), - new Promise((resolve) => { - const timer = setTimeout(resolve, SHELL_GUARD_DISPOSE_REAP_MS); - if (typeof timer.unref === "function") timer.unref(); - }), - ]); + 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, From 0465ede3a52b760d0d70cc8a24ebf8a1435d05a6 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 08:50:23 -0700 Subject: [PATCH 08/14] Surface subagent posix dispose failures --- src/subagent/dispose.ts | 6 +----- src/subagent/index.test.ts | 15 +++++++++++++++ src/subagent/run.ts | 6 ------ 3 files changed, 16 insertions(+), 11 deletions(-) diff --git a/src/subagent/dispose.ts b/src/subagent/dispose.ts index fcf9c928e..5061abe09 100644 --- a/src/subagent/dispose.ts +++ b/src/subagent/dispose.ts @@ -120,9 +120,5 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput) } catch { // ignore } - try { - await input.posixTools.dispose(); - } catch { - // LSP shutdown can fail when several sub-agents exit together. - } + await input.posixTools.dispose(); } diff --git a/src/subagent/index.test.ts b/src/subagent/index.test.ts index 2f82b3872..2e708512b 100644 --- a/src/subagent/index.test.ts +++ b/src/subagent/index.test.ts @@ -72,6 +72,21 @@ describe("sub-agent teardown", () => { expect(disposeCount).toBe(2); }); + 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("spawn registry tracks in-flight plugin tool calls", async () => { const { plugin, snapshot } = createSubAgentSpawnRegistryPlugin(); expect(plugin.middleware).toBeDefined(); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index 2794b0ba6..a81324190 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -1106,12 +1106,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) { From d88e375f98da35d5997a82f806f3cd2c0a7df457 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 09:29:31 -0700 Subject: [PATCH 09/14] Fail persist close_agent when child dispose throws A leftover child after posix reap must fail persist close and parent toolset dispose instead of looking like a successful shutdown. A wedged close still times out as shutdown. --- CHANGELOG.md | 2 + src/agent/fleet-verbs-mount.test.ts | 32 +++++++++ src/subagent/run-persist-close.test.ts | 97 ++++++++++++++++++++++++++ src/subagent/run.ts | 11 ++- src/subagent/session-store.test.ts | 13 ++++ src/subagent/session-store.ts | 17 +++-- 6 files changed, 164 insertions(+), 8 deletions(-) create mode 100644 src/subagent/run-persist-close.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 0a686ea30..8acbd8b81 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -51,6 +51,8 @@ parallel copies under `docs/` or `scripts/notes/`. At cut time: rename - 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. ### Changed diff --git a/src/agent/fleet-verbs-mount.test.ts b/src/agent/fleet-verbs-mount.test.ts index b822cf01a..191b99b36 100644 --- a/src/agent/fleet-verbs-mount.test.ts +++ b/src/agent/fleet-verbs-mount.test.ts @@ -53,6 +53,38 @@ 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 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/subagent/run-persist-close.test.ts b/src/subagent/run-persist-close.test.ts new file mode 100644 index 000000000..d7256ae66 --- /dev/null +++ b/src/subagent/run-persist-close.test.ts @@ -0,0 +1,97 @@ +/** + * 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/); + }, + ), + ); + }); +}); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index a81324190..415d54593 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -1123,15 +1123,19 @@ async function runSubAgentInner( 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")); + let disposeError: unknown; 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. - }); + }).then( + () => undefined, + (err: unknown) => { + disposeError = err; + }, + ); await Promise.race([ teardown, new Promise((resolve) => setTimeout(resolve, deadlineMs)), @@ -1140,6 +1144,7 @@ async function runSubAgentInner( // for a persisted session (see runController.dispose's doc); now that // this session is actually closing, tear it down for real. runController.dispose(); + if (disposeError !== undefined) throw disposeError; }; // Interrupt only fires interruptController — never runController/ // close, so it cannot hit the close()-ordering wedge documented in diff --git a/src/subagent/session-store.test.ts b/src/subagent/session-store.test.ts index 918304681..4ea00516d 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -587,6 +587,19 @@ describe("CL-6943 reusable worker sessions", () => { 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 33c021142..32d62b659 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -1226,12 +1226,17 @@ 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. + let closeError: unknown; await Promise.race([ - close(deadlineMs).catch((err: unknown) => { - log.warn("session close raced deadline: {error}", { - error: err instanceof Error ? err.message : String(err), - }); - }), + close(deadlineMs).then( + () => undefined, + (err: unknown) => { + closeError = err; + log.warn("session close raced deadline: {error}", { + error: err instanceof Error ? err.message : String(err), + }); + }, + ), new Promise((resolve) => setTimeout(resolve, deadlineMs)), ]); if (keepFailed) { @@ -1243,6 +1248,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 +1272,7 @@ export function createSubAgentSessionStore( deliverHandles.delete(id); runInFlight.delete(id); pruneCompleted(); + if (closeError !== undefined) throw closeError; return "shutdown"; }, From ad8261906330b5c3e8ab21570a7c1e981dc0b4a9 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 09:55:50 -0700 Subject: [PATCH 10/14] Fail parent dispose when persist workers leave children --- src/agent/fleet-verbs-mount.test.ts | 74 +++++++++++++++++++++++ src/agent/tools.ts | 2 +- src/exec/runner.ts | 9 +-- src/subagent/lifecycle-tools.test.ts | 39 ++++++++++++ src/subagent/lifecycle-tools.ts | 18 +++++- src/subagent/retain-salvage.test.ts | 8 +-- src/subagent/session-store.test.ts | 35 +++++++++++ src/subagent/session-store.ts | 63 ++++++++++++++----- src/tui/runner/exit.ts | 9 ++- src/tui/runner/shutdown.ts | 10 ++- src/tui/runner/wiring.ts | 4 +- src/tui/runtime-shutdown.test.ts | 40 ++++++++++-- tests/unit/subagent-session-store.test.ts | 4 +- 13 files changed, 274 insertions(+), 41 deletions(-) diff --git a/src/agent/fleet-verbs-mount.test.ts b/src/agent/fleet-verbs-mount.test.ts index 191b99b36..578901b9b 100644 --- a/src/agent/fleet-verbs-mount.test.ts +++ b/src/agent/fleet-verbs-mount.test.ts @@ -85,6 +85,80 @@ describe("primary fleet verb mount", () => { 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 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 de2c0787e..168f24fa3 100644 --- a/src/agent/tools.ts +++ b/src/agent/tools.ts @@ -984,7 +984,7 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise { const fleetSessions = fleetSessionsForDispose; if (fleetSessions !== undefined) { - fleetSessions.cancelAll("parent session closed"); + await fleetSessions.cancelAll("parent session closed"); for (const session of [...fleetSessions.list()].reverse()) { await fleetSessions.closeOne(session.id, DEFAULT_CLOSE_DEADLINE_MS); } diff --git a/src/exec/runner.ts b/src/exec/runner.ts index 87cce5897..7eb73369b 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -128,9 +128,10 @@ 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`. Once-only per runtime object so the send - * path, `finally`, and signal host cannot double-dispose. + * the primary agent and dispose the toolset. `cancelAll` is awaited so a + * leftover-child throw is visible; hang-forever close is still deadline-bounded. + * Once-only per runtime object so the send path, `finally`, and signal host + * cannot double-dispose. */ const execDisposeInFlight = new WeakMap>(); @@ -164,7 +165,7 @@ async function runExecDispose(args: { }): Promise { const failures: unknown[] = []; try { - args.subAgentSessions?.cancelAll("Session closed"); + await args.subAgentSessions?.cancelAll("Session closed"); } catch (err) { failures.push(err); } diff --git a/src/subagent/lifecycle-tools.test.ts b/src/subagent/lifecycle-tools.test.ts index 0731bbb4a..e2572af65 100644 --- a/src/subagent/lifecycle-tools.test.ts +++ b/src/subagent/lifecycle-tools.test.ts @@ -90,6 +90,45 @@ describe("close_agent", () => { 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); + }); }); describe("resume_agent", () => { diff --git a/src/subagent/lifecycle-tools.ts b/src/subagent/lifecycle-tools.ts index 4eb6b753a..1ed45eabe 100644 --- a/src/subagent/lifecycle-tools.ts +++ b/src/subagent/lifecycle-tools.ts @@ -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 a960ee8ca..fcb561f11 100644 --- a/src/subagent/retain-salvage.test.ts +++ b/src/subagent/retain-salvage.test.ts @@ -24,7 +24,7 @@ describe("retained session lifecycle", () => { expect(outcome.ok).toBe(false); }); - test("cancelAll does not close retained completed sessions", () => { + test("cancelAll does not close retained completed sessions", async () => { const store = createSubAgentSessionStore({ maxCompleted: 5 }); const s = store.start({ description: "worker", @@ -37,7 +37,7 @@ describe("retained session lifecycle", () => { closed = true; }); store.complete(s.id, "done"); - const cancelled = store.cancelAll("parent stop"); + const cancelled = await store.cancelAll("parent stop"); console.log("cancelAll returned:", cancelled, "| close invoked:", closed); expect(closed).toBe(true); }); @@ -67,7 +67,7 @@ describe("retained session lifecycle", () => { 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 +86,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/session-store.test.ts b/src/subagent/session-store.test.ts index 4ea00516d..a1452abdf 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -574,6 +574,41 @@ describe("CL-6943 reusable worker sessions", () => { expect(store.resumeOne("missing", "more")).toEqual({ ok: false, status: "not_found" }); }); + 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 is bounded by its deadline when the registered close hangs forever", 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 32d62b659..49d31e68c 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -22,6 +22,26 @@ 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 { + let closeError: unknown; + await Promise.race([ + close(deadlineMs).then( + () => undefined, + (err: unknown) => { + closeError = err; + log.warn("session close raced deadline: {error}", { + error: err instanceof Error ? err.message : String(err), + }); + }, + ), + new Promise((resolve) => setTimeout(resolve, deadlineMs)), + ]); + if (closeError !== undefined) throw closeError; +} + export type SubAgentSessionStatus = "running" | "done" | "failed" | "cancelled"; /** @@ -194,8 +214,10 @@ 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 reject, hang-forever + // resolves without throwing. 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; @@ -1227,18 +1249,11 @@ export function createSubAgentSessionStore( // close that does not honor its own deadline argument — a wedged // descendant must not hang the whole close_agent call. let closeError: unknown; - await Promise.race([ - close(deadlineMs).then( - () => undefined, - (err: unknown) => { - closeError = err; - log.warn("session close raced deadline: {error}", { - error: err instanceof Error ? err.message : String(err), - }); - }, - ), - new Promise((resolve) => setTimeout(resolve, deadlineMs)), - ]); + try { + await invokeCloseBounded(close, deadlineMs); + } catch (err: unknown) { + closeError = err; + } if (keepFailed) { // fail() already stamped failed; invoke leftover teardown without // rewriting that to shutdown. @@ -1488,7 +1503,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). @@ -1500,10 +1515,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", @@ -1512,6 +1534,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/tui/runner/exit.ts b/src/tui/runner/exit.ts index 2a7fd8d04..e72e080a7 100644 --- a/src/tui/runner/exit.ts +++ b/src/tui/runner/exit.ts @@ -43,18 +43,20 @@ 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 { @@ -471,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 diff --git a/src/tui/runner/shutdown.ts b/src/tui/runner/shutdown.ts index e74eaeb42..20b2a3598 100644 --- a/src/tui/runner/shutdown.ts +++ b/src/tui/runner/shutdown.ts @@ -1,6 +1,6 @@ export interface RuntimeShutdownDeps { disposeHost: () => void; - cancelWorkers: () => void; + cancelWorkers: () => void | Promise; closeAgent: () => Promise; disposeToolset: () => Promise; } @@ -22,6 +22,7 @@ export function createRuntimeShutdown(deps: RuntimeShutdownDeps): () => Promise< started = true; const failures: unknown[] = []; + let cancelWorkersResult: void | Promise = undefined; try { deps.disposeHost(); @@ -29,12 +30,17 @@ export function createRuntimeShutdown(deps: RuntimeShutdownDeps): () => Promise< failures.push(err); } try { - deps.cancelWorkers(); + cancelWorkersResult = deps.cancelWorkers(); } catch (err) { failures.push(err); } completion = (async () => { + try { + if (cancelWorkersResult !== undefined) await cancelWorkersResult; + } catch (err) { + failures.push(err); + } try { await deps.closeAgent(); } catch (err) { diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index c60851a39..dc7f48d8a 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -128,8 +128,8 @@ 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(), diff --git a/src/tui/runtime-shutdown.test.ts b/src/tui/runtime-shutdown.test.ts index 9f56ceb9e..901d39aee 100644 --- a/src/tui/runtime-shutdown.test.ts +++ b/src/tui/runtime-shutdown.test.ts @@ -7,7 +7,9 @@ describe("runtime shutdown", () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, @@ -25,7 +27,9 @@ describe("runtime shutdown", () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, @@ -46,7 +50,9 @@ describe("runtime shutdown", () => { calls.push("host"); throw new Error("renderer failure"); }, - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, @@ -63,7 +69,9 @@ describe("runtime shutdown", () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, @@ -77,6 +85,26 @@ describe("runtime shutdown", () => { expect(calls).toEqual(["host", "workers", "agent", "toolset"]); }); + 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", "workers", "agent", "toolset"]); + }); + test("awaits an async toolset dispose before resolving", async () => { const calls: string[] = []; let resolveToolset!: () => void; @@ -85,7 +113,9 @@ describe("runtime shutdown", () => { }); const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), - cancelWorkers: () => calls.push("workers"), + cancelWorkers: () => { + calls.push("workers"); + }, closeAgent: async () => { calls.push("agent"); }, diff --git a/tests/unit/subagent-session-store.test.ts b/tests/unit/subagent-session-store.test.ts index b91e09a9b..7f0bf5922 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"); From 90e009e62dbfff8d68610b17fbb593df0cb0d287 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 11:52:20 -0700 Subject: [PATCH 11/14] Fail leftover exec dispose without skipping sibling teardown --- CHANGELOG.md | 3 ++ src/agent/fleet-verbs-mount.test.ts | 52 +++++++++++++++++++++++++++++ src/agent/tools.ts | 27 +++++++++++++-- src/exec/runner.ts | 10 +++--- src/subagent/retain-salvage.test.ts | 9 ++--- tests/unit/exec/runner.test.ts | 24 +++++++++++++ 6 files changed, 111 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8acbd8b81..a1c3a3d9f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -53,6 +53,9 @@ parallel copies under `docs/` or `scripts/notes/`. At cut time: rename 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/src/agent/fleet-verbs-mount.test.ts b/src/agent/fleet-verbs-mount.test.ts index 578901b9b..11ce7112c 100644 --- a/src/agent/fleet-verbs-mount.test.ts +++ b/src/agent/fleet-verbs-mount.test.ts @@ -122,6 +122,58 @@ describe("primary fleet verb mount", () => { 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"); diff --git a/src/agent/tools.ts b/src/agent/tools.ts index 168f24fa3..35a46771e 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,20 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise { + const failures: unknown[] = []; const fleetSessions = fleetSessionsForDispose; if (fleetSessions !== undefined) { - await 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()]); @@ -1000,8 +1016,13 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise { try { await disposeExecRuntime({ agent, toolset, subAgentSessions }); } catch (err: unknown) { - logger.debug("disposeExecRuntime failed: {error}", { - error: formatCaughtError(err), - }); - if (result !== undefined && result.exitCode === 0) { + 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/subagent/retain-salvage.test.ts b/src/subagent/retain-salvage.test.ts index fcb561f11..fc541779a 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", async () => { + 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 = await store.cancelAll("parent stop"); - console.log("cancelAll returned:", cancelled, "| close invoked:", closed); + await store.cancelAll("parent stop"); expect(closed).toBe(true); }); @@ -63,7 +59,6 @@ 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); }); diff --git a/tests/unit/exec/runner.test.ts b/tests/unit/exec/runner.test.ts index a2d2d0776..1ea35a3d6 100644 --- a/tests/unit/exec/runner.test.ts +++ b/tests/unit/exec/runner.test.ts @@ -298,6 +298,12 @@ describe("runExec", () => { 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"), @@ -347,6 +353,9 @@ describe("runExec", () => { 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(); }, @@ -356,6 +365,7 @@ describe("runExec", () => { }, ); } finally { + process.stderr.write = origWrite; if (previous !== null) setActiveRun(previous); else clearActiveRun(); rmSync(cwd, { recursive: true, force: true }); @@ -415,6 +425,20 @@ describe("disposeExecRuntime", () => { expect(calls).toEqual(["agent", "toolset"]); }); + 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("rejects when toolset dispose fails", async () => { await expect( disposeExecRuntime({ From e2634607865b340b19f2ba3b3e60416c73964235 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 12:34:49 -0700 Subject: [PATCH 12/14] Reap posix children before waiting on agent close A hung agent.close used to run before process-group reap, so teardown could report success while detached run_shell children were still live. Dispose first, fail a close deadline instead of succeeding, and clear the two-second host timer when dispose wins. --- src/agent/tools.ts | 16 +++--- src/exec/runner.ts | 30 +++++----- src/index.ts | 5 +- src/subagent/dispose.ts | 40 +++++++++++++- src/subagent/index.test.ts | 27 +++++++++ src/subagent/lifecycle-tools.test.ts | 5 +- src/subagent/lifecycle-tools.ts | 4 +- src/subagent/run-persist-close.test.ts | 55 +++++++++++++++++++ src/subagent/run.ts | 50 ++++++++--------- src/subagent/session-store.test.ts | 5 +- src/subagent/session-store.ts | 29 ++++------ src/subagent/types.ts | 5 +- src/tui/runner/shutdown.ts | 32 ++++------- src/tui/runtime-shutdown.test.ts | 44 ++++++++++++--- .../exec-shutdown-reap/simulate-reap.ts | 6 +- tests/unit/exec/runner.test.ts | 31 ++++++++++- 16 files changed, 272 insertions(+), 112 deletions(-) diff --git a/src/agent/tools.ts b/src/agent/tools.ts index 35a46771e..3e1ccb5bb 100644 --- a/src/agent/tools.ts +++ b/src/agent/tools.ts @@ -990,6 +990,14 @@ 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) { try { @@ -1013,14 +1021,6 @@ 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"); - try { - await posixTools.dispose(); - } catch (err: unknown) { - failures.push(err); - } await disposeWebSearchClients(); rethrowToolsetDisposeFailures(failures); })(); diff --git a/src/exec/runner.ts b/src/exec/runner.ts index 3a7b06fc2..de673a351 100644 --- a/src/exec/runner.ts +++ b/src/exec/runner.ts @@ -127,11 +127,11 @@ export function execUserFailureMessage( } /** - * Headless analogue of TUI `runtime-shutdown`: abort live workers, then close - * the primary agent and dispose the toolset. `cancelAll` is awaited so a - * leftover-child throw is visible; hang-forever close is still deadline-bounded. - * Once-only per runtime object so the send path, `finally`, and signal host - * cannot double-dispose. + * 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. */ const execDisposeInFlight = new WeakMap>(); @@ -164,6 +164,16 @@ async function runExecDispose(args: { subAgentSessions: Pick | null; }): Promise { const failures: unknown[] = []; + if (args.toolset !== null) { + 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) { @@ -179,16 +189,6 @@ async function runExecDispose(args: { failures.push(err); } } - if (args.toolset !== null) { - try { - await args.toolset.dispose(); - } catch (err: unknown) { - logger.debug("toolset.dispose during exec finally failed: {error}", { - error: formatCaughtError(err), - }); - failures.push(err); - } - } rethrowExecDisposeFailures(failures); } diff --git a/src/index.ts b/src/index.ts index a1a44eaa5..a609a636d 100644 --- a/src/index.ts +++ b/src/index.ts @@ -123,11 +123,12 @@ 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) => { - const timer = setTimeout(() => { + timer = setTimeout(() => { reject(new Error(`runtime teardown exceeded ${RUNTIME_TEARDOWN_DEADLINE_MS}ms`)); }, RUNTIME_TEARDOWN_DEADLINE_MS); if (typeof timer.unref === "function") timer.unref(); @@ -137,6 +138,8 @@ async function awaitActiveDisposeHost(context: string): Promise { process.stderr.write( `host dispose failed ${context}: ${disposeErr instanceof Error ? disposeErr.message : String(disposeErr)}\n`, ); + } finally { + if (timer !== undefined) clearTimeout(timer); } } diff --git a/src/subagent/dispose.ts b/src/subagent/dispose.ts index 5061abe09..9ac72b4b6 100644 --- a/src/subagent/dispose.ts +++ b/src/subagent/dispose.ts @@ -45,9 +45,39 @@ export const DEFAULT_CLOSE_DEADLINE_MS = 30_000; * 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). " + + "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`); +} + export interface SubAgentSpawnSnapshot { inFlightToolCalls: number; inFlightByTool: Readonly>; @@ -110,6 +140,12 @@ 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.posixTools.dispose(); + } catch (err: unknown) { + posixError = err; + } try { await input.agent?.close(); } catch { @@ -120,5 +156,5 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput) } catch { // ignore } - await input.posixTools.dispose(); + if (posixError !== undefined) throw posixError; } diff --git a/src/subagent/index.test.ts b/src/subagent/index.test.ts index 2e708512b..17b27f489 100644 --- a/src/subagent/index.test.ts +++ b/src/subagent/index.test.ts @@ -72,6 +72,33 @@ 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 () => { diff --git a/src/subagent/lifecycle-tools.test.ts b/src/subagent/lifecycle-tools.test.ts index e2572af65..57f2d4a8a 100644 --- a/src/subagent/lifecycle-tools.test.ts +++ b/src/subagent/lifecycle-tools.test.ts @@ -86,9 +86,10 @@ 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 () => { diff --git a/src/subagent/lifecycle-tools.ts b/src/subagent/lifecycle-tools.ts index 1ed45eabe..9c92c5690 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", diff --git a/src/subagent/run-persist-close.test.ts b/src/subagent/run-persist-close.test.ts index d7256ae66..3ba3705c5 100644 --- a/src/subagent/run-persist-close.test.ts +++ b/src/subagent/run-persist-close.test.ts @@ -94,4 +94,59 @@ describe("persist close_agent leftover dispose", () => { ), ); }); + + 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); + }, + ), + ); + }); }); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index 415d54593..1a41a7737 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, @@ -1117,38 +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")); - let disposeError: unknown; - const teardown = disposeSubAgentSession({ - signal: runController.signal, - ...(closeOnAbort !== undefined ? { closeOnAbort } : {}), - agent, - ...(streamPromise !== undefined ? { streamPromise } : {}), - posixTools, - }).then( - () => undefined, - (err: unknown) => { - disposeError = err; - }, - ); - 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(); - if (disposeError !== undefined) throw disposeError; + 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 a1452abdf..9facbbf13 100644 --- a/src/subagent/session-store.test.ts +++ b/src/subagent/session-store.test.ts @@ -609,15 +609,14 @@ describe("CL-6943 reusable worker sessions", () => { expect(cancelThrew === false && closeThrew === false && closeStatus === "shutdown").toBe(false); }); - test("closeOne is bounded by its deadline when the registered close hangs forever", async () => { + 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); }); diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index 49d31e68c..7d1dda952 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, @@ -26,20 +26,14 @@ async function invokeCloseBounded( close: (deadlineMs?: number) => Promise, deadlineMs: number, ): Promise { - let closeError: unknown; - await Promise.race([ - close(deadlineMs).then( - () => undefined, - (err: unknown) => { - closeError = err; - log.warn("session close raced deadline: {error}", { - error: err instanceof Error ? err.message : String(err), - }); - }, - ), - new Promise((resolve) => setTimeout(resolve, deadlineMs)), - ]); - if (closeError !== undefined) throw closeError; + 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"; @@ -215,8 +209,9 @@ export interface SubAgentSessionStore { // session was cancelled; false if missing or already terminal. cancel(id: string, reason?: string): boolean; // Cancel every running session. Closes retained workers with the same - // deadline race as closeOne: leftover-child throws reject, hang-forever - // resolves without throwing. Returns the ids that transitioned to cancelled. + // 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. diff --git a/src/subagent/types.ts b/src/subagent/types.ts index 11769ef18..efe7e3935 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/shutdown.ts b/src/tui/runner/shutdown.ts index 20b2a3598..7ccb94a65 100644 --- a/src/tui/runner/shutdown.ts +++ b/src/tui/runner/shutdown.ts @@ -14,40 +14,30 @@ function rethrowShutdownFailures(failures: unknown[]): void { /** 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; - - const failures: unknown[] = []; - let cancelWorkersResult: void | Promise = undefined; - - try { - deps.disposeHost(); - } catch (err) { - failures.push(err); - } - try { - cancelWorkersResult = deps.cancelWorkers(); - } catch (err) { - failures.push(err); - } + if (completion !== undefined) return completion; completion = (async () => { + const failures: unknown[] = []; try { - if (cancelWorkersResult !== undefined) await cancelWorkersResult; + deps.disposeHost(); } catch (err) { failures.push(err); } try { - await deps.closeAgent(); + await deps.disposeToolset(); } catch (err) { failures.push(err); } try { - await deps.disposeToolset(); + await deps.cancelWorkers(); + } catch (err) { + failures.push(err); + } + try { + await deps.closeAgent(); } catch (err) { failures.push(err); } diff --git a/src/tui/runtime-shutdown.test.ts b/src/tui/runtime-shutdown.test.ts index 901d39aee..4489250bd 100644 --- a/src/tui/runtime-shutdown.test.ts +++ b/src/tui/runtime-shutdown.test.ts @@ -3,7 +3,7 @@ import { describe, expect, test } from "bun:test"; import { createRuntimeShutdown } from "./runner/shutdown.js"; describe("runtime shutdown", () => { - test("restores the terminal, cancels workers, closes the primary agent, and disposes the toolset", async () => { + test("restores the terminal, disposes the toolset, then closes the primary agent", async () => { const calls: string[] = []; const shutdown = createRuntimeShutdown({ disposeHost: () => calls.push("host"), @@ -20,7 +20,7 @@ describe("runtime shutdown", () => { await shutdown(); - expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); test("runs teardown only once when exit and a signal race", async () => { @@ -40,7 +40,7 @@ describe("runtime shutdown", () => { await Promise.all([shutdown(), shutdown()]); - expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); test("still runs remaining legs and rejects when host disposal fails", async () => { @@ -62,7 +62,7 @@ describe("runtime shutdown", () => { }); await expect(shutdown()).rejects.toThrow("renderer failure"); - expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); test("rejects when toolset dispose throws after other legs ran", async () => { @@ -82,7 +82,7 @@ describe("runtime shutdown", () => { }); await expect(shutdown()).rejects.toThrow("plugin dispose failed"); - expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); test("rejects when async cancelWorkers throws leftover children", async () => { @@ -102,7 +102,7 @@ describe("runtime shutdown", () => { }); await expect(shutdown()).rejects.toThrow(/still live after 2000ms reap/); - expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + expect(calls).toEqual(["host", "toolset", "workers", "agent"]); }); test("awaits an async toolset dispose before resolving", async () => { @@ -127,9 +127,37 @@ describe("runtime shutdown", () => { const pending = shutdown(); await Promise.resolve(); - expect(calls).toEqual(["host", "workers", "agent"]); + expect(calls).toEqual(["host"]); resolveToolset(); await pending; - expect(calls).toEqual(["host", "workers", "agent", "toolset"]); + 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"]); }); }); diff --git a/tests/fixtures/exec-shutdown-reap/simulate-reap.ts b/tests/fixtures/exec-shutdown-reap/simulate-reap.ts index a2b06664e..559c3c5a9 100644 --- a/tests/fixtures/exec-shutdown-reap/simulate-reap.ts +++ b/tests/fixtures/exec-shutdown-reap/simulate-reap.ts @@ -50,9 +50,13 @@ const countedDispose = async (): Promise => { }; const toolset = { dispose: countedDispose }; +const hungClose = + exitPath === "crash" || exitPath === "signal" + ? { close: () => new Promise(() => undefined) } + : null; const host = (): Promise => disposeExecRuntime({ - agent: null, + agent: hungClose, toolset, subAgentSessions: null, }); diff --git a/tests/unit/exec/runner.test.ts b/tests/unit/exec/runner.test.ts index 1ea35a3d6..d29856ba6 100644 --- a/tests/unit/exec/runner.test.ts +++ b/tests/unit/exec/runner.test.ts @@ -400,7 +400,7 @@ 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 () => { @@ -422,7 +422,34 @@ describe("disposeExecRuntime", () => { await Promise.all([disposeExecRuntime(args), disposeExecRuntime(args)]); - expect(calls).toEqual(["agent", "toolset"]); + 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 () => { From 1c1dc74f7ab03a90a510049f3e86c8e51c928210 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 15:00:18 -0700 Subject: [PATCH 13/14] Surface leftover dispose when agent close hangs --- src/exec/runner.ts | 3 +- src/subagent/dispose.ts | 22 +++++++++- src/subagent/index.test.ts | 32 ++++++++++++++ src/subagent/run-persist-close.test.ts | 58 ++++++++++++++++++++++++++ src/tui/runner/shutdown.ts | 4 +- src/tui/runtime-shutdown.test.ts | 29 +++++++++++++ tests/unit/exec/runner.test.ts | 32 ++++++++++++++ 7 files changed, 176 insertions(+), 4 deletions(-) diff --git a/src/exec/runner.ts b/src/exec/runner.ts index de673a351..5a8f55ae5 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"; @@ -181,7 +182,7 @@ async function runExecDispose(args: { } if (args.agent !== null) { try { - await args.agent.close(); + await awaitCloseWithoutHidingLeftover(args.agent.close(), failures[0]); } catch (err: unknown) { logger.debug("agent.close during exec finally failed: {error}", { error: formatCaughtError(err), diff --git a/src/subagent/dispose.ts b/src/subagent/dispose.ts index 9ac72b4b6..75f4e7eb5 100644 --- a/src/subagent/dispose.ts +++ b/src/subagent/dispose.ts @@ -78,6 +78,24 @@ export async function awaitBoundedTeardown( 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; inFlightByTool: Readonly>; @@ -147,12 +165,12 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput) posixError = err; } try { - await input.agent?.close(); + await awaitCloseWithoutHidingLeftover(input.agent?.close() ?? Promise.resolve(), posixError); } catch { // ignore } try { - await input.streamPromise; + await awaitCloseWithoutHidingLeftover(input.streamPromise ?? Promise.resolve(), posixError); } catch { // ignore } diff --git a/src/subagent/index.test.ts b/src/subagent/index.test.ts index 17b27f489..8c4603e1c 100644 --- a/src/subagent/index.test.ts +++ b/src/subagent/index.test.ts @@ -114,6 +114,38 @@ describe("sub-agent teardown", () => { ).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(); diff --git a/src/subagent/run-persist-close.test.ts b/src/subagent/run-persist-close.test.ts index 3ba3705c5..8816fbf96 100644 --- a/src/subagent/run-persist-close.test.ts +++ b/src/subagent/run-persist-close.test.ts @@ -149,4 +149,62 @@ describe("persist close_agent leftover dispose", () => { ), ); }); + + 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/tui/runner/shutdown.ts b/src/tui/runner/shutdown.ts index 7ccb94a65..eb2c9b0e0 100644 --- a/src/tui/runner/shutdown.ts +++ b/src/tui/runner/shutdown.ts @@ -1,3 +1,5 @@ +import { awaitCloseWithoutHidingLeftover } from "../../subagent/dispose.js"; + export interface RuntimeShutdownDeps { disposeHost: () => void; cancelWorkers: () => void | Promise; @@ -37,7 +39,7 @@ export function createRuntimeShutdown(deps: RuntimeShutdownDeps): () => Promise< failures.push(err); } try { - await deps.closeAgent(); + await awaitCloseWithoutHidingLeftover(deps.closeAgent(), failures[0]); } catch (err) { failures.push(err); } diff --git a/src/tui/runtime-shutdown.test.ts b/src/tui/runtime-shutdown.test.ts index 4489250bd..c55453337 100644 --- a/src/tui/runtime-shutdown.test.ts +++ b/src/tui/runtime-shutdown.test.ts @@ -160,4 +160,33 @@ describe("runtime shutdown", () => { 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/tests/unit/exec/runner.test.ts b/tests/unit/exec/runner.test.ts index d29856ba6..f24400af8 100644 --- a/tests/unit/exec/runner.test.ts +++ b/tests/unit/exec/runner.test.ts @@ -466,6 +466,38 @@ describe("disposeExecRuntime", () => { ).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({ From 9a7454006bc1b5fea167f86912b63c7fe36c3e80 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Wed, 9 Sep 2026 19:38:56 -0700 Subject: [PATCH 14/14] Refuse queued run_shell after shell-guard dispose Latch dispose so overlapping calls join one reap, and refuse queued shells that would spawn after the guard is gone. --- src/plugins/shell-guard-plugin.test.ts | 100 +++++++++++++++++++++++++ src/plugins/shell-guard-plugin.ts | 27 ++++++- 2 files changed, 124 insertions(+), 3 deletions(-) diff --git a/src/plugins/shell-guard-plugin.test.ts b/src/plugins/shell-guard-plugin.test.ts index f28278cdb..c41b2c8e2 100644 --- a/src/plugins/shell-guard-plugin.test.ts +++ b/src/plugins/shell-guard-plugin.test.ts @@ -782,6 +782,106 @@ describe("shellGuardPlugin", () => { } }); + 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, diff --git a/src/plugins/shell-guard-plugin.ts b/src/plugins/shell-guard-plugin.ts index 1a5fd6f1a..95a7da70b 100644 --- a/src/plugins/shell-guard-plugin.ts +++ b/src/plugins/shell-guard-plugin.ts @@ -250,8 +250,12 @@ 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; @@ -395,12 +399,19 @@ export function shellGuardPlugin( 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, @@ -492,6 +503,7 @@ export function shellGuardPlugin( }, signal, liveChildren, + () => disposed, ); const parsed = parsePwdProbeOutput(output); if (perCallCwdRaw === undefined && parsed.finalCwd !== undefined) { @@ -523,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)) { @@ -579,7 +595,12 @@ export function shellGuardPlugin( }, dispose: () => { if (disposal !== undefined) return disposal; - disposal = reapLiveChildren(liveChildren); + disposed = true; + disposal = (async () => { + await reapLiveChildren(liveChildren); + await shellChain; + await reapLiveChildren(liveChildren); + })(); return disposal; }, };