/** * agent-runner.ts — Enterprise Agent Execution Engine * * Core execution engine that creates sessions, runs agents, collects results. * Enhanced with: * - Resource quotas (token budgets, time limits, tool limits) * - Circuit breaker for model calls * - Structured error classification * - Graceful degradation strategies * - Comprehensive telemetry and metrics */ import type { Api, Model } from "@earendil-works/pi-ai"; import type { ExtensionContext } from "@earendil-works/pi-coding-agent"; import { type AgentSession, type AgentSessionEvent, createAgentSession, DefaultResourceLoader, type ExtensionAPI, getAgentDir, SessionManager, SettingsManager, } from "@earendil-works/pi-coding-agent"; import { isAbortError } from "./abort-wait.js"; import { getPromptCompressionLevel, isFreeModelsOnly } from "./agent-registry.js"; import { runAdversarialValidation, } from "./agent-runner-validator.js"; import { type EffectiveConfig, getAgentConfig, getConfig, getMemoryToolNames, getReadOnlyMemoryToolNames, getToolNamesForType } from "./agent-types.js"; import { buildCompactionSnapshot, type CompactionSnapshot } from "./compaction-snapshot.js"; import { buildParentContext, extractText } from "./context.js"; import { resolveCtxInjectionForAgent } from "./context-mode-bridge.js"; import { DEFAULT_AGENTS } from "./default-agents.js"; import { buildEnvFromContext } from "./env-context.js"; import { type AgentHandoff, parseHandoff, renderHandoffForParent } from "./handoff.js"; import { type HookRegistry, normalizeHookResponse } from "./hooks.js"; import { logger } from "./logger.js"; import { buildMemoryBlock, buildReadOnlyMemoryBlock } from "./memory.js"; import { buildAgentPrompt, type PromptExtras } from "./prompts.js"; import { loadSettings } from "./settings.js"; import { preloadSkills } from "./skill-loader.js"; import { type AgentAbortReason, type AgentOutcome, type AgentRunnerErrorCode, deriveAgentOutcome, } from "./spend.js"; import { emitTelemetry } from "./telemetry.js"; import type { SubagentType, ThinkingLevel, ValidationResult } from "./types.js"; import { hasValidators, } from "./validators.js"; // ============================================================================ // Constants & Error Types // ============================================================================ /** Names of tools registered by this extension that subagents must NOT inherit. */ const EXCLUDED_TOOL_NAMES: ReadonlySet = new Set([ "Agent", "get_subagent_result", "steer_subagent", ]); /** * Default max turns. undefined = unlimited. * Data-driven default (2026-08-10, 268-session eval): workers without an * explicit budget ran to the 10min duration quota (DEFAULT_MAX_DURATION_MS) * and were aborted with no final report — 40/43 subagent aborts in the * eval were exactly that. A finite default turns the duration kill into a * soft-limit steer (with graceTurns of headroom for the report) so workers * always return an end report instead of dying silently. * 30 turns ≈ 4-8 min at typical model turn cost; oversize tasks should * pass an explicit max_turns instead of relying on the ceiling. */ let defaultMaxTurns: number | undefined = 30; /** Additional turns allowed after the soft limit steer message. */ let graceTurns = 5; /** * Max revision turns after a blocking `subagent:end` hook. * 0 = fail closed immediately (no revision). Fresh-install default keeps * end hooks observational unless a quality gate opts into revisions. */ let maxEndHookRevisions = 0; /** Resource quota defaults. */ const DEFAULT_MAX_TOKENS = 500_000; const DEFAULT_MAX_DURATION_MS = 600_000; // 10 minutes const DEFAULT_MAX_TOOL_CALLS = 100; /** Circuit breaker defaults. */ const CB_FAILURE_THRESHOLD = 5; const CB_RECOVERY_TIMEOUT_MS = 30_000; export class AgentRunnerError extends Error { constructor( message: string, public readonly code: AgentRunnerErrorCode, public readonly context?: Record, ) { super(message); this.name = "AgentRunnerError"; } } // ============================================================================ // Configuration // ============================================================================ export function normalizeMaxTurns(n: number | undefined): number | undefined { if (n == null || typeof n !== "number" || Number.isNaN(n) || !Number.isFinite(n) || n === 0) return undefined; return Math.max(1, Math.floor(n)); } export function getDefaultMaxTurns(): number | undefined { return defaultMaxTurns; } export function setDefaultMaxTurns(n: number | undefined): void { defaultMaxTurns = normalizeMaxTurns(n); } export function getGraceTurns(): number { return graceTurns; } export function setGraceTurns(n: number): void { graceTurns = Math.max(1, n); } export function getMaxEndHookRevisions(): number { return maxEndHookRevisions; } /** Clamp end-hook revision budget to a safe finite 0..10 range (NaN/Inf → 0). */ export function clampMaxEndHookRevisions(n: number): number { return Number.isFinite(n) ? Math.max(0, Math.min(10, Math.trunc(n))) : 0; } export function setMaxEndHookRevisions(n: number): void { maxEndHookRevisions = clampMaxEndHookRevisions(n); } // ============================================================================ // Circuit Breaker for Model Calls // ============================================================================ /** * True when a thrown error or soft assistant errorMessage looks like a model / * provider transport failure (auth, rate limit, unavailable). Used so the * circuit breaker trips on 401s during session.prompt — not on bad cwd / config * failures from createAgentSession (CHE-17). */ export function isModelTransportFailure(err: unknown): boolean { if (isAbortError(err)) return false; const msg = err instanceof Error ? err.message : String(err); return isModelFailureMessage(msg); } export function isModelFailureMessage(message: string): boolean { return ( isModelAuthFailure(message) || /\b403\b/.test(message) || /\b429\b/.test(message) || /rate[\s_-]?limit/i.test(message) || /model\s+(?:unavailable|not\s+found)/i.test(message) || /overloaded/i.test(message) || /insufficient[\s_-]?quota/i.test(message) ); } /** * Auth-shaped model failures that justify one automatic retry on the * session-default (parent) model — e.g. Explore's pinned haiku returning * PAID_MODEL_AUTH_REQUIRED / 401 while the parent model still works (CHE-15). */ export function isModelAuthFailure(message: string): boolean { return ( /PAID_MODEL_AUTH_REQUIRED/i.test(message) || /\b401\b/.test(message) || /auth(?:entication|orization)?\s+(?:failed|required|error|denied)/i.test(message) || /User not found/i.test(message) ); } /** True when a pi-ai Model has zero cost on every dimension (free tier). */ export function isFreeModel(model: Model): boolean { const cost = (model as unknown as { cost?: { input?: number; output?: number; cacheRead?: number; cacheWrite?: number } }).cost; if (!cost) return false; return (cost.input ?? 0) === 0 && (cost.output ?? 0) === 0 && (cost.cacheRead ?? 0) === 0 && (cost.cacheWrite ?? 0) === 0; } /** Build ordered fallback candidates excluding the failed model. */ function getFallbackModels( failed: Model, registry: { getAvailable?(): Model[]; getAll(): Model[] }, parent: Model | undefined, freeOnly: boolean, ): Model[] { const available = registry.getAvailable?.() ?? registry.getAll(); const filtered = freeOnly ? available.filter(isFreeModel) : available; const failedKey = modelKey(failed); const byKey = new Map(filtered.map((m) => [modelKey(m), m] as const)); byKey.delete(failedKey); const out: Model[] = []; if (parent) { const pk = modelKey(parent); const parentFreeOk = !freeOnly || isFreeModel(parent); if (pk !== failedKey && parentFreeOk) { // Parent is always a candidate even if not in the available list (e.g. host model) if (byKey.has(pk)) byKey.delete(pk); out.push(parent); } } for (const m of filtered) { const k = modelKey(m); if (byKey.has(k)) out.push(byKey.get(k)!); } return out.slice(0, 5); } function modelKey(model: Model): string { return `${model.provider}/${model.id}`; } /** * Backoff before retrying on a rate-limit failure. Honors a `retry-after: * Ns` hint in the error message when present (capped at 10s); otherwise * waits a short fixed delay so we do not instantly burn the next candidate. */ async function rateLimitBackoff(errorMessage: string): Promise { const isRateLimit = /\b429\b|rate[\s_-]?limit/i.test(errorMessage); if (!isRateLimit) return; let ms = 1500; const hint = errorMessage.match(/retry[-\s]?after[:\s]*(\d+)/i); if (hint) { ms = Math.min(10_000, Number.parseInt(hint[1], 10) * 1000); } logger.debug("Rate-limit backoff before fallback retry", { ms }); await new Promise((resolve) => setTimeout(resolve, ms)); } class ModelCircuitBreaker { private failures = 0; private lastFailureAt = 0; private state: "closed" | "open" | "half-open" = "closed"; /** Throw if the breaker is open and the recovery window has not elapsed. */ assertAllow(): void { if (this.state === "open") { if (Date.now() - this.lastFailureAt > CB_RECOVERY_TIMEOUT_MS) { this.state = "half-open"; } else { throw new AgentRunnerError( "Model circuit breaker is OPEN — too many consecutive failures", "model_unavailable", { failures: this.failures, lastFailure: this.lastFailureAt }, ); } } } recordSuccess(): void { if (this.state === "half-open") { this.state = "closed"; } this.failures = 0; } recordFailure(): void { this.failures++; this.lastFailureAt = Date.now(); if (this.failures >= CB_FAILURE_THRESHOLD) { this.state = "open"; } } /** * Run `fn` under the breaker. By default every rejection counts; pass * `countFailure` to ignore non-model errors (e.g. AbortError). */ call( fn: () => Promise, options?: { countFailure?: (err: unknown) => boolean }, ): Promise { this.assertAllow(); const shouldCount = options?.countFailure ?? (() => true); return fn().then( (result) => { this.recordSuccess(); return result; }, (err) => { if (shouldCount(err)) this.recordFailure(); throw err; }, ); } getState(): { state: string; failures: number; lastFailureAt: number } { return { state: this.state, failures: this.failures, lastFailureAt: this.lastFailureAt }; } /** Test helper — clear breaker state between cases. */ reset(): void { this.failures = 0; this.lastFailureAt = 0; this.state = "closed"; } } const globalCircuitBreaker = new ModelCircuitBreaker(); /** * Run a model prompt under the circuit breaker. * * 401 / PAID_MODEL_AUTH_REQUIRED usually surface as a soft assistant * `stopReason: "error"` after prompt resolves (not a thrown exception), so we * inspect getLastAssistantError afterwards. Thrown model failures count too; * AbortError does not (CHE-17). */ async function promptWithCircuitBreaker( session: AgentSession, text: string, options?: { allowWhenOpen?: boolean }, ): Promise { if (!options?.allowWhenOpen) globalCircuitBreaker.assertAllow(); try { await session.prompt(text); } catch (err) { if (isModelTransportFailure(err)) globalCircuitBreaker.recordFailure(); throw err; } const softError = getLastAssistantError(session); if (softError && isModelFailureMessage(softError)) { globalCircuitBreaker.recordFailure(); } else { globalCircuitBreaker.recordSuccess(); } } // ============================================================================ // Model Resolution // ============================================================================ let _cachedRegistry: unknown = null; let _cachedKeys: Set | null = null; function getAvailableKeys( registry: { getAvailable?(): Model[] }, ): Set | undefined { if (registry === _cachedRegistry && _cachedKeys) return _cachedKeys; const available = registry.getAvailable?.(); if (!available) return undefined; _cachedKeys = new Set(available.map((m) => `${m.provider}/${m.id}`)); _cachedRegistry = registry; return _cachedKeys; } /** * Resolve the effective configured model string for a spawned agent. * * A `subagentModel` setting overrides the agent's own configured model: * - `"inherit"` → `undefined`, so `resolveDefaultModel` falls back to the * session-default (parent) model. This is the escape hatch when a built-in * read-only agent's pinned model is unreachable (e.g. provider 401). * - any other non-empty string → used verbatim as a `provider/modelId`. * - `undefined` → the agent's own configured model is used (current behavior). */ export function resolveConfiguredModel( subagentModelSetting: string | undefined, agentModel: string | undefined, ): string | undefined { if (subagentModelSetting === "inherit") return undefined; return subagentModelSetting ?? agentModel; } function resolveDefaultModel( parentModel: Model | undefined, registry: { find(provider: string, modelId: string): Model | undefined; getAvailable?(): Model[]; }, configModel?: string, ): Model | undefined { if (configModel) { const slashIdx = configModel.indexOf("/"); if (slashIdx !== -1) { const provider = configModel.slice(0, slashIdx); const modelId = configModel.slice(slashIdx + 1); const availableKeys = getAvailableKeys(registry); const isAvailable = (p: string, id: string) => !availableKeys || availableKeys.has(`${p}/${id}`); const found = registry.find(provider, modelId); if (found && isAvailable(provider, modelId)) return found; } } return parentModel; } // ============================================================================ // Types // ============================================================================ export interface ToolActivity { type: "start" | "end"; toolName: string; } export interface ResourceQuotas { /** Max total tokens (input + output) before hard stop. */ maxTokens?: number; /** Max execution duration in ms. */ maxDurationMs?: number; /** Max number of tool calls. */ maxToolCalls?: number; } export interface RunOptions { pi: ExtensionAPI; agentId?: string; model?: Model; maxTurns?: number; signal?: AbortSignal; isolated?: boolean; inheritContext?: boolean; thinkingLevel?: ThinkingLevel; cwd?: string; onToolActivity?: (activity: ToolActivity) => void; onTextDelta?: (delta: string, fullText: string) => void; onSessionCreated?: (session: AgentSession) => void; onTurnEnd?: (turnCount: number) => void; onAssistantUsage?: (usage: { input: number; output: number; cacheWrite: number }) => void; onCompaction?: (info: CompactionSnapshot) => void; skipValidators?: boolean; onValidationComplete?: (results: ValidationResult[]) => void; currentLevel?: number; levelLimit?: number; parentConfig?: EffectiveConfig; partitions?: readonly string[]; /** Short correlation id (8 hex chars). Set by AgentManager at spawn. */ correlationId?: string; hooks?: HookRegistry; spawnedAt?: number; onContextBuilt?: (timestamp: number) => void; /** Resource quotas for this run. */ quotas?: ResourceQuotas; /** * Override the module-level `maxEndHookRevisions` setting for this run. * `0` = fail closed on `subagent:end` block with no revision turn. */ maxEndHookRevisions?: number; } export interface RunResult { responseText: string; session: AgentSession; aborted: boolean; /** True when the run was aborted because it exceeded its duration quota. */ timedOut?: boolean; steered: boolean; validationResults?: ValidationResult[]; validated?: boolean; handoff?: AgentHandoff; /** Execution metrics. */ metrics: RunMetrics; /** * Set when the run ended because the model/provider errored on its final * turn (no thrown exception). Callers that key off status should treat a * non-empty `error` as a failure even though `session.prompt` resolved. */ error?: string; /** * Structured reason when an internal budget gate (token/tool/duration/turn * quota) stopped the run. Absent for normal completion and external stops, * so a `blocked_budget` outcome is derivable rather than guessed from * error strings (R4 / AE2). */ abortReason?: AgentAbortReason; /** Explicit outcome contract (R4): executed | blocked_budget | not_executed. */ outcome?: AgentOutcome; /** Reason for the outcome — the structured abort message or a partial-progress note. */ outcomeReason?: string; } export interface RunMetrics { durationMs: number; turns: number; toolCalls: number; tokensIn: number; tokensOut: number; tokensCacheWrite: number; contextBuiltAt?: number; latencyToFirstTokenMs?: number; } // ============================================================================ // Response Collection // ============================================================================ function collectResponseText(session: AgentSession) { let text = ""; const unsubscribe = session.subscribe((event: AgentSessionEvent) => { if (event.type === "message_start") { text = ""; } if (event.type === "message_update" && event.assistantMessageEvent.type === "text_delta") { text += event.assistantMessageEvent.delta; } }); return { getText: () => text, unsubscribe }; } function getLastAssistantText(session: AgentSession): string { for (let i = session.messages.length - 1; i >= 0; i--) { const msg = session.messages[i]; if (msg.role !== "assistant") continue; const text = extractText(msg.content).trim(); if (text) return text; } return ""; } /** * Return the most recent assistant turn's model/provider error, if any. * * The host surfaces a model or provider failure (401 auth, unavailable model, * rate limit, etc.) as an AssistantMessage with `stopReason: "error"` and an * `errorMessage`. The run loop does not throw for these — the session simply * ends — so without this check the error is silently dropped and the agent * appears to complete with empty output. */ function getLastAssistantError(session: AgentSession): string | undefined { for (let i = session.messages.length - 1; i >= 0; i--) { const msg = session.messages[i]; if (msg.role !== "assistant") continue; // Only the most recent assistant turn reflects the run's final outcome. // An earlier errored turn is stale once a later assistant turn exists // (e.g. a retry/revision that produced an empty-but-non-error turn), so we // do not scan past the latest assistant message. return msg.stopReason === "error" && msg.errorMessage ? msg.errorMessage : undefined; } return undefined; } function forwardAbortSignal( session: AgentSession, signal?: AbortSignal, onExternalAbort?: () => void, ): () => void { if (!signal) return () => {}; const onAbort = () => { // Mark the run aborted first so prompt()'s AbortError is treated as a // graceful stop (CHE-19). Always abort the session even if the callback // throws — cancellation must propagate. try { onExternalAbort?.(); } finally { session.abort(); } }; if (signal.aborted) { onAbort(); return () => {}; } signal.addEventListener("abort", onAbort, { once: true }); return () => signal.removeEventListener("abort", onAbort); } // ============================================================================ // Deferred Context // ============================================================================ export function buildEffectivePrompt( ctx: ExtensionContext, prompt: string, options: RunOptions, ): string { if (!options.inheritContext) return prompt; const parentContext = buildParentContext(ctx); const builtAt = Date.now(); options.onContextBuilt?.(builtAt); const spawnedAgo = options.spawnedAt ? builtAt - options.spawnedAt : 0; logger.debug("Context built after spawn", { agentId: options.agentId ?? "unknown", spawnedAgo }); if (!parentContext) return prompt; return parentContext + prompt; } // ============================================================================ // Main Agent Runner // ============================================================================ export async function runAgent( ctx: ExtensionContext, type: SubagentType, prompt: string, options: RunOptions, ): Promise { const startTime = performance.now(); const quotas = { maxTokens: options.quotas?.maxTokens ?? DEFAULT_MAX_TOKENS, maxDurationMs: options.quotas?.maxDurationMs ?? DEFAULT_MAX_DURATION_MS, maxToolCalls: options.quotas?.maxToolCalls ?? DEFAULT_MAX_TOOL_CALLS, }; // Check duration quota early. Returns true once the quota is exceeded so the // subscriber can abort the session gracefully. We intentionally do NOT throw // here: this runs inside a sync session event subscriber invoked by the host // emitter, and a throw escapes the run loop and surfaces as an uncaught // exception that kills the whole host process. Quota exhaustion is a normal // runtime outcome and must be handled like the token/tool-call quotas below // (abort + flag + return). let durationQuotaTripped = false; const checkDurationQuota = (): boolean => { if (durationQuotaTripped) return true; const elapsedMs = performance.now() - startTime; if (elapsedMs > quotas.maxDurationMs) { durationQuotaTripped = true; return true; } return false; }; const config = getConfig(type, options.parentConfig, options.partitions); const agentConfig = getAgentConfig(type); // Early exit: check level limit const currentLevel = options.currentLevel ?? 0; const depthLimit = options.levelLimit ?? 5; if (currentLevel >= depthLimit) { throw new AgentRunnerError( `Max agent depth reached (${currentLevel}/${depthLimit})`, "depth_exceeded", { currentLevel, depthLimit }, ); } // Telemetry emitTelemetry("agent:spawned", { type, parentType: options.parentConfig ? type : undefined, depth: currentLevel, budget: options.maxTurns, }); const effectiveCwd = options.cwd ?? ctx.cwd; const env = buildEnvFromContext(options.pi) ?? { isGitRepo: false, branch: "", platform: process.platform, }; const parentSystemPrompt = ctx.getSystemPrompt(); // Resolve extensions/skills const extensions = options.isolated ? false : config.extensions; const skills = options.isolated ? false : config.skills; const extras: PromptExtras = {}; if (Array.isArray(skills)) { const loaded = preloadSkills(skills, effectiveCwd); if (loaded.length > 0) extras.skillBlocks = loaded; } let toolNames = getToolNamesForType(type); // Persistent memory if (agentConfig?.memory) { const existingNames = new Set(toolNames); const denied = agentConfig.disallowedTools ? new Set(agentConfig.disallowedTools) : undefined; const effectivelyHas = (name: string) => existingNames.has(name) && !denied?.has(name); const hasWriteTools = effectivelyHas("write") || effectivelyHas("edit"); if (hasWriteTools) { const extraNames = getMemoryToolNames(existingNames); if (extraNames.length > 0) toolNames = [...toolNames, ...extraNames]; extras.memoryBlock = buildMemoryBlock(agentConfig.name, agentConfig.memory, effectiveCwd, agentConfig.maxMemoryLines); } else { const extraNames = getReadOnlyMemoryToolNames(existingNames); if (extraNames.length > 0) toolNames = [...toolNames, ...extraNames]; extras.memoryBlock = buildReadOnlyMemoryBlock(agentConfig.name, agentConfig.memory, effectiveCwd, agentConfig.maxMemoryLines); } } // Parent permission inheritance const allowedTools = new Set(config.builtinToolNames); toolNames = toolNames.filter((t) => allowedTools.has(t)); // Build system prompt const compressionLevel = agentConfig?.promptCompressionLevel ?? getPromptCompressionLevel(); let systemPrompt: string; if (agentConfig) { systemPrompt = buildAgentPrompt(agentConfig, effectiveCwd, env, parentSystemPrompt, extras, compressionLevel); } else { const fallback = DEFAULT_AGENTS.get("general-purpose"); if (!fallback) { throw new AgentRunnerError( `No fallback config available for unknown type "${type}"`, "unknown", ); } systemPrompt = buildAgentPrompt({ ...fallback, name: type }, effectiveCwd, env, parentSystemPrompt, extras, compressionLevel); } // Context-mode injection — only for agents that opt in via useContextMode, // and only when getConfig kept ctx_* tools in the effective allowlist. const ctxInjection = resolveCtxInjectionForAgent(agentConfig?.useContextMode); if (ctxInjection) { const allowed = new Set(config.builtinToolNames); // Exclude tools the agent denies via disallowedTools, matching the active-tool // filter below (Lines 614-634). Otherwise an opted-in agent that denies a ctx_* // tool would still get the context-mode prompt/log for a tool the session removes. const denied = agentConfig?.disallowedTools ? new Set(agentConfig.disallowedTools) : undefined; const injectable = ctxInjection.toolAllowList.filter( (t) => allowed.has(t) && !toolNames.includes(t) && !denied?.has(t), ); if (injectable.length > 0) { systemPrompt = `${systemPrompt}\n\n${ctxInjection.systemPromptAddition}`; toolNames = [...toolNames, ...injectable]; logger.debug("context-mode tools injected", { agentId: options.agentId ?? "unknown" }); } } const noSkills = skills === false || Array.isArray(skills); const agentDir = getAgentDir(); const loader = new DefaultResourceLoader({ cwd: effectiveCwd, agentDir, noExtensions: extensions === false, noSkills, noPromptTemplates: true, noThemes: true, noContextFiles: true, systemPromptOverride: () => systemPrompt, appendSystemPromptOverride: () => [], }); await loader.reload(); // Resolve model with circuit breaker. A `subagentModel` setting can override // the agent's configured model ("inherit" falls back to the session-default // model; "provider/modelId" pins a specific one). Lets read-only agents move // off a broken configured model without a code change. const configuredModel = resolveConfiguredModel( loadSettings(effectiveCwd).subagentModel, agentConfig?.model, ); let model = options.model ?? resolveDefaultModel( ctx.model, ctx.modelRegistry, configuredModel, ); if (!model) { throw new AgentRunnerError( "No model available for agent execution", "model_unavailable", ); } // Session-only gate: when free-only is on, swap a paid model for a free one // before any prompt is sent. Prefers the parent model when it is free, then // the first available free model. This keeps the session within the free tier // without needing to rewrite persisted subagentModel settings. if (isFreeModelsOnly() && !isFreeModel(model!)) { const freeOnlyParent = ctx.model && isFreeModel(ctx.model) ? ctx.model : undefined; if (freeOnlyParent && modelKey(freeOnlyParent) !== modelKey(model!)) { logger.warn("Free-only mode: swapping paid model for free parent", { agentId: options.agentId ?? "unknown", from: modelKey(model!), to: modelKey(freeOnlyParent), }); model = freeOnlyParent; } else { const avail = ctx.modelRegistry.getAvailable?.() ?? ctx.modelRegistry.getAll(); const freeAlt = avail.find((m) => isFreeModel(m) && modelKey(m) !== modelKey(model!)); if (freeAlt) { logger.warn("Free-only mode: swapping paid model for free alternative", { agentId: options.agentId ?? "unknown", from: modelKey(model!), to: modelKey(freeAlt), }); model = freeAlt; } else { logger.warn("Free-only mode: no free alternative found, using paid model", { agentId: options.agentId ?? "unknown", model: modelKey(model!), }); } } } const thinkingLevel = options.thinkingLevel ?? agentConfig?.thinking; const sessionOpts: Parameters[0] = { cwd: effectiveCwd, agentDir, sessionManager: SessionManager.inMemory(effectiveCwd), settingsManager: SettingsManager.create(effectiveCwd, agentDir), // Pi 0.80.8 replaced CreateAgentSessionOptions.modelRegistry with the async // modelRuntime. Extensions only receive the synchronous ModelRegistry // facade (ctx.modelRegistry), not the underlying ModelRuntime, so we let // the SDK build its default runtime from the same agentDir (auth.json + // models.json). The subagent shares the host agentDir, so credentials and // the model catalog match, and `model` is passed explicitly above. model, tools: toolNames, // Keep orchestration tools out of the initial active set (Pi 0.81+). excludeTools: [...EXCLUDED_TOOL_NAMES], resourceLoader: loader, }; if (thinkingLevel) { sessionOpts.thinkingLevel = thinkingLevel; } // Hook: subagent:start if (options.hooks) { const hookResult = await options.hooks.dispatch("subagent:start", options.agentId ?? "unknown", { type, model: `${model.provider}/${model.id}`, quotas, }); const startDecision = normalizeHookResponse(hookResult); if (startDecision.action === "block") { throw new AgentRunnerError(startDecision.reason ?? "Blocked by hook", "aborted", { hook: "subagent:start", ...(startDecision.feedback !== undefined ? { feedback: startDecision.feedback } : {}), }); } } const effectivePrompt = buildEffectivePrompt(ctx, prompt, options); // Session creation is NOT under the model circuit breaker: bad cwd / invalid // config must not trip a global spawn freeze. Model failures are counted on // prompt() (incl. soft 401 assistant errors) — see promptWithCircuitBreaker (CHE-17). const { session } = await createAgentSession(sessionOpts); const baseSessionName = agentConfig?.name ?? type; session.setSessionName( options.agentId ? `${baseSessionName}#${options.agentId.slice(0, 8)}` : baseSessionName, ); // Tool filtering const disallowedSet = agentConfig?.disallowedTools ? new Set(agentConfig.disallowedTools) : undefined; if (extensions !== false) { const builtinToolNameSet = new Set(toolNames); const activeTools = session.getActiveToolNames().filter((t) => { if (EXCLUDED_TOOL_NAMES.has(t)) return false; if (disallowedSet?.has(t)) return false; if (builtinToolNameSet.has(t)) return true; if (Array.isArray(extensions)) { return extensions.some((ext) => t.startsWith(ext) || t.includes(ext)); } return true; }); session.setActiveToolsByName(activeTools); } else if (disallowedSet) { const activeTools = session.getActiveToolNames().filter((t) => !disallowedSet.has(t)); session.setActiveToolsByName(activeTools); } await session.bindExtensions({ onError: (err) => { options.onToolActivity?.({ type: "end", toolName: `extension-error:${err.extensionPath}`, }); }, }); options.onSessionCreated?.(session); // Turn tracking and quotas let turnCount = 0; let toolCallCount = 0; let tokensIn = 0; let tokensOut = 0; let tokensCacheWrite = 0; let latencyToFirstToken: number | undefined; const maxTurns = normalizeMaxTurns(options.maxTurns ?? agentConfig?.maxTurns ?? defaultMaxTurns); let softLimitReached = false; let aborted = false; let timedOut = false; // Structured reason for internal budget aborts (R4). Set at the abort site // below so the outcome is derived from a deliberate signal, not from // matching error strings. External stops leave it undefined. let abortReason: AgentAbortReason | undefined; let currentMessageText = ""; const unsubTurns = session.subscribe((event: AgentSessionEvent) => { // Quota checks — guard re-entry: once timed out *or* externally aborted, // stop re-aborting on every subsequent event (session.abort() can emit // synchronously). Also prevents an external cancel from being mis-labeled // as timedOut when duration has already elapsed (CHE-19 / CodeRabbit). if (timedOut || aborted) return; if (checkDurationQuota()) { const elapsedMs = performance.now() - startTime; logger.warn(`Duration quota exceeded`, { agentId: options.agentId, elapsedMs, maxDurationMs: quotas.maxDurationMs, }); timedOut = true; abortReason = { kind: "duration_quota", message: `Duration quota exceeded (${quotas.maxDurationMs}ms)` }; session.abort(); aborted = true; return; } const totalTokens = tokensIn + tokensOut; if (totalTokens > quotas.maxTokens) { logger.warn(`Token quota exceeded`, { agentId: options.agentId, totalTokens, maxTokens: quotas.maxTokens }); abortReason = { kind: "token_quota", message: `Token quota exceeded (${totalTokens}/${quotas.maxTokens} tokens)` }; session.abort(); aborted = true; return; } if (event.type === "turn_end") { options.hooks ?.dispatch("turn:end", options.agentId ?? "unknown") .catch((err) => { logger.debug(`Hook dispatch error: ${err instanceof Error ? err.message : String(err)}`); }); turnCount++; options.onTurnEnd?.(turnCount); if (maxTurns != null) { if (!softLimitReached && turnCount >= maxTurns) { softLimitReached = true; session.steer("You have reached your turn budget. Wrap up NOW: write your final end report — findings, what is done, what is blocked, exact file paths/commit SHAs. Do not start new work."); } else if (softLimitReached && turnCount >= maxTurns + graceTurns) { aborted = true; abortReason = { kind: "turn_budget", message: `Turn budget exhausted (${turnCount}/${maxTurns} turns + ${graceTurns} grace turns)`, }; session.abort(); } } } if (event.type === "turn_start") { options.hooks ?.dispatch("turn:start", options.agentId ?? "unknown") .catch((err) => { logger.debug(`Hook dispatch error: ${err instanceof Error ? err.message : String(err)}`); }); } if (event.type === "message_start") { currentMessageText = ""; if (latencyToFirstToken === undefined) { latencyToFirstToken = performance.now() - startTime; } } if (event.type === "message_update" && event.assistantMessageEvent.type === "text_delta") { currentMessageText += event.assistantMessageEvent.delta; options.onTextDelta?.(event.assistantMessageEvent.delta, currentMessageText); } if (event.type === "tool_execution_start") { toolCallCount++; if (toolCallCount > quotas.maxToolCalls) { logger.warn(`Tool call quota exceeded`, { agentId: options.agentId, toolCallCount, maxToolCalls: quotas.maxToolCalls }); abortReason = { kind: "tool_quota", message: `Tool call quota exceeded (${toolCallCount}/${quotas.maxToolCalls} tool calls)` }; session.abort(); aborted = true; return; } options.onToolActivity?.({ type: "start", toolName: event.toolName }); } if (event.type === "tool_execution_end") { options.onToolActivity?.({ type: "end", toolName: event.toolName }); } if (event.type === "message_end" && event.message.role === "assistant") { const msg = event.message as { usage?: { input?: number; output?: number; cacheWrite?: number } }; const u = msg.usage; if (u) { tokensIn += u.input ?? 0; tokensOut += u.output ?? 0; tokensCacheWrite += u.cacheWrite ?? 0; options.onAssistantUsage?.({ input: u.input ?? 0, output: u.output ?? 0, cacheWrite: u.cacheWrite ?? 0, }); } } if (event.type === "compaction_end") { const snapshot = buildCompactionSnapshot(event); options.onCompaction?.(snapshot); if (!event.aborted) { options.hooks ?.dispatch("compaction:end", options.agentId ?? "unknown", { reason: event.reason, tokensBefore: snapshot.tokensBefore, tokensAfter: snapshot.tokensAfter, reductionPercent: snapshot.reductionPercent, aborted: false, willRetry: snapshot.willRetry, }) .catch((err) => { logger.debug(`Hook dispatch error: ${err instanceof Error ? err.message : String(err)}`); }); } else { logger.debug( `Compaction aborted for ${options.agentId ?? "unknown"}: ${event.errorMessage ?? event.reason}` + (event.willRetry ? " (will retry)" : ""), ); } } if ( event.type === "summarization_retry_scheduled" || event.type === "summarization_retry_attempt_start" || event.type === "summarization_retry_finished" ) { const attempt = event.type === "summarization_retry_scheduled" ? event.attempt : undefined; const maxAttempts = event.type === "summarization_retry_scheduled" ? event.maxAttempts : undefined; const detail = event.type === "summarization_retry_scheduled" ? `: ${event.errorMessage}` : event.type === "summarization_retry_attempt_start" ? ` source=${event.source}` : ""; logger.debug( `Summarization retry (${event.type}) for ${options.agentId ?? "unknown"}` + (attempt !== undefined ? ` attempt=${attempt}/${maxAttempts}` : "") + detail, ); emitTelemetry("subagent:summarization_retry", { agentId: options.agentId ?? "unknown", phase: event.type, attempt, maxAttempts, }); } if (event.type === "compaction_start") { options.hooks ?.dispatch("compaction:start", options.agentId ?? "unknown", { reason: event.reason, }) .catch((err) => { logger.debug(`Hook dispatch error: ${err instanceof Error ? err.message : String(err)}`); }); } }); const collector = collectResponseText(session); const cleanupAbort = forwardAbortSignal(session, options.signal, () => { aborted = true; }); let gatedResponseText = ""; let exhaustedModelError: string | undefined; try { let currentModelForRetry: Model = model!; let soft: string | undefined; try { await promptWithCircuitBreaker(session, effectivePrompt); gatedResponseText = collector.getText().trim() || getLastAssistantText(session); soft = !aborted && !options.signal?.aborted ? getLastAssistantError(session) : undefined; } catch (err) { if ((timedOut || aborted || options.signal?.aborted) && isAbortError(err)) { aborted = true; gatedResponseText = collector.getText().trim() || getLastAssistantText(session); soft = undefined; } else if (isModelTransportFailure(err) && (!options.model || isFreeModelsOnly()) && !aborted && !timedOut && !options.signal?.aborted) { soft = err instanceof Error ? err.message : String(err); gatedResponseText = ""; } else { throw err; } } // Generalized fallback (CHE-15 extended): 401/auth, 429/rate-limit, // 403, unavailable, overloaded, quota — retry transparently on a // different model instead of failing the spawn. Tries parent first, // then up to 4 more available models; when free-only is on, only // zero-cost models are considered. Skips when caller pinned options.model // or output was already produced. const shouldFallback = (!options.model || isFreeModelsOnly()) && !gatedResponseText && !!soft && isModelFailureMessage(soft); if (shouldFallback) { const candidates = getFallbackModels(currentModelForRetry, ctx.modelRegistry, ctx.model, isFreeModelsOnly()); let lastErr: string | undefined = soft; for (const candidate of candidates) { if (modelKey(candidate) === modelKey(currentModelForRetry)) continue; logger.warn("Model failed; retrying with fallback model", { agentId: options.agentId ?? "unknown", from: modelKey(currentModelForRetry), to: modelKey(candidate), error: lastErr, freeOnly: isFreeModelsOnly(), }); try { if (lastErr) await rateLimitBackoff(lastErr); await session.setModel(candidate); currentModelForRetry = candidate; await promptWithCircuitBreaker(session, effectivePrompt, { allowWhenOpen: true }); gatedResponseText = collector.getText().trim() || getLastAssistantText(session); const nextErr = !aborted && !options.signal?.aborted ? getLastAssistantError(session) : undefined; if (gatedResponseText || !nextErr || !isModelFailureMessage(nextErr)) { break; } lastErr = nextErr; } catch (err) { if ((timedOut || aborted || options.signal?.aborted) && isAbortError(err)) { aborted = true; break; } if (isModelTransportFailure(err)) { lastErr = err instanceof Error ? err.message : String(err); continue; } throw err; } } if (!gatedResponseText && !aborted && lastErr) exhaustedModelError = lastErr; } else if (!gatedResponseText && soft) { exhaustedModelError = soft; } if (options.hooks) { const revisionBudget = clampMaxEndHookRevisions( options.maxEndHookRevisions ?? maxEndHookRevisions, ); let attempt = 1; const maxAttempts = revisionBudget + 1; while (true) { // Resolve end status from the *current* session state before publishing. // A model/provider failure (stopReason "error", empty content) must not // advertise status "completed" — that lied to hooks/telemetry even though // the manager later stored the run as failed. const endModelError = getLastAssistantError(session); const endStatus = aborted ? "aborted" : softLimitReached ? "steered" : (!gatedResponseText && endModelError) ? "error" : "completed"; const endDecision = normalizeHookResponse( await options.hooks.dispatch("subagent:end", options.agentId ?? "unknown", { status: endStatus, ...(endStatus === "error" && endModelError ? { error: endModelError } : {}), tokensIn, tokensOut, turns: turnCount, responseText: gatedResponseText, attempt, maxAttempts, }), ); if (endDecision.action !== "block") break; if (attempt > revisionBudget) { throw new AgentRunnerError( endDecision.reason ?? "Blocked by hook", "aborted", { hook: "subagent:end", attempt, maxAttempts, ...(endDecision.feedback !== undefined ? { feedback: endDecision.feedback } : {}), }, ); } const revisionPrompt = endDecision.feedback?.trim() || "Your previous output was rejected by a quality gate. Revise and improve it."; await promptWithCircuitBreaker(session, revisionPrompt); gatedResponseText = collector.getText().trim() || getLastAssistantText(session); if (!getLastAssistantError(session)) exhaustedModelError = undefined; attempt++; } } } catch (err) { // A duration-quota (or token/tool) abort — or an external AbortSignal — // makes the active session.prompt() reject with an AbortError. That is the // *expected* outcome of a graceful stop — swallow it so runAgent resolves // with { aborted, timedOut } instead of throwing and masking the stop as an // error (CHE-19 for external signals; quota path already set aborted=true). if ((timedOut || aborted || options.signal?.aborted) && isAbortError(err)) { aborted = true; // fall through to the graceful return path below } else { options.hooks ?.dispatch("subagent:error", options.agentId ?? "unknown", { error: err instanceof Error ? err.message : String(err), }) .catch((err2) => { logger.debug(`Hook dispatch error: ${err2 instanceof Error ? err2.message : String(err2)}`); }); throw err; } } finally { unsubTurns(); collector.unsubscribe(); cleanupAbort(); } let responseText = gatedResponseText || collector.getText().trim() || getLastAssistantText(session); // Capture the raw-output signal BEFORE the fail-loud end-report substitution // below — the outcome contract (R4) must know whether the agent actually // produced text, not whether diagnostics were synthesized afterwards. const hasRawOutput = responseText.length > 0; let runError: string | undefined; const modelError = !aborted && !options.signal?.aborted ? (getLastAssistantError(session) ?? exhaustedModelError) : undefined; if (!responseText && modelError) { // The run produced no text because the model/provider errored on its final // turn (e.g. 401 auth, unavailable model). Surface it instead of returning // an empty "completed" result that hides the failure from the caller, and // record it on the result so the manager finalizes the agent as failed. logger.warn("Subagent stopped with a model/provider error", { agentId: options.agentId ?? "unknown", error: modelError, }); responseText = `Agent stopped with a model/provider error: ${modelError}`; runError = modelError; options.hooks ?.dispatch("subagent:error", options.agentId ?? "unknown", { error: modelError }) .catch((hookErr) => { logger.debug(`Hook dispatch error: ${hookErr instanceof Error ? hookErr.message : String(hookErr)}`); }); } if (!responseText && !modelError) { // Agent completed without producing text. Surface a machine-readable // end report so get_subagent_result never returns an uninformative // empty result. This covers cases where the agent finished silently // after a soft-limit steer or normal turn end. const durationMs = Math.round(performance.now() - startTime); responseText = `[pi-agent-orchestrator] Agent completed without producing output.\nType: ${type}\nTurns: ${turnCount}\nTool calls: ${toolCallCount}\nDuration: ${durationMs}ms\nStatus: ${aborted ? "aborted" : softLimitReached ? "steered" : "completed"}\n\nIf this was a route-discovery task, prefer direct shell/docs lookup.`; logger.warn("Subagent completed with empty output", { agentId: options.agentId ?? "unknown", type, turns: turnCount, toolCalls: toolCallCount, durationMs, }); } // Explicit outcome contract (R4): classify how the run ended so callers // never read a budget cut or a silent run as a successful empty result. const derived = deriveAgentOutcome({ aborted, abortReason, hasOutput: hasRawOutput, executedWork: toolCallCount > 0 || hasRawOutput, }); const duration = performance.now() - startTime; // Structured handoff parsing let handoff: AgentHandoff | undefined; if (!aborted && !options.signal?.aborted && agentConfig?.handoff) { const parsed = parseHandoff(responseText); if (parsed) { handoff = parsed; responseText = renderHandoffForParent(parsed); } } // Adversarial validation (extracted to agent-runner-validator.ts) let validationResults: ValidationResult[] | undefined; let validated: boolean | undefined; if (!aborted && !options.signal?.aborted && !options.skipValidators && hasValidators(agentConfig)) { const result = await runAdversarialValidation( session, ctx, responseText, agentConfig, options.agentId ?? "unknown", { pi: options.pi, model: options.model, signal: options.signal, hooks: options.hooks, onToolActivity: options.onToolActivity, onAssistantUsage: options.onAssistantUsage, onCompaction: options.onCompaction, onValidationComplete: options.onValidationComplete, runAgent, resumeAgent, }, ); responseText = result.responseText; validationResults = result.validationResults; validated = result.validated; } // Telemetry emitTelemetry("agent:completed", { type, duration, validatorResults: validationResults?.map((r) => ({ passed: r.passed, summary: r.summary })), }); const metrics: RunMetrics = { durationMs: duration, turns: turnCount, toolCalls: toolCallCount, tokensIn, tokensOut, tokensCacheWrite, latencyToFirstTokenMs: latencyToFirstToken, }; return { responseText, session, aborted, timedOut, steered: softLimitReached, validationResults, validated, handoff, metrics, error: runError, abortReason, outcome: derived.outcome, outcomeReason: derived.reason, }; } // ============================================================================ // Resume Agent // ============================================================================ export async function resumeAgent( session: AgentSession, prompt: string, options: { agentId?: string; hooks?: HookRegistry; onToolActivity?: (activity: ToolActivity) => void; onAssistantUsage?: (usage: { input: number; output: number; cacheWrite: number }) => void; onCompaction?: (info: CompactionSnapshot) => void; signal?: AbortSignal; inheritContext?: boolean; ctx?: ExtensionContext; } = {}, ): Promise { const collector = collectResponseText(session); const cleanupAbort = forwardAbortSignal(session, options.signal); const agentId = options.agentId ?? "unknown"; const unsubEvents = (options.onToolActivity || options.onAssistantUsage || options.onCompaction || options.hooks) ? session.subscribe((event: AgentSessionEvent) => { if (event.type === "tool_execution_start") options.onToolActivity?.({ type: "start", toolName: event.toolName }); if (event.type === "tool_execution_end") options.onToolActivity?.({ type: "end", toolName: event.toolName }); if (event.type === "message_end" && event.message.role === "assistant") { const msg = event.message as { usage?: { input?: number; output?: number; cacheWrite?: number } }; const u = msg.usage; if (u) options.onAssistantUsage?.({ input: u.input ?? 0, output: u.output ?? 0, cacheWrite: u.cacheWrite ?? 0 }); } if (event.type === "compaction_start") { options.hooks ?.dispatch("compaction:start", agentId, { reason: event.reason }) .catch((err) => { logger.debug(`Hook dispatch error: ${err instanceof Error ? err.message : String(err)}`); }); } if (event.type === "compaction_end") { const snapshot = buildCompactionSnapshot(event); options.onCompaction?.(snapshot); if (!event.aborted) { options.hooks ?.dispatch("compaction:end", agentId, { reason: event.reason, tokensBefore: snapshot.tokensBefore, tokensAfter: snapshot.tokensAfter, reductionPercent: snapshot.reductionPercent, aborted: false, willRetry: snapshot.willRetry, }) .catch((err) => { logger.debug(`Hook dispatch error: ${err instanceof Error ? err.message : String(err)}`); }); } } if ( event.type === "summarization_retry_scheduled" || event.type === "summarization_retry_attempt_start" || event.type === "summarization_retry_finished" ) { const attempt = event.type === "summarization_retry_scheduled" ? event.attempt : undefined; const maxAttempts = event.type === "summarization_retry_scheduled" ? event.maxAttempts : undefined; logger.debug(`Summarization retry (${event.type}) for ${agentId}`); emitTelemetry("subagent:summarization_retry", { agentId, phase: event.type, attempt, maxAttempts, }); } }) : () => {}; let effectivePrompt = prompt; if (options.inheritContext && options.ctx) { const parentContext = buildParentContext(options.ctx); if (parentContext) { effectivePrompt = parentContext + prompt; } } try { await promptWithCircuitBreaker(session, effectivePrompt); } finally { collector.unsubscribe(); unsubEvents(); cleanupAbort(); } return collector.getText().trim() || getLastAssistantText(session); } // ============================================================================ // Steering // ============================================================================ export async function steerAgent(session: AgentSession, message: string): Promise { await session.steer(message); } // ============================================================================ // Conversation Serialization // ============================================================================ export function getAgentConversation(session: AgentSession): string { const parts: string[] = []; for (const msg of session.messages) { if (msg.role === "user") { const text = typeof msg.content === "string" ? msg.content : extractText(msg.content); if (text.trim()) parts.push(`[User]: ${text.trim()}`); } else if (msg.role === "assistant") { const textParts: string[] = []; const toolCalls: string[] = []; for (const c of msg.content) { if (c.type === "text" && c.text) textParts.push(c.text); else if (c.type === "toolCall") toolCalls.push(` Tool: ${c.name ?? "unknown"}`); } if (textParts.length > 0) parts.push(`[Assistant]: ${textParts.join("\n")}`); if (toolCalls.length > 0) parts.push(`[Tool Calls]:\n${toolCalls.join("\n")}`); } else if (msg.role === "toolResult") { const text = extractText(msg.content); const truncated = text.length > 200 ? `${text.slice(0, 200)}...` : text; parts.push(`[Tool Result (${msg.toolName})]: ${truncated}`); } } return parts.join("\n\n"); } // ============================================================================ // Utility Exports // ============================================================================ export { globalCircuitBreaker };