Skip to content

Commit 75df5a7

Browse files
committed
Surface leftover dispose when agent close hangs
1 parent 96ddba5 commit 75df5a7

7 files changed

Lines changed: 176 additions & 4 deletions

File tree

src/exec/runner.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import {
2222
type SubAgentSessionStore,
2323
} from "../subagent/index.js";
2424
import { getProcessAdmissionQueue } from "../subagent/admission.js";
25+
import { awaitCloseWithoutHidingLeftover } from "../subagent/dispose.js";
2526
import type { ContextStore, InferenceSource, InboundMessage } from "@intx/types/runtime";
2627
import { OPERATOR_ORIGINATED_FLAG } from "../agent/message-provenance.js";
2728
import { loadAgentProfiles } from "../agent/profiles.js";
@@ -181,7 +182,7 @@ async function runExecDispose(args: {
181182
}
182183
if (args.agent !== null) {
183184
try {
184-
await args.agent.close();
185+
await awaitCloseWithoutHidingLeftover(args.agent.close(), failures[0]);
185186
} catch (err: unknown) {
186187
logger.debug("agent.close during exec finally failed: {error}", {
187188
error: formatCaughtError(err),

src/subagent/dispose.ts

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,24 @@ export async function awaitBoundedTeardown(
7878
if (timedOut) throw new Error(`session close exceeded ${deadlineMs}ms`);
7979
}
8080

81+
/**
82+
* Always start `close`. Await it only when posix/toolset dispose already
83+
* succeeded; a leftover throw must not wait unbounded on a hung close.
84+
*/
85+
export async function awaitCloseWithoutHidingLeftover(
86+
close: Promise<unknown>,
87+
leftover: unknown,
88+
): Promise<void> {
89+
if (leftover === undefined) {
90+
await close;
91+
return;
92+
}
93+
void close.then(
94+
() => undefined,
95+
() => undefined,
96+
);
97+
}
98+
8199
export interface SubAgentSpawnSnapshot {
82100
inFlightToolCalls: number;
83101
inFlightByTool: Readonly<Record<string, number>>;
@@ -147,12 +165,12 @@ export async function disposeSubAgentSession(input: SubAgentSessionDisposeInput)
147165
posixError = err;
148166
}
149167
try {
150-
await input.agent?.close();
168+
await awaitCloseWithoutHidingLeftover(input.agent?.close() ?? Promise.resolve(), posixError);
151169
} catch {
152170
// ignore
153171
}
154172
try {
155-
await input.streamPromise;
173+
await awaitCloseWithoutHidingLeftover(input.streamPromise ?? Promise.resolve(), posixError);
156174
} catch {
157175
// ignore
158176
}

src/subagent/index.test.ts

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -114,6 +114,38 @@ describe("sub-agent teardown", () => {
114114
).rejects.toThrow(/still live after 2000ms reap/);
115115
});
116116

117+
test("disposeSubAgentSession surfaces leftover posix dispose when agent.close hangs", async () => {
118+
const posixTools = {
119+
dispose: async () => {
120+
throw new Error("1 shell child process still live after 2000ms reap");
121+
},
122+
};
123+
let closeStarted = false;
124+
const pending = disposeSubAgentSession({
125+
agent: {
126+
close: () => {
127+
closeStarted = true;
128+
return new Promise<void>(() => {});
129+
},
130+
},
131+
posixTools,
132+
});
133+
const result = await Promise.race([
134+
pending.then(
135+
() => ({ kind: "resolved" as const }),
136+
(err: unknown) => ({ kind: "rejected" as const, err }),
137+
),
138+
new Promise<{ kind: "timeout" }>((resolve) => {
139+
setTimeout(() => resolve({ kind: "timeout" }), 200);
140+
}),
141+
]);
142+
expect(closeStarted).toBe(true);
143+
expect(result.kind).toBe("rejected");
144+
if (result.kind !== "rejected") throw new Error("expected leftover reject");
145+
expect(result.err).toBeInstanceOf(Error);
146+
expect((result.err as Error).message).toMatch(/still live after 2000ms reap/);
147+
});
148+
117149
test("spawn registry tracks in-flight plugin tool calls", async () => {
118150
const { plugin, snapshot } = createSubAgentSpawnRegistryPlugin();
119151
expect(plugin.middleware).toBeDefined();

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

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,4 +149,62 @@ describe("persist close_agent leftover dispose", () => {
149149
),
150150
);
151151
});
152+
153+
test("onAgentReady close surfaces leftover posix dispose when agent.close hangs", async () => {
154+
const cwd = await mkdtemp(join(tmpdir(), "corbits-persist-close-leftover-hang-"));
155+
let closeStarted = false;
156+
157+
await withMockedModuleDuring(
158+
import.meta.resolve("@intx/tools-posix"),
159+
(real: typeof import("@intx/tools-posix")) => ({
160+
...real,
161+
createPosixTools: (opts: Parameters<typeof real.createPosixTools>[0]) =>
162+
Object.assign(real.createPosixTools(opts), {
163+
dispose: async () => {
164+
throw new Error("1 shell child process still live after 2000ms reap");
165+
},
166+
}),
167+
}),
168+
async () =>
169+
withMockedModuleDuring(
170+
import.meta.resolve("../agent/live-tool-dispatch.js"),
171+
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
172+
...real,
173+
createAgentWithLiveToolDispatch: async () =>
174+
({
175+
...stubAgent(),
176+
close: () => {
177+
closeStarted = true;
178+
return new Promise<void>(() => {});
179+
},
180+
}) as unknown as Awaited<ReturnType<typeof real.createAgentWithLiveToolDispatch>>,
181+
}),
182+
async () => {
183+
const { runSubAgent } = await import("./run.js");
184+
let handles:
185+
| {
186+
close: (deadlineMs?: number) => Promise<void>;
187+
}
188+
| undefined;
189+
const params: RunSubAgentParams = {
190+
cwd,
191+
workdirBase: join(cwd, ".ctx"),
192+
permissionGate,
193+
provider: { providerName: "test", baseURL: "http://localhost", model: "test-model" },
194+
description: "persist close leftover hung close probe",
195+
prompt: "finish the first turn",
196+
persist: true,
197+
onAgentReady: (h) => {
198+
handles = h;
199+
},
200+
};
201+
const result = await runSubAgent(params);
202+
expect(result.agentRetained).toBe(true);
203+
if (handles === undefined) throw new Error("onAgentReady never fired");
204+
await expect(handles.close(200)).rejects.toThrow(/still live after 2000ms reap/);
205+
expect(closeStarted).toBe(true);
206+
},
207+
),
208+
);
209+
});
152210
});

