Skip to content

Commit 96ddba5

Browse files
committed
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.
1 parent 4866888 commit 96ddba5

16 files changed

Lines changed: 272 additions & 112 deletions

‎src/agent/tools.ts‎

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -990,6 +990,14 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
990990
mcpAbortController.abort(new Error("MCP toolset disposed"));
991991
disposal = (async () => {
992992
const failures: unknown[] = [];
993+
// Kill every live background process group before the posix teardown so
994+
// /clear, interrupt, and reload cannot leave orphans behind.
995+
backgroundShells.disposeAll("session closed");
996+
try {
997+
await posixTools.dispose();
998+
} catch (err: unknown) {
999+
failures.push(err);
1000+
}
9931001
const fleetSessions = fleetSessionsForDispose;
9941002
if (fleetSessions !== undefined) {
9951003
try {
@@ -1013,14 +1021,6 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
10131021
[...connectedClients.values()].map((client) => client.close().catch(() => undefined)),
10141022
);
10151023
connectedClients.clear();
1016-
// Kill every live background process group before the posix teardown so
1017-
// /clear, interrupt, and reload cannot leave orphans behind.
1018-
backgroundShells.disposeAll("session closed");
1019-
try {
1020-
await posixTools.dispose();
1021-
} catch (err: unknown) {
1022-
failures.push(err);
1023-
}
10241024
await disposeWebSearchClients();
10251025
rethrowToolsetDisposeFailures(failures);
10261026
})();

‎src/exec/runner.ts‎

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -127,11 +127,11 @@ export function execUserFailureMessage(
127127
}
128128

129129
/**
130-
* Headless analogue of TUI `runtime-shutdown`: abort live workers, then close
131-
* the primary agent and dispose the toolset. `cancelAll` is awaited so a
132-
* leftover-child throw is visible; hang-forever close is still deadline-bounded.
133-
* Once-only per runtime object so the send path, `finally`, and signal host
134-
* cannot double-dispose.
130+
* Headless analogue of TUI `runtime-shutdown`: dispose the toolset (posix
131+
* process-group reap) before waiting on agent.close so a hung close cannot
132+
* skip killing detached run_shell children. `cancelAll` is awaited so a
133+
* leftover-child throw is visible. Once-only per runtime object so the send
134+
* path, `finally`, and signal host cannot double-dispose.
135135
*/
136136
const execDisposeInFlight = new WeakMap<object, Promise<void>>();
137137

@@ -164,6 +164,16 @@ async function runExecDispose(args: {
164164
subAgentSessions: Pick<SubAgentSessionStore, "cancelAll"> | null;
165165
}): Promise<void> {
166166
const failures: unknown[] = [];
167+
if (args.toolset !== null) {
168+
try {
169+
await args.toolset.dispose();
170+
} catch (err: unknown) {
171+
logger.debug("toolset.dispose during exec finally failed: {error}", {
172+
error: formatCaughtError(err),
173+
});
174+
failures.push(err);
175+
}
176+
}
167177
try {
168178
await args.subAgentSessions?.cancelAll("Session closed");
169179
} catch (err) {
@@ -179,16 +189,6 @@ async function runExecDispose(args: {
179189
failures.push(err);
180190
}
181191
}
182-
if (args.toolset !== null) {
183-
try {
184-
await args.toolset.dispose();
185-
} catch (err: unknown) {
186-
logger.debug("toolset.dispose during exec finally failed: {error}", {
187-
error: formatCaughtError(err),
188-
});
189-
failures.push(err);
190-
}
191-
}
192192
rethrowExecDisposeFailures(failures);
193193
}
194194

‎src/index.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -123,11 +123,12 @@ export const RUNTIME_TEARDOWN_DEADLINE_MS = 2_000;
123123
async function awaitActiveDisposeHost(context: string): Promise<void> {
124124
const dispose = getActiveDisposeHost();
125125
if (dispose === null) return;
126+
let timer: ReturnType<typeof setTimeout> | undefined;
126127
try {
127128
await Promise.race([
128129
Promise.resolve(dispose()),
129130
new Promise<never>((_, reject) => {
130-
const timer = setTimeout(() => {
131+
timer = setTimeout(() => {
131132
reject(new Error(`runtime teardown exceeded ${RUNTIME_TEARDOWN_DEADLINE_MS}ms`));
132133
}, RUNTIME_TEARDOWN_DEADLINE_MS);
133134
if (typeof timer.unref === "function") timer.unref();
@@ -137,6 +138,8 @@ async function awaitActiveDisposeHost(context: string): Promise<void> {
137138
process.stderr.write(
138139
`host dispose failed ${context}: ${disposeErr instanceof Error ? disposeErr.message : String(disposeErr)}\n`,
139140
);
141+
} finally {
142+
if (timer !== undefined) clearTimeout(timer);
140143
}
141144
}
142145

‎src/subagent/dispose.ts‎

Lines changed: 38 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,9 +45,39 @@ export const DEFAULT_CLOSE_DEADLINE_MS = 30_000;
4545
* tracked in a global registry.
4646
*/
4747
export const SUBAGENT_PLUGIN_SPAWN_TEARDOWN_LIMITS =
48-
"Per sub-agent session Corbits Code runs agent.close(), drains in-flight tool middleware (best-effort), then posixTools.dispose() (LSP and plugin dispose callbacks). " +
48+
"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. " +
4949
"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.";
5050

51+
/** Fail a hung close instead of resolving as successful teardown. */
52+
export async function awaitBoundedTeardown(
53+
teardown: Promise<void>,
54+
deadlineMs: number,
55+
): Promise<void> {
56+
let teardownError: unknown;
57+
let timedOut = false;
58+
let timer: ReturnType<typeof setTimeout> | undefined;
59+
try {
60+
await Promise.race([
61+
teardown.then(
62+
() => undefined,
63+
(err: unknown) => {
64+
teardownError = err;
65+
},
66+
),
67+
new Promise<void>((resolve) => {
68+
timer = setTimeout(() => {
69+
timedOut = true;
70+
resolve();
71+
}, deadlineMs);
72+
}),
73+
]);
74+
} finally {
75+
if (timer !== undefined) clearTimeout(timer);
76+
}
77+
if (teardownError !== undefined) throw teardownError;
78+
if (timedOut) throw new Error(`session close exceeded ${deadlineMs}ms`);
79+
}
80+
5181
export interface SubAgentSpawnSnapshot {
5282
inFlightToolCalls: number;
5383
inFlightByTool: Readonly<Record<string, number>>;
@@ -110,6 +140,12 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput)
110140
if (input.signal !== undefined && input.closeOnAbort !== undefined) {
111141
input.signal.removeEventListener("abort", input.closeOnAbort);
112142
}
143+
let posixError: unknown;
144+
try {
145+
await input.posixTools.dispose();
146+
} catch (err: unknown) {
147+
posixError = err;
148+
}
113149
try {
114150
await input.agent?.close();
115151
} catch {
@@ -120,5 +156,5 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput)
120156
} catch {
121157
// ignore
122158
}
123-
await input.posixTools.dispose();
159+
if (posixError !== undefined) throw posixError;
124160
}

‎src/subagent/index.test.ts‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,33 @@ describe("sub-agent teardown", () => {
7272
expect(disposeCount).toBe(2);
7373
});
7474

75+
test("disposeSubAgentSession reaps posix tools before waiting on agent.close", async () => {
76+
const order: string[] = [];
77+
let releaseClose!: () => void;
78+
const closeGate = new Promise<void>((resolve) => {
79+
releaseClose = resolve;
80+
});
81+
const pending = disposeSubAgentSession({
82+
agent: {
83+
close: async () => {
84+
order.push("close-start");
85+
await closeGate;
86+
order.push("close-end");
87+
},
88+
},
89+
posixTools: {
90+
dispose: async () => {
91+
order.push("posix");
92+
},
93+
},
94+
});
95+
await new Promise((resolve) => setTimeout(resolve, 20));
96+
expect(order).toEqual(["posix", "close-start"]);
97+
releaseClose();
98+
await pending;
99+
expect(order).toEqual(["posix", "close-start", "close-end"]);
100+
});
101+
75102
test("disposeSubAgentSession does not treat a throwing posix dispose as success", async () => {
76103
const posixTools = {
77104
dispose: async () => {

‎src/subagent/lifecycle-tools.test.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -86,9 +86,10 @@ describe("close_agent", () => {
8686
// Exercise the store directly with a short deadline (the tool itself
8787
// uses the real ~30s bound, which would make this test slow).
8888
const started = Date.now();
89-
const childStatus = await sessions.closeOne(wedgedChild.id, 25);
89+
await expect(sessions.closeOne(wedgedChild.id, 25)).rejects.toThrow(
90+
/session close exceeded 25ms/,
91+
);
9092
expect(Date.now() - started).toBeLessThan(500);
91-
expect(childStatus).toBe("shutdown");
9293
});
9394

9495
test("closes remaining siblings after a leftover-child throw, then fails", async () => {

‎src/subagent/lifecycle-tools.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,8 @@ export const closeAgentToolDefinition: ToolDefinition = {
4141
description:
4242
"Permanently close a worker session by agent_id, closing its descendants first. Bounded " +
4343
`by a ~${Math.round(DEFAULT_CLOSE_DEADLINE_MS / 1000)}s cleanup deadline per session so a wedged worker cannot hang ` +
44-
"this call — a session that misses the deadline is still marked shutdown; its teardown just " +
45-
"keeps running in the background. Unblocks any in-flight wait_agents on these ids immediately with " +
44+
"this call — a session that misses the deadline is still marked shutdown and the call fails " +
45+
"instead of reporting success while children may still be live. Unblocks any in-flight wait_agents on these ids immediately with " +
4646
"status 'interrupted'. Closing is permanent: a closed session cannot be resumed.",
4747
inputSchema: {
4848
type: "object",

‎src/subagent/run-persist-close.test.ts‎

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,4 +94,59 @@ describe("persist close_agent leftover dispose", () => {
9494
),
9595
);
9696
});
97+
98+
test("onAgentReady close reaps posix tools before a hung agent.close and fails the deadline", async () => {
99+
const cwd = await mkdtemp(join(tmpdir(), "corbits-persist-close-hung-"));
100+
let posixDisposed = false;
101+
102+
await withMockedModuleDuring(
103+
import.meta.resolve("@intx/tools-posix"),
104+
(real: typeof import("@intx/tools-posix")) => ({
105+
...real,
106+
createPosixTools: (opts: Parameters<typeof real.createPosixTools>[0]) =>
107+
Object.assign(real.createPosixTools(opts), {
108+
dispose: async () => {
109+
posixDisposed = true;
110+
},
111+
}),
112+
}),
113+
async () =>
114+
withMockedModuleDuring(
115+
import.meta.resolve("../agent/live-tool-dispatch.js"),
116+
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
117+
...real,
118+
createAgentWithLiveToolDispatch: async () =>
119+
({
120+
...stubAgent(),
121+
close: () => new Promise<void>(() => {}),
122+
}) as unknown as Awaited<ReturnType<typeof real.createAgentWithLiveToolDispatch>>,
123+
}),
124+
async () => {
125+
const { runSubAgent } = await import("./run.js");
126+
let handles:
127+
| {
128+
close: (deadlineMs?: number) => Promise<void>;
129+
}
130+
| undefined;
131+
const params: RunSubAgentParams = {
132+
cwd,
133+
workdirBase: join(cwd, ".ctx"),
134+
permissionGate,
135+
provider: { providerName: "test", baseURL: "http://localhost", model: "test-model" },
136+
description: "persist close hung close probe",
137+
prompt: "finish the first turn",
138+
persist: true,
139+
onAgentReady: (h) => {
140+
handles = h;
141+
},
142+
};
143+
const result = await runSubAgent(params);
144+
expect(result.agentRetained).toBe(true);
145+
if (handles === undefined) throw new Error("onAgentReady never fired");
146+
await expect(handles.close(50)).rejects.toThrow(/session close exceeded 50ms/);
147+
expect(posixDisposed).toBe(true);
148+
},
149+
),
150+
);
151+
});
97152
});

‎src/subagent/run.ts‎

Lines changed: 23 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -130,6 +130,7 @@ import {
130130
disposeSubAgentSession,
131131
isSubAgentCancelError,
132132
DEFAULT_CLOSE_DEADLINE_MS,
133+
awaitBoundedTeardown,
133134
} from "./dispose.js";
134135
import {
135136
createFleetMailbox,
@@ -1117,38 +1118,33 @@ async function runSubAgentInner(
11171118
// Hand the caller a bounded, idempotent close it can call at any time
11181119
// (close_agent) — independent of whether this run ends up retained.
11191120
// Aborting first stops a still-running turn before tearing down; on an
1120-
// already-finished turn the abort is a no-op. The timeout races teardown
1121-
// itself so a wedged descendant cannot hang the caller — see dispose.ts
1122-
// for the close()-ordering issue that can stall it.
1121+
// already-finished turn the abort is a no-op. posix dispose/reap runs
1122+
// before waiting on agent.close so a wedged close cannot skip killing
1123+
// detached run_shell children. The deadline abandons a hung close and
1124+
// fails rather than reporting success while children may still be live.
11231125
if (params.onAgentReady !== undefined) {
11241126
const boundedClose = async (deadlineMs = DEFAULT_CLOSE_DEADLINE_MS): Promise<void> => {
11251127
if (!runController.signal.aborted) runController.abort(new Error("closed by close_agent"));
1126-
let disposeError: unknown;
1127-
const teardown = disposeSubAgentSession({
1128-
signal: runController.signal,
1129-
...(closeOnAbort !== undefined ? { closeOnAbort } : {}),
1130-
agent,
1131-
...(streamPromise !== undefined ? { streamPromise } : {}),
1132-
posixTools,
1133-
}).then(
1134-
() => undefined,
1135-
(err: unknown) => {
1136-
disposeError = err;
1137-
},
1138-
);
1139-
await Promise.race([
1140-
teardown,
1141-
new Promise<void>((resolve) => setTimeout(resolve, deadlineMs)),
1142-
]);
1143-
// The finally block kept the parent-abort forwarding listener alive
1144-
// for a persisted session (see runController.dispose's doc); now that
1145-
// this session is actually closing, tear it down for real.
1146-
runController.dispose();
1147-
if (disposeError !== undefined) throw disposeError;
1128+
try {
1129+
await awaitBoundedTeardown(
1130+
disposeSubAgentSession({
1131+
signal: runController.signal,
1132+
...(closeOnAbort !== undefined ? { closeOnAbort } : {}),
1133+
agent,
1134+
...(streamPromise !== undefined ? { streamPromise } : {}),
1135+
posixTools,
1136+
}),
1137+
deadlineMs,
1138+
);
1139+
} finally {
1140+
// The finally block kept the parent-abort forwarding listener alive
1141+
// for a persisted session (see runController.dispose's doc); now that
1142+
// this session is actually closing, tear it down for real.
1143+
runController.dispose();
1144+
}
11481145
};
11491146
// Interrupt only fires interruptController — never runController/
1150-
// close, so it cannot hit the close()-ordering wedge documented in
1151-
// dispose.ts.
1147+
// close, so it cannot hang teardown on a wedged agent.close.
11521148
const interrupt = (): void => {
11531149
if (!interruptController.signal.aborted) {
11541150
interruptController.abort(new Error("interrupted by interrupt_agent"));

‎src/subagent/session-store.test.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -609,15 +609,14 @@ describe("CL-6943 reusable worker sessions", () => {
609609
expect(cancelThrew === false && closeThrew === false && closeStatus === "shutdown").toBe(false);
610610
});
611611

612-
test("closeOne is bounded by its deadline when the registered close hangs forever", async () => {
612+
test("closeOne fails a hung close instead of reporting shutdown success", async () => {
613613
const store = createSubAgentSessionStore();
614614
const session = store.start({ description: "d", agentId: "a", brief: "b", retained: true });
615615
store.registerClose(session.id, () => new Promise<void>(() => {})); // never resolves
616616

617617
const started = Date.now();
618-
const status = await store.closeOne(session.id, 25);
618+
await expect(store.closeOne(session.id, 25)).rejects.toThrow(/session close exceeded 25ms/);
619619
expect(Date.now() - started).toBeLessThan(500);
620-
expect(status).toBe("shutdown");
621620
expect(store.get(session.id)?.lifecycleStatus).toBe("shutdown");
622621
expect(store.get(session.id)?.retained).toBe(false);
623622
});

0 commit comments

Comments
 (0)