// Completion watch: the parent-side arc that drives an adopted Completion // operation through Completion supervision, Completion handoff, and one // Completion delivery attempt, ending in registry release. Foreground and // background calls cross the same arc as Completion delivery destinations; // presentation and the running registry remain parent-runtime concerns. // See CONTEXT.md "Completion watch" and "Completion delivery destination". import { closeHerdrSurface, inspectHerdrAgent, sendHerdrAgentEscape } from "./herdr.ts"; import { superviseCompletion, type CompletionControlOutcome } from "./completion.ts"; import type { CompletionDeliveryManager } from "./completion-delivery.ts"; import { createNativeSupervisionAdapter } from "./native-supervision.ts"; import type { NativeAgentObservation } from "./native-supervision.ts"; import type { CompletionOperation } from "./completion-sidecar.ts"; import { createOperationReference } from "./operation-identity.ts"; import { buildCompletionSteerMessage, buildForegroundSubagentResult, cancelledCompletionHandoff, failedCompletionHandoff, prepareCompletionHandoff, type CompletionHandoffContext, type CompletionHandoffResult, type CompletionSteerMessage, type SubagentToolResult, } from "./completion-handoff.ts"; import { attachWatcher, interruptAndCloseSubagent, noteOutcomeObserved, observeSupervision, operationForRunningSubagent, requestSubagentCancellation, resolveCancelledWatchCleanup, settleOutcome, type RunningSubagent, } from "./running-registry.ts"; import { freezeRuntimePlan } from "./runtime-routing.ts"; import type { DescendantHandoffRequest } from "./nested-lifecycle.ts"; /** Parent-runtime collaborators the Completion watch arc crosses once per parent runtime. */ export interface CompletionWatchRuntime { /** Parent-wide owner of the one-shot Completion delivery attempt. */ manager: CompletionDeliveryManager; /** Refresh live presentation after lifecycle or activity changes. */ onLiveStateChange(): void; /** Admit a handoff for delivery; nested-lifecycle cancellation may suppress it. */ admitHandoff(running: RunningSubagent, result: CompletionHandoffResult): DescendantHandoffRequest | void; /** Settle the handoff after its single delivery attempt (consumed or lost). */ settleDelivery( running: RunningSubagent, request: DescendantHandoffRequest, consumed: boolean, ): void; /** Report a destination failure without changing the Child Completion outcome. */ onDeliveryFailure(running: RunningSubagent, result: CompletionHandoffResult, error: string): void; /** Remove the entry from the parent registry and refresh presentation. */ release(running: RunningSubagent): void; } /** The Direct parent's current agent loop: the tool call's own result and signal. */ export interface ForegroundCompletionDestination { kind: "foreground"; /** Parent tool-call signal; abort maps to watcher cancellation or nested cascade. */ signal: AbortSignal; } /** A continuation steered into the parent conversation. */ export interface BackgroundCompletionDestination { kind: "background"; /** Deliver a background Completion steer to the parent conversation. */ steer(payload: CompletionSteerMessage): unknown; } /** Where the one Completion delivery attempt goes. See CONTEXT.md "Completion delivery destination". */ export type CompletionWatchDestination = | ForegroundCompletionDestination | BackgroundCompletionDestination; /** Native transport seam for Completion watch; tests substitute a deterministic adapter. */ export interface NativeWatchAdapter { inspect(surface: string): Promise; close(surface: string): void; interrupt(surface: string): void; } const productionNativeWatchAdapter: NativeWatchAdapter = { inspect: inspectHerdrAgent, close: closeHerdrSurface, interrupt: sendHerdrAgentEscape, }; export function createCompletionHandoffContext( running: RunningSubagent, ): CompletionHandoffContext { const runtimePlan = freezeRuntimePlan(running.runtimePlan); return Object.freeze({ operationId: running.id, name: running.name, task: running.task, ...(running.agent ? { agent: running.agent } : {}), artifacts: running.artifacts, startedAt: running.startTime, ...(runtimePlan ? { runtimePlan } : {}), }); } /** * Watch a launched subagent until it exits. Completion supervision and * lifecycle remain in the parent runtime; Completion handoff owns enrichment. */ export async function watchCompletionOperation( running: RunningSubagent, signal: AbortSignal, runtime: CompletionWatchRuntime, native: NativeWatchAdapter = productionNativeWatchAdapter, ): Promise { const { surface, startTime } = running; const operation = operationForRunningSubagent(running); const handoffContext = createCompletionHandoffContext(running); try { const outcome: CompletionControlOutcome = await superviseCompletion(signal, { ...operation, native: createNativeSupervisionAdapter(() => native.inspect(surface)), observationSink: (observation, observedAt) => { // One tick: native projection, activity enrichment, lifecycle fold. observeSupervision(running, observation, observedAt); runtime.onLiveStateChange(); }, }); const detectedAt = Date.now(); noteOutcomeObserved(running, outcome, detectedAt); if (outcome.kind === "cancelled") { try { native.close(surface); } catch {} return cancelledCompletionHandoff( handoffContext, Math.floor((detectedAt - startTime) / 1000), running.runtimePlan, ); } runtime.onLiveStateChange(); const elapsed = Math.floor((detectedAt - startTime) / 1000); // Completion evidence is authoritative. Surface cleanup is best effort and // must not replace a valid completed or failed outcome. Close // before preparing the sidecar-backed result so a completed surface cannot // keep consuming a Herdr slot while delivery continues. try { native.close(surface); } catch {} const result = prepareCompletionHandoff(handoffContext, outcome, elapsed); settleOutcome(running, result); return result; } catch (err: any) { if (signal.aborted) { // Cancellation is control flow, not a child failure. The registry's // cleanup policy decides whether this watch still owns the surface. if (resolveCancelledWatchCleanup(running) === "close-surface") { interruptAndCloseSubagent(running, native.interrupt, native.close); } } else { try { native.close(surface); } catch {} } if (signal.aborted) { // A cancelled watch changed no lifecycle state; refresh as-is. runtime.onLiveStateChange(); return cancelledCompletionHandoff( handoffContext, Math.floor((Date.now() - startTime) / 1000), running.runtimePlan, ); } const failure = { errorMessage: err?.message ?? String(err), failureCategory: "supervision-error" as const, }; const result = failedCompletionHandoff( handoffContext, failure, Math.floor((Date.now() - startTime) / 1000), ); // The failed handoff settles the lifecycle with the same supervision // error it reports, keeping one writer for terminal lifecycle states; // the refresh must observe the settled lifecycle, so it follows. settleOutcome(running, result); runtime.onLiveStateChange(); return result; } } function operationForHandoffContext(context: CompletionHandoffContext): CompletionOperation { if (!context.artifacts) { throw new Error("Completion handoff requires operation artifacts."); } return createOperationReference(context.operationId, context.artifacts); } function reportDeliveryFailure( runtime: CompletionWatchRuntime, running: RunningSubagent, result: CompletionHandoffResult, error: string, ): void { try { runtime.onDeliveryFailure(running, result, error); } catch { // Delivery failure reporting is parent-side enrichment; it must not prevent // the operation from being released after its one delivery attempt. } } function linkAbortSignal( source: AbortSignal, onAbort: () => void, ): () => void { const abort = () => onAbort(); if (source.aborted) onAbort(); else source.addEventListener("abort", abort, { once: true }); return () => source.removeEventListener("abort", abort); } /** Terminal status of the one delivery attempt, normalized for the arc. */ type WatchDeliveryStatus = "delivered" | "suppressed" | "cancelled" | "failed" | "skipped"; interface WatchDeliveryOutcome { status: WatchDeliveryStatus; /** Delivery failure error; reported without changing the Child Completion outcome. */ failure?: string; } /** One delivery attempt's outcome plus the foreground tool result it produced. */ interface WatchDelivery { outcome: WatchDeliveryOutcome; toolResult?: SubagentToolResult; } /** The one foreground delivery attempt through the delivery manager. */ async function deliverForegroundHandoff( context: CompletionHandoffContext, result: CompletionHandoffResult, signal: AbortSignal, manager: CompletionDeliveryManager, ): Promise { const handoff = Object.freeze(buildForegroundSubagentResult(context, result)); const delivery = await manager.deliver( { operation: operationForHandoffContext(context), payload: handoff, destination: () => handoff as SubagentToolResult, signal, }, ); if (delivery.status === "delivered") { return { outcome: { status: "delivered" }, toolResult: delivery.value }; } if (delivery.status === "cancelled") { return { outcome: { status: "cancelled" }, toolResult: buildForegroundSubagentResult( context, cancelledCompletionHandoff(context, result.elapsed, result.runtimePlan), ), }; } if (delivery.status === "failed") { return { outcome: { status: "failed", failure: delivery.error }, toolResult: { ...handoff, content: [{ type: "text", text: `${handoff.content[0]?.text ?? ""}\n\nCompletion delivery failed: ${delivery.error}`, }], details: { ...handoff.details, deliveryError: delivery.error }, }, }; } return { outcome: { status: "suppressed" }, // Suppression is the parent-side cancellation cutover. A foreground // caller must not receive the completed Child payload after that cutover. toolResult: buildForegroundSubagentResult( context, cancelledCompletionHandoff(context, result.elapsed, result.runtimePlan), ), }; } /** The one background delivery attempt through the delivery manager. */ async function deliverBackgroundHandoff( context: CompletionHandoffContext, result: CompletionHandoffResult, steer: (payload: CompletionSteerMessage) => unknown, manager: CompletionDeliveryManager, ): Promise { const payload = Object.freeze(buildCompletionSteerMessage(context, result)); const delivery = await manager.deliver({ operation: operationForHandoffContext(context), payload, destination: (delivered) => steer(delivered), }); if (delivery.status === "failed") return { status: "failed", failure: delivery.error }; return { status: delivery.status }; } /** The foreground destination resolves with a tool result; background resolves without one. */ type DeliverWatchedResult = D extends ForegroundCompletionDestination ? SubagentToolResult : undefined; /** * The one Completion watch arc: supervision, admission, the single delivery * attempt, settlement, and registry release. Foreground and background calls * differ only in their destination adapter; the handoff admission and its * settlement live in the nested lifecycle, not here. */ async function deliverWatchedCompletion( running: RunningSubagent, runtime: CompletionWatchRuntime, destination: D, watcher: AbortController, native: NativeWatchAdapter, ): Promise> { const result = await watchCompletionOperation(running, watcher.signal, runtime, native); const cancelledResult = result.error === "cancelled"; const request = cancelledResult ? "accepted" : runtime.admitHandoff(running, result) ?? "accepted"; // A cancelled watch or an unadmitted handoff never crosses the delivery // manager; the foreground destination still answers its waiting tool call. const deliverable = !cancelledResult && request === "accepted"; const context = createCompletionHandoffContext(running); let delivery: WatchDelivery; if (destination.kind === "foreground") { const payload = deliverable ? result : cancelledCompletionHandoff(context, result.elapsed, result.runtimePlan); delivery = deliverable ? await deliverForegroundHandoff(context, payload, destination.signal, runtime.manager) : { outcome: { status: "skipped" }, toolResult: buildForegroundSubagentResult(context, payload) }; } else { if (cancelledResult && running.userCancelRequested) { // A user-requested cancellation (the cancel tool) terminates in a // background steer: the tool acknowledged locally, the arc reports the // terminal cancellation once the child's cancellation acknowledgement is // observed. Derived cascade and shutdown cancellations remain control // flow and deliver nothing. await deliverBackgroundHandoff(context, result, destination.steer, runtime.manager); delivery = { outcome: { status: "cancelled" } }; } else { delivery = deliverable ? { outcome: await deliverBackgroundHandoff(context, result, destination.steer, runtime.manager) } : { outcome: { status: "skipped" } }; } } const { outcome, toolResult } = delivery; if (!cancelledResult) { // Cancellation is control flow: no handoff was admitted, so there is // nothing to settle, deliver, or report. if (outcome.status === "failed") { reportDeliveryFailure(runtime, running, result, outcome.failure ?? "delivery failed"); } // The handoff's own lifecycle consumes a delivered result or rejects a // lost one; a duplicate admission leaves it standing. runtime.settleDelivery( running, request, outcome.status === "delivered", ); } // The watch arc releases after Completion evidence or cancellation evidence; // a foreground caller may already have received its cancellation result. runtime.release(running); // The tool result is only ever unset for background destinations. return toolResult as DeliverWatchedResult; } /** * Wait for a Foreground subagent call: the tool call's own Promise resolves * with the destination's result. The parent signal is translated here, at * the destination adapter — it returns control promptly while the watcher's * structured cancellation drain remains active. */ export async function waitForForegroundCompletion( running: RunningSubagent, destination: ForegroundCompletionDestination, runtime: CompletionWatchRuntime, native: NativeWatchAdapter = productionNativeWatchAdapter, ): Promise { const watcher = new AbortController(); attachWatcher(running, watcher); let resolveCancellation!: (result: SubagentToolResult) => void; const cancellation = new Promise((resolve) => { resolveCancellation = resolve; }); const unlinkAbort = linkAbortSignal( destination.signal, () => { // Every adopted record cancels through the structured lifecycle path; // the Descendant cancellation cascade owns surface reclamation while // the Child self-terminates off its cancellation request. Return the // foreground cancellation result immediately, but keep this watcher // alive so the Completion arc can observe the Child's acknowledgement // before releasing its operation artifacts and registry entry. requestSubagentCancellation(running, native.interrupt); const context = createCompletionHandoffContext(running); resolveCancellation(buildForegroundSubagentResult( context, cancelledCompletionHandoff( context, Math.max(0, Math.floor((Date.now() - running.startTime) / 1000)), running.runtimePlan, ), )); }, ); try { // A parent-requested foreground cancellation is control flow, not a // reason to abort Completion supervision. The race lets the tool return // promptly while the same watch arc continues in the background and // releases only after cancellation evidence or terminal Native failure. const watch = deliverWatchedCompletion(running, runtime, destination, watcher, native); return await Promise.race([watch, cancellation]); } finally { unlinkAbort(); } } /** * Start the one background Completion operation used by Fresh Subagents. * Foreground calls wait on the same arc above; only the destination differs. */ export function startBackgroundCompletionWatch( running: RunningSubagent, destination: BackgroundCompletionDestination, runtime: CompletionWatchRuntime, native: NativeWatchAdapter = productionNativeWatchAdapter, ): void { const watcher = new AbortController(); attachWatcher(running, watcher); void deliverWatchedCompletion(running, runtime, destination, watcher, native); }