diff --git a/extensions/subagents/docs/design-plan.md b/extensions/subagents/docs/design-plan.md index 76556a50..007dbc0a 100644 --- a/extensions/subagents/docs/design-plan.md +++ b/extensions/subagents/docs/design-plan.md @@ -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). diff --git a/extensions/subagents/index.test.ts b/extensions/subagents/index.test.ts index a52114ce..257dc896 100644 --- a/extensions/subagents/index.test.ts +++ b/extensions/subagents/index.test.ts @@ -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) { @@ -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, [ { @@ -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 }, }, ]); }); diff --git a/extensions/subagents/index.ts b/extensions/subagents/index.ts index 59d8fd3a..a5db7af8 100644 --- a/extensions/subagents/index.ts +++ b/extensions/subagents/index.ts @@ -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, @@ -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) => @@ -249,7 +248,7 @@ export function createSubagentResultDispatcher( display: false, details, }, - resultDeliveryOptions(wake), + { deliverAs: "followUp", triggerTurn: true }, ); }; } @@ -322,8 +321,14 @@ export default function (pi: ExtensionAPI) { let requestWidgetRender: (() => void) | undefined; let navigationLayerRegistered = false; let dashboardOpen = false; - const resultDelivery = createDeferredResultDelivery(); const dispatchResults = createSubagentResultDispatcher(pi); + const resultDelivery = createSubagentResultDelivery({ + 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, @@ -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 @@ -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) => { @@ -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"); diff --git a/extensions/subagents/result-delivery.test.ts b/extensions/subagents/result-delivery.test.ts index 44aaab25..38c910a4 100644 --- a/extensions/subagents/result-delivery.test.ts +++ b/extensions/subagents/result-delivery.test.ts @@ -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"]]); }); diff --git a/extensions/subagents/src/result-delivery.ts b/extensions/subagents/src/result-delivery.ts index 9c7bb6da..52eb2739 100644 --- a/extensions/subagents/src/result-delivery.ts +++ b/extensions/subagents/src/result-delivery.ts @@ -1,17 +1,62 @@ -export function createDeferredResultDelivery() { +export interface SubagentResultDeliveryOptions { + /** 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( + options: SubagentResultDeliveryOptions, +) { const pending = new Map(); + 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) { 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();