import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { buildAdversaryFollowUpTask, buildAdversaryTask, extractUserGoal, formatAdversaryVerdict, type AdversaryVerdict, } from "./adversary.ts"; import { listPendingWorkLabels } from "./background-work.ts"; import { type OrchestrationConfig, type ResolvedOrchestrationConfig, resolveOrchestrationConfig } from "./config.ts"; import { requestAdversaryReview, type ReviewProgress } from "./delegation.ts"; import { formatElapsedStatus, type ReviewStatus } from "./elapsed-status.ts"; import { bootstrapMetaFiles, buildOrchestrationPrompt, type MetaBootstrap, metaDirFor } from "./meta.ts"; export interface OrchestrationRuntime { dispose(): void; } export interface PendingWorkSnapshot { active: boolean; summary?: string; } export interface OrchestrationDependencies { /** * Report whether the Main Agent is waiting on known work. When it is, the * harness schedules a fallback wake instead of starting a review. */ inspectPendingWork?: () => PendingWorkSnapshot; /** Override the adversary launcher. Defaults to the pi-subagents delegation protocol. */ launchAdversary?: (input: { cwd: string; signal: AbortSignal; fresh: boolean }) => Promise; } export interface TimerDriver { setTimeout(callback: () => void, delayMs: number): unknown; clearTimeout(handle: unknown): void; setInterval(callback: () => void, delayMs: number): unknown; clearInterval(handle: unknown): void; } /** * Timers are unref'd so they never keep the Node event loop alive. A referenced * repeating timer makes non-interactive runs (`pi -p`) hang after the answer is * produced instead of exiting. */ function unref(handle: unknown): unknown { (handle as { unref?: () => void })?.unref?.(); return handle; } const defaultTimers: TimerDriver = { setTimeout: (callback, delayMs) => unref(setTimeout(callback, delayMs)), clearTimeout: (handle) => clearTimeout(handle as ReturnType), setInterval: (callback, delayMs) => unref(setInterval(callback, delayMs)), clearInterval: (handle) => clearInterval(handle as ReturnType), }; function lastAssistantStopReason(messages: readonly unknown[] | undefined): string | undefined { if (!Array.isArray(messages)) return undefined; for (let index = messages.length - 1; index >= 0; index -= 1) { const message = messages[index]; if (!message || typeof message !== "object") continue; const candidate = message as { role?: unknown; stopReason?: unknown }; if (candidate.role !== "assistant") continue; return typeof candidate.stopReason === "string" ? candidate.stopReason : undefined; } return undefined; } /** * Register the continuous orchestration harness on a Pi extension API. * * The Main Agent receives a static autonomous-work instruction with no mention * of a reviewer. When it settles with no pending work, an independent adversary * reviews the result and any remaining work returns as a plain user follow-up. */ export function registerContinuousOrchestration( pi: ExtensionAPI, config?: OrchestrationConfig, timers: TimerDriver = defaultTimers, dependencies: OrchestrationDependencies = {}, ): OrchestrationRuntime { const resolved: ResolvedOrchestrationConfig = resolveOrchestrationConfig(config); if (!resolved.enabled) return { dispose() {} }; let disposed = false; let wakeTimer: unknown; let reviewController: AbortController | undefined; let settledVerdict: AdversaryVerdict | undefined; let lastAgentStopReason: string | undefined; let metaBootstrap: MetaBootstrap | undefined; let metaNoticePending = false; let sessionActive = false; let sessionCwd = ""; let uiContext: { cwd: string; hasUI: boolean; ui?: unknown; sessionManager?: unknown } | undefined; // Adversary session reuse: one persistent reviewer across every review. let adversaryRunId: string | undefined; let adversaryParentSessionId: string | undefined; let originalUserGoal: string | undefined; let promptStartedAt: number | undefined; let finishedAt: number | undefined; let elapsedTimer: unknown; let reviewStatus: ReviewStatus | undefined; const setStatus = (value: string | undefined) => { const ctx = uiContext as | { hasUI?: boolean; ui?: { setStatus?: (key: string, value: string | undefined) => void } } | undefined; if (!ctx?.hasUI) return; try { ctx.ui?.setStatus?.("orchestration-elapsed", value); } catch { // The interactive UI may disappear while the session is shutting down. } }; const renderElapsed = () => { if (promptStartedAt === undefined) return; setStatus( formatElapsedStatus({ startedAt: promptStartedAt, now: Date.now(), ...(finishedAt !== undefined ? { finishedAt } : {}), ...(reviewStatus ? { review: reviewStatus } : {}), }), ); }; const stopElapsedTimer = () => { if (elapsedTimer === undefined) return; timers.clearInterval(elapsedTimer); elapsedTimer = undefined; }; const hasUI = (): boolean => Boolean((uiContext as { hasUI?: boolean } | undefined)?.hasUI); const ensureElapsedTimer = () => { if (elapsedTimer !== undefined || finishedAt !== undefined) return; // The ticking clock only exists to repaint a live status line. Headless // runs (`pi -p`) have nothing to repaint, so starting a repeating timer // there would only add wakeups. if (!hasUI()) return; elapsedTimer = timers.setInterval(renderElapsed, 1_000); }; const clearElapsed = () => { stopElapsedTimer(); setStatus(undefined); promptStartedAt = undefined; finishedAt = undefined; reviewStatus = undefined; }; const beginElapsed = () => { if (promptStartedAt === undefined || finishedAt !== undefined) { promptStartedAt = Date.now(); finishedAt = undefined; } reviewStatus = undefined; ensureElapsedTimer(); renderElapsed(); }; const finishElapsed = () => { if (promptStartedAt === undefined) return; finishedAt = Date.now(); reviewStatus = undefined; stopElapsedTimer(); renderElapsed(); }; const cancelWake = () => { if (wakeTimer === undefined) return; timers.clearTimeout(wakeTimer); wakeTimer = undefined; }; const cancelReview = () => { const controller = reviewController; reviewController = undefined; controller?.abort(); }; const clearSettledVerdict = () => { settledVerdict = undefined; }; const currentSessionId = (): string | undefined => { const manager = (uiContext as { sessionManager?: { getSessionId?: () => string | null } } | undefined) ?.sessionManager; return manager?.getSessionId?.() ?? undefined; }; const inspectPendingWork = (): PendingWorkSnapshot => { try { if (dependencies.inspectPendingWork) return dependencies.inspectPendingWork(); const sessionId = currentSessionId(); if (!sessionId) { return { active: true, summary: "Pending-work state unavailable: no active Pi session ID." }; } const labels = listPendingWorkLabels(sessionId); // No registry means no provider ever registered background work, which is // the normal state for a plain session rather than a lost signal. if (labels === undefined) return { active: false }; return { active: labels.length > 0, ...(labels.length > 0 ? { summary: `Pending work: ${labels.join(", ")}` } : {}), }; } catch (error) { // Fail closed: an unreadable pending-work state must not be treated as // "nothing is running", or a review could start mid-flight. return { active: true, summary: `Pending-work inspection failed: ${error instanceof Error ? error.message : String(error)}`, }; } }; const schedulePendingWake = (pending: PendingWorkSnapshot) => { cancelWake(); if (disposed || !sessionActive) return; wakeTimer = timers.setTimeout(() => { wakeTimer = undefined; if (disposed || !sessionActive) return; const current = inspectPendingWork(); const summary = current.summary ?? pending.summary; void pi.sendMessage( { customType: "orchestration-pending-wake", content: current.active ? `Known work is still pending. Inspect its current state and handle any stall or result.${summary ? `\n\n${summary}` : ""}` : `Previously pending work is no longer active. Inspect its result and continue.${summary ? `\n\n${summary}` : ""}`, display: false, details: { reason: "pending-work", idleWakeMs: resolved.idleWakeMs }, }, { triggerTurn: true }, ); }, resolved.idleWakeMs); }; const defaultLaunchAdversary = async (input: { cwd: string; signal: AbortSignal; fresh: boolean; }): Promise => { const sessionManager = (uiContext as { sessionManager?: { getSessionId?: () => string | null; getBranch?: () => readonly unknown[] } } | undefined) ?.sessionManager; const parentSessionId = sessionManager?.getSessionId?.() ?? undefined; if (originalUserGoal === undefined && sessionManager?.getBranch) { originalUserGoal = extractUserGoal(sessionManager.getBranch()); } // The public delegation protocol cannot resume a prior child session, so // every review runs a fresh reviewer. After the first review we still send // the incremental follow-up task, which points the reviewer at the current // notes and changed evidence instead of repeating a full review. const isFollowUp = !input.fresh && Boolean(adversaryRunId) && adversaryParentSessionId === parentSessionId; const startedAt = Date.now(); reviewStatus = { startedAt, session: "fresh" }; promptStartedAt ??= startedAt; finishedAt = undefined; ensureElapsedTimer(); renderElapsed(); const onProgress = (progress: ReviewProgress) => { if (!reviewStatus) return; reviewStatus = { ...reviewStatus, ...(progress.turns !== undefined ? { turns: progress.turns } : {}), ...(progress.toolCount !== undefined ? { toolCount: progress.toolCount } : {}), ...(progress.tokens !== undefined ? { tokens: progress.tokens } : {}), ...(progress.currentTool ? { currentTool: progress.currentTool } : {}), }; renderElapsed(); }; const review = await requestAdversaryReview(pi, { agent: resolved.adversaryAgent, task: isFollowUp ? buildAdversaryFollowUpTask() : buildAdversaryTask(originalUserGoal), cwd: input.cwd, signal: input.signal, ackTimeoutMs: resolved.delegationAckTimeoutMs, onProgress, }); if (review.runId) { adversaryRunId = review.runId; adversaryParentSessionId = parentSessionId; } reviewStatus = undefined; return review.verdict; }; const launchAdversarialReview = async (): Promise => { if (disposed || !sessionActive || reviewController || !sessionCwd) return; const launch = dependencies.launchAdversary ?? defaultLaunchAdversary; const controller = new AbortController(); reviewController = controller; try { const result = await launch({ cwd: sessionCwd, signal: controller.signal, fresh: resolved.freshAdversaryEachReview }); if (disposed || !sessionActive || reviewController !== controller || controller.signal.aborted) return; reviewController = undefined; if (result.verdict === "complete") { settledVerdict = result; finishElapsed(); return; } void pi.sendUserMessage(formatAdversaryVerdict(result), { deliverAs: "followUp" }); } catch (error) { if (disposed || !sessionActive || reviewController !== controller || controller.signal.aborted) return; reviewController = undefined; finishElapsed(); // This runs only from agent_settled. Explicitly disable turn triggering // so verifier infrastructure failures cannot steer the agent into a loop. void pi.sendMessage( { customType: "orchestration-verification-unavailable", content: `Independent verification could not complete: ${error instanceof Error ? error.message : String(error)}`, display: true, details: { reason: "verification-unavailable" }, }, { triggerTurn: false }, ); } }; pi.on("session_start", (_event, ctx) => { if (disposed) return; sessionActive = true; sessionCwd = ctx.cwd; uiContext = ctx as unknown as typeof uiContext; metaBootstrap = bootstrapMetaFiles(ctx.cwd); metaNoticePending = true; clearElapsed(); cancelWake(); cancelReview(); clearSettledVerdict(); lastAgentStopReason = undefined; originalUserGoal = undefined; }); pi.on("before_agent_start", (event, ctx) => { if (disposed) return; uiContext = ctx as unknown as typeof uiContext; if (!metaBootstrap || sessionCwd !== ctx.cwd) { metaBootstrap = bootstrapMetaFiles(ctx.cwd); metaNoticePending = true; sessionCwd = ctx.cwd; } const notice = metaNoticePending ? metaBootstrap : undefined; metaNoticePending = false; const metaDir = metaDirFor(ctx.cwd); const createdPaths = notice?.created.map((name) => `\`${metaDir}/${name}\``) ?? []; const existingPaths = notice?.existing.map((name) => `\`${metaDir}/${name}\``) ?? []; const noticeParts = [ createdPaths.length > 0 ? `Durable work files created: ${createdPaths.join(", ")}.` : "", existingPaths.length > 0 ? `Existing durable work files: ${existingPaths.join(", ")}.` : "", "Use these files as durable memory when useful. Work head down until no known task remains. Before ending your work, update them with the actual result and evidence. You may create additional clearly named Markdown files in the same directory when useful for note-taking.", ].filter(Boolean); return { systemPrompt: `${event.systemPrompt}\n\n${buildOrchestrationPrompt(ctx.cwd)}`, ...(notice ? { message: { customType: "autonomous-work", content: noticeParts.join(" "), display: false, }, } : {}), }; }); pi.on("input", (event) => { if (disposed) return; // A real user interruption cancels review and restarts the elapsed clock. if (event.source !== "extension") { cancelWake(); cancelReview(); clearSettledVerdict(); beginElapsed(); } }); pi.on("agent_start", () => { if (disposed) return; cancelWake(); cancelReview(); clearSettledVerdict(); }); pi.on("agent_end", (event, ctx) => { if (disposed) return; sessionCwd = ctx.cwd; uiContext = ctx as unknown as typeof uiContext; lastAgentStopReason = lastAssistantStopReason(event.messages); }); pi.on("agent_settled", async (_event, ctx) => { if (disposed || !sessionActive) return; sessionCwd = ctx.cwd; uiContext = ctx as unknown as typeof uiContext; // The loop keeps a live session working: it answers, reviews, then feeds // remaining work back as another user turn. A one-shot run (`pi -p`) has no // session to continue, so driving more turns there would stall the run. if (!hasUI()) return; const stopReason = lastAgentStopReason; lastAgentStopReason = undefined; // Escape and provider/agent failures are user or host recovery outcomes, // not completion claims. Never review or relaunch work after them. if (stopReason === "aborted" || stopReason === "error") { finishElapsed(); return; } const pending = inspectPendingWork(); if (pending.active) { schedulePendingWake(pending); return; } // agent_settled fires after Pi has exhausted retry, compaction, and queued // continuation handling, so the review sees the stable final workspace. await launchAdversarialReview(); const result = settledVerdict; settledVerdict = undefined; if (!result || disposed || !sessionActive) return; void pi.sendMessage( { customType: "autonomous-completion", content: formatAdversaryVerdict(result), display: true, details: { verdict: result.verdict }, }, { triggerTurn: false }, ); }); pi.on("session_shutdown", () => { if (disposed) return; sessionActive = false; metaBootstrap = undefined; metaNoticePending = false; lastAgentStopReason = undefined; clearElapsed(); cancelWake(); cancelReview(); clearSettledVerdict(); }); return { dispose() { if (disposed) return; disposed = true; sessionActive = false; metaBootstrap = undefined; metaNoticePending = false; lastAgentStopReason = undefined; clearElapsed(); cancelWake(); cancelReview(); clearSettledVerdict(); }, }; }