import * as crypto from "node:crypto"; import * as fsSync from "node:fs"; import * as fs from "node:fs/promises"; import * as path from "node:path"; import { scheduler } from "node:timers/promises"; import { isOfficialAnthropicApiUrl } from "@oh-my-pi/pi-catalog/compat/anthropic"; import type { Effort } from "@oh-my-pi/pi-catalog/effort"; import { isVertexExpressOpenAIUrl, isVertexRawPredictUrl, resolveVertexEndpointHost } from "@oh-my-pi/pi-catalog/hosts"; import { defaultSupportedEffort, mapEffortToAnthropicAdaptiveEffort, mapEffortToGoogleThinkingLevel, requireSupportedEffort, resolveWireModelId, } from "@oh-my-pi/pi-catalog/model-thinking"; import { providerEntries } from "@oh-my-pi/pi-catalog/compat/providers"; import { CODEX_BASE_URL } from "@oh-my-pi/pi-catalog/wire/codex"; import { $env, $pickenv, getProviderInFlightRoot, isEnoent, logger, untilAborted } from "@oh-my-pi/pi-utils"; import { getCustomApi } from "./api-registry"; import { createAuthRetryKeyState, isApiKeyResolver, resolvedApiKeyBearer, resolveNextAuthRetryKey } from "./auth-retry"; import type { OAuthRequestIdentity } from "./auth/types"; import * as AIError from "./error"; import { ProviderHttpError } from "./error"; import { isConcurrencyCapExclusion, isUsageLimitOutcome } from "./error/rate-limit"; import type { BedrockOptions } from "./providers/amazon-bedrock"; import type { AnthropicOptions } from "./providers/anthropic"; import type { AppleFoundationModelsOptions } from "./providers/apple-foundation-models"; import type { CursorOptions } from "./providers/cursor"; import type { DevinOptions } from "./providers/devin"; import { type FactoryDroidOptions, streamFactoryDroid } from "./providers/factory-droid"; import { streamGitLabDuo } from "./providers/gitlab-duo"; import { type GitLabDuoWorkflowOptions, streamGitLabDuoWorkflow } from "./providers/gitlab-duo-workflow"; import type { GoogleOptions } from "./providers/google"; import { getVertexAccessToken } from "./providers/google-auth"; import type { GoogleGeminiCliOptions } from "./providers/google-gemini-cli"; import type { GoogleVertexOptions } from "./providers/google-vertex"; import { streamKimi } from "./providers/kimi"; import type { OllamaChatOptions } from "./providers/ollama"; import type { OpenAICompletionsOptions } from "./providers/openai-completions"; import { streamPiNative } from "./providers/pi-native-client"; import { streamSynthetic } from "./providers/synthetic"; import { streamAnthropic, streamAppleFoundationModels, streamAzureOpenAIResponses, streamBedrock, streamCursor, streamDevin, streamGoogle, streamGoogleGeminiCli, streamGoogleVertex, streamOllama, streamOpenAICodexResponses, streamOpenAICompletions, streamOpenAIResponses, } from "./providers/register-builtins"; import { getProviderDefinition, PROVIDER_REGISTRY } from "./registry"; import type { Api, AssistantMessage, AssistantMessageEvent, Context, FetchImpl, Model, OptionsForApi, SimpleStreamOptions, StreamOptions, ThinkingBudgets, ToolChoice, } from "./types"; import { getHeaderCaseInsensitive, resolveCacheRetention } from "./utils"; import { AssistantMessageEventStream } from "./utils/event-stream"; import { isFoundryEnabled } from "./utils/foundry"; import { applyGlyphCodec } from "./utils/glyph-codec"; import { wrapLeakedThinkingStream } from "./utils/leaked-thinking-stream"; import { withThinkingLoopGuard } from "./utils/thinking-loop"; import { withTransportFetch } from "./utils/transport-fetch"; function isGoogleVertexAuthenticatedModel(model: Model): boolean { return ( model.provider === "google-vertex" && ((model.api === "openai-completions" && isVertexExpressOpenAIUrl(model.baseUrl)) || (model.api === "anthropic-messages" && isVertexRawPredictUrl(model.baseUrl))) ); } /** * Whether {@link model} is an official first-party endpoint whose stream needs * no leaked-thinking healing — the official Anthropic API and the official * OpenAI / OpenAI-Codex endpoints return structured thinking blocks and never * leak reasoning idioms into the visible text channel. * * The gate is provider id **and** official endpoint URL: pointing * `provider: "anthropic"` (or `openai`) at a custom proxy via `models.yml` * still routes through {@link wrapLeakedThinkingStream}, since a third-party * gateway may well leak. URL checks are strict (exact origin / path boundary * or parsed hostname) — a substring match would accept lookalikes like * `https://api.openai.com.evil/`. Anthropic Foundry (`CLAUDE_CODE_USE_FOUNDRY`) * redirects an empty `baseUrl` to `FOUNDRY_BASE_URL`, so the check runs against * that effective endpoint — exempt only when it resolves to the official host. */ function isLeakedThinkingHealExempt(model: Model): boolean { switch (model.provider) { case "anthropic": { // Mirror resolveAnthropicBaseUrl's effective endpoint: Foundry redirects // an empty baseUrl to FOUNDRY_BASE_URL; otherwise an explicit non-official // model.baseUrl wins, then the ANTHROPIC_BASE_URL gateway fallback, then // the official default. Exempt only when the effective endpoint is official. if (isFoundryEnabled()) { const foundry = $env.FOUNDRY_BASE_URL?.trim(); if (foundry) return isOfficialAnthropicApiUrl(foundry); } if (model.baseUrl && !isOfficialAnthropicApiUrl(model.baseUrl)) return false; return isOfficialAnthropicApiUrl($env.ANTHROPIC_BASE_URL?.trim() || model.baseUrl); } case "openai": return isOfficialOpenAIApiUrl(model.baseUrl); case "openai-codex": return isOfficialCodexApiUrl(model.baseUrl); default: return false; } } /** Strict official-OpenAI endpoint check; missing baseUrl defaults to `api.openai.com`. */ function isOfficialOpenAIApiUrl(baseUrl: string | undefined): boolean { if (!baseUrl) return true; try { return new URL(baseUrl).hostname === "api.openai.com"; } catch { return false; } } const OFFICIAL_CODEX_URL = new URL(CODEX_BASE_URL); /** Strict official-Codex endpoint check; exact origin or a path boundary after {@link CODEX_BASE_URL}. */ export function isOfficialCodexApiUrl(baseUrl: string | undefined): boolean { if (!baseUrl) return true; try { const candidate = new URL(baseUrl); const candidatePath = candidate.pathname.replace(/\/+$/, ""); return ( candidate.origin === OFFICIAL_CODEX_URL.origin && (candidatePath === OFFICIAL_CODEX_URL.pathname || candidatePath.startsWith(`${OFFICIAL_CODEX_URL.pathname}/`)) ); } catch { return false; } } /** * Apply live leaked-thinking healing unless {@link model} is an official * first-party endpoint ({@link isLeakedThinkingHealExempt}), which emits * structured thinking and needs no healing. */ function healLeakedThinking(model: Model, inner: AssistantMessageEventStream): AssistantMessageEventStream { return isLeakedThinkingHealExempt(model) ? inner : wrapLeakedThinkingStream(inner); } type ProviderInFlightLease = { path: string; stopHeartbeat: () => Promise; }; type ProviderInFlightLeaseInfo = { pid: number; timestamp: number; token: string; }; type ProviderInFlightStaleLock = { token: string } | { mtimeMs: number }; type ProviderInFlightLockIdentity = { dev: number; ino: number; birthtimeMs: number }; const PROVIDER_INFLIGHT_LOCK_STALE_MS = 10_000; const PROVIDER_INFLIGHT_LEASE_STALE_MS = 30_000; const PROVIDER_INFLIGHT_HEARTBEAT_MS = 5_000; const PROVIDER_INFLIGHT_SIGNAL_FALLBACK_MS = 250; const PROVIDER_INFLIGHT_HEARTBEAT_FLUSH_TIMEOUT_MS = 1_000; const PROVIDER_INFLIGHT_RELEASE_TIMEOUT_MS = 5_000; const PROVIDER_INFLIGHT_LOCK_RETRY_INITIAL_MS = 25; const PROVIDER_INFLIGHT_LOCK_RETRY_MAX_MS = 250; const PROVIDER_INFLIGHT_LOCK_RETRY_BUDGET_MS = 3_000; let configuredProviderMaxInFlightRequests: Record = {}; let providerInFlightRootOverride: string | undefined; let providerInFlightHeartbeatMsOverride: number | undefined; let providerInFlightHeartbeatFlushTimeoutMsOverride: number | undefined; let providerInFlightHeartbeatWriterOverride: | ((writeProviderInFlightInfo: () => Promise) => Promise) | undefined; let providerInFlightLeaseRemoverOverride: ((leasePath: string) => Promise) | undefined; let providerInFlightWaitObserverOverride: ((provider: string) => void) | undefined; let providerInFlightLockCreatedObserverOverride: ((lockDir: string) => Promise) | undefined; let providerInFlightLockIdentifiedObserverOverride: ((lockDir: string) => Promise) | undefined; let providerInFlightLockMkdirOverride: ((lockDir: string) => Promise) | undefined; let providerInFlightLockPlatformOverride: NodeJS.Platform | undefined; let providerInFlightLockRetryTimings: { budgetMs?: number; initialDelayMs?: number; maxDelayMs?: number } | undefined; export function configureProviderMaxInFlightRequests(limits: Record | undefined): void { configuredProviderMaxInFlightRequests = limits ?? {}; } function resolveProviderInFlightLimit( provider: string, options?: Pick, ): number | undefined { const limits = options?.maxInFlightRequests ?? configuredProviderMaxInFlightRequests; const value = limits[provider]; if (typeof value !== "number" || !Number.isFinite(value) || value <= 0) return undefined; return Math.max(1, Math.floor(value)); } function providerInFlightRoot(): string { if (providerInFlightRootOverride) return providerInFlightRootOverride; return getProviderInFlightRoot(); } function providerInFlightSegment(provider: string): string { return Bun.SHA256.hash(provider, "base64url"); } function providerInFlightDir(provider: string): string { return path.join(providerInFlightRoot(), providerInFlightSegment(provider)); } function providerInFlightSignalPath(provider: string): string { return path.join(providerInFlightDir(provider), ".wakeup"); } function providerInFlightLockDir(provider: string): string { return `${providerInFlightDir(provider)}.lock`; } // `process.kill(pid, 0)` may throw for permission/sandbox reasons even when a // process exists. Treat non-ESRCH failures as alive; timestamp expiry still // reaps leases whose heartbeat stopped. function isProcessAlive(pid: number): boolean { try { process.kill(pid, 0); return true; } catch (error) { return (error as NodeJS.ErrnoException).code !== "ESRCH"; } } async function readProviderInFlightInfo(infoPath: string): Promise { try { const content = await fs.readFile(infoPath, "utf-8"); const parsed = JSON.parse(content) as Partial; if (typeof parsed.pid !== "number" || typeof parsed.timestamp !== "number" || typeof parsed.token !== "string") { return null; } return { pid: parsed.pid, timestamp: parsed.timestamp, token: parsed.token }; } catch { return null; } } async function writeProviderInFlightInfo(dir: string, token: string): Promise { const info: ProviderInFlightLeaseInfo = { pid: process.pid, timestamp: Date.now(), token }; const infoPath = path.join(dir, "info.json"); const tempPath = path.join(dir, `.info-${process.pid}-${crypto.randomUUID()}.tmp`); try { // Unlike Bun.write, fs.writeFile does not recreate a lease directory that // was removed while a timed-out heartbeat was still pending. await fs.writeFile(tempPath, JSON.stringify(info), "utf8"); await fs.rename(tempPath, infoPath); } catch (error) { await fs.rm(tempPath, { force: true }).catch(() => {}); throw error; } } async function isProviderInFlightDirStale(dir: string, staleMs: number): Promise { const info = await readProviderInFlightInfo(path.join(dir, "info.json")); if (info) { if (!isProcessAlive(info.pid)) return true; return Date.now() - info.timestamp > staleMs; } try { const stat = await fs.stat(path.join(dir, "info.json")); return Date.now() - stat.mtimeMs > staleMs; } catch (error) { if (!isEnoent(error)) throw error; } try { const stat = await fs.stat(dir); return Date.now() - stat.mtimeMs > staleMs; } catch (error) { if (isEnoent(error)) return false; throw error; } } async function readProviderInFlightStaleLock(lockDir: string): Promise { const infoPath = path.join(lockDir, "info.json"); const info = await readProviderInFlightInfo(infoPath); if (info) return isProcessAlive(info.pid) ? null : { token: info.token }; try { const stat = await fs.stat(lockDir); return Date.now() - stat.mtimeMs > PROVIDER_INFLIGHT_LOCK_STALE_MS ? { mtimeMs: stat.mtimeMs } : null; } catch (error) { if (isEnoent(error)) return null; throw error; } } async function readProviderInFlightLockIdentity(lockDir: string): Promise { const stat = await fs.stat(lockDir); return { dev: stat.dev, ino: stat.ino, birthtimeMs: stat.birthtimeMs }; } function isSameProviderInFlightLock( current: ProviderInFlightLockIdentity, expected: ProviderInFlightLockIdentity, ): boolean { if (current.dev !== expected.dev) return false; if (current.ino !== 0 || expected.ino !== 0) return current.ino === expected.ino; return current.birthtimeMs === expected.birthtimeMs; } async function releaseProviderInFlightStaleLock(lockDir: string, stale: ProviderInFlightStaleLock): Promise { if ("token" in stale) { await releaseProviderInFlightLock(lockDir, stale.token); return; } const infoPath = path.join(lockDir, "info.json"); if (await readProviderInFlightInfo(infoPath)) return; try { const stat = await fs.stat(lockDir); if (stat.mtimeMs !== stale.mtimeMs || Date.now() - stat.mtimeMs <= PROVIDER_INFLIGHT_LOCK_STALE_MS) return; await fs.rm(lockDir, { recursive: true, force: true }); } catch {} } // Best-effort token-checked release. A token mismatch means another process has // already replaced the lock, so the fresh lock must be left intact. async function releaseProviderInFlightLock(lockDir: string, token: string): Promise { try { const info = await readProviderInFlightInfo(path.join(lockDir, "info.json")); if (!info || info.token !== token) return; await fs.rm(lockDir, { recursive: true, force: true }); } catch {} } async function releaseProviderInFlightLockDirIfSame( lockDir: string, identity: ProviderInFlightLockIdentity, ): Promise { try { if (await readProviderInFlightInfo(path.join(lockDir, "info.json"))) return; const current = await readProviderInFlightLockIdentity(lockDir); if (!isSameProviderInFlightLock(current, identity)) return; await fs.rm(lockDir, { recursive: true, force: true }); } catch {} } // On Windows, mkdir on a lock directory another process is concurrently deleting // (delete-pending while a watcher, antivirus, or indexer still holds it) fails // with ERROR_ACCESS_DENIED — mapped to EPERM/EACCES — instead of EEXIST, so // ordinary contention looks like a permission error. On win32 only, back off and // retry briefly; if the failure persists past the budget, rethrow the original // error so a real permission problem still surfaces. Every other outcome // (success, EEXIST, any other code) exits immediately, keeping non-win32 and // non-transient behavior unchanged. function isTransientProviderInFlightLockMkdirError(error: NodeJS.ErrnoException): boolean { if ((providerInFlightLockPlatformOverride ?? process.platform) !== "win32") return false; return error.code === "EPERM" || error.code === "EACCES"; } async function mkdirProviderInFlightLockDir(lockDir: string, signal?: AbortSignal): Promise { const mkdir = providerInFlightLockMkdirOverride ?? ((dir: string) => fs.mkdir(dir)); let failures = 0; let since = Date.now(); while (true) { try { await mkdir(lockDir); return; } catch (error) { if (!isTransientProviderInFlightLockMkdirError(error as NodeJS.ErrnoException)) throw error; const now = Date.now(); if (failures === 0) since = now; const budgetMs = providerInFlightLockRetryTimings?.budgetMs ?? PROVIDER_INFLIGHT_LOCK_RETRY_BUDGET_MS; if (now - since >= budgetMs) throw error; const initialMs = providerInFlightLockRetryTimings?.initialDelayMs ?? PROVIDER_INFLIGHT_LOCK_RETRY_INITIAL_MS; const maxMs = providerInFlightLockRetryTimings?.maxDelayMs ?? PROVIDER_INFLIGHT_LOCK_RETRY_MAX_MS; const delayMs = Math.min(maxMs, initialMs * 2 ** failures); failures++; await untilAborted(signal, Bun.sleep(delayMs)); } } } async function acquireProviderInFlightLock(provider: string, signal?: AbortSignal): Promise<() => Promise> { const lockDir = providerInFlightLockDir(provider); await fs.mkdir(path.dirname(lockDir), { recursive: true }); while (true) { if (signal?.aborted) throw signal.reason ?? new AIError.AbortError("Provider request aborted before dispatch"); try { await mkdirProviderInFlightLockDir(lockDir, signal); await providerInFlightLockCreatedObserverOverride?.(lockDir); let lockIdentity: ProviderInFlightLockIdentity; try { lockIdentity = await readProviderInFlightLockIdentity(lockDir); } catch (error) { if (isEnoent(error)) continue; throw error; } const token = crypto.randomUUID(); try { await providerInFlightLockIdentifiedObserverOverride?.(lockDir); await writeProviderInFlightInfo(lockDir, token); } catch (error) { await releaseProviderInFlightLockDirIfSame(lockDir, lockIdentity); if (isEnoent(error)) continue; throw error; } return async () => { await releaseProviderInFlightLock(lockDir, token); }; } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; } const staleLock = await readProviderInFlightStaleLock(lockDir); if (staleLock) { await releaseProviderInFlightStaleLock(lockDir, staleLock); await signalProviderInFlightWaiters(provider); continue; } await waitForProviderInFlightSignal(provider, signal); } } async function cleanupProviderInFlightLeases(providerDir: string): Promise { let active = 0; let entries: string[]; try { entries = await fs.readdir(providerDir); } catch (error) { if (isEnoent(error)) return 0; throw error; } for (const entry of entries) { const leaseDir = path.join(providerDir, entry); let isDirectory = false; try { isDirectory = (await fs.stat(leaseDir)).isDirectory(); } catch (error) { if (isEnoent(error)) continue; throw error; } if (!isDirectory) continue; if (await isProviderInFlightDirStale(leaseDir, PROVIDER_INFLIGHT_LEASE_STALE_MS)) { await fs.rm(leaseDir, { recursive: true, force: true }); continue; } active++; } return active; } async function tryAcquireProviderInFlightLease( provider: string, limit: number, signal?: AbortSignal, ): Promise { const releaseLock = await acquireProviderInFlightLock(provider, signal); try { const dir = providerInFlightDir(provider); await fs.mkdir(dir, { recursive: true }); const active = await cleanupProviderInFlightLeases(dir); if (active >= limit) return null; const leaseDir = path.join(dir, `${process.pid}-${Date.now()}-${crypto.randomUUID()}`); const token = crypto.randomUUID(); try { await fs.mkdir(leaseDir); await writeProviderInFlightInfo(leaseDir, token); } catch (error) { await removeProviderInFlightLeaseDir(leaseDir).catch(() => {}); throw error; } let heartbeatActive = true; let heartbeatFlush = Promise.resolve(); const touchHeartbeat = () => { if (!heartbeatActive) return; heartbeatFlush = heartbeatFlush .then(async () => { if (!heartbeatActive) return; const write = () => { if (!heartbeatActive) return Promise.resolve(); return writeProviderInFlightInfo(leaseDir, token); }; if (providerInFlightHeartbeatWriterOverride) { await providerInFlightHeartbeatWriterOverride(write); } else { await write(); } }) .catch(() => {}); }; const heartbeat = setInterval( touchHeartbeat, providerInFlightHeartbeatMsOverride ?? PROVIDER_INFLIGHT_HEARTBEAT_MS, ); heartbeat.unref?.(); return { path: leaseDir, stopHeartbeat: () => { heartbeatActive = false; clearInterval(heartbeat); return heartbeatFlush; }, }; } finally { await releaseLock(); } } async function signalProviderInFlightWaitersInDir(dir: string): Promise { try { await fs.mkdir(dir, { recursive: true }); await Bun.write(path.join(dir, ".wakeup"), String(Date.now())); } catch {} } async function signalProviderInFlightWaiters(provider: string): Promise { await signalProviderInFlightWaitersInDir(providerInFlightDir(provider)); } function waitForProviderInFlightSignal(provider: string, signal?: AbortSignal): Promise { if (signal?.aborted) return Promise.reject(signal.reason ?? new AIError.AbortError("Provider request aborted before dispatch")); const signalPath = providerInFlightSignalPath(provider); providerInFlightWaitObserverOverride?.(provider); const waitStarted = Date.now(); const { promise, resolve, reject } = Promise.withResolvers(); let settled = false; let watcher: fsSync.FSWatcher | undefined; const timer = setTimeout(() => finish(resolve), PROVIDER_INFLIGHT_SIGNAL_FALLBACK_MS); const finish = (settle: () => void) => { if (settled) return; settled = true; clearTimeout(timer); watcher?.close(); signal?.removeEventListener("abort", onAbort); settle(); }; const onAbort = () => { finish(() => reject(signal?.reason ?? new AIError.AbortError("Provider request aborted before dispatch"))); }; signal?.addEventListener("abort", onAbort, { once: true }); try { watcher = fsSync.watch(providerInFlightDir(provider), (_event, filename) => { if (filename === ".wakeup" || filename === null) { finish(resolve); } }); void fs.stat(signalPath).then( stat => { if (stat.mtimeMs >= waitStarted) finish(resolve); }, error => { if (!isEnoent(error)) finish(resolve); }, ); } catch { // Filesystem notifications are best-effort across platforms; the fallback // timer keeps stale-lock/lease cleanup progressing if an event is dropped. } return promise; } async function removeProviderInFlightLeaseDir(leasePath: string): Promise { for (let attempt = 0; attempt < 3; attempt++) { try { await fs.rm(leasePath, { recursive: true, force: true }); return; } catch (error) { if (isEnoent(error)) return; const code = (error as NodeJS.ErrnoException).code; if (attempt < 2 && (code === "EBUSY" || code === "ENOTEMPTY" || code === "EPERM")) { await Bun.sleep(25); continue; } throw error; } } } // Signal into the lease's OWN provider directory (derived from `lease.path`) // rather than recomputing it from the current root. A release that lands after // the in-flight root has been repointed (only the test seam does that) must not // write `.wakeup` into an unrelated provider directory. async function releaseProviderInFlightLease(lease: ProviderInFlightLease): Promise { const heartbeatFlush = lease.stopHeartbeat(); const flushTimeout = Promise.withResolvers<"timeout">(); const flushTimer = setTimeout( () => flushTimeout.resolve("timeout"), providerInFlightHeartbeatFlushTimeoutMsOverride ?? PROVIDER_INFLIGHT_HEARTBEAT_FLUSH_TIMEOUT_MS, ); flushTimer.unref?.(); try { const outcome = await Promise.race([heartbeatFlush.then(() => "flushed" as const), flushTimeout.promise]); if (outcome === "timeout") { logger.warn("Provider in-flight heartbeat flush timed out; forcing lease cleanup", { path: lease.path }); } } finally { clearTimeout(flushTimer); } const releaseTimeout = Promise.withResolvers(); const releaseTimer = setTimeout( () => releaseTimeout.reject(new Error("Provider in-flight lease cleanup timed out")), PROVIDER_INFLIGHT_RELEASE_TIMEOUT_MS, ); releaseTimer.unref?.(); try { const removeLease = providerInFlightLeaseRemoverOverride ?? removeProviderInFlightLeaseDir; await Promise.race([removeLease(lease.path), releaseTimeout.promise]); } finally { clearTimeout(releaseTimer); } // Wake-up is an optimization: waiters also poll every 250 ms. Do not let a // notification-file stall keep a completed provider request open. void signalProviderInFlightWaitersInDir(path.dirname(lease.path)); } async function acquireProviderInFlightSlot( provider: string, limit: number | undefined, signal?: AbortSignal, ): Promise<() => Promise> { if (limit === undefined) return async () => {}; let loggedWait = false; while (true) { if (signal?.aborted) throw signal.reason ?? new AIError.AbortError("Provider request aborted before dispatch"); const lease = await tryAcquireProviderInFlightLease(provider, limit, signal); if (lease) return () => releaseProviderInFlightLease(lease); if (!loggedWait) { loggedWait = true; logger.debug("Provider in-flight limit blocked request", { provider, limit }); } await waitForProviderInFlightSignal(provider, signal); } } export const __providerInFlightForTesting = { setRoot(root: string | undefined): void { providerInFlightRootOverride = root; }, setHeartbeatTimings(timings: { heartbeatMs?: number; heartbeatFlushTimeoutMs?: number } | undefined): void { providerInFlightHeartbeatMsOverride = timings?.heartbeatMs; providerInFlightHeartbeatFlushTimeoutMsOverride = timings?.heartbeatFlushTimeoutMs; }, setHeartbeatWriter(writer: ((writeProviderInFlightInfo: () => Promise) => Promise) | undefined): void { providerInFlightHeartbeatWriterOverride = writer; }, setLeaseRemover(remover: ((leasePath: string) => Promise) | undefined): void { providerInFlightLeaseRemoverOverride = remover; }, setWaitObserver(observer: ((provider: string) => void) | undefined): void { providerInFlightWaitObserverOverride = observer; }, setLockCreatedObserver(observer: ((lockDir: string) => Promise) | undefined): void { providerInFlightLockCreatedObserverOverride = observer; }, setLockIdentifiedObserver(observer: ((lockDir: string) => Promise) | undefined): void { providerInFlightLockIdentifiedObserverOverride = observer; }, setLockMkdirOverride(mkdir: ((lockDir: string) => Promise) | undefined): void { providerInFlightLockMkdirOverride = mkdir; }, setLockPlatformOverride(platform: NodeJS.Platform | undefined): void { providerInFlightLockPlatformOverride = platform; }, setLockRetryTimings(timings: { budgetMs?: number; initialDelayMs?: number; maxDelayMs?: number } | undefined): void { providerInFlightLockRetryTimings = timings; }, providerDir(provider: string): string { return providerInFlightDir(provider); }, lockDir(provider: string): string { return providerInFlightLockDir(provider); }, async captureStaleLockRelease(provider: string): Promise<(() => Promise) | null> { const lockDir = providerInFlightLockDir(provider); const stale = await readProviderInFlightStaleLock(lockDir); if (!stale) return null; return () => releaseProviderInFlightStaleLock(lockDir, stale); }, async captureLockDirRelease(provider: string): Promise<(() => Promise) | null> { const lockDir = providerInFlightLockDir(provider); try { const identity = await readProviderInFlightLockIdentity(lockDir); return () => releaseProviderInFlightLockDirIfSame(lockDir, identity); } catch { return null; } }, }; function withProviderInFlightLimit>( model: Model, options: TOptions | undefined, dispatch: () => AssistantMessageEventStream, ): AssistantMessageEventStream { // Leaked-thinking healing folds in here — the one shared provider-dispatch // chokepoint — so the loop guard (which wraps this) sees healed events and all // provider exits are covered by one wrap. Official first-party providers are // exempt (see `healLeakedThinking`); healing is otherwise idempotent. const limit = resolveProviderInFlightLimit(model.provider, options); if (limit === undefined) return healLeakedThinking(model, dispatch()); const outer = new AssistantMessageEventStream(); void (async () => { let release: (() => Promise) | undefined; let releasePromise: Promise | undefined; const releaseOnce = () => { if (!release) return Promise.resolve(); releasePromise ??= release(); return releasePromise; }; const releaseBestEffort = async () => { try { await releaseOnce(); } catch (releaseError) { // The lease has stopped heartbeating and stale cleanup will reap it // within PROVIDER_INFLIGHT_LEASE_STALE_MS. Until then, its slot may // remain unavailable and waiters rely on the fallback poll. // Never replace a completed response or the provider's original error // with a coordination-directory cleanup failure. logger.warn("Provider in-flight permit release failed", { provider: model.provider, error: String(releaseError), }); } }; try { const startedWaitingAt = Date.now(); release = await acquireProviderInFlightSlot(model.provider, limit, options?.signal); if (Date.now() - startedWaitingAt >= PROVIDER_INFLIGHT_SIGNAL_FALLBACK_MS) { logger.debug("Provider in-flight limit wait completed", { provider: model.provider, limit }); } if (options?.signal?.aborted) { throw options.signal.reason ?? new AIError.AbortError("Provider request aborted before dispatch"); } const inner = healLeakedThinking(model, dispatch()); let terminalEvent: AssistantMessageEvent | undefined; for await (const event of inner) { if (event.type === "done" || event.type === "error") { terminalEvent = event; break; } outer.push(event); if (outer.done) { await releaseBestEffort(); return; } } const result = await inner.result(); // Releasing the permit is part of request completion. Publishing the // result first lets an immediate follow-up turn contend with its own // still-live lease, which is particularly costly on Windows. await releaseBestEffort(); if (!outer.done) { if (terminalEvent) outer.push(terminalEvent); else outer.end(result); } } catch (error) { await releaseBestEffort(); if (!outer.done) outer.fail(error); } })(); return outer; } function createVertexAuthenticatedFetch(options: StreamOptions | undefined): FetchImpl { const baseFetch = options?.fetch ?? fetch; const vertexFetch = async (input: string | URL | Request, init?: RequestInit): Promise => { const token = await getVertexAccessToken({ signal: options?.signal, fetch: baseFetch }); const headers = new Headers(init?.headers); headers.set("Authorization", `Bearer ${token}`); const rewritten = resolveVertexRequest(input); const url = rewritten instanceof Request ? rewritten.url : rewritten.toString(); if (isVertexRawPredictUrl(url)) { const bodyText = await readVertexRequestBody(rewritten, init); const transformed = transformVertexAnthropicBody(bodyText); return baseFetch(url, { ...init, method: init?.method ?? (rewritten instanceof Request ? rewritten.method : "POST"), headers, body: transformed, }); } return baseFetch(rewritten, { ...init, headers }); }; return Object.assign(vertexFetch, baseFetch.preconnect ? { preconnect: baseFetch.preconnect } : {}); } async function readVertexRequestBody(input: string | URL | Request, init: RequestInit | undefined): Promise { if (input instanceof Request) return input.clone().text(); const body = init?.body; if (typeof body === "string") return body; if (body instanceof Uint8Array) return new TextDecoder().decode(body); if (body instanceof ArrayBuffer) return new TextDecoder().decode(body); return ""; } // Vertex Claude rejects the standard Anthropic body shape: the `model` field // is encoded in the URL path and `anthropic_version: "vertex-2023-10-16"` is // required in the JSON body instead of the `anthropic-version` HTTP header. function transformVertexAnthropicBody(bodyText: string): string { if (!bodyText) return bodyText; try { const payload = JSON.parse(bodyText) as Record; delete payload.model; payload.anthropic_version = "vertex-2023-10-16"; return JSON.stringify(payload); } catch { return bodyText; } } function resolveVertexRequest(input: string | URL | Request): string | URL | Request { const project = $env.GOOGLE_CLOUD_PROJECT || $env.GCP_PROJECT || $env.GCLOUD_PROJECT; const location = $env.GOOGLE_VERTEX_LOCATION || $env.GOOGLE_CLOUD_LOCATION || $env.VERTEX_LOCATION; if (!project || !location) return input; const rewriteUrl = (url: string): string => { const hasPlaceholder = url.includes("{project}") || url.includes("{location}") || url.includes("%7Bproject%7D") || url.includes("%7Blocation%7D"); const host = resolveVertexEndpointHost(location); const rewritten = hasPlaceholder ? url .replace("https://{location}-aiplatform.googleapis.com", `https://${host}`) .replace("https://%7Blocation%7D-aiplatform.googleapis.com", `https://${host}`) .replaceAll("{project}", encodeURIComponent(project)) .replaceAll("%7Bproject%7D", encodeURIComponent(project)) .replaceAll("{location}", encodeURIComponent(location)) .replaceAll("%7Blocation%7D", encodeURIComponent(location)) : url; return rewritten.replace(":streamRawPredict/v1/messages", ":streamRawPredict"); }; if (input instanceof Request) { const rewrittenUrl = rewriteUrl(input.url); return rewrittenUrl === input.url ? input : new Request(rewrittenUrl, input); } if (input instanceof URL) { const rewrittenUrl = rewriteUrl(input.toString()); return rewrittenUrl === input.toString() ? input : new URL(rewrittenUrl); } return rewriteUrl(input); } type KeyResolver = string | (() => string | undefined); const LEGACY_ENV_KEYS: Record = { // Non-provider / search-tool keys and API-name keys not modeled as registry provider defs. "azure-openai-responses": "AZURE_OPENAI_API_KEY", jina: "JINA_API_KEY", brave: "BRAVE_API_KEY", tinyfish: "TINYFISH_API_KEY", firecrawl: "FIRECRAWL_API_KEY", }; /** * Env fallbacks derived from the catalog provider entries (`env` in * `providers/.kdl`) — the single source for plain provider env-var names. * Registry defs override with computed resolvers (Foundry/ADC/Bedrock * probes); legacy non-provider keys merge last. */ const CATALOG_ENTRY_ENV_KEYS = Object.values(providerEntries()).flatMap(provider => { const envVars = provider.envVars; if (!envVars || envVars.length === 0) return []; const resolver: KeyResolver = envVars.length === 1 ? envVars[0] : () => $pickenv(...envVars); return [[provider.id, resolver] as [string, KeyResolver]]; }); const serviceProviderMap: Record = { ...Object.fromEntries(CATALOG_ENTRY_ENV_KEYS), ...Object.fromEntries( PROVIDER_REGISTRY.flatMap(provider => provider.envKeys != null ? [[provider.id, provider.envKeys] as [string, KeyResolver]] : [], ), ), ...LEGACY_ENV_KEYS, }; /** * Get API key for provider from known environment variables, e.g. OPENAI_API_KEY. * * Will not return API keys for providers that require OAuth tokens. * Checks Bun.env, then cwd/.env, then ~/.env. */ export function getEnvApiKey(provider: string): string | undefined { const resolver = serviceProviderMap[provider]; if (typeof resolver === "string") { return $env[resolver]; } return resolver?.(); } /** * Name of the environment variable that backs `getEnvApiKey` for a provider, * when that provider maps to a single named variable (e.g. `github-copilot` → * `COPILOT_GITHUB_TOKEN`). Returns undefined for providers whose env fallback * is computed (multi-var pickers, Vertex ADC / Bedrock probes, …) since no * single variable name describes the source. */ export function getEnvApiKeyName(provider: string): string | undefined { const resolver = serviceProviderMap[provider]; return typeof resolver === "string" ? resolver : undefined; } /** * Enumerate every provider that has an env-var fallback for `getEnvApiKey`. * Used by `omp auth-broker migrate --include-env` to discover env-sourced keys * that should be uploaded to the broker. */ export function listProvidersWithEnvKey(): string[] { return Object.keys(serviceProviderMap); } function withResolvedModelHeaders( model: Model, signal: AbortSignal | undefined, run: (resolvedModel: Model) => AssistantMessageEventStream, ): AssistantMessageEventStream { const resolveHeaders = model.resolveHeaders; if (!resolveHeaders) return run(model); const outer = new AssistantMessageEventStream(); void (async () => { try { const headers = await untilAborted(signal, () => resolveHeaders(signal)); signal?.throwIfAborted(); const inner = run({ ...model, resolveHeaders: undefined, headers: headers ? { ...headers } : undefined }); for await (const event of inner) { outer.push(event); if (outer.done) return; } if (!outer.done) outer.end(await inner.result()); } catch (error) { outer.fail(error); } })(); return outer; } export function stream( model: Model, context: Context, options?: OptionsForApi, ): AssistantMessageEventStream { if (model.resolveHeaders) { return withResolvedModelHeaders(model, options?.signal, resolvedModel => stream(resolvedModel, context, options)); } if (!model.requiresGlyphTokenization) { return withThinkingLoopGuard(model, options, opts => withProviderInFlightLimit(model, opts, () => streamDispatch(model, context, opts)), ); } const codec = applyGlyphCodec(context); const execHandlers = options?.execHandlers; const wireOptions: OptionsForApi | undefined = execHandlers === undefined ? options : { ...options, execHandlers: codec.wrapCursorExecHandlers(execHandlers) }; return codec.wrap( withThinkingLoopGuard(model, wireOptions, opts => withProviderInFlightLimit(model, opts, () => streamDispatch(model, codec.context, opts)), ), ); } function streamDispatch( model: Model, context: Context, options?: OptionsForApi, ): AssistantMessageEventStream { const requestOptions = withSupportedSamplingParams( model, withTransportFetch(model, (options || {}) as StreamOptions), ) as OptionsForApi; assertExplicitOpenAIResponsesPromptCacheSupport(model, requestOptions); // Check custom API registry first (extension-provided APIs like "vertex-claude-api") const customApiProvider = getCustomApi(model.api); if (customApiProvider) { return customApiProvider.stream(model, context, requestOptions as StreamOptions); } if (model.provider === "gitlab-duo") { const apiKey = requestOptions.apiKey || getEnvApiKey(model.provider); if (!apiKey) { throw new AIError.MissingApiKeyError(model.provider); } return streamGitLabDuo(model, context, { ...(requestOptions as SimpleStreamOptions), apiKey, }); } if (model.api === "gitlab-duo-agent") { const apiKey = (requestOptions as StreamOptions | undefined)?.apiKey || getEnvApiKey(model.provider); if (!apiKey) { throw new AIError.MissingApiKeyError(model.provider); } return streamGitLabDuoWorkflow(model as Model<"gitlab-duo-agent">, context, { ...(requestOptions as StreamOptions | undefined), apiKey, } as GitLabDuoWorkflowOptions); } // Vertex AI and Bedrock Converse authenticate outside the generic API-key path. if (model.api === "google-vertex") { return streamGoogleVertex(model as Model<"google-vertex">, context, requestOptions as GoogleVertexOptions); } if (model.api === "bedrock-converse-stream") { return streamBedrock(model as Model<"bedrock-converse-stream">, context, requestOptions as BedrockOptions); } if (model.api === "factory-droid-agent") { return streamFactoryDroid(model as Model<"factory-droid-agent">, context, requestOptions as FactoryDroidOptions); } const providerDefinition = getProviderDefinition(model.provider); const requestModel = providerDefinition?.prepareModel?.(model) ?? model; const prepared = providerDefinition?.prepareRequest?.(requestModel, requestOptions as StreamOptions); const providerModel = prepared?.model ?? requestModel; const preparedOptions = prepared?.options ?? (requestOptions as StreamOptions); const apiKey = preparedOptions.apiKey || getEnvApiKey(providerModel.provider); if (!apiKey) { throw new AIError.MissingApiKeyError(providerModel.provider); } const providerOptions = isGoogleVertexAuthenticatedModel(providerModel) ? { ...preparedOptions, apiKey: "vertex-adc", fetch: createVertexAuthenticatedFetch(preparedOptions), } : { ...preparedOptions, apiKey }; const api: Api = providerModel.api; switch (api) { case "anthropic-messages": { const anthropicOptions = providerOptions as AnthropicOptions; return streamAnthropic(providerModel as Model<"anthropic-messages">, context, { ...anthropicOptions, isOAuth: anthropicOptions.isOAuth ?? providerModel.isOAuth, }); } case "openrouter": { const useResponses = $env.PI_OPENROUTER_RESPONSES !== "0"; if (useResponses) { return streamOpenAIResponses( providerModel as Model<"openai-responses">, context, providerOptions as OptionsForApi<"openai-responses">, ); } return streamOpenAICompletions( providerModel as Model<"openai-completions">, context, providerOptions as OptionsForApi<"openai-completions">, ); } case "openai-completions": return streamOpenAICompletions( providerModel as Model<"openai-completions">, context, providerOptions as OptionsForApi<"openai-completions">, ); case "openai-responses": return streamOpenAIResponses( providerModel as Model<"openai-responses">, context, providerOptions as OptionsForApi<"openai-responses">, ); case "azure-openai-responses": return streamAzureOpenAIResponses( providerModel as Model<"azure-openai-responses">, context, providerOptions as OptionsForApi<"azure-openai-responses">, ); case "openai-codex-responses": return streamOpenAICodexResponses( providerModel as Model<"openai-codex-responses">, context, providerOptions as OptionsForApi<"openai-codex-responses">, ); case "google-generative-ai": return streamGoogle(providerModel as Model<"google-generative-ai">, context, providerOptions); case "google-gemini-cli": return streamGoogleGeminiCli( providerModel as Model<"google-gemini-cli">, context, providerOptions as GoogleGeminiCliOptions, ); case "ollama-chat": return streamOllama(providerModel as Model<"ollama-chat">, context, providerOptions as OllamaChatOptions); case "cursor-agent": return streamCursor(providerModel as Model<"cursor-agent">, context, providerOptions as CursorOptions); case "devin-agent": return streamDevin(providerModel as Model<"devin-agent">, context, providerOptions as DevinOptions); case "apple-foundation-models": return streamAppleFoundationModels( providerModel as Model<"apple-foundation-models">, context, providerOptions as AppleFoundationModelsOptions, ); default: throw new AIError.ConfigurationError(`Unhandled API: ${api}`); } } /** Maximum guarded attempts for a detected thinking loop. */ const THINKING_LOOP_MAX_ATTEMPTS = 3; const THINKING_LOOP_RETRY_BASE_DELAY_MS = 500; const THINKING_LOOP_RETRY_MAX_DELAY_MS = 8_000; function isRetryableThinkingLoop(message: AssistantMessage): boolean { return ( message.stopReason === "error" && message.content.length === 0 && AIError.is(message.errorId, AIError.Flag.ThinkingLoop) ); } /** * Resolve a completion, re-sampling a thinking-loop stall for at most * {@link THINKING_LOOP_MAX_ATTEMPTS} guarded attempts. The loop guard raises an * empty `stopReason: "error"` stall; after the budget is spent that error is * returned unchanged. Detection is never disabled as a fallback, because an * unguarded retry can consume the remaining output budget and persist runaway * content. Non-stall results, including genuine errors, return immediately. A * caller abort during backoff propagates so cancellation surfaces as an abort, * never a stale stall result. */ async function resolveWithThinkingLoopRetries( signal: AbortSignal | undefined, dispatch: () => AssistantMessageEventStream, onAttempt?: (message: AssistantMessage) => void, ): Promise { const dispatchAttempt = async (): Promise => { const response = dispatch(); for await (const _event of response) { // Completion callers do not consume deltas; drain them as they arrive to avoid retaining the response history. } const message = await response.result(); onAttempt?.(message); return message; }; let message = await dispatchAttempt(); let thinkingLoopRetry = isRetryableThinkingLoop(message); for (let attempt = 1; thinkingLoopRetry && attempt < THINKING_LOOP_MAX_ATTEMPTS; attempt += 1) { // A caller abort surfaces as a thrown abort (never the stall, which would // misclassify as a 502): throwIfAborted before backoff, and scheduler.wait // rejects if the abort lands mid-delay. signal?.throwIfAborted(); const delay = Math.min(THINKING_LOOP_RETRY_BASE_DELAY_MS * 2 ** (attempt - 1), THINKING_LOOP_RETRY_MAX_DELAY_MS); await scheduler.wait(delay, { signal }); message = await dispatchAttempt(); thinkingLoopRetry = isRetryableThinkingLoop(message); } if (thinkingLoopRetry) signal?.throwIfAborted(); return message; } export async function complete( model: Model, context: Context, options?: OptionsForApi, ): Promise { return resolveWithThinkingLoopRetries(options?.signal, () => stream(model, context, options)); } type AuthRetryFailure = { error: unknown; bufferedEvents: AssistantMessageEvent[]; terminalEvent?: Extract; }; function extractStatusFromAssistantError(message: AssistantMessage): number | undefined { if (message.errorStatus !== undefined) return message.errorStatus; if (!message.errorMessage) return undefined; return AIError.status({ message: message.errorMessage }); } function isRetryableUpstreamError( model: Model, error: unknown, status: number | undefined, message: string | undefined, ): boolean { if (AIError.isAuthRetryableError(error)) return true; // 401 means the credential is bad; 403 is its valid-token twin (access // denied by plan, model policy, or org restriction — a sibling account may // not share it). Explicit account-scoped policy errors such as Codex // `cyber_policy` are likewise rotatable. The exact ChatGPT-account model // denial is rotatable only when its provider and requested model match. // Usage-limit phrasing (Codex's // "You have hit your ChatGPT usage limit", Anthropic's "usage_limit_reached", // Google's "resource_exhausted", OpenAI's "insufficient_quota") and 429s // without transient rate-limit wording mean this account is parked but a // sibling credential can usually pick the request up. Both are rotatable // via `onAuthError` — the auth-gateway maps hard auth failures to // `invalidateCredentialMatching` and temporary account constraints to a // credential block. Transient 429s ("Too many requests", per-minute caps) // classify as RATE_LIMIT_EXCEEDED in `parseRateLimitReason` and stay in the // provider's own backoff layer instead of burning siblings. if (AIError.isCodexChatGPTAccountPolicyError(error, model.provider, model.id)) return true; if (status === 401 || (status === 403 && !isConcurrencyCapExclusion(status, message))) return true; return isUsageLimitOutcome(status, message); } function createAssistantAuthError(message: AssistantMessage): Error { const text = message.errorMessage ?? "Provider authentication failed"; const status = extractStatusFromAssistantError(message); const error = status === undefined ? new AIError.ProviderResponseError(text, { kind: "runtime" }) : new ProviderHttpError(text, status); return typeof message.errorId === "number" ? AIError.attach(error, message.errorId) : error; } function contextualizeAuthRetryError(model: Model, error: unknown): unknown { if ( !error || typeof error !== "object" || !AIError.isCodexChatGPTAccountPolicyError(error, model.provider, model.id) ) { return error; } return AIError.attach(error, AIError.create(AIError.Flag.AccountPolicy | AIError.Flag.ContentBlocked)); } function emitBufferedEvents(stream: AssistantMessageEventStream, events: AssistantMessageEvent[]): void { for (const event of events) { stream.push(event); } } function withInferenceSessionId(options?: SimpleStreamOptions): SimpleStreamOptions { if (options?.sessionId) return options; return { ...options, sessionId: crypto.randomUUID() }; } type SamplingOptions = Pick< StreamOptions, "temperature" | "topP" | "topK" | "minP" | "presencePenalty" | "repetitionPenalty" | "frequencyPenalty" >; /** * Drop explicit sampling parameters before any provider builds its payload * when the model's resolved `compat.supportsSamplingParams` is `false`. The * catalog's class rules assign that per model lineage on every compat record, * and an explicit compat override still wins, so this one check covers every * provider. */ function withSupportedSamplingParams(model: Model, options: T): T { if ( options.temperature === undefined && options.topP === undefined && options.topK === undefined && options.minP === undefined && options.presencePenalty === undefined && options.repetitionPenalty === undefined && options.frequencyPenalty === undefined ) { return options; } const compat = model.compat; if (!compat || !("supportsSamplingParams" in compat) || compat.supportsSamplingParams !== false) return options; const supported = { ...options }; delete supported.temperature; delete supported.topP; delete supported.topK; delete supported.minP; delete supported.presencePenalty; delete supported.repetitionPenalty; delete supported.frequencyPenalty; return supported; } export function streamSimple( model: Model, context: Context, options?: SimpleStreamOptions, ): AssistantMessageEventStream { const sessionOptions = withInferenceSessionId(options); if (!model.requiresGlyphTokenization) { return streamSimpleRequest(model, context, sessionOptions); } const codec = applyGlyphCodec(context); const execHandlers = sessionOptions.cursorExecHandlers ?? sessionOptions.execHandlers; const wrappedExecHandlers = execHandlers === undefined ? undefined : codec.wrapCursorExecHandlers(execHandlers); const wireOptions = wrappedExecHandlers === undefined ? sessionOptions : { ...sessionOptions, execHandlers: wrappedExecHandlers, cursorExecHandlers: wrappedExecHandlers, }; return codec.wrap(streamSimpleRequest(model, codec.context, wireOptions)); } /** * Forward a model-configured `User-Agent` override across the pi-native wire. * The model itself never crosses the wire — the client sends only `modelId` * and the gateway resolves its own model — so without this the gateway's * resolved Bedrock model always sends the default `omp/` UA even * when the client's local model config set an override. Only the single * header is forwarded, not the rest of `model.headers` (which may carry * unrelated local config), and only when the caller hasn't already set their * own `User-Agent` — a per-call header still wins, matching `streamBedrock`'s * own caller-headers precedence. */ function forwardBedrockUserAgent( modelHeaders: Record | undefined, callerHeaders: Record | undefined, ): Record | undefined { if (getHeaderCaseInsensitive(callerHeaders, "user-agent") !== undefined) return callerHeaders; const modelUserAgent = getHeaderCaseInsensitive(modelHeaders, "user-agent"); return modelUserAgent === undefined ? callerHeaders : { ...callerHeaders, "User-Agent": modelUserAgent }; } function streamSimpleRequest( model: Model, context: Context, options?: SimpleStreamOptions, ): AssistantMessageEventStream { const requestOptions = withSupportedSamplingParams( model, withTransportFetch(model, (options || {}) as SimpleStreamOptions), ); const apiKeyResolver = isApiKeyResolver(requestOptions?.apiKey) ? requestOptions.apiKey : undefined; if (apiKeyResolver) { const outer = new AssistantMessageEventStream(); const signal = requestOptions?.signal; // One inner attempt against a resolved key, or against the Bedrock AWS // credential chain when its optional resolver has no stored bearer key. // Retryable auth failures are buffered until replay is safe. const runAttempt = async ( apiKey?: string, credentialId?: number, oauthIdentity?: OAuthRequestIdentity, ): Promise => { const bufferedEvents: AssistantMessageEvent[] = []; let emittedReplayUnsafeEvent = false; const flushBuffered = (): void => { emitBufferedEvents(outer, bufferedEvents); bufferedEvents.length = 0; }; try { const attemptOptions = { ...requestOptions, apiKey, credentialId, oauthIdentity }; const inner = streamSimpleRequest(model, context, attemptOptions); for await (const event of inner) { if (credentialId !== undefined) { if ("partial" in event) event.partial.credentialId = credentialId; else if (event.type === "done") event.message.credentialId = credentialId; else event.error.credentialId = credentialId; } if (!emittedReplayUnsafeEvent && event.type === "start") { bufferedEvents.push(event); continue; } if ( !emittedReplayUnsafeEvent && event.type === "error" && isRetryableUpstreamError( model, event.error, extractStatusFromAssistantError(event.error), event.error.errorMessage, ) ) { return { error: contextualizeAuthRetryError(model, createAssistantAuthError(event.error)), bufferedEvents, terminalEvent: event, }; } flushBuffered(); emittedReplayUnsafeEvent = true; outer.push(event); if (outer.done) return undefined; } flushBuffered(); if (!outer.done) { const result = await inner.result(); if (credentialId !== undefined) result.credentialId = credentialId; outer.end(result); } } catch (error) { if ( !emittedReplayUnsafeEvent && isRetryableUpstreamError( model, error, AIError.status(error), error instanceof Error ? error.message : undefined, ) ) { return { error: contextualizeAuthRetryError(model, error), bufferedEvents }; } flushBuffered(); outer.fail(error); } return undefined; }; const emitFailure = (failure: AuthRetryFailure): void => { emitBufferedEvents(outer, failure.bufferedEvents); if (failure.terminalEvent) { outer.push(failure.terminalEvent); } else { outer.fail(failure.error); } }; void (async () => { let lastKey: string | undefined; let credentialId: number | undefined; let oauthIdentity: OAuthRequestIdentity | undefined; try { const resolved = await apiKeyResolver({ lastChance: false, error: undefined, signal }); lastKey = resolvedApiKeyBearer(resolved); credentialId = typeof resolved === "string" ? undefined : resolved?.credentialId; oauthIdentity = typeof resolved === "string" ? undefined : resolved?.oauthIdentity; } catch (error) { // A thrown resolver is a broker/OAuth/network failure, not a missing // key — surface the cause instead of masking it as "No API key". outer.fail( new AIError.ConfigurationError( `Failed to resolve API key for provider ${model.provider}: ${error instanceof Error ? error.message : String(error)}`, { cause: error }, ), ); return; } if (lastKey === undefined) { if (getProviderDefinition(model.provider)?.allowsMissingApiKey) { const failure = await runAttempt(); if (failure) emitFailure(failure); return; } outer.fail(new AIError.MissingApiKeyError(model.provider)); return; } const retryState = createAuthRetryKeyState(lastKey); let failure = await runAttempt(lastKey, credentialId, oauthIdentity); if (!failure) return; while (true) { // Caller aborted between attempts: don't mint a fresh token or fire // another doomed request — emit the captured failure instead. if (signal?.aborted) break; let nextCredentialId: number | undefined; let nextOAuthIdentity: OAuthRequestIdentity | undefined; const nextKey = await resolveNextAuthRetryKey( retryState, apiKeyResolver, failure.error, signal, resolved => { nextCredentialId = typeof resolved === "string" ? undefined : resolved?.credentialId; nextOAuthIdentity = typeof resolved === "string" ? undefined : resolved?.oauthIdentity; }, ); if (nextKey === undefined) break; const next = await runAttempt(nextKey, nextCredentialId, nextOAuthIdentity); if (!next) return; failure = next; } emitFailure(failure); })(); return outer; } if (model.resolveHeaders) { return withResolvedModelHeaders(model, requestOptions.signal, resolvedModel => streamSimpleRequest(resolvedModel, context, requestOptions), ); } // Pi-native transport short-circuits the per-provider dispatch entirely: // the gateway resolves provider + credential server-side, so we don't // need an `apiKey` from `getEnvApiKey` here — `options.apiKey` carries // the gateway bearer instead. Comes BEFORE the custom-API check so // extension-registered APIs can't accidentally override a configured // pi-native transport. if (model.transport === "pi-native") { return withThinkingLoopGuard(model, requestOptions, opts => withProviderInFlightLimit(model, opts, () => { const nativeOptions = model.api === "bedrock-converse-stream" ? { ...opts, guardrailIdentifier: model.guardrailIdentifier ?? opts?.guardrailIdentifier, guardrailVersion: model.guardrailVersion ?? opts?.guardrailVersion, guardrailTrace: model.guardrailTrace ?? opts?.guardrailTrace, // The model itself never crosses the wire — the client sends only // `modelId` and the gateway resolves its own model — so the model's // tags must be flattened in here or they are lost entirely. Per-call // entries win per key; the merged map then wins per key over the // gateway-resolved model's own tags in its `streamBedrock`. requestMetadata: model.requestMetadata || opts?.requestMetadata ? { ...model.requestMetadata, ...opts?.requestMetadata } : undefined, headers: forwardBedrockUserAgent(model.headers, opts?.headers), } : opts; return streamPiNative(model, context, nativeOptions); }), ); } // Check custom API registry (extension-provided APIs) const customApiProvider = getCustomApi(model.api); if (customApiProvider) { return withThinkingLoopGuard(model, requestOptions, opts => withProviderInFlightLimit(model, opts, () => customApiProvider.streamSimple(model, context, opts)), ); } // Vertex AI uses Application Default Credentials, not API keys if (model.api === "google-vertex") { const providerOptions = mapOptionsForApi(model, requestOptions, undefined); return stream(model, context, providerOptions); } else if (model.api === "bedrock-converse-stream") { // Bedrock doesn't have any API keys instead it sources credentials from standard AWS env variables or from given AWS profile. const providerOptions = mapOptionsForApi(model, requestOptions, undefined); return stream(model, context, providerOptions); } else if (getProviderDefinition(model.provider)?.allowsMissingApiKey) { const providerOptions = mapOptionsForApi( model, requestOptions, typeof requestOptions.apiKey === "string" ? requestOptions.apiKey : getEnvApiKey(model.provider), ); return stream(model, context, providerOptions); } // The resolver form is handled by the wrapper above; only a static string // key reaches this point. const apiKey = (typeof requestOptions?.apiKey === "string" ? requestOptions.apiKey : undefined) || getEnvApiKey(model.provider); if (!apiKey) { throw new AIError.MissingApiKeyError(model.provider); } // GitLab Duo - wraps Anthropic/OpenAI behind GitLab AI Gateway direct access tokens if (model.provider === "gitlab-duo") { return withThinkingLoopGuard(model, requestOptions, opts => withProviderInFlightLimit(model, opts, () => streamGitLabDuo(model, context, { ...opts, apiKey, }), ), ); } // GitLab Duo Workflow - IDE workflow protocol + WebSocket action bridge if (model.api === "gitlab-duo-agent") { // Does not route through withProviderInFlightLimit, so heal explicitly. return withThinkingLoopGuard(model, requestOptions, opts => healLeakedThinking( model, streamGitLabDuoWorkflow(model as Model<"gitlab-duo-agent">, context, { ...opts, apiKey, }), ), ); } // Kimi Code - route to dedicated handler that wraps OpenAI or Anthropic API if (model.provider === "kimi-code") { // streamKimi handles openai/anthropic format mapping internally, but the // mandatory-reasoning clamp is a request-shaping concern owned here: K3's // `supports_thinking_type: "only"` endpoint rejects disabled/omitted // thinking, so clamp disabled requests to the lowest supported effort // (mirrors the mapOptionsForApi path every other provider takes). const kimiOptions = normalizeMandatoryReasoningOptions(model, requestOptions); return withThinkingLoopGuard(model, kimiOptions, opts => withProviderInFlightLimit(model, opts, () => streamKimi(model as Model<"openai-completions">, context, { ...opts, apiKey, format: opts?.kimiApiFormat, }), ), ); } // Synthetic - route to dedicated handler that wraps OpenAI or Anthropic API if (model.provider === "synthetic") { // Pass raw SimpleStreamOptions - streamSynthetic handles mapping internally. return withThinkingLoopGuard(model, requestOptions, opts => withProviderInFlightLimit(model, opts, () => streamSynthetic(model as Model<"openai-completions">, context, { ...opts, apiKey, format: opts?.syntheticApiFormat ?? "openai", }), ), ); } const providerModel = getProviderDefinition(model.provider)?.prepareModel?.(model) ?? model; const providerOptions = mapOptionsForApi(providerModel, requestOptions, apiKey); return stream(providerModel, context, providerOptions); } export async function completeSimple( model: Model, context: Context, options?: SimpleStreamOptions & { /** Receives every completed result, including results retried by the thinking-loop guard. */ onAttempt?: (message: AssistantMessage) => void; }, ): Promise { const { onAttempt, ...streamOptions } = options ?? {}; const sessionOptions = withInferenceSessionId(streamOptions); return resolveWithThinkingLoopRetries( options?.signal, () => streamSimple(model, context, sessionOptions), onAttempt, ); } const MIN_OUTPUT_TOKENS = 1024; // Fallback total output cap for models whose catalog entry has no maxTokens. const OUTPUT_CAP_WHEN_UNKNOWN = 64_000; function maxTokensWithThinkingBudget( baseMaxTokens: number | undefined, modelMaxTokens: number | null, thinkingBudget: number, ): number { const uncappedMaxTokens = baseMaxTokens === undefined ? OUTPUT_CAP_WHEN_UNKNOWN : baseMaxTokens + thinkingBudget; return Math.min(uncappedMaxTokens, modelMaxTokens ?? Number.POSITIVE_INFINITY); } export const OUTPUT_FALLBACK_BUFFER = 4000; const ANTHROPIC_USE_INTERLEAVED_THINKING = Bun.env.PI_NO_INTERLEAVED_THINKING !== "1"; export const ANTHROPIC_THINKING: Record = { minimal: 1024, low: 4096, medium: 8192, high: 16384, xhigh: 32768, max: 32768, }; const GOOGLE_THINKING: Record = { minimal: 1024, low: 4096, medium: 8192, high: 16384, xhigh: 24575, max: 32768, }; const BEDROCK_CLAUDE_THINKING: Record = { minimal: 1024, low: 2048, medium: 8192, high: 16384, xhigh: 16384, max: 32768, }; function resolveBedrockThinkingBudget( model: Model<"bedrock-converse-stream">, options?: SimpleStreamOptions, ): { budget: number; level: Effort } | null { if (!options?.reasoning || !model.reasoning || options.disableReasoning || options.forceReasoningOff) return null; const level = requireSupportedEffort(model, options.reasoning); const budget = options.thinkingBudgets?.[level] ?? BEDROCK_CLAUDE_THINKING[level]; return { budget, level }; } export function mapAnthropicToolChoice(choice?: ToolChoice): AnthropicOptions["toolChoice"] { if (!choice) return undefined; if (typeof choice === "string") { if (choice === "required") return "any"; if (choice === "auto" || choice === "none" || choice === "any") return choice; return undefined; } if (choice.type === "tool") { return choice.name ? { type: "tool", name: choice.name } : undefined; } if (choice.type === "function") { const name = "function" in choice ? choice.function?.name : choice.name; return name ? { type: "tool", name } : undefined; } return undefined; } export function mapGoogleToolChoice( choice?: ToolChoice, ): GoogleOptions["toolChoice"] | GoogleGeminiCliOptions["toolChoice"] | GoogleVertexOptions["toolChoice"] { if (!choice) return undefined; if (typeof choice === "string") { if (choice === "required") return "any"; if (choice === "auto" || choice === "none" || choice === "any") return choice; return undefined; } // Named-tool routing on Google: emit an `ANY`-mode allow-list of one entry, // mirroring the Anthropic mapper that returns `{type: "tool", name}`. if (choice.type === "tool") { return choice.name ? { mode: "ANY", allowedFunctionNames: [choice.name] } : undefined; } if (choice.type === "function") { const name = "function" in choice ? choice.function?.name : choice.name; return name ? { mode: "ANY", allowedFunctionNames: [name] } : undefined; } return undefined; } function mapOpenAiToolChoice(choice?: ToolChoice): OpenAICompletionsOptions["toolChoice"] { if (!choice) return undefined; if (typeof choice === "string") { if (choice === "any") return "required"; if (choice === "auto" || choice === "none" || choice === "required") return choice; return undefined; } if (choice.type === "tool") { return choice.name ? { type: "function", function: { name: choice.name } } : undefined; } if (choice.type === "function") { const name = "function" in choice ? choice.function?.name : choice.name; return name ? { type: "function", function: { name } } : undefined; } return undefined; } type ReasoningEffortMapCompat = { reasoningEffortMap?: Partial>; }; function getCompatReasoningEffortMap( model: Model, ): Partial> | undefined { const compat = model.compat; if (compat === undefined || typeof compat !== "object" || !("reasoningEffortMap" in compat)) { return undefined; } return (compat as ReasoningEffortMapCompat).reasoningEffortMap; } function resolveSupportedMappedReasoningEffort( model: Model, reasoning: Effort, ): Effort | undefined { const mapped = getCompatReasoningEffortMap(model)?.[reasoning]; if (!mapped) return undefined; const mappedEffort = mapped as Effort; return model.thinking?.efforts.includes(mappedEffort) ? mappedEffort : undefined; } function resolveOpenAiReasoningEffort( model: Model, options?: SimpleStreamOptions, ): Effort | undefined { const reasoning = options?.reasoning; if (!reasoning || !model.reasoning) return undefined; // Models that reason natively but expose no effort dial carry // `thinking: undefined` (baked at build time from // `compat.supportsReasoningEffort: false` on openai-responses*). The // wire-side omitReasoningEffort gate (stream.ts) is the actual strip; returning // undefined here avoids a redundant requireSupportedEffort throw that would // defeat the gate and surface a confusing "Compaction failed: Thinking effort // high is not supported by..." to the user. if (!model.thinking) return undefined; if (model.thinking.efforts.includes(reasoning)) return reasoning; const mappedReasoning = resolveSupportedMappedReasoningEffort(model, reasoning); if (mappedReasoning) return mappedReasoning; if (getCompatReasoningEffortMap(model)?.[reasoning] !== undefined) return reasoning; if (model.thinking.effortMap?.[reasoning] !== undefined) return reasoning; return requireSupportedEffort(model, reasoning); } function resolveGoogleThinkingOff(model: Model): NonNullable { const thinking: NonNullable = { enabled: false }; if (!model.reasoning || !model.thinking) return thinking; if (model.thinking.mode === "budget" && (!model.thinking.requiresEffort || model.thinking.suppressWhenOff)) { thinking.budgetTokens = 0; } else if (model.thinking.mode === "google-level" && model.thinking.suppressWhenOff) { thinking.level = "MINIMAL"; } return thinking; } const castApi = (api: OptionsForApi): OptionsForApi => api as OptionsForApi; /** * Mandatory-reasoning endpoints (`thinking.requiresEffort`) reject disabled * or omitted thinking ("Reasoning is mandatory for this endpoint and cannot * be disabled") — clamp to the lowest supported effort instead. * `suppressWhenOff` models handle off provider-side via explicit wire * suppression. Collapsed pairs interplay: pair derivation strips member * flags (off routes to a bare SKU that CAN disable), while identity backfill * re-flags pairs whose logical id is itself mandatory (Gemini 3.x) — there * the clamp wins and the floored effort routes to the thinking SKU. */ function normalizeMandatoryReasoningOptions( model: Model, options?: SimpleStreamOptions, ): SimpleStreamOptions | undefined { if ( !model.reasoning || !model.thinking?.requiresEffort || model.thinking.suppressWhenOff || (options?.reasoning !== undefined && !options.disableReasoning && !options.forceReasoningOff) ) { return options; } const floor = defaultSupportedEffort(model); if (floor === undefined) return options; return { ...options, reasoning: floor, disableReasoning: undefined, forceReasoningOff: undefined }; } function supportsExplicitOpenAIResponsesPromptCache(compat: unknown): boolean { return ( typeof compat === "object" && compat !== null && "supportsPromptCacheBreakpoints" in compat && compat.supportsPromptCacheBreakpoints === true ); } function isOpenAIResponsesPromptCacheSurface(model: Model): boolean { return ( model.api === "openai-responses" || model.api === "azure-openai-responses" || (model.api === "openrouter" && $env.PI_OPENROUTER_RESPONSES !== "0") ); } function assertExplicitOpenAIResponsesPromptCacheSupport( model: Model, options?: StreamOptions, ): void { if ( model.transport === "pi-native" || resolveCacheRetention(options?.cacheRetention) === "none" || options?.promptCache?.mode !== "explicit" || !isOpenAIResponsesPromptCacheSurface(model) || supportsExplicitOpenAIResponsesPromptCache(model.compat) ) { return; } throw new AIError.ConfigurationError( `OpenAI explicit prompt caching is unsupported for ${model.provider}/${model.id}; enable compat.supportsPromptCacheBreakpoints only for a compatible endpoint.`, ); } function mapOptionsForApi( model: Model, rawOptions?: SimpleStreamOptions, apiKey?: string, ): OptionsForApi { const options = normalizeMandatoryReasoningOptions(model, rawOptions); const simpleProviderOptions = getProviderDefinition(model.provider)?.mapSimpleOptions?.(options ?? {}); const base = { temperature: options?.temperature, topP: options?.topP, topK: options?.topK, minP: options?.minP, presencePenalty: options?.presencePenalty, repetitionPenalty: options?.repetitionPenalty, maxTokens: options?.maxTokens ?? model.maxTokens ?? undefined, signal: options?.signal, apiKey: apiKey ?? (typeof options?.apiKey === "string" ? options.apiKey : undefined), credentialId: options?.credentialId, oauthIdentity: options?.oauthIdentity, cacheRetention: options?.cacheRetention, headers: options?.headers, initiatorOverride: options?.initiatorOverride, maxRetryDelayMs: options?.maxRetryDelayMs, metadata: options?.metadata, taskBudget: options?.taskBudget, sessionId: options?.sessionId, promptCacheKey: options?.promptCacheKey, streamFirstEventTimeoutMs: options?.streamFirstEventTimeoutMs, streamIdleTimeoutMs: options?.streamIdleTimeoutMs, codexSseMaxAttempts: options?.codexSseMaxAttempts, providerSessionState: options?.providerSessionState, liveSteering: options?.liveSteering, maxInFlightRequests: options?.maxInFlightRequests, toolNamespacesInfo: options?.toolNamespacesInfo, onPayload: options?.onPayload, onResponse: options?.onResponse, onSseEvent: options?.onSseEvent, execHandlers: options?.execHandlers, fetch: options?.fetch, fallbacks: options?.fallbacks, acceptEmptyResponse: options?.acceptEmptyResponse, anthropicPrefixMismatchBehavior: options?.anthropicPrefixMismatchBehavior, anthropicCompaction: options?.anthropicCompaction, anthropicSlowMode: options?.anthropicSlowMode, userProfileId: options?.userProfileId, ...simpleProviderOptions, }; switch (model.api) { case "anthropic-messages": { // Explicitly disable thinking when reasoning is not specified, the caller // disabled it, an external scratchpad replaces it, or the model doesn't // support it. These SimpleStreamOptions flags never reach AnthropicOptions // on their own, so fold them into thinkingEnabled here (mandatory-reasoning // models already clamp them away in normalizeMandatoryReasoningOptions). const reasoning = options?.reasoning; if (!reasoning || !model.reasoning || options?.disableReasoning || options?.forceReasoningOff) { return castApi<"anthropic-messages">({ ...base, requestModelId: resolveWireModelId(model, undefined), thinkingEnabled: false, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, serviceTier: options?.serviceTier, }); } let thinkingBudget = options.thinkingBudgets?.[reasoning] ?? ANTHROPIC_THINKING[reasoning]; if (thinkingBudget <= 0) { return castApi<"anthropic-messages">({ ...base, requestModelId: resolveWireModelId(model, undefined), thinkingEnabled: false, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, serviceTier: options?.serviceTier, }); } const thinkingMode = model.thinking?.mode; const effort = thinkingMode === "anthropic-adaptive" || thinkingMode === "anthropic-budget-effort" ? mapEffortToAnthropicAdaptiveEffort(model, reasoning) : undefined; // A caller's maxTokens is the output it asked for, but thinking spends the // same max_tokens: adaptive thinking can use all of it and leave no answer. // Give thinking its budget on top, as the budget-only path below does. An // uncapped request keeps the provider default. const maxTokensWithThinking = base.maxTokens === undefined ? undefined : maxTokensWithThinkingBudget(base.maxTokens, model.maxTokens, thinkingBudget); // For Opus 4.6+ and Sonnet 4.6+: use adaptive thinking with effort level // For older models: use budget-based thinking if (thinkingMode === "anthropic-adaptive") { return castApi<"anthropic-messages">({ ...base, maxTokens: maxTokensWithThinking, requestModelId: resolveWireModelId(model, reasoning), thinkingEnabled: true, effort, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, serviceTier: options?.serviceTier, }); } if (ANTHROPIC_USE_INTERLEAVED_THINKING) { if ( model.maxTokens !== null && model.maxTokens !== undefined && model.maxTokens < thinkingBudget + OUTPUT_FALLBACK_BUFFER ) { thinkingBudget = model.maxTokens - OUTPUT_FALLBACK_BUFFER; } if (thinkingBudget >= ANTHROPIC_THINKING.minimal) { return castApi<"anthropic-messages">({ ...base, maxTokens: maxTokensWithThinking, requestModelId: resolveWireModelId(model, reasoning), thinkingEnabled: true, thinkingBudgetTokens: thinkingBudget, effort, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, serviceTier: options?.serviceTier, }); } } // Caller's maxTokens is desired output, so add thinking budget on top. With no caller/model cap, use a finite total fallback. const maxTokens = maxTokensWithThinkingBudget(base.maxTokens, model.maxTokens, thinkingBudget); // Keep the provider's output buffer after thinking, reducing the // budget before its wire-level clamp could fall below the API minimum. if (maxTokens < thinkingBudget + OUTPUT_FALLBACK_BUFFER) { thinkingBudget = maxTokens - OUTPUT_FALLBACK_BUFFER; } // If thinking budget is too low, disable thinking if (thinkingBudget < ANTHROPIC_THINKING.minimal) { return castApi<"anthropic-messages">({ ...base, requestModelId: resolveWireModelId(model, undefined), thinkingEnabled: false, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, serviceTier: options?.serviceTier, }); } else { return castApi<"anthropic-messages">({ ...base, maxTokens, requestModelId: resolveWireModelId(model, reasoning), thinkingEnabled: true, thinkingBudgetTokens: thinkingBudget, effort, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, serviceTier: options?.serviceTier, }); } } case "bedrock-converse-stream": { const bedrockBase: BedrockOptions = { ...base, // Explicit reasoning-off must fold here like the anthropic-messages // branch: the provider gates thinking only on `reasoning`, and the // budget path below must not inflate a capped request for thinking // that was turned off. reasoning: options?.disableReasoning || options?.forceReasoningOff ? undefined : options?.reasoning, thinkingBudgets: options?.thinkingBudgets, toolChoice: mapAnthropicToolChoice(options?.toolChoice), thinkingDisplay: options?.hideThinkingSummary ? "omitted" : undefined, guardrailIdentifier: model.guardrailIdentifier ?? options?.guardrailIdentifier, guardrailVersion: model.guardrailVersion ?? options?.guardrailVersion, guardrailTrace: model.guardrailTrace ?? options?.guardrailTrace, requestMetadata: options?.requestMetadata, }; // Adaptive Claude shares max_tokens between thinking and the answer, like // the anthropic-messages adaptive path: a caller's cap is the output it // wants, so add the effort's budget on top. Uncapped requests keep the // provider default. if (model.thinking?.mode === "anthropic-adaptive") { const reasoning = bedrockBase.reasoning; const budget = reasoning ? (options?.thinkingBudgets?.[reasoning] ?? BEDROCK_CLAUDE_THINKING[reasoning]) : 0; if (!model.reasoning || bedrockBase.maxTokens === undefined || budget <= 0) { return castApi<"bedrock-converse-stream">(bedrockBase); } return castApi<"bedrock-converse-stream">({ ...bedrockBase, maxTokens: maxTokensWithThinkingBudget(bedrockBase.maxTokens, model.maxTokens, budget), }); } // Effort mode sends effort directly, no budget_tokens — skip budget inflation. if (model.thinking?.mode === "effort") { return castApi<"bedrock-converse-stream">(bedrockBase); } const budgetInfo = resolveBedrockThinkingBudget(model as Model<"bedrock-converse-stream">, options); if (!budgetInfo) return bedrockBase as OptionsForApi; let maxTokens = bedrockBase.maxTokens ?? model.maxTokens ?? OUTPUT_CAP_WHEN_UNKNOWN; let thinkingBudgets = bedrockBase.thinkingBudgets; if (maxTokens <= budgetInfo.budget) { const desiredMaxTokens = Math.min( model.maxTokens ?? Number.POSITIVE_INFINITY, budgetInfo.budget + MIN_OUTPUT_TOKENS, ); if (desiredMaxTokens > maxTokens) { maxTokens = desiredMaxTokens; } } if (maxTokens <= budgetInfo.budget) { const adjustedBudget = Math.max(0, maxTokens - MIN_OUTPUT_TOKENS); thinkingBudgets = { ...thinkingBudgets, [budgetInfo.level]: adjustedBudget }; } return castApi<"bedrock-converse-stream">({ ...bedrockBase, maxTokens, thinkingBudgets }); } case "openrouter": { const useResponses = $env.PI_OPENROUTER_RESPONSES !== "0"; if (useResponses) { return castApi<"openai-responses">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, reasoningSummary: options?.hideThinkingSummary ? null : undefined, openrouterVariant: options?.openrouterVariant, maxTokensExplicit: rawOptions?.maxTokens !== undefined, disableReasoning: options?.disableReasoning, // Forwarded, not folded: the Responses record reads both flags // itself (`applyResponsesCompatPolicy`). forceReasoningOff: options?.forceReasoningOff, textVerbosity: options?.textVerbosity, promptCache: options?.promptCache, statefulResponses: options?.statefulResponses, }); } return castApi<"openai-completions">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), // `OpenAICompletionsOptions` carries no forceReasoningOff; fold it. disableReasoning: options?.disableReasoning || options?.forceReasoningOff, toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, openrouterVariant: options?.openrouterVariant, maxTokensExplicit: rawOptions?.maxTokens !== undefined, promptCache: options?.promptCache, }); } case "openai-completions": return castApi<"openai-completions">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), // `OpenAICompletionsOptions` carries no forceReasoningOff; fold it. disableReasoning: options?.disableReasoning || options?.forceReasoningOff, toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, openrouterVariant: options?.openrouterVariant, maxTokensExplicit: rawOptions?.maxTokens !== undefined, promptCache: options?.promptCache, }); case "openai-responses": return castApi<"openai-responses">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, reasoningSummary: options?.hideThinkingSummary ? null : undefined, openrouterVariant: options?.openrouterVariant, maxTokensExplicit: rawOptions?.maxTokens !== undefined, disableReasoning: options?.disableReasoning, forceReasoningOff: options?.forceReasoningOff, textVerbosity: options?.textVerbosity, promptCache: options?.promptCache, statefulResponses: options?.statefulResponses, }); case "azure-openai-responses": return castApi<"azure-openai-responses">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, reasoningSummary: options?.hideThinkingSummary ? null : undefined, promptCache: options?.promptCache, statefulResponses: options?.statefulResponses, disableReasoning: options?.disableReasoning || options?.forceReasoningOff, forceReasoningOff: options?.forceReasoningOff, }); case "openai-codex-responses": return castApi<"openai-codex-responses">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, preferWebsockets: options?.preferWebsockets, codexCompaction: options?.codexCompaction, reasoningSummary: options?.hideThinkingSummary ? null : undefined, textVerbosity: options?.textVerbosity, forceReasoningOff: options?.forceReasoningOff, }); case "google-generative-ai": { // Explicitly disable thinking when reasoning is absent, unsupported, or // replaced by the caller's external scratchpad. Gemini defaults thinking on. const reasoning = options?.reasoning; if (!reasoning || !model.reasoning || options?.disableReasoning || options?.forceReasoningOff) { return castApi<"google-generative-ai">({ ...base, serviceTier: options?.serviceTier, thinking: resolveGoogleThinkingOff(model), toolChoice: mapGoogleToolChoice(options?.toolChoice), cachedContent: options?.cachedContent, }); } const googleModel = model as Model<"google-generative-ai">; const effort = requireSupportedEffort(googleModel, reasoning); // Gemini 3+ models use thinkingLevel exclusively instead of thinkingBudget. // https://ai.google.dev/gemini-api/docs/thinking#set-budget if (googleModel.thinking?.mode === "google-level") { return castApi<"google-generative-ai">({ ...base, serviceTier: options?.serviceTier, thinking: { enabled: true, level: mapEffortToGoogleThinkingLevel(effort, googleModel), }, hideThinkingSummary: options?.hideThinkingSummary, toolChoice: mapGoogleToolChoice(options?.toolChoice), cachedContent: options?.cachedContent, }); } return castApi<"google-generative-ai">({ ...base, thinking: { enabled: true, budgetTokens: getGoogleBudget(googleModel, effort, options?.thinkingBudgets), }, hideThinkingSummary: options?.hideThinkingSummary, toolChoice: mapGoogleToolChoice(options?.toolChoice), cachedContent: options?.cachedContent, }); } case "google-gemini-cli": { const reasoning = options?.reasoning; const toolChoice = mapGoogleToolChoice(options?.toolChoice); if (reasoning && model.reasoning && !options?.disableReasoning && !options?.forceReasoningOff) { const effort = requireSupportedEffort(model, reasoning); // Gemini 3+ models use thinkingLevel instead of thinkingBudget if (model.thinking?.mode === "google-level") { return castApi<"google-gemini-cli">({ ...base, requestModelId: resolveWireModelId(model, effort), thinking: { enabled: true, level: mapEffortToGoogleThinkingLevel(effort, model), }, hideThinkingSummary: options?.hideThinkingSummary, toolChoice, antigravityEndpointMode: options?.antigravityEndpointMode, }); } let thinkingBudget = options.thinkingBudgets?.[effort] ?? model.thinking?.effortBudgets?.[effort] ?? GOOGLE_THINKING[effort]; // Caller's maxTokens is desired output, so add thinking budget on top. With no caller/model cap, use a finite total fallback. const maxTokens = maxTokensWithThinkingBudget(base.maxTokens, model.maxTokens, thinkingBudget); // If not enough room for thinking + output, reduce thinking budget if (maxTokens <= thinkingBudget) { thinkingBudget = Math.max(0, maxTokens - MIN_OUTPUT_TOKENS); } if (thinkingBudget > 0) { return castApi<"google-gemini-cli">({ ...base, maxTokens, requestModelId: resolveWireModelId(model, effort), thinking: { enabled: true, budgetTokens: thinkingBudget }, hideThinkingSummary: options?.hideThinkingSummary, toolChoice, antigravityEndpointMode: options?.antigravityEndpointMode, }); } // Budget clamped to zero — fall through to the thinking-off path. } const thinking: GoogleGeminiCliOptions["thinking"] = { enabled: false }; if (model.reasoning && model.thinking?.suppressWhenOff) { // CCA re-applies the per-id baked server default when the config // is omitted; suppression must be explicit on the wire. thinking.suppress = model.thinking.mode === "google-level" ? { level: "MINIMAL" } : { budget: 0 }; } return castApi<"google-gemini-cli">({ ...base, requestModelId: resolveWireModelId(model, undefined), thinking, toolChoice, antigravityEndpointMode: options?.antigravityEndpointMode, }); } case "google-vertex": { // Explicitly disable thinking when reasoning is absent, unsupported, or // replaced by the caller's external scratchpad. const reasoning = options?.reasoning; if (!reasoning || !model.reasoning || options?.disableReasoning || options?.forceReasoningOff) { return castApi<"google-vertex">({ ...base, serviceTier: options?.serviceTier, thinking: resolveGoogleThinkingOff(model), toolChoice: mapGoogleToolChoice(options?.toolChoice), cachedContent: options?.cachedContent, }); } const vertexModel = model as Model<"google-vertex">; const effort = requireSupportedEffort(vertexModel, reasoning); const geminiModel = vertexModel as unknown as Model<"google-generative-ai">; if (geminiModel.thinking?.mode === "google-level") { return castApi<"google-vertex">({ ...base, serviceTier: options?.serviceTier, thinking: { enabled: true, level: mapEffortToGoogleThinkingLevel(effort, model), }, hideThinkingSummary: options?.hideThinkingSummary, toolChoice: mapGoogleToolChoice(options?.toolChoice), cachedContent: options?.cachedContent, }); } return castApi<"google-vertex">({ ...base, serviceTier: options?.serviceTier, thinking: { enabled: true, budgetTokens: getGoogleBudget(geminiModel, effort, options?.thinkingBudgets), }, hideThinkingSummary: options?.hideThinkingSummary, toolChoice: mapGoogleToolChoice(options?.toolChoice), cachedContent: options?.cachedContent, }); } case "ollama-chat": return castApi<"ollama-chat">({ ...base, reasoning: resolveOpenAiReasoningEffort(model, options), disableReasoning: options?.disableReasoning, toolChoice: options?.toolChoice, }); case "cursor-agent": { const execHandlers = options?.cursorExecHandlers ?? options?.execHandlers; const onToolResult = options?.cursorOnToolResult ?? execHandlers?.onToolResult; const cursorModel = model as Model<"cursor-agent">; const effort = options?.reasoning && !options.disableReasoning && !options.forceReasoningOff && cursorModel.reasoning ? requireSupportedEffort(cursorModel, options.reasoning) : undefined; return castApi<"cursor-agent">({ ...base, execHandlers, onToolResult, externalToolExecutor: options?.cursorExternalToolExecutor, wireModelId: resolveWireModelId(cursorModel, effort), }); } case "factory-droid-agent": { const factoryModel = model as Model<"factory-droid-agent">; const reasoning = options?.reasoning && !options.disableReasoning && !options.forceReasoningOff ? requireSupportedEffort(factoryModel, options.reasoning) : undefined; return castApi<"factory-droid-agent">({ ...base, // The wrapper resolves native defaults for the selected OAuth // account; do not turn the discovery scope's cap into a caller cap. maxTokens: options?.maxTokens, reasoning, disableReasoning: options?.disableReasoning || options?.forceReasoningOff, hideThinkingSummary: options?.hideThinkingSummary, textVerbosity: options?.textVerbosity, serviceTier: options?.serviceTier, toolChoice: options?.toolChoice, }); } case "gitlab-duo-agent": return castApi<"gitlab-duo-agent">({ ...base, cwd: options?.cwd, toolChoice: options?.toolChoice, }); case "apple-foundation-models": return castApi<"apple-foundation-models">({ ...base, toolChoice: options?.toolChoice, reasoning: options?.disableReasoning || options?.forceReasoningOff ? undefined : options?.reasoning, }); case "devin-agent": { const devinModel = model as Model<"devin-agent">; const effort = options?.reasoning && !options.disableReasoning ? requireSupportedEffort(devinModel, options.reasoning) : undefined; return castApi<"devin-agent">({ ...base, chatModelUid: resolveWireModelId(devinModel, effort), }); } default: throw new AIError.ConfigurationError(`Unhandled API in mapOptionsForApi: ${model.api}`); } } function getGoogleBudget( model: Model<"google-generative-ai">, effort: Effort, customBudgets?: ThinkingBudgets, ): number { requireSupportedEffort(model, effort); // Custom budgets take precedence if provided for this level if (customBudgets?.[effort] !== undefined) { return customBudgets[effort]!; } // See https://ai.google.dev/gemini-api/docs/thinking#set-budget const resolvedBudget = model.thinking?.effortBudgets?.[effort]; if (resolvedBudget !== undefined) return resolvedBudget; // Unknown model - use dynamic return -1; }