import { readCompletionSidecar, validateOperation, type CompletionOperation, type CompletionSidecar, type CompletionSidecarResult, } from "./completion-sidecar.ts"; import { isNativeAgentObservation, type NativeAgentObservation, type NativeSupervisionAdapter, } from "./native-supervision.ts"; import { hasCancellationRequest, readCancellationSidecar, } from "./cancellation-sidecar.ts"; const ABORT_MESSAGE = "Aborted while waiting for subagent to finish"; const COMPLETION_POLL_INTERVAL_MS = 1_000; const DISAPPEARANCE_GRACE_MS = 500; const DISAPPEARANCE_CHECK_INTERVAL_MS = 25; /** Give a transient Windows reader/AV lock a bounded extra window to clear. */ const TRANSIENT_ARTIFACT_GRACE_MS = 5_000; /** Stable machine-readable reasons for a failed Completion outcome. */ export type CompletionFailureCategory = | "child-reported-error" | "missing-evidence" | "invalid-evidence" | "supervision-error"; export type CompletionOutcome = | { kind: "completed"; result?: CompletionSidecarResult } | { kind: "failed"; failureCategory: CompletionFailureCategory; errorMessage: string; result?: CompletionSidecarResult; }; export type CompletionControlOutcome = CompletionOutcome | { kind: "cancelled" }; /** * Parent-facing presentation data derived from a Completion outcome. * * Completion supervision deliberately returns CompletionOutcome. This shape * exists at the delivery boundary because tool results and lifecycle history * still expose exit-code and failure fields to the parent session. */ export interface CompletionResult { reason: "done" | "error"; exitCode: number; errorMessage?: string; failureCategory?: CompletionFailureCategory; } /** Convert the supervision seam's outcome into parent delivery fields. */ export function completionResultForOutcome(outcome: CompletionOutcome): CompletionResult { if (outcome.kind === "completed") return { reason: "done", exitCode: 0 }; return { reason: "error", exitCode: 1, errorMessage: outcome.errorMessage, failureCategory: outcome.failureCategory, }; } export interface CompletionSupervisionOptions extends CompletionOperation { /** * The transport adapter; no pane or Herdr details cross this boundary. * Omit it when only Completion evidence is being supervised; the operation * resolves from the disappearance grace path exactly as if every observation * were missing. */ native?: NativeSupervisionAdapter; /** Receives every native observation as non-authoritative lifecycle input. */ observationSink?: (observation: NativeAgentObservation, observedAt: number) => void; } function outcomeForSidecar(sidecar: CompletionSidecar): CompletionOutcome { const result = sidecar.result ? { result: sidecar.result } : {}; if (sidecar.type === "error") { return { kind: "failed", failureCategory: "child-reported-error", errorMessage: sidecar.errorMessage, ...result, }; } return { kind: "completed", ...result }; } function completionArtifact( operation: CompletionOperation, ): | { kind: "missing" } | { kind: "valid"; outcome: CompletionOutcome } | { kind: "invalid"; error: string } | { kind: "transient"; error: string } { const result = readCompletionSidecar(operation); if (result.kind === "missing") return result; if (result.kind === "invalid" || result.kind === "transient") return result; return { kind: "valid", outcome: outcomeForSidecar(result.sidecar) }; } function invalidEvidenceOutcome(error: string): CompletionOutcome { return { kind: "failed", failureCategory: "invalid-evidence", errorMessage: `Invalid subagent completion sidecar: ${error}`, }; } function missingEvidenceOutcome(): CompletionOutcome { return { kind: "failed", failureCategory: "missing-evidence", errorMessage: "Native agent was missing before completion evidence was recorded.", }; } function supervisionErrorOutcome(error: unknown): CompletionOutcome { return { kind: "failed", failureCategory: "supervision-error", errorMessage: error instanceof Error ? error.message : String(error), }; } function abortableDelay( milliseconds: number, signal: AbortSignal, delay = defaultAbortableDelay, ): Promise { if (signal.aborted) return Promise.reject(new Error(ABORT_MESSAGE)); return delay(milliseconds, signal); } function defaultAbortableDelay(milliseconds: number, signal: AbortSignal): Promise { return new Promise((resolve, reject) => { let settled = false; let timer: ReturnType | undefined; const cleanup = () => { if (timer !== undefined) clearTimeout(timer); signal.removeEventListener("abort", onAbort); }; const onAbort = () => { if (settled) return; settled = true; cleanup(); reject(new Error(ABORT_MESSAGE)); }; timer = setTimeout(() => { if (settled) return; settled = true; cleanup(); resolve(); }, Math.max(0, milliseconds)); signal.addEventListener("abort", onAbort, { once: true }); }); } function abortablePromise(operation: Promise, signal: AbortSignal): Promise { return new Promise((resolve, reject) => { let settled = false; const cleanup = () => signal.removeEventListener("abort", onAbort); const finish = (callback: () => void) => { if (settled) return; settled = true; cleanup(); callback(); }; const onAbort = () => finish(() => reject(new Error(ABORT_MESSAGE))); if (signal.aborted) { onAbort(); return; } signal.addEventListener("abort", onAbort, { once: true }); operation.then( (value) => finish(() => resolve(value)), (error) => finish(() => reject(error)), ); }); } function callObservationSink( options: CompletionSupervisionOptions, observation: NativeAgentObservation, observedAt: number, ): void { // Lifecycle reporting is enrichment. A broken sink must never turn valid // Completion evidence into a supervision failure. try { options.observationSink?.(observation, observedAt); } catch { // Deliberately ignored: Completion evidence is authoritative. } } async function waitForCancellationDrain( signal: AbortSignal, operation: CompletionOperation, ): Promise { for (;;) { const artifact = completionArtifact(operation); if (artifact.kind === "valid") return artifact.outcome; if (artifact.kind === "invalid") return invalidEvidenceOutcome(artifact.error); const cancellation = readCancellationSidecar(operation); if (cancellation.kind === "valid") return { kind: "cancelled" }; if (cancellation.kind === "invalid") return supervisionErrorOutcome(cancellation.error); // A cancellation request is not Completion evidence. Keep the operation // alive until the Child acknowledges its descendant drain while Native // supervision remains unavailable or uncertain. await abortableDelay(DISAPPEARANCE_CHECK_INTERVAL_MS, signal); } } async function waitForDisappearanceArtifacts( signal: AbortSignal, operation: CompletionOperation, ): Promise { const readArtifact = () => completionArtifact(operation); const readControl = (): CompletionControlOutcome | "transient" | null => { const cancellation = readCancellationSidecar(operation); if (cancellation.kind === "valid") return { kind: "cancelled" }; if (cancellation.kind === "invalid") return supervisionErrorOutcome(cancellation.error); if (cancellation.kind === "transient") return "transient"; return null; }; const immediate = readArtifact(); if (immediate.kind === "valid") return immediate.outcome; if (immediate.kind === "invalid") return invalidEvidenceOutcome(immediate.error); const immediateControl = readControl(); if (immediateControl && immediateControl !== "transient") return immediateControl; const deadline = Date.now() + DISAPPEARANCE_GRACE_MS; let transientDeadline = immediate.kind === "transient" || immediateControl === "transient" ? Date.now() + TRANSIENT_ARTIFACT_GRACE_MS : undefined; while (Date.now() < deadline || (transientDeadline != null && Date.now() < transientDeadline)) { const now = Date.now(); const activeDeadlines = [deadline, transientDeadline] .filter((end): end is number => end != null && end > now); const remaining = Math.max(0, Math.min(...activeDeadlines.map((end) => end - now))); await abortableDelay(Math.min(DISAPPEARANCE_CHECK_INTERVAL_MS, remaining), signal); const result = readArtifact(); if (result.kind === "valid") return result.outcome; if (result.kind === "invalid") return invalidEvidenceOutcome(result.error); if (result.kind === "transient" && transientDeadline == null) { transientDeadline = Date.now() + TRANSIENT_ARTIFACT_GRACE_MS; } const control = readControl(); if (control && control !== "transient") return control; if (control === "transient" && transientDeadline == null) { transientDeadline = Date.now() + TRANSIENT_ARTIFACT_GRACE_MS; } } // Include evidence published exactly at the grace boundary. A transient // read is treated as no evidence, but never as malformed evidence. const final = readArtifact(); if (final.kind === "valid") return final.outcome; if (final.kind === "invalid") return invalidEvidenceOutcome(final.error); const finalControl = readControl(); return finalControl === "transient" ? null : finalControl; } async function runCompletionSupervision( signal: AbortSignal, options: CompletionSupervisionOptions, ): Promise { validateOperation(options); const adapter = options.native; for (;;) { if (signal.aborted) throw new Error(ABORT_MESSAGE); // Read before querying Native supervision. The sidecar is the completion // authority; a Herdr observation is only a lifecycle/status observation. let artifact: ReturnType; try { artifact = completionArtifact(options); } catch (error) { return supervisionErrorOutcome(error); } if (artifact.kind === "valid") return artifact.outcome; if (artifact.kind === "invalid") return invalidEvidenceOutcome(artifact.error); const cancellation = readCancellationSidecar(options); if (cancellation.kind === "valid") return { kind: "cancelled" }; if (cancellation.kind === "invalid") return supervisionErrorOutcome(cancellation.error); // Without a native adapter there is no live surface to observe; treat // the operation like a missing observation and let the grace path decide. if (!adapter) { const racedCompletion = await waitForDisappearanceArtifacts(signal, options); if (racedCompletion) return racedCompletion; if (hasCancellationRequest(options)) { // Native disappearance confirms that there is no surface left to // reclaim. Treat the parent-requested cancellation as control flow; // a late acknowledgement cannot be supervised after this boundary. return { kind: "cancelled" }; } return missingEvidenceOutcome(); } let observation: NativeAgentObservation; try { observation = await abortablePromise(adapter.observe(signal), signal); } catch (error) { if (signal.aborted) throw new Error(ABORT_MESSAGE); if (hasCancellationRequest(options)) { return waitForCancellationDrain(signal, options); } // Adapters should translate transient transport loss to an explicit // `unavailable` observation. An unexpected adapter exception is a // supervision failure rather than an indefinitely retried child state. return supervisionErrorOutcome(error); } if (!isNativeAgentObservation(observation)) { return supervisionErrorOutcome("Native supervision returned an invalid observation."); } const observedAt = observation.kind === "present" && observation.observedAt != null ? observation.observedAt : Date.now(); callObservationSink(options, observation, observedAt); // A child may publish evidence while Native supervision is in flight. // Re-read before the fixed poll delay so fresh evidence is delivered // without waiting for another transport query. try { artifact = completionArtifact(options); } catch (error) { return supervisionErrorOutcome(error); } if (artifact.kind === "valid") return artifact.outcome; if (artifact.kind === "invalid") return invalidEvidenceOutcome(artifact.error); const cancellationAfterObservation = readCancellationSidecar(options); if (cancellationAfterObservation.kind === "valid") return { kind: "cancelled" }; if (cancellationAfterObservation.kind === "invalid") { return supervisionErrorOutcome(cancellationAfterObservation.error); } if (observation.kind === "missing") { // Native-agent closure and atomic sidecar publication are separate // operations. A bounded internal grace period lets publication win. const racedCompletion = await waitForDisappearanceArtifacts(signal, options); if (racedCompletion) return racedCompletion; if (hasCancellationRequest(options)) { // Native disappearance confirms that there is no surface left to // reclaim. Treat the parent-requested cancellation as control flow; // a late acknowledgement cannot be supervised after this boundary. return { kind: "cancelled" }; } return missingEvidenceOutcome(); } await abortableDelay(COMPLETION_POLL_INTERVAL_MS, signal); } } /** * Supervise one Completion operation until its authoritative evidence is * available. Polling, disappearance grace, and sidecar handling are private * policy; callers provide only the operation, native adapter, sink, and signal. */ export function superviseCompletion( signal: AbortSignal, options: CompletionSupervisionOptions, ): Promise { return runCompletionSupervision(signal, options); }