// Descendant continuation: parent-side delivery of Descendant results as // continuation turns. The caller decides mode: a Direct parent's nested // results queue into one Descendant continuation batch whose flush re-checks // the Cancellation cutover; everything else delivers directly. Each // Parent Runtime ID owns an independent reload-safe queue. See CONTEXT.md // "Descendant continuation batch". import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { CompletionDeliverySuppressedError } from "./completion-delivery.ts"; import { getProcessParentRuntimeId, type ParentRuntimeId } from "./operation-identity.ts"; import { retainAcrossReload } from "./reload-boundary.ts"; import type { NestedLifecycleCoordinator } from "./nested-lifecycle.ts"; /** Shared-runtime key: queue state must survive extension module reload. */ const DESCENDANT_CONTINUATION_KEY = Symbol.for("pi-subagents/descendant-continuation"); interface PendingDescendantSteer { payload: any; resolve: (value: unknown) => void; reject: (error: unknown) => void; } /** Queue state retained while an extension module is reloaded. */ interface DescendantContinuationState { pending: PendingDescendantSteer[]; flushScheduled: boolean; /** Most recent parent API for this Parent Runtime ID. */ lastApi: ExtensionAPI | undefined; } type ContinuationStateStore = Map; const states: ContinuationStateStore = retainAcrossReload( DESCENDANT_CONTINUATION_KEY, () => new Map(), ); function stateFor(parentRuntimeId: ParentRuntimeId): DescendantContinuationState { const existing = states.get(parentRuntimeId); if (existing) return existing; const created: DescendantContinuationState = { pending: [], flushScheduled: false, lastApi: undefined, }; states.set(parentRuntimeId, created); return created; } /** A Fresh descendant process batches continuations; parents deliver directly. */ /** * Resolve the parent delivery API. The latest instantiated API for this * Parent Runtime ID wins, so a watcher closure held across /reload delivers * through the reloaded parent; the fallback is the API that queued the batch. */ function resolveParentApi( state: DescendantContinuationState, fallback: ExtensionAPI | undefined, ): ExtensionAPI { const api = state.lastApi ?? fallback; if (!api) throw new Error("Descendant continuation requires a parent extension API."); return api; } function sendDescendantSteerBatch( state: DescendantContinuationState, fallbackApi: ExtensionAPI | undefined, coordinator: NestedLifecycleCoordinator, ): void { state.flushScheduled = false; const pending = state.pending; state.pending = []; if (pending.length === 0) return; // A queued steer may outlive the operation that accepted it. The // Cancellation cutover is checked again at flush time so it cannot reopen // a cancelled parent after the handoff was prepared. if (coordinator.isCancellationRequested()) { const error = new CompletionDeliverySuppressedError( "Descendant continuation suppressed by Cancellation cutover.", ); for (const entry of pending) entry.reject(error); return; } const first = pending[0].payload; const payload = pending.length === 1 ? first : { ...first, content: pending.map((entry) => entry.payload.content).join("\n\n"), details: { ...first.details, descendantBatch: pending.map((entry) => ({ customType: entry.payload.customType, ...entry.payload.details, })), }, }; try { const accepted = resolveParentApi(state, fallbackApi).sendMessage( payload, { triggerTurn: true, deliverAs: "steer" }, ); for (const entry of pending) entry.resolve(accepted); } catch (error) { for (const entry of pending) entry.reject(error); } } function scheduleDescendantSteerFlush( state: DescendantContinuationState, fallbackApi: ExtensionAPI | undefined, coordinator: NestedLifecycleCoordinator, ): void { if (state.flushScheduled) return; state.flushScheduled = true; // The flush delivers through the latest noted parent API for this runtime: // a batch queued before /reload is flushed by the reloaded parent API. setTimeout(() => sendDescendantSteerBatch(state, fallbackApi, coordinator), 0); } /** * Deliver one steer directly through the latest parent API. Used by callers * outside a nested-lifecycle context, and for steers that are not descendant * continuations. */ export function sendSteerDirect( pi: ExtensionAPI, payload: any, parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): Promise { const state = stateFor(parentRuntimeId); return Promise.resolve( resolveParentApi(state, pi).sendMessage( payload, { triggerTurn: true, deliverAs: "steer" }, ), ); } /** * Queue one descendant result steer into the Descendant continuation batch; * simultaneous results coalesce into a single continuation turn. Queues are * scoped by the owning Parent Runtime ID so unrelated parents cannot merge. */ export function queueDescendantContinuation( pi: ExtensionAPI, payload: any, coordinator: NestedLifecycleCoordinator, ): Promise { const state = stateFor(coordinator.parentRuntimeId); // The parent runtime records the current API at instantiation. A queue may // be invoked by a watcher closure captured before /reload, so never let a // stale fallback overwrite an already-noted current API. if (!state.lastApi) state.lastApi = pi; return new Promise((resolve, reject) => { state.pending.push({ payload, resolve, reject }); scheduleDescendantSteerFlush(state, pi, coordinator); }); } /** Reschedule the flush after a /reload; the queue itself never moved. */ export function rescheduleFlushAfterReload( pi: ExtensionAPI, coordinator: NestedLifecycleCoordinator, ): void { const state = stateFor(coordinator.parentRuntimeId); if (state.pending.length > 0) { // The reloaded module flushes the surviving queue through its own API. state.lastApi = pi; scheduleDescendantSteerFlush(state, pi, coordinator); } } /** True while continuations are queued for this Parent Runtime ID. */ export function hasPendingDescendantContinuations( parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): boolean { return stateFor(parentRuntimeId).pending.length > 0; } /** * Note the current parent API for one Parent Runtime ID. Called on every * extension instantiation, so queued batches survive /reload without crossing * into another parent runtime. */ export function noteParentApi( pi: ExtensionAPI, parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): void { stateFor(parentRuntimeId).lastApi = pi; } /** Reject every queued continuation for one parent runtime. */ export function rejectPendingDescendantContinuations( error: Error, parentRuntimeId: ParentRuntimeId = getProcessParentRuntimeId(), ): void { const state = stateFor(parentRuntimeId); const pending = state.pending; state.pending = []; state.flushScheduled = false; for (const entry of pending) entry.reject(error); }