src/tui/runner/shutdown.ts

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import { awaitCloseWithoutHidingLeftover } from "../../subagent/dispose.js";
2+
13
export interface RuntimeShutdownDeps {
24
disposeHost: () => void;
35
cancelWorkers: () => void | Promise<void>;
@@ -37,7 +39,7 @@ export function createRuntimeShutdown(deps: RuntimeShutdownDeps): () => Promise<
3739
failures.push(err);
3840
}
3941
try {
40-
await deps.closeAgent();
42+
await awaitCloseWithoutHidingLeftover(deps.closeAgent(), failures[0]);
4143
} catch (err) {
4244
failures.push(err);
4345
}

src/tui/runtime-shutdown.test.ts

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -160,4 +160,33 @@ describe("runtime shutdown", () => {
160160
await pending;
161161
expect(calls).toEqual(["host", "toolset", "workers", "agent"]);
162162
});
163+
164+
test("surfaces leftover toolset dispose when agent.close hangs", async () => {
165+
let closeStarted = false;
166+
const shutdown = createRuntimeShutdown({
167+
disposeHost: () => undefined,
168+
cancelWorkers: () => undefined,
169+
closeAgent: () => {
170+
closeStarted = true;
171+
return new Promise<void>(() => {});
172+
},
173+
disposeToolset: async () => {
174+
throw new Error("1 shell child process still live after 2000ms reap");
175+
},
176+
});
177+
const result = await Promise.race([
178+
shutdown().then(
179+
() => ({ kind: "resolved" as const }),
180+
(err: unknown) => ({ kind: "rejected" as const, err }),
181+
),
182+
new Promise<{ kind: "timeout" }>((resolve) => {
183+
setTimeout(() => resolve({ kind: "timeout" }), 200);
184+
}),
185+
]);
186+
expect(closeStarted).toBe(true);
187+
expect(result.kind).toBe("rejected");
188+
if (result.kind !== "rejected") throw new Error("expected leftover reject");
189+
expect(result.err).toBeInstanceOf(Error);
190+
expect((result.err as Error).message).toMatch(/still live after 2000ms reap/);
191+
});
163192
});

tests/unit/exec/runner.test.ts

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -466,6 +466,38 @@ describe("disposeExecRuntime", () => {
466466
).rejects.toThrow(/still live after 2000ms reap/);
467467
});
468468

469+
test("surfaces leftover toolset dispose when agent.close hangs", async () => {
470+
let closeStarted = false;
471+
const pending = disposeExecRuntime({
472+
agent: {
473+
close: () => {
474+
closeStarted = true;
475+
return new Promise<void>(() => {});
476+
},
477+
},
478+
toolset: {
479+
dispose: async () => {
480+
throw new Error("1 shell child process still live after 2000ms reap");
481+
},
482+
},
483+
subAgentSessions: null,
484+
});
485+
const result = await Promise.race([
486+
pending.then(
487+
() => ({ kind: "resolved" as const }),
488+
(err: unknown) => ({ kind: "rejected" as const, err }),
489+
),
490+
new Promise<{ kind: "timeout" }>((resolve) => {
491+
setTimeout(() => resolve({ kind: "timeout" }), 200);
492+
}),
493+
]);
494+
expect(closeStarted).toBe(true);
495+
expect(result.kind).toBe("rejected");
496+
if (result.kind !== "rejected") throw new Error("expected leftover reject");
497+
expect(result.err).toBeInstanceOf(Error);
498+
expect((result.err as Error).message).toMatch(/still live after 2000ms reap/);
499+
});
500+
469501
test("rejects when toolset dispose fails", async () => {
470502
await expect(
471503
disposeExecRuntime({

0 commit comments

Comments
 (0)