Skip to content

Commit 05ad41f

Browse files
committed
Defer spawn_agent worktree cleanup while the session is retained
A retained or interrupt_agent keep-alive session still needs its cwd for followup_task. Cleaning the worktree in the run's finally removed that path too early; reclaim it from the close handle instead, matching run.ts's persisting gate.
1 parent 85c7716 commit 05ad41f

2 files changed

Lines changed: 199 additions & 13 deletions

File tree

src/subagent/agent-fleet.ts

Lines changed: 36 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -592,6 +592,26 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
592592
}
593593
: undefined;
594594

595+
// Aligns with run.ts: (persist && turnSucceeded) || interruptedKeepAlive.
596+
// When true, the worktree stays until close_agent / eviction calls the
597+
// wrapped close below; otherwise the run's finally reclaims it immediately.
598+
let keepWorktreeAlive = false;
599+
const reclaimWorktree = async (): Promise<void> => {
600+
if (worktreeCwd === undefined) return;
601+
const path = worktreeCwd;
602+
worktreeCwd = undefined;
603+
try {
604+
await cleanupSubAgentWorktree(deps.cwd, path, {
605+
stashBaseline: worktreeStashBaseline,
606+
...(worktreeHeadAtCreate !== undefined ? { headAtCreate: worktreeHeadAtCreate } : {}),
607+
});
608+
} catch (err: unknown) {
609+
log.error("spawn_agent worktree cleanup failed: {error}", {
610+
error: err instanceof Error ? err.message : String(err),
611+
});
612+
}
613+
};
614+
595615
const params: RunSubAgentParams = {
596616
// Name the trace directory after the session-store id so the
597617
// descendant-scoping check behind read_agent_trace can resolve this
@@ -636,9 +656,19 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
636656
: {}),
637657
// Keep the session open after a clean completion, and hand the
638658
// store a bounded close for close_agent to call later.
659+
// Worktree cleanup is deferred until that close when the session
660+
// stays alive for followup (agentRetained / interrupt keep-alive) —
661+
// matching run.ts's persisting gate so followup_task does not hit a
662+
// removed cwd.
639663
persist: true,
640664
onAgentReady: ({ close, interrupt, followup, deliver }) => {
641-
deps.sessions.registerClose(session.id, close);
665+
deps.sessions.registerClose(session.id, async (deadlineMs) => {
666+
try {
667+
await close(deadlineMs);
668+
} finally {
669+
await reclaimWorktree();
670+
}
671+
});
642672
deps.sessions.registerInterrupt(session.id, interrupt);
643673
deps.sessions.registerFollowup(session.id, followup);
644674
deps.sessions.registerDeliver(session.id, deliver);
@@ -662,6 +692,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
662692
// "completed" status. Still terminalize fleetRecords so a waiter
663693
// that never saw interrupt_agent (or raced it) cannot hang.
664694
if (result.interrupted === true) {
695+
keepWorktreeAlive = true;
665696
deps.fleetRecords.interrupt(session.id, result.report);
666697
return;
667698
}
@@ -676,8 +707,10 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
676707
// its agent first, so the store must not treat it as resumable
677708
// just because retained:true was requested at spawn.
678709
// complete() no-ops when status is already cancelled.
710+
const agentRetained = result.agentRetained === true;
711+
if (agentRetained) keepWorktreeAlive = true;
679712
deps.sessions.complete(session.id, result.report, {
680-
agentRetained: result.agentRetained === true,
713+
agentRetained,
681714
...(result.stopReason !== undefined ? { stopReason: result.stopReason } : {}),
682715
});
683716
})
@@ -695,15 +728,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
695728
status: deps.sessions.get(session.id)?.status ?? "completed",
696729
duration_ms: Date.now() - startedAt,
697730
});
698-
if (worktreeCwd === undefined) return;
699-
void cleanupSubAgentWorktree(deps.cwd, worktreeCwd, {
700-
stashBaseline: worktreeStashBaseline,
701-
...(worktreeHeadAtCreate !== undefined ? { headAtCreate: worktreeHeadAtCreate } : {}),
702-
}).catch((err: unknown) => {
703-
log.error("spawn_agent worktree cleanup failed: {error}", {
704-
error: err instanceof Error ? err.message : String(err),
705-
});
706-
});
731+
if (!keepWorktreeAlive) void reclaimWorktree();
707732
});
708733

709734
return fleetResult(call.id, JSON.stringify({ agent_id: session.id, status: "running" }));

src/subagent/spawn-agent-worktree.test.ts

Lines changed: 163 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,14 @@
11
import { afterEach, describe, expect, test } from "bun:test";
22
import { execFile } from "node:child_process";
3-
import { mkdtemp, rm, writeFile } from "node:fs/promises";
3+
import { access, mkdtemp, rm, writeFile } from "node:fs/promises";
44
import { tmpdir } from "node:os";
55
import { join } from "node:path";
66
import { promisify } from "node:util";
77

88
import { createFleetRecords, createSpawnAgentTool } from "./agent-fleet.js";
99
import { createSubAgentSessionStore } from "./session-store.js";
1010
import { createPermissionGate } from "../permission/gate.js";
11-
import type { RunSubAgentParams } from "./types.js";
11+
import type { RunSubAgentParams, RunSubAgentResult } from "./types.js";
1212

1313
const run = promisify(execFile);
1414

@@ -44,6 +44,26 @@ async function makeRepo(): Promise<string> {
4444
return dir;
4545
}
4646

47+
async function pathExists(path: string): Promise<boolean> {
48+
try {
49+
await access(path);
50+
return true;
51+
} catch {
52+
return false;
53+
}
54+
}
55+
56+
function deferred<T>(): {
57+
promise: Promise<T>;
58+
resolve: (v: T) => void;
59+
} {
60+
let resolve!: (v: T) => void;
61+
const promise = new Promise<T>((res) => {
62+
resolve = res;
63+
});
64+
return { promise, resolve };
65+
}
66+
4767
describe("spawn_agent worktree isolation", () => {
4868
test("propagates a fresh worktree path as the worker cwd", async () => {
4969
const repo = await makeRepo();
@@ -113,4 +133,145 @@ describe("spawn_agent worktree isolation", () => {
113133
expect(result.isError).toBe(true);
114134
expect(ran).toBe(false);
115135
});
136+
137+
test("defers worktree cleanup while the session is retained for followup", async () => {
138+
const repo = await makeRepo();
139+
tempDirs.push(repo);
140+
const workdirBase = await mkdtemp(join(tmpdir(), "corbits-workdir-"));
141+
tempDirs.push(workdirBase);
142+
143+
const settle = deferred<RunSubAgentResult>();
144+
let workerCwd: string | undefined;
145+
const sessions = createSubAgentSessionStore();
146+
const tool = createSpawnAgentTool({
147+
permissionGate: testPermissionGate,
148+
cwd: repo,
149+
getWorkdirBase: () => workdirBase,
150+
provider,
151+
useWorktree: true,
152+
run: async (params) => {
153+
workerCwd = params.cwd;
154+
params.onAgentReady?.({
155+
close: async () => {},
156+
interrupt: () => {},
157+
followup: async () => "",
158+
deliver: () => {},
159+
});
160+
return settle.promise;
161+
},
162+
sessions,
163+
fleetRecords: createFleetRecords(),
164+
});
165+
if (tool.kind !== "full") throw new Error("expected full tool");
166+
const spawned = await tool.handler(
167+
{
168+
id: "retain-wt",
169+
name: "spawn_agent",
170+
arguments: { description: "keep alive", prompt: "Do the work", intent: "explore" },
171+
},
172+
new AbortController().signal,
173+
);
174+
const content = typeof spawned.content === "string" ? spawned.content : "";
175+
const agentId = (JSON.parse(content) as { agent_id: string }).agent_id;
176+
177+
settle.resolve({ report: "## Summary\nDone.", agentRetained: true });
178+
await new Promise((resolve) => setTimeout(resolve, 50));
179+
180+
expect(workerCwd).toBeDefined();
181+
expect(await pathExists(workerCwd!)).toBe(true);
182+
183+
await sessions.closeOne(agentId, 1000);
184+
await new Promise((resolve) => setTimeout(resolve, 50));
185+
expect(await pathExists(workerCwd!)).toBe(false);
186+
});
187+
188+
test("defers worktree cleanup while the session is interrupted for followup", async () => {
189+
const repo = await makeRepo();
190+
tempDirs.push(repo);
191+
const workdirBase = await mkdtemp(join(tmpdir(), "corbits-workdir-"));
192+
tempDirs.push(workdirBase);
193+
194+
const settle = deferred<RunSubAgentResult>();
195+
let workerCwd: string | undefined;
196+
const sessions = createSubAgentSessionStore();
197+
const tool = createSpawnAgentTool({
198+
permissionGate: testPermissionGate,
199+
cwd: repo,
200+
getWorkdirBase: () => workdirBase,
201+
provider,
202+
useWorktree: true,
203+
run: async (params) => {
204+
workerCwd = params.cwd;
205+
params.onAgentReady?.({
206+
close: async () => {},
207+
interrupt: () => {},
208+
followup: async () => "",
209+
deliver: () => {},
210+
});
211+
return settle.promise;
212+
},
213+
sessions,
214+
fleetRecords: createFleetRecords(),
215+
});
216+
if (tool.kind !== "full") throw new Error("expected full tool");
217+
const spawned = await tool.handler(
218+
{
219+
id: "interrupt-wt",
220+
name: "spawn_agent",
221+
arguments: { description: "interrupt me", prompt: "Do the work", intent: "explore" },
222+
},
223+
new AbortController().signal,
224+
);
225+
const content = typeof spawned.content === "string" ? spawned.content : "";
226+
const agentId = (JSON.parse(content) as { agent_id: string }).agent_id;
227+
228+
settle.resolve({
229+
report: "## Summary\nStopped.\n## Findings\npartial\n## Blockers\ninterrupted\n## Paths\n",
230+
interrupted: true,
231+
});
232+
await new Promise((resolve) => setTimeout(resolve, 50));
233+
234+
expect(workerCwd).toBeDefined();
235+
expect(await pathExists(workerCwd!)).toBe(true);
236+
237+
await sessions.closeOne(agentId, 1000);
238+
await new Promise((resolve) => setTimeout(resolve, 50));
239+
expect(await pathExists(workerCwd!)).toBe(false);
240+
});
241+
242+
test("reclaims the worktree immediately when the agent is not retained", async () => {
243+
const repo = await makeRepo();
244+
tempDirs.push(repo);
245+
const workdirBase = await mkdtemp(join(tmpdir(), "corbits-workdir-"));
246+
tempDirs.push(workdirBase);
247+
248+
let workerCwd: string | undefined;
249+
const tool = createSpawnAgentTool({
250+
permissionGate: testPermissionGate,
251+
cwd: repo,
252+
getWorkdirBase: () => workdirBase,
253+
provider,
254+
useWorktree: true,
255+
run: async (params) => {
256+
workerCwd = params.cwd;
257+
// Salvage / non-persist path: no agentRetained flag.
258+
return { report: "## Summary\nSalvaged." };
259+
},
260+
sessions: createSubAgentSessionStore(),
261+
fleetRecords: createFleetRecords(),
262+
});
263+
if (tool.kind !== "full") throw new Error("expected full tool");
264+
await tool.handler(
265+
{
266+
id: "no-retain-wt",
267+
name: "spawn_agent",
268+
arguments: { description: "one shot", prompt: "Do the work", intent: "explore" },
269+
},
270+
new AbortController().signal,
271+
);
272+
await new Promise((resolve) => setTimeout(resolve, 50));
273+
274+
expect(workerCwd).toBeDefined();
275+
expect(await pathExists(workerCwd!)).toBe(false);
276+
});
116277
});

0 commit comments

Comments
 (0)