Skip to content

Commit ca98cf6

Browse files
committed
Fail parent dispose when persist workers leave children
1 parent c9ead31 commit ca98cf6

13 files changed

Lines changed: 274 additions & 41 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
@@ -984,7 +984,7 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
984984
disposal = (async () => {
985985
const fleetSessions = fleetSessionsForDispose;
986986
if (fleetSessions !== undefined) {
987-
fleetSessions.cancelAll("parent session closed");
987+
await fleetSessions.cancelAll("parent session closed");
988988
for (const session of [...fleetSessions.list()].reverse()) {
989989
await fleetSessions.closeOne(session.id, DEFAULT_CLOSE_DEADLINE_MS);
990990
}

‎src/exec/runner.ts‎

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

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

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

‎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
@@ -574,6 +574,41 @@ describe("CL-6943 reusable worker sessions", () => {
574574
expect(store.resumeOne("missing", "more")).toEqual({ ok: false, status: "not_found" });
575575
});
576576

577+
test("cancelAll then closeOne does not swallow a leftover-child throw as shutdown success", async () => {
578+
const store = createSubAgentSessionStore();
579+
const session = store.start({
580+
description: "d",
581+
agentId: "a",
582+
brief: "b",
583+
retained: true,
584+
});
585+
store.registerClose(session.id, async () => {
586+
throw new Error("1 shell child process still live after 2000ms reap");
587+
});
588+
store.complete(session.id, "done", { agentRetained: true });
589+
590+
const leftover = /still live after 2000ms reap/;
591+
let cancelThrew = false;
592+
try {
593+
await store.cancelAll("parent stop");
594+
} catch (err) {
595+
expect(err).toBeInstanceOf(Error);
596+
expect((err as Error).message).toMatch(leftover);
597+
cancelThrew = true;
598+
}
599+
let closeThrew = false;
600+
let closeStatus: string | undefined;
601+
try {
602+
closeStatus = await store.closeOne(session.id, 1000);
603+
} catch (err) {
604+
expect(err).toBeInstanceOf(Error);
605+
expect((err as Error).message).toMatch(leftover);
606+
closeThrew = true;
607+
}
608+
expect(cancelThrew || closeThrew).toBe(true);
609+
expect(cancelThrew === false && closeThrew === false && closeStatus === "shutdown").toBe(false);
610+
});
611+
577612
test("closeOne is bounded by its deadline when the registered close hangs forever", async () => {
578613
const store = createSubAgentSessionStore();
579614
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;
@@ -1227,18 +1249,11 @@ export function createSubAgentSessionStore(
12271249
// close that does not honor its own deadline argument — a wedged
12281250
// descendant must not hang the whole close_agent call.
12291251
let closeError: unknown;
1230-
await Promise.race([
1231-
close(deadlineMs).then(
1232-
() => undefined,
1233-
(err: unknown) => {
1234-
closeError = err;
1235-
log.warn("session close raced deadline: {error}", {
1236-
error: err instanceof Error ? err.message : String(err),
1237-
});
1238-
},
1239-
),
1240-
new Promise<void>((resolve) => setTimeout(resolve, deadlineMs)),
1241-
]);
1252+
try {
1253+
await invokeCloseBounded(close, deadlineMs);
1254+
} catch (err: unknown) {
1255+
closeError = err;
1256+
}
12421257
if (keepFailed) {
12431258
// fail() already stamped failed; invoke leftover teardown without
12441259
// rewriting that to shutdown.
@@ -1488,7 +1503,7 @@ export function createSubAgentSessionStore(
14881503
return cancelSession(id, reason);
14891504
},
14901505

1491-
cancelAll(reason = DEFAULT_CANCEL_REASON): string[] {
1506+
async cancelAll(reason = DEFAULT_CANCEL_REASON): Promise<string[]> {
14921507
// Snapshot before cancelSession: markCancelled clears retained, and a
14931508
// resumed retained worker is strip-live so the first loop would otherwise
14941509
// skip the close-handle pass (CL-7001).
@@ -1500,10 +1515,17 @@ export function createSubAgentSessionStore(
15001515
for (const session of running) {
15011516
if (cancelSession(session.id, reason)) cancelled.push(session.id);
15021517
}
1518+
const pendingCloses: Promise<void>[] = [];
15031519
for (const id of retainedIds) {
15041520
const session = sessions.get(id);
15051521
if (session === undefined || session.lifecycle.state === "shutdown") continue;
1506-
releaseHandles(id);
1522+
const close = closeHandles.get(id);
1523+
cancelAskInternal(id, "session handles released");
1524+
closeHandles.delete(id);
1525+
cancelHandles.delete(id);
1526+
interruptHandles.delete(id);
1527+
followupHandles.delete(id);
1528+
deliverHandles.delete(id);
15071529
mutate(id, (s) => {
15081530
s.lifecycle = {
15091531
state: "shutdown",
@@ -1512,6 +1534,15 @@ export function createSubAgentSessionStore(
15121534
};
15131535
s.retained = false;
15141536
});
1537+
if (close !== undefined) {
1538+
pendingCloses.push(invokeCloseBounded(close, DEFAULT_CLOSE_DEADLINE_MS));
1539+
}
1540+
}
1541+
const results = await Promise.allSettled(pendingCloses);
1542+
const failures = results.flatMap((r) => (r.status === "rejected" ? [r.reason] : []));
1543+
if (failures.length === 1) throw failures[0];
1544+
if (failures.length > 1) {
1545+
throw new AggregateError(failures, "session cancelAll close failed");
15151546
}
15161547
return cancelled;
15171548
},

0 commit comments

Comments
 (0)