From 0ef2b03bd06597ee14ed2faaba05ed28fbc31628 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 11:49:58 -0700 Subject: [PATCH 1/7] fix(providers): retry transient failures on OpenAI Responses and Gemini --- apps/sim/lib/core/errors/provider-quota.ts | 19 ++ apps/sim/lib/embeddings/client.ts | 21 +- .../conversation-generation-coverage.test.ts | 27 ++- apps/sim/providers/gemini/core.test.ts | 25 +- apps/sim/providers/gemini/core.ts | 62 +++-- .../providers/gemini/streaming-tool-loop.ts | 21 +- .../openai/core.transport-phase.test.ts | 33 ++- apps/sim/providers/openai/core.ts | 27 ++- apps/sim/providers/retry.test.ts | 219 ++++++++++++++++++ apps/sim/providers/retry.ts | 144 ++++++++++++ apps/sim/providers/transport.ts | 2 + apps/sim/providers/typesafe/transport.ts | 10 +- 12 files changed, 537 insertions(+), 73 deletions(-) create mode 100644 apps/sim/lib/core/errors/provider-quota.ts create mode 100644 apps/sim/providers/retry.test.ts create mode 100644 apps/sim/providers/retry.ts diff --git a/apps/sim/lib/core/errors/provider-quota.ts b/apps/sim/lib/core/errors/provider-quota.ts new file mode 100644 index 00000000000..6ba9437c9f9 --- /dev/null +++ b/apps/sim/lib/core/errors/provider-quota.ts @@ -0,0 +1,19 @@ +/** + * True when a provider rejection body reports an exhausted balance rather than a rate + * limit. OpenAI returns 429 for both, but only a rate limit reopens: a spent account + * stands until someone adds credit, so retrying it cannot succeed. + */ +export function isQuotaExhaustionBody(errorText: string): boolean { + try { + const body = JSON.parse(errorText) as { error?: { type?: string; code?: string } } + const type = body.error?.type + const code = body.error?.code + return ( + type === 'insufficient_quota' || + code === 'insufficient_quota' || + code === 'credit_balance_exhausted' + ) + } catch { + return false + } +} diff --git a/apps/sim/lib/embeddings/client.ts b/apps/sim/lib/embeddings/client.ts index 7369b40e12c..6f7fd9e0028 100644 --- a/apps/sim/lib/embeddings/client.ts +++ b/apps/sim/lib/embeddings/client.ts @@ -11,6 +11,7 @@ import { wireFallback, } from '@/lib/core/config/env-capabilities' import { isHosted } from '@/lib/core/config/env-flags' +import { isQuotaExhaustionBody } from '@/lib/core/errors/provider-quota' import { ProviderQuotaExhaustedError, recordProviderCooldown, @@ -254,26 +255,6 @@ export function isBYOKEmbeddingCredentialRejection(error: unknown): error is Emb ) } -/** - * True when a rejection body reports an exhausted balance rather than a rate - * limit. OpenAI returns 429 for both, but only a rate limit reopens: a spent - * account stands until someone adds credit, so retrying it cannot succeed. - */ -function isQuotaExhaustionBody(errorText: string): boolean { - try { - const body = JSON.parse(errorText) as { error?: { type?: string; code?: string } } - const type = body.error?.type - const code = body.error?.code - return ( - type === 'insufficient_quota' || - code === 'insufficient_quota' || - code === 'credit_balance_exhausted' - ) - } catch { - return false - } -} - /** Reads a bounded provider body for internal diagnostics and quota classification. */ async function readEmbeddingErrorBody(response: Response, signal?: AbortSignal): Promise { try { diff --git a/apps/sim/providers/conversation-generation-coverage.test.ts b/apps/sim/providers/conversation-generation-coverage.test.ts index 2a9f4d0d7d5..fa1b25d9d4f 100644 --- a/apps/sim/providers/conversation-generation-coverage.test.ts +++ b/apps/sim/providers/conversation-generation-coverage.test.ts @@ -21,6 +21,7 @@ function insideNamedAncestor(node: ts.Node, name: string): boolean { } function preparedPayload(node: ts.Node): boolean { + if (ts.isIdentifier(node)) return preparedBinding(node) return ( ts.isAwaitExpression(node) && ts.isCallExpression(node.expression) && @@ -28,6 +29,30 @@ function preparedPayload(node: ts.Node): boolean { ) } +/** + * A `const` in an enclosing block bound to a prepared payload. Retried sends prepare once + * and replay the binding, since preparing can itself call a model to compact history. + */ +function preparedBinding(identifier: ts.Identifier): boolean { + for (let scope = identifier.parent; scope; scope = scope.parent) { + if (!ts.isBlock(scope) && !ts.isSourceFile(scope)) continue + for (const statement of scope.statements) { + if (!ts.isVariableStatement(statement)) continue + if (!(statement.declarationList.flags & ts.NodeFlags.Const)) continue + for (const declaration of statement.declarationList.declarations) { + if ( + declaration.name.getText() === identifier.text && + declaration.initializer && + preparedPayload(declaration.initializer) + ) { + return true + } + } + } + } + return false +} + describe('provider generation context coverage', () => { it('guards every model SDK send, shared stream callback, retry and finalizer', () => { const uncovered: string[] = [] @@ -60,7 +85,7 @@ describe('provider generation context coverage', () => { file.endsWith('/openai-compat/streaming-tool-loop.ts') ) { argument = 0 - } else if (callee === 'JSON.stringify' && insideNamedAncestor(node, 'postOnce')) { + } else if (callee === 'JSON.stringify' && insideNamedAncestor(node, 'post')) { argument = 0 } if (argument !== undefined) { diff --git a/apps/sim/providers/gemini/core.test.ts b/apps/sim/providers/gemini/core.test.ts index 19ab6d4b9cf..1dbfde9461c 100644 --- a/apps/sim/providers/gemini/core.test.ts +++ b/apps/sim/providers/gemini/core.test.ts @@ -1,4 +1,8 @@ -import type { GenerateContentParameters, GenerateContentResponse } from '@google/genai' +import { + ApiError, + type GenerateContentParameters, + type GenerateContentResponse, +} from '@google/genai' import { providersMock } from '@sim/testing/mocks/providers.mock' import { providersConversationHistoryMock, @@ -142,6 +146,25 @@ describe('Vertex Gemini request compatibility', () => { } ) + /** The SDK's own retry is off (it would replace every error message), so ours must run. */ + it('replays a transient server error before reporting it', async () => { + vi.useFakeTimers() + try { + const generateContent = vi + .fn() + .mockRejectedValueOnce(new ApiError({ message: 'Internal error', status: 500 })) + .mockResolvedValue(textTurn()) + + const pending = run('vertex/gemini-3.8-flash', generateContent) + await vi.runAllTimersAsync() + + await expect(pending).resolves.toMatchObject({ content: 'answer' }) + expect(generateContent).toHaveBeenCalledTimes(2) + } finally { + vi.useRealTimers() + } + }) + it('prices Vertex-only catalog entries using the namespaced model ID', async () => { const generateContent = vi.fn().mockResolvedValue(textTurn()) const result = await run('vertex/gemini-3.8-flash', generateContent) diff --git a/apps/sim/providers/gemini/core.ts b/apps/sim/providers/gemini/core.ts index 29d18e0c91a..922ca3898cd 100644 --- a/apps/sim/providers/gemini/core.ts +++ b/apps/sim/providers/gemini/core.ts @@ -37,6 +37,7 @@ import { supportsDisablingGemini25Thinking, } from '@/providers/google/utils' import { getModelCapabilities, isKnownModelId } from '@/providers/models' +import { withProviderRetry } from '@/providers/retry' import { executeProviderTool } from '@/providers/runtime-context' import { createSettledAgentEventStream } from '@/providers/stream-events' import { createStreamingExecution } from '@/providers/streaming-execution' @@ -962,6 +963,11 @@ export async function executeGeminiRequest( } const logger = createLogger(providerType === 'google' ? 'GoogleProvider' : 'VertexProvider') + const retryOptions = { + logger, + label: providerType === 'google' ? 'Gemini' : 'Vertex AI', + abortSignal: request.abortSignal, + } logger.info(`Preparing ${providerType} Gemini request`, { model, @@ -1148,12 +1154,14 @@ export async function executeGeminiRequest( if (shouldStream) { logger.info('Handling Gemini streaming response') - const streamGenerator = await ai.models.generateContentStream( - await prepareConversationGeneration(request, 'gemini', { - model, - contents, - config: geminiConfig, - }) + const streamPayload = await prepareConversationGeneration(request, 'gemini', { + model, + contents, + config: geminiConfig, + }) + const streamGenerator = await withProviderRetry( + () => ai.models.generateContentStream(streamPayload), + retryOptions ) const firstResponseTime = Date.now() - initialCallTime @@ -1202,12 +1210,14 @@ export async function executeGeminiRequest( } // Non-streaming request - const response = await ai.models.generateContent( - await prepareConversationGeneration(request, 'gemini', { - model, - contents, - config: geminiConfig, - }) + const responsePayload = await prepareConversationGeneration(request, 'gemini', { + model, + contents, + config: geminiConfig, + }) + const response = await withProviderRetry( + () => ai.models.generateContent(responsePayload), + retryOptions ) if (!extractAllFunctionCallParts(response.candidates?.[0]).length) { await captureProviderConversationStep( @@ -1256,12 +1266,14 @@ export async function executeGeminiRequest( } const finalStartTime = Date.now() - const finalResponse = await ai.models.generateContent( - await prepareConversationGeneration(request, 'gemini', { - model, - contents: currentState.contents, - config: finalConfig, - }) + const finalResponsePayload = await prepareConversationGeneration(request, 'gemini', { + model, + contents: currentState.contents, + config: finalConfig, + }) + const finalResponse = await withProviderRetry( + () => ai.models.generateContent(finalResponsePayload), + retryOptions ) if (!extractAllFunctionCallParts(finalResponse.candidates?.[0]).length) { await captureProviderConversationStep( @@ -1375,12 +1387,14 @@ export async function executeGeminiRequest( /** Resolve the final turn, then project its settled answer when streaming was requested. */ const nextModelStartTime = Date.now() - const nextResponse = await ai.models.generateContent( - await prepareConversationGeneration(request, 'gemini', { - model, - contents: state.contents, - config: nextConfig, - }) + const nextResponsePayload = await prepareConversationGeneration(request, 'gemini', { + model, + contents: state.contents, + config: nextConfig, + }) + const nextResponse = await withProviderRetry( + () => ai.models.generateContent(nextResponsePayload), + retryOptions ) if (!extractAllFunctionCallParts(nextResponse.candidates?.[0]).length) { await captureProviderConversationStep( diff --git a/apps/sim/providers/gemini/streaming-tool-loop.ts b/apps/sim/providers/gemini/streaming-tool-loop.ts index 4d3e22c2e65..ec62002395b 100644 --- a/apps/sim/providers/gemini/streaming-tool-loop.ts +++ b/apps/sim/providers/gemini/streaming-tool-loop.ts @@ -38,6 +38,7 @@ import { convertUsageMetadata, ensureStructResponse, } from '@/providers/google/utils' +import { withProviderRetry } from '@/providers/retry' import { executeProviderTool } from '@/providers/runtime-context' import type { AgentStreamEvent, ToolCallEndStatus } from '@/providers/stream-events' import { @@ -309,15 +310,17 @@ export function createGeminiStreamingToolLoopStream( ) const modelStart = Date.now() - const streamGenerator = await ai.models.generateContentStream( - await prepareConversationGeneration(request, 'gemini', { - model, - contents, - config: { - ...turnConfig, - abortSignal: loopAbortController.signal, - }, - }) + const turnPayload = await prepareConversationGeneration(request, 'gemini', { + model, + contents, + config: { + ...turnConfig, + abortSignal: loopAbortController.signal, + }, + }) + const streamGenerator = await withProviderRetry( + () => ai.models.generateContentStream(turnPayload), + { logger, label: 'Gemini', abortSignal: loopAbortController.signal } ) const drained = await drainGeminiTurn( diff --git a/apps/sim/providers/openai/core.transport-phase.test.ts b/apps/sim/providers/openai/core.transport-phase.test.ts index 15411fb0a67..d8f9aaa474f 100644 --- a/apps/sim/providers/openai/core.transport-phase.test.ts +++ b/apps/sim/providers/openai/core.transport-phase.test.ts @@ -40,6 +40,9 @@ function timeoutError() { return new DOMException('The operation timed out.', 'TimeoutError') } +/** Pins an error response to one attempt, so these cases exercise presentation, not retry. */ +const NOT_RETRYABLE = new Headers({ 'x-should-retry': 'false' }) + const COMPLETED = { id: 'resp_1', status: 'completed', @@ -116,7 +119,7 @@ describe('OpenAI transport phase annotation', () => { const apiError = { ok: false, status: 429, - headers: new Headers(), + headers: NOT_RETRYABLE, text: () => Promise.resolve(JSON.stringify({ error: { message: 'Rate limit reached' } })), } @@ -129,7 +132,7 @@ describe('OpenAI transport phase annotation', () => { const htmlError = { ok: false, status: 502, - headers: new Headers(), + headers: NOT_RETRYABLE, text: () => Promise.resolve(`${'x'.repeat(5000)}`), } @@ -180,7 +183,7 @@ describe('OpenAI transport phase annotation', () => { const unreadable = { ok: false, status: 502, - headers: new Headers(), + headers: NOT_RETRYABLE, text: () => Promise.reject(timeoutError()), } @@ -209,6 +212,30 @@ describe('OpenAI transport phase annotation', () => { expect(error.message).toMatch(/elapsedMs=\d+/) }) + /** The incident this path once had: a single transient 500 failed the whole run. */ + it('replays a transient server error before reporting it', async () => { + vi.useFakeTimers() + try { + const fetchMock = vi + .fn() + .mockResolvedValueOnce({ ok: false, status: 500, headers: new Headers() }) + .mockResolvedValue({ + ok: true, + status: 200, + headers: new Headers(), + json: () => Promise.resolve(COMPLETED), + }) + + const pending = run(fetchMock) + await vi.runAllTimersAsync() + + await expect(pending).resolves.toMatchObject({ content: 'ok' }) + expect(fetchMock).toHaveBeenCalledTimes(2) + } finally { + vi.useRealTimers() + } + }) + it('leaves a healthy response entirely unaffected', async () => { const fetchMock = vi.fn().mockResolvedValue({ ok: true, diff --git a/apps/sim/providers/openai/core.ts b/apps/sim/providers/openai/core.ts index 9cc91967d3b..4a17df95ca8 100644 --- a/apps/sim/providers/openai/core.ts +++ b/apps/sim/providers/openai/core.ts @@ -23,6 +23,7 @@ import { buildOpenAIUsageTokens, createOpenAIUsageAccumulator, } from '@/providers/openai/usage' +import { fetchWithProviderRetry } from '@/providers/retry' import { executeProviderTool } from '@/providers/runtime-context' import { createStreamingExecution } from '@/providers/streaming-execution' import { isAbortError, parseToolArguments } from '@/providers/streaming-tool-loop-shared' @@ -428,19 +429,27 @@ export async function executeResponsesProviderRequest( * headers is named on the streaming paths too — they call * {@link fetchResponsesWithSummaryFallback} directly and never reach `postResponses`, * which is where the annotation used to live. + * + * The body is prepared once, outside the retry: preparing it can compact the + * conversation with a model call of its own, which a replayed send must not repeat. */ - const postOnce = async ( + const post = async ( payload: Record, abortSignal: AbortSignal | undefined, startedAt: number ): Promise => { + const body = JSON.stringify(await prepareConversationGeneration(request, 'responses', payload)) try { - return await fetchImpl(config.endpoint, { - method: 'POST', - headers: config.headers, - body: JSON.stringify(await prepareConversationGeneration(request, 'responses', payload)), - signal: abortSignal, - }) + return await fetchWithProviderRetry( + () => + fetchImpl(config.endpoint, { + method: 'POST', + headers: config.headers, + body, + signal: abortSignal, + }), + { logger, label: config.providerLabel, abortSignal } + ) } catch (error) { throw annotateTransportFailure(error, 'awaiting-response-headers', startedAt) } @@ -454,7 +463,7 @@ export async function executeResponsesProviderRequest( const body = reasoningSummariesUnavailable ? (stripReasoningSummary(requestedBody) ?? requestedBody) : requestedBody - const response = await postOnce(body, abortSignal, startedAt) + const response = await post(body, abortSignal, startedAt) if (response.ok) return response const message = await parseErrorResponse(response, startedAt) @@ -470,7 +479,7 @@ export async function executeResponsesProviderRequest( `${config.providerLabel} rejected reasoning summaries (organization not verified); retrying without summary`, { model: config.modelName } ) - const retryResponse = await postOnce(strippedBody, abortSignal, startedAt) + const retryResponse = await post(strippedBody, abortSignal, startedAt) if (!retryResponse.ok) { const retryMessage = await parseErrorResponse(retryResponse, startedAt) throw new Error( diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts new file mode 100644 index 00000000000..dcb090786c2 --- /dev/null +++ b/apps/sim/providers/retry.test.ts @@ -0,0 +1,219 @@ +/** + * Failure modes of the shared provider retry policy. A lost retry fails a run on a transient + * upstream fault; a wrong one re-bills a completion or hammers a spent account. + */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { fetchWithProviderRetry, withProviderRetry } from '@/providers/retry' +import { PROVIDER_MAX_RETRIES } from '@/providers/transport' + +const logger = { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() } as never +const MAX_ATTEMPTS = PROVIDER_MAX_RETRIES + 1 + +function reply(status: number, body = '', headers: Record = {}) { + return () => Promise.resolve(new Response(body || null, { status, headers })) +} + +function connectionReset() { + return Object.assign(new Error('The socket connection was closed unexpectedly'), { + code: 'ECONNRESET', + }) +} + +/** Runs the call to completion while fake timers drive every backoff sleep. */ +async function settle(promise: Promise): Promise { + const outcome = promise.then( + (value) => ({ value }), + (error: unknown) => ({ error }) + ) + await vi.runAllTimersAsync() + const result = await outcome + if ('error' in result) throw result.error + return result.value +} + +describe('fetchWithProviderRetry', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('retries a transient 5xx and returns the later success', async () => { + const send = vi.fn().mockImplementationOnce(reply(500)).mockImplementation(reply(200, 'ok')) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + expect(await response.text()).toBe('ok') + expect(send).toHaveBeenCalledTimes(2) + }) + + it.each([408, 409, 429, 502, 503])('treats %i as transient', async (status) => { + const send = vi.fn().mockImplementationOnce(reply(status)).mockImplementation(reply(200)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + }) + + it('stops after the retry budget and hands back the last response with its body intact', async () => { + const send = vi.fn().mockImplementation(reply(500, '{"error":{"message":"server error"}}')) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(send).toHaveBeenCalledTimes(MAX_ATTEMPTS) + expect(response.status).toBe(500) + expect(await response.text()).toContain('server error') + }) + + it('never replays a request the provider rejected as invalid', async () => { + const send = vi.fn().mockImplementation(reply(400)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(400) + expect(send).toHaveBeenCalledTimes(1) + }) + + /** OpenAI answers 429 for a spent balance too, and that one never reopens on its own. */ + it('does not retry an exhausted quota, and leaves its body readable for the caller', async () => { + const body = JSON.stringify({ error: { type: 'insufficient_quota', message: 'No credit' } }) + const send = vi.fn().mockImplementation(reply(429, body)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(send).toHaveBeenCalledTimes(1) + expect(await response.text()).toContain('No credit') + }) + + it('obeys the server when it says a 5xx must not be retried', async () => { + const send = vi.fn().mockImplementation(reply(500, '', { 'x-should-retry': 'false' })) + + await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(send).toHaveBeenCalledTimes(1) + }) + + it('obeys the server when it says a 4xx is safe to retry', async () => { + const send = vi + .fn() + .mockImplementationOnce(reply(400, '', { 'x-should-retry': 'true' })) + .mockImplementation(reply(200)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + }) + + it('waits the retry-after-ms the server asked for before replaying', async () => { + const send = vi + .fn() + .mockImplementationOnce(reply(429, '', { 'retry-after-ms': '4000', 'retry-after': '1' })) + .mockImplementation(reply(200)) + + const pending = fetchWithProviderRetry(send, { logger, label: 'OpenAI' }) + await vi.advanceTimersByTimeAsync(3900) + expect(send).toHaveBeenCalledTimes(1) + + await vi.advanceTimersByTimeAsync(200) + expect(send).toHaveBeenCalledTimes(2) + await pending + }) + + it('retries a dropped connection', async () => { + const send = vi.fn().mockRejectedValueOnce(connectionReset()).mockImplementation(reply(200)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + }) + + it('does not retry an unrecognised failure', async () => { + const send = vi.fn().mockRejectedValue(new Error('body serialization failed')) + + await expect(settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' }))).rejects.toThrow( + 'body serialization failed' + ) + expect(send).toHaveBeenCalledTimes(1) + }) + + it('stops as soon as the caller aborts during backoff', async () => { + const controller = new AbortController() + const send = vi.fn().mockImplementation(reply(503)) + + const pending = fetchWithProviderRetry(send, { + logger, + label: 'OpenAI', + abortSignal: controller.signal, + }) + const outcome = pending.catch((error: unknown) => error) + await vi.advanceTimersByTimeAsync(0) + controller.abort() + + expect(await outcome).toMatchObject({ name: 'AbortError' }) + expect(send).toHaveBeenCalledTimes(1) + }) + + it('surfaces the real failure of a request the caller already aborted, without a retry', async () => { + const controller = new AbortController() + const send = vi.fn().mockImplementation(() => { + controller.abort() + return Promise.reject(connectionReset()) + }) + + await expect( + settle( + fetchWithProviderRetry(send, { + logger, + label: 'OpenAI', + abortSignal: controller.signal, + }) + ) + ).rejects.toMatchObject({ code: 'ECONNRESET' }) + expect(send).toHaveBeenCalledTimes(1) + }) +}) + +describe('withProviderRetry', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + vi.useRealTimers() + }) + + function sdkError(status: number) { + return Object.assign(new Error(`got status: ${status}`), { status }) + } + + it('retries an SDK error that carries a transient status', async () => { + const operation = vi.fn().mockRejectedValueOnce(sdkError(503)).mockResolvedValue('answer') + + await expect(settle(withProviderRetry(operation, { logger, label: 'Gemini' }))).resolves.toBe( + 'answer' + ) + expect(operation).toHaveBeenCalledTimes(2) + }) + + it('surfaces the original SDK error once retries are exhausted', async () => { + const final = sdkError(500) + const operation = vi.fn().mockRejectedValue(final) + + await expect(settle(withProviderRetry(operation, { logger, label: 'Gemini' }))).rejects.toBe( + final + ) + expect(operation).toHaveBeenCalledTimes(MAX_ATTEMPTS) + }) + + it('does not retry an SDK error that carries a permanent status', async () => { + const operation = vi.fn().mockRejectedValue(sdkError(400)) + + await expect(settle(withProviderRetry(operation, { logger, label: 'Gemini' }))).rejects.toThrow( + 'status: 400' + ) + expect(operation).toHaveBeenCalledTimes(1) + }) +}) diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts new file mode 100644 index 00000000000..a3802fee064 --- /dev/null +++ b/apps/sim/providers/retry.ts @@ -0,0 +1,144 @@ +import type { Logger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { interruptibleSleep } from '@sim/utils/helpers' +import { isRecordLike } from '@sim/utils/object' +import { backoffWithJitter, parseRetryAfter } from '@sim/utils/retry' +import { isQuotaExhaustionBody } from '@/lib/core/errors/provider-quota' +import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure' +import { + consumeOrCancelBody, + DEFAULT_MAX_ERROR_BODY_BYTES, + readResponseTextWithLimit, +} from '@/lib/core/utils/stream-limits' +import { PROVIDER_MAX_RETRIES } from '@/providers/transport' + +/** + * Retry policy for provider calls that no vendor SDK retries for us: the raw-fetch + * Responses path and SDKs whose own retry is unusable. + * + * It mirrors what `openai`, `@anthropic-ai/sdk`, and the AI SDK agree on, so a model + * behaves the same whichever transport serves it: {@link PROVIDER_MAX_RETRIES} replays, + * jittered exponential backoff, the server's `x-should-retry` and `retry-after-ms` / + * `retry-after` obeyed, and a caller's abort never retried. Only the wait for response + * headers is covered — once a body is streaming, nothing is replayed. + * + * One deliberate divergence: a 429 reporting an exhausted balance is not retried. The + * SDKs replay it, but a spent account does not reopen within a backoff window. + */ + +export interface ProviderRetryOptions { + logger: Logger + /** Provider name for the retry log line. */ + label: string + abortSignal?: AbortSignal +} + +/** Statuses every vendor SDK treats as transient: timeout, lock conflict, rate limit, server fault. */ +export function isRetryableProviderStatus(status: number): boolean { + return status === 408 || status === 409 || status === 429 || status >= 500 +} + +/** The server's requested delay: OpenAI's precise `retry-after-ms`, then standard `retry-after`. */ +export function providerRetryAfterMs(headers: Headers): number | null { + const precise = Number.parseFloat(headers.get('retry-after-ms') ?? '') + if (Number.isFinite(precise) && precise >= 0) return precise + return parseRetryAfter(headers.get('retry-after')) +} + +/** + * A dropped or refused connection: Node's `fetch` raises a `TypeError`, Bun's an `Error` + * carrying the syscall code. + */ +function isRetryableTransportFailure(error: unknown): boolean { + return error instanceof TypeError || isRetryableInfrastructureError(error) +} + +/** Reads a clone, so the caller can still read the body of a response that is not retried. */ +async function isQuotaExhaustedResponse(response: Response): Promise { + const body = await readResponseTextWithLimit(response.clone(), { + maxBytes: DEFAULT_MAX_ERROR_BODY_BYTES, + label: 'provider error response', + }).catch(() => '') + return isQuotaExhaustionBody(body) +} + +async function shouldRetryResponse(response: Response): Promise { + const directive = response.headers.get('x-should-retry') + if (directive === 'true') return true + if (directive === 'false') return false + if (!isRetryableProviderStatus(response.status)) return false + return response.status !== 429 || !(await isQuotaExhaustedResponse(response)) +} + +async function waitBeforeRetry( + attempt: number, + retryAfterMs: number | null, + reason: string, + { logger, label, abortSignal }: ProviderRetryOptions +): Promise { + const delayMs = backoffWithJitter(attempt, retryAfterMs) + logger.warn(`${label} request failed (${reason}); retrying`, { + attempt, + maxRetries: PROVIDER_MAX_RETRIES, + delayMs: Math.round(delayMs), + }) + await interruptibleSleep(delayMs, abortSignal) + abortSignal?.throwIfAborted() +} + +/** + * Sends a raw-fetch provider request, replaying it on a transient failure. + * + * Resolves with the final response — successful, not retryable, or out of retries — so + * the caller reads and reports its body exactly as it would a single attempt's. Rejects + * only with the last transport failure. + */ +export async function fetchWithProviderRetry( + send: () => Promise, + options: ProviderRetryOptions +): Promise { + for (let attempt = 1; ; attempt++) { + const canRetry = attempt <= PROVIDER_MAX_RETRIES + let response: Response + try { + response = await send() + } catch (error) { + if (!canRetry || options.abortSignal?.aborted || !isRetryableTransportFailure(error)) { + throw error + } + await waitBeforeRetry(attempt, null, getErrorMessage(error), options) + continue + } + if (response.ok || !canRetry || !(await shouldRetryResponse(response))) return response + await consumeOrCancelBody(response) + await waitBeforeRetry( + attempt, + providerRetryAfterMs(response.headers), + `HTTP ${response.status}`, + options + ) + } +} + +/** + * Runs a vendor SDK call, replaying it when it rejects with a transient HTTP status or a + * dropped connection. For SDKs that raise non-OK responses as errors carrying `status`. + */ +export async function withProviderRetry( + operation: () => Promise, + options: ProviderRetryOptions +): Promise { + for (let attempt = 1; ; attempt++) { + try { + return await operation() + } catch (error) { + const status = isRecordLike(error) && typeof error.status === 'number' ? error.status : null + const retryable = + status === null ? isRetryableTransportFailure(error) : isRetryableProviderStatus(status) + if (attempt > PROVIDER_MAX_RETRIES || options.abortSignal?.aborted || !retryable) { + throw error + } + await waitBeforeRetry(attempt, null, getErrorMessage(error), options) + } + } +} diff --git a/apps/sim/providers/transport.ts b/apps/sim/providers/transport.ts index 3755d3cf8ce..420fb926de5 100644 --- a/apps/sim/providers/transport.ts +++ b/apps/sim/providers/transport.ts @@ -37,6 +37,8 @@ export const PROVIDER_HEADERS_TIMEOUT_MS = 600_000 * is non-idempotent and carries no idempotency key, and on the non-streaming path * the response only exists once the generation has already been billed — so a * replay re-bills completed work, multiplied by every turn of the tool loop. + * + * Calls no vendor SDK retries apply the same budget through `@/providers/retry`. */ export const PROVIDER_MAX_RETRIES = 2 diff --git a/apps/sim/providers/typesafe/transport.ts b/apps/sim/providers/typesafe/transport.ts index 137caf295ad..9f74d3edc68 100644 --- a/apps/sim/providers/typesafe/transport.ts +++ b/apps/sim/providers/typesafe/transport.ts @@ -1,7 +1,8 @@ import { interruptibleSleep } from '@sim/utils/helpers' -import { backoffWithJitter, parseRetryAfter } from '@sim/utils/retry' +import { backoffWithJitter } from '@sim/utils/retry' import { stringifyBoundedJson } from '@/lib/core/utils/bounded-json' import { consumeOrCancelBody, readResponseJsonWithLimit } from '@/lib/core/utils/stream-limits' +import { isRetryableProviderStatus, providerRetryAfterMs } from '@/providers/retry' import { PROVIDER_HEADERS_TIMEOUT_MS, PROVIDER_MAX_RETRIES } from '@/providers/transport' import type { buildJevBody } from '@/providers/typesafe/schema' @@ -43,10 +44,7 @@ export async function requestJevEvaluation( }) if (!response.ok) { await consumeOrCancelBody(response) - throw new TypeSafeHttpError( - response.status, - parseRetryAfter(response.headers.get('retry-after')) - ) + throw new TypeSafeHttpError(response.status, providerRetryAfterMs(response.headers)) } return await readResponseJsonWithLimit(response, { maxBytes: MAX_EVALUATION_RESPONSE_BYTES, @@ -57,7 +55,7 @@ export async function requestJevEvaluation( abortSignal?.throwIfAborted() const retryable = error instanceof TypeSafeHttpError - ? error.status === 408 || error.status === 429 || error.status >= 500 + ? isRetryableProviderStatus(error.status) : !response || timeout.aborted || error instanceof TypeError if (!retryable || attempt >= PROVIDER_MAX_RETRIES) throw error const retryAfterMs = error instanceof TypeSafeHttpError ? error.retryAfterMs : null From 41c378660780f5c697cf6fa6f3c3bf42411c4995 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 12:02:38 -0700 Subject: [PATCH 2/7] fix(providers): read 429 bodies without a tee, honor Gemini RetryInfo, resolve shadowed payloads --- .../conversation-generation-coverage.test.ts | 38 +++++++---- apps/sim/providers/gemini/core.ts | 2 + .../providers/gemini/streaming-tool-loop.ts | 8 ++- apps/sim/providers/google/utils.test.ts | 40 ++++++++++++ apps/sim/providers/google/utils.ts | 27 +++++++- apps/sim/providers/retry.test.ts | 47 ++++++++++++++ apps/sim/providers/retry.ts | 63 ++++++++++++++----- 7 files changed, 194 insertions(+), 31 deletions(-) diff --git a/apps/sim/providers/conversation-generation-coverage.test.ts b/apps/sim/providers/conversation-generation-coverage.test.ts index fa1b25d9d4f..fda584bdf34 100644 --- a/apps/sim/providers/conversation-generation-coverage.test.ts +++ b/apps/sim/providers/conversation-generation-coverage.test.ts @@ -30,24 +30,38 @@ function preparedPayload(node: ts.Node): boolean { } /** - * A `const` in an enclosing block bound to a prepared payload. Retried sends prepare once - * and replay the binding, since preparing can itself call a model to compact history. + * Whether the nearest declaration of `identifier` is a `const` bound to a prepared payload. + * Retried sends prepare once and replay the binding, since preparing can itself call a model + * to compact history. The nearest declaration decides, so an inner shadow must be prepared too. */ function preparedBinding(identifier: ts.Identifier): boolean { + const declares = (name: ts.BindingName) => name.getText() === identifier.text for (let scope = identifier.parent; scope; scope = scope.parent) { + if (ts.isFunctionLike(scope) && scope.parameters.some((parameter) => declares(parameter.name))) + return false + if ( + ts.isCatchClause(scope) && + scope.variableDeclaration && + declares(scope.variableDeclaration.name) + ) + return false + if ( + (ts.isForStatement(scope) || ts.isForInStatement(scope) || ts.isForOfStatement(scope)) && + scope.initializer && + ts.isVariableDeclarationList(scope.initializer) && + scope.initializer.declarations.some((declaration) => declares(declaration.name)) + ) + return false if (!ts.isBlock(scope) && !ts.isSourceFile(scope)) continue for (const statement of scope.statements) { if (!ts.isVariableStatement(statement)) continue - if (!(statement.declarationList.flags & ts.NodeFlags.Const)) continue - for (const declaration of statement.declarationList.declarations) { - if ( - declaration.name.getText() === identifier.text && - declaration.initializer && - preparedPayload(declaration.initializer) - ) { - return true - } - } + const declaration = statement.declarationList.declarations.find((d) => declares(d.name)) + if (!declaration) continue + return ( + Boolean(statement.declarationList.flags & ts.NodeFlags.Const) && + declaration.initializer !== undefined && + preparedPayload(declaration.initializer) + ) } } return false diff --git a/apps/sim/providers/gemini/core.ts b/apps/sim/providers/gemini/core.ts index 922ca3898cd..17970aa8aaa 100644 --- a/apps/sim/providers/gemini/core.ts +++ b/apps/sim/providers/gemini/core.ts @@ -32,6 +32,7 @@ import { ensureStructResponse, extractAllFunctionCallParts, extractTextContent, + geminiRetryDelayMs, mapToThinkingBudget, mapToThinkingLevel, supportsDisablingGemini25Thinking, @@ -967,6 +968,7 @@ export async function executeGeminiRequest( logger, label: providerType === 'google' ? 'Gemini' : 'Vertex AI', abortSignal: request.abortSignal, + retryAfterMs: geminiRetryDelayMs, } logger.info(`Preparing ${providerType} Gemini request`, { diff --git a/apps/sim/providers/gemini/streaming-tool-loop.ts b/apps/sim/providers/gemini/streaming-tool-loop.ts index ec62002395b..c57bce99fde 100644 --- a/apps/sim/providers/gemini/streaming-tool-loop.ts +++ b/apps/sim/providers/gemini/streaming-tool-loop.ts @@ -37,6 +37,7 @@ import { cleanSchemaForGemini, convertUsageMetadata, ensureStructResponse, + geminiRetryDelayMs, } from '@/providers/google/utils' import { withProviderRetry } from '@/providers/retry' import { executeProviderTool } from '@/providers/runtime-context' @@ -320,7 +321,12 @@ export function createGeminiStreamingToolLoopStream( }) const streamGenerator = await withProviderRetry( () => ai.models.generateContentStream(turnPayload), - { logger, label: 'Gemini', abortSignal: loopAbortController.signal } + { + logger, + label: 'Gemini', + abortSignal: loopAbortController.signal, + retryAfterMs: geminiRetryDelayMs, + } ) const drained = await drainGeminiTurn( diff --git a/apps/sim/providers/google/utils.test.ts b/apps/sim/providers/google/utils.test.ts index b415a4993c1..561b50b283b 100644 --- a/apps/sim/providers/google/utils.test.ts +++ b/apps/sim/providers/google/utils.test.ts @@ -1,9 +1,11 @@ +import { ApiError } from '@google/genai' import { describe, expect, it, vi } from 'vitest' import { setNativeConversationMessage } from '@/providers/conversation-metadata' import { convertToGeminiFormat, convertUsageMetadata, createReadableStreamFromGeminiStream, + geminiRetryDelayMs, mapToThinkingBudget, } from '@/providers/google/utils' import type { AgentStreamEvent } from '@/providers/stream-events' @@ -585,3 +587,41 @@ describe('createReadableStreamFromGeminiStream', () => { ) }) }) + +/** The SDK sets `ApiError.message` to the JSON error body, so this is the wire shape it sees. */ +describe('geminiRetryDelayMs', () => { + function rateLimited(retryDelay: string) { + return new ApiError({ + status: 429, + message: JSON.stringify({ + error: { + code: 429, + status: 'RESOURCE_EXHAUSTED', + details: [ + { '@type': 'type.googleapis.com/google.rpc.QuotaFailure', violations: [] }, + { '@type': 'type.googleapis.com/google.rpc.RetryInfo', retryDelay }, + ], + }, + }), + }) + } + + it.each([ + ['31s', 31_000], + ['0.5s', 500], + ])('reads a RetryInfo delay of %s', (retryDelay, expected) => { + expect(geminiRetryDelayMs(rateLimited(retryDelay))).toBe(expected) + }) + + it('has no delay when the error carries no RetryInfo', () => { + const error = new ApiError({ + status: 503, + message: JSON.stringify({ error: { code: 503, status: 'UNAVAILABLE' } }), + }) + expect(geminiRetryDelayMs(error)).toBeNull() + }) + + it('has no delay for a transport failure', () => { + expect(geminiRetryDelayMs(new TypeError('fetch failed'))).toBeNull() + }) +}) diff --git a/apps/sim/providers/google/utils.ts b/apps/sim/providers/google/utils.ts index 4f315edb035..18617c222b4 100644 --- a/apps/sim/providers/google/utils.ts +++ b/apps/sim/providers/google/utils.ts @@ -1,4 +1,5 @@ import { + ApiError, type Candidate, type Content, type FunctionCall, @@ -14,7 +15,7 @@ import { } from '@google/genai' import { createLogger } from '@sim/logger' import { toError } from '@sim/utils/errors' -import { isRecordLike } from '@sim/utils/object' +import { isRecordLike, toArray, toRecord } from '@sim/utils/object' import { buildGeminiMessageParts } from '@/providers/attachments' import { captureProviderConversationStep } from '@/providers/conversation-history' import { @@ -29,6 +30,30 @@ import { trackForcedToolUsage } from '@/providers/utils' const logger = createLogger('GoogleUtils') +const RETRY_INFO_TYPE = 'type.googleapis.com/google.rpc.RetryInfo' + +/** + * The delay a Gemini rejection asks for before a retry. The SDK surfaces no response + * headers: `ApiError.message` is the JSON error body, and a rate limit carries a + * `google.rpc.RetryInfo` detail whose `retryDelay` is a protobuf Duration such as `"31s"`. + */ +export function geminiRetryDelayMs(error: unknown): number | null { + if (!(error instanceof ApiError)) return null + let body: unknown + try { + body = JSON.parse(error.message) + } catch { + return null + } + for (const detail of toArray(toRecord(toRecord(body).error).details)) { + const record = toRecord(detail) + if (record['@type'] !== RETRY_INFO_TYPE || typeof record.retryDelay !== 'string') continue + const seconds = /^(\d+(?:\.\d+)?)s$/.exec(record.retryDelay) + if (seconds) return Number(seconds[1]) * 1000 + } + return null +} + /** * Ensures a value is a valid object for Gemini's functionResponse.response field. * Gemini's API requires functionResponse.response to be a google.protobuf.Struct, diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts index dcb090786c2..ab2df3c5af7 100644 --- a/apps/sim/providers/retry.test.ts +++ b/apps/sim/providers/retry.test.ts @@ -88,6 +88,38 @@ describe('fetchWithProviderRetry', () => { expect(await response.text()).toContain('No credit') }) + it('retries a rate limit whose body is too large to classify', async () => { + const send = vi + .fn() + .mockImplementationOnce(reply(429, 'x'.repeat(256 * 1024))) + .mockImplementation(reply(200)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + }) + + it('stops reading a rate-limit body that never finishes once the caller aborts', async () => { + const controller = new AbortController() + const endless = new ReadableStream({ + start(stream) { + stream.enqueue(new TextEncoder().encode('{"error":')) + }, + }) + const send = vi.fn().mockResolvedValue(new Response(endless, { status: 429 })) + + const outcome = fetchWithProviderRetry(send, { + logger, + label: 'OpenAI', + abortSignal: controller.signal, + }).catch((error: unknown) => error) + await vi.advanceTimersByTimeAsync(0) + controller.abort() + + expect(await outcome).toMatchObject({ name: 'AbortError' }) + expect(send).toHaveBeenCalledTimes(1) + }) + it('obeys the server when it says a 5xx must not be retried', async () => { const send = vi.fn().mockImplementation(reply(500, '', { 'x-should-retry': 'false' })) @@ -216,4 +248,19 @@ describe('withProviderRetry', () => { ) expect(operation).toHaveBeenCalledTimes(1) }) + + it('waits the delay the caller reads off the SDK error before replaying', async () => { + const operation = vi.fn().mockRejectedValueOnce(sdkError(429)).mockResolvedValue('answer') + + const pending = withProviderRetry(operation, { + logger, + label: 'Gemini', + retryAfterMs: () => 7000, + }) + await vi.advanceTimersByTimeAsync(6900) + expect(operation).toHaveBeenCalledTimes(1) + + await vi.advanceTimersByTimeAsync(200) + await expect(pending).resolves.toBe('answer') + }) }) diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts index a3802fee064..080fc4e25a6 100644 --- a/apps/sim/providers/retry.ts +++ b/apps/sim/providers/retry.ts @@ -31,6 +31,8 @@ export interface ProviderRetryOptions { /** Provider name for the retry log line. */ label: string abortSignal?: AbortSignal + /** The delay an SDK error asks for, for SDKs that carry it in the error rather than headers. */ + retryAfterMs?: (error: unknown) => number | null } /** Statuses every vendor SDK treats as transient: timeout, lock conflict, rate limit, server fault. */ @@ -53,21 +55,41 @@ function isRetryableTransportFailure(error: unknown): boolean { return error instanceof TypeError || isRetryableInfrastructureError(error) } -/** Reads a clone, so the caller can still read the body of a response that is not retried. */ -async function isQuotaExhaustedResponse(response: Response): Promise { - const body = await readResponseTextWithLimit(response.clone(), { - maxBytes: DEFAULT_MAX_ERROR_BODY_BYTES, - label: 'provider error response', - }).catch(() => '') - return isQuotaExhaustionBody(body) -} - -async function shouldRetryResponse(response: Response): Promise { +/** + * The response to hand back when a failed one must not be replayed, or `null` to replay it. + * + * Telling a spent balance from a rate limit means reading a 429's body. It is read + * directly, bounded and under the caller's signal, and a 429 that is not replayed comes + * back rebuilt from that text. A `clone()` would tee the stream, and cancelling one branch + * of a tee settles only once the other is cancelled too — an oversized body would hang. + */ +async function finalResponse( + response: Response, + abortSignal: AbortSignal | undefined +): Promise { const directive = response.headers.get('x-should-retry') - if (directive === 'true') return true - if (directive === 'false') return false - if (!isRetryableProviderStatus(response.status)) return false - return response.status !== 429 || !(await isQuotaExhaustedResponse(response)) + if (directive === 'true') return null + if (directive === 'false' || !isRetryableProviderStatus(response.status)) return response + if (response.status !== 429) return null + + let body: string + try { + body = await readResponseTextWithLimit(response, { + maxBytes: DEFAULT_MAX_ERROR_BODY_BYTES, + label: 'provider error response', + signal: abortSignal, + }) + } catch { + abortSignal?.throwIfAborted() + /** An unreadable or oversized body cannot be a quota error, so it stays a rate limit. */ + return null + } + if (!isQuotaExhaustionBody(body)) return null + return new Response(body, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }) } async function waitBeforeRetry( @@ -109,8 +131,10 @@ export async function fetchWithProviderRetry( await waitBeforeRetry(attempt, null, getErrorMessage(error), options) continue } - if (response.ok || !canRetry || !(await shouldRetryResponse(response))) return response - await consumeOrCancelBody(response) + if (response.ok || !canRetry) return response + const final = await finalResponse(response, options.abortSignal) + if (final) return final + if (!response.bodyUsed) await consumeOrCancelBody(response) await waitBeforeRetry( attempt, providerRetryAfterMs(response.headers), @@ -138,7 +162,12 @@ export async function withProviderRetry( if (attempt > PROVIDER_MAX_RETRIES || options.abortSignal?.aborted || !retryable) { throw error } - await waitBeforeRetry(attempt, null, getErrorMessage(error), options) + await waitBeforeRetry( + attempt, + options.retryAfterMs?.(error) ?? null, + getErrorMessage(error), + options + ) } } } From f6648f649fc57bcf9ce023ec2b70a555300ac0d4 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 12:17:15 -0700 Subject: [PATCH 3/7] improvement(providers): surface long server delays, narrow transport retries, pin SDK retry budgets --- apps/sim/providers/anthropic/index.ts | 6 +- apps/sim/providers/azure-anthropic/index.ts | 2 + apps/sim/providers/azure-openai/index.ts | 2 + .../conversation-generation-coverage.test.ts | 5 +- apps/sim/providers/retry.test.ts | 58 +++++++++-- apps/sim/providers/retry.ts | 98 +++++++++++++------ apps/sim/providers/transport.ts | 6 +- apps/sim/providers/typesafe/transport.ts | 10 +- 8 files changed, 143 insertions(+), 44 deletions(-) diff --git a/apps/sim/providers/anthropic/index.ts b/apps/sim/providers/anthropic/index.ts index 80129eb3cfc..7beb53c79e5 100644 --- a/apps/sim/providers/anthropic/index.ts +++ b/apps/sim/providers/anthropic/index.ts @@ -4,6 +4,7 @@ import type { StreamingExecution } from '@/executor/types' import { executeAnthropicProviderRequest } from '@/providers/anthropic/core' import { getCachedProviderClient } from '@/providers/client-cache' import { getProviderDefaultModel, getProviderModels } from '@/providers/models' +import { PROVIDER_MAX_RETRIES } from '@/providers/transport' import type { ProviderConfig, ProviderRequest, ProviderResponse } from '@/providers/types' const logger = createLogger('AnthropicProvider') @@ -24,7 +25,10 @@ export const anthropicProvider: ProviderConfig = { providerLabel: 'Anthropic', createClient: (apiKey) => { const cacheKey = `anthropic::${apiKey}` - return getCachedProviderClient(cacheKey, () => new Anthropic({ apiKey })) + return getCachedProviderClient( + cacheKey, + () => new Anthropic({ apiKey, maxRetries: PROVIDER_MAX_RETRIES }) + ) }, logger, }) diff --git a/apps/sim/providers/azure-anthropic/index.ts b/apps/sim/providers/azure-anthropic/index.ts index 8431204cd73..7737dcd4abc 100644 --- a/apps/sim/providers/azure-anthropic/index.ts +++ b/apps/sim/providers/azure-anthropic/index.ts @@ -6,6 +6,7 @@ import type { StreamingExecution } from '@/executor/types' import { executeAnthropicProviderRequest } from '@/providers/anthropic/core' import { getCachedProviderClient } from '@/providers/client-cache' import { getProviderDefaultModel, getProviderModels } from '@/providers/models' +import { PROVIDER_MAX_RETRIES } from '@/providers/transport' import type { ProviderConfig, ProviderRequest, ProviderResponse } from '@/providers/types' const logger = createLogger('AzureAnthropicProvider') @@ -79,6 +80,7 @@ export const azureAnthropicProvider: ProviderConfig = { new Anthropic({ baseURL, apiKey, + maxRetries: PROVIDER_MAX_RETRIES, ...(pinnedFetch ? { fetch: pinnedFetch } : {}), defaultHeaders: { 'anthropic-version': anthropicVersion, diff --git a/apps/sim/providers/azure-openai/index.ts b/apps/sim/providers/azure-openai/index.ts index 180410c9bb4..e4092d00163 100644 --- a/apps/sim/providers/azure-openai/index.ts +++ b/apps/sim/providers/azure-openai/index.ts @@ -48,6 +48,7 @@ import { createStreamingExecution } from '@/providers/streaming-execution' import { isAbortError, parseToolArguments } from '@/providers/streaming-tool-loop-shared' import { adaptOpenAIChatToolSchema } from '@/providers/tool-schema-adapter' import { enrichLastModelSegmentFromChatCompletions } from '@/providers/trace-enrichment' +import { openAICompatTransport } from '@/providers/transport' import type { FunctionCallResponse, ProviderConfig, @@ -97,6 +98,7 @@ async function executeChatCompletionsRequest( apiKey: request.apiKey!, apiVersion: azureApiVersion, endpoint: azureEndpoint, + ...openAICompatTransport(), ...(pinnedFetch ? { fetch: pinnedFetch } : {}), }) diff --git a/apps/sim/providers/conversation-generation-coverage.test.ts b/apps/sim/providers/conversation-generation-coverage.test.ts index fda584bdf34..37ef9075e02 100644 --- a/apps/sim/providers/conversation-generation-coverage.test.ts +++ b/apps/sim/providers/conversation-generation-coverage.test.ts @@ -35,7 +35,10 @@ function preparedPayload(node: ts.Node): boolean { * to compact history. The nearest declaration decides, so an inner shadow must be prepared too. */ function preparedBinding(identifier: ts.Identifier): boolean { - const declares = (name: ts.BindingName) => name.getText() === identifier.text + const declares = (name: ts.BindingName): boolean => + ts.isIdentifier(name) + ? name.text === identifier.text + : name.elements.some((element) => !ts.isOmittedExpression(element) && declares(element.name)) for (let scope = identifier.parent; scope; scope = scope.parent) { if (ts.isFunctionLike(scope) && scope.parameters.some((parameter) => declares(parameter.name))) return false diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts index ab2df3c5af7..fe8ede1c634 100644 --- a/apps/sim/providers/retry.test.ts +++ b/apps/sim/providers/retry.test.ts @@ -154,19 +154,42 @@ describe('fetchWithProviderRetry', () => { await pending }) - it('retries a dropped connection', async () => { - const send = vi.fn().mockRejectedValueOnce(connectionReset()).mockImplementation(reply(200)) + /** A daily quota or provider-wide pause outlasts any backoff; the fallback model should run now. */ + it('does not wait out a server delay longer than the retry window', async () => { + const send = vi.fn().mockImplementation(reply(429, '', { 'retry-after': '120' })) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(429) + expect(send).toHaveBeenCalledTimes(1) + }) + + it.each([ + ['Bun connection reset', connectionReset()], + [ + 'Bun refused connection', + Object.assign(new Error('Unable to connect'), { code: 'ConnectionRefused' }), + ], + [ + 'Node network failure', + new TypeError('fetch failed', { cause: Object.assign(new Error(), { code: 'ECONNRESET' }) }), + ], + ])('retries a %s', async (_, failure) => { + const send = vi.fn().mockRejectedValueOnce(failure).mockImplementation(reply(200)) const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) expect(response.status).toBe(200) }) - it('does not retry an unrecognised failure', async () => { - const send = vi.fn().mockRejectedValue(new Error('body serialization failed')) + it.each([ + ['an unrecognised failure', new Error('body serialization failed')], + ['a TypeError from building the request', new TypeError('Invalid URL')], + ])('does not retry %s', async (_, failure) => { + const send = vi.fn().mockRejectedValue(failure) - await expect(settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' }))).rejects.toThrow( - 'body serialization failed' + await expect(settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' }))).rejects.toBe( + failure ) expect(send).toHaveBeenCalledTimes(1) }) @@ -263,4 +286,27 @@ describe('withProviderRetry', () => { await vi.advanceTimersByTimeAsync(200) await expect(pending).resolves.toBe('answer') }) + + it('surfaces an SDK error at once when its delay is longer than the retry window', async () => { + const quota = sdkError(429) + const operation = vi.fn().mockRejectedValue(quota) + + await expect( + settle( + withProviderRetry(operation, { logger, label: 'Gemini', retryAfterMs: () => 3_600_000 }) + ) + ).rejects.toBe(quota) + expect(operation).toHaveBeenCalledTimes(1) + }) + + /** A malformed request fails inside the SDK before any fetch; replaying cannot fix it. */ + it('does not retry a TypeError the SDK raised while building the request', async () => { + const bug = new TypeError("Cannot use 'in' operator to search for 'functionDeclarations'") + const operation = vi.fn().mockRejectedValue(bug) + + await expect(settle(withProviderRetry(operation, { logger, label: 'Gemini' }))).rejects.toBe( + bug + ) + expect(operation).toHaveBeenCalledTimes(1) + }) }) diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts index 080fc4e25a6..8587f4f96f4 100644 --- a/apps/sim/providers/retry.ts +++ b/apps/sim/providers/retry.ts @@ -1,8 +1,26 @@ +/** + * Retry policy for provider calls that no vendor SDK retries for us: the raw-fetch + * Responses path and SDKs whose own retry is unusable. + * + * It mirrors what `openai`, `@anthropic-ai/sdk`, and the AI SDK agree on, so a model + * behaves the same whichever transport serves it: {@link PROVIDER_MAX_RETRIES} replays, + * jittered exponential backoff, the server's `x-should-retry` and `retry-after-ms` / + * `retry-after` obeyed, and a caller's abort never retried. Only the wait for response + * headers is covered — once a body is streaming, nothing is replayed. + * + * Two deliberate divergences, both so a block's fallback model runs instead of waiting on + * a failure that will not clear: a 429 reporting an exhausted balance is not retried, and + * neither is a failure whose requested delay exceeds {@link MAX_RETRY_AFTER_MS}. + * + * @packageDocumentation + */ + import type { Logger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' import { interruptibleSleep } from '@sim/utils/helpers' import { isRecordLike } from '@sim/utils/object' import { backoffWithJitter, parseRetryAfter } from '@sim/utils/retry' +import { truncate } from '@sim/utils/string' import { isQuotaExhaustionBody } from '@/lib/core/errors/provider-quota' import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure' import { @@ -13,18 +31,17 @@ import { import { PROVIDER_MAX_RETRIES } from '@/providers/transport' /** - * Retry policy for provider calls that no vendor SDK retries for us: the raw-fetch - * Responses path and SDKs whose own retry is unusable. - * - * It mirrors what `openai`, `@anthropic-ai/sdk`, and the AI SDK agree on, so a model - * behaves the same whichever transport serves it: {@link PROVIDER_MAX_RETRIES} replays, - * jittered exponential backoff, the server's `x-should-retry` and `retry-after-ms` / - * `retry-after` obeyed, and a caller's abort never retried. Only the wait for response - * headers is covered — once a body is streaming, nothing is replayed. - * - * One deliberate divergence: a 429 reporting an exhausted balance is not retried. The - * SDKs replay it, but a spent account does not reopen within a backoff window. + * The longest server-requested delay waited out before a replay. A longer one — a daily + * quota, a provider-wide pause — outlasts the retry budget, so the failure surfaces at once. + * Matches the backoff ceiling. */ +const MAX_RETRY_AFTER_MS = 30_000 + +/** Node's `fetch` rejects with this `TypeError` message when no response arrived. */ +const FETCH_FAILED_MESSAGES = new Set(['fetch failed', 'failed to fetch']) + +/** Bun's `fetch` rejects with an `Error` carrying one of these codes when no response arrived. */ +const BUN_CONNECTION_ERROR_CODES = new Set(['ConnectionRefused', 'ConnectionClosed']) export interface ProviderRetryOptions { logger: Logger @@ -44,15 +61,26 @@ export function isRetryableProviderStatus(status: number): boolean { export function providerRetryAfterMs(headers: Headers): number | null { const precise = Number.parseFloat(headers.get('retry-after-ms') ?? '') if (Number.isFinite(precise) && precise >= 0) return precise - return parseRetryAfter(headers.get('retry-after')) + return parseRetryAfter(headers.get('retry-after'), Number.POSITIVE_INFINITY) +} + +/** Whether a requested delay is short enough to wait out; no delay means backoff decides. */ +export function isWithinRetryWindow(retryAfterMs: number | null): boolean { + return retryAfterMs === null || retryAfterMs <= MAX_RETRY_AFTER_MS } /** - * A dropped or refused connection: Node's `fetch` raises a `TypeError`, Bun's an `Error` - * carrying the syscall code. + * A request that never got a response. Node rejects with `TypeError('fetch failed')` carrying + * the socket error as `cause`; Bun with an `Error` carrying the connection code. Any other + * `TypeError` is a malformed request, which a replay cannot fix. */ function isRetryableTransportFailure(error: unknown): boolean { - return error instanceof TypeError || isRetryableInfrastructureError(error) + if (error instanceof TypeError) { + return error.cause !== undefined && FETCH_FAILED_MESSAGES.has(error.message.toLowerCase()) + } + const code = isRecordLike(error) ? error.code : undefined + if (typeof code === 'string' && BUN_CONNECTION_ERROR_CODES.has(code)) return true + return isRetryableInfrastructureError(error) } /** @@ -63,14 +91,15 @@ function isRetryableTransportFailure(error: unknown): boolean { * back rebuilt from that text. A `clone()` would tee the stream, and cancelling one branch * of a tee settles only once the other is cancelled too — an oversized body would hang. */ -async function finalResponse( +async function nonRetryableResponse( response: Response, abortSignal: AbortSignal | undefined ): Promise { const directive = response.headers.get('x-should-retry') - if (directive === 'true') return null - if (directive === 'false' || !isRetryableProviderStatus(response.status)) return response - if (response.status !== 429) return null + if (directive === 'false') return response + if (directive !== 'true' && !isRetryableProviderStatus(response.status)) return response + if (!isWithinRetryWindow(providerRetryAfterMs(response.headers))) return response + if (directive === 'true' || response.status !== 429) return null let body: string try { @@ -95,11 +124,12 @@ async function finalResponse( async function waitBeforeRetry( attempt: number, retryAfterMs: number | null, - reason: string, + failure: Record, { logger, label, abortSignal }: ProviderRetryOptions ): Promise { const delayMs = backoffWithJitter(attempt, retryAfterMs) - logger.warn(`${label} request failed (${reason}); retrying`, { + logger.warn(`${label} request failed; retrying`, { + ...failure, attempt, maxRetries: PROVIDER_MAX_RETRIES, delayMs: Math.round(delayMs), @@ -108,6 +138,11 @@ async function waitBeforeRetry( abortSignal?.throwIfAborted() } +/** Bounded, so an SDK that folds the whole error body into its message cannot flood the log. */ +function describeError(error: unknown): string { + return truncate(getErrorMessage(error), 200) +} + /** * Sends a raw-fetch provider request, replaying it on a transient failure. * @@ -128,17 +163,17 @@ export async function fetchWithProviderRetry( if (!canRetry || options.abortSignal?.aborted || !isRetryableTransportFailure(error)) { throw error } - await waitBeforeRetry(attempt, null, getErrorMessage(error), options) + await waitBeforeRetry(attempt, null, { error: describeError(error) }, options) continue } if (response.ok || !canRetry) return response - const final = await finalResponse(response, options.abortSignal) + const final = await nonRetryableResponse(response, options.abortSignal) if (final) return final if (!response.bodyUsed) await consumeOrCancelBody(response) await waitBeforeRetry( attempt, providerRetryAfterMs(response.headers), - `HTTP ${response.status}`, + { status: response.status, requestId: response.headers.get('x-request-id') }, options ) } @@ -159,15 +194,16 @@ export async function withProviderRetry( const status = isRecordLike(error) && typeof error.status === 'number' ? error.status : null const retryable = status === null ? isRetryableTransportFailure(error) : isRetryableProviderStatus(status) - if (attempt > PROVIDER_MAX_RETRIES || options.abortSignal?.aborted || !retryable) { + const retryAfterMs = options.retryAfterMs?.(error) ?? null + if ( + attempt > PROVIDER_MAX_RETRIES || + options.abortSignal?.aborted || + !retryable || + !isWithinRetryWindow(retryAfterMs) + ) { throw error } - await waitBeforeRetry( - attempt, - options.retryAfterMs?.(error) ?? null, - getErrorMessage(error), - options - ) + await waitBeforeRetry(attempt, retryAfterMs, { status, error: describeError(error) }, options) } } } diff --git a/apps/sim/providers/transport.ts b/apps/sim/providers/transport.ts index 420fb926de5..d882cfe9936 100644 --- a/apps/sim/providers/transport.ts +++ b/apps/sim/providers/transport.ts @@ -2,8 +2,8 @@ * Transport policy for provider requests, in one place. * * These pin what the vendor SDKs already default to, so an SDK bump cannot silently - * move production behaviour. For 16 of 18 providers this is a no-op. Groq and Cerebras - * are the exception and are marked at {@link PROVIDER_HEADERS_TIMEOUT_MS}. + * move production behaviour. For every client but Groq's and Cerebras's this is a no-op; + * those two are the exception and are marked at {@link PROVIDER_HEADERS_TIMEOUT_MS}. * * What these deliberately do NOT do: bound a stalled stream. Bun's `fetch` is native * and appears to impose a socket-scoped idle wall of roughly 300s, reduced by however @@ -17,7 +17,7 @@ * Time-to-headers budget for a single attempt, matching `openai@7`'s own * `DEFAULT_TIMEOUT`. * - * Behaviour-preserving for the 16 providers already on that client. It is a deliberate + * Behaviour-preserving for the providers already on that client. It is a deliberate * divergence for **Groq and Cerebras**, whose SDKs default to 60s: on a non-streaming * call headers do not arrive until the generation completes, so 60s caps every * generation at a minute and then retries it twice, re-billing. The cost of the raise is diff --git a/apps/sim/providers/typesafe/transport.ts b/apps/sim/providers/typesafe/transport.ts index 9f74d3edc68..658267367f0 100644 --- a/apps/sim/providers/typesafe/transport.ts +++ b/apps/sim/providers/typesafe/transport.ts @@ -2,7 +2,11 @@ import { interruptibleSleep } from '@sim/utils/helpers' import { backoffWithJitter } from '@sim/utils/retry' import { stringifyBoundedJson } from '@/lib/core/utils/bounded-json' import { consumeOrCancelBody, readResponseJsonWithLimit } from '@/lib/core/utils/stream-limits' -import { isRetryableProviderStatus, providerRetryAfterMs } from '@/providers/retry' +import { + isRetryableProviderStatus, + isWithinRetryWindow, + providerRetryAfterMs, +} from '@/providers/retry' import { PROVIDER_HEADERS_TIMEOUT_MS, PROVIDER_MAX_RETRIES } from '@/providers/transport' import type { buildJevBody } from '@/providers/typesafe/schema' @@ -57,8 +61,10 @@ export async function requestJevEvaluation( error instanceof TypeSafeHttpError ? isRetryableProviderStatus(error.status) : !response || timeout.aborted || error instanceof TypeError - if (!retryable || attempt >= PROVIDER_MAX_RETRIES) throw error const retryAfterMs = error instanceof TypeSafeHttpError ? error.retryAfterMs : null + if (!retryable || attempt >= PROVIDER_MAX_RETRIES || !isWithinRetryWindow(retryAfterMs)) { + throw error + } await interruptibleSleep(backoffWithJitter(attempt + 1, retryAfterMs), abortSignal) } } From 263308884e09d81754a83b326da8750816b24e65 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 12:24:32 -0700 Subject: [PATCH 4/7] fix(providers): classify transport failures by socket error code, not message --- apps/sim/providers/retry.test.ts | 12 +++++++++++- apps/sim/providers/retry.ts | 13 ++++--------- 2 files changed, 15 insertions(+), 10 deletions(-) diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts index fe8ede1c634..fd7327f730c 100644 --- a/apps/sim/providers/retry.test.ts +++ b/apps/sim/providers/retry.test.ts @@ -171,9 +171,19 @@ describe('fetchWithProviderRetry', () => { Object.assign(new Error('Unable to connect'), { code: 'ConnectionRefused' }), ], [ - 'Node network failure', + 'Node connection reset', new TypeError('fetch failed', { cause: Object.assign(new Error(), { code: 'ECONNRESET' }) }), ], + [ + 'Node refused connection', + new TypeError('fetch failed', { + cause: Object.assign(new AggregateError([]), { code: 'ECONNREFUSED' }), + }), + ], + [ + 'TypeError caused by a socket error, whatever its message', + new TypeError('network error', { cause: Object.assign(new Error(), { code: 'ECONNRESET' }) }), + ], ])('retries a %s', async (_, failure) => { const send = vi.fn().mockRejectedValueOnce(failure).mockImplementation(reply(200)) diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts index 8587f4f96f4..3dc5f535904 100644 --- a/apps/sim/providers/retry.ts +++ b/apps/sim/providers/retry.ts @@ -37,9 +37,6 @@ import { PROVIDER_MAX_RETRIES } from '@/providers/transport' */ const MAX_RETRY_AFTER_MS = 30_000 -/** Node's `fetch` rejects with this `TypeError` message when no response arrived. */ -const FETCH_FAILED_MESSAGES = new Set(['fetch failed', 'failed to fetch']) - /** Bun's `fetch` rejects with an `Error` carrying one of these codes when no response arrived. */ const BUN_CONNECTION_ERROR_CODES = new Set(['ConnectionRefused', 'ConnectionClosed']) @@ -70,14 +67,12 @@ export function isWithinRetryWindow(retryAfterMs: number | null): boolean { } /** - * A request that never got a response. Node rejects with `TypeError('fetch failed')` carrying - * the socket error as `cause`; Bun with an `Error` carrying the connection code. Any other - * `TypeError` is a malformed request, which a replay cannot fix. + * A request that never got a response, recognised by its socket error code rather than its + * message: Node rejects with `TypeError('fetch failed')` carrying the code on its `cause`, + * Bun with an `Error` carrying it directly. An error with no such code — a `TypeError` from + * a malformed request, say — is not a network fault, and a replay cannot fix it. */ function isRetryableTransportFailure(error: unknown): boolean { - if (error instanceof TypeError) { - return error.cause !== undefined && FETCH_FAILED_MESSAGES.has(error.message.toLowerCase()) - } const code = isRecordLike(error) ? error.code : undefined if (typeof code === 'string' && BUN_CONNECTION_ERROR_CODES.has(code)) return true return isRetryableInfrastructureError(error) From 8c36f1c859433c0d65d6e08cec1e636f7c0f28ae Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 12:31:54 -0700 Subject: [PATCH 5/7] fix(providers): cancel a discarded retry body instead of draining it --- apps/sim/providers/retry.test.ts | 16 ++++++++++++++++ apps/sim/providers/retry.ts | 4 ++-- 2 files changed, 18 insertions(+), 2 deletions(-) diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts index fd7327f730c..8de356f5b50 100644 --- a/apps/sim/providers/retry.test.ts +++ b/apps/sim/providers/retry.test.ts @@ -88,6 +88,22 @@ describe('fetchWithProviderRetry', () => { expect(await response.text()).toContain('No credit') }) + it('discards a retryable response whose body never finishes instead of waiting on it', async () => { + const endless = new ReadableStream({ + start(stream) { + stream.enqueue(new TextEncoder().encode('')) + }, + }) + const send = vi + .fn() + .mockResolvedValueOnce(new Response(endless, { status: 503 })) + .mockImplementation(reply(200)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + }) + it('retries a rate limit whose body is too large to classify', async () => { const send = vi .fn() diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts index 3dc5f535904..e33e81add62 100644 --- a/apps/sim/providers/retry.ts +++ b/apps/sim/providers/retry.ts @@ -24,7 +24,6 @@ import { truncate } from '@sim/utils/string' import { isQuotaExhaustionBody } from '@/lib/core/errors/provider-quota' import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure' import { - consumeOrCancelBody, DEFAULT_MAX_ERROR_BODY_BYTES, readResponseTextWithLimit, } from '@/lib/core/utils/stream-limits' @@ -164,7 +163,8 @@ export async function fetchWithProviderRetry( if (response.ok || !canRetry) return response const final = await nonRetryableResponse(response, options.abortSignal) if (final) return final - if (!response.bodyUsed) await consumeOrCancelBody(response) + /** Cancelled, not drained, as the SDKs do: draining a body that stalls would hang the loop. */ + if (!response.bodyUsed) await response.body?.cancel().catch(() => {}) await waitBeforeRetry( attempt, providerRetryAfterMs(response.headers), From b6b8b9d03956d76b08b6b07679d4829457185519 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 12:40:29 -0700 Subject: [PATCH 6/7] fix(providers): bound how long a 429 body may take to classify --- apps/sim/providers/retry.test.ts | 16 ++++++++++++++++ apps/sim/providers/retry.ts | 18 ++++++++++++++---- 2 files changed, 30 insertions(+), 4 deletions(-) diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts index 8de356f5b50..d517b00a952 100644 --- a/apps/sim/providers/retry.test.ts +++ b/apps/sim/providers/retry.test.ts @@ -104,6 +104,22 @@ describe('fetchWithProviderRetry', () => { expect(response.status).toBe(200) }) + it('treats a rate limit whose body never finishes as a rate limit, without a caller signal', async () => { + const endless = new ReadableStream({ + start(stream) { + stream.enqueue(new TextEncoder().encode('{"error":')) + }, + }) + const send = vi + .fn() + .mockResolvedValueOnce(new Response(endless, { status: 429 })) + .mockImplementation(reply(200)) + + const response = await settle(fetchWithProviderRetry(send, { logger, label: 'OpenAI' })) + + expect(response.status).toBe(200) + }) + it('retries a rate limit whose body is too large to classify', async () => { const send = vi .fn() diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts index e33e81add62..daa4d72afcf 100644 --- a/apps/sim/providers/retry.ts +++ b/apps/sim/providers/retry.ts @@ -36,6 +36,12 @@ import { PROVIDER_MAX_RETRIES } from '@/providers/transport' */ const MAX_RETRY_AFTER_MS = 30_000 +/** + * How long a 429's body may take to arrive for quota classification. A real error body comes + * with the headers; one that stalls past this is treated as an ordinary rate limit. + */ +const QUOTA_BODY_READ_TIMEOUT_MS = 5_000 + /** Bun's `fetch` rejects with an `Error` carrying one of these codes when no response arrived. */ const BUN_CONNECTION_ERROR_CODES = new Set(['ConnectionRefused', 'ConnectionClosed']) @@ -81,8 +87,8 @@ function isRetryableTransportFailure(error: unknown): boolean { * The response to hand back when a failed one must not be replayed, or `null` to replay it. * * Telling a spent balance from a rate limit means reading a 429's body. It is read - * directly, bounded and under the caller's signal, and a 429 that is not replayed comes - * back rebuilt from that text. A `clone()` would tee the stream, and cancelling one branch + * directly, bounded in size and time and under the caller's signal, and a 429 that is not + * replayed comes back rebuilt from that text. A `clone()` would tee the stream, and cancelling one branch * of a tee settles only once the other is cancelled too — an oversized body would hang. */ async function nonRetryableResponse( @@ -95,17 +101,21 @@ async function nonRetryableResponse( if (!isWithinRetryWindow(providerRetryAfterMs(response.headers))) return response if (directive === 'true' || response.status !== 429) return null + const deadline = new AbortController() + const timer = setTimeout(() => deadline.abort(), QUOTA_BODY_READ_TIMEOUT_MS) let body: string try { body = await readResponseTextWithLimit(response, { maxBytes: DEFAULT_MAX_ERROR_BODY_BYTES, label: 'provider error response', - signal: abortSignal, + signal: abortSignal ? AbortSignal.any([abortSignal, deadline.signal]) : deadline.signal, }) } catch { abortSignal?.throwIfAborted() - /** An unreadable or oversized body cannot be a quota error, so it stays a rate limit. */ + /** An unreadable, oversized, or stalled body cannot be a quota error; it stays a rate limit. */ return null + } finally { + clearTimeout(timer) } if (!isQuotaExhaustionBody(body)) return null return new Response(body, { From e180404654e364141ab8c459f0dfef75997d1a0e Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 28 Sep 2026 12:47:28 -0700 Subject: [PATCH 7/7] fix(providers): limit transport retries to network socket codes --- .../sim/lib/core/errors/retryable-infrastructure.ts | 13 +++++++++++++ apps/sim/providers/retry.test.ts | 5 +++++ apps/sim/providers/retry.ts | 4 ++-- 3 files changed, 20 insertions(+), 2 deletions(-) diff --git a/apps/sim/lib/core/errors/retryable-infrastructure.ts b/apps/sim/lib/core/errors/retryable-infrastructure.ts index 8513969dd8b..73051cb11c0 100644 --- a/apps/sim/lib/core/errors/retryable-infrastructure.ts +++ b/apps/sim/lib/core/errors/retryable-infrastructure.ts @@ -97,6 +97,19 @@ export function isRetryableInfrastructureError(error: unknown): boolean { return Boolean(describeRetryableInfrastructureError(error)) } +/** + * A network-level failure only — a dropped, refused, or timed-out socket anywhere in the + * cause chain. Narrower than {@link isRetryableInfrastructureError}: database and application + * codes are excluded, so a caller replaying a billed request never mistakes one for a socket. + */ +export function isRetryableNetworkError(error: unknown): boolean { + return getErrorChain(error).some( + (candidate) => + (typeof candidate.code === 'string' && RETRYABLE_NETWORK_ERROR_CODES.has(candidate.code)) || + (typeof candidate.errno === 'string' && RETRYABLE_NETWORK_ERROR_CODES.has(candidate.errno)) + ) +} + /** * A retryable infrastructure failure raised strictly BEFORE the guarded * operation performed any effect (no workflow block ran, no mutation diff --git a/apps/sim/providers/retry.test.ts b/apps/sim/providers/retry.test.ts index d517b00a952..919997b0cea 100644 --- a/apps/sim/providers/retry.test.ts +++ b/apps/sim/providers/retry.test.ts @@ -227,6 +227,11 @@ describe('fetchWithProviderRetry', () => { it.each([ ['an unrecognised failure', new Error('body serialization failed')], ['a TypeError from building the request', new TypeError('Invalid URL')], + [ + 'an application error code', + Object.assign(new Error('overloaded'), { code: 'RESOURCE_EXHAUSTED' }), + ], + ['a database error code', Object.assign(new Error('admin shutdown'), { code: '57P01' })], ])('does not retry %s', async (_, failure) => { const send = vi.fn().mockRejectedValue(failure) diff --git a/apps/sim/providers/retry.ts b/apps/sim/providers/retry.ts index daa4d72afcf..846bbed02e2 100644 --- a/apps/sim/providers/retry.ts +++ b/apps/sim/providers/retry.ts @@ -22,7 +22,7 @@ import { isRecordLike } from '@sim/utils/object' import { backoffWithJitter, parseRetryAfter } from '@sim/utils/retry' import { truncate } from '@sim/utils/string' import { isQuotaExhaustionBody } from '@/lib/core/errors/provider-quota' -import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure' +import { isRetryableNetworkError } from '@/lib/core/errors/retryable-infrastructure' import { DEFAULT_MAX_ERROR_BODY_BYTES, readResponseTextWithLimit, @@ -80,7 +80,7 @@ export function isWithinRetryWindow(retryAfterMs: number | null): boolean { function isRetryableTransportFailure(error: unknown): boolean { const code = isRecordLike(error) ? error.code : undefined if (typeof code === 'string' && BUN_CONNECTION_ERROR_CODES.has(code)) return true - return isRetryableInfrastructureError(error) + return isRetryableNetworkError(error) } /**