Skip to content

Commit dff5aa4

Browse files
committed
fix(subagents): wake the parent for deferred results
Results that settled while the parent was busy were deliberately flushed as nextTurn at agent_settled. Pi only consumes that queue on another user prompt, so fire-and-forget results could not fulfill their documented automatic re-invocation contract. Flush the pending batch as one follow-up at the authoritative parent boundary, while preserving wait consumption, exactly-once delivery, synchronous retry, and aborted-turn suppression. Closes #47
1 parent 4eced51 commit dff5aa4

5 files changed

Lines changed: 311 additions & 78 deletions

File tree

extensions/subagents/docs/design-plan.md

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -71,13 +71,17 @@ the parent conversation.
7171
### 1.3 Result delivery back to the parent
7272

7373
- When a child settles **unconsumed**, `onSettled` defers it into a tiny
74-
`createDeferredResultDelivery` buffer (defer/consume/drain/clear keyed by id).
74+
`createSubagentResultDelivery` buffer (defer/consume/clear keyed by id).
7575
- Flush happens when the parent goes idle: immediately if `sessionContext.isIdle()`,
76-
otherwise on the parent's `agent_settled` event. A later `subagent_wait` can still
77-
consume a deferred result before flush (that's why it is a buffer, not an immediate
78-
send).
79-
- Delivery = `pi.sendMessage({ customType: "subagent-result", content, display: true,
80-
details: { id, title, status } }, { deliverAs: "followUp", triggerTurn: true })`.
76+
otherwise on the parent's authoritative `agent_settled` event. Busy-period results
77+
are batched into one automatic wake-up so none can be stranded until the user's next
78+
prompt. A later `subagent_wait` can still consume a deferred result before flush.
79+
If the parent run was aborted, delivery uses `nextTurn` rather than resurrecting
80+
work the user explicitly stopped.
81+
- Normal delivery = `pi.sendMessage({ customType: "subagent-result", content,
82+
display: false, details: { id, title, status } },
83+
{ deliverAs: "followUp", triggerTurn: true })`; a separate session entry renders the
84+
report at its actual completion point.
8185
Content is built by `buildSubagentResultMessage` (`Subagent sa-N "title"
8286
finished/failed.` + optional `Error:` line + output truncated to 24KB/600 lines with a
8387
pointer to the child session file for the full transcript).

extensions/subagents/index.test.ts

Lines changed: 87 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,12 @@ import type {
99
ExtensionContext,
1010
} from "@earendil-works/pi-coding-agent";
1111
import { PLAN_MODE_CHANNEL } from "../shared/plan-mode-state.ts";
12-
import subagents, { createSubagentResultDispatcher } from "./index.ts";
12+
import subagents, {
13+
createSubagentResultDispatcher,
14+
registerSubagentResultDeliveryLifecycle,
15+
} from "./index.ts";
1316

