Skip to content
19 changes: 19 additions & 0 deletions apps/sim/lib/core/errors/provider-quota.ts
Original file line number Diff line number Diff line change
@@ -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
}
}
13 changes: 13 additions & 0 deletions apps/sim/lib/core/errors/retryable-infrastructure.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
21 changes: 1 addition & 20 deletions apps/sim/lib/embeddings/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<string> {
try {
Expand Down
6 changes: 5 additions & 1 deletion apps/sim/providers/anthropic/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand All @@ -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,
})
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/providers/azure-anthropic/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down Expand Up @@ -79,6 +80,7 @@ export const azureAnthropicProvider: ProviderConfig = {
new Anthropic({
baseURL,
apiKey,
maxRetries: PROVIDER_MAX_RETRIES,
...(pinnedFetch ? { fetch: pinnedFetch } : {}),
defaultHeaders: {
'anthropic-version': anthropicVersion,
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/providers/azure-openai/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -97,6 +98,7 @@ async function executeChatCompletionsRequest(
apiKey: request.apiKey!,
apiVersion: azureApiVersion,
endpoint: azureEndpoint,
...openAICompatTransport(),
...(pinnedFetch ? { fetch: pinnedFetch } : {}),
})

Expand Down
44 changes: 43 additions & 1 deletion apps/sim/providers/conversation-generation-coverage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,55 @@ 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) &&
node.expression.expression.getText() === 'prepareConversationGeneration'
)
}

/**
* 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): 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
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
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
}

describe('provider generation context coverage', () => {
it('guards every model SDK send, shared stream callback, retry and finalizer', () => {
const uncovered: string[] = []
Expand Down Expand Up @@ -60,7 +102,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) {
Expand Down
25 changes: 24 additions & 1 deletion apps/sim/providers/gemini/core.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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)
Expand Down
64 changes: 40 additions & 24 deletions apps/sim/providers/gemini/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,11 +32,13 @@ import {
ensureStructResponse,
extractAllFunctionCallParts,
extractTextContent,
geminiRetryDelayMs,
mapToThinkingBudget,
mapToThinkingLevel,
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'
Expand Down Expand Up @@ -962,6 +964,12 @@ export async function executeGeminiRequest(
}

const logger = createLogger(providerType === 'google' ? 'GoogleProvider' : 'VertexProvider')
const retryOptions = {
logger,
label: providerType === 'google' ? 'Gemini' : 'Vertex AI',
abortSignal: request.abortSignal,
retryAfterMs: geminiRetryDelayMs,
}

logger.info(`Preparing ${providerType} Gemini request`, {
model,
Expand Down Expand Up @@ -1148,12 +1156,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

Expand Down Expand Up @@ -1202,12 +1212,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),
Comment thread
waleedlatif1 marked this conversation as resolved.
retryOptions
)
if (!extractAllFunctionCallParts(response.candidates?.[0]).length) {
await captureProviderConversationStep(
Expand Down Expand Up @@ -1256,12 +1268,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(
Expand Down Expand Up @@ -1375,12 +1389,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(
Expand Down
27 changes: 18 additions & 9 deletions apps/sim/providers/gemini/streaming-tool-loop.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,9 @@ import {
cleanSchemaForGemini,
convertUsageMetadata,
ensureStructResponse,
geminiRetryDelayMs,
} from '@/providers/google/utils'
import { withProviderRetry } from '@/providers/retry'
import { executeProviderTool } from '@/providers/runtime-context'
import type { AgentStreamEvent, ToolCallEndStatus } from '@/providers/stream-events'
import {
Expand Down Expand Up @@ -309,15 +311,22 @@ 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,
retryAfterMs: geminiRetryDelayMs,
}
)

const drained = await drainGeminiTurn(
Expand Down
Loading
Loading