/** * Proxy Handler — protocol-translating reverse proxy for OpenAI-compatible APIs. * * Receives Anthropic Messages API requests from the CLI, transforms them to * OpenAI Chat Completions or Responses API format, forwards to the upstream * provider, and transforms the response back to Anthropic format. * * Derived from cc-switch (https://github.com/farion1231/cc-switch) * Original work by Jason Young, MIT License */ import { getOpenAIPolicyError } from '../../services/openaiAuth/policyError.js' import { buildOpenaiEndpoint } from './openaiEndpoint.js' import { normalizeAnthropicBaseUrl } from '../../services/api/anthropicBaseUrl.js' import { createGunzip, createInflate } from 'node:zlib' import { ProviderService } from '../services/providerService.js' import { diagnosticsService } from '../services/diagnosticsService.js' import type { ProviderAuthStrategy } from '../types/provider.js' import { resolvePromptCacheKey } from './promptCacheKey.js' import { getOpenAIChatImageContentMode, isOpenAIChatImageRejection, rememberOpenAIChatTextOnlyModel, } from './openaiChatImageSupport.js' import { anthropicToOpenaiChat, type OpenAIChatImageContentMode } from './transform/anthropicToOpenaiChat.js' import { anthropicToOpenaiResponses } from './transform/anthropicToOpenaiResponses.js' import { RequestCompatibilityError, resolveRequestCompatibility, type RequestCompatibilityOptions } from './transform/requestCompatibility.js' import { ProtocolTraceObserver, observeProtocolStream, type ProtocolTraceTransport } from './protocolTrace.js' import { OUTPUT_BUDGET_SOURCE_HEADER } from '../../services/api/outputBudget.js' import { hoistToolResultMediaForCompatibility, shouldHoistNestedToolResultMedia } from './transform/anthropicMediaHoist.js' import { openaiChatToAnthropic } from './transform/openaiChatToAnthropic.js' import { openaiResponsesToAnthropic } from './transform/openaiResponsesToAnthropic.js' import { openaiChatStreamToAnthropic } from './streaming/openaiChatStreamToAnthropic.js' import { openaiResponsesStreamToAnthropic } from './streaming/openaiResponsesStreamToAnthropic.js' import type { AnthropicRequest } from './transform/types.js' import { getProxyFetchOptions } from '../../utils/proxy.js' import { getNetworkProxyFetchOptions, loadNetworkSettings, type NetworkSettings, } from '../services/networkSettings.js' import { normalizeModelStringForAPI } from '../../utils/model/model.js' import { createTraceCallId, createTraceBodySnapshot, TRACE_STREAM_CAPTURE_BYTES, traceCaptureService, type TraceBodySnapshot, type TraceProviderInfo, } from '../services/traceCaptureService.js' import { resolveModelReasoningProfile } from '../../shared/modelReasoning.js' import { resolveModelApiFormat } from '../../shared/modelApiFormats.js' import { applyUpstreamHeaders, resolveUpstreamHeaders } from './upstreamHeaders.js' const providerService = new ProviderService() type ProxyFetchOptions = ReturnType // `decompress` is a Bun fetch option absent from the DOM RequestInit type. type UpstreamRequestInit = RequestInit & ProxyFetchOptions & { decompress?: boolean } type ProxyTraceContext = { sessionId: string provider: TraceProviderInfo anthropicRequest: AnthropicRequest protocolTrace?: ProtocolTraceObserver } const TRACE_RECORDED_ERROR_MARKER = Symbol('cc-haha-trace-recorded-error') // Per-context dedup for failures that rethrow a value that cannot carry a // marker (stream errors may be any value, e.g. a string from // `controller.error('...')`). The marker above still covers Error objects for // paths that only see object throws. const recordedTraceErrorContexts = new WeakSet() function markTraceErrorRecorded(error: unknown): void { if (error && typeof error === 'object') { try { Object.defineProperty(error, TRACE_RECORDED_ERROR_MARKER, { value: true, enumerable: false, }) } catch { // Best effort only; proxy error handling must not depend on trace metadata. } } } function wasTraceErrorRecorded(error: unknown): boolean { return Boolean(error && typeof error === 'object' && (error as Record)[TRACE_RECORDED_ERROR_MARKER]) } function createTimeoutController(timeoutMs: number): { signal: AbortSignal clear: () => void } { const controller = new AbortController() const timer = setTimeout(() => { controller.abort(new DOMException('The operation timed out.', 'TimeoutError')) }, timeoutMs) return { signal: controller.signal, clear: () => clearTimeout(timer), } } async function fetchUpstreamWithTimeout( url: string, init: Omit, timeoutMs: number, isStream: boolean, ): Promise { if (!isStream) { return fetch(url, { ...init, signal: AbortSignal.timeout(timeoutMs), }) } // For streaming requests, this timeout should only cover the connection and // response headers. Keeping the signal alive aborts long generations mid-body. const timeout = createTimeoutController(timeoutMs) try { return await fetch(url, { ...init, signal: timeout.signal, }) } finally { timeout.clear() } } export function withStreamIdleTimeout( upstream: ReadableStream, timeoutMs: number, ): ReadableStream { let reader: ReadableStreamDefaultReader | null = null let timer: ReturnType | null = null const clearIdleTimer = () => { if (timer) { clearTimeout(timer) timer = null } } return new ReadableStream({ async start(controller) { reader = upstream.getReader() let timedOut = false const armIdleTimer = () => { clearIdleTimer() timer = setTimeout(() => { timedOut = true void reader?.cancel('stream idle timeout').catch(() => undefined) controller.error(new Error(`Upstream stream idle timeout after ${timeoutMs}ms`)) }, timeoutMs) } try { armIdleTimer() while (true) { const { done, value } = await reader.read() if (done) break if (timedOut) break controller.enqueue(value) armIdleTimer() } clearIdleTimer() if (!timedOut) controller.close() } catch (err) { clearIdleTimer() if (!timedOut) controller.error(err) } finally { reader?.releaseLock() reader = null } }, cancel(reason) { clearIdleTimer() return reader?.cancel(reason) }, }) } export async function handleProxyRequest(req: Request, url: URL): Promise { const providerMatch = url.pathname.match(/^\/proxy\/providers\/([^/]+)\/v1\/messages$/) const providerId = providerMatch ? decodeURIComponent(providerMatch[1]!) : undefined const isActiveProxyPath = url.pathname === '/proxy/v1/messages' // Only handle POST /proxy/v1/messages or POST /proxy/providers/:providerId/v1/messages if (req.method !== 'POST' || (!isActiveProxyPath && !providerMatch)) { return Response.json( { error: 'Not Found', message: 'Proxy only handles POST /proxy/v1/messages and POST /proxy/providers/:providerId/v1/messages', }, { status: 404 }, ) } // Read active/default provider config or an explicitly-scoped provider config. const config = await providerService.getProviderForProxy(providerId) if (!config) { return Response.json( { type: 'error', error: { type: 'invalid_request_error', message: providerId ? `Provider "${providerId}" is not configured for proxy` : 'No active provider configured for proxy', }, }, { status: 400 }, ) } // Parse request body (needed by both the anthropic-compatible path and the // OpenAI-transforming paths). let body: AnthropicRequest try { body = (await req.json()) as AnthropicRequest } catch { return Response.json( { type: 'error', error: { type: 'invalid_request_error', message: 'Invalid JSON in request body' } }, { status: 400 }, ) } body = { ...body, model: normalizeModelStringForAPI(body.model), } const isStream = body.stream === true const baseUrl = config.baseUrl.replace(/\/+$/, '') const networkSettings = await loadNetworkSettings() const traceContext = buildProxyTraceContext(req, config, body) const inboundSessionId = req.headers.get('x-claude-code-session-id') const promptCacheKey = resolvePromptCacheKey(body, inboundSessionId) const requestOptions: RequestCompatibilityOptions = { requestCompatibility: config.requestCompatibility, budgetSource: req.headers.get(OUTPUT_BUDGET_SOURCE_HEADER) === 'default' ? 'default' : 'explicit', } // Gateways that bind the wire format to the URL path announce their exceptions // on the preset; a model matching no rule keeps the provider's own format. // Resolved per request because one provider record can serve several formats. const modelApiFormat = resolveModelApiFormat(config.modelApiFormats, body.model) const apiFormat = modelApiFormat ?? config.apiFormat const upstreamHeaders = resolveUpstreamHeaders(config.upstreamHeaders, { sessionId: inboundSessionId }) try { if (apiFormat === 'anthropic') { // Anthropic-format providers normally connect directly to the upstream // endpoint (see providerRuntimeEnv). Only providers that explicitly opt out // of nested tool-result media (supportsNestedToolResultMedia=false) route // through the proxy so images/documents can be lifted out of tool_result // before the request reaches an endpoint that would drop them. // // That reasoning is about the provider's *own* format, so the guard only // applies without a per-model override: an OpenAI-format provider whose // model routes to anthropic still reaches the proxy, and its upstream is a // native Messages endpoint that accepts nested media unchanged. if (modelApiFormat === undefined && config.supportsNestedToolResultMedia) { return Response.json( { type: 'error', error: { type: 'invalid_request_error', message: providerId ? `Provider "${providerId}" uses anthropic format — proxy not needed` : 'Active provider uses anthropic format — proxy not needed', }, }, { status: 400 }, ) } const hoistNestedMedia = shouldHoistNestedToolResultMedia({ providerApiFormat: config.apiFormat, resolvedApiFormat: apiFormat, supportsNestedToolResultMedia: config.supportsNestedToolResultMedia, }) return await handleAnthropicCompatible(body, baseUrl, config.apiKey, config.authStrategy, req.headers, isStream, networkSettings, traceContext, upstreamHeaders, hoistNestedMedia) } if (apiFormat === 'openai_chat') { return await handleOpenaiChat(body, baseUrl, config.apiKey, isStream, networkSettings, traceContext, requestOptions, upstreamHeaders) } return await handleOpenaiResponses(body, baseUrl, config.apiKey, isStream, networkSettings, traceContext, promptCacheKey, requestOptions, upstreamHeaders) } catch (err) { if (traceContext && !wasTraceErrorRecorded(err) && !recordedTraceErrorContexts.has(traceContext)) { void recordProxyTrace({ context: traceContext, model: body.model, upstreamUrl: baseUrl, upstreamRequest: null, startedAt: new Date().toISOString(), startedAtMs: Date.now(), error: err, }).catch(() => {}) } console.error('[Proxy] Upstream request failed:', err) const policyError = getOpenAIPolicyError(err) return Response.json( { type: 'error', error: policyError ? { type: 'permission_error', ...policyError } : { type: err instanceof RequestCompatibilityError ? 'invalid_request_error' : 'api_error', message: err instanceof Error ? err.message : String(err), }, }, { status: policyError ? 403 : err instanceof RequestCompatibilityError ? 400 : 502 }, ) } } /** * Build the upstream auth headers for an anthropic-format provider, matching * the strategy semantics that providerRuntimeEnv normally encodes into the * CLI environment (see buildProviderAuthEnv). An empty key omits the auth * header entirely, same as the direct path where the SDK sends no * credential rather than an empty one. */ function buildAnthropicAuthHeaders( apiKey: string, authStrategy: ProviderAuthStrategy, ): Record { switch (authStrategy) { case 'auth_token': case 'auth_token_empty_api_key': return apiKey ? { Authorization: `Bearer ${apiKey}` } : {} case 'dual_same_token': return apiKey ? { 'x-api-key': apiKey, Authorization: `Bearer ${apiKey}` } : {} case 'dual_dummy': return { 'x-api-key': 'dummy', Authorization: 'Bearer dummy' } case 'api_key': default: return apiKey ? { 'x-api-key': apiKey } : {} } } // Connection-management headers that must not be forwarded between hops (they // describe the client↔proxy hop, not the proxy↔upstream hop). Per RFC 9110 the // `Connection` header may also name additional connection-specific headers, // which are added to the deny set dynamically. const HOP_BY_HOP_HEADERS = new Set([ 'connection', 'keep-alive', 'proxy-authenticate', 'proxy-authorization', 'te', 'trailer', 'transfer-encoding', 'upgrade', 'proxy-connection', ]) const INTERNAL_CLIENT_HEADERS = new Set([ OUTPUT_BUDGET_SOURCE_HEADER, 'x-claude-code-session-id', 'x-claude-remote-container-id', 'x-claude-remote-session-id', 'x-client-app', ]) function isInternalClientHeader(name: string, value: string): boolean { if (INTERNAL_CLIENT_HEADERS.has(name)) return true if (name === 'x-app') return value === 'cli' return name === 'user-agent' && /^claude-cli\/[^\s]+\s+\(/i.test(value) } function hopByHopDenySet(headers: Headers): Set { const deny = new Set(HOP_BY_HOP_HEADERS) const connection = headers.get('connection') if (connection) { for (const token of connection.split(',')) { const trimmed = token.trim().toLowerCase() if (trimmed) deny.add(trimmed) } } return deny } /** * Copy a response's entity headers minus hop-by-hop headers. The upstream's * `Connection`-scoped headers describe the proxy↔upstream hop and must not * leak into the proxy↔client hop. */ /** * Parse a `Content-Encoding` header into the codecs to unwind, in decoding * order. `Content-Encoding` lists encodings in application order, so decoding * unwinds them in reverse; `identity` is a no-op. */ function parseContentEncodings(contentEncoding: string | undefined): string[] { return (contentEncoding ?? '') .split(',') .map(encoding => encoding.trim().toLowerCase()) .filter(encoding => encoding !== '' && encoding !== 'identity') .reverse() } /** Codecs the trace decoder can unwind. Anything else (br, zstd, stacked * combinations) is marked unavailable instead of decoding raw bytes as UTF-8. */ const SUPPORTED_TRACE_CODECS = new Set(['gzip', 'x-gzip', 'deflate']) /** * Decode captured upstream bytes for trace storage. The passthrough keeps the * raw bytes (`decompress: false`) so Content-Encoding/Length stay valid for * the client, but the trace should store readable text — decompress a copy * when the upstream compressed the body. * * Unknown encodings (for example `br`) and failed decompression of a known * codec are both marked unavailable — trace capture must never fail the * request, but it must not store bytes decoded as UTF-8 either. */ function decodeTraceBytes(bytes: Uint8Array, contentEncoding: string | undefined): string { const encodings = parseContentEncodings(contentEncoding) // An unknown codec (br, zstd, …) or a stacked combination cannot be // unwound — mark the trace unavailable instead of storing compressed bytes // decoded as UTF-8, matching the streaming branch. if (encodings.some(codec => !SUPPORTED_TRACE_CODECS.has(codec))) { return '[trace body unavailable: unsupported content encoding]' } // Copy into an ArrayBuffer-backed view: Bun's sync decompressors require // Uint8Array (a view over a non-shared buffer). let data: Uint8Array = new Uint8Array(bytes) for (const encoding of encodings) { try { if (encoding === 'gzip' || encoding === 'x-gzip') { data = Bun.gunzipSync(data) } else if (encoding === 'deflate') { data = Bun.inflateSync(data) } else { return '[trace body unavailable: unsupported content encoding]' } } catch { // The header names a codec the decoder supports, but the body is not // valid for it (corrupt member, truncated stream). Decoding whatever // was unwound so far as UTF-8 would store binary garbage — mark the // trace unavailable instead. return '[trace body unavailable: decompression failed]' } } return new TextDecoder().decode(data) } function decodeTraceResponseBody(bytes: ArrayBuffer, headers: Headers | undefined): string { return decodeTraceBytes(new Uint8Array(bytes), headers?.get('content-encoding') ?? undefined) } function stripHopByHopHeaders(headers: Headers): Headers { const deny = hopByHopDenySet(headers) const stripped = new Headers() for (const [name, value] of headers.entries()) { if (!deny.has(name.toLowerCase())) stripped.set(name, value) } return stripped } /** * Forward an Anthropic Messages request to an anthropic-format upstream after * lifting media out of nested tool results (provider opted out of nested media). * The wire format stays Anthropic; only the media placement changes. Protocol * headers and error responses are passed through so SDK classification and * retry behavior are preserved. */ async function handleAnthropicCompatible( body: AnthropicRequest, baseUrl: string, apiKey: string, authStrategy: ProviderAuthStrategy, incomingHeaders: Headers, isStream: boolean, networkSettings: NetworkSettings, traceContext: ProxyTraceContext | null, upstreamHeaders: Record = {}, hoistNestedMedia = true, ): Promise { // The media hoist is a compatibility rewrite for endpoints that drop media // nested in tool_result. A native Messages endpoint accepts it unchanged, so it // only runs when the provider actually asked for it. const transformed = hoistNestedMedia ? hoistToolResultMediaForCompatibility(body) : body const url = `${normalizeAnthropicBaseUrl(baseUrl)}/v1/messages` const proxyOptions = getNetworkProxyFetchOptions(networkSettings, url) const headers: Record = { 'Content-Type': 'application/json', ...buildAnthropicAuthHeaders(apiKey, authStrategy), } // Preserve protocol and custom headers from the incoming request // (anthropic-version is required; anthropic-beta and custom headers such as // those injected via ANTHROPIC_CUSTOM_HEADERS carry real semantics for the // upstream endpoint). Hop-by-hop and auth headers are not forwarded. const deny = hopByHopDenySet(incomingHeaders) for (const [name, value] of incomingHeaders.entries()) { const lower = name.toLowerCase() if (deny.has(lower) || isInternalClientHeader(lower, value)) continue // The local proxy's authority must not replace the upstream host. if (lower === 'host' || lower === 'content-type' || lower === 'content-length') continue if (lower === 'x-api-key' || lower === 'authorization') continue if (value) headers[name] = value } // Preset-declared headers go last so they win over the pass-through above — // notably the session id, which is internal-client-filtered on the way in and // would otherwise never reach a gateway that requires it. The resolver drops // anything that would shadow the request framing or its credential. applyUpstreamHeaders(headers, upstreamHeaders) const traceHeaders = Object.fromEntries( Object.entries(headers).map(([name, value]) => { const lower = name.toLowerCase() return [ name, lower === 'content-type' || lower === 'anthropic-version' || lower === 'anthropic-beta' ? value : '[redacted]', ] }), ) const startedAtMs = Date.now() const startedAt = new Date(startedAtMs).toISOString() const traceCallId = traceContext ? startProxyTraceCall({ context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: traceHeaders, startedAt, }) : undefined // Close the pending trace started above when the upstream call fails, so the // caller's unified error handling does not record a second trace for the // same request. const recordTraceError = (err: unknown): void => { if (!traceContext) return recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: traceHeaders, startedAt, startedAtMs, error: err, }) markTraceErrorRecorded(err) recordedTraceErrorContexts.add(traceContext) } let upstream: Response try { upstream = await fetchUpstreamWithTimeout(url, { method: 'POST', headers, body: JSON.stringify(transformed), // Keep the raw bytes: Bun decompresses by default, which would leave a // decompressed body behind the upstream Content-Encoding/Length headers // when forwarding the response unchanged. decompress: false, ...proxyOptions, }, networkSettings.aiRequestTimeoutMs, isStream) } catch (err) { recordTraceError(err) console.error('[Proxy] Upstream anthropic request failed:', err) return Response.json( { type: 'error', error: { type: 'api_error', message: err instanceof Error ? err.message : String(err), }, }, { status: 502 }, ) } try { if (!upstream.ok) { // Pass the upstream error body and headers through unchanged so the SDK // keeps error classification (authentication_error, rate_limit_error, …), // request_id, and retry-after semantics. A body read failure here closes // the pending trace and surfaces to the caller's unified error handling // (structured 502) like any other upstream failure. const errBody = await upstream.arrayBuffer() if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: traceHeaders, startedAt, startedAtMs, responseStatus: upstream.status, upstreamResponseBody: decodeTraceResponseBody(errBody, upstream.headers), responseHeaders: upstream.headers, }) } return new Response(errBody, { status: upstream.status, headers: stripHopByHopHeaders(upstream.headers), }) } if (isStream) { if (!upstream.body) { if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: traceHeaders, startedAt, startedAtMs, error: new Error('Upstream returned no body for stream'), }) } return Response.json( { type: 'error', error: { type: 'api_error', message: 'Upstream returned no body for stream' } }, { status: 502 }, ) } // Keep SSE framing headers while passing through request/rate-limit // metadata from the upstream (request_id, ratelimit-*, custom headers). const responseHeaders = stripHopByHopHeaders(upstream.headers) responseHeaders.set('Content-Type', 'text/event-stream') responseHeaders.set('Cache-Control', 'no-cache') responseHeaders.set('Connection', 'keep-alive') const anthropicStream = withStreamIdleTimeout(upstream.body, networkSettings.aiRequestTimeoutMs) const tracedStream = traceContext ? captureTraceStream(anthropicStream, async (bodySnapshot, error, protocolTraceEnd) => { await recordProxyTrace({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: traceHeaders, startedAt, startedAtMs, responseStatus: 200, responseBodySnapshot: bodySnapshot, protocolTraceEnd, responseHeaders: upstream.headers, ...(error ? { error } : {}), }) }, upstream.headers.get('content-encoding') ?? undefined) : anthropicStream return new Response(tracedStream, { status: 200, headers: responseHeaders, }) } // Byte-for-byte passthrough: re-serializing the body would invalidate // Content-Length/ETag entity headers from the upstream. const responseBody = await upstream.arrayBuffer() if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: traceHeaders, startedAt, startedAtMs, responseStatus: upstream.status, upstreamResponseBody: decodeTraceResponseBody(responseBody, upstream.headers), responseHeaders: upstream.headers, }) } return new Response(responseBody, { status: upstream.status, headers: stripHopByHopHeaders(upstream.headers), }) } catch (err) { // A body read failure closes the pending trace with the original call id // (so no trace stays pending and no second trace is created), then // rethrows so the caller returns the same structured 502 as any other // upstream failure. recordTraceError(err) throw err } } async function handleOpenaiChat( body: AnthropicRequest, baseUrl: string, apiKey: string, isStream: boolean, networkSettings: NetworkSettings, traceContext: ProxyTraceContext | null, requestOptions: RequestCompatibilityOptions = {}, upstreamHeaders: Record = {}, imageContentModeOverride?: OpenAIChatImageContentMode, ): Promise { const knownDeepSeekHost = shouldUseDeepSeekReasoningCompat(baseUrl) const reasoningProfile = resolveModelReasoningProfile(body.model, 'openai_chat') const url = buildOpenaiEndpoint(baseUrl, 'chat/completions') const imageContentMode = imageContentModeOverride ?? getOpenAIChatImageContentMode(url, body.model) const transformed = anthropicToOpenaiChat(body, { ...requestOptions, roundTripReasoningContent: knownDeepSeekHost || reasoningProfile?.family === 'deepseek-v4', passThinkingToggle: knownDeepSeekHost, imageContentMode, }) if (traceContext) { traceContext.protocolTrace = new ProtocolTraceObserver('openai_chat', transformed, resolveRequestCompatibility(body, { ...requestOptions, protocol: 'openai_chat' }).outputBudget) } // Preset-declared headers first: `Authorization` is applied last so a preset can // never shadow the credential, and the resolver already drops framing headers. const upstreamRequestHeaders: Record = {} applyUpstreamHeaders(upstreamRequestHeaders, upstreamHeaders) upstreamRequestHeaders['Content-Type'] = 'application/json' upstreamRequestHeaders.Authorization = `Bearer ${apiKey}` const proxyOptions = getNetworkProxyFetchOptions(networkSettings, url) const startedAtMs = Date.now() const startedAt = new Date(startedAtMs).toISOString() const traceCallId = traceContext ? startProxyTraceCall({ context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, }) : undefined let upstream: Response try { upstream = await fetchUpstreamWithTimeout(url, { method: 'POST', headers: upstreamRequestHeaders, body: JSON.stringify(transformed), ...proxyOptions, }, networkSettings.aiRequestTimeoutMs, isStream) } catch (err) { if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, error: err, }) markTraceErrorRecorded(err) } throw err } if (!upstream.ok) { const errText = await upstream.text().catch(() => '') if (imageContentMode === 'vision' && isOpenAIChatImageRejection(upstream.status, errText, transformed)) { // The upstream refused the request before generating anything, so the // client has seen no output and resending is safe. The model is only // remembered once the same request succeeds without images, so an // unrelated failure of the resend cannot disable images for good. if (traceContext) { recordProxyTraceInBackground({ context: traceContext, callId: traceCallId, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: upstream.status, upstreamResponseBody: errText, responseHeaders: upstream.headers, }) } const resent = await handleOpenaiChat( body, baseUrl, apiKey, isStream, networkSettings, traceContext, requestOptions, upstreamHeaders, 'text_only', ) if (resent.ok) { rememberOpenAIChatTextOnlyModel(url, body.model) void diagnosticsService.recordEvent({ type: 'openai_chat_image_input_disabled', severity: 'info', summary: `Upstream rejected image input for ${body.model}; images are omitted for this model for 30 minutes`, details: { model: body.model, upstreamUrl: url, httpStatus: upstream.status, upstreamError: errText.slice(0, 500) }, }) } return resent } let policyError = null try { policyError = getOpenAIPolicyError(JSON.parse(errText)) } catch { // Unstructured upstream failures keep their existing error classification. } const errorBody = { type: 'error', error: policyError ? { type: 'permission_error', ...policyError } : { type: 'api_error', message: `Upstream returned HTTP ${upstream.status}: ${errText.slice(0, 500)}`, }, } if (traceContext) { recordProxyTraceInBackground({ context: traceContext, callId: traceCallId, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: upstream.status, upstreamResponseBody: errText, anthropicResponseBody: errorBody, responseHeaders: upstream.headers, }) } return Response.json( errorBody, { status: policyError ? 403 : upstream.status }, ) } if (isStream) { if (!upstream.body) { if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, error: new Error('Upstream returned no body for stream'), }) } return Response.json( { type: 'error', error: { type: 'api_error', message: 'Upstream returned no body for stream' } }, { status: 502 }, ) } const upstreamBody = withStreamIdleTimeout(upstream.body, networkSettings.aiRequestTimeoutMs) const observedBody = traceContext?.protocolTrace ? observeProtocolStream(upstreamBody, traceContext.protocolTrace) : upstreamBody const anthropicStream = openaiChatStreamToAnthropic(observedBody, body.model) const tracedStream = traceContext ? captureTraceStream(anthropicStream, async (bodySnapshot, error, protocolTraceEnd) => { await recordProxyTrace({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: 200, responseBodySnapshot: bodySnapshot, protocolTraceEnd, responseHeaders: upstream.headers, ...(error ? { error } : {}), }) }) : anthropicStream return new Response(tracedStream, { status: 200, headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', Connection: 'keep-alive', }, }) } // Non-streaming const responseBody = await upstream.json() const policyError = getOpenAIPolicyError(responseBody) const anthropicResponse = policyError ? { type: 'error', error: { type: 'permission_error', ...policyError } } : openaiChatToAnthropic(responseBody, body.model) if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: 200, upstreamResponseBody: responseBody, anthropicResponseBody: anthropicResponse, responseHeaders: upstream.headers, }) } return Response.json(anthropicResponse, { status: policyError ? 403 : 200 }) } function shouldUseDeepSeekReasoningCompat(baseUrl: string): boolean { return ( /(^|[./-])deepseek([./-]|$)/i.test(baseUrl) || /(^|[./-])opencode\.ai([:/]|$)/i.test(baseUrl) ) } async function handleOpenaiResponses( body: AnthropicRequest, baseUrl: string, apiKey: string, isStream: boolean, networkSettings: NetworkSettings, traceContext: ProxyTraceContext | null, promptCacheKey?: string, requestOptions: RequestCompatibilityOptions = {}, upstreamHeaders: Record = {}, ): Promise { const transformed = anthropicToOpenaiResponses(body, { ...requestOptions, cacheKey: promptCacheKey }) if (traceContext) { traceContext.protocolTrace = new ProtocolTraceObserver('openai_responses', transformed, resolveRequestCompatibility(body, { ...requestOptions, protocol: 'openai_responses' }).outputBudget) } const url = buildOpenaiEndpoint(baseUrl, 'responses') // Preset-declared headers first: `Authorization` is applied last so a preset can // never shadow the credential, and the resolver already drops framing headers. const upstreamRequestHeaders: Record = {} applyUpstreamHeaders(upstreamRequestHeaders, upstreamHeaders) upstreamRequestHeaders['Content-Type'] = 'application/json' upstreamRequestHeaders.Authorization = `Bearer ${apiKey}` const proxyOptions = getNetworkProxyFetchOptions(networkSettings, url) const startedAtMs = Date.now() const startedAt = new Date(startedAtMs).toISOString() const traceCallId = traceContext ? startProxyTraceCall({ context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, }) : undefined let upstream: Response try { upstream = await fetchUpstreamWithTimeout(url, { method: 'POST', headers: upstreamRequestHeaders, body: JSON.stringify(transformed), ...proxyOptions, }, networkSettings.aiRequestTimeoutMs, isStream) } catch (err) { if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, error: err, }) markTraceErrorRecorded(err) } throw err } if (!upstream.ok) { const errText = await upstream.text().catch(() => '') let policyError = null try { policyError = getOpenAIPolicyError(JSON.parse(errText)) } catch { // Unstructured upstream failures keep their existing error classification. } const errorBody = { type: 'error', error: policyError ? { type: 'permission_error', ...policyError } : { type: 'api_error', message: `Upstream returned HTTP ${upstream.status}: ${errText.slice(0, 500)}`, }, } if (traceContext) { recordProxyTraceInBackground({ context: traceContext, callId: traceCallId, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: upstream.status, upstreamResponseBody: errText, anthropicResponseBody: errorBody, responseHeaders: upstream.headers, }) } return Response.json( errorBody, { status: policyError ? 403 : upstream.status }, ) } if (isStream) { if (!upstream.body) { if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, error: new Error('Upstream returned no body for stream'), }) } return Response.json( { type: 'error', error: { type: 'api_error', message: 'Upstream returned no body for stream' } }, { status: 502 }, ) } const upstreamBody = withStreamIdleTimeout(upstream.body, networkSettings.aiRequestTimeoutMs) const observedBody = traceContext?.protocolTrace ? observeProtocolStream(upstreamBody, traceContext.protocolTrace) : upstreamBody const anthropicStream = openaiResponsesStreamToAnthropic(observedBody, body.model) const tracedStream = traceContext ? captureTraceStream(anthropicStream, async (bodySnapshot, error, protocolTraceEnd) => { await recordProxyTrace({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: 200, responseBodySnapshot: bodySnapshot, protocolTraceEnd, responseHeaders: upstream.headers, ...(error ? { error } : {}), }) }) : anthropicStream return new Response(tracedStream, { status: 200, headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', Connection: 'keep-alive', }, }) } // Non-streaming const responseBody = await upstream.json() const policyError = getOpenAIPolicyError(responseBody) const anthropicResponse = policyError ? { type: 'error', error: { type: 'permission_error', ...policyError } } : openaiResponsesToAnthropic(responseBody, body.model) if (traceContext) { recordProxyTraceInBackground({ callId: traceCallId, context: traceContext, model: body.model, upstreamUrl: url, upstreamRequest: transformed, requestHeaders: upstreamRequestHeaders, startedAt, startedAtMs, responseStatus: 200, upstreamResponseBody: responseBody, anthropicResponseBody: anthropicResponse, responseHeaders: upstream.headers, }) } return Response.json(anthropicResponse, { status: policyError ? 403 : 200 }) } function buildProxyTraceContext( req: Request, config: { id: string; name: string; apiFormat: string }, anthropicRequest: AnthropicRequest, ): ProxyTraceContext | null { const sessionId = req.headers.get('x-claude-code-session-id')?.trim() if (!sessionId) return null return { sessionId, provider: { id: config.id, name: config.name, format: config.apiFormat, }, anthropicRequest, } } function createProxyTraceRequestBody(context: ProxyTraceContext, upstreamRequest: unknown): Record { return upstreamRequest ? { anthropic: context.anthropicRequest, upstream: upstreamRequest, } : { anthropic: context.anthropicRequest, } } function startProxyTraceCall({ context, model, upstreamUrl, upstreamRequest, requestHeaders, startedAt, }: { context: ProxyTraceContext model: string upstreamUrl: string upstreamRequest: unknown requestHeaders: Record startedAt: string }): string { const callId = createTraceCallId() void traceCaptureService.recordCall({ id: callId, sessionId: context.sessionId, source: 'proxy', provider: context.provider, model, status: 'pending', startedAt, request: { method: 'POST', url: upstreamUrl, headers: requestHeaders, bodySnapshot: createTraceBodySnapshot({ pending: true, note: 'proxy request body captured on call completion', }), }, metadata: { phase: 'upstream_fetch_started', ...(context.protocolTrace ? { protocolTrace: context.protocolTrace.snapshot() } : {}), }, }) void traceCaptureService.recordEvent({ sessionId: context.sessionId, callId, source: 'proxy', provider: context.provider, model, timestamp: startedAt, phase: 'upstream_fetch_started', severity: 'info', title: 'Upstream fetch started', metadata: { url: upstreamUrl, }, }) return callId } type RecordProxyTraceInput = { callId?: string context: ProxyTraceContext model: string upstreamUrl: string upstreamRequest: unknown requestHeaders?: Record startedAt: string startedAtMs: number responseStatus?: number upstreamResponseBody?: unknown anthropicResponseBody?: unknown responseBodySnapshot?: TraceBodySnapshot responseHeaders?: Headers error?: unknown protocolTraceEnd?: Exclude } function recordProxyTraceInBackground(input: RecordProxyTraceInput): void { void recordProxyTrace(input).catch(() => {}) } async function recordProxyTrace({ callId, context, model, upstreamUrl, upstreamRequest, requestHeaders, startedAt, startedAtMs, responseStatus, upstreamResponseBody, anthropicResponseBody, responseBodySnapshot, responseHeaders, error, protocolTraceEnd, }: RecordProxyTraceInput): Promise { const completedAt = new Date().toISOString() const requestBody = createProxyTraceRequestBody(context, upstreamRequest) const responseBody = anthropicResponseBody === undefined && upstreamResponseBody === undefined ? undefined : { ...(upstreamResponseBody !== undefined ? { upstream: upstreamResponseBody } : {}), ...(anthropicResponseBody !== undefined ? { anthropic: anthropicResponseBody } : {}), } const observer = context.protocolTrace if (observer) { if (upstreamResponseBody !== undefined) { if (typeof upstreamResponseBody !== 'string') observer.observeJson(upstreamResponseBody) else if (upstreamResponseBody.length <= 64 * 1024) { try { observer.observeJson(JSON.parse(upstreamResponseBody)) } catch { /* Unstructured HTTP error. */ } } observer.finish(error || (responseStatus ?? 200) >= 400 ? 'error' : 'non_stream') } else if (protocolTraceEnd) { observer.finish(protocolTraceEnd) } else if (error) { observer.finish('error') } } await traceCaptureService.recordCall({ ...(callId ? { id: callId } : {}), sessionId: context.sessionId, source: 'proxy', provider: context.provider, model, startedAt, completedAt, durationMs: Date.now() - startedAtMs, request: { method: 'POST', url: upstreamUrl, headers: requestHeaders, body: requestBody, }, ...(responseStatus !== undefined ? { response: { status: responseStatus, headers: responseHeaders, ...(responseBodySnapshot ? { bodySnapshot: responseBodySnapshot } : { body: responseBody }), }, } : {}), ...(error ? { error } : {}), metadata: { phase: error ? 'upstream_fetch_failed' : 'upstream_fetch_completed', ...(observer ? { protocolTrace: { ...observer.snapshot(), ...(protocolTraceEnd ? { delivery: protocolTraceEnd } : {}), }, } : {}), }, }) await traceCaptureService.recordEvent({ sessionId: context.sessionId, ...(callId ? { callId } : {}), source: 'proxy', provider: context.provider, model, timestamp: completedAt, phase: error ? 'upstream_fetch_failed' : 'upstream_fetch_completed', severity: error ? 'error' : responseStatus !== undefined && responseStatus >= 400 ? 'warning' : 'info', title: error ? 'Upstream fetch failed' : 'Upstream fetch completed', message: error instanceof Error ? error.message : error ? String(error) : undefined, metadata: { status: responseStatus, url: upstreamUrl, }, }) } function captureTraceStream( stream: ReadableStream, onComplete: (snapshot: TraceBodySnapshot, error?: unknown, end?: Exclude) => Promise, contentEncoding?: string, ): ReadableStream { // The Anthropic passthrough forwards raw upstream bytes (`decompress: // false`), so a compressed SSE body would otherwise be stored as binary // garbage. Decompress the *trace copy* while it streams — the capture cap // applies to decoded output, so truncation, client cancellation, or an // upstream error leave a readable plain-text prefix instead of an // unterminated gzip member, and a highly compressible body cannot blow past // the memory cap before it is counted. const chunks: Uint8Array[] = [] let bytes = 0 let truncated = false let finalized = false let reader: ReadableStreamDefaultReader | null = null const captureDecoded = (chunk: Uint8Array) => { bytes += chunk.byteLength if (bytes <= TRACE_STREAM_CAPTURE_BYTES) { chunks.push(chunk) } else { truncated = true } } const finalize = async (error?: unknown, end: Exclude = 'eof') => { if (finalized) return finalized = true const joined = new Uint8Array(chunks.reduce((total, chunk) => total + chunk.byteLength, 0)) let offset = 0 for (const chunk of chunks) { joined.set(chunk, offset) offset += chunk.byteLength } const snapshot = createTraceBodySnapshot( unsupportedEncoding ? '[trace body unavailable: unsupported content encoding]' : unexpectedDecompressionFailure ? '[trace body unavailable: decompression failed]' : new TextDecoder().decode(joined), { alreadyTruncated: truncated }, ) await onComplete(snapshot, error, end).catch(() => {}) } // The streaming branch decodes a *single* known codec (gzip/x-gzip or // deflate). Stacked or unknown encodings cannot be unwound here — mark the // trace unavailable instead of storing compressed bytes decoded as UTF-8. // The buffered path still unwinds every codec via decodeTraceBytes. const encodings = parseContentEncodings(contentEncoding) const singleKnownCodec = encodings.length === 1 && SUPPORTED_TRACE_CODECS.has(encodings[0]!) const unsupportedEncoding = encodings.length > 0 && !singleKnownCodec const decompressor = singleKnownCodec ? encodings[0] === 'deflate' ? createInflate() : createGunzip() : null // node:zlib streams honor the Writable backpressure contract: write() // returns false when the writable buffer is full and the caller must wait // for 'drain' before writing more. The trace copy is a side channel, but it // still must not buffer an unbounded amount of compressed input. let decompressorFailed = false let decompressorEnded = false // Explicitly ended by design (client cancel, capture cap, upstream read // error): an end() on an unterminated gzip member then errors as expected // and the decoded plain-text prefix is kept. Only a body that fails to // decompress while ending normally marks the trace unavailable — an error // from zlib is delivered asynchronously, so "ended before the error" is not // a reliable signal, but the ending path itself is. let activelyEnded = false let unexpectedDecompressionFailure = false let decompressEnded: Promise = Promise.resolve() if (decompressor) { decompressor.on('data', captureDecoded) // Partial data is already captured; an error mid-stream must not surface // beyond the trace copy. decompressor.on('error', () => { if (!activelyEnded) { unexpectedDecompressionFailure = true } decompressorFailed = true }) decompressEnded = new Promise(resolve => { decompressor.on('end', resolve) decompressor.on('error', resolve) }) } // Resolve when the decompressor is ready for more input. An errored stream // never drains and rejects further writes, so failure also resolves — the // loop must stop feeding it afterwards. An end() from the cancel/cap path // may finish through 'finish'/'close' without ever emitting 'drain', so the // waiter settles on any of the terminal events or the read loop would hang // with the upstream reader lock never released. const waitForDecompressorDrain = (): Promise => { if (!decompressor || decompressorFailed || decompressorEnded) return Promise.resolve() return new Promise(resolve => { const settle = () => { decompressor.off('drain', settle) decompressor.off('error', settle) decompressor.off('finish', settle) decompressor.off('close', settle) resolve() } decompressor.once('drain', settle) decompressor.once('error', settle) decompressor.once('finish', settle) decompressor.once('close', settle) }) } // End the decompressor gracefully instead of destroying it: destroy() // would drop data already written but not yet flushed as 'data', so a // cancelled or errored stream would lose the decoded prefix. end() flushes // what was received; an unterminated gzip member then errors, which // resolves decompressEnded through the error branch. const finishDecompressor = async () => { if (!decompressor) return if (!decompressorEnded) { decompressorEnded = true try { decompressor.end() } catch { decompressor.destroy() return } } await decompressEnded } return new ReadableStream({ async start(controller) { reader = stream.getReader() try { while (true) { const { done, value } = await reader.read() if (done) break controller.enqueue(value) if (decompressor) { if (!decompressorFailed && !decompressorEnded) { if (truncated) { // The decoded trace already exceeded the capture cap: stop // feeding the decompressor so the trace side work cannot // grow without bound or delay upstream cancellation. end() // flushes the bytes already accepted, preserving the // captured plain-text prefix. decompressorEnded = true activelyEnded = true decompressor.end() } else if (!decompressor.write(value)) { await waitForDecompressorDrain() } } } else if (!unsupportedEncoding) { captureDecoded(value) } } controller.close() await finishDecompressor() void finalize() } catch (err) { controller.error(err) activelyEnded = true await finishDecompressor() void finalize(err, 'error') } finally { reader?.releaseLock() reader = null } }, async cancel(reason) { const error = reason instanceof Error ? reason : new Error(reason ? `Stream cancelled: ${String(reason)}` : 'Stream cancelled') activelyEnded = true await finishDecompressor() void finalize(error, 'cancelled') await reader?.cancel(reason).catch(() => undefined) }, }) }