Skip to content

Commit 1f6c66c

Browse files
committed
Fail parent dispose when persist workers leave children
1 parent 3d2976b commit 1f6c66c

13 files changed

Lines changed: 273 additions & 40 deletions

src/agent/fleet-verbs-mount.test.ts

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,80 @@ describe("primary fleet verb mount", () => {
8585
await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/);
8686
});
8787

88+
test("createAgentToolset dispose rejects when a retained completed persist worker leaves children", async () => {
89+
const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-"));
90+
const { createAgentToolset } = await import("./tools.js");
91+
const permissionGate = {
92+
check: async () => ({ allowed: true }),
93+
getSkipPermissions: () => false,
94+
} as never;
95+
const sessions = createSubAgentSessionStore();
96+
const worker = sessions.start({
97+
description: "d",
98+
agentId: "a",
99+
brief: "b",
100+
retained: true,
101+
});
102+
sessions.registerClose(worker.id, async () => {
103+
throw new Error("1 shell child process still live after 2000ms reap");
104+
});
105+
sessions.complete(worker.id, "done", { agentRetained: true });
106+
107+
const toolset = await createAgentToolset({
108+
cwd,
109+
permissionGate,
110+
onOperatorGate: async () => ({ kind: "option", index: 0 }),
111+
subAgent: {
112+
provider: {
113+
providerName: "test",
114+
baseURL: "http://127.0.0.1:0",
115+
model: "test-model",
116+
},
117+
getWorkdirBase: () => cwd,
118+
sessions,
119+
},
120+
});
121+
122+
await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/);
123+
});
124+
125+
test("createAgentToolset dispose rejects when a retained running persist worker leaves children", async () => {
126+
const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-"));
127+
const { createAgentToolset } = await import("./tools.js");
128+
const permissionGate = {
129+
check: async () => ({ allowed: true }),
130+
getSkipPermissions: () => false,
131+
} as never;
132+
const sessions = createSubAgentSessionStore();
133+
const worker = sessions.start({
134+
description: "d",
135+
agentId: "a",
136+
brief: "b",
137+
retained: true,
138+
});
139+
sessions.markRunning(worker.id);
140+
sessions.registerClose(worker.id, async () => {
141+
throw new Error("1 shell child process still live after 2000ms reap");
142+
});
143+
144+
const toolset = await createAgentToolset({
145+
cwd,
146+
permissionGate,
147+
onOperatorGate: async () => ({ kind: "option", index: 0 }),
148+
subAgent: {
149+
provider: {
150+
providerName: "test",
151+
baseURL: "http://127.0.0.1:0",
152+
model: "test-model",
153+
},
154+
getWorkdirBase: () => cwd,
155+
sessions,
156+
},
157+
});
158+
159+
await expect(toolset.dispose()).rejects.toThrow(/still live after 2000ms reap/);
160+
});
161+
88162
test("createAgentToolset omits fleet verbs when subAgent is not set", async () => {
89163
const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-"));
90164
const { createAgentToolset } = await import("./tools.js");

src/agent/tools.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -945,7 +945,7 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
945945
disposal = (async () => {
946946
const fleetSessions = fleetSessionsForDispose;
947947
if (fleetSessions !== undefined) {
948-
fleetSessions.cancelAll("parent session closed");
948+
await fleetSessions.cancelAll("parent session closed");
949949
for (const session of [...fleetSessions.list()].reverse()) {
950950
await fleetSessions.closeOne(session.id, DEFAULT_CLOSE_DEADLINE_MS);
951951
}

src/exec/runner.ts

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -127,9 +127,10 @@ export function execUserFailureMessage(
127127

128128
/**
129129
* Headless analogue of TUI `runtime-shutdown`: abort live workers, then close
130-
* the primary agent and dispose the toolset. `cancelAll` is fire-and-forget —
131-
* it does not serialize `closeOne`. Once-only per runtime object so the send
132-
* path, `finally`, and signal host cannot double-dispose.
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.
133134
*/
134135
const execDisposeInFlight = new WeakMap<object, Promise<void>>();
135136

@@ -163,7 +164,7 @@ async function runExecDispose(args: {
163164
}): Promise<void> {
164165
const failures: unknown[] = [];
165166
try {
166-
args.subAgentSessions?.cancelAll("Session closed");
167+
await args.subAgentSessions?.cancelAll("Session closed");
167168
} catch (err) {
168169
failures.push(err);
169170
}

src/subagent/lifecycle-tools.test.ts

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,45 @@ describe("close_agent", () => {
9090
expect(Date.now() - started).toBeLessThan(500);
9191
expect(childStatus).toBe("shutdown");
9292
});
93+
94+
test("closes remaining siblings after a leftover-child throw, then fails", async () => {
95+
const sessions = createSubAgentSessionStore();
96+
const parent = sessions.start({ description: "parent", agentId: "a", brief: "b" });
97+
const leftover = sessions.start({
98+
description: "leftover",
99+
agentId: "a",
100+
brief: "b",
101+
parentSessionId: parent.id,
102+
});
103+
const sibling = sessions.start({
104+
description: "sibling",
105+
agentId: "a",
106+
brief: "b",
107+
parentSessionId: parent.id,
108+
});
109+
const closedOrder: string[] = [];
110+
sessions.registerClose(leftover.id, async () => {
111+
closedOrder.push(leftover.id);
112+
throw new Error("1 shell child process still live after 2000ms reap");
113+
});
114+
sessions.registerClose(sibling.id, async () => {
115+
closedOrder.push(sibling.id);
116+
});
117+
sessions.registerClose(parent.id, async () => {
118+
closedOrder.push(parent.id);
119+
});
120+
121+
const closeAgent = createCloseAgentTool({
122+
sessions,
123+
fleetRecords: createFleetMailbox(sessions),
124+
});
125+
await expect(callTool(closeAgent, { target: parent.id })).rejects.toThrow(
126+
/still live after 2000ms reap/,
127+
);
128+
expect(closedOrder).toContain(leftover.id);
129+
expect(closedOrder).toContain(sibling.id);
130+
expect(closedOrder).toContain(parent.id);
131+
});
93132
});
94133

95134
describe("resume_agent", () => {

src/subagent/lifecycle-tools.ts

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -184,14 +184,28 @@ export function createCloseAgentTool(deps: CloseAgentToolDeps): AgentTool {
184184
.map((s) => ({ id: s.id, parentSessionId: s.parentSessionId }));
185185
const order = descendantsClosingOrder(nodes, target);
186186
const closed: { agent_id: string; status: AgentLifecycleStatus }[] = [];
187+
const failures: unknown[] = [];
187188
for (const id of order) {
188189
// Terminalize the wait mailbox before teardown. closeOne flips strip
189190
// status to "cancelled", which kills the soft-interrupt fallback that
190191
// still requires status === "running" — without this, in-flight
191192
// wait_agents hangs until timeout.
192193
deps.fleetRecords.interrupt(id);
193-
const status = await deps.sessions.closeOne(id, DEFAULT_CLOSE_DEADLINE_MS);
194-
closed.push({ agent_id: id, status });
194+
try {
195+
const status = await deps.sessions.closeOne(id, DEFAULT_CLOSE_DEADLINE_MS);
196+
closed.push({ agent_id: id, status });
197+
} catch (err: unknown) {
198+
failures.push(err);
199+
const after = deps.sessions.get(id);
200+
closed.push({
201+
agent_id: id,
202+
status: after === undefined ? "not_found" : after.lifecycleStatus,
203+
});
204+
}
205+
}
206+
if (failures.length === 1) throw failures[0];
207+
if (failures.length > 1) {
208+
throw new AggregateError(failures, "close_agent leftover dispose failed");
195209
}
196210
const own = closed.find((c) => c.agent_id === target);
197211
return lifecycleResult(

src/subagent/retain-salvage.test.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ describe("retained session lifecycle", () => {
2424
expect(outcome.ok).toBe(false);
2525
});
2626

27-
test("cancelAll does not close retained completed sessions", () => {
27+
test("cancelAll does not close retained completed sessions", async () => {
2828
const store = createSubAgentSessionStore({ maxCompleted: 5 });
2929
const s = store.start({
3030
description: "worker",
@@ -37,7 +37,7 @@ describe("retained session lifecycle", () => {
3737
closed = true;
3838
});
3939
store.complete(s.id, "done");
40-
const cancelled = store.cancelAll("parent stop");
40+
const cancelled = await store.cancelAll("parent stop");
4141
console.log("cancelAll returned:", cancelled, "| close invoked:", closed);
4242
expect(closed).toBe(true);
4343
});
@@ -67,7 +67,7 @@ describe("retained session lifecycle", () => {
6767
expect(store.list().length).toBeLessThanOrEqual(3);
6868
});
6969

70-
test("a genuinely retained clean completion IS resumable, and cancelAll releases it", () => {
70+
test("a genuinely retained clean completion IS resumable, and cancelAll releases it", async () => {
7171
const store = createSubAgentSessionStore({ maxCompleted: 5 });
7272
const s = store.start({
7373
description: "worker",
@@ -86,7 +86,7 @@ describe("retained session lifecycle", () => {
8686
store.registerFollowup(s.id, async () => "next");
8787
expect(store.resumeOne(s.id, "more").ok).toBe(true);
8888
expect(closed).toBe(false);
89-
store.cancelAll("parent stop");
89+
await store.cancelAll("parent stop");
9090
expect(closed).toBe(true);
9191
});
9292

src/subagent/session-store.test.ts

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -534,6 +534,41 @@ describe("CL-6943 reusable worker sessions", () => {
534534
expect(store.resumeOne("missing", "more")).toEqual({ ok: false, status: "not_found" });
535535
});
536536

537+
test("cancelAll then closeOne does not swallow a leftover-child throw as shutdown success", async () => {
538+
const store = createSubAgentSessionStore();
539+
const session = store.start({
540+
description: "d",
541+
agentId: "a",
542+
brief: "b",
543+
retained: true,
544+
});
545+
store.registerClose(session.id, async () => {
546+
throw new Error("1 shell child process still live after 2000ms reap");
547+
});
548+
store.complete(session.id, "done", { agentRetained: true });
549+
550+
const leftover = /still live after 2000ms reap/;
551+
let cancelThrew = false;
552+
try {
553+
await store.cancelAll("parent stop");
554+
} catch (err) {
555+
expect(err).toBeInstanceOf(Error);
556+
expect((err as Error).message).toMatch(leftover);
557+
cancelThrew = true;
558+
}
559+
let closeThrew = false;
560+
let closeStatus: string | undefined;
561+
try {
562+
closeStatus = await store.closeOne(session.id, 1000);
563+
} catch (err) {
564+
expect(err).toBeInstanceOf(Error);
565+
expect((err as Error).message).toMatch(leftover);
566+
closeThrew = true;
567+
}
568+
expect(cancelThrew || closeThrew).toBe(true);
569+
expect(cancelThrew === false && closeThrew === false && closeStatus === "shutdown").toBe(false);
570+
});
571+
537572
test("closeOne is bounded by its deadline when the registered close hangs forever", async () => {
538573
const store = createSubAgentSessionStore();
539574
const session = store.start({ description: "d", agentId: "a", brief: "b", retained: true });

src/subagent/session-store.ts

Lines changed: 47 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,26 @@ import type { AdmissionQueue, AdmissionStatus } from "./admission.js";
2222

2323
const log = getLogger([LOG_NAMESPACE_ROOT, "subagent", "session-store"]);
2424

25+
async function invokeCloseBounded(
26+
close: (deadlineMs?: number) => Promise<void>,
27+
deadlineMs: number,
28+
): Promise<void> {
29+
let closeError: unknown;
30+
await Promise.race([
31+
close(deadlineMs).then(
32+
() => undefined,
33+
(err: unknown) => {
34+
closeError = err;
35+
log.warn("session close raced deadline: {error}", {
36+
error: err instanceof Error ? err.message : String(err),
37+
});
38+
},
39+
),
40+
new Promise<void>((resolve) => setTimeout(resolve, deadlineMs)),
41+
]);
42+
if (closeError !== undefined) throw closeError;
43+
}
44+
2545
export type SubAgentSessionStatus = "running" | "done" | "failed" | "cancelled";
2646

2747
/**
@@ -194,8 +214,10 @@ export interface SubAgentSessionStore {
194214
// Abort a running session and mark it cancelled. Returns true when a running
195215
// session was cancelled; false if missing or already terminal.
196216
cancel(id: string, reason?: string): boolean;
197-
// Cancel every running session. Returns the ids that transitioned.
198-
cancelAll(reason?: string): string[];
217+
// Cancel every running session. Closes retained workers with the same
218+
// deadline race as closeOne: leftover-child throws reject, hang-forever
219+
// resolves without throwing. Returns the ids that transitioned to cancelled.
220+
cancelAll(reason?: string): Promise<string[]>;
199221
// CL-6943: flips a "pending_init" session to "running" once its agent
200222
// object actually exists. No-op on an unknown id or one already past init.
201223
markRunning(id: string): void;
@@ -1219,18 +1241,11 @@ export function createSubAgentSessionStore(
12191241
// close that does not honor its own deadline argument — a wedged
12201242
// descendant must not hang the whole close_agent call.
12211243
let closeError: unknown;
1222-
await Promise.race([
1223-
close(deadlineMs).then(
1224-
() => undefined,
1225-
(err: unknown) => {
1226-
closeError = err;
1227-
log.warn("session close raced deadline: {error}", {
1228-
error: err instanceof Error ? err.message : String(err),
1229-
});
1230-
},
1231-
),
1232-
new Promise<void>((resolve) => setTimeout(resolve, deadlineMs)),
1233-
]);
1244+
try {
1245+
await invokeCloseBounded(close, deadlineMs);
1246+
} catch (err: unknown) {
1247+
closeError = err;
1248+
}
12341249
if (keepFailed) {
12351250
// fail() already stamped failed; invoke leftover teardown without
12361251
// rewriting that to shutdown.
@@ -1465,7 +1480,7 @@ export function createSubAgentSessionStore(
14651480
return cancelSession(id, reason);
14661481
},
14671482

1468-
cancelAll(reason = DEFAULT_CANCEL_REASON): string[] {
1483+
async cancelAll(reason = DEFAULT_CANCEL_REASON): Promise<string[]> {
14691484
// Snapshot before cancelSession: markCancelled clears retained, and a
14701485
// resumed retained worker is strip-live so the first loop would otherwise
14711486
// skip the close-handle pass (CL-7001).
@@ -1477,10 +1492,17 @@ export function createSubAgentSessionStore(
14771492
for (const session of running) {
14781493
if (cancelSession(session.id, reason)) cancelled.push(session.id);
14791494
}
1495+
const pendingCloses: Promise<void>[] = [];
14801496
for (const id of retainedIds) {
14811497
const session = sessions.get(id);
14821498
if (session === undefined || session.lifecycle.state === "shutdown") continue;
1483-
releaseHandles(id);
1499+
const close = closeHandles.get(id);
1500+
cancelAskInternal(id, "session handles released");
1501+
closeHandles.delete(id);
1502+
cancelHandles.delete(id);
1503+
interruptHandles.delete(id);
1504+
followupHandles.delete(id);
1505+
deliverHandles.delete(id);
14841506
mutate(id, (s) => {
14851507
s.lifecycle = {
14861508
state: "shutdown",
@@ -1489,6 +1511,15 @@ export function createSubAgentSessionStore(
14891511
};
14901512
s.retained = false;
14911513
});
1514+
if (close !== undefined) {
1515+
pendingCloses.push(invokeCloseBounded(close, DEFAULT_CLOSE_DEADLINE_MS));
1516+
}
1517+
}
1518+
const results = await Promise.allSettled(pendingCloses);
1519+
const failures = results.flatMap((r) => (r.status === "rejected" ? [r.reason] : []));
1520+
if (failures.length === 1) throw failures[0];
1521+
if (failures.length > 1) {
1522+
throw new AggregateError(failures, "session cancelAll close failed");
14921523
}
14931524
return cancelled;
14941525
},

src/tui/runner/exit.ts

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -440,13 +440,16 @@ export async function createRunLifecycle(
440440
// and that path never moved to OpenTUI.
441441
services.emitter.emit("session.clear");
442442
// Cancel live workers before rotation so /clear does not leave orphaned
443-
// child reactors burning tokens under the old session id.
444-
services.subAgentSessions.cancelAll("Session cleared");
443+
// child reactors burning tokens under the old session id. Start immediately
444+
// (do not wait on the op queue) and await on the rotation path so a leftover
445+
// throw is not an unhandled rejection.
446+
const cancelledWorkers = services.subAgentSessions.cancelAll("Session cleared");
445447
// Backend rotation is always enqueued regardless of contention; the queue
446448
// serialises it behind any in-progress op. Sub-agents nest under the new
447449
// session automatically because getWorkdirBase reads the live sessionId.
448450
void enqueueOp(async () => {
449451
try {
452+
await cancelledWorkers;
450453
// Tear the old agent down and dispose the recorder before workdir is
451454
// repointed: the pump can deliver stray deltas until the stream
452455
// settles, and a dead cycle's partial must land in the session that

0 commit comments

Comments
 (0)