Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 8 additions & 6 deletions extensions/subagents/docs/design-plan.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,13 +71,15 @@ the parent conversation.
### 1.3 Result delivery back to the parent

- When a child settles **unconsumed**, `onSettled` defers it into a tiny
`createDeferredResultDelivery` buffer (defer/consume/drain/clear keyed by id).
`createSubagentResultDelivery` buffer (defer/consume/clear keyed by id).
- Flush happens when the parent goes idle: immediately if `sessionContext.isIdle()`,
otherwise on the parent's `agent_settled` event. A later `subagent_wait` can still
consume a deferred result before flush (that's why it is a buffer, not an immediate
send).
- Delivery = `pi.sendMessage({ customType: "subagent-result", content, display: true,
details: { id, title, status } }, { deliverAs: "followUp", triggerTurn: true })`.
otherwise on the parent's authoritative `agent_settled` event. Busy-period results
are batched into one automatic wake-up so none can be stranded until the user's next
prompt. A later `subagent_wait` can still consume a deferred result before flush.
- Normal delivery = `pi.sendMessage({ customType: "subagent-result", content,
display: false, details: { id, title, status } },
{ deliverAs: "followUp", triggerTurn: true })`; a separate session entry renders the
report at its actual completion point.
Content is built by `buildSubagentResultMessage` (`Subagent sa-N "title"
finished/failed.` + optional `Error:` line + output truncated to 24KB/600 lines with a
pointer to the child session file for the full transcript).
Expand Down
47 changes: 22 additions & 25 deletions extensions/subagents/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import type {
import { PLAN_MODE_CHANNEL } from "../shared/plan-mode-state.ts";
import subagents, { createSubagentResultDispatcher } from "./index.ts";

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

dispatch(
[
{
id: "sa-3",
origin: "model",
backend: "pi",
title: "investigate plan mode",
prompt: "inspect",
cwd: process.cwd(),
status: "done",
createdAt: 0,
settledAt: 1_000,
meta: { backend: "pi" },
usage: {},
transcript: [],
liveTools: [],
queued: [],
finalText: "report",
turns: 1,
},
],
false,
);
dispatch([
{
id: "sa-3",
origin: "model",
backend: "pi",
title: "investigate plan mode",
prompt: "inspect",
cwd: process.cwd(),
status: "done",
createdAt: 0,
settledAt: 1_000,
meta: { backend: "pi" },
usage: {},
transcript: [],
liveTools: [],
queued: [],
finalText: "report",
turns: 1,
},
]);

assert.deepEqual(events, [
{
Expand Down Expand Up @@ -74,7 +71,7 @@ test("deferred subagent results render before a hidden next-turn injection", ()
status: "done",
},
},
options: { deliverAs: "nextTurn" },
options: { deliverAs: "followUp", triggerTurn: true },
},
]);
});
Expand Down
44 changes: 13 additions & 31 deletions extensions/subagents/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -95,8 +95,7 @@ import {
SUBAGENT_WAIT_PARAMETER_DESCRIPTIONS,
SUBAGENT_WAIT_TOOL_DESCRIPTION,
} from "./src/prompt.ts";
import { createDeferredResultDelivery } from "./src/result-delivery.ts";
import { resultDeliveryOptions } from "../background-terminals/src/result-delivery.ts";
import { createSubagentResultDelivery } from "./src/result-delivery.ts";
import {
effectiveChildToolAllowlist,
resolveStandaloneChildProjectTrust,
Expand Down Expand Up @@ -210,7 +209,7 @@ export function createSubagentResultDispatcher(
pi: ExtensionAPI,
outputFor: (snap: SubagentSnapshot) => string = truncatedOutput,
) {
return (snaps: readonly SubagentSnapshot[], wake: boolean) => {
return (snaps: readonly SubagentSnapshot[]) => {
if (snaps.length === 0) return;
const content = snaps
.map((snap) =>
Expand Down Expand Up @@ -249,7 +248,7 @@ export function createSubagentResultDispatcher(
display: false,
details,
},
resultDeliveryOptions(wake),
{ deliverAs: "followUp", triggerTurn: true },
);
};
}
Expand Down Expand Up @@ -322,8 +321,14 @@ export default function (pi: ExtensionAPI) {
let requestWidgetRender: (() => void) | undefined;
let navigationLayerRegistered = false;
let dashboardOpen = false;
const resultDelivery = createDeferredResultDelivery<SubagentSnapshot>();
const dispatchResults = createSubagentResultDispatcher(pi);
const resultDelivery = createSubagentResultDelivery<SubagentSnapshot>({
isIdle: () => sessionContext?.isIdle() === true,
// Every unconsumed fire-and-forget result must reach the parent. The
// delivery coordinator batches results that settled while it was busy.
deliver: dispatchResults,
});
pi.on("agent_settled", () => resultDelivery.parentSettled());
const hideLifecycleTools = () =>
patchOwnedTools(pi, "subagents", {
disable: OPENPI_TOOL_SURFACE.subagents.deferred,
Expand Down Expand Up @@ -436,25 +441,6 @@ export default function (pi: ExtensionAPI) {
navigationLayerRegistered = true;
};

/**
* `wake` decides whether this costs the model a turn. A subagent that
* settled while the model sits idle is the result it is waiting on. A
* backlog that piled up while it worked is not: waking once per stale
* subagent forces a turn each, and the model can only answer "that one
* already finished". `nextTurn` still enters context with the user's next
* message, without demanding a reply.
*/
const deliverResults = (
snaps: readonly SubagentSnapshot[],
wake: boolean,
) => {
dispatchResults(snaps, wake);
};

const flushResults = (wake: boolean) => {
deliverResults(resultDelivery.drain(), wake);
};

const deliverBtwResult = (snap: SubagentSnapshot) => {
// appendEntry is a synchronous SessionManager operation and emits an
// entry_appended event, so it is safe while the parent is streaming and
Expand Down Expand Up @@ -500,10 +486,10 @@ export default function (pi: ExtensionAPI) {
// subagent_wait can consume it before agent_settled flushes follow-ups.
// Defer a copy: the live snapshot keeps mutating if the subagent is
// restarted before the deferred result flushes.
// The delivery coordinator closes both sides of the wake-up race: it
// flushes now if the parent is already idle, otherwise the parent's next
// agent_settled edge rechecks this same pending Map.
resultDelivery.defer({ ...snap, meta: { ...snap.meta } });
// Settled while the model sits idle: it has nothing else in flight, so
// this is the result it is waiting on — wake it.
if (sessionContext?.isIdle()) flushResults(true);
};

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

// These settled while the model was working on something else, so they go
// into context without forcing a turn per stale subagent.
pi.on("agent_settled", () => flushResults(false));

pi.on("session_shutdown", async () => {
if (navigationLayerRegistered) {
removeEditorLayer(pi, "subagents");
Expand Down
119 changes: 109 additions & 10 deletions extensions/subagents/result-delivery.test.ts
Original file line number Diff line number Diff line change
@@ -1,27 +1,126 @@
import assert from "node:assert/strict";
import test from "node:test";
import { createDeferredResultDelivery } from "./src/result-delivery.ts";
import { createSubagentResultDelivery } from "./src/result-delivery.ts";

test("a result consumed by a later wait is not delivered", () => {
const delivery = createDeferredResultDelivery<{
interface DeliveredBatch {
results: readonly { id: string; output?: string }[];
}

function harness(initialIdle: boolean) {
let idle = initialIdle;
const deliveries: DeliveredBatch[] = [];
const delivery = createSubagentResultDelivery<{
id: string;
output: string;
}>();
output?: string;
}>({
isIdle: () => idle,
deliver: (results) => deliveries.push({ results }),
});
return {
delivery,
deliveries,
setIdle(value: boolean) {
idle = value;
},
};
}

test("a result consumed by a later wait is not delivered", () => {
const { delivery, deliveries } = harness(false);

delivery.defer({ id: "sa-1", output: "done" });
delivery.consume(["sa-1"]);
delivery.parentSettled();

assert.deepEqual(deliveries, []);
});

test("an idle child settlement wakes the parent immediately", () => {
const { delivery, deliveries } = harness(true);

delivery.defer({ id: "sa-1" });

assert.deepEqual(delivery.drain(), []);
assert.deepEqual(deliveries, [{ results: [{ id: "sa-1" }] }]);
});

test("unconsumed results are delivered once in settlement order", () => {
const delivery = createDeferredResultDelivery<{ id: string }>();
test("a busy child settlement wakes when the parent settles", () => {
const { delivery, deliveries, setIdle } = harness(false);

delivery.defer({ id: "sa-1" });
assert.deepEqual(deliveries, []);

setIdle(true);
delivery.parentSettled();

assert.deepEqual(deliveries, [{ results: [{ id: "sa-1" }] }]);
});

test("the parent boundary still delivers if another extension started a run", () => {
const { delivery, deliveries } = harness(false);

delivery.defer({ id: "sa-1" });
// Pi has already emitted parent agent_settled, but an earlier extension
// handler started another run before this extension's handler executes.
delivery.parentSettled();

assert.deepEqual(deliveries, [{ results: [{ id: "sa-1" }] }]);
});

test("busy results batch in settlement order into one parent wake-up", () => {
const { delivery, deliveries } = harness(false);
const first = { id: "sa-1" };
const second = { id: "sa-2" };

delivery.defer(first);
delivery.defer(second);
delivery.parentSettled();

assert.deepEqual(deliveries, [{ results: [first, second] }]);
});

test("the child-settled and parent-settled edges cannot double deliver", () => {
const { delivery, deliveries, setIdle } = harness(false);

delivery.defer({ id: "sa-1" });
setIdle(true);
delivery.parentSettled();
delivery.parentSettled();

assert.deepEqual(deliveries, [{ results: [{ id: "sa-1" }] }]);
});

test("a child settling during the wake turn waits for its boundary", () => {
const { delivery, deliveries, setIdle } = harness(true);

delivery.defer({ id: "sa-1" });
setIdle(false);
delivery.defer({ id: "sa-2" });
assert.deepEqual(deliveries, [{ results: [{ id: "sa-1" }] }]);

delivery.parentSettled();
assert.deepEqual(deliveries, [
{ results: [{ id: "sa-1" }] },
{ results: [{ id: "sa-2" }] },
]);
});

test("a synchronous delivery failure restores the batch in order", () => {
let attempts = 0;
const delivered: string[][] = [];
const delivery = createSubagentResultDelivery<{ id: string }>({
isIdle: () => false,
deliver(results) {
attempts++;
if (attempts === 1) throw new Error("session switching");
delivered.push(results.map(({ id }) => id));
},
});

delivery.defer({ id: "sa-1" });
delivery.defer({ id: "sa-2" });
assert.throws(() => delivery.parentSettled(), /session switching/);
delivery.defer({ id: "sa-3" });
delivery.parentSettled();

assert.deepEqual(delivery.drain(), [first, second]);
assert.deepEqual(delivery.drain(), []);
assert.deepEqual(delivered, [["sa-1", "sa-2", "sa-3"]]);
});
55 changes: 50 additions & 5 deletions extensions/subagents/src/result-delivery.ts
Original file line number Diff line number Diff line change
@@ -1,17 +1,62 @@
export function createDeferredResultDelivery<T extends { id: string }>() {
export interface SubagentResultDeliveryOptions<T> {
/** True only when the parent has no run or queued continuation in flight. */
readonly isIdle: () => boolean;
/** Deliver one drained batch and wake the parent. */
readonly deliver: (results: readonly T[]) => void;
}

/**
* One-shot result delivery for fire-and-forget subagents.
*
* The tool contract promises that a settled child re-invokes the parent. A
* child that settles while the parent is busy therefore remains retractable
* until the parent's `agent_settled` event, but it must never be downgraded to
* a `nextTurn` message that needs another user prompt. There are two symmetric
* wake-up edges so no ordering can lose the notification:
*
* 1. child settles after the parent became idle -> `defer` flushes now;
* 2. parent settles after the child -> `parentSettled` flushes the batch.
*
* The parent boundary wakes even if an earlier extension handler has already
* started another turn: Pi queues the follow-up into that active run.
*
* The Map is the one-shot gate: `subagent_wait` may consume a result before it
* is delivered, and whichever path drains first prevents duplicate delivery.
*/
export function createSubagentResultDelivery<T extends { id: string }>(
options: SubagentResultDeliveryOptions<T>,
) {
const pending = new Map<string, T>();

const flush = () => {
if (pending.size === 0) return;
const results = [...pending.values()];
pending.clear();
try {
options.deliver(results);
} catch (error) {
// A synchronous session teardown may reject append/send. Preserve the
// original batch ahead of anything deferred re-entrantly while delivery
// ran, so a later boundary can retry without loss or reordering.
const current = [...pending.values()];
pending.clear();
for (const result of results) pending.set(result.id, result);
for (const result of current) pending.set(result.id, result);
throw error;
}
};

return {
defer(result: T) {
pending.set(result.id, result);
if (options.isIdle()) flush();
},
consume(ids: Iterable<string>) {
for (const id of ids) pending.delete(id);
},
drain() {
const results = [...pending.values()];
pending.clear();
return results;
/** Flush at the authoritative parent boundary. */
parentSettled() {
flush();
},
clear() {
pending.clear();
Expand Down
Loading