Skip to content

Commit aed3ec0

Browse files
committed
Report cancelled and interrupted workers truthfully
1 parent 1760102 commit aed3ec0

6 files changed

Lines changed: 226 additions & 6 deletions

File tree

docs/TELEMETRY.md

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,10 @@ a truthy value (`1`, `true`, …) to restore per-call spans for debugging.
109109
Leaf `runSubAgent` workers do not emit `$ai_*`; worker rollups travel on
110110
`subagent_end` instead. Both TUI and exec install the same turn observer, so a
111111
worker ending during an active parent turn carries that turn's `parent_trace_id`.
112+
Pre-progress operator aborts settle with `status=cancelled` and
113+
`stop_reason=cancelled` even when the worker promise rejects. An interrupt that
114+
keeps a worker resumable settles with `status=interrupted` and the same
115+
`stop_reason=cancelled`; terminal events never report a still-running status.
112116

113117
A deterministic representative fixture uses 10 parent turns with 80 parent tool
114118
calls and 4 workers totaling 24 turns and 96 tool calls. The former per-call and

src/subagent/agent-fleet.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -549,9 +549,16 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
549549
const finalizeEnd = (setupFailed = false): void => {
550550
if (endFinalized) return;
551551
endFinalized = true;
552+
const terminalSession = deps.sessions.get(session.id);
553+
const status =
554+
terminalSession?.status === "cancelled"
555+
? "cancelled"
556+
: terminalSession?.lifecycleStatus === "interrupted"
557+
? "interrupted"
558+
: (terminalSession?.status ?? "completed");
552559
captureSubagentEnd(telemetry, {
553560
agentName,
554-
status: setupFailed ? "failed" : (deps.sessions.get(session.id)?.status ?? "completed"),
561+
status: setupFailed ? "failed" : status,
555562
durationMs: Date.now() - startedAt,
556563
model: settlement?.model ?? provider.model,
557564
stopReason: setupFailed ? "setup_error" : (settlement?.terminal_reason ?? "error"),
@@ -731,6 +738,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
731738
// that never saw interrupt_agent (or raced it) cannot hang.
732739
if (result.interrupted === true) {
733740
keepWorktreeAlive = true;
741+
deps.sessions.interruptOne(session.id);
734742
deps.fleetRecords.interrupt(session.id, result.report);
735743
return;
736744
}

src/subagent/run-settlement.test.ts

Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,3 +120,65 @@ test("rejected workers settle prior rollups with the latest observed model", asy
120120
});
121121
expect(Object.isFrozen(settlement)).toBe(true);
122122
});
123+
124+
test("pre-progress cancellation settles as cancelled without changing rejection", async () => {
125+
const cwd = await mkdtemp(join(tmpdir(), "corbits-run-cancelled-"));
126+
const controller = new AbortController();
127+
const originalError = new DOMException("operator cancelled", "AbortError");
128+
controller.abort(originalError);
129+
let settlement: Readonly<SubAgentRunSettlement> | undefined;
130+
131+
const caught = await withMockedModuleDuring(
132+
import.meta.resolve("../agent/live-tool-dispatch.js"),
133+
(real: typeof import("../agent/live-tool-dispatch.js")) => ({
134+
...real,
135+
createAgentWithLiveToolDispatch: async () => ({
136+
send: async () => {
137+
throw new Error("send must not start after cancellation");
138+
},
139+
stream: () => (async function* () {})(),
140+
deliver: () => {},
141+
close: async () => {},
142+
setSource: () => {},
143+
setSources: () => {},
144+
history: async () => [],
145+
checkpoints: async () => [],
146+
readAt: async () => [],
147+
blobReader: {},
148+
}),
149+
}),
150+
async () => {
151+
const { runSubAgent } = await import("./run.js");
152+
try {
153+
await runSubAgent({
154+
cwd,
155+
workdirBase: join(cwd, ".ctx"),
156+
permissionGate,
157+
provider: {
158+
providerName: "initial",
159+
baseURL: "http://localhost",
160+
model: "initial-model",
161+
},
162+
description: "cancelled settlement probe",
163+
prompt: "do not start",
164+
signal: controller.signal,
165+
onRunSettled: (summary) => {
166+
settlement = summary;
167+
},
168+
});
169+
} catch (error) {
170+
return error;
171+
}
172+
throw new Error("expected runSubAgent to reject");
173+
},
174+
);
175+
176+
expect(caught).toBe(originalError);
177+
expect(settlement).toMatchObject({
178+
turn_count: 0,
179+
error_count: 1,
180+
model: "initial-model",
181+
terminal_reason: "cancelled",
182+
});
183+
expect(Object.isFrozen(settlement)).toBe(true);
184+
});