14-
test("deferred subagent results render before a hidden next-turn injection", () => {
17+
test("subagent results render before the hidden wake-up message", () => {
1518
const events: unknown[] = [];
1619
const pi = {
1720
appendEntry(customType: string, data: unknown) {
@@ -23,29 +26,26 @@ test("deferred subagent results render before a hidden next-turn injection", ()
2326
} as unknown as ExtensionAPI;
2427
const dispatch = createSubagentResultDispatcher(pi, () => "report");
2528

26-
dispatch(
27-
[
28-
{
29-
id: "sa-3",
30-
origin: "model",
31-
backend: "pi",
32-
title: "investigate plan mode",
33-
prompt: "inspect",
34-
cwd: process.cwd(),
35-
status: "done",
36-
createdAt: 0,
37-
settledAt: 1_000,
38-
meta: { backend: "pi" },
39-
usage: {},
40-
transcript: [],
41-
liveTools: [],
42-
queued: [],
43-
finalText: "report",
44-
turns: 1,
45-
},
46-
],
47-
false,
48-
);
29+
dispatch([
30+
{
31+
id: "sa-3",
32+
origin: "model",
33+
backend: "pi",
34+
title: "investigate plan mode",
35+
prompt: "inspect",
36+
cwd: process.cwd(),
37+
status: "done",
38+
createdAt: 0,
39+
settledAt: 1_000,
40+
meta: { backend: "pi" },
41+
usage: {},
42+
transcript: [],
43+
liveTools: [],
44+
queued: [],
45+
finalText: "report",
46+
turns: 1,
47+
},
48+
]);
4949

5050
assert.deepEqual(events, [
5151
{
@@ -74,9 +74,70 @@ test("deferred subagent results render before a hidden next-turn injection", ()
7474
status: "done",
7575
},
7676
},
77-
options: { deliverAs: "nextTurn" },
77+
options: { deliverAs: "followUp", triggerTurn: true },
7878
},
7979
]);
80+
81+
events.length = 0;
82+
dispatch(
83+
[
84+
{
85+
id: "sa-4",
86+
origin: "model",
87+
backend: "pi",
88+
title: "aborted parent",
89+
prompt: "inspect",
90+
cwd: process.cwd(),
91+
status: "done",
92+
createdAt: 0,
93+
settledAt: 1_000,
94+
meta: { backend: "pi" },
95+
usage: {},
96+
transcript: [],
97+
liveTools: [],
98+
queued: [],
99+
finalText: "report",
100+
turns: 1,
101+
},
102+
],
103+
false,
104+
);
105+
assert.deepEqual((events.at(-1) as { options: unknown }).options, {
106+
deliverAs: "nextTurn",
107+
});
108+
});
109+
110+
test("the parent lifecycle wakes deferred results unless the run was aborted", () => {
111+
const handlers = new Map<string, (event: never) => void>();
112+
const pi = {
113+
on(event: string, handler: (event: never) => void) {
114+
handlers.set(event, handler);
115+
},
116+
} as unknown as ExtensionAPI;
117+
const settled: Array<boolean | undefined> = [];
118+
registerSubagentResultDeliveryLifecycle(pi, {
119+
parentSettled: (aborted) => settled.push(aborted),
120+
});
121+
122+
handlers.get("agent_start")?.({ type: "agent_start" } as never);
123+
handlers.get("agent_end")?.({
124+
type: "agent_end",
125+
messages: [{ role: "assistant", stopReason: "stop" }],
126+
} as never);
127+
handlers.get("agent_settled")?.({ type: "agent_settled" } as never);
128+
129+
handlers.get("agent_start")?.({ type: "agent_start" } as never);
130+
handlers.get("agent_end")?.({
131+
type: "agent_end",
132+
messages: [
133+
{ role: "assistant", stopReason: "stop" },
134+
{ role: "toolResult" },
135+
{ role: "assistant", stopReason: "aborted" },
136+
],
137+
} as never);
138+
handlers.get("agent_settled")?.({ type: "agent_settled" } as never);
139+
140+
assert.deepEqual(settled, [false, true]);
80141
});
81142

82143
test("the visible subagent result entry renders the completed report", () => {

extensions/subagents/index.ts

Lines changed: 43 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -95,8 +95,7 @@ import {
9595
SUBAGENT_WAIT_PARAMETER_DESCRIPTIONS,
9696
SUBAGENT_WAIT_TOOL_DESCRIPTION,
9797
} from "./src/prompt.ts";
98-
import { createDeferredResultDelivery } from "./src/result-delivery.ts";
99-
import { resultDeliveryOptions } from "../background-terminals/src/result-delivery.ts";
98+
import { createSubagentResultDelivery } from "./src/result-delivery.ts";
10099
import {
101100
effectiveChildToolAllowlist,
102101
resolveStandaloneChildProjectTrust,
@@ -210,7 +209,7 @@ export function createSubagentResultDispatcher(
210209
pi: ExtensionAPI,
211210
outputFor: (snap: SubagentSnapshot) => string = truncatedOutput,
212211
) {
213-
return (snaps: readonly SubagentSnapshot[], wake: boolean) => {
212+
return (snaps: readonly SubagentSnapshot[], wake = true) => {
214213
if (snaps.length === 0) return;
215214
const content = snaps
216215
.map((snap) =>
@@ -249,11 +248,37 @@ export function createSubagentResultDispatcher(
249248
display: false,
250249
details,
251250
},
252-
resultDeliveryOptions(wake),
251+
wake
252+
? { deliverAs: "followUp", triggerTurn: true }
253+
: { deliverAs: "nextTurn" },
253254
);
254255
};
255256
}
256257

258+
export function registerSubagentResultDeliveryLifecycle(
259+
pi: ExtensionAPI,
260+
delivery: { parentSettled(aborted?: boolean): void },
261+
) {
262+
let parentRunAborted = false;
263+
264+
pi.on("agent_start", () => {
265+
parentRunAborted = false;
266+
});
267+
pi.on("agent_end", (event) => {
268+
for (let index = event.messages.length - 1; index >= 0; index--) {
269+
const message = event.messages[index];
270+
if (message?.role !== "assistant") continue;
271+
parentRunAborted = message.stopReason === "aborted";
272+
break;
273+
}
274+
});
275+
pi.on("agent_settled", () => delivery.parentSettled(parentRunAborted));
276+
277+
return () => {
278+
parentRunAborted = false;
279+
};
280+
}
281+
257282
type SubagentResultTheme = Parameters<MessageRenderer>[2];
258283

259284
function renderSubagentResult(
@@ -322,8 +347,17 @@ export default function (pi: ExtensionAPI) {
322347
let requestWidgetRender: (() => void) | undefined;
323348
let navigationLayerRegistered = false;
324349
let dashboardOpen = false;
325-
const resultDelivery = createDeferredResultDelivery<SubagentSnapshot>();
326350
const dispatchResults = createSubagentResultDispatcher(pi);
351+
const resultDelivery = createSubagentResultDelivery<SubagentSnapshot>({
352+
isIdle: () => sessionContext?.isIdle() === true,
353+
// Every unconsumed fire-and-forget result must reach the parent. The
354+
// delivery coordinator batches results that settled while it was busy.
355+
deliver: dispatchResults,
356+
});
357+
const resetResultDeliveryLifecycle = registerSubagentResultDeliveryLifecycle(
358+
pi,
359+
resultDelivery,
360+
);
327361
const hideLifecycleTools = () =>
328362
patchOwnedTools(pi, "subagents", {
329363
disable: OPENPI_TOOL_SURFACE.subagents.deferred,
@@ -436,25 +470,6 @@ export default function (pi: ExtensionAPI) {
436470
navigationLayerRegistered = true;
437471
};
438472

439-
/**
440-
* `wake` decides whether this costs the model a turn. A subagent that
441-
* settled while the model sits idle is the result it is waiting on. A
442-
* backlog that piled up while it worked is not: waking once per stale
443-
* subagent forces a turn each, and the model can only answer "that one
444-
* already finished". `nextTurn` still enters context with the user's next
445-
* message, without demanding a reply.
446-
*/
447-
const deliverResults = (
448-
snaps: readonly SubagentSnapshot[],
449-
wake: boolean,
450-
) => {
451-
dispatchResults(snaps, wake);
452-
};
453-
454-
const flushResults = (wake: boolean) => {
455-
deliverResults(resultDelivery.drain(), wake);
456-
};
457-
458473
const deliverBtwResult = (snap: SubagentSnapshot) => {
459474
// appendEntry is a synchronous SessionManager operation and emits an
460475
// entry_appended event, so it is safe while the parent is streaming and
@@ -500,10 +515,10 @@ export default function (pi: ExtensionAPI) {
500515
// subagent_wait can consume it before agent_settled flushes follow-ups.
501516
// Defer a copy: the live snapshot keeps mutating if the subagent is
502517
// restarted before the deferred result flushes.
518+
// The delivery coordinator closes both sides of the wake-up race: it
519+
// flushes now if the parent is already idle, otherwise the parent's next
520+
// agent_settled edge rechecks this same pending Map.
503521
resultDelivery.defer({ ...snap, meta: { ...snap.meta } });
504-
// Settled while the model sits idle: it has nothing else in flight, so
505-
// this is the result it is waiting on — wake it.
506-
if (sessionContext?.isIdle()) flushResults(true);
507522
};
508523

509524
pi.on("session_start", (_event, ctx) => {
@@ -530,10 +545,6 @@ export default function (pi: ExtensionAPI) {
530545
managerPromise?.then(updateStatus).catch(() => undefined);
531546
});
532547

533-
// These settled while the model was working on something else, so they go
534-
// into context without forcing a turn per stale subagent.
535-
pi.on("agent_settled", () => flushResults(false));
536-
537548
pi.on("session_shutdown", async () => {
538549
if (navigationLayerRegistered) {
539550
removeEditorLayer(pi, "subagents");
@@ -555,6 +566,7 @@ export default function (pi: ExtensionAPI) {
555566
requestWidgetRender = undefined;
556567
stripState.focused = false;
557568
dashboardOpen = false;
569+
resetResultDeliveryLifecycle();
558570
const closing = runtime;
559571
runtime = undefined;
560572
managerPromise = undefined;

0 commit comments

Comments
 (0)