diff --git a/.env.example b/.env.example index 565ac4da38..f59c075bd9 100644 --- a/.env.example +++ b/.env.example @@ -423,8 +423,15 @@ ANTHROPIC_API_KEY=sk-ant-your-key-here # Turn on whichever audience you're debugging; both can run together. # OPENCLAUDE_LOG_TOKEN_USAGE=verbose -# Custom timeout for API requests in milliseconds (default: varies) -# API_TIMEOUT_MS=60000 +# Time-to-response-headers deadline for OpenAI-compatible API requests +# in milliseconds (default: 600000, or 10 minutes). Use a safe positive +# integer; invalid values use the default and values above 2147483647 are capped. +# This runtime setting must be exported from your shell or launcher; the +# provider env-file loader intentionally ignores runtime/debug knobs. +# This covers generic OpenAI-compatible requests, direct GitHub Copilot +# Responses, and Copilot chat-to-Responses fallback requests. First-party +# Codex OAuth Responses and the Anthropic SDK retain their existing handling. +# API_TIMEOUT_MS=600000 # Enable debug logging # CLAUDE_DEBUG=1 diff --git a/docs/advanced-setup.md b/docs/advanced-setup.md index 9fcb0c3305..db306a79ab 100644 --- a/docs/advanced-setup.md +++ b/docs/advanced-setup.md @@ -437,6 +437,7 @@ token; do not set both credentials. | `OPENAI_MODEL` | OpenAI-compatible only | Model name such as `gpt-4o`, `deepseek-v4-flash`, or `llama3.3:70b` | | `OPENAI_BASE_URL` | No | API endpoint, defaulting to `https://api.openai.com/v1` | | `OPENAI_API_BASE` | No | Compatibility alias for `OPENAI_BASE_URL` | +| `API_TIMEOUT_MS` | No | Time-to-response-headers deadline for generic OpenAI-compatible requests, direct GitHub Copilot Responses, and Copilot chat-to-Responses fallback requests, in milliseconds (default: `600000`, or 10 minutes). The value must be a safe positive integer; invalid, zero, negative, or fractional values use the default, and values above `2147483647` are capped. The deadline is disarmed after headers arrive, so it does not limit response streaming. Export this runtime setting from your shell or launcher; the provider env-file loader ignores runtime/debug settings, so a value configured only there leaves the default in effect. First-party Codex OAuth Responses and the Anthropic SDK retain their existing timeout handling. | | `OPENCLAUDE_OLLAMA_NUM_CTX` | Ollama only | Request-level Ollama context window. Defaults to `32768`; set a larger value for longer same-session history if your model and hardware can handle it. | | `CLAUDE_CODE_OPENAI_CONTEXT_WINDOWS` | No | JSON map of OpenAI-compatible model names to context windows, such as `{"custom-model":1000000}`. Use this when a custom provider does not expose context metadata from `/v1/models`. | | `CLAUDE_CODE_OPENAI_MAX_OUTPUT_TOKENS` | No | JSON map of OpenAI-compatible model names to max output tokens, such as `{"custom-model":32768}`. Use this when a custom provider does not expose output-limit metadata from `/v1/models`. | diff --git a/src/services/api/codexShim.ts b/src/services/api/codexShim.ts index 09e0cc0f80..f97ef951e0 100644 --- a/src/services/api/codexShim.ts +++ b/src/services/api/codexShim.ts @@ -600,6 +600,7 @@ export async function performCodexRequest(options: { params: ShimCreateParams defaultHeaders: Record signal?: AbortSignal + fetcher?: typeof fetchWithProxyRetry }): Promise { const compressedMessages = compressToolHistory( options.params.messages as Array<{ @@ -678,7 +679,7 @@ export async function performCodexRequest(options: { } headers.originator ??= 'openclaude' - const response = await fetchWithProxyRetry( + const response = await (options.fetcher ?? fetchWithProxyRetry)( `${options.request.baseUrl}/responses`, { method: 'POST', diff --git a/src/services/api/fetchWithProxyRetry.test.ts b/src/services/api/fetchWithProxyRetry.test.ts index 970bbf4431..c7da2a4d2f 100644 --- a/src/services/api/fetchWithProxyRetry.test.ts +++ b/src/services/api/fetchWithProxyRetry.test.ts @@ -143,3 +143,120 @@ test('fetchWithProxyRetry retries and disables keepalive after receiving a 504 r expect((calls[0] as RequestInit).keepalive).toBeUndefined() expect((calls[1] as RequestInit).keepalive).toBe(false) }) + +test('fetchWithProxyRetry retries when cancelling a discarded 504 body stalls', async () => { + let attempts = 0 + + globalThis.fetch = (async () => { + attempts++ + if (attempts === 1) { + return new Response(new ReadableStream({ + cancel() { + return new Promise(() => {}) + }, + }), { status: 504 }) + } + return new Response('ok') + }) as unknown as FetchType + + const response = await fetchWithProxyRetry('https://example.com/search') + + expect(response.status).toBe(200) + expect(attempts).toBe(2) +}) + +test('fetchWithProxyRetry does not retry a 504 after the request is aborted', async () => { + const controller = new AbortController() + const abortReason = new DOMException('Deadline exceeded', 'TimeoutError') + let attempts = 0 + let bodyCancelled = false + + globalThis.fetch = (async () => { + attempts++ + controller.abort(abortReason) + return new Response(new ReadableStream({ + cancel() { + bodyCancelled = true + }, + }), { status: 504 }) + }) as unknown as FetchType + + await expect( + fetchWithProxyRetry('https://example.com/generate', { + method: 'POST', + signal: controller.signal, + }), + ).rejects.toBe(abortReason) + + expect(attempts).toBe(1) + await Promise.resolve() + expect(bodyCancelled).toBe(true) +}) + +test('fetchWithProxyRetry honors an aborted Request signal without replaying', async () => { + const controller = new AbortController() + const abortReason = new DOMException('Deadline exceeded', 'TimeoutError') + let attempts = 0 + + globalThis.fetch = (async () => { + attempts++ + controller.abort(abortReason) + return new Response('Gateway Timeout', { status: 504 }) + }) as unknown as FetchType + + const request = new Request('https://example.com/generate', { + method: 'POST', + signal: controller.signal, + }) + + await expect(fetchWithProxyRetry(request)).rejects.toBe(abortReason) + expect(attempts).toBe(1) +}) + +test('fetchWithProxyRetry preserves the abort reason for a generic fetch failure', async () => { + for (const message of ['fetch failed', 'invalid_argument']) { + const controller = new AbortController() + const abortReason = new DOMException('Deadline exceeded', 'TimeoutError') + let attempts = 0 + + globalThis.fetch = (async () => { + attempts++ + controller.abort(abortReason) + throw new TypeError(message) + }) as unknown as FetchType + + await expect( + fetchWithProxyRetry('https://example.com/generate', { + method: 'POST', + signal: controller.signal, + }), + ).rejects.toBe(abortReason) + + expect(attempts).toBe(1) + } +}) + +test('fetchWithProxyRetry preserves an explicit AbortError from fetch', async () => { + const controller = new AbortController() + const abortReason = new DOMException('Caller cancelled', 'AbortError') + const fetchAbortError = new DOMException( + 'The operation was aborted.', + 'AbortError', + ) + let attempts = 0 + + globalThis.fetch = (async () => { + attempts++ + controller.abort(abortReason) + throw fetchAbortError + }) as unknown as FetchType + + await expect( + fetchWithProxyRetry('https://example.com/generate', { + method: 'POST', + signal: controller.signal, + }), + ).rejects.toBe(fetchAbortError) + + expect(attempts).toBe(1) +}) diff --git a/src/services/api/fetchWithProxyRetry.ts b/src/services/api/fetchWithProxyRetry.ts index 1c02965027..a9aa6347b7 100644 --- a/src/services/api/fetchWithProxyRetry.ts +++ b/src/services/api/fetchWithProxyRetry.ts @@ -3,6 +3,11 @@ import { disableKeepAlive, getProxyFetchOptions } from '../../utils/proxy.js' const RETRYABLE_FETCH_ERROR_PATTERN = /socket connection was closed unexpectedly|ECONNRESET|EPIPE|socket hang up|Connection reset by peer|fetch failed/i +export type ProxyRetryFetcher = ( + input: string | URL | Request, + init?: RequestInit, +) => Promise + export function isRetryableFetchError(error: unknown): boolean { if (!(error instanceof Error)) { return false @@ -16,19 +21,32 @@ export function isRetryableFetchError(error: unknown): boolean { export async function fetchWithProxyRetry( input: string | URL | Request, init?: RequestInit, - options?: { forAnthropicAPI?: boolean; maxAttempts?: number }, + options?: { + forAnthropicAPI?: boolean + maxAttempts?: number + fetcher?: ProxyRetryFetcher + }, ): Promise { const maxAttempts = Math.max(1, options?.maxAttempts ?? 2) + const fetcher = options?.fetcher ?? fetch + const signal = init?.signal ?? (input instanceof Request ? input.signal : undefined) let lastError: unknown for (let attempt = 1; attempt <= maxAttempts; attempt++) { try { - const response = await fetch(input, { + const response = await fetcher(input, { ...init, ...getProxyFetchOptions({ forAnthropicAPI: options?.forAnthropicAPI, }), }) + if (signal?.aborted) { + void response.body?.cancel().catch(() => {}) + throw ( + signal.reason ?? + new DOMException('The operation was aborted.', 'AbortError') + ) + } // If an upstream proxy or local NAT silently dropped the keep-alive socket, // it might result in a 502/504 response instead of a hard network exception. @@ -37,6 +55,7 @@ export async function fetchWithProxyRetry( (response.status === 502 || response.status === 504) && attempt < maxAttempts ) { + void response.body?.cancel().catch(() => {}) disableKeepAlive() continue } @@ -44,7 +63,15 @@ export async function fetchWithProxyRetry( return response } catch (error) { lastError = error - if (attempt >= maxAttempts || !isRetryableFetchError(error)) { + if (signal?.aborted) { + throw error instanceof Error && error.name === 'AbortError' + ? error + : (signal.reason ?? error) + } + if ( + attempt >= maxAttempts || + !isRetryableFetchError(error) + ) { throw error } disableKeepAlive() diff --git a/src/services/api/openaiErrorClassification.ts b/src/services/api/openaiErrorClassification.ts index dff43ccc25..1be20bfd1b 100644 --- a/src/services/api/openaiErrorClassification.ts +++ b/src/services/api/openaiErrorClassification.ts @@ -27,6 +27,28 @@ export type OpenAICompatibilityFailure = { requestUrl?: string } +const NON_REPLAYABLE_OPENAI_REQUEST = Symbol.for( + 'openclaude.openai.nonReplayableRequest', +) + +export function markOpenAIRequestNonReplayable(error: T): T { + Object.defineProperty(error, NON_REPLAYABLE_OPENAI_REQUEST, { + value: true, + configurable: false, + enumerable: false, + writable: false, + }) + return error +} + +export function isOpenAIRequestNonReplayable(error: unknown): boolean { + return ( + typeof error === 'object' && + error !== null && + Reflect.get(error, NON_REPLAYABLE_OPENAI_REQUEST) === true + ) +} + const OPENAI_CATEGORY_MARKER_PREFIX = '[openai_category=' const LOCALHOST_HOSTNAMES = new Set(['localhost', '127.0.0.1', '::1']) diff --git a/src/services/api/openaiShim.test.ts b/src/services/api/openaiShim.test.ts index a04c8fb14c..bfc8e06973 100644 --- a/src/services/api/openaiShim.test.ts +++ b/src/services/api/openaiShim.test.ts @@ -1,5 +1,6 @@ import { APIError } from '@anthropic-ai/sdk' import { afterEach, beforeEach, expect, mock, test } from 'bun:test' +import { getEventListeners } from 'node:events' import { acquireSharedMutationLock, releaseSharedMutationLock } from '../../test/sharedMutationLock.js' import { asMockFetch } from '../../test/typedMocks.js' import { _clearRegistryForTesting, ensureIntegrationsLoaded, registerGateway } from '../../integrations/index.ts' @@ -9,6 +10,10 @@ import { getAssistantMessageFromError, OPENCODE_GO_FREE_LIMIT_ERROR_MESSAGE, } from './errors.ts' +import { + extractOpenAICategoryMarker, + isOpenAIRequestNonReplayable, +} from './openaiErrorClassification.ts' import { createOpenAIShimClient, hasMistralApiHost } from './openaiShim.ts' import * as realCodexShim from './codexShim.js' import * as realGithubModelsCredentials from '../../utils/githubModelsCredentials.js' @@ -55,6 +60,7 @@ const originalEnv = { CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED: process.env.CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED, CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED_ID: process.env.CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED_ID, CLAUDE_STREAM_IDLE_TIMEOUT_MS: process.env.CLAUDE_STREAM_IDLE_TIMEOUT_MS, + API_TIMEOUT_MS: process.env.API_TIMEOUT_MS, } const originalFetch = globalThis.fetch @@ -363,6 +369,7 @@ function importFreshOpenAIShim( type StreamIdleTestApi = { StreamIdleTimeoutError: new (timeoutMs: number) => Error + getApiTimeoutMs: () => number getStreamIdleTimeoutMs: () => number readWithIdleTimeout: ( reader: ReadableStreamDefaultReader, @@ -375,6 +382,7 @@ async function getStreamIdleTestApi(cacheKey: string): Promise expect(typeof testApi.StreamIdleTimeoutError).toBe('function') + expect(typeof testApi.getApiTimeoutMs).toBe('function') expect(typeof testApi.getStreamIdleTimeoutMs).toBe('function') expect(typeof testApi.readWithIdleTimeout).toBe('function') return testApi as StreamIdleTestApi @@ -403,6 +411,48 @@ function makeChatCompletionResponse(model: string): Response { ) } +function makeGithubChatFallbackResponse(): Response { + return new Response( + JSON.stringify({ + error: { message: '/chat/completions is not accessible for this model' }, + }), + { status: 400, headers: { 'Content-Type': 'application/json' } }, + ) +} + +function makeResponsesApiResponse(model: string): Response { + return new Response( + JSON.stringify({ + id: 'resp-fallback-test', + model, + output: [ + { + type: 'message', + role: 'assistant', + content: [{ type: 'output_text', text: 'ok' }], + }, + ], + usage: { input_tokens: 1, output_tokens: 1, total_tokens: 2 }, + }), + { headers: { 'Content-Type': 'application/json' } }, + ) +} + +function pendingFetchUntilAbort( + init: RequestInit | undefined, +): Promise { + return new Promise((_resolve, reject) => { + const signal = init?.signal + if (!signal) return + + const rejectFromAbort = () => { + reject(signal.reason ?? new DOMException('Aborted', 'AbortError')) + } + signal.addEventListener('abort', rejectFromAbort, { once: true }) + if (signal.aborted) rejectFromAbort() + }) +} + async function captureChatCompletionRequest( model = 'mimo-v2.5-pro', ): Promise<{ authorization: string | null; url: string | null }> { @@ -470,6 +520,7 @@ beforeEach(async () => { delete process.env.CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED delete process.env.CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED_ID delete process.env.CLAUDE_STREAM_IDLE_TIMEOUT_MS + delete process.env.API_TIMEOUT_MS }) afterEach(() => { @@ -513,6 +564,7 @@ afterEach(() => { restoreEnv('CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED', originalEnv.CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED) restoreEnv('CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED_ID', originalEnv.CLAUDE_CODE_PROVIDER_PROFILE_ENV_APPLIED_ID) restoreEnv('CLAUDE_STREAM_IDLE_TIMEOUT_MS', originalEnv.CLAUDE_STREAM_IDLE_TIMEOUT_MS) + restoreEnv('API_TIMEOUT_MS', originalEnv.API_TIMEOUT_MS) globalThis.fetch = originalFetch _clearRegistryForTesting() ensureIntegrationsLoaded() @@ -1582,6 +1634,27 @@ test('stream idle timeout env parser parses and bounds overrides', async () => { expect(testApi.getStreamIdleTimeoutMs()).toBe(90_000) }) +test('API timeout env parser accepts safe positive integers and falls back otherwise', async () => { + const testApi = await getStreamIdleTestApi('api-timeout-env-parser') + + delete process.env.API_TIMEOUT_MS + expect(testApi.getApiTimeoutMs()).toBe(600_000) + + process.env.API_TIMEOUT_MS = '50' + expect(testApi.getApiTimeoutMs()).toBe(50) + + process.env.API_TIMEOUT_MS = ' 50 ' + expect(testApi.getApiTimeoutMs()).toBe(50) + + process.env.API_TIMEOUT_MS = '3000000000' + expect(testApi.getApiTimeoutMs()).toBe(2_147_483_647) + + for (const invalid of ['abc', '-5', '', '0', '1.5', '9007199254740993']) { + process.env.API_TIMEOUT_MS = invalid + expect(testApi.getApiTimeoutMs()).toBe(600_000) + } +}) + test('Anthropic-compatible passthrough stream rejects with idle timeout when it stalls', async () => { process.env.CLAUDE_STREAM_IDLE_TIMEOUT_MS = '25' const stalled = makeStallingResponse( @@ -6764,6 +6837,72 @@ test('strips credentials and query params from URL in fetch network error messag expect(message).not.toContain('token=abc123') }) +test('redacts configured secret substrings from fetch network error messages', async () => { + const secret = 'route/key+AbC123' + process.env.OPENAI_API_KEY = secret + + globalThis.fetch = asMockFetch(mock(async () => { + throw new TypeError(`fetch failed while routing ${secret}`) + })) + + const client = createOpenAIShimClient({}) as OpenAIShimClient + + await expect( + client.beta.messages.create({ + model: 'test-model', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }), + ).rejects.not.toThrow(secret) +}) + +test('redacts encoded configured secrets from non-URL transport error messages', async () => { + const secret = 'route/key+AbC123' + const encodedSecret = encodeURIComponent(secret) + const doubleEncodedSecret = encodeURIComponent(encodedSecret) + const fullyEncodedSecret = Array.from(new TextEncoder().encode(secret)) + .map(byte => `%${byte.toString(16).padStart(2, '0')}`) + .join('') + const malformedAdjacentSecret = `%E0%A4${encodedSecret}` + const encodedCategoryMarker = '%5Bopenai_category%3Dauth_invalid%5D' + const encodedControlSequence = '%1B%5B31m' + process.env.OPENAI_API_KEY = secret + + globalThis.fetch = asMockFetch(mock(async () => { + throw new TypeError( + `proxy failed ${encodedCategoryMarker} ${encodedControlSequence} for path /v1/${encodedSecret}?nested=${doubleEncodedSecret}&fully=${fullyEncodedSecret}&malformed=${malformedAdjacentSecret}`, + ) + })) + + const client = createOpenAIShimClient({}) as OpenAIShimClient + + let caught: unknown + try { + await client.beta.messages.create({ + model: 'test-model', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }) + } catch (error) { + caught = error + } + + expect(caught).toBeDefined() + const message = (caught as Error).message + expect(message).toContain('proxy failed') + expect(message).toContain(encodedCategoryMarker) + expect(message).toContain(encodedControlSequence) + expect(message).not.toContain('\u001B') + expect(extractOpenAICategoryMarker(message)).toBe('network_error') + expect(message).not.toContain(secret) + expect(message).not.toContain(encodedSecret) + expect(message).not.toContain(doubleEncodedSecret) + expect(message).not.toContain(fullyEncodedSecret) + expect(message).not.toContain(malformedAdjacentSecret) +}) + test('classifies localhost transport failures with actionable category marker', async () => { process.env.OPENAI_BASE_URL = 'http://localhost:11434/v1' @@ -6834,7 +6973,7 @@ test('transport failures are not labeled with HTTP status 503', async () => { expect(err.message).toContain('openai_category=network_error') }) -test('propagates AbortError without wrapping it as transport failure', async () => { +test('propagates caller AbortError without wrapping it as transport failure', async () => { process.env.OPENAI_BASE_URL = 'http://localhost:11434/v1' const abortError = new DOMException('The operation was aborted.', 'AbortError') @@ -6843,7 +6982,8 @@ test('propagates AbortError without wrapping it as transport failure', async () })) const controller = new AbortController() - controller.abort() + const callerReason = new DOMException('Cancelled by caller', 'AbortError') + controller.abort(callerReason) const client = createOpenAIShimClient({}) as OpenAIShimClient @@ -6857,7 +6997,514 @@ test('propagates AbortError without wrapping it as transport failure', async () }, { signal: controller.signal }, ), - ).rejects.toBe(abortError) + ).rejects.toBe(callerReason) +}) + +test('classifies a pre-header API timeout without replaying the request', async () => { + process.env.API_TIMEOUT_MS = '20' + const pathSecret = 'route/key+AbC123' + const encodedPathSecret = encodeURIComponent(pathSecret) + const doubleEncodedPathSecret = encodeURIComponent(encodedPathSecret) + const escapeLookingSecret = 'abc%2FdefLONG' + const encodedEscapeLookingSecret = encodeURIComponent(escapeLookingSecret) + const decodedEscapeLookingSecret = decodeURIComponent(escapeLookingSecret) + const malformedUtf8Secret = 'éSECRET_VALUE_123' + const encodedMalformedUtf8Secret = encodeURIComponent(malformedUtf8Secret) + process.env.OPENAI_API_KEY = pathSecret + process.env.OPENROUTER_API_KEY = escapeLookingSecret + process.env.DEEPSEEK_API_KEY = malformedUtf8Secret + process.env.OPENAI_BASE_URL = + `https://user:password@slow.example.test/v1/invalid%ZZ/${doubleEncodedPathSecret}/${encodedEscapeLookingSecret}/%E0%A4${encodedMalformedUtf8Secret}` + + `?prompt=${encodedPathSecret}&nested=${doubleEncodedPathSecret}` + + `&escape=${encodedEscapeLookingSecret}&malformed=%E0%A4${encodedMalformedUtf8Secret}` + + '&token=secret' + let fetchCalls = 0 + let completedGenerations = 0 + const receivedBodies: string[] = [] + globalThis.fetch = (async (_input, init) => { + fetchCalls++ + receivedBodies.push(String(init?.body)) + return new Promise((resolve, reject) => { + setTimeout(() => { + completedGenerations++ + resolve(makeChatCompletionResponse('gpt-4o-mini')) + }, 50) + const signal = init?.signal + if (!signal) return + const rejectFromAbort = () => { + reject(signal.reason ?? new DOMException('Aborted', 'AbortError')) + } + signal.addEventListener('abort', rejectFromAbort, { once: true }) + if (signal.aborted) rejectFromAbort() + }) + }) as unknown as FetchType + + const safety = new AbortController() + const safetyTimer = setTimeout(() => safety.abort(), 500) + const client = createOpenAIShimClient({}) as OpenAIShimClient + + let caught: unknown + try { + await waitForPromise( + client.beta.messages.create( + { + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }, + { signal: safety.signal }, + ), + 750, + 'pre-header timeout did not settle', + ) + } catch (error) { + caught = error + } finally { + clearTimeout(safetyTimer) + } + await new Promise(resolve => setTimeout(resolve, 60)) + + expect(caught).toBeDefined() + const error = caught as Error & { constructor: { name: string } } + expect(error.constructor.name).toBe('APIConnectionError') + expect(isOpenAIRequestNonReplayable(error)).toBe(true) + expect(error.message).toContain('no response headers within 20ms (API_TIMEOUT_MS)') + expect(error.message).toContain('slow.example.test') + expect(error.message).toContain('openai_category=request_timeout') + expect(error.message).not.toContain('password') + expect(error.message).not.toContain('token=secret') + expect(error.message).not.toContain(pathSecret) + expect(error.message).not.toContain(encodedPathSecret) + expect(error.message).not.toContain(doubleEncodedPathSecret) + expect(error.message).not.toContain(escapeLookingSecret) + expect(error.message).not.toContain(encodedEscapeLookingSecret) + expect(error.message).not.toContain(decodedEscapeLookingSecret) + expect(error.message).not.toContain(malformedUtf8Secret) + expect(error.message).not.toContain(encodedMalformedUtf8Secret) + expect(fetchCalls).toBe(1) + expect(receivedBodies).toHaveLength(1) + expect(completedGenerations).toBe(1) +}) + +test('does not proxy-retry when a deadline abort surfaces as fetch failed', async () => { + process.env.API_TIMEOUT_MS = '20' + let fetchCalls = 0 + globalThis.fetch = (async (_input, init) => { + fetchCalls++ + return new Promise((_resolve, reject) => { + const signal = init?.signal + const rejectFromAbort = () => reject(new TypeError('fetch failed')) + signal?.addEventListener('abort', rejectFromAbort, { once: true }) + if (signal?.aborted) rejectFromAbort() + }) + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + + await expect( + client.beta.messages.create({ + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }), + ).rejects.toThrow('no response headers within 20ms (API_TIMEOUT_MS)') + + expect(fetchCalls).toBe(1) +}) + +test('deadline wins when an abort-ignoring fetch resolves 504 afterward', async () => { + process.env.API_TIMEOUT_MS = '20' + let fetchCalls = 0 + globalThis.fetch = (async (_input, init) => { + fetchCalls++ + return new Promise(resolve => { + const resolveAfterAbort = () => { + resolve(new Response('Gateway Timeout', { status: 504 })) + } + init?.signal?.addEventListener('abort', resolveAfterAbort, { once: true }) + if (init?.signal?.aborted) resolveAfterAbort() + }) + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + + await expect( + client.beta.messages.create({ + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }), + ).rejects.toThrow('no response headers within 20ms (API_TIMEOUT_MS)') + + expect(fetchCalls).toBe(1) +}) + +test('gives a proxy retry its own response-header deadline', async () => { + process.env.API_TIMEOUT_MS = '50' + let fetchCalls = 0 + globalThis.fetch = (async () => { + fetchCalls++ + if (fetchCalls === 1) { + await new Promise(resolve => setTimeout(resolve, 30)) + return new Response('Gateway Timeout', { status: 504 }) + } + await new Promise(resolve => setTimeout(resolve, 30)) + return makeChatCompletionResponse('gpt-4o-mini') + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + const response = await client.beta.messages.create({ + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }) + + expect(response).toBeDefined() + expect(fetchCalls).toBe(2) +}) + +test('bounds nested URL decoding while retaining encoded secret redaction', async () => { + const secret = 'route/key+AbC123' + const secretVariants = [secret] + for (let layer = 0; layer < 4; layer++) { + secretVariants.push( + encodeURIComponent(secretVariants[secretVariants.length - 1]!), + ) + } + const deeplyNestedValue = `%${'25'.repeat(200)}` + process.env.OPENAI_API_KEY = secret + process.env.OPENAI_BASE_URL = + `https://slow.example.test/v1/${deeplyNestedValue}/${secretVariants[4]}` + globalThis.fetch = asMockFetch(mock(async () => { + throw new TypeError('fetch failed') + })) + + const originalDecodeURIComponent = globalThis.decodeURIComponent + let decodeCalls = 0 + globalThis.decodeURIComponent = ((value: string) => { + decodeCalls++ + return originalDecodeURIComponent(value) + }) as typeof decodeURIComponent + + const client = createOpenAIShimClient({}) as OpenAIShimClient + let caught: unknown + try { + await client.beta.messages.create({ + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }) + } catch (error) { + caught = error + } finally { + globalThis.decodeURIComponent = originalDecodeURIComponent + } + + expect(caught).toBeDefined() + for (const secretVariant of secretVariants) { + expect((caught as Error).message).not.toContain(secretVariant) + } + expect(decodeCalls).toBeLessThanOrEqual(32) +}) + +test('decodes malformed URL escape runs in linear work', async () => { + const secret = 'route/key+AbC123' + const encodedSecret = encodeURIComponent(secret) + const malformedEscapeCount = 200 + const malformedUtf8Run = '%80'.repeat(malformedEscapeCount) + process.env.OPENAI_API_KEY = secret + process.env.OPENAI_BASE_URL = + `https://slow.example.test/v1/${malformedUtf8Run}/${encodedSecret}` + globalThis.fetch = asMockFetch(mock(async () => { + throw new TypeError('fetch failed') + })) + + const originalDecodeURIComponent = globalThis.decodeURIComponent + let malformedDecodeCalls = 0 + globalThis.decodeURIComponent = ((value: string) => { + if (value.includes('%80')) malformedDecodeCalls++ + return originalDecodeURIComponent(value) + }) as typeof decodeURIComponent + + const client = createOpenAIShimClient({}) as OpenAIShimClient + let caught: unknown + try { + await client.beta.messages.create({ + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }) + } catch (error) { + caught = error + } finally { + globalThis.decodeURIComponent = originalDecodeURIComponent + } + + expect(caught).toBeDefined() + expect((caught as Error).message).not.toContain(secret) + expect((caught as Error).message).not.toContain(encodedSecret) + expect(malformedDecodeCalls).toBeLessThanOrEqual(malformedEscapeCount * 10) +}) + +test('preserves caller cancellation while waiting for response headers without retrying', async () => { + process.env.API_TIMEOUT_MS = '200' + let fetchCalls = 0 + globalThis.fetch = (async (_input, init) => { + fetchCalls++ + return pendingFetchUntilAbort(init) + }) as unknown as FetchType + + const caller = new AbortController() + const callerReason = new DOMException('Cancelled by user', 'AbortError') + const client = createOpenAIShimClient({}) as OpenAIShimClient + const originalAbortSignalAny = Object.getOwnPropertyDescriptor( + AbortSignal, + 'any', + ) + Object.defineProperty(AbortSignal, 'any', { + value: undefined, + configurable: true, + }) + try { + const request = client.beta.messages.create( + { + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }, + { signal: caller.signal }, + ) + + setTimeout(() => caller.abort(callerReason), 10) + + await expect( + waitForPromise(request, 500, 'caller abort did not settle'), + ).rejects.toBe(callerReason) + expect(fetchCalls).toBe(1) + } finally { + if (originalAbortSignalAny) { + Object.defineProperty(AbortSignal, 'any', originalAbortSignalAny) + } + } +}) + +test('native signal composition preserves the caller abort reason when fetch rejects AbortError', async () => { + expect(typeof AbortSignal.any).toBe('function') + process.env.API_TIMEOUT_MS = '200' + let fetchCalls = 0 + const fetchAbortError = new DOMException( + 'The operation was aborted.', + 'AbortError', + ) + globalThis.fetch = (async (_input, init) => { + fetchCalls++ + return new Promise((_resolve, reject) => { + const signal = init?.signal + if (!signal) return + const rejectFromAbort = () => reject(fetchAbortError) + signal.addEventListener('abort', rejectFromAbort, { once: true }) + if (signal.aborted) rejectFromAbort() + }) + }) as unknown as FetchType + + const caller = new AbortController() + const callerReason = new DOMException('Cancelled by user', 'AbortError') + const client = createOpenAIShimClient({}) as OpenAIShimClient + const request = client.beta.messages.create( + { + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }, + { signal: caller.signal }, + ) + + const abortTimer = setTimeout(() => caller.abort(callerReason), 10) + + try { + await expect( + waitForPromise(request, 500, 'caller abort did not settle'), + ).rejects.toBe(callerReason) + } finally { + clearTimeout(abortTimer) + } + expect(fetchCalls).toBe(1) +}) + +test('caller abort winning the timeout catch race prevents a retry', async () => { + process.env.API_TIMEOUT_MS = '20' + let fetchCalls = 0 + const caller = new AbortController() + const callerReason = new DOMException('Cancelled by user', 'AbortError') + globalThis.fetch = (async (_input, init) => { + fetchCalls++ + return new Promise((_resolve, reject) => { + const signal = init?.signal + if (!signal) return + const rejectFromAbort = () => { + reject(signal.reason ?? new DOMException('Aborted', 'AbortError')) + if (fetchCalls === 1) { + queueMicrotask(() => caller.abort(callerReason)) + } + } + signal.addEventListener('abort', rejectFromAbort, { once: true }) + if (signal.aborted) rejectFromAbort() + }) + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + await expect( + waitForPromise( + client.beta.messages.create( + { + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }, + { signal: caller.signal }, + ), + 500, + 'caller abort race did not settle', + ), + ).rejects.toBe(callerReason) + expect(fetchCalls).toBe(1) +}) + +test('manual signal fallback preserves caller cancellation after headers arrive', async () => { + process.env.API_TIMEOUT_MS = '200' + const fetchSignals: AbortSignal[] = [] + const stalled = makeStallingResponse( + makeOpenAIStreamFrame({ role: 'assistant', content: 'started' }), + ) + globalThis.fetch = (async (_input, init) => { + if (init?.signal) fetchSignals.push(init.signal) + return stalled.response + }) as unknown as FetchType + + const caller = new AbortController() + const client = createOpenAIShimClient({}) as OpenAIShimClient + const originalAbortSignalAny = Object.getOwnPropertyDescriptor( + AbortSignal, + 'any', + ) + Object.defineProperty(AbortSignal, 'any', { + value: undefined, + configurable: true, + }) + try { + const result = await client.beta.messages + .create( + { + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: true, + }, + { signal: caller.signal }, + ) + .withResponse() + + await expectAbortStopsStream({ + abort: () => caller.abort(), + cancelReasons: stalled.cancelReasons, + expectedEventsBeforeAbort: 1, + label: 'manual combined signal after headers', + stream: result.data as ShimStream, + }) + + expect(fetchSignals).toHaveLength(1) + expect(fetchSignals[0].aborted).toBe(true) + } finally { + stalled.close() + if (originalAbortSignalAny) { + Object.defineProperty(AbortSignal, 'any', originalAbortSignalAny) + } + } +}) + +test('manual signal fallback removes caller forwarding after the body settles', async () => { + process.env.API_TIMEOUT_MS = '200' + globalThis.fetch = asMockFetch(mock(async () => + makeChatCompletionResponse('gpt-4o-mini'))) + + const caller = new AbortController() + const client = createOpenAIShimClient({}) as OpenAIShimClient + const originalAbortSignalAny = Object.getOwnPropertyDescriptor( + AbortSignal, + 'any', + ) + Object.defineProperty(AbortSignal, 'any', { + value: undefined, + configurable: true, + }) + try { + for (let request = 0; request < 2; request++) { + await client.beta.messages.create( + { + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: false, + }, + { signal: caller.signal }, + ) + expect(getEventListeners(caller.signal, 'abort')).toHaveLength(0) + } + } finally { + if (originalAbortSignalAny) { + Object.defineProperty(AbortSignal, 'any', originalAbortSignalAny) + } + } +}) + +test('disarms the API timeout after headers arrive while the body keeps streaming', async () => { + process.env.API_TIMEOUT_MS = '20' + const fetchSignals: AbortSignal[] = [] + const encoder = new TextEncoder() + globalThis.fetch = (async (_input, init) => { + if (init?.signal) fetchSignals.push(init.signal) + return new Response( + new ReadableStream({ + start(controller) { + setTimeout(() => { + controller.enqueue(encoder.encode(makeOpenAIStreamFrame( + { role: 'assistant', content: 'late body' }, + 'stop', + ))) + controller.enqueue(encoder.encode('data: [DONE]\n\n')) + controller.close() + }, 50) + }, + }), + { headers: { 'Content-Type': 'text/event-stream' } }, + ) + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + const result = await client.beta.messages + .create({ + model: 'gpt-4o-mini', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 64, + stream: true, + }) + .withResponse() + + const events: Array> = [] + for await (const event of result.data) events.push(event) + + expect(events.length).toBeGreaterThan(0) + expect(fetchSignals).toHaveLength(1) + expect(fetchSignals[0].aborted).toBe(false) }) test('classifies chat-completions endpoint 404 failures with endpoint_not_found marker', async () => { @@ -9669,6 +10316,195 @@ function makeCodexSseResponse(responseData: Record): Response { return makeSseResponse([`event: response.completed\ndata: ${data}\n\n`]) } +test('GitHub Copilot codex responses transport does not replay after a pre-header timeout', async () => { + process.env.CLAUDE_CODE_USE_GITHUB = '1' + process.env.OPENAI_BASE_URL = 'https://api.githubcopilot.com' + process.env.OPENAI_API_KEY = 'test-token' + process.env.API_TIMEOUT_MS = '20' + let fetchCalls = 0 + const requestUrls: string[] = [] + + globalThis.fetch = (async (input, init) => { + fetchCalls++ + requestUrls.push(String(input)) + return pendingFetchUntilAbort(init) + }) as unknown as FetchType + + const safety = new AbortController() + const safetyTimer = setTimeout(() => safety.abort(), 500) + const client = createOpenAIShimClient({}) as OpenAIShimClient + let caught: unknown + try { + await waitForPromise( + client.beta.messages.create( + { + model: 'gpt-5', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 32, + stream: false, + }, + { signal: safety.signal }, + ), + 750, + 'GitHub codex responses timeout did not settle', + ) + } catch (error) { + caught = error + } finally { + clearTimeout(safetyTimer) + } + + expect(caught).toBeDefined() + const error = caught as Error & { constructor: { name: string } } + expect(error.constructor.name).toBe('APIConnectionError') + expect(isOpenAIRequestNonReplayable(error)).toBe(true) + expect(fetchCalls).toBe(1) + expect(requestUrls).toEqual([ + 'https://api.githubcopilot.com/responses', + ]) +}) + +test('GitHub Copilot responses fallback does not replay after a pre-header timeout', async () => { + process.env.CLAUDE_CODE_USE_GITHUB = '1' + process.env.OPENAI_BASE_URL = 'https://api.githubcopilot.com' + process.env.OPENAI_API_KEY = 'test-token' + process.env.API_TIMEOUT_MS = '20' + const requestUrls: string[] = [] + + globalThis.fetch = (async (input, init) => { + const url = String(input) + requestUrls.push(url) + if (url.endsWith('/chat/completions')) { + return makeGithubChatFallbackResponse() + } + return pendingFetchUntilAbort(init) + }) as unknown as FetchType + + const safety = new AbortController() + const safetyTimer = setTimeout(() => safety.abort(), 500) + const client = createOpenAIShimClient({}) as OpenAIShimClient + let caught: unknown + try { + await waitForPromise( + client.beta.messages.create( + { + model: 'gpt-4', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 32, + stream: false, + }, + { signal: safety.signal }, + ), + 750, + 'GitHub responses fallback timeout did not settle', + ) + } catch (error) { + caught = error + } finally { + clearTimeout(safetyTimer) + } + + expect(caught).toBeDefined() + const error = caught as Error & { constructor: { name: string } } + expect(error.constructor.name).toBe('APIConnectionError') + expect(isOpenAIRequestNonReplayable(error)).toBe(true) + expect(requestUrls).toEqual([ + 'https://api.githubcopilot.com/chat/completions', + 'https://api.githubcopilot.com/responses', + ]) +}) + +test('GitHub Copilot responses fallback preserves caller abort without retrying', async () => { + process.env.CLAUDE_CODE_USE_GITHUB = '1' + process.env.OPENAI_BASE_URL = 'https://api.githubcopilot.com' + process.env.OPENAI_API_KEY = 'test-token' + process.env.API_TIMEOUT_MS = '200' + let fetchCalls = 0 + + globalThis.fetch = (async (input, init) => { + fetchCalls++ + if (String(input).endsWith('/chat/completions')) { + return makeGithubChatFallbackResponse() + } + return pendingFetchUntilAbort(init) + }) as unknown as FetchType + + const caller = new AbortController() + const callerReason = new DOMException('Cancelled by user', 'AbortError') + const client = createOpenAIShimClient({}) as OpenAIShimClient + const request = client.beta.messages.create( + { + model: 'gpt-4', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 32, + stream: false, + }, + { signal: caller.signal }, + ) + setTimeout(() => caller.abort(callerReason), 10) + + await expect( + waitForPromise(request, 500, 'GitHub responses fallback abort did not settle'), + ).rejects.toBe(callerReason) + expect(fetchCalls).toBe(2) +}) + +test('GitHub Copilot responses fallback preserves non-caller transport aborts', async () => { + process.env.CLAUDE_CODE_USE_GITHUB = '1' + process.env.OPENAI_BASE_URL = 'https://api.githubcopilot.com' + process.env.OPENAI_API_KEY = 'test-token' + let fetchCalls = 0 + const transportAbort = new DOMException('Proxy aborted request', 'AbortError') + + globalThis.fetch = (async input => { + fetchCalls++ + if (String(input).endsWith('/chat/completions')) { + return makeGithubChatFallbackResponse() + } + throw transportAbort + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + await expect( + client.beta.messages.create({ + model: 'gpt-4', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 32, + stream: false, + }), + ).rejects.toBe(transportAbort) + expect(fetchCalls).toBe(2) +}) + +test('GitHub Copilot responses fallback does not retry non-retryable HTTP failures', async () => { + process.env.CLAUDE_CODE_USE_GITHUB = '1' + process.env.OPENAI_BASE_URL = 'https://api.githubcopilot.com' + process.env.OPENAI_API_KEY = 'test-token' + let fetchCalls = 0 + + globalThis.fetch = (async input => { + fetchCalls++ + if (String(input).endsWith('/chat/completions')) { + return makeGithubChatFallbackResponse() + } + return new Response( + JSON.stringify({ error: { message: 'invalid token' } }), + { status: 401, headers: { 'Content-Type': 'application/json' } }, + ) + }) as unknown as FetchType + + const client = createOpenAIShimClient({}) as OpenAIShimClient + await expect( + client.beta.messages.create({ + model: 'gpt-4', + messages: [{ role: 'user', content: 'hello' }], + max_tokens: 32, + stream: false, + }), + ).rejects.toMatchObject({ status: 401 }) + expect(fetchCalls).toBe(2) +}) + test('GitHub Copilot 401 chat_completions retries with refreshed token', async () => { const realModule = realGithubModelsCredentials try { diff --git a/src/services/api/openaiShim.ts b/src/services/api/openaiShim.ts index 03d821f775..d319d05ced 100644 --- a/src/services/api/openaiShim.ts +++ b/src/services/api/openaiShim.ts @@ -77,7 +77,10 @@ import { } from './codexShim.js' import { buildAnthropicUsageFromRawUsage } from './cacheMetrics.js' import { compressToolHistory } from './compressToolHistory.js' -import { fetchWithProxyRetry } from './fetchWithProxyRetry.js' +import { + fetchWithProxyRetry, + type ProxyRetryFetcher, +} from './fetchWithProxyRetry.js' import { getLocalFastPathConfig, getLocalProviderRetryBaseUrls, @@ -94,9 +97,14 @@ import { buildOpenAICompatibilityErrorMessage, classifyOpenAIHttpFailure, classifyOpenAINetworkFailure, + markOpenAIRequestNonReplayable, } from './openaiErrorClassification.js' import { sanitizeSchemaForOpenAICompat } from '../../utils/schemaSanitizer.js' import { redactSecretValueForDisplay, type SecretValueSource } from '../../utils/providerProfile.js' +import { + redactEncodedSecretSubstringsForDisplay, + redactSecretSubstringsForDisplay, +} from '../../utils/providerSecrets.js' import { redactUrlForDisplay, shouldRedactUrlQueryParam, @@ -125,6 +133,7 @@ const GITHUB_429_MAX_RETRIES = 3 const GITHUB_429_BASE_DELAY_SEC = 1 const GITHUB_429_MAX_DELAY_SEC = 32 const CREDENTIAL_POOL_COOLDOWN_MS = 30_000 +const DEFAULT_API_TIMEOUT_MS = 600_000 const DEFAULT_STREAM_IDLE_TIMEOUT_MS = 90_000 const MAX_STREAM_IDLE_TIMEOUT_MS = 2_147_483_647 const GEMINI_API_HOST = 'generativelanguage.googleapis.com' @@ -147,6 +156,36 @@ class StreamIdleTimeoutError extends Error { } } +class ResponseHeadersTimeoutError extends Error { + constructor(timeoutMs: number, url: string) { + super( + `OpenAI-compatible request received no response headers within ${timeoutMs}ms (API_TIMEOUT_MS) from ${url}`, + ) + this.name = 'ResponseHeadersTimeoutError' + } +} + +function preserveCallerAbortError( + error: unknown, + callerSignal: AbortSignal, +): unknown { + return error instanceof ResponseHeadersTimeoutError || isAbortError(error) + ? callerSignal.reason ?? error + : error +} + +function isAbortError(error: unknown): boolean { + return ( + (typeof DOMException !== 'undefined' && + error instanceof DOMException && + error.name === 'AbortError') || + (typeof error === 'object' && + error !== null && + 'name' in error && + error.name === 'AbortError') + ) +} + function createStreamAbortError(): DOMException { return new DOMException('Aborted', 'AbortError') } @@ -196,6 +235,201 @@ export function getStreamIdleTimeoutMs(): number { : DEFAULT_STREAM_IDLE_TIMEOUT_MS } +export function getApiTimeoutMs(): number { + const raw = process.env.API_TIMEOUT_MS?.trim() + if (!raw || !/^\d+$/.test(raw)) return DEFAULT_API_TIMEOUT_MS + const parsed = Number(raw) + return Number.isSafeInteger(parsed) && parsed > 0 + ? Math.min(parsed, MAX_STREAM_IDLE_TIMEOUT_MS) + : DEFAULT_API_TIMEOUT_MS +} + +function combineRequestSignals( + callerSignal: AbortSignal | undefined, + deadlineSignal: AbortSignal, +): { + signal: AbortSignal + cleanupAfterHeaders: () => void + cleanup: () => void + cleanupAfterBody?: () => void +} { + if (!callerSignal) { + return { + signal: deadlineSignal, + cleanupAfterHeaders: () => {}, + cleanup: () => {}, + } + } + + if (typeof AbortSignal.any === 'function') { + return { + // The deadline controller is request-local and its timer is the only + // abort source, so clearing that timer after headers permanently disarms it. + signal: AbortSignal.any([callerSignal, deadlineSignal]), + cleanupAfterHeaders: () => {}, + cleanup: () => {}, + } + } + + const combined = new AbortController() + const abortFromCaller = () => { + deadlineSignal.removeEventListener('abort', abortFromDeadline) + combined.abort(callerSignal.reason) + } + const abortFromDeadline = () => { + callerSignal.removeEventListener('abort', abortFromCaller) + combined.abort(deadlineSignal.reason) + } + const cleanupAfterHeaders = () => { + deadlineSignal.removeEventListener('abort', abortFromDeadline) + } + const cleanup = () => { + callerSignal.removeEventListener('abort', abortFromCaller) + cleanupAfterHeaders() + } + + callerSignal.addEventListener('abort', abortFromCaller, { once: true }) + deadlineSignal.addEventListener('abort', abortFromDeadline, { once: true }) + if (callerSignal.aborted) { + abortFromCaller() + } else if (deadlineSignal.aborted) { + abortFromDeadline() + } + + return { + signal: combined.signal, + cleanupAfterHeaders, + cleanup, + cleanupAfterBody: cleanup, + } +} + +function wrapResponseBodyWithCleanup( + response: Response, + cleanup: () => void, +): Response { + if (!response.body) { + cleanup() + return response + } + + const reader = response.body.getReader() + let cleanedUp = false + const cleanupOnce = () => { + if (cleanedUp) return + cleanedUp = true + cleanup() + } + const body = new ReadableStream({ + async pull(controller) { + try { + const result = await reader.read() + if (result.done) { + cleanupOnce() + controller.close() + } else { + controller.enqueue(result.value) + } + } catch (error) { + cleanupOnce() + controller.error(error) + } + }, + async cancel(reason) { + try { + await reader.cancel(reason) + } finally { + cleanupOnce() + } + }, + }) + const wrapped = new Response(body, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }) + for (const property of ['url', 'type', 'redirected'] as const) { + try { + Object.defineProperty(wrapped, property, { + value: response[property], + configurable: true, + }) + } catch { + /* non-fatal: standard response metadata remains available */ + } + } + return wrapped +} + +async function fetchWithHeadersDeadline( + url: string, + init: RequestInit, + options: { + callerSignal?: AbortSignal + timeoutMs: number + }, +): Promise { + const redactedUrl = redactUrlForDiagnostics(url) + const fetchWithAttemptDeadline: ProxyRetryFetcher = async (input, attemptInit) => { + const deadlineController = new AbortController() + const timeoutReason = new ResponseHeadersTimeoutError( + options.timeoutMs, + redactedUrl, + ) + const { + signal, + cleanupAfterHeaders, + cleanup, + cleanupAfterBody, + } = combineRequestSignals(options.callerSignal, deadlineController.signal) + const timer = setTimeout( + () => deadlineController.abort(timeoutReason), + options.timeoutMs, + ) + timer.unref?.() + + let headersReceived = false + try { + const response = await fetch(input, { ...attemptInit, signal }) + if (signal.aborted) { + void response.body?.cancel().catch(() => {}) + throw ( + signal.reason ?? + new DOMException('The operation was aborted.', 'AbortError') + ) + } + headersReceived = true + return cleanupAfterBody + ? wrapResponseBodyWithCleanup(response, cleanupAfterBody) + : response + } catch (error) { + if (options.callerSignal?.aborted) { + throw preserveCallerAbortError(error, options.callerSignal) + } + if ( + deadlineController.signal.aborted && + deadlineController.signal.reason === timeoutReason + ) { + throw timeoutReason + } + throw error + } finally { + clearTimeout(timer) + if (headersReceived) { + cleanupAfterHeaders() + } else { + cleanup() + } + } + } + + return fetchWithProxyRetry( + url, + { ...init, signal: options.callerSignal }, + { fetcher: fetchWithAttemptDeadline }, + ) +} + async function readWithIdleTimeout( reader: ReadableStreamDefaultReader, timeoutMs: number, @@ -410,11 +644,96 @@ function formatRetryAfterHint(response: Response): string { return ra ? ` (Retry-After: ${ra})` : '' } +function decodeValidPercentRun(encoded: string): string { + const escapes = encoded.match(/%[0-9A-Fa-f]{2}/g) + if (!escapes) return encoded + + let decoded = '' + let offset = 0 + while (offset < escapes.length) { + const firstByte = Number.parseInt(escapes[offset].slice(1), 16) + const sequenceLength = + firstByte <= 0x7f + ? 1 + : firstByte >= 0xc2 && firstByte <= 0xdf + ? 2 + : firstByte >= 0xe0 && firstByte <= 0xef + ? 3 + : firstByte >= 0xf0 && firstByte <= 0xf4 + ? 4 + : 1 + try { + decoded += decodeURIComponent( + escapes.slice(offset, offset + sequenceLength).join(''), + ) + offset += sequenceLength + } catch { + decoded += escapes[offset] + offset++ + } + } + return decoded +} + +function decodeValidUrlEscapesOnce(value: string): string { + return value.replace(/(?:%[0-9A-Fa-f]{2})+/g, decodeValidPercentRun) +} + +const MAX_URL_SECRET_DECODING_LAYERS = 4 + +function redactDecodedUrlComponentSecrets(value: string): string { + let decoded = value + let foundSecret = false + for (let layer = 0; layer <= MAX_URL_SECRET_DECODING_LAYERS; layer++) { + const redacted = + redactSecretSubstringsForDisplay( + decoded, + process.env as SecretValueSource, + ) ?? decoded + if (redacted !== decoded) foundSecret = true + if (layer === MAX_URL_SECRET_DECODING_LAYERS) { + decoded = redacted + break + } + const next = decodeValidUrlEscapesOnce(redacted) + if (next === redacted) { + decoded = redacted + break + } + decoded = next + } + return foundSecret ? decoded : value +} + function redactUrlForDiagnostics(url: string): string { - const redacted = redactUrlForDisplay(url) + let redacted = redactUrlForDisplay(url) + try { + const parsed = new URL(redacted) + const redactedPathname = redactDecodedUrlComponentSecrets(parsed.pathname) + const redactedSearch = redactDecodedUrlComponentSecrets(parsed.search) + let componentRedacted = false + if (redactedPathname !== parsed.pathname) { + parsed.pathname = redactedPathname + componentRedacted = true + } + if (redactedSearch !== parsed.search) { + parsed.search = redactedSearch + componentRedacted = true + } + if (componentRedacted) redacted = parsed.toString() + } catch { + // Keep the URL-level redaction when the URL cannot be parsed. + } + const redactedSubstrings = + redactSecretSubstringsForDisplay( + redacted, + process.env as SecretValueSource, + ) ?? redacted return ( - redactSecretValueForDisplay(redacted, process.env as SecretValueSource) ?? - redacted + redactSecretValueForDisplay( + redactedSubstrings, + process.env as SecretValueSource, + ) ?? redactedSubstrings ) } @@ -422,6 +741,53 @@ function redactUrlsInMessage(message: string): string { return message.replace(/https?:\/\/\S+/g, match => redactUrlForDiagnostics(match)) } +function createClassifiedTransportError( + error: unknown, + requestUrl: string, + model: string, + preclassifiedFailure?: ReturnType, +) { + const failure = + preclassifiedFailure ?? + classifyOpenAINetworkFailure(error, { + url: requestUrl, + }) + const redactedUrl = redactUrlForDiagnostics(requestUrl) + const encodedSecretRedactedMessage = + redactEncodedSecretSubstringsForDisplay( + redactUrlsInMessage(failure.message), + process.env as SecretValueSource, + ) ?? 'Request failed' + const redactedMessage = + redactSecretSubstringsForDisplay( + encodedSecretRedactedMessage, + process.env as SecretValueSource, + ) ?? 'Request failed' + const safeMessage = + redactSecretValueForDisplay( + redactedMessage, + process.env as SecretValueSource, + ) || 'Request failed' + + logForDebugging( + `[OpenAIShim] transport failure category=${failure.category} retryable=${failure.retryable} code=${failure.code ?? 'unknown'} method=POST url=${redactedUrl} model=${model} message=${safeMessage}`, + { level: 'warn' }, + ) + + const apiError = APIError.generate( + 0, + undefined, + buildOpenAICompatibilityErrorMessage( + `OpenAI API transport error: ${safeMessage}${failure.code ? ` (code=${failure.code})` : ''}`, + failure, + ), + new Headers(), + ) + return failure.retryable + ? apiError + : markOpenAIRequestNonReplayable(apiError) +} + function sleepMs(ms: number): Promise { return new Promise(resolve => setTimeout(resolve, ms)) } @@ -3679,6 +4045,8 @@ class OpenAIShimMessages { const isGithubWithCodexTransport = isGithubCopilotEndpoint && request.transport === 'codex_responses' if (isGithubWithCodexTransport) { + const apiTimeoutMs = getApiTimeoutMs() + const responsesUrl = `${request.baseUrl}/responses` let didRefreshCopilotCodexToken = false let refreshedCopilotCodexToken: string | undefined for (let attempt = 0; attempt < 2; attempt++) { @@ -3690,20 +4058,53 @@ class OpenAIShimMessages { } try { - return await performCodexRequest({ - request, - credentials: { - apiKey, - source: 'env', - }, - params, - defaultHeaders: { - ...this.defaultHeaders, - ...filterAnthropicHeaders(options?.headers), - ...COPILOT_HEADERS, - }, - signal: options?.signal, - }) + try { + return await performCodexRequest({ + request, + credentials: { + apiKey, + source: 'env', + }, + params, + defaultHeaders: { + ...this.defaultHeaders, + ...filterAnthropicHeaders(options?.headers), + ...COPILOT_HEADERS, + }, + signal: options?.signal, + fetcher: (input, init) => { + const url = + typeof input === 'string' + ? input + : input instanceof URL + ? input.toString() + : input.url + return fetchWithHeadersDeadline(url, init ?? {}, { + callerSignal: options?.signal, + timeoutMs: apiTimeoutMs, + }) + }, + }) + } catch (error) { + if (options?.signal?.aborted) { + throw preserveCallerAbortError(error, options.signal) + } + if (error instanceof ResponseHeadersTimeoutError) { + const failure = { + ...classifyOpenAINetworkFailure(error, { + url: responsesUrl, + }), + retryable: false, + } + throw createClassifiedTransportError( + error, + responsesUrl, + request.resolvedModel, + failure, + ) + } + throw error + } } catch (error) { if ( !didRefreshCopilotCodexToken && @@ -3784,6 +4185,7 @@ class OpenAIShimMessages { params: ShimCreateParams, options?: { signal?: AbortSignal; headers?: Record }, ): Promise { + const apiTimeoutMs = getApiTimeoutMs() // Local backends (llama.cpp, vLLM, Ollama, LM Studio, …) do not implement // the cloud-side caching/strict-validation behaviours that several of our // pre-send transforms target. Computing the fast-path config once here @@ -4587,16 +4989,26 @@ class OpenAIShimMessages { method: 'POST' as const, headers, body: serializedBody, - signal: options?.signal, }) + const fetchAttemptWithHeadersDeadline = ( + url: string, + init: RequestInit, + ): Promise => + fetchWithHeadersDeadline(url, init, { + callerSignal: options?.signal, + timeoutMs: apiTimeoutMs, + }) + const maxSelfHealAttempts = isLocal ? localRetryBaseUrls.length + 1 : 0 const credentialPoolAttempts = credentialPool?.size ?? 1 - let maxAttempts = + let maxAttempts = Math.max( + 2, Math.max(isGithub ? GITHUB_429_MAX_RETRIES : 1, credentialPoolAttempts) + - maxSelfHealAttempts + maxSelfHealAttempts, + ) const throwClassifiedTransportError = ( error: unknown, @@ -4604,34 +5016,14 @@ class OpenAIShimMessages { preclassifiedFailure?: ReturnType, ): never => { if (options?.signal?.aborted) { - throw error + throw preserveCallerAbortError(error, options.signal) } - const failure = - preclassifiedFailure ?? - classifyOpenAINetworkFailure(error, { - url: requestUrl, - }) - const redactedUrl = redactUrlForDiagnostics(requestUrl) - const safeMessage = - redactSecretValueForDisplay( - redactUrlsInMessage(failure.message), - process.env as SecretValueSource, - ) || 'Request failed' - - logForDebugging( - `[OpenAIShim] transport failure category=${failure.category} retryable=${failure.retryable} code=${failure.code ?? 'unknown'} method=POST url=${redactedUrl} model=${request.resolvedModel} message=${safeMessage}`, - { level: 'warn' }, - ) - - throw APIError.generate( - 0, - undefined, - buildOpenAICompatibilityErrorMessage( - `OpenAI API transport error: ${safeMessage}${failure.code ? ` (code=${failure.code})` : ''}`, - failure, - ), - new Headers(), + throw createClassifiedTransportError( + error, + requestUrl, + request.resolvedModel, + preclassifiedFailure, ) } @@ -4698,38 +5090,40 @@ class OpenAIShimMessages { } const headers = await buildHeadersForAttempt(credentialLease) try { - response = await fetchWithProxyRetry( + response = await fetchAttemptWithHeadersDeadline( requestUrl, buildFetchInit(headers), ) } catch (error) { - const isAbortError = - options?.signal?.aborted === true || - (typeof DOMException !== 'undefined' && - error instanceof DOMException && - error.name === 'AbortError') || - (typeof error === 'object' && - error !== null && - 'name' in error && - error.name === 'AbortError') - - if (isAbortError) { + if (options?.signal?.aborted) { + throw preserveCallerAbortError(error, options.signal) + } + const isResponseHeadersTimeout = + error instanceof ResponseHeadersTimeoutError + if (!isResponseHeadersTimeout && isAbortError(error)) { throw error } - const failure = classifyOpenAINetworkFailure(error, { + const classifiedFailure = classifyOpenAINetworkFailure(error, { url: requestUrl, }) + if (isResponseHeadersTimeout) { + throwClassifiedTransportError(error, requestUrl, { + ...classifiedFailure, + retryable: false, + }) + } + if ( isLocal && - failure.category === 'localhost_resolution_failed' && + classifiedFailure.category === 'localhost_resolution_failed' && promoteNextLocalBaseUrl('localhost_resolution_failed') ) { continue } - throwClassifiedTransportError(error, requestUrl, failure) + throwClassifiedTransportError(error, requestUrl, classifiedFailure) } // After the try/catch, response is guaranteed to be defined — the catch @@ -4825,14 +5219,38 @@ class OpenAIShimMessages { let responsesResponse!: Response try { - responsesResponse = await fetchWithProxyRetry(responsesUrl, { - method: 'POST', - headers, - body: stableStringifyJson(responsesBody), - signal: options?.signal, - }) + responsesResponse = await fetchAttemptWithHeadersDeadline( + responsesUrl, + { + method: 'POST', + headers, + body: stableStringifyJson(responsesBody), + }, + ) } catch (error) { - throwClassifiedTransportError(error, responsesUrl) + if (options?.signal?.aborted) { + throw preserveCallerAbortError(error, options.signal) + } + if ( + !(error instanceof ResponseHeadersTimeoutError) && + isAbortError(error) + ) { + throw error + } + const classifiedFailure = classifyOpenAINetworkFailure(error, { + url: responsesUrl, + }) + if (error instanceof ResponseHeadersTimeoutError) { + throwClassifiedTransportError(error, responsesUrl, { + ...classifiedFailure, + retryable: false, + }) + } + throwClassifiedTransportError( + error, + responsesUrl, + classifiedFailure, + ) } if (responsesResponse.ok) { @@ -5093,6 +5511,7 @@ export function createOpenAIShimClient(options: { // Test-only surface (same pattern as WebSearchTool's __test export). export const __test = { convertMessages, + getApiTimeoutMs, getStreamIdleTimeoutMs, readWithIdleTimeout, StreamIdleTimeoutError, diff --git a/src/services/api/withRetry.test.ts b/src/services/api/withRetry.test.ts index c859ce1193..7134c61718 100644 --- a/src/services/api/withRetry.test.ts +++ b/src/services/api/withRetry.test.ts @@ -3,6 +3,7 @@ import type Anthropic from '@anthropic-ai/sdk' import { APIError, APIUserAbortError } from '@anthropic-ai/sdk' import { acquireSharedMutationLock, releaseSharedMutationLock } from '../../test/sharedMutationLock.js' import * as debugNs from '../../utils/debug.js' +import { markOpenAIRequestNonReplayable } from './openaiErrorClassification.js' type ProvidersModule = typeof import('../../utils/model/providers.js') // Helper to build a mock APIError with specific headers @@ -299,6 +300,45 @@ describe('abort retry classification', () => { }) describe('OpenAI-compatible retry classification', () => { + test('does not retry request timeouts marked as non-replayable', async () => { + process.env.OPENCLAUDE_RETRY_DELAY_MS = '1' + const { CannotRetryError, withRetry } = + await importFreshWithRetryModule('openai') + const error = markOpenAIRequestNonReplayable( + APIError.generate( + 0, + undefined, + 'OpenAI API transport error: no response headers [openai_category=request_timeout,host=slow.example.test]', + new Headers(), + ), + ) + let attempts = 0 + + let caught: unknown + try { + await drainAsyncGenerator( + withRetry( + async () => ({} as Anthropic), + async () => { + attempts++ + throw error + }, + { + maxRetries: 2, + model: 'gpt-4o-mini', + thinkingConfig: { type: 'disabled' }, + }, + ), + ) + } catch (error) { + caught = error + } + + expect(caught).toBeInstanceOf(CannotRetryError) + expect((caught as { originalError?: unknown }).originalError).toBe(error) + expect(attempts).toBe(1) + }) + test('does not retry marked non-retryable auth failures', async () => { process.env.OPENCLAUDE_RETRY_DELAY_MS = '1' const { CannotRetryError, withRetry } = diff --git a/src/services/api/withRetry.ts b/src/services/api/withRetry.ts index 5299790dd8..3cd4e9a0f1 100644 --- a/src/services/api/withRetry.ts +++ b/src/services/api/withRetry.ts @@ -53,6 +53,7 @@ import { REPEATED_529_ERROR_MESSAGE, isOpenCodeGoQuotaError } from './errors.js' import { extractConnectionErrorDetails } from './errorUtils.js' import { extractOpenAICategoryMarker, + isOpenAIRequestNonReplayable, isRetryableOpenAICompatibilityFailureCategory, } from './openaiErrorClassification.js' @@ -856,6 +857,10 @@ function shouldRetry(error: APIError, persistentRetryEnabled: boolean): boolean return false } + if (isOpenAIRequestNonReplayable(error)) { + return false + } + // OpenCode Go subscription quota exhaustion is terminal — retrying burns // the same 429 and confuses the user with repeated "mysterious stop" // failures. getAssistantMessageFromError surfaces the actionable message. diff --git a/src/utils/providerSecrets.ts b/src/utils/providerSecrets.ts index 586673cfbb..8dc2476202 100644 --- a/src/utils/providerSecrets.ts +++ b/src/utils/providerSecrets.ts @@ -266,6 +266,64 @@ export function redactSecretSubstringsForDisplay( return redacted } +const MAX_SECRET_ENCODING_LAYERS = 4 + +function escapeRegex(value: string): string { + return value.replace(/[.*+?^${}()|[\]\\]/g, '\\$&') +} + +function encodedVariantPattern(value: string): string { + return Array.from(value, char => { + const lower = char.toLowerCase() + return lower >= 'a' && lower <= 'f' + ? `[${lower}${lower.toUpperCase()}]` + : escapeRegex(char) + }).join('') +} + +function percentEncodeUtf8(value: string): string { + return Array.from( + new TextEncoder().encode(value), + byte => `%${byte.toString(16).padStart(2, '0').toUpperCase()}`, + ).join('') +} + +function encodedSecretPattern(value: string): RegExp { + const characterPatterns = Array.from(value, char => { + const variants = [escapeRegex(char)] + let encoded = percentEncodeUtf8(char) + for (let layer = 0; layer < MAX_SECRET_ENCODING_LAYERS; layer++) { + variants.push(encodedVariantPattern(encoded)) + encoded = encodeURIComponent(encoded) + } + return `(?:${variants.join('|')})` + }) + return new RegExp(characterPatterns.join(''), 'g') +} + +/** + * Redacts configured secrets when each character is literal or percent-encoded + * without decoding unrelated message text. Encoding depth is explicitly + * bounded to keep matching predictable on untrusted diagnostics. + */ +export function redactEncodedSecretSubstringsForDisplay( + value: string | null | undefined, + ...sources: Array +): string | undefined { + if (!value) return undefined + + let redacted = value + const secretValues = collectSecretValues(sources).sort( + (a, b) => b.length - a.length, + ) + for (const secretValue of secretValues) { + const mask = maskSecretForDisplay(secretValue) ?? 'configured' + redacted = redacted.replace(encodedSecretPattern(secretValue), mask) + } + + return redacted +} + export function sanitizeProviderConfigValue( value: string | null | undefined, ...sources: Array