src/subagent/run.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -383,6 +383,7 @@ export async function runSubAgent(params: RunSubAgentParams): Promise<RunSubAgen
383383
return result;
384384
} catch (error) {
385385
errorCount = 1;
386+
if (isSubAgentCancelError(error, params.signal)) terminalReason = "cancelled";
386387
throw error;
387388
} finally {
388389
try {

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

Lines changed: 103 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,14 @@ async function pathExists(path: string): Promise<boolean> {
6767
}
6868
}
6969

70+
async function waitFor(predicate: () => boolean | Promise<boolean>): Promise<void> {
71+
for (let attempt = 0; attempt < 500; attempt++) {
72+
if (await predicate()) return;
73+
await new Promise((resolve) => setTimeout(resolve, 1));
74+
}
75+
throw new Error("condition was not reached");
76+
}
77+
7078
function deferred<T>(): {
7179
promise: Promise<T>;
7280
resolve: (v: T) => void;
@@ -170,6 +178,62 @@ describe("spawn_agent worktree isolation", () => {
170178
expect(typeof ends[0]?.properties.duration_ms).toBe("number");
171179
});
172180

181+
test("pairs pre-progress cancellation with a cancelled terminal event", async () => {
182+
const repo = await makeRepo();
183+
tempDirs.push(repo);
184+
const { telemetry, events } = telemetryCapture();
185+
const sessions = createSubAgentSessionStore();
186+
const tool = createSpawnAgentTool({
187+
permissionGate: testPermissionGate,
188+
cwd: repo,
189+
getWorkdirBase: () => repo,
190+
provider,
191+
telemetry,
192+
run: async (params) => {
193+
params.onRunSettled?.({
194+
turn_count: 0,
195+
input_tokens: 0,
196+
output_tokens: 0,
197+
cache_read_tokens: 0,
198+
cache_write_tokens: 0,
199+
reasoning_tokens: 0,
200+
tool_call_count: 0,
201+
tool_error_count: 0,
202+
error_count: 1,
203+
duration_ms: 1,
204+
model: "test-model",
205+
terminal_reason: "cancelled",
206+
});
207+
const error = new Error("aborted");
208+
error.name = "AbortError";
209+
throw error;
210+
},
211+
sessions,
212+
fleetRecords: createFleetRecords(),
213+
});
214+
if (tool.kind !== "full") throw new Error("expected full tool");
215+
216+
const result = await tool.handler(
217+
{
218+
id: "cancelled-spawn",
219+
name: "spawn_agent",
220+
arguments: { description: "cancelled", prompt: "Do the work", intent: "explore" },
221+
},
222+
new AbortController().signal,
223+
);
224+
225+
expect(result.isError).not.toBe(true);
226+
await waitFor(() => events.some((event) => event.event === "subagent_end"));
227+
expect(sessions.list()[0]?.status).toBe("cancelled");
228+
expect(events.filter((event) => event.event === "subagent_start")).toHaveLength(1);
229+
const ends = events.filter((event) => event.event === "subagent_end");
230+
expect(ends).toHaveLength(1);
231+
expect(ends[0]?.properties).toMatchObject({
232+
status: "cancelled",
233+
stop_reason: "cancelled",
234+
});
235+
});
236+
173237
test("defers worktree cleanup while the session is retained for followup", async () => {
174238
const repo = await makeRepo();
175239
tempDirs.push(repo);
@@ -229,13 +293,17 @@ describe("spawn_agent worktree isolation", () => {
229293

230294
const settle = deferred<RunSubAgentResult>();
231295
let workerCwd: string | undefined;
296+
let settlementCount = 0;
297+
let settlementWasFrozen = false;
298+
const { telemetry, events } = telemetryCapture();
232299
const sessions = createSubAgentSessionStore();
233300
const tool = createSpawnAgentTool({
234301
permissionGate: testPermissionGate,
235302
cwd: repo,
236303
getWorkdirBase: () => workdirBase,
237304
provider,
238305
useWorktree: true,
306+
telemetry,
239307
run: async (params) => {
240308
workerCwd = params.cwd;
241309
params.onAgentReady?.({
@@ -244,7 +312,25 @@ describe("spawn_agent worktree isolation", () => {
244312
followup: async () => "",
245313
deliver: () => {},
246314
});
247-
return settle.promise;
315+
const result = await settle.promise;
316+
const summary = Object.freeze({
317+
turn_count: 0,
318+
input_tokens: 0,
319+
output_tokens: 0,
320+
cache_read_tokens: 0,
321+
cache_write_tokens: 0,
322+
reasoning_tokens: 0,
323+
tool_call_count: 0,
324+
tool_error_count: 0,
325+
error_count: 0,
326+
duration_ms: 1,
327+
model: "test-model",
328+
terminal_reason: "cancelled" as const,
329+
});
330+
settlementCount += 1;
331+
settlementWasFrozen = Object.isFrozen(summary);
332+
params.onRunSettled?.(summary);
333+
return result;
248334
},
249335
sessions,
250336
fleetRecords: createFleetRecords(),
@@ -263,10 +349,20 @@ describe("spawn_agent worktree isolation", () => {
263349

264350
settle.resolve({
265351
report: "## Summary\nStopped.\n## Findings\npartial\n## Blockers\ninterrupted\n## Paths\n",
352+
stopReason: "cancelled",
266353
interrupted: true,
267354
});
268-
await new Promise((resolve) => setTimeout(resolve, 50));
355+
await waitFor(() => events.some((event) => event.event === "subagent_end"));
269356

357+
expect(settlementCount).toBe(1);
358+
expect(settlementWasFrozen).toBe(true);
359+
expect(sessions.get(agentId)?.lifecycleStatus).toBe("interrupted");
360+
const ends = events.filter((event) => event.event === "subagent_end");
361+
expect(ends).toHaveLength(1);
362+
expect(ends[0]?.properties).toMatchObject({
363+
status: "interrupted",
364+
stop_reason: "cancelled",
365+
});
270366
expect(workerCwd).toBeDefined();
271367
expect(await pathExists(workerCwd!)).toBe(true);
272368

@@ -305,9 +401,11 @@ describe("spawn_agent worktree isolation", () => {
305401
},
306402
new AbortController().signal,
307403
);
308-
await new Promise((resolve) => setTimeout(resolve, 50));
404+
await waitFor(() => workerCwd !== undefined);
405+
if (workerCwd === undefined) throw new Error("worker cwd was not captured");
406+
const completedWorkerCwd = workerCwd;
407+
await waitFor(async () => !(await pathExists(completedWorkerCwd)));
309408

310-
expect(workerCwd).toBeDefined();
311-
expect(await pathExists(workerCwd!)).toBe(false);
409+
expect(await pathExists(completedWorkerCwd)).toBe(false);
312410
});
313411
});

src/subagent/task-tool-worktree.test.ts

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,53 @@ describe("createTaskTool worktree isolation", () => {
178178
expect(typeof ends[0]?.properties.duration_ms).toBe("number");
179179
});
180180

181+
test("pairs pre-progress cancellation with a cancelled terminal event", async () => {
182+
const repo = await makeRepo();
183+
tempDirs.push(repo);
184+
const { telemetry, events } = telemetryCapture();
185+
const tool = createTaskTool({
186+
permissionGate: testPermissionGate,
187+
cwd: repo,
188+
getWorkdirBase: () => repo,
189+
provider,
190+
telemetry,
191+
run: async (params) => {
192+
params.onRunSettled?.({
193+
turn_count: 0,
194+
input_tokens: 0,
195+
output_tokens: 0,
196+
cache_read_tokens: 0,
197+
cache_write_tokens: 0,
198+
reasoning_tokens: 0,
199+
tool_call_count: 0,
200+
tool_error_count: 0,
201+
error_count: 1,
202+
duration_ms: 1,
203+
model: "test-model",
204+
terminal_reason: "cancelled",
205+
});
206+
const error = new Error("aborted");
207+
error.name = "AbortError";
208+
throw error;
209+
},
210+
});
211+
212+
const result = await callTask(tool, {
213+
description: "cancelled job",
214+
prompt: "Do the work",
215+
intent: "explore",
216+
});
217+
218+
expect(result).toContain("cancelled by operator");
219+
expect(events.filter((event) => event.event === "subagent_start")).toHaveLength(1);
220+
const ends = events.filter((event) => event.event === "subagent_end");
221+
expect(ends).toHaveLength(1);
222+
expect(ends[0]?.properties).toMatchObject({
223+
status: "cancelled",
224+
stop_reason: "cancelled",
225+
});
226+
});
227+
181228
test("preserves a worktree the sub-agent left dirty, with a notice in the report", async () => {
182229
const repo = await makeRepo();
183230
tempDirs.push(repo);

0 commit comments

Comments
 (0)