/** * Extension loaded into sub-agents. * - Records child activity for parent-side lifecycle status * - Publishes Completion evidence automatically at Pi's `agent_settled` boundary */ import { SUBAGENT_ENV_ACTIVITY_FILE, SUBAGENT_ENV_CANCELLATION_FILE, SUBAGENT_ENV_CANCELLATION_REQUEST_FILE, SUBAGENT_ENV_COMPLETION_FILE, SUBAGENT_ENV_OPERATION_ID, } from "./operation-env.ts"; import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; import { createSubagentActivityRecorder } from "./activity.ts"; import { publishCompletionSidecar, type CompletionOperation, type CompletionSidecarOutcome, type CompletionSidecarRuntime, } from "./completion-sidecar.ts"; import { isThinkingLevel, type ThinkingLevel } from "./runtime-routing.ts"; import { createOperationArtifactsReference } from "./operation-artifacts.ts"; import { isTransientFileError } from "./file-retry.ts"; import { createOperationReference } from "./operation-identity.ts"; import { hasCancellationRequest, publishCancellationSidecar, } from "./cancellation-sidecar.ts"; import { getNestedLifecycleCoordinator, type CompletionPublicationRequest, } from "./nested-lifecycle.ts"; const MAX_CANCELLATION_ACK_ATTEMPTS = 3; const MAX_COMPLETION_RETRY_ATTEMPTS = 8; const COMPLETION_RETRY_DELAY_MS = 100; export interface SubagentErrorInfo { errorMessage: string; stopReason: "error"; } /** * If the last assistant message in the turn ended with `stopReason: "error"` * (typically auto-retry exhausted on an overload / rate limit / server error), * return its error info so the parent orchestrator can surface a clear * failure instead of silently treating the run as completed. * * Returns `null` when the latest assistant turn completed normally or was * aborted by the user (handled separately by the legacy predicate above). */ export function findLatestAssistantError( messages: any[] | undefined, ): SubagentErrorInfo | null { if (!messages) return null; for (let i = messages.length - 1; i >= 0; i--) { const msg = messages[i]; if (msg?.role !== "assistant") continue; if (msg.stopReason !== "error") return null; const raw = typeof msg.errorMessage === "string" ? msg.errorMessage.trim() : ""; return { errorMessage: raw || "Subagent agent loop ended with stopReason=error (no errorMessage field).", stopReason: "error", }; } return null; } function latestAssistantMessage(messages: any[] | undefined): any | undefined { if (!messages) return undefined; for (let i = messages.length - 1; i >= 0; i--) { if (messages[i]?.role === "assistant") return messages[i]; } return undefined; } function assistantText(message: any | undefined): string | undefined { const content = message?.content; const textBlocks = typeof content === "string" ? [content] : Array.isArray(content) ? content .filter((block: any) => block?.type === "text" && typeof block.text === "string") .map((block: any) => block.text as string) : []; const summary = textBlocks.join("\n").trim(); return summary || undefined; } function assistantSummary(messages: any[] | undefined, latest: any | undefined): string { const text = assistantText(latest); if (text) return text; const errorMessage = typeof latest?.errorMessage === "string" ? latest.errorMessage.trim() : ""; if (latest?.stopReason === "error") { return `Subagent error: ${errorMessage || "agent loop ended with stopReason=error (no errorMessage field)."}`; } if (messages) { for (let i = messages.length - 1; i >= 0; i--) { const message = messages[i]; if (message?.role !== "assistant") continue; const summary = assistantText(message); if (summary) return summary; } } return "Sub-agent exited without output"; } function observedRuntime( message: any | undefined, thinking: ThinkingLevel | undefined, ): CompletionSidecarRuntime | undefined { const runtime: CompletionSidecarRuntime = {}; if (typeof message?.provider === "string" && message.provider.trim()) { runtime.provider = message.provider; } const model = typeof message?.responseModel === "string" && message.responseModel.trim() ? message.responseModel : typeof message?.model === "string" && message.model.trim() ? message.model : typeof message?.modelId === "string" && message.modelId.trim() ? message.modelId : undefined; if (model) runtime.modelId = model; const messageThinking = typeof message?.thinkingLevel === "string" ? message.thinkingLevel : typeof message?.thinking === "string" ? message.thinking : undefined; const actualThinking = thinking ?? messageThinking; if (actualThinking && isThinkingLevel(actualThinking)) runtime.thinking = actualThinking; return Object.keys(runtime).length > 0 ? runtime : undefined; } function childThinkingLevel(pi: ExtensionAPI): ThinkingLevel | undefined { try { const thinking = pi.getThinkingLevel(); return isThinkingLevel(thinking) ? thinking : undefined; } catch { return undefined; } } /** Build the self-contained terminal publication for the latest child run. */ export function buildCompletionSidecar( messages: any[] | undefined, thinking?: ThinkingLevel, ): CompletionSidecarOutcome { const errorInfo = findLatestAssistantError(messages); const latest = latestAssistantMessage(messages); const runtime = observedRuntime(latest, thinking); const result = { summary: assistantSummary(messages, latest), ...(runtime ? { runtime } : {}), }; return errorInfo ? { type: "error", ...errorInfo, result } : { type: "done", result }; } export type SettledRunOutcome = "done" | "error" | "aborted"; /** Reconstruct the operation reference from the Ephemeral Child environment. */ function operationFromEnvironment(): CompletionOperation | undefined { const operationId = process.env[SUBAGENT_ENV_OPERATION_ID]?.trim(); if (!operationId) return undefined; const completion = process.env[SUBAGENT_ENV_COMPLETION_FILE]?.trim(); const activity = process.env[SUBAGENT_ENV_ACTIVITY_FILE]?.trim(); const cancellationRequest = process.env[SUBAGENT_ENV_CANCELLATION_REQUEST_FILE]?.trim(); const cancellation = process.env[SUBAGENT_ENV_CANCELLATION_FILE]?.trim(); if (!completion || !activity || !cancellationRequest || !cancellation) return undefined; const artifacts = createOperationArtifactsReference(operationId, { completion, activity, cancellationRequest, cancellation, }); return createOperationReference(operationId, artifacts); } /** Classify the last low-level run at Pi's settled boundary. */ export function classifySettledRun(messages: any[] | undefined): SettledRunOutcome { if (messages) { for (let i = messages.length - 1; i >= 0; i--) { const message = messages[i]; if (message?.role !== "assistant") continue; return message.stopReason === "aborted" ? "aborted" : message.stopReason === "error" ? "error" : "done"; } } return "done"; } function publishSubagentCompletion(outcome: CompletionSidecarOutcome): void { const operation = operationFromEnvironment(); if (!operation) { throw new Error("Subagent completion publication requires an operation identity."); } publishCompletionSidecar(operation, outcome); } export default function (pi: ExtensionAPI) { const recorder = createSubagentActivityRecorder({ runningChildId: process.env[SUBAGENT_ENV_OPERATION_ID], activityFile: process.env[SUBAGENT_ENV_ACTIVITY_FILE], // The parent allocator creates the private namespace before Native // startup. Never recreate it after release if a late child event arrives. createParentDirectory: false, }); let latestRunMessages: any[] | undefined; let completionPublished = false; let cancellationAcknowledgementStarted = false; let cancellationAcknowledgementAttempts = 0; let cancellationWatch: ReturnType | undefined; let completionRetry: ReturnType | undefined; let completionRetryAttempts = 0; const nestedLifecycle = getNestedLifecycleCoordinator(); function stopCompletionRetry(): void { if (!completionRetry) return; clearTimeout(completionRetry); completionRetry = undefined; } function scheduleCompletionRetry(outcome: CompletionSidecarOutcome): void { if ( completionPublished || completionRetry || completionRetryAttempts >= MAX_COMPLETION_RETRY_ATTEMPTS ) return; completionRetryAttempts += 1; completionRetry = setTimeout(() => { completionRetry = undefined; if (completionPublished || nestedLifecycle.isCancellationRequested()) return; try { const publication = publishCompletionOnce(outcome); if (publication === "published" || publication === "already-requested") { recorder.agentSettledDone(); } } catch (error) { // A transient lock is retried asynchronously so an agent_settled event // cannot lose otherwise valid Completion evidence. Permanent protocol // failures remain bounded and are reported by parent supervision. if (isTransientFileError(error)) scheduleCompletionRetry(outcome); } }, COMPLETION_RETRY_DELAY_MS); (completionRetry as any).unref?.(); } function publishCompletionOnce(outcome: CompletionSidecarOutcome): CompletionPublicationRequest { if (completionPublished) return "already-requested"; return nestedLifecycle.requestCompletion(() => { try { publishSubagentCompletion(outcome); completionPublished = true; stopCompletionRetry(); } catch (error) { if (isTransientFileError(error)) scheduleCompletionRetry(outcome); throw error; } }); } function acknowledgeCancellation(): void { const operation = operationFromEnvironment(); if (!operation || cancellationAcknowledgementStarted) return; const request = nestedLifecycle.requestCancellation(); if (request === "already-completed") return; stopCompletionRetry(); if (cancellationAcknowledgementAttempts >= MAX_CANCELLATION_ACK_ATTEMPTS) return; cancellationAcknowledgementAttempts += 1; cancellationAcknowledgementStarted = true; void nestedLifecycle.waitForDescendantDrain().then(() => { try { publishCancellationSidecar(operation); } catch (error) { // Keep the cancellation state active when an acknowledgement cannot // be written. A short-lived contention gets a small number of bounded // publication attempts; a persistent lock or protocol failure does // not re-arm the watcher indefinitely. if ( isTransientFileError(error) && cancellationAcknowledgementAttempts < MAX_CANCELLATION_ACK_ATTEMPTS ) { cancellationAcknowledgementStarted = false; } } }); } /** * A Child can be self-settled while it waits for a descendant. In that * state no new Pi turn is guaranteed, so agent_settled/session_shutdown may * never observe a parent cancellation request. Poll the durable request * while the Child process is alive; acknowledgement still waits for the * nested coordinator's drain before publishing cancellation evidence. */ function startCancellationWatch(): void { if (!operationFromEnvironment() || cancellationWatch) return; cancellationWatch = setInterval(() => { const operation = operationFromEnvironment(); if (operation && hasCancellationRequest(operation)) acknowledgeCancellation(); }, 50); (cancellationWatch as any).unref?.(); } function stopCancellationWatch(): void { if (!cancellationWatch) return; clearInterval(cancellationWatch); cancellationWatch = undefined; } pi.on("session_start", () => { recorder.sessionStart(); startCancellationWatch(); }); pi.on("input", () => { recorder.input(); }); pi.on("before_agent_start", () => { recorder.beforeAgentStart(); }); pi.on("agent_start", () => { // A descendant result may have queued a new turn. Consume all results // accepted since the previous turn as one Descendant continuation batch; // the new run must settle before a deferred parent Completion can win. nestedLifecycle.continuationStarted(); // A new low-level run supersedes any earlier run that may have ended in a // retry or compaction retry. Only the latest run may be classified at // agent_settled. latestRunMessages = undefined; recorder.agentStart(); }); pi.on("agent_end", (event) => { // agent_end is an intermediate low-level boundary. Pi may retry, compact // and retry, or process queued follow-up input after it fires, so cache its // state without publishing Completion evidence here. latestRunMessages = (event as any).messages as any[] | undefined; recorder.agentEndWaiting(); }); pi.on("agent_settled", () => { if (completionPublished) return; // An aborted final run establishes the Cancellation cutover. It is // control flow rather than Completion evidence, and the coordinator's // registered parent-side actions cascade cancellation to descendants. if (classifySettledRun(latestRunMessages) === "aborted") { const operation = operationFromEnvironment(); if (operation && hasCancellationRequest(operation)) acknowledgeCancellation(); return; } try { // Provider errors are published only after the final settled run, so a // transient error followed by a retry cannot win over the retry result. const publication = publishCompletionOnce( buildCompletionSidecar(latestRunMessages, childThinkingLevel(pi)), ); if (publication === "published" || publication === "already-requested") { recorder.agentSettledDone(); } } catch { // Best effort — the parent watcher will report missing Completion // evidence if the sidecar cannot be written. } }); pi.on("turn_start", (event) => { recorder.turnStart((event as any).turnIndex); }); pi.on("turn_end", (event) => { recorder.turnEnd((event as any).turnIndex); }); pi.on("before_provider_request", () => { recorder.beforeProviderRequest(); }); pi.on("after_provider_response", () => { recorder.afterProviderResponse(); }); pi.on("message_update", (event) => { recorder.messageUpdate((event as any).assistantMessageEvent?.type); }); pi.on("tool_execution_start", (event) => { recorder.toolExecutionStart((event as any).toolCallId, (event as any).toolName); }); pi.on("tool_call", (event) => { recorder.toolCall((event as any).toolCallId, (event as any).toolName); }); pi.on("tool_execution_update", (event) => { recorder.toolExecutionUpdate((event as any).toolCallId, (event as any).toolName); }); pi.on("tool_result", (event) => { recorder.toolResult((event as any).toolCallId, (event as any).toolName); }); pi.on("tool_execution_end", (event) => { recorder.toolExecutionEnd((event as any).toolCallId, (event as any).toolName); }); pi.on("session_shutdown", (event) => { const operation = operationFromEnvironment(); if ( (event as any).reason !== "reload" && operation && hasCancellationRequest(operation) ) { acknowledgeCancellation(); } stopCancellationWatch(); // Pi may emit session_shutdown immediately after agent_settled. Keep a // pending transient publication alive through a normal quit; otherwise a // valid Completion result can be lost just as the child closes. if ((event as any).reason === "quit") (completionRetry as any)?.ref?.(); recorder.sessionShutdown((event as any).reason); }); }