Skip to content

Commit 4381ce5

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 7a2ab30 commit 4381ce5

16 files changed

Lines changed: 269 additions & 109 deletions

src/agent/tools.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -951,6 +951,11 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
951951
mcpAbortController.abort(new Error("MCP toolset disposed"));
952952
disposal = (async () => {
953953
const failures: unknown[] = [];
954+
try {
955+
await posixTools.dispose();
956+
} catch (err: unknown) {
957+
failures.push(err);
958+
}
954959
const fleetSessions = fleetSessionsForDispose;
955960
if (fleetSessions !== undefined) {
956961
try {
@@ -974,11 +979,6 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
974979
[...connectedClients.values()].map((client) => client.close().catch(() => undefined)),
975980
);
976981
connectedClients.clear();
977-
try {
978-
await posixTools.dispose();
979-
} catch (err: unknown) {
980-
failures.push(err);
981-
}
982982
await disposeWebSearchClients();
983983
rethrowToolsetDisposeFailures(failures);
984984
})();

src/exec/runner.ts

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

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

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

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
@@ -122,6 +122,7 @@ import {
122122
disposeSubAgentSession,
123123
isSubAgentCancelError,
124124
DEFAULT_CLOSE_DEADLINE_MS,
125+
awaitBoundedTeardown,
125126
} from "./dispose.js";
126127
import {
127128
createFleetMailbox,
@@ -1089,38 +1090,33 @@ async function runSubAgentInner(
10891090
// Hand the caller a bounded, idempotent close it can call at any time
10901091
// (close_agent) — independent of whether this run ends up retained.
10911092
// Aborting first stops a still-running turn before tearing down; on an
1092-
// already-finished turn the abort is a no-op. The timeout races teardown
1093-
// itself so a wedged descendant cannot hang the caller — see dispose.ts
1094-
// for the close()-ordering issue that can stall it.
1093+
// already-finished turn the abort is a no-op. posix dispose/reap runs
1094+
// before waiting on agent.close so a wedged close cannot skip killing
1095+
// detached run_shell children. The deadline abandons a hung close and
1096+
// fails rather than reporting success while children may still be live.
10951097
if (params.onAgentReady !== undefined) {
10961098
const boundedClose = async (deadlineMs = DEFAULT_CLOSE_DEADLINE_MS): Promise<void> => {
10971099
if (!runController.signal.aborted) runController.abort(new Error("closed by close_agent"));
1098-
let disposeError: unknown;
1099-
const teardown = disposeSubAgentSession({
1100-
signal: runController.signal,
1101-
...(closeOnAbort !== undefined ? { closeOnAbort } : {}),
1102-
agent,
1103-
...(streamPromise !== undefined ? { streamPromise } : {}),
1104-
posixTools,
1105-
}).then(
1106-
() => undefined,
1107-
(err: unknown) => {
1108-
disposeError = err;
1109-
},
1110-
);
1111-
await Promise.race([
1112-
teardown,
1113-
new Promise<void>((resolve) => setTimeout(resolve, deadlineMs)),
1114-
]);
1115-
// The finally block kept the parent-abort forwarding listener alive
1116-
// for a persisted session (see runController.dispose's doc); now that
1117-
// this session is actually closing, tear it down for real.
1118-
runController.dispose();
1119-
if (disposeError !== undefined) throw disposeError;
1100+
try {
1101+
await awaitBoundedTeardown(
1102+
disposeSubAgentSession({
1103+
signal: runController.signal,
1104+
...(closeOnAbort !== undefined ? { closeOnAbort } : {}),
1105+
agent,
1106+
...(streamPromise !== undefined ? { streamPromise } : {}),
1107+
posixTools,
1108+
}),
1109+
deadlineMs,
1110+
);
1111+
} finally {
1112+
// The finally block kept the parent-abort forwarding listener alive
1113+
// for a persisted session (see runController.dispose's doc); now that
1114+
// this session is actually closing, tear it down for real.
1115+
runController.dispose();
1116+
}
11201117
};
11211118
// Interrupt only fires interruptController — never runController/
1122-
// close, so it cannot hit the close()-ordering wedge documented in
1123-
// dispose.ts.
1119+
// close, so it cannot hang teardown on a wedged agent.close.
11241120
const interrupt = (): void => {
11251121
if (!interruptController.signal.aborted) {
11261122
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
@@ -569,15 +569,14 @@ describe("CL-6943 reusable worker sessions", () => {
569569
expect(cancelThrew === false && closeThrew === false && closeStatus === "shutdown").toBe(false);
570570
});
571571

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

577577
const started = Date.now();
578-
const status = await store.closeOne(session.id, 25);
578+
await expect(store.closeOne(session.id, 25)).rejects.toThrow(/session close exceeded 25ms/);
579579
expect(Date.now() - started).toBeLessThan(500);
580-
expect(status).toBe("shutdown");
581580
expect(store.get(session.id)?.lifecycleStatus).toBe("shutdown");
582581
expect(store.get(session.id)?.retained).toBe(false);
583582
});

0 commit comments

Comments
 (0)