import { createHash, randomUUID } from "node:crypto"; import type { OrchestratorConfig, ThinkingLevel } from "../src/core/config.js"; import type { McpCompletionCandidate } from "./routing.js"; import { classifyMcpProviderFailure, isUncertainMcpProviderFailure, type McpProviderDefiniteFailureCode, type McpProviderUncertainFailureCode, } from "./failureCodes.js"; export type ModelRole = "planner" | "judge"; export interface CompletionRequest { config: OrchestratorConfig; role: ModelRole; prompt: string; signal?: AbortSignal; } export interface RoutedCompletionRequest extends CompletionRequest { candidates: readonly McpCompletionCandidate[]; /** Frozen before the first durable attempt reservation or provider call. */ routingDecision?: RoutedCompletionDecision; /** Reject candidate output before accepting it, allowing an eligible fallback. */ validateText?: (text: string) => void; /** Called before each paid provider attempt. Callback failures stop fallback. */ beforeAttempt?: (attempt: RoutedCompletionAttempt) => void | Promise; /** Called after each provider attempt. Callback failures stop fallback. */ afterAttempt?: (result: RoutedCompletionAttemptResult) => void | Promise; } export interface RoutedCompletionAttempt { attempt: number; /** Stable, server-derived identity for the exact paid provider request. */ providerRequestRef: string; routingDecision: RoutedCompletionDecision; identity: { provider: string; model: string; family?: string }; thinking: ThinkingLevel; requestedOutputTokens?: number; estimatedCostUsd?: number; } export type RoutedCompletionAttemptResult = | (RoutedCompletionAttempt & { outcome: "succeeded" }) | (RoutedCompletionAttempt & { outcome: "failed"; failureCode: McpProviderDefiniteFailureCode }) | (RoutedCompletionAttempt & { outcome: "unknown"; failureCode: McpProviderUncertainFailureCode }); export interface RoutedCompletionResult { text: string; selectedIndex: number; fallbackHistory: Array<{ identity: string; reason: McpProviderDefiniteFailureCode }>; } export interface RoutedCompletionDecision { decisionId: string; policyVersion: string; policyDigest: string; configDigest: string; candidatesDigest: string; } const DEFAULT_LLM_TIMEOUT_MS = 120_000; const MAX_LLM_RESPONSE_BYTES = 2 * 1024 * 1024; const ANTHROPIC_DEFAULT_MAX_TOKENS = 8192; const ANTHROPIC_MAX_THINKING_MAX_TOKENS = 16384; let fakeLlmWarningPrinted = false; export async function completeWithRole({ config, role, prompt, signal }: CompletionRequest): Promise { return (await completeRouted({ config, role, prompt, signal, candidates: [config.roles[role]] })).text; } export async function completeRouted({ config, role, prompt, signal, candidates, routingDecision: suppliedRoutingDecision, validateText, beforeAttempt, afterAttempt, }: RoutedCompletionRequest): Promise { const routingDecision = suppliedRoutingDecision ?? freezeRoutingDecision( config, candidates, config.routing.version, `compat-${randomUUID()}`, ); const fallbackHistory: RoutedCompletionResult["fallbackHistory"] = []; for (let index = 0; index < candidates.length; index += 1) { const roleConfig = candidates[index]!; const attempt: RoutedCompletionAttempt = { attempt: index + 1, providerRequestRef: providerRequestReference(config, role, prompt, roleConfig, index + 1), routingDecision: structuredClone(routingDecision), identity: { provider: roleConfig.provider, model: roleConfig.model, ...(roleConfig.family === undefined ? {} : { family: roleConfig.family }), }, thinking: roleConfig.thinking, ...(roleConfig.requestedOutputTokens === undefined ? {} : { requestedOutputTokens: roleConfig.requestedOutputTokens }), ...(roleConfig.estimatedCostUsd === undefined ? {} : { estimatedCostUsd: roleConfig.estimatedCostUsd }), }; // These callbacks are outside the provider catch boundary deliberately: // a failed durable reservation/settlement must never trigger another call. await beforeAttempt?.(attempt); let text: string; try { text = await completeCandidate(config, role, roleConfig, prompt, signal); if (validateText) { try { validateText(text); } catch { throw new Error("Candidate output failed required schema validation"); } } } catch (error) { const safeSummary = sanitizeError(error, Object.values(config.mcp.providers).flatMap((provider) => provider.apiKey ? [provider.apiKey] : [])); const failureCode = classifyMcpProviderFailure(safeSummary); if (isUncertainMcpProviderFailure(failureCode)) { await afterAttempt?.({ ...attempt, outcome: "unknown", failureCode }); throw new Error(`Provider outcome is uncertain for ${roleConfig.provider}/${roleConfig.model}: ${safeSummary}`); } await afterAttempt?.({ ...attempt, outcome: "failed", failureCode }); fallbackHistory.push({ identity: `${roleConfig.provider}/${roleConfig.model}`, reason: failureCode }); if (signal?.aborted || index === candidates.length - 1) { throw new Error(`All eligible MCP completion candidates failed: ${roleConfig.provider}/${roleConfig.model}: ${safeSummary}`); } continue; } await afterAttempt?.({ ...attempt, outcome: "succeeded" }); return { text, selectedIndex: index, fallbackHistory }; } throw new Error("No MCP completion candidates were supplied"); } /** Freeze non-secret routing authority before any provider-attempt callback. */ export function freezeRoutingDecision( config: OrchestratorConfig, candidates: readonly McpCompletionCandidate[], policyVersion: string, decisionId: string, ): RoutedCompletionDecision { const providers = Object.fromEntries(Object.entries(config.mcp.providers).map(([name, provider]) => [ name, { baseUrl: provider.baseUrl, api: provider.api }, ])); return { decisionId, policyVersion, policyDigest: sha256(canonicalJson(config.routing)), configDigest: sha256(canonicalJson({ roles: config.roles, providers, models: config.mcp.models })), candidatesDigest: sha256(canonicalJson(candidates)), }; } /** * Hash only canonical, non-secret inputs that determine a provider call. The * reference is intentionally independent from the client mutation request ID: * one mutation may contain several paid fallback attempts. */ export function providerRequestReference( config: OrchestratorConfig, role: ModelRole, prompt: string, candidate: McpCompletionCandidate, attempt: number, ): string { const provider = config.mcp.providers[candidate.provider]; const canonical = { version: "mcp-provider-request-v1", operation: role === "planner" ? "plan" : "judge", role, attempt, promptHash: sha256(prompt), provider: { name: candidate.provider, api: provider?.api ?? null, baseUrl: provider?.baseUrl ?? null, }, candidate: { model: candidate.model, family: candidate.family ?? null, thinking: candidate.thinking, requestedOutputTokens: candidate.requestedOutputTokens ?? null, maxOutputTokens: candidate.maxOutputTokens ?? null, }, }; return sha256(JSON.stringify(canonical)); } function sha256(value: string): string { return createHash("sha256").update(value).digest("hex"); } function canonicalJson(value: unknown): string { return JSON.stringify(canonicalize(value)); } function canonicalize(value: unknown): unknown { if (Array.isArray(value)) return value.map(canonicalize); if (value !== null && typeof value === "object") { return Object.fromEntries(Object.entries(value as Record) .filter(([, item]) => item !== undefined) .sort(([left], [right]) => left < right ? -1 : left > right ? 1 : 0) .map(([key, item]) => [key, canonicalize(item)])); } return value; } async function completeCandidate(config: OrchestratorConfig, role: ModelRole, roleConfig: McpCompletionCandidate, prompt: string, signal?: AbortSignal): Promise { const provider = config.mcp.providers[roleConfig.provider]; if (!provider) { throw new Error(`No MCP provider configured for role ${role} provider ${roleConfig.provider}`); } if (!provider.apiKey) { throw new Error(`Missing API key for MCP provider ${roleConfig.provider}; set mcp.providers.${roleConfig.provider}.apiKey`); } const apiKey = provider.apiKey; if (process.env.AI_ORCH_FAKE_LLM === "1") { if (!fakeLlmWarningPrinted) { console.error("[ai-orchestrator-mcp] AI_ORCH_FAKE_LLM=1 is active; returning fake planner/judge responses."); fakeLlmWarningPrinted = true; } return fakeCompletion(role, prompt); } const text = await (async () => { switch (provider.api) { case "anthropic-messages": return anthropicMessages(provider.baseUrl, apiKey, roleConfig, prompt, signal); case "openai-responses": return openaiResponses(provider.baseUrl, apiKey, roleConfig, prompt, signal); case "openai-completions": return openaiCompletions(provider.baseUrl, apiKey, roleConfig, prompt, signal); default: throw new Error(`Unsupported MCP provider API for ${roleConfig.provider}: ${provider.api}`); } })(); if (text.trim().length === 0) { throw new Error(`LLM provider ${roleConfig.provider} returned an empty completion for role ${role}`); } return text; } function sanitizeError(error: unknown, secrets: string[]): string { const message = redactSecrets(error instanceof Error ? error.message : String(error), secrets); const httpStatus = /^LLM request failed \((\d+)/.exec(message)?.[1]; if (httpStatus) return `LLM request failed (${httpStatus})`; if (message.startsWith("LLM response was not valid JSON")) return "LLM response was not valid JSON"; if (message.startsWith("OpenAI response failed")) return "OpenAI response failed"; if (message.startsWith("OpenAI response was incomplete")) return "OpenAI response was incomplete"; if (message.startsWith("Candidate output failed")) return "candidate output failed required schema validation"; if (message.includes("timed out")) return "LLM request timed out"; if (message.includes("aborted by client")) return "LLM request aborted by client"; if (message.includes("truncated") || message.includes("token limit")) return "LLM response was truncated"; if (message.startsWith("No MCP provider configured")) return "MCP provider is not configured"; if (message.startsWith("Missing API key")) return "MCP provider API key is missing"; if (message.startsWith("Unsupported MCP provider API")) return "MCP provider API is unsupported"; if (message.includes("returned an empty completion")) return "LLM provider returned an empty completion"; if (message.includes("response exceeded")) return "LLM response exceeded the size limit"; return "MCP candidate failed"; } function fakeCompletion(role: ModelRole, prompt: string): string { if (role === "planner") { return [ "1. Inspect the relevant files and existing tests.", "2. Implement the requested behavior with minimal, focused changes.", "3. Run the detected project tests and fix any failures.", ].join("\n"); } if (prompt.includes("read-only DEBUG checker")) { return JSON.stringify({ rootCauseCategory: "implementation-defect", confidence: "high", summary: "The fake checker rejection identifies a local implementation defect.", repairScope: ["src"], validationRequirements: ["verification-tests"], topologyAssessment: "preserve", }); } const verdict = process.env.AI_ORCH_FAKE_LLM_VERDICT === "approve" ? "approve" : "reject"; return JSON.stringify({ verdict, reasons: `Fake judge ${verdict} for MCP protocol tests.`, ...(verdict === "reject" ? { requiredFixes: "Address the fake failing condition before retrying." } : {}), }); } async function anthropicMessages(baseUrl: string, apiKey: string, role: McpCompletionCandidate, prompt: string, signal?: AbortSignal): Promise { const allocation = anthropicTokenAllocation(role); const json = await postJson(endpoint(baseUrl, "/messages"), { "content-type": "application/json", "x-api-key": apiKey, "anthropic-version": "2023-06-01", }, { model: role.model, max_tokens: allocation.maxTokens, messages: [{ role: "user", content: prompt }], ...(allocation.thinkingBudget > 0 ? { thinking: { type: "enabled", budget_tokens: allocation.thinkingBudget } } : {}), }, signal, [apiKey]); const response = json as { content?: Array<{ type?: string; text?: string }>; stop_reason?: string }; if (response.stop_reason === "max_tokens") { throw new Error("Anthropic response was truncated at max_tokens; increase token budget or reduce prompt size"); } return (response.content ?? []).map((part) => (part.type === "text" ? part.text ?? "" : "")).join("\n").trim(); } async function openaiResponses(baseUrl: string, apiKey: string, role: McpCompletionCandidate, prompt: string, signal?: AbortSignal): Promise { const json = await postJson(endpoint(baseUrl, "/responses"), { "authorization": `Bearer ${apiKey}`, "content-type": "application/json", }, { model: role.model, input: prompt, max_output_tokens: completionTokenLimit(role, 8192), ...openaiResponsesReasoning(role.thinking), }, signal, [apiKey]); const response = json as { output_text?: string; status?: string; incomplete_details?: unknown; error?: unknown; output?: Array<{ content?: Array<{ text?: string; type?: string }> }>; }; if (response.status === "incomplete") { throw new Error(`OpenAI response was incomplete: ${redactSecrets(JSON.stringify(response.incomplete_details ?? {}), [apiKey])}`); } if (response.status === "failed") { throw new Error(`OpenAI response failed: ${redactSecrets(JSON.stringify(response.error ?? {}), [apiKey])}`); } if (typeof response.output_text === "string") { return response.output_text.trim(); } return (response.output ?? []) .flatMap((item) => item.content ?? []) .map((part) => (part.type === "output_text" ? part.text ?? "" : "")) .join("\n") .trim(); } async function openaiCompletions(baseUrl: string, apiKey: string, role: McpCompletionCandidate, prompt: string, signal?: AbortSignal): Promise { const reasoning = openaiCompletionsReasoning(role.thinking); const json = await postJson(endpoint(baseUrl, "/chat/completions"), { "authorization": `Bearer ${apiKey}`, "content-type": "application/json", }, { model: role.model, messages: [{ role: "user", content: prompt }], ...(reasoning.reasoning_effort ? { max_completion_tokens: completionTokenLimit(role, 8192), ...reasoning } : { max_tokens: completionTokenLimit(role, 4096), temperature: 0 }), }, signal, [apiKey]); const choice = (json as { choices?: Array<{ finish_reason?: string; message?: { content?: string } }> }).choices?.[0]; if (choice?.finish_reason === "length") { throw new Error("OpenAI chat completion was truncated at the token limit"); } return choice?.message?.content?.trim() ?? ""; } async function postJson( url: string, headers: Record, body: unknown, signal?: AbortSignal, secrets: string[] = [], ): Promise { const timeout = withTimeout(signal, DEFAULT_LLM_TIMEOUT_MS); try { const response = await fetch(url, { method: "POST", headers, body: JSON.stringify(body), signal: timeout.signal, redirect: "error" }); const text = await readBoundedResponseText(response); if (!response.ok) { throw new Error( `LLM request failed (${response.status} ${response.statusText}): ${redactSecrets(text, secrets).slice(0, 1000)}`, ); } if (!text) return {}; try { return JSON.parse(text); } catch { throw new Error(`LLM response was not valid JSON: ${redactSecrets(text, secrets).slice(0, 1000)}`); } } catch (error) { if (timeout.timedOut()) { throw new Error(`LLM request timed out after ${DEFAULT_LLM_TIMEOUT_MS}ms`); } if (signal?.aborted) { throw new Error("LLM request aborted by client"); } throw error; } finally { timeout.cleanup(); } } async function readBoundedResponseText(response: Response): Promise { const declaredLength = Number(response.headers.get("content-length")); if (Number.isFinite(declaredLength) && declaredLength > MAX_LLM_RESPONSE_BYTES) { await response.body?.cancel(); throw new Error(`LLM response exceeded ${MAX_LLM_RESPONSE_BYTES} bytes`); } if (!response.body) return ""; const reader = response.body.getReader(); const decoder = new TextDecoder(); let bytes = 0; let text = ""; while (true) { const { done, value } = await reader.read(); if (done) break; bytes += value.byteLength; if (bytes > MAX_LLM_RESPONSE_BYTES) { await reader.cancel(); throw new Error(`LLM response exceeded ${MAX_LLM_RESPONSE_BYTES} bytes`); } text += decoder.decode(value, { stream: true }); } return text + decoder.decode(); } function withTimeout(signal: AbortSignal | undefined, timeoutMs: number): { signal: AbortSignal; cleanup: () => void; timedOut: () => boolean; } { const controller = new AbortController(); let didTimeOut = false; const abortFromCaller = (): void => controller.abort(); const timer = setTimeout(() => { didTimeOut = true; controller.abort(); }, timeoutMs); if (signal?.aborted) { controller.abort(); } else { signal?.addEventListener("abort", abortFromCaller, { once: true }); } return { signal: controller.signal, cleanup: () => { clearTimeout(timer); signal?.removeEventListener("abort", abortFromCaller); }, timedOut: () => didTimeOut, }; } function redactSecrets(text: string, secrets: string[]): string { let redacted = text; for (const secret of secrets) { if (secret.length === 0) continue; redacted = redacted.split(secret).join("[redacted]"); } return redacted; } function endpoint(baseUrl: string, suffix: string): string { const url = new URL(baseUrl); const path = url.pathname.replace(/\/+$/, ""); url.pathname = path.endsWith(suffix) ? path : `${path}${suffix}`; return url.toString(); } function openaiResponsesReasoning(thinking: ThinkingLevel): Record { const effort = openaiEffort(thinking); return effort ? { reasoning: { effort } } : {}; } function openaiCompletionsReasoning(thinking: ThinkingLevel): Record { const effort = openaiEffort(thinking); return effort ? { reasoning_effort: effort } : {}; } function openaiEffort(thinking: ThinkingLevel): "minimal" | "low" | "medium" | "high" | undefined { if (thinking === "off") { return undefined; } return thinking === "minimal" ? "minimal" : thinking === "low" ? "low" : thinking === "medium" ? "medium" : "high"; } function anthropicTokenAllocation(role: McpCompletionCandidate): { maxTokens: number; thinkingBudget: number } { const totalLimit = Math.max(1, Math.min(role.maxOutputTokens ?? anthropicMaxTokens(role.thinking), anthropicMaxTokens(role.thinking))); const requestedVisible = Math.max(1, Math.min(role.requestedOutputTokens ?? Math.min(4_096, totalLimit), totalLimit)); if (role.thinking === "off" || role.thinking === "minimal" || role.thinking === "low") { return { maxTokens: requestedVisible, thinkingBudget: 0 }; } const budgetByLevel: Record = { off: 0, minimal: 0, low: 0, medium: 1024, high: 2048, xhigh: 4096, max: 8192, }; const availableThinking = Math.max(0, totalLimit - requestedVisible); const thinkingBudget = availableThinking >= 1_024 ? Math.min(budgetByLevel[role.thinking], availableThinking) : 0; return { maxTokens: requestedVisible + thinkingBudget, thinkingBudget }; } function completionTokenLimit(role: McpCompletionCandidate, apiLimit: number): number { return Math.max(1, Math.min(role.maxOutputTokens ?? apiLimit, role.requestedOutputTokens ?? apiLimit, apiLimit)); } function anthropicMaxTokens(thinking: ThinkingLevel): number { return thinking === "max" ? ANTHROPIC_MAX_THINKING_MAX_TOKENS : ANTHROPIC_DEFAULT_MAX_TOKENS; }