Skip to content

Commit d66933c

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 606a86a commit d66933c

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
@@ -586,6 +586,26 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
586586
}
587587
: undefined;
588588

589+
// Aligns with run.ts: (persist && turnSucceeded) || interruptedKeepAlive.
590+
// When true, the worktree stays until close_agent / eviction calls the
591+
// wrapped close below; otherwise the run's finally reclaims it immediately.
592+
let keepWorktreeAlive = false;
593+
const reclaimWorktree = async (): Promise<void> => {
594+
if (worktreeCwd === undefined) return;
595+
const path = worktreeCwd;
596+
worktreeCwd = undefined;
597+
try {
598+
await cleanupSubAgentWorktree(deps.cwd, path, {
599+
stashBaseline: worktreeStashBaseline,
600+
...(worktreeHeadAtCreate !== undefined ? { headAtCreate: worktreeHeadAtCreate } : {}),
601+
});
602+
} catch (err: unknown) {
603+
log.error("spawn_agent worktree cleanup failed: {error}", {
604+
error: err instanceof Error ? err.message : String(err),
605+
});
606+
}
607+
};
608+
589609
const params: RunSubAgentParams = {
590610
// Name the trace directory after the session-store id so the
591611
// descendant-scoping check behind read_agent_trace can resolve this
@@ -630,9 +650,19 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
630650
: {}),
631651
// Keep the session open after a clean completion, and hand the
632652
// store a bounded close for close_agent to call later.
653+
// Worktree cleanup is deferred until that close when the session
654+
// stays alive for followup (agentRetained / interrupt keep-alive) —
655+
// matching run.ts's persisting gate so followup_task does not hit a
656+
// removed cwd.
633657
persist: true,
634658
onAgentReady: ({ close, interrupt, followup, deliver }) => {
635-
deps.sessions.registerClose(session.id, close);
659+
deps.sessions.registerClose(session.id, async (deadlineMs) => {
660+
try {
661+
await close(deadlineMs);
662+
} finally {
663+
await reclaimWorktree();
664+
}
665+
});
636666
deps.sessions.registerInterrupt(session.id, interrupt);
637667
deps.sessions.registerFollowup(session.id, followup);
638668
deps.sessions.registerDeliver(session.id, deliver);
@@ -656,6 +686,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
656686
// "completed" status. Still terminalize fleetRecords so a waiter
657687
// that never saw interrupt_agent (or raced it) cannot hang.
658688
if (result.interrupted === true) {
689+
keepWorktreeAlive = true;
659690
deps.fleetRecords.interrupt(session.id, result.report);
660691
return;
661692
}
@@ -670,8 +701,10 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
670701
// its agent first, so the store must not treat it as resumable
671702
// just because retained:true was requested at spawn.
672703
// complete() no-ops when status is already cancelled.
704+
const agentRetained = result.agentRetained === true;
705+
if (agentRetained) keepWorktreeAlive = true;
673706
deps.sessions.complete(session.id, result.report, {
674-
agentRetained: result.agentRetained === true,
707+
agentRetained,
675708
...(result.stopReason !== undefined ? { stopReason: result.stopReason } : {}),
676709
});
677710
})
@@ -689,15 +722,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
689722
status: deps.sessions.get(session.id)?.status ?? "completed",
690723
duration_ms: Date.now() - startedAt,
691724
});
692-
if (worktreeCwd === undefined) return;
693-
void cleanupSubAgentWorktree(deps.cwd, worktreeCwd, {
694-
stashBaseline: worktreeStashBaseline,
695-
...(worktreeHeadAtCreate !== undefined ? { headAtCreate: worktreeHeadAtCreate } : {}),
696-
}).catch((err: unknown) => {
697-
log.error("spawn_agent worktree cleanup failed: {error}", {
698-
error: err instanceof Error ? err.message : String(err),
699-
});
700-
});
725+
if (!keepWorktreeAlive) void reclaimWorktree();
701726
});
702727

703728
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)