import { existsSync, readFileSync, rmSync } from "node:fs"; import { isNonEmptyString, isString } from "./type-guards.ts"; const ABORT_MESSAGE = "Aborted while waiting for subagent to finish"; const TERMINAL_SENTINEL = /__SUBAGENT_DONE_(\d+)__/; export interface CompletionResult { reason: "done" | "ping" | "sentinel" | "error"; exitCode: number; ping?: { name: string; message: string }; errorMessage?: string; } export interface CompletionOptions { intervalMs: number; readTerminalTail: () => Promise; inspectPane?: () => Promise; /** Bounded artifact grace after explicit pane disappearance. Default: 500ms. */ paneDisappearanceGraceMs?: number; onPaneInspection?: ( inspection: import("./lifecycle.ts").PaneInspection, observedAt: number, ) => void; sessionFile?: string; onTick?: (elapsedSeconds: number) => void; /** Wait for a local file wake-up or a scheduled reconciliation. */ waitForNextCheck?: (signal: AbortSignal) => Promise<"wake" | "reconcile">; /** Drain local task evidence before any terminal fallback. */ onLocalEvidence?: () => void; } export function interpretExitSidecar(payload: any): CompletionResult { if (payload?.type === "ping") { return { reason: "ping", exitCode: 0, ping: { name: isString(payload.name) ? payload.name : "subagent", message: isString(payload.message) ? payload.message : "", }, }; } if (payload?.type === "error") { const errorMessage = isNonEmptyString(payload.errorMessage) ? payload.errorMessage : "Subagent exited with stopReason=error (no errorMessage in sidecar)."; return { reason: "error", exitCode: 1, errorMessage }; } if (payload?.type === "done") return { reason: "done", exitCode: 0 }; return { reason: "error", exitCode: 1, errorMessage: "Invalid subagent completion sidecar: unsupported payload type.", }; } function consumeExitSidecar( sessionFile: string | undefined, ): CompletionResult | null { if (!sessionFile) return null; const exitFile = `${sessionFile}.exit`; if (!existsSync(exitFile)) return null; try { const result = interpretExitSidecar( JSON.parse(readFileSync(exitFile, "utf8")), ); rmSync(exitFile, { force: true }); return result; } catch { // The child may still be writing the file. Retry on the next polling cycle. return null; } } function terminalExitCode(screen: string): number | null { const match = screen.match(TERMINAL_SENTINEL); return match ? Number.parseInt(match[1], 10) : null; } function completionArtifact( options: CompletionOptions, ): CompletionResult | null { return consumeExitSidecar(options.sessionFile); } async function waitForDelayedSidecar( signal: AbortSignal, options: CompletionOptions, ): Promise { const immediate = completionArtifact(options); if (immediate) return immediate; const graceMs = Math.max(0, options.paneDisappearanceGraceMs ?? 500); const deadline = Date.now() + graceMs; while (Date.now() < deadline) { const remaining = deadline - Date.now(); await abortableDelay(Math.min(25, remaining), signal); const result = completionArtifact(options); if (result) return result; } return null; } function abortableDelay( milliseconds: number, signal: AbortSignal, ): Promise { if (signal.aborted) return Promise.reject(new Error(ABORT_MESSAGE)); return new Promise((resolve, reject) => { const onAbort = () => { clearTimeout(timer); reject(new Error(ABORT_MESSAGE)); }; const timer = setTimeout(() => { signal.removeEventListener("abort", onAbort); resolve(); }, milliseconds); signal.addEventListener("abort", onAbort, { once: true }); }); } export async function waitForCompletion( signal: AbortSignal, options: CompletionOptions, ): Promise { const startedAt = Date.now(); // Coordinated supervision waits for its first shared epoch; legacy callers // retain the immediate terminal probe. let reconcile = !options.waitForNextCheck; for (;;) { if (signal.aborted) throw new Error(ABORT_MESSAGE); const sidecarResult = consumeExitSidecar(options.sessionFile); if (sidecarResult) return sidecarResult; options.onLocalEvidence?.(); if (reconcile) try { const exitCode = terminalExitCode(await options.readTerminalTail()); if (exitCode !== null) { // The shell sentinel can become readable while the child publishes its // authoritative semantic record. Always recheck after the asynchronous // terminal read; a zero exit can race a ping or error just like a failure. options.onLocalEvidence?.(); const racedCompletion = completionArtifact(options); if (racedCompletion) return racedCompletion; if (exitCode !== 0) { const delayedCompletion = await waitForDelayedSidecar( signal, options, ); if (delayedCompletion) return delayedCompletion; } return { reason: "sentinel", exitCode }; } } catch { // Terminal reads are only sentinel/output probes; Herdr status is polled // independently below, even when terminal reads succeed. } if (reconcile && options.inspectPane) { let inspection: import("./lifecycle.ts").PaneInspection; try { inspection = await options.inspectPane(); } catch { inspection = { kind: "unavailable", error: "inspectPane threw" }; } const observedAt = Date.now(); options.onPaneInspection?.(inspection, observedAt); if (inspection.kind === "missing") { // Drain concurrent persistent events before accepting pane fallback. options.onLocalEvidence?.(); // Pane closure and atomic artifact publication are separate operations. // Allow a short bounded grace window before declaring evidence lost. const racedCompletion = await waitForDelayedSidecar(signal, options); if (racedCompletion) return racedCompletion; return { reason: "error", exitCode: 1, errorMessage: "Subagent pane disappeared before completion evidence was recorded.", }; } } options.onTick?.(Math.floor((Date.now() - startedAt) / 1000)); if (options.waitForNextCheck) { reconcile = (await options.waitForNextCheck(signal)) === "reconcile"; } else { await abortableDelay(options.intervalMs, signal); reconcile = true; } } }