diff --git a/changelog.d/fixes/chatcore-legresult-narrowing.md b/changelog.d/fixes/chatcore-legresult-narrowing.md new file mode 100644 index 0000000000..ae8c7f084d --- /dev/null +++ b/changelog.d/fixes/chatcore-legresult-narrowing.md @@ -0,0 +1 @@ +- Restore the API-route typecheck gate: the non-streaming leg result lost its discriminated-union narrowing after the server-owned tool loop reassignment, producing 13 new TS2339 diagnostics in `chatCore.ts`. diff --git a/config/quality/eslint-suppressions.json b/config/quality/eslint-suppressions.json index e10d1f6551..2074040010 100644 --- a/config/quality/eslint-suppressions.json +++ b/config/quality/eslint-suppressions.json @@ -215,11 +215,6 @@ "count": 1 } }, - "open-sse/handlers/chatCore.ts": { - "@typescript-eslint/no-unused-vars": { - "count": 26 - } - }, "open-sse/handlers/chatCore/executorHelpers.ts": { "@typescript-eslint/no-unused-vars": { "count": 1 diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 63395b47ad..96ff6a260d 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -111,7 +111,7 @@ import { import { recoverAnthropicThinkingSignature } from "./chatCore/thinkingSignatureRecovery.ts"; import { runProviderExecutionPipeline } from "./chatCore/providerExecutionPipeline.ts"; import { runNonStreamingProviderLeg } from "./chatCore/nonStreamingProviderLeg.ts"; -import type { ChatCoreErrorResult } from "@/lib/skills/toolLoopTypes.ts"; +import type { NonStreamingProviderLegResult } from "@/lib/skills/toolLoopTypes.ts"; import { applyServerOwnedToolLoopIfNeeded, derivePostInjectionRequestIdentity, @@ -258,10 +258,7 @@ import { resolveResilienceSettings, isStreamRecoveryExplicitlyConfigured, } from "@/lib/resilience/settings"; -import { - classifyProviderError, - PROVIDER_ERROR_TYPES, -} from "../services/errorClassifier.ts"; +import { classifyProviderError, PROVIDER_ERROR_TYPES } from "../services/errorClassifier.ts"; import { updateProviderConnection, getProviderConnectionById } from "@/lib/db/providers"; import { wasRefreshTokenRotated } from "@omniroute/open-sse/services/refreshSerializer.ts"; import { connectionHasExtraKeys } from "../services/apiKeyRotator.ts"; @@ -3626,778 +3623,666 @@ export async function handleChatCore({ let pipelineRecovered = false; if (stream) { - try { - const pipelineOutcome = await runProviderExecutionPipeline({ - policy: { - allowAccountRotation: !managedLease && comboStrategy !== "context-relay", - allowModelFallback: true, - expectedConnectionId: managedLease - ? String(getCurrentConnectionId() || connectionId || "") || undefined - : undefined, - }, - target: { - provider, - requestedModel: effectiveModel, - sourceFormat, - targetFormat, - stream, - }, - connection: { - initialConnectionId: String(getCurrentConnectionId() || connectionId || ""), - getCurrentConnectionId: () => getCurrentConnectionId() || undefined, - getCredentials: () => (credentials || {}) as Record, - replaceCredentials: (next) => { - Object.assign(credentials, next); - }, - onCredentialsRefreshed: async () => {}, - assertManagedLeaseFence: (id) => { - assertManagedLeaseFence(id); - }, - getProviderCredentials, - }, - wire: { - body: translatedBody as Record, - currentModel, - triedModels, - setBodyAndModel: (body, model) => { - translatedBody = body as typeof translatedBody; - currentModel = model; - triedModels.add(model); - }, - }, - state: { - updatePendingStage: (stage, data) => { - updatePendingScope(pendingScope, { stage, ...(data || {}) }); - }, - recordRateLimitHeaders: updateFromHeaders, - recordRateLimitBody: updateFromResponseBody, - writeTerminalStatus, - persistConnectionPatch: updateProviderConnection, - setConnectionRateLimitedUntil: async (id, untilMs) => { - const { setConnectionRateLimitUntil } = await import("@/lib/db/providers"); - setConnectionRateLimitUntil(id, untilMs); - }, - lockModel, - recordAntigravityQuotaState: recordCoreOwnedAntigravityQuotaState, - markAccountSemaphoreBlocked: (key) => { - markAccountSemaphoreBlocked(key, Date.now() + 60_000); - }, - isolateProbeFailures: () => shouldIsolateProbeFailures(), - onCodexScopeRateLimited: async (params) => { - await markCodexScopeRateLimited({ - failedConnectionId: params.failedConnectionId, - model: params.model, - rateLimitedUntil: params.rateLimitedUntil, - credentials: (params.credentials || credentials) as { - connectionId?: string | null; - providerSpecificData?: unknown; - }, - }); - }, - onClearSessionAffinity: () => { - const key = - sessionAffinityKey || - extractSessionAffinityKey(body, clientRawRequest?.headers) || - null; - if (!key) return; - try { - deleteSessionAccountAffinity(key, "codex"); - } catch { - // best-effort - } - }, - onAuditAccountRotation: (params) => { - logAuditEvent({ - action: params.action, - actor: apiKeyInfo?.name || "system", - target: params.newConnectionId, - details: { - failed_connection_id: params.failedConnectionId, - new_connection_id: params.newConnectionId, - attempt: params.attempt, - retry_after_ms: params.retryAfterMs, - }, - }); - }, - }, - sendProviderAttempt: (modelToCall, allowDedup) => executeProviderRequest(modelToCall, allowDedup), - }); - - pipelineRecovered = true; - currentModel = pipelineOutcome.model; - if (pipelineOutcome.kind === "error") { - providerResponse = pipelineOutcome.result.response; - providerUrl = ""; - providerHeaders = normalizeHeaders(pipelineOutcome.result.response.headers); - finalBody = translatedBody; - } else { - const result = { - response: pipelineOutcome.response, - url: pipelineOutcome.url, - headers: pipelineOutcome.headers, - transformedBody: pipelineOutcome.transformedBody, - }; - providerResponse = result.response; - providerUrl = result.url; - providerHeaders = result.headers; - finalBody = providerRequestCapture.body(result.transformedBody); - } - const responseConnectionId = getCurrentConnectionId(); - effectiveServiceTier = resolveEffectiveServiceTier(finalBody); - claudePromptCacheLogMeta = buildClaudePromptCacheLogMeta( - targetFormat, - finalBody, - providerHeaders, - clientRawRequest?.headers - ); - - // Log target request (final request to provider) - reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingScope(pendingScope, { - providerRequest: finalBody, - providerUrl, - stage: "provider_response_started", - }); - // Update rate limiter from response headers (learn limits dynamically) - updateFromHeaders( - provider, - responseConnectionId, - providerResponse.headers, - providerResponse.status, - model - ); - - // Store rate-limit headers for quota saturation signals try { - const { storeRateLimitHeaders } = await import("@/lib/quota/saturationSignals"); - storeRateLimitHeaders( - responseConnectionId, - provider, - providerResponse.headers as Record + const pipelineOutcome = await runProviderExecutionPipeline({ + policy: { + allowAccountRotation: !managedLease && comboStrategy !== "context-relay", + allowModelFallback: true, + expectedConnectionId: managedLease + ? String(getCurrentConnectionId() || connectionId || "") || undefined + : undefined, + }, + target: { + provider, + requestedModel: effectiveModel, + sourceFormat, + targetFormat, + stream, + }, + connection: { + initialConnectionId: String(getCurrentConnectionId() || connectionId || ""), + getCurrentConnectionId: () => getCurrentConnectionId() || undefined, + getCredentials: () => (credentials || {}) as Record, + replaceCredentials: (next) => { + Object.assign(credentials, next); + }, + onCredentialsRefreshed: async () => {}, + assertManagedLeaseFence: (id) => { + assertManagedLeaseFence(id); + }, + getProviderCredentials, + }, + wire: { + body: translatedBody as Record, + currentModel, + triedModels, + setBodyAndModel: (body, model) => { + translatedBody = body as typeof translatedBody; + currentModel = model; + triedModels.add(model); + }, + }, + state: { + updatePendingStage: (stage, data) => { + updatePendingScope(pendingScope, { stage, ...(data || {}) }); + }, + recordRateLimitHeaders: updateFromHeaders, + recordRateLimitBody: updateFromResponseBody, + writeTerminalStatus, + persistConnectionPatch: updateProviderConnection, + setConnectionRateLimitedUntil: async (id, untilMs) => { + const { setConnectionRateLimitUntil } = await import("@/lib/db/providers"); + setConnectionRateLimitUntil(id, untilMs); + }, + lockModel, + recordAntigravityQuotaState: recordCoreOwnedAntigravityQuotaState, + markAccountSemaphoreBlocked: (key) => { + markAccountSemaphoreBlocked(key, Date.now() + 60_000); + }, + isolateProbeFailures: () => shouldIsolateProbeFailures(), + onCodexScopeRateLimited: async (params) => { + await markCodexScopeRateLimited({ + failedConnectionId: params.failedConnectionId, + model: params.model, + rateLimitedUntil: params.rateLimitedUntil, + credentials: (params.credentials || credentials) as { + connectionId?: string | null; + providerSpecificData?: unknown; + }, + }); + }, + onClearSessionAffinity: () => { + const key = + sessionAffinityKey || + extractSessionAffinityKey(body, clientRawRequest?.headers) || + null; + if (!key) return; + try { + deleteSessionAccountAffinity(key, "codex"); + } catch { + // best-effort + } + }, + onAuditAccountRotation: (params) => { + logAuditEvent({ + action: params.action, + actor: apiKeyInfo?.name || "system", + target: params.newConnectionId, + details: { + failed_connection_id: params.failedConnectionId, + new_connection_id: params.newConnectionId, + attempt: params.attempt, + retry_after_ms: params.retryAfterMs, + }, + }); + }, + }, + sendProviderAttempt: (modelToCall, allowDedup) => + executeProviderRequest(modelToCall, allowDedup), + }); + + pipelineRecovered = true; + currentModel = pipelineOutcome.model; + if (pipelineOutcome.kind === "error") { + providerResponse = pipelineOutcome.result.response; + providerUrl = ""; + providerHeaders = normalizeHeaders(pipelineOutcome.result.response.headers); + finalBody = translatedBody; + } else { + const result = { + response: pipelineOutcome.response, + url: pipelineOutcome.url, + headers: pipelineOutcome.headers, + transformedBody: pipelineOutcome.transformedBody, + }; + providerResponse = result.response; + providerUrl = result.url; + providerHeaders = result.headers; + finalBody = providerRequestCapture.body(result.transformedBody); + } + const responseConnectionId = getCurrentConnectionId(); + effectiveServiceTier = resolveEffectiveServiceTier(finalBody); + claudePromptCacheLogMeta = buildClaudePromptCacheLogMeta( + targetFormat, + finalBody, + providerHeaders, + clientRawRequest?.headers ); - } catch { - // fail-open: saturation signal is best-effort - } - } catch (error) { - trackPendingRequest(model, provider, connectionId, false); - if (isManagedLeaseFenceError(error)) return managedLeaseFenceErrorResult(error); - if (isSemaphoreCapacityError(error)) { + + // Log target request (final request to provider) + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); + updatePendingScope(pendingScope, { + providerRequest: finalBody, + providerUrl, + stage: "provider_response_started", + }); + // Update rate limiter from response headers (learn limits dynamically) + updateFromHeaders( + provider, + responseConnectionId, + providerResponse.headers, + providerResponse.status, + model + ); + + // Store rate-limit headers for quota saturation signals + try { + const { storeRateLimitHeaders } = await import("@/lib/quota/saturationSignals"); + storeRateLimitHeaders( + responseConnectionId, + provider, + providerResponse.headers as Record + ); + } catch { + // fail-open: saturation signal is best-effort + } + } catch (error) { + trackPendingRequest(model, provider, connectionId, false); + if (isManagedLeaseFenceError(error)) return managedLeaseFenceErrorResult(error); + if (isSemaphoreCapacityError(error)) { + appendRequestLog({ + model, + provider, + connectionId, + status: `FAILED ${error.code}`, + }).catch(() => {}); + const failureMessage = error.message || "Semaphore timeout"; + persistAttemptLogs({ + status: HTTP_STATUS.RATE_LIMITED, + error: failureMessage, + providerRequest: finalBody || translatedBody, + clientResponse: buildErrorBody(HTTP_STATUS.RATE_LIMITED, failureMessage), + claudeCacheMeta: claudePromptCacheLogMeta, + cacheSource: "upstream", + }); + persistFailureUsage(HTTP_STATUS.RATE_LIMITED, error.code); + const result = stream + ? createStreamingErrorResult(HTTP_STATUS.RATE_LIMITED, failureMessage, error.code) + : createErrorResult(HTTP_STATUS.RATE_LIMITED, failureMessage); + return { + ...result, + errorType: "account_semaphore_capacity", + errorCode: error.code, + }; + } + // abort(reason) can reject with a raw string lacking `name`/`status`; classify + // it through isLocalStreamLifecycleError so it maps to 499 rather than the + // 502 provider-failure default. + const isRequestAborted = isLocalStreamLifecycleError(error); + // #8376: proxyFetch tags unreachable transport failures so they remain + // distinguishable from ordinary provider 5xx responses. + const isProxyUnreachableFailure = + !isRequestAborted && (error as { errorCode?: unknown })?.errorCode === "proxy_unreachable"; + const errorCode = getUpstreamErrorIdentifier(error); + const localRateLimitFailure = localLimiterErrors.getClientSafeLocalRateLimitError(error); + const failureStatus = isRequestAborted + ? 499 + : isProxyUnreachableFailure + ? HTTP_STATUS.BAD_GATEWAY + : localRateLimitFailure + ? localRateLimitFailure.status + : error.name === "TimeoutError" || error.name === "BodyTimeoutError" + ? HTTP_STATUS.GATEWAY_TIMEOUT + : error.status && typeof error.status === "number" + ? error.status + : HTTP_STATUS.BAD_GATEWAY; + const failureMessage = isRequestAborted + ? "Request aborted" + : formatProviderError(localRateLimitFailure ?? error, provider, model, failureStatus); + const upstreamErrorCode = + localRateLimitFailure?.code ?? + (isProxyUnreachableFailure ? "proxy_unreachable" : errorCode); + // Tag our own deadline timeouts (fetch-start TimeoutError / body BodyTimeoutError, + // both surfaced as a 504) as "upstream_timeout" so the cooldown layer can tell a + // slow-but-not-failed request apart from a real provider 5xx. (Antigravity already + // tags its pre-response timeout via the code below.) + const isOwnDeadlineTimeout = + failureStatus === HTTP_STATUS.GATEWAY_TIMEOUT && + (error.name === "TimeoutError" || error.name === "BodyTimeoutError"); + const upstreamErrorType = + upstreamErrorCode === ANTIGRAVITY_PRE_RESPONSE_TIMEOUT_CODE || isOwnDeadlineTimeout + ? "upstream_timeout" + : failureStatus === 401 + ? "authentication_error" + : undefined; appendRequestLog({ model, provider, connectionId, - status: `FAILED ${error.code}`, + status: `FAILED ${failureStatus}`, }).catch(() => {}); - const failureMessage = error.message || "Semaphore timeout"; persistAttemptLogs({ - status: HTTP_STATUS.RATE_LIMITED, + status: failureStatus, error: failureMessage, providerRequest: finalBody || translatedBody, - clientResponse: buildErrorBody(HTTP_STATUS.RATE_LIMITED, failureMessage), + // On a client-abort (AbortError), the client already disconnected before + // we ever got here — this body is what we WOULD have sent, not what was + // actually delivered. Logging it as `clientResponse` is misleading (the + // dashboard reads that field as "what the client received"), so omit it + // for this case; `error` above already records the failure reason. + clientResponse: + error.name === "AbortError" ? undefined : buildErrorBody(failureStatus, failureMessage), claudeCacheMeta: claudePromptCacheLogMeta, cacheSource: "upstream", }); - persistFailureUsage(HTTP_STATUS.RATE_LIMITED, error.code); - const result = stream - ? createStreamingErrorResult(HTTP_STATUS.RATE_LIMITED, failureMessage, error.code) - : createErrorResult(HTTP_STATUS.RATE_LIMITED, failureMessage); - return { - ...result, - errorType: "account_semaphore_capacity", - errorCode: error.code, - }; - } - // abort(reason) can reject with a raw string lacking `name`/`status`; classify - // it through isLocalStreamLifecycleError so it maps to 499 rather than the - // 502 provider-failure default. - const isRequestAborted = isLocalStreamLifecycleError(error); - // #8376: proxyFetch tags unreachable transport failures so they remain - // distinguishable from ordinary provider 5xx responses. - const isProxyUnreachableFailure = - !isRequestAborted && (error as { errorCode?: unknown })?.errorCode === "proxy_unreachable"; - const errorCode = getUpstreamErrorIdentifier(error); - const localRateLimitFailure = localLimiterErrors.getClientSafeLocalRateLimitError(error); - const failureStatus = isRequestAborted - ? 499 - : isProxyUnreachableFailure - ? HTTP_STATUS.BAD_GATEWAY - : localRateLimitFailure - ? localRateLimitFailure.status - : error.name === "TimeoutError" || error.name === "BodyTimeoutError" - ? HTTP_STATUS.GATEWAY_TIMEOUT - : error.status && typeof error.status === "number" - ? error.status - : HTTP_STATUS.BAD_GATEWAY; - const failureMessage = isRequestAborted - ? "Request aborted" - : formatProviderError(localRateLimitFailure ?? error, provider, model, failureStatus); - const upstreamErrorCode = - localRateLimitFailure?.code ?? (isProxyUnreachableFailure ? "proxy_unreachable" : errorCode); - // Tag our own deadline timeouts (fetch-start TimeoutError / body BodyTimeoutError, - // both surfaced as a 504) as "upstream_timeout" so the cooldown layer can tell a - // slow-but-not-failed request apart from a real provider 5xx. (Antigravity already - // tags its pre-response timeout via the code below.) - const isOwnDeadlineTimeout = - failureStatus === HTTP_STATUS.GATEWAY_TIMEOUT && - (error.name === "TimeoutError" || error.name === "BodyTimeoutError"); - const upstreamErrorType = - upstreamErrorCode === ANTIGRAVITY_PRE_RESPONSE_TIMEOUT_CODE || isOwnDeadlineTimeout - ? "upstream_timeout" - : failureStatus === 401 - ? "authentication_error" - : undefined; - appendRequestLog({ - model, - provider, - connectionId, - status: `FAILED ${failureStatus}`, - }).catch(() => {}); - persistAttemptLogs({ - status: failureStatus, - error: failureMessage, - providerRequest: finalBody || translatedBody, - // On a client-abort (AbortError), the client already disconnected before - // we ever got here — this body is what we WOULD have sent, not what was - // actually delivered. Logging it as `clientResponse` is misleading (the - // dashboard reads that field as "what the client received"), so omit it - // for this case; `error` above already records the failure reason. - clientResponse: - error.name === "AbortError" ? undefined : buildErrorBody(failureStatus, failureMessage), - claudeCacheMeta: claudePromptCacheLogMeta, - cacheSource: "upstream", - }); - if (isRequestAborted) { - streamController.handleError(error); - return createErrorResult(499, "Request aborted"); - } - const persistentErrorCode = projectFailureUsageErrorCode({ - statusCode: failureStatus, - message: failureMessage, - errorCode: - upstreamErrorCode || (error instanceof Error && error.name ? error.name : "upstream_error"), - errorType: upstreamErrorType, - }); - persistFailureUsage(failureStatus, persistentErrorCode); - console.log(`${COLORS.red}[ERROR] ${failureMessage}${COLORS.reset}`); - if (stream && upstreamErrorCode) { - const result = createStreamingErrorResult( + if (isRequestAborted) { + streamController.handleError(error); + return createErrorResult(499, "Request aborted"); + } + const persistentErrorCode = projectFailureUsageErrorCode({ + statusCode: failureStatus, + message: failureMessage, + errorCode: + upstreamErrorCode || + (error instanceof Error && error.name ? error.name : "upstream_error"), + errorType: upstreamErrorType, + }); + persistFailureUsage(failureStatus, persistentErrorCode); + console.log(`${COLORS.red}[ERROR] ${failureMessage}${COLORS.reset}`); + if (stream && upstreamErrorCode) { + const result = createStreamingErrorResult( + failureStatus, + failureMessage, + upstreamErrorCode, + upstreamErrorType + ); + localLimiterErrors.markTrustedLocalRateLimitResponse(result.response, error); + return { + ...result, + errorType: upstreamErrorType, + errorCode: upstreamErrorCode, + }; + } + const result = createErrorResult( failureStatus, failureMessage, + null, upstreamErrorCode, upstreamErrorType ); localLimiterErrors.markTrustedLocalRateLimitResponse(result.response, error); - return { - ...result, - errorType: upstreamErrorType, - errorCode: upstreamErrorCode, - }; + return result; } - const result = createErrorResult( - failureStatus, - failureMessage, - null, - upstreamErrorCode, - upstreamErrorType - ); - localLimiterErrors.markTrustedLocalRateLimitResponse(result.response, error); - return result; - } - let upstreamErrorParsed = false; - let parsedStatusCode = providerResponse.status; - let parsedMessage = ""; - let parsedRetryAfterMs: number | null = null; - let upstreamErrorBody: unknown = null; + let upstreamErrorParsed = false; + let parsedStatusCode = providerResponse.status; + let parsedMessage = ""; + let parsedRetryAfterMs: number | null = null; + let upstreamErrorBody: unknown = null; - // Track whether stream_options was present and stripped — if so, 401/403 after - // that may be from the modification rather than a genuine auth failure, so we - // skip the credential refresh attempt in that case. - const hadStreamOptions = - targetFormat === FORMATS.OPENAI_RESPONSES && "stream_options" in translatedBody; - if (hadStreamOptions) { - delete translatedBody.stream_options; - } + // Track whether stream_options was present and stripped — if so, 401/403 after + // that may be from the modification rather than a genuine auth failure, so we + // skip the credential refresh attempt in that case. + const hadStreamOptions = + targetFormat === FORMATS.OPENAI_RESPONSES && "stream_options" in translatedBody; + if (hadStreamOptions) { + delete translatedBody.stream_options; + } - // Handle 401/403 - try token refresh using executor - // T-PROBE: probe-origin failures never attempt the refresh — a probe must - // not consume a rotating refresh token nor persist an "expired" - // deactivation on refresh failure (#9817). The 401/403 then flows into - // the normal providerFailure classification (record-only in probe mode). - if ( - (providerResponse.status === HTTP_STATUS.UNAUTHORIZED || - providerResponse.status === HTTP_STATUS.FORBIDDEN) && - !hadStreamOptions && // Skip refresh if failure may be from stream_options removal, not auth - !(await shouldIsolateProbeFailures()) - ) { - // Fix A: wrap refreshCredentials in runWithOnPersist so the persist callback - // executes INSIDE the per-connection mutex held by getAccessToken. This makes - // [network refresh + DB write + outer-state mutation] one atomic step and - // prevents concurrent requests from reading a stale refreshToken before the - // DB has been updated (refresh_token_reused on Codex/OpenAI). - // - // Not every executor routes refresh through getAccessToken (e.g. github.ts - // calls refreshCopilotToken directly). When the persistFn doesn't fire from - // inside getAccessToken, we still need to do the credentials mutation + user - // callback after refreshCredentials returns. The `persistFnRan` flag tracks - // which path executed so we don't double-fire (race-prone) or skip (regression). - // Front 3: remember the refresh_token we are about to present so that, if the - // refresh fails as unrecoverable, we can tell a genuine death apart from a - // stale-token reuse that a concurrent/sibling refresh already rotated past. - const attemptedRefreshToken = - typeof credentials?.refreshToken === "string" ? credentials.refreshToken : null; - let persistFnRan = false; - const persistFn = onCredentialsRefreshed - ? async (refreshResult: Record) => { - persistFnRan = true; - // Mutate the shared credentials object so subsequent executor calls - // in this request see the new tokens. Runs INSIDE the mutex. - Object.assign(credentials, refreshResult); - await onCredentialsRefreshed(refreshResult); + // Handle 401/403 - try token refresh using executor + // T-PROBE: probe-origin failures never attempt the refresh — a probe must + // not consume a rotating refresh token nor persist an "expired" + // deactivation on refresh failure (#9817). The 401/403 then flows into + // the normal providerFailure classification (record-only in probe mode). + if ( + (providerResponse.status === HTTP_STATUS.UNAUTHORIZED || + providerResponse.status === HTTP_STATUS.FORBIDDEN) && + !hadStreamOptions && // Skip refresh if failure may be from stream_options removal, not auth + !(await shouldIsolateProbeFailures()) + ) { + // Fix A: wrap refreshCredentials in runWithOnPersist so the persist callback + // executes INSIDE the per-connection mutex held by getAccessToken. This makes + // [network refresh + DB write + outer-state mutation] one atomic step and + // prevents concurrent requests from reading a stale refreshToken before the + // DB has been updated (refresh_token_reused on Codex/OpenAI). + // + // Not every executor routes refresh through getAccessToken (e.g. github.ts + // calls refreshCopilotToken directly). When the persistFn doesn't fire from + // inside getAccessToken, we still need to do the credentials mutation + user + // callback after refreshCredentials returns. The `persistFnRan` flag tracks + // which path executed so we don't double-fire (race-prone) or skip (regression). + // Front 3: remember the refresh_token we are about to present so that, if the + // refresh fails as unrecoverable, we can tell a genuine death apart from a + // stale-token reuse that a concurrent/sibling refresh already rotated past. + const attemptedRefreshToken = + typeof credentials?.refreshToken === "string" ? credentials.refreshToken : null; + let persistFnRan = false; + const persistFn = onCredentialsRefreshed + ? async (refreshResult: Record) => { + persistFnRan = true; + // Mutate the shared credentials object so subsequent executor calls + // in this request see the new tokens. Runs INSIDE the mutex. + Object.assign(credentials, refreshResult); + await onCredentialsRefreshed(refreshResult); + } + : undefined; + + // #4038: build a compare-and-swap reread so getAccessToken can skip the persist if a + // concurrent writer (sibling request / HealthCheck / replica) already rotated this + // connection's refresh_token past the one we presented — overwriting would revert it + // and revoke the token family. No connectionId ⇒ no guard (behavior unchanged). + const casConnectionId = + typeof credentials?.connectionId === "string" ? credentials.connectionId.trim() : ""; + const casReread = casConnectionId + ? async () => { + const latest = await getProviderConnectionById(casConnectionId); + return typeof latest?.refreshToken === "string" ? latest.refreshToken : null; + } + : null; + + const newCredentials = (await refreshWithRetry( + () => + runWithCasGuard( + casReread ? { expectedRefreshToken: attemptedRefreshToken, reread: casReread } : null, + () => runWithOnPersist(persistFn, () => executor.refreshCredentials(credentials, log)) + ), + 3, + log, + provider // Explicitly pass the provider to avoid universally tripping the "unknown" circuit breaker + )) as null | { + accessToken?: string; + copilotToken?: string; + }; + + if (newCredentials?.accessToken || newCredentials?.copilotToken) { + log?.info?.("TOKEN", `${provider?.toUpperCase()} | refreshed`); + + // Fall back to post-mutex mutation only for executors that don't route + // through getAccessToken (and therefore never fire onPersist). For + // executors that DO route through it (Codex, Claude, Gemini, etc.) the + // mutation already happened atomically inside the mutex. + if (!persistFnRan) { + Object.assign(credentials, newCredentials); + if (onCredentialsRefreshed) { + await onCredentialsRefreshed(newCredentials); + } } - : undefined; - // #4038: build a compare-and-swap reread so getAccessToken can skip the persist if a - // concurrent writer (sibling request / HealthCheck / replica) already rotated this - // connection's refresh_token past the one we presented — overwriting would revert it - // and revoke the token family. No connectionId ⇒ no guard (behavior unchanged). - const casConnectionId = - typeof credentials?.connectionId === "string" ? credentials.connectionId.trim() : ""; - const casReread = casConnectionId - ? async () => { - const latest = await getProviderConnectionById(casConnectionId); - return typeof latest?.refreshToken === "string" ? latest.refreshToken : null; + // Retry with new credentials — model + extra headers follow translatedBody.model so they + // stay aligned if this block ever runs after a path that mutates body.model (e.g. fallback). + try { + const retryModelId = String(translatedBody.model || effectiveModel); + assertManagedLeaseFence(getExecutionConnectionId(getExecutionCredentials())); + const retryResult = normalizeExecutorResult( + await runWithCapture(providerRequestCapture, () => + executor.execute({ + model: retryModelId, + body: translatedBody, + stream: upstreamStream, + credentials: getExecutionCredentials(), + signal: streamController.signal, + log, + extendedContext, + upstreamExtraHeaders: buildUpstreamHeadersForExecute(retryModelId), + clientHeaders: buildExecutorClientHeaders(clientRawRequest?.headers, userAgent), + clientResponseFormat, + onCredentialsRefreshed, + skipUpstreamRetry: isCombo, + contextEditing: { enabled: contextEditingEnabled }, + }) + ) + ); + + if (retryResult.response.ok) { + providerResponse = retryResult.response; + providerUrl = retryResult.url; + providerHeaders = new Headers(retryResult.headers || {}); + finalBody = providerRequestCapture.body(retryResult.transformedBody); + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); + updatePendingScope(pendingScope, { + providerRequest: finalBody, + providerUrl, + stage: "provider_response_started", + }); + upstreamErrorParsed = false; // Reset since new response is OK + } else { + providerResponse = retryResult.response; + upstreamErrorParsed = false; // Let it be parsed downstream + } + } catch (retryErr) { + if (isManagedLeaseFenceError(retryErr)) return managedLeaseFenceErrorResult(retryErr); + // Refresh succeeded but the retry leg failed (network blip, AbortError, + // executor throw). Don't swallow — the operator-visible signal "the user + // saw 401 even though auth was actually fixed" is much more confusing + // than the original 401 alone. Surface at error level with sanitization. + log?.error?.( + "TOKEN", + `${provider?.toUpperCase()} | retry after refresh failed: ${sanitizeErrorMessage(retryErr)}` + ); } - : null; - - const newCredentials = (await refreshWithRetry( - () => - runWithCasGuard( - casReread ? { expectedRefreshToken: attemptedRefreshToken, reread: casReread } : null, - () => runWithOnPersist(persistFn, () => executor.refreshCredentials(credentials, log)) - ), - 3, - log, - provider // Explicitly pass the provider to avoid universally tripping the "unknown" circuit breaker - )) as null | { - accessToken?: string; - copilotToken?: string; - }; - - if (newCredentials?.accessToken || newCredentials?.copilotToken) { - log?.info?.("TOKEN", `${provider?.toUpperCase()} | refreshed`); - - // Fall back to post-mutex mutation only for executors that don't route - // through getAccessToken (and therefore never fire onPersist). For - // executors that DO route through it (Codex, Claude, Gemini, etc.) the - // mutation already happened atomically inside the mutex. - if (!persistFnRan) { - Object.assign(credentials, newCredentials); - if (onCredentialsRefreshed) { - await onCredentialsRefreshed(newCredentials); + } else { + log?.warn?.("TOKEN", `${provider?.toUpperCase()} | refresh failed`); + if (isUnrecoverableRefreshError(newCredentials) && onCredentialsRefreshed) { + // Front 3 (reuse-race tolerance): before deactivating, re-read the DB. + // If a sibling/concurrent refresh already rotated this connection's + // refresh_token (common for Codex/OpenAI under one shared Auth0 client), + // the failure we saw was a stale-token reuse — the account is healthy + // with the newer token, so keep it active instead of killing it. + let alreadyRotated = false; + if (typeof connectionId === "string" && connectionId && attemptedRefreshToken) { + try { + const latest = await getProviderConnectionById(connectionId); + if (wasRefreshTokenRotated(attemptedRefreshToken, latest?.refreshToken)) { + alreadyRotated = true; + log?.warn?.( + "TOKEN", + `${provider.toUpperCase()} | refresh_token already rotated by a concurrent refresh — keeping connection active` + ); + } + } catch { + // DB read failed — fall through to the safe default (deactivate). + } + } + if (!alreadyRotated) { + await onCredentialsRefreshed({ testStatus: "expired", isActive: false }); + } } } + } - // Retry with new credentials — model + extra headers follow translatedBody.model so they - // stay aligned if this block ever runs after a path that mutates body.model (e.g. fallback). - try { - const retryModelId = String(translatedBody.model || effectiveModel); - assertManagedLeaseFence(getExecutionConnectionId(getExecutionCredentials())); - const retryResult = normalizeExecutorResult( - await runWithCapture(providerRequestCapture, () => - executor.execute({ - model: retryModelId, - body: translatedBody, - stream: upstreamStream, - credentials: getExecutionCredentials(), - signal: streamController.signal, - log, - extendedContext, - upstreamExtraHeaders: buildUpstreamHeadersForExecute(retryModelId), - clientHeaders: buildExecutorClientHeaders(clientRawRequest?.headers, userAgent), - clientResponseFormat, - onCredentialsRefreshed, - skipUpstreamRetry: isCombo, - contextEditing: { enabled: contextEditingEnabled }, - }) - ) + // Check provider response - return error info for fallback handling + providerFailure: if (!providerResponse.ok) { + trackPendingRequest(model, provider, connectionId, false); + + let statusCode = providerResponse.status; + let message = ""; + let retryAfterMs: number | null = null; + let upstreamErrorCode: string | undefined; + let upstreamErrorType: string | undefined; + + if (upstreamErrorParsed) { + statusCode = parsedStatusCode; + message = parsedMessage; + retryAfterMs = parsedRetryAfterMs; + } else { + const details = await parseUpstreamError(providerResponse, provider); + statusCode = details.statusCode; + message = details.message; + retryAfterMs = details.retryAfterMs; + upstreamErrorBody = details.responseBody; + upstreamErrorCode = details.errorCode as string | undefined; + upstreamErrorType = details.errorType as string | undefined; + } + + // Gateways like agentrouter misstate temporary quota exhaustion as 403/400, + // which downstream classification treats as AUTH_ERROR and clients like + // Claude Code treat as permanent. Restate to 429 (+ synthetic Retry-After) + // BEFORE any classification so both the fallback engine and the surfaced + // client status see a retryable error. Registry-scoped per provider. + const restatement = applyStatusRestatement({ + provider, + status: statusCode, + message, + body: upstreamErrorBody, + retryAfterMs, + }); + if (restatement.ruleId) { + statusCode = restatement.status; + retryAfterMs = restatement.retryAfterMs; + log?.info?.( + "STATUS_RESTATE", + `${provider} ${restatement.fromStatus}→${statusCode} (${restatement.ruleId})` ); + } - if (retryResult.response.ok) { - providerResponse = retryResult.response; - providerUrl = retryResult.url; - providerHeaders = new Headers(retryResult.headers || {}); - finalBody = providerRequestCapture.body(retryResult.transformedBody); + const signatureRecovery = pipelineRecovered + ? { attempted: false, succeeded: false, execution: null, error: null, recoveryBody: null } + : await recoverAnthropicThinkingSignature({ + provider, + statusCode, + message, + body: translatedBody, + execute: async (recoveryBody) => { + translatedBody = recoveryBody as typeof translatedBody; + return executeProviderRequest(currentModel, false); + }, + parseError: (response) => parseUpstreamError(response, provider), + }); + if (!pipelineRecovered && signatureRecovery.attempted && signatureRecovery.execution) { + providerResponse = signatureRecovery.execution.response; + if (signatureRecovery.succeeded) { + providerUrl = signatureRecovery.execution.url; + providerHeaders = signatureRecovery.execution.headers; + finalBody = providerRequestCapture.body(signatureRecovery.execution.transformedBody); reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); updatePendingScope(pendingScope, { providerRequest: finalBody, providerUrl, stage: "provider_response_started", }); - upstreamErrorParsed = false; // Reset since new response is OK - } else { - providerResponse = retryResult.response; - upstreamErrorParsed = false; // Let it be parsed downstream - } - } catch (retryErr) { - if (isManagedLeaseFenceError(retryErr)) return managedLeaseFenceErrorResult(retryErr); - // Refresh succeeded but the retry leg failed (network blip, AbortError, - // executor throw). Don't swallow — the operator-visible signal "the user - // saw 401 even though auth was actually fixed" is much more confusing - // than the original 401 alone. Surface at error level with sanitization. - log?.error?.( - "TOKEN", - `${provider?.toUpperCase()} | retry after refresh failed: ${sanitizeErrorMessage(retryErr)}` - ); - } - } else { - log?.warn?.("TOKEN", `${provider?.toUpperCase()} | refresh failed`); - if (isUnrecoverableRefreshError(newCredentials) && onCredentialsRefreshed) { - // Front 3 (reuse-race tolerance): before deactivating, re-read the DB. - // If a sibling/concurrent refresh already rotated this connection's - // refresh_token (common for Codex/OpenAI under one shared Auth0 client), - // the failure we saw was a stale-token reuse — the account is healthy - // with the newer token, so keep it active instead of killing it. - let alreadyRotated = false; - if (typeof connectionId === "string" && connectionId && attemptedRefreshToken) { - try { - const latest = await getProviderConnectionById(connectionId); - if (wasRefreshTokenRotated(attemptedRefreshToken, latest?.refreshToken)) { - alreadyRotated = true; - log?.warn?.( - "TOKEN", - `${provider.toUpperCase()} | refresh_token already rotated by a concurrent refresh — keeping connection active` - ); - } - } catch { - // DB read failed — fall through to the safe default (deactivate). - } - } - if (!alreadyRotated) { - await onCredentialsRefreshed({ testStatus: "expired", isActive: false }); + log?.info?.( + "THINKING_SIGNATURE", + `Recovered ${provider}/${currentModel} after one historical-thinking retry` + ); + } else if (signatureRecovery.error) { + statusCode = signatureRecovery.error.statusCode; + message = signatureRecovery.error.message; + retryAfterMs = signatureRecovery.error.retryAfterMs; + upstreamErrorBody = signatureRecovery.error.responseBody; + upstreamErrorCode = signatureRecovery.error.errorCode as string | undefined; + upstreamErrorType = signatureRecovery.error.errorType as string | undefined; } } - } - } - // Check provider response - return error info for fallback handling - providerFailure: if (!providerResponse.ok) { - trackPendingRequest(model, provider, connectionId, false); + if (signatureRecovery.succeeded) break providerFailure; - let statusCode = providerResponse.status; - let message = ""; - let retryAfterMs: number | null = null; - let upstreamErrorCode: string | undefined; - let upstreamErrorType: string | undefined; - - if (upstreamErrorParsed) { - statusCode = parsedStatusCode; - message = parsedMessage; - retryAfterMs = parsedRetryAfterMs; - } else { - const details = await parseUpstreamError(providerResponse, provider); - statusCode = details.statusCode; - message = details.message; - retryAfterMs = details.retryAfterMs; - upstreamErrorBody = details.responseBody; - upstreamErrorCode = details.errorCode as string | undefined; - upstreamErrorType = details.errorType as string | undefined; - } - - // Gateways like agentrouter misstate temporary quota exhaustion as 403/400, - // which downstream classification treats as AUTH_ERROR and clients like - // Claude Code treat as permanent. Restate to 429 (+ synthetic Retry-After) - // BEFORE any classification so both the fallback engine and the surfaced - // client status see a retryable error. Registry-scoped per provider. - const restatement = applyStatusRestatement({ - provider, - status: statusCode, - message, - body: upstreamErrorBody, - retryAfterMs, - }); - if (restatement.ruleId) { - statusCode = restatement.status; - retryAfterMs = restatement.retryAfterMs; - log?.info?.( - "STATUS_RESTATE", - `${provider} ${restatement.fromStatus}→${statusCode} (${restatement.ruleId})` - ); - } - - const signatureRecovery = pipelineRecovered - ? { attempted: false, succeeded: false, execution: null, error: null, recoveryBody: null } - : await recoverAnthropicThinkingSignature({ - provider, - statusCode, - message, - body: translatedBody, - execute: async (recoveryBody) => { - translatedBody = recoveryBody as typeof translatedBody; - return executeProviderRequest(currentModel, false); - }, - parseError: (response) => parseUpstreamError(response, provider), - }); - if (!pipelineRecovered && signatureRecovery.attempted && signatureRecovery.execution) { - providerResponse = signatureRecovery.execution.response; - if (signatureRecovery.succeeded) { - providerUrl = signatureRecovery.execution.url; - providerHeaders = signatureRecovery.execution.headers; - finalBody = providerRequestCapture.body(signatureRecovery.execution.transformedBody); - reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingScope(pendingScope, { - providerRequest: finalBody, - providerUrl, - stage: "provider_response_started", + // #10281 — tiny-budget reasoning probes (e.g. Claude Code's `/model` check + // sends `max_tokens: 1`): the model burns the whole budget on thinking, and + // some upstreams (e.g. api.cline.bot for deepseek-v4-flash) answer the empty + // outcome with a 5xx ("empty response content") instead of a truncated 200. + // Answer such probes with a valid truncated response rather than relaying the + // upstream failure — which would also mark the connection unavailable and + // poison fallback/cooldown bookkeeping for a request that is only a probe. + if ( + !stream && + isTinyBudgetReasoningProbe({ model: currentModel, body: finalBody || translatedBody }) && + isEmptyContentUpstreamFailure(statusCode, message) + ) { + providerResponse = buildReasoningProbeTruncatedResponse({ + model: currentModel, + maxTokens: toPositiveInteger( + (finalBody || translatedBody)?.max_tokens ?? + (finalBody || translatedBody)?.max_completion_tokens + ), + requestId: skillRequestId, }); - log?.info?.( - "THINKING_SIGNATURE", - `Recovered ${provider}/${currentModel} after one historical-thinking retry` + log?.warn?.( + "PROBE", + `Reasoning probe (max_tokens < ${REASONING_BUFFER_MIN_TRIGGER}) answered with truncated 200 — upstream reported "${message}"` ); - } else if (signatureRecovery.error) { - statusCode = signatureRecovery.error.statusCode; - message = signatureRecovery.error.message; - retryAfterMs = signatureRecovery.error.retryAfterMs; - upstreamErrorBody = signatureRecovery.error.responseBody; - upstreamErrorCode = signatureRecovery.error.errorCode as string | undefined; - upstreamErrorType = signatureRecovery.error.errorType as string | undefined; + break providerFailure; } - } - if (signatureRecovery.succeeded) break providerFailure; - - // #10281 — tiny-budget reasoning probes (e.g. Claude Code's `/model` check - // sends `max_tokens: 1`): the model burns the whole budget on thinking, and - // some upstreams (e.g. api.cline.bot for deepseek-v4-flash) answer the empty - // outcome with a 5xx ("empty response content") instead of a truncated 200. - // Answer such probes with a valid truncated response rather than relaying the - // upstream failure — which would also mark the connection unavailable and - // poison fallback/cooldown bookkeeping for a request that is only a probe. - if ( - !stream && - isTinyBudgetReasoningProbe({ model: currentModel, body: finalBody || translatedBody }) && - isEmptyContentUpstreamFailure(statusCode, message) - ) { - providerResponse = buildReasoningProbeTruncatedResponse({ - model: currentModel, - maxTokens: toPositiveInteger( - (finalBody || translatedBody)?.max_tokens ?? - (finalBody || translatedBody)?.max_completion_tokens - ), - requestId: skillRequestId, - }); - log?.warn?.( - "PROBE", - `Reasoning probe (max_tokens < ${REASONING_BUFFER_MIN_TRIGGER}) answered with truncated 200 — upstream reported "${message}"` - ); - break providerFailure; - } - - // T06/T10/T36: classify provider errors and persist terminal account states. - let errorType = classifyProviderError(statusCode, message, provider); - if (statusCode === 429 && isModelScope()) { - const decision = classifyModelScope429(message, normalizeHeaders(providerResponse.headers)); - errorType = - decision.kind === "quota_exhausted" - ? PROVIDER_ERROR_TYPES.QUOTA_EXHAUSTED - : PROVIDER_ERROR_TYPES.RATE_LIMITED; - log?.warn?.( - "MODELSCOPE_429", - `${decision.kind} (model remaining: ${decision.snapshot.modelRemaining ?? "unknown"}, total remaining: ${decision.snapshot.totalRemaining ?? "unknown"})` - ); - } - // Classifiers and recovery paths above consume the raw provider wording. - // Project a separate value only at persistent connection-state boundaries. - const persistentMessage = sanitizeErrorMessage(message) || "Provider request failed"; - const errorConnectionId = getCurrentConnectionId(); - if (errorConnectionId && errorType) { - try { - if (errorType === PROVIDER_ERROR_TYPES.FORBIDDEN) { - { - const probeIsolated = await shouldIsolateProbeFailures(); - await writeTerminalStatus( - errorConnectionId, - { - testStatus: "banned", - isActive: false, - lastError: persistentMessage, - lastErrorType: errorType, - errorCode: String(statusCode), - }, - probeIsolated ? "probe" : "production" - ); - if (probeIsolated) { - console.warn( - `[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active` - ); - } else { - console.warn( - `[provider] Node ${errorConnectionId} banned (${statusCode}) — disabling permanently` - ); - } - } - } else if (errorType === PROVIDER_ERROR_TYPES.ACCOUNT_DEACTIVATED) { - // T-PROBE: probe-origin failures (test-all) never deactivate — - // record but stay active; Plan A (extra keys) stays first so the - // real path keeps its existing priority (#9817). - // Plan A: if connection has extra API keys, don't disable — only the failing key is affected. - // Single-key connections still get disabled as before. - if ( - connectionHasExtraKeys( - errorConnectionId, - (credentials?.providerSpecificData as Record | undefined) - ?.extraApiKeys as string[] | undefined - ) - ) { - await updateProviderConnection(errorConnectionId, { - lastErrorType: errorType, - lastError: persistentMessage, - errorCode: statusCode, - }); - console.warn( - `[provider] Node ${errorConnectionId} account deactivated (${statusCode}) — has extra keys, keeping connection active` - ); - } else { - const probeIsolated2 = await shouldIsolateProbeFailures(); - await writeTerminalStatus( - errorConnectionId, - { - testStatus: "deactivated", - isActive: false, - lastError: persistentMessage, - lastErrorType: errorType, - errorCode: String(statusCode), - }, - probeIsolated2 ? "probe" : "production" - ); - if (probeIsolated2) { - console.warn( - `[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active` - ); - } else { - console.warn( - `[provider] Node ${errorConnectionId} account deactivated (${statusCode}) — disabling permanently` - ); - } - } - } else if (errorType === PROVIDER_ERROR_TYPES.QUOTA_EXHAUSTED) { - { - const probeIsolated3 = await shouldIsolateProbeFailures(); - if (probeIsolated3) { + // T06/T10/T36: classify provider errors and persist terminal account states. + let errorType = classifyProviderError(statusCode, message, provider); + if (statusCode === 429 && isModelScope()) { + const decision = classifyModelScope429(message, normalizeHeaders(providerResponse.headers)); + errorType = + decision.kind === "quota_exhausted" + ? PROVIDER_ERROR_TYPES.QUOTA_EXHAUSTED + : PROVIDER_ERROR_TYPES.RATE_LIMITED; + log?.warn?.( + "MODELSCOPE_429", + `${decision.kind} (model remaining: ${decision.snapshot.modelRemaining ?? "unknown"}, total remaining: ${decision.snapshot.totalRemaining ?? "unknown"})` + ); + } + // Classifiers and recovery paths above consume the raw provider wording. + // Project a separate value only at persistent connection-state boundaries. + const persistentMessage = sanitizeErrorMessage(message) || "Provider request failed"; + const errorConnectionId = getCurrentConnectionId(); + if (errorConnectionId && errorType) { + try { + if (errorType === PROVIDER_ERROR_TYPES.FORBIDDEN) { + { + const probeIsolated = await shouldIsolateProbeFailures(); await writeTerminalStatus( errorConnectionId, { - testStatus: "credits_exhausted", + testStatus: "banned", + isActive: false, lastError: persistentMessage, lastErrorType: errorType, errorCode: String(statusCode), }, - "probe" + probeIsolated ? "probe" : "production" ); - console.warn( - `[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active` - ); - } else { - // Kimi's 403 says "billing cycle" for both an exhausted subscription and a - // temporary request window. Read its official usage endpoint before making - // the connection terminal: a non-zero Weekly quota plus an empty Ratelimit - // window must recover automatically at the reported reset time. - let kimiRateLimitResetAt: string | null = null; - if (provider === "kimi-coding") { - try { - const { fetchAndPersistProviderLimits } = - await import("@/lib/usage/providerLimits"); - const { usage } = await fetchAndPersistProviderLimits( - errorConnectionId, - "manual" - ); - kimiRateLimitResetAt = getKimiTemporaryRateLimitResetAt(usage); - } catch { - // Preserve the existing quota handling when Kimi's usage endpoint is unavailable. - } - } - - // Providers with per-model quotas — lock the model only, not the connection - let quotaCooldownMs = kimiRateLimitResetAt - ? Math.max(new Date(kimiRateLimitResetAt).getTime() - Date.now(), 0) - : retryAfterMs || COOLDOWN_MS.rateLimit; - const deferAntigravityQuotaStateToCaller = shouldDeferAntigravityQuotaStateToCaller( - provider, - typeof onStreamFailure === "function" - ); - const isAntigravityQuotaFamily = shouldDeferAntigravityQuotaStateToCaller( - provider, - true - ); - let coreOwnedAntigravityLockout: { - cooldownMs: number; - failureCount: number; - } | null = null; - if (isAntigravityQuotaFamily && !deferAntigravityQuotaStateToCaller) { - const quotaErrorText = - typeof upstreamErrorBody === "string" - ? upstreamErrorBody - : upstreamErrorBody == null - ? message - : JSON.stringify(upstreamErrorBody); - coreOwnedAntigravityLockout = await recordCoreOwnedAntigravityQuotaState({ - provider, - connectionId: errorConnectionId, - model, - status: statusCode, - errorText: quotaErrorText, - headers: providerResponse.headers, - }); - quotaCooldownMs = coreOwnedAntigravityLockout.cooldownMs; - } - const accountSemaphoreKey = resolveAccountSemaphoreKey({ - provider, - model: currentModel, - connectionId: errorConnectionId, - credentials, - }); - if (accountSemaphoreKey && !deferAntigravityQuotaStateToCaller) { - markAccountSemaphoreBlocked(accountSemaphoreKey, quotaCooldownMs); - } - if (deferAntigravityQuotaStateToCaller) { - // Defer both model and account-semaphore cooldowns to - // markAccountUnavailable, where header/body provenance and the - // configured maxCooldownMs are available. Direct consumers such - // as Responses pass no owner callback and retain core ownership. - } else if (coreOwnedAntigravityLockout) { + if (probeIsolated) { console.warn( - `[provider] Node ${errorConnectionId} Antigravity model quota exhausted (${statusCode}) for ${model} - ${Math.ceil(coreOwnedAntigravityLockout.cooldownMs / 1000)}s (failureCount=${coreOwnedAntigravityLockout.failureCount}, owner=core)` - ); - } else if (kimiRateLimitResetAt) { - await updateProviderConnection(errorConnectionId, { - testStatus: "unavailable", - rateLimitedUntil: kimiRateLimitResetAt, - backoffLevel: 0, - lastErrorType: PROVIDER_ERROR_TYPES.RATE_LIMITED, - lastError: persistentMessage, - errorCode: statusCode, - }); - console.warn( - `[provider] Node ${errorConnectionId} Kimi request window exhausted (${statusCode}) — retrying after ${kimiRateLimitResetAt}` - ); - } else if (isModelScope() && errorConnectionId) { - lockModel(provider, errorConnectionId, model, "quota_exhausted", quotaCooldownMs); - console.warn( - `[provider] Node ${errorConnectionId} ModelScope model quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (connection stays active)` - ); - } else if ( - lockModelIfPerModelQuota( - provider, - errorConnectionId, - model, - "quota_exhausted", - quotaCooldownMs - ) - ) { - const quotaScope = getQuotaScopeLabelForProvider(provider, model); - console.warn( - `[provider] Node ${errorConnectionId} ${quotaScope}-only quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (cooldown_scope=${quotaScope}, ttl_source=${retryAfterMs ? "upstream" : "inferred"}, connection stays active)` + `[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active` ); } else { + console.warn( + `[provider] Node ${errorConnectionId} banned (${statusCode}) — disabling permanently` + ); + } + } + } else if (errorType === PROVIDER_ERROR_TYPES.ACCOUNT_DEACTIVATED) { + // T-PROBE: probe-origin failures (test-all) never deactivate — + // record but stay active; Plan A (extra keys) stays first so the + // real path keeps its existing priority (#9817). + // Plan A: if connection has extra API keys, don't disable — only the failing key is affected. + // Single-key connections still get disabled as before. + if ( + connectionHasExtraKeys( + errorConnectionId, + (credentials?.providerSpecificData as Record | undefined) + ?.extraApiKeys as string[] | undefined + ) + ) { + await updateProviderConnection(errorConnectionId, { + lastErrorType: errorType, + lastError: persistentMessage, + errorCode: statusCode, + }); + console.warn( + `[provider] Node ${errorConnectionId} account deactivated (${statusCode}) — has extra keys, keeping connection active` + ); + } else { + const probeIsolated2 = await shouldIsolateProbeFailures(); + await writeTerminalStatus( + errorConnectionId, + { + testStatus: "deactivated", + isActive: false, + lastError: persistentMessage, + lastErrorType: errorType, + errorCode: String(statusCode), + }, + probeIsolated2 ? "probe" : "production" + ); + if (probeIsolated2) { + console.warn( + `[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active` + ); + } else { + console.warn( + `[provider] Node ${errorConnectionId} account deactivated (${statusCode}) — disabling permanently` + ); + } + } + } else if (errorType === PROVIDER_ERROR_TYPES.QUOTA_EXHAUSTED) { + { + const probeIsolated3 = await shouldIsolateProbeFailures(); + if (probeIsolated3) { await writeTerminalStatus( errorConnectionId, { @@ -4406,176 +4291,314 @@ export async function handleChatCore({ lastErrorType: errorType, errorCode: String(statusCode), }, - "production" + "probe" ); console.warn( - `[provider] Node ${errorConnectionId} exhausted quota (${statusCode})` + `[provider] Node ${errorConnectionId} probe ${errorType} (${statusCode}) — connection stays active` ); + } else { + // Kimi's 403 says "billing cycle" for both an exhausted subscription and a + // temporary request window. Read its official usage endpoint before making + // the connection terminal: a non-zero Weekly quota plus an empty Ratelimit + // window must recover automatically at the reported reset time. + let kimiRateLimitResetAt: string | null = null; + if (provider === "kimi-coding") { + try { + const { fetchAndPersistProviderLimits } = + await import("@/lib/usage/providerLimits"); + const { usage } = await fetchAndPersistProviderLimits( + errorConnectionId, + "manual" + ); + kimiRateLimitResetAt = getKimiTemporaryRateLimitResetAt(usage); + } catch { + // Preserve the existing quota handling when Kimi's usage endpoint is unavailable. + } + } + + // Providers with per-model quotas — lock the model only, not the connection + let quotaCooldownMs = kimiRateLimitResetAt + ? Math.max(new Date(kimiRateLimitResetAt).getTime() - Date.now(), 0) + : retryAfterMs || COOLDOWN_MS.rateLimit; + const deferAntigravityQuotaStateToCaller = shouldDeferAntigravityQuotaStateToCaller( + provider, + typeof onStreamFailure === "function" + ); + const isAntigravityQuotaFamily = shouldDeferAntigravityQuotaStateToCaller( + provider, + true + ); + let coreOwnedAntigravityLockout: { + cooldownMs: number; + failureCount: number; + } | null = null; + if (isAntigravityQuotaFamily && !deferAntigravityQuotaStateToCaller) { + const quotaErrorText = + typeof upstreamErrorBody === "string" + ? upstreamErrorBody + : upstreamErrorBody == null + ? message + : JSON.stringify(upstreamErrorBody); + coreOwnedAntigravityLockout = await recordCoreOwnedAntigravityQuotaState({ + provider, + connectionId: errorConnectionId, + model, + status: statusCode, + errorText: quotaErrorText, + headers: providerResponse.headers, + }); + quotaCooldownMs = coreOwnedAntigravityLockout.cooldownMs; + } + const accountSemaphoreKey = resolveAccountSemaphoreKey({ + provider, + model: currentModel, + connectionId: errorConnectionId, + credentials, + }); + if (accountSemaphoreKey && !deferAntigravityQuotaStateToCaller) { + markAccountSemaphoreBlocked(accountSemaphoreKey, quotaCooldownMs); + } + if (deferAntigravityQuotaStateToCaller) { + // Defer both model and account-semaphore cooldowns to + // markAccountUnavailable, where header/body provenance and the + // configured maxCooldownMs are available. Direct consumers such + // as Responses pass no owner callback and retain core ownership. + } else if (coreOwnedAntigravityLockout) { + console.warn( + `[provider] Node ${errorConnectionId} Antigravity model quota exhausted (${statusCode}) for ${model} - ${Math.ceil(coreOwnedAntigravityLockout.cooldownMs / 1000)}s (failureCount=${coreOwnedAntigravityLockout.failureCount}, owner=core)` + ); + } else if (kimiRateLimitResetAt) { + await updateProviderConnection(errorConnectionId, { + testStatus: "unavailable", + rateLimitedUntil: kimiRateLimitResetAt, + backoffLevel: 0, + lastErrorType: PROVIDER_ERROR_TYPES.RATE_LIMITED, + lastError: persistentMessage, + errorCode: statusCode, + }); + console.warn( + `[provider] Node ${errorConnectionId} Kimi request window exhausted (${statusCode}) — retrying after ${kimiRateLimitResetAt}` + ); + } else if (isModelScope() && errorConnectionId) { + lockModel(provider, errorConnectionId, model, "quota_exhausted", quotaCooldownMs); + console.warn( + `[provider] Node ${errorConnectionId} ModelScope model quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (connection stays active)` + ); + } else if ( + lockModelIfPerModelQuota( + provider, + errorConnectionId, + model, + "quota_exhausted", + quotaCooldownMs + ) + ) { + const quotaScope = getQuotaScopeLabelForProvider(provider, model); + console.warn( + `[provider] Node ${errorConnectionId} ${quotaScope}-only quota exhausted (${statusCode}) for ${model} - ${Math.ceil(quotaCooldownMs / 1000)}s (cooldown_scope=${quotaScope}, ttl_source=${retryAfterMs ? "upstream" : "inferred"}, connection stays active)` + ); + } else { + await writeTerminalStatus( + errorConnectionId, + { + testStatus: "credits_exhausted", + lastError: persistentMessage, + lastErrorType: errorType, + errorCode: String(statusCode), + }, + "production" + ); + console.warn( + `[provider] Node ${errorConnectionId} exhausted quota (${statusCode})` + ); + } + } // close probeIsolated3 else + } + } else if (errorType === PROVIDER_ERROR_TYPES.UNAUTHORIZED) { + // Normal 401 (token/session auth issue): keep account active for refresh/re-auth. + await updateProviderConnection(errorConnectionId, { + lastErrorType: errorType, + lastError: persistentMessage, + errorCode: statusCode, + }); + } else if (errorType === PROVIDER_ERROR_TYPES.OAUTH_INVALID_TOKEN) { + // OAuth 401 with invalid credentials - token refresh can recover + await updateProviderConnection(errorConnectionId, { + lastErrorType: errorType, + lastError: persistentMessage, + errorCode: statusCode, + }); + console.warn( + `[provider] Node ${errorConnectionId} OAuth token invalid (${statusCode}) — token refresh available` + ); + } else if (errorType === PROVIDER_ERROR_TYPES.PROJECT_ROUTE_ERROR) { + // Cloud Code 403 with stale project: not a ban, keep account active. + await updateProviderConnection(errorConnectionId, { + lastErrorType: errorType, + lastError: persistentMessage, + errorCode: statusCode, + }); + console.warn( + `[provider] Node ${errorConnectionId} project routing error (${statusCode}) — not banning` + ); + } else if (errorType === PROVIDER_ERROR_TYPES.GEO_BLOCKED) { + // Google regional-availability refusal (e.g. "User location is not + // supported for the API use."). Account-independent and non-terminal: + // exclude the connection for the cooldown window so routing moves to + // other accounts instead of re-selecting this one on every request, + // and never mark it banned/expired. It becomes usable again once + // egress is routed through a supported-region proxy. + const geoCooldownMs = COOLDOWN_MS.geoBlocked ?? 24 * 60 * 60 * 1000; + await updateProviderConnection(errorConnectionId, { + lastErrorType: errorType, + lastError: persistentMessage, + errorCode: statusCode, + }); + // T-PROBE: the 24h exclusion is a routing mutation — a probe must + // not push a connection into a day-long cooldown (#9817). + if (!(await shouldIsolateProbeFailures())) { + try { + const { setConnectionRateLimitUntil } = await import("@/lib/db/providers"); + setConnectionRateLimitUntil(errorConnectionId, Date.now() + geoCooldownMs); + } catch { + // DB write failure must never break the fallback loop } - } // close probeIsolated3 else - } - } else if (errorType === PROVIDER_ERROR_TYPES.UNAUTHORIZED) { - // Normal 401 (token/session auth issue): keep account active for refresh/re-auth. - await updateProviderConnection(errorConnectionId, { - lastErrorType: errorType, - lastError: persistentMessage, - errorCode: statusCode, - }); - } else if (errorType === PROVIDER_ERROR_TYPES.OAUTH_INVALID_TOKEN) { - // OAuth 401 with invalid credentials - token refresh can recover - await updateProviderConnection(errorConnectionId, { - lastErrorType: errorType, - lastError: persistentMessage, - errorCode: statusCode, - }); - console.warn( - `[provider] Node ${errorConnectionId} OAuth token invalid (${statusCode}) — token refresh available` - ); - } else if (errorType === PROVIDER_ERROR_TYPES.PROJECT_ROUTE_ERROR) { - // Cloud Code 403 with stale project: not a ban, keep account active. - await updateProviderConnection(errorConnectionId, { - lastErrorType: errorType, - lastError: persistentMessage, - errorCode: statusCode, - }); - console.warn( - `[provider] Node ${errorConnectionId} project routing error (${statusCode}) — not banning` - ); - } else if (errorType === PROVIDER_ERROR_TYPES.GEO_BLOCKED) { - // Google regional-availability refusal (e.g. "User location is not - // supported for the API use."). Account-independent and non-terminal: - // exclude the connection for the cooldown window so routing moves to - // other accounts instead of re-selecting this one on every request, - // and never mark it banned/expired. It becomes usable again once - // egress is routed through a supported-region proxy. - const geoCooldownMs = COOLDOWN_MS.geoBlocked ?? 24 * 60 * 60 * 1000; - await updateProviderConnection(errorConnectionId, { - lastErrorType: errorType, - lastError: persistentMessage, - errorCode: statusCode, - }); - // T-PROBE: the 24h exclusion is a routing mutation — a probe must - // not push a connection into a day-long cooldown (#9817). - if (!(await shouldIsolateProbeFailures())) { + } + console.warn( + `[provider] Node ${errorConnectionId} geo-blocked (${statusCode}) — excluded for ${Math.ceil(geoCooldownMs / 1000)}s, trying other accounts` + ); + } else if (errorType === PROVIDER_ERROR_TYPES.GCP_PROJECT_REQUIRED) { + // Antigravity BYOP: the account must Bring Its Own GCP Project. + // Account-specific and fixable by entering a Project ID — never a + // model lockout, never a ban. Exclude the connection for the + // cooldown window so selection prefers sibling accounts; the 422 + // body carries the actionable message when no sibling is available. + const byopCooldownMs = COOLDOWN_MS.gcpProjectRequired ?? 24 * 60 * 60 * 1000; + await updateProviderConnection(errorConnectionId, { + lastErrorType: errorType, + lastError: persistentMessage, + errorCode: statusCode, + }); try { const { setConnectionRateLimitUntil } = await import("@/lib/db/providers"); - setConnectionRateLimitUntil(errorConnectionId, Date.now() + geoCooldownMs); + setConnectionRateLimitUntil(errorConnectionId, Date.now() + byopCooldownMs); } catch { - // DB write failure must never break the fallback loop + // best-effort — never break the error path + } + console.warn( + `[provider] Node ${errorConnectionId} GCP project required (${statusCode}) — excluded for ${Math.ceil(byopCooldownMs / 1000)}s, routing to other accounts (enter a Project ID to restore)` + ); + } else if (errorType === PROVIDER_ERROR_TYPES.MODEL_NOT_FOUND) { + // 404 — model/endpoint does not exist upstream. Lock the model so the + // retry/backoff loop stops hammering the dead endpoint (which would + // otherwise degenerate into a 429 rate-limit storm). Connection stays + // active since only the specific model is unavailable. (#6827) + const notFoundCooldownMs = COOLDOWN_MS.notFound; + // T-PROBE: the model lockout is a routing mutation — a probe must + // not lock a model for the cooldown window (#9817). + if (!(await shouldIsolateProbeFailures())) { + lockModel( + provider, + errorConnectionId, + currentModel, + "model_not_found", + notFoundCooldownMs + ); + console.warn( + `[provider] Node ${errorConnectionId} model not found (${statusCode}) for ${currentModel} - locking model for ${Math.ceil(notFoundCooldownMs / 1000)}s (connection stays active)` + ); } } - console.warn( - `[provider] Node ${errorConnectionId} geo-blocked (${statusCode}) — excluded for ${Math.ceil(geoCooldownMs / 1000)}s, trying other accounts` - ); - } else if (errorType === PROVIDER_ERROR_TYPES.GCP_PROJECT_REQUIRED) { - // Antigravity BYOP: the account must Bring Its Own GCP Project. - // Account-specific and fixable by entering a Project ID — never a - // model lockout, never a ban. Exclude the connection for the - // cooldown window so selection prefers sibling accounts; the 422 - // body carries the actionable message when no sibling is available. - const byopCooldownMs = COOLDOWN_MS.gcpProjectRequired ?? 24 * 60 * 60 * 1000; - await updateProviderConnection(errorConnectionId, { - lastErrorType: errorType, - lastError: persistentMessage, - errorCode: statusCode, - }); - try { - const { setConnectionRateLimitUntil } = await import("@/lib/db/providers"); - setConnectionRateLimitUntil(errorConnectionId, Date.now() + byopCooldownMs); - } catch { - // best-effort — never break the error path - } - console.warn( - `[provider] Node ${errorConnectionId} GCP project required (${statusCode}) — excluded for ${Math.ceil(byopCooldownMs / 1000)}s, routing to other accounts (enter a Project ID to restore)` - ); - } else if (errorType === PROVIDER_ERROR_TYPES.MODEL_NOT_FOUND) { - // 404 — model/endpoint does not exist upstream. Lock the model so the - // retry/backoff loop stops hammering the dead endpoint (which would - // otherwise degenerate into a 429 rate-limit storm). Connection stays - // active since only the specific model is unavailable. (#6827) - const notFoundCooldownMs = COOLDOWN_MS.notFound; - // T-PROBE: the model lockout is a routing mutation — a probe must - // not lock a model for the cooldown window (#9817). - if (!(await shouldIsolateProbeFailures())) { - lockModel( - provider, - errorConnectionId, - currentModel, - "model_not_found", - notFoundCooldownMs - ); - console.warn( - `[provider] Node ${errorConnectionId} model not found (${statusCode}) for ${currentModel} - locking model for ${Math.ceil(notFoundCooldownMs / 1000)}s (connection stays active)` - ); - } + } catch { + // Best-effort state update; request flow should continue with fallback handling. } - } catch { - // Best-effort state update; request flow should continue with fallback handling. } - } - appendRequestLog({ - model, - provider, - connectionId: errorConnectionId, - status: `FAILED ${statusCode}`, - }).catch(() => {}); + appendRequestLog({ + model, + provider, + connectionId: errorConnectionId, + status: `FAILED ${statusCode}`, + }).catch(() => {}); - const errMsg = formatProviderError(new Error(message), provider, model, statusCode); - console.log(`${COLORS.red}[ERROR] ${errMsg}${COLORS.reset}`); + const errMsg = formatProviderError(new Error(message), provider, model, statusCode); + console.log(`${COLORS.red}[ERROR] ${errMsg}${COLORS.reset}`); - // Log Antigravity retry time if available - if (retryAfterMs && provider === "antigravity") { - const retrySeconds = Math.ceil(retryAfterMs / 1000); - log?.debug?.("RETRY", `Antigravity quota reset in ${retrySeconds}s (${retryAfterMs}ms)`); - } + // Log Antigravity retry time if available + if (retryAfterMs && provider === "antigravity") { + const retrySeconds = Math.ceil(retryAfterMs / 1000); + log?.debug?.("RETRY", `Antigravity quota reset in ${retrySeconds}s (${retryAfterMs}ms)`); + } - // Log error with full request body for debugging - reqLogger.logError(new Error(message), finalBody || translatedBody); - reqLogger.logProviderResponse( - providerResponse.status, - providerResponse.statusText, - providerResponse.headers, - upstreamErrorBody - ); + // Log error with full request body for debugging + reqLogger.logError(new Error(message), finalBody || translatedBody); + reqLogger.logProviderResponse( + providerResponse.status, + providerResponse.statusText, + providerResponse.headers, + upstreamErrorBody + ); - // Update rate limiter from error response headers - updateFromHeaders(provider, errorConnectionId, providerResponse.headers, statusCode, model); - if (errorConnectionId && upstreamErrorBody !== null && upstreamErrorBody !== undefined) { - updateFromResponseBody(provider, errorConnectionId, upstreamErrorBody, statusCode, model); - } + // Update rate limiter from error response headers + updateFromHeaders(provider, errorConnectionId, providerResponse.headers, statusCode, model); + if (errorConnectionId && upstreamErrorBody !== null && upstreamErrorBody !== undefined) { + updateFromResponseBody(provider, errorConnectionId, upstreamErrorBody, statusCode, model); + } - // ── T5: Intra-family model fallback ────────────────────────────────────── - // Before returning a model-unavailable error upstream, try sibling models - // from the same family. This keeps the request alive on the same account - // instead of failing the entire combo. - if (!pipelineRecovered && isModelUnavailableError(statusCode, message, provider)) { - const nextModel = getNextFamilyFallback(currentModel, triedModels, provider); - if (nextModel) { - triedModels.add(nextModel); - currentModel = nextModel; - translatedBody.model = nextModel; - log?.info?.("MODEL_FALLBACK", `${model} unavailable (${statusCode}) → trying ${nextModel}`); - // Re-execute with the fallback model - try { - const fallbackResult = await executeProviderRequest(nextModel, false); - if (fallbackResult.response.ok) { - providerResponse = fallbackResult.response; - providerUrl = fallbackResult.url; - providerHeaders = fallbackResult.headers; - finalBody = providerRequestCapture.body(fallbackResult.transformedBody); - reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingScope(pendingScope, { - providerRequest: finalBody, - providerUrl, - stage: "provider_response_started", - }); - // Continue processing with the fallback response — skip error return - log?.info?.("MODEL_FALLBACK", `Serving ${nextModel} as fallback for ${model}`); - // Jump to streaming/non-streaming handling below - // We fall through by NOT returning here - } else { - // Fallback also failed — return original error + // ── T5: Intra-family model fallback ────────────────────────────────────── + // Before returning a model-unavailable error upstream, try sibling models + // from the same family. This keeps the request alive on the same account + // instead of failing the entire combo. + if (!pipelineRecovered && isModelUnavailableError(statusCode, message, provider)) { + const nextModel = getNextFamilyFallback(currentModel, triedModels, provider); + if (nextModel) { + triedModels.add(nextModel); + currentModel = nextModel; + translatedBody.model = nextModel; + log?.info?.( + "MODEL_FALLBACK", + `${model} unavailable (${statusCode}) → trying ${nextModel}` + ); + // Re-execute with the fallback model + try { + const fallbackResult = await executeProviderRequest(nextModel, false); + if (fallbackResult.response.ok) { + providerResponse = fallbackResult.response; + providerUrl = fallbackResult.url; + providerHeaders = fallbackResult.headers; + finalBody = providerRequestCapture.body(fallbackResult.transformedBody); + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); + updatePendingScope(pendingScope, { + providerRequest: finalBody, + providerUrl, + stage: "provider_response_started", + }); + // Continue processing with the fallback response — skip error return + log?.info?.("MODEL_FALLBACK", `Serving ${nextModel} as fallback for ${model}`); + // Jump to streaming/non-streaming handling below + // We fall through by NOT returning here + } else { + // Fallback also failed — return original error + persistAttemptLogs({ + status: statusCode, + error: errMsg, + providerRequest: finalBody || translatedBody, + providerResponse: upstreamErrorBody, + clientResponse: buildErrorBody(statusCode, errMsg), + cacheSource: "upstream", + }); + persistFailureUsage(statusCode, "model_unavailable"); + return createErrorResult( + statusCode, + errMsg, + retryAfterMs, + upstreamErrorCode, + upstreamErrorType, + upstreamErrorBody, + { passthrough: sourceFormat === FORMATS.CLAUDE } + ); + } + } catch { persistAttemptLogs({ status: statusCode, error: errMsg, @@ -4595,7 +4618,7 @@ export async function handleChatCore({ { passthrough: sourceFormat === FORMATS.CLAUDE } ); } - } catch { + } else { persistAttemptLogs({ status: statusCode, error: errMsg, @@ -4615,56 +4638,59 @@ export async function handleChatCore({ { passthrough: sourceFormat === FORMATS.CLAUDE } ); } - } else { - persistAttemptLogs({ - status: statusCode, - error: errMsg, - providerRequest: finalBody || translatedBody, - providerResponse: upstreamErrorBody, - clientResponse: buildErrorBody(statusCode, errMsg), - cacheSource: "upstream", - }); - persistFailureUsage(statusCode, "model_unavailable"); - return createErrorResult( - statusCode, - errMsg, - retryAfterMs, - upstreamErrorCode, - upstreamErrorType, - upstreamErrorBody, - { passthrough: sourceFormat === FORMATS.CLAUDE } + } else if (isContextOverflowError(statusCode, message)) { + const familyCandidates = getModelFamily(currentModel, provider).filter( + (m) => m !== currentModel && !triedModels.has(m) ); - } - } else if (isContextOverflowError(statusCode, message)) { - const familyCandidates = getModelFamily(currentModel, provider).filter( - (m) => m !== currentModel && !triedModels.has(m) - ); - const nextModel = - findLargerContextModel(currentModel, familyCandidates, provider) ?? - getNextFamilyFallback(currentModel, triedModels, provider); - if (nextModel) { - triedModels.add(nextModel); - currentModel = nextModel; - translatedBody.model = nextModel; - log?.info?.("CONTEXT_OVERFLOW_FALLBACK", `${model} context overflow → trying ${nextModel}`); - try { - const fallbackResult = await executeProviderRequest(nextModel, false); - if (fallbackResult.response.ok) { - providerResponse = fallbackResult.response; - providerUrl = fallbackResult.url; - providerHeaders = fallbackResult.headers; - finalBody = providerRequestCapture.body(fallbackResult.transformedBody); - reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingScope(pendingScope, { - providerRequest: finalBody, - providerUrl, - stage: "provider_response_started", - }); - log?.info?.( - "CONTEXT_OVERFLOW_FALLBACK", - `Serving ${nextModel} as fallback for ${model}` - ); - } else { + const nextModel = + findLargerContextModel(currentModel, familyCandidates, provider) ?? + getNextFamilyFallback(currentModel, triedModels, provider); + if (nextModel) { + triedModels.add(nextModel); + currentModel = nextModel; + translatedBody.model = nextModel; + log?.info?.( + "CONTEXT_OVERFLOW_FALLBACK", + `${model} context overflow → trying ${nextModel}` + ); + try { + const fallbackResult = await executeProviderRequest(nextModel, false); + if (fallbackResult.response.ok) { + providerResponse = fallbackResult.response; + providerUrl = fallbackResult.url; + providerHeaders = fallbackResult.headers; + finalBody = providerRequestCapture.body(fallbackResult.transformedBody); + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); + updatePendingScope(pendingScope, { + providerRequest: finalBody, + providerUrl, + stage: "provider_response_started", + }); + log?.info?.( + "CONTEXT_OVERFLOW_FALLBACK", + `Serving ${nextModel} as fallback for ${model}` + ); + } else { + persistAttemptLogs({ + status: statusCode, + error: errMsg, + providerRequest: finalBody || translatedBody, + providerResponse: upstreamErrorBody, + clientResponse: buildErrorBody(statusCode, errMsg), + cacheSource: "upstream", + }); + persistFailureUsage(statusCode, "context_overflow"); + return createErrorResult( + statusCode, + errMsg, + retryAfterMs, + upstreamErrorCode, + upstreamErrorType, + upstreamErrorBody, + { passthrough: sourceFormat === FORMATS.CLAUDE } + ); + } + } catch { persistAttemptLogs({ status: statusCode, error: errMsg, @@ -4684,7 +4710,7 @@ export async function handleChatCore({ { passthrough: sourceFormat === FORMATS.CLAUDE } ); } - } catch { + } else { persistAttemptLogs({ status: statusCode, error: errMsg, @@ -4713,7 +4739,14 @@ export async function handleChatCore({ clientResponse: buildErrorBody(statusCode, errMsg), cacheSource: "upstream", }); - persistFailureUsage(statusCode, "context_overflow"); + persistFailureUsage(statusCode, `upstream_${statusCode}`); + + // Emergency budget fallback is orchestrated exclusively by the routing layer + // (src/sse/handlers/chat.ts), which resolves credentials FOR the emergency + // provider through account selection. The executor-level hop that used to + // live here re-sent the FAILING provider's credentials to the emergency + // provider's endpoint (e.g. the OpenAI API key to integrate.api.nvidia.com) + // — a cross-provider credential leak that also never succeeded upstream. return createErrorResult( statusCode, errMsg, @@ -4724,41 +4757,18 @@ export async function handleChatCore({ { passthrough: sourceFormat === FORMATS.CLAUDE } ); } - } else { - persistAttemptLogs({ - status: statusCode, - error: errMsg, - providerRequest: finalBody || translatedBody, - providerResponse: upstreamErrorBody, - clientResponse: buildErrorBody(statusCode, errMsg), - cacheSource: "upstream", - }); - persistFailureUsage(statusCode, `upstream_${statusCode}`); - - // Emergency budget fallback is orchestrated exclusively by the routing layer - // (src/sse/handlers/chat.ts), which resolves credentials FOR the emergency - // provider through account selection. The executor-level hop that used to - // live here re-sent the FAILING provider's credentials to the emergency - // provider's endpoint (e.g. the OpenAI API key to integrate.api.nvidia.com) - // — a cross-provider credential leak that also never succeeded upstream. - return createErrorResult( - statusCode, - errMsg, - retryAfterMs, - upstreamErrorCode, - upstreamErrorType, - upstreamErrorBody, - { passthrough: sourceFormat === FORMATS.CLAUDE } - ); + // ── End T5 ─────────────────────────────────────────────────────────────── } - // ── End T5 ─────────────────────────────────────────────────────────────── - } } // Non-streaming response if (!stream) { try { - const runNonStreamingPipeline = async ({ policy, model: pipelineModel, translatedBody: wireBody }) => { + const runNonStreamingPipeline = async ({ + policy, + model: pipelineModel, + translatedBody: wireBody, + }) => { translatedBody = wireBody as typeof translatedBody; currentModel = pipelineModel; triedModels.add(pipelineModel); @@ -4852,341 +4862,486 @@ export async function handleChatCore({ sendProviderAttempt: (modelToCall, allowDedup) => executeProviderRequest(modelToCall, allowDedup), }); - }; + }; - let toolLoopRan = false; - let toolLoopUsage = null; - let legResult = await runNonStreamingProviderLeg({ - phase: "initial", - sourceBody: (body || {}) as Record, - expectedConnectionId: managedLease + let toolLoopRan = false; + let toolLoopUsage = null; + let legResult = await runNonStreamingProviderLeg({ + phase: "initial", + sourceBody: (body || {}) as Record, + expectedConnectionId: managedLease + ? String(getCurrentConnectionId() || connectionId || "") || undefined + : undefined, + allowAccountRotation: !managedLease && comboStrategy !== "context-relay", + allowModelFallback: true, + executeProviderRequest: (modelToCall, allowDedup) => + executeProviderRequest(modelToCall, allowDedup), + runProviderExecution: runNonStreamingPipeline, + setRequestWireState: ({ translatedBody: nextBody, effectiveModel: nextModel }) => { + translatedBody = nextBody as typeof translatedBody; + currentModel = nextModel; + triedModels.add(nextModel); + }, + sourceFormat, + targetFormat, + clientResponseFormat, + provider, + model: effectiveModel, + connectionId: String(getCurrentConnectionId() || connectionId || ""), + getCurrentConnectionId: () => getCurrentConnectionId() || undefined, + effectiveModel: currentModel, + translatedBody: translatedBody as Record, + toolNameMap, + requestToolIdentityMap, + reasoningCacheScope, + clientHeaders: clientRawRequest?.headers ?? null, + isClaudeCodeCompatible, + log, + }); + + if (legResult.kind === "error") { + const err = legResult.result; + const captured = providerRequestCapture.latest?.() ?? null; + finalBody = captured?.body ?? finalBody ?? translatedBody; + if (captured) { + reqLogger.logTargetRequest(captured.url, captured.headers, captured.body); + } + reqLogger.logError(new Error(err.error || "Provider request failed"), finalBody); + const isNetworkThrow = Boolean(err.originalError); + if (err.response && !isNetworkThrow) { + reqLogger.logProviderResponse( + err.status, + err.response.statusText || "Error", + err.response.headers, + err.response + ); + } + appendRequestLog({ + model, + provider, + connectionId, + status: `FAILED ${err.status}`, + }).catch(() => {}); + persistAttemptLogs({ + status: err.status, + error: err.error || "Provider request failed", + providerRequest: finalBody || translatedBody, + providerResponse: isNetworkThrow ? undefined : err.response, + clientResponse: buildErrorBody(err.status, err.error || "Provider request failed"), + cacheSource: "upstream", + }); + persistFailureUsage(err.status, err.errorCode || `upstream_${err.status}`); + trackPendingRequest(model, provider, connectionId, false); + return err; + } + + pipelineRecovered = true; + const expectedConn = managedLease ? String(getCurrentConnectionId() || connectionId || "") || undefined - : undefined, - allowAccountRotation: !managedLease && comboStrategy !== "context-relay", - allowModelFallback: true, - executeProviderRequest: (modelToCall, allowDedup) => - executeProviderRequest(modelToCall, allowDedup), - runProviderExecution: runNonStreamingPipeline, - setRequestWireState: ({ translatedBody: nextBody, effectiveModel: nextModel }) => { - translatedBody = nextBody as typeof translatedBody; - currentModel = nextModel; - triedModels.add(nextModel); - }, - sourceFormat, - targetFormat, - clientResponseFormat, - provider, - model: effectiveModel, - connectionId: String(getCurrentConnectionId() || connectionId || ""), - getCurrentConnectionId: () => getCurrentConnectionId() || undefined, - effectiveModel: currentModel, - translatedBody: translatedBody as Record, - toolNameMap, - requestToolIdentityMap, - reasoningCacheScope, - clientHeaders: clientRawRequest?.headers ?? null, - isClaudeCodeCompatible, - log, - }); + : undefined; + const loopApply = await applyServerOwnedToolLoopIfNeeded({ + enabled: isServerOwnedToolLoopEnabled(), + stream, + isResponsesEndpoint, + sourceFormat, + initialLeg: legResult, + sourceBody: (body || {}) as Record, + skillsModelId: getSkillsModelIdForFormat(sourceFormat), + executionContext: { + apiKeyId: memoryOwnerId || "local", + sessionId: pipelineSessionId, + requestId: skillRequestId, + requestIdentity: derivePostInjectionRequestIdentity({ + apiKeyId: memoryOwnerId || "local", + headers: clientRawRequest?.headers ?? null, + skillRequestId, + postInjectionBody: (body || {}) as Record, + }), + builtinToolNames: injectionResult.builtinToolNames, + injectedCustomSkillNames: injectionResult.injectedCustomSkillNames, + customSkillExecutionEnabled: + Boolean(memoryOwnerId) && memorySettings?.skillsEnabled === true, + executionFenceEnabled: true, + provider, + model: effectiveModel, + }, + abortSignal: clientRawRequest?.signal, + expectedConnectionId: expectedConn, + followUpLeg: async (nextSourceBody) => { + translatedBody = translateRequest( + sourceFormat, + targetFormat, + model, + { ...nextSourceBody }, + false, + credentials, + provider, + reqLogger, + { + normalizeToolCallId: getModelNormalizeToolCallId( + provider || "", + model || "", + sourceFormat + ), + preserveDeveloperRole: getModelPreserveOpenAIDeveloperRole( + provider || "", + model || "", + sourceFormat + ), + preserveCacheControl, + signatureNamespace: connectionId, + copilotClient: copilotCompatibleReasoning, + reasoningCacheScope, + } + ); + return runNonStreamingProviderLeg( + followUpLegInput( + { + executeProviderRequest: (modelToCall, allowDedup) => + executeProviderRequest(modelToCall, allowDedup), + runProviderExecution: runNonStreamingPipeline, + setRequestWireState: ({ translatedBody: nextBody, effectiveModel: nextModel }) => { + translatedBody = nextBody as typeof translatedBody; + currentModel = nextModel; + triedModels.add(nextModel); + }, + sourceFormat, + targetFormat, + clientResponseFormat, + provider, + model: effectiveModel, + connectionId: String(getCurrentConnectionId() || connectionId || ""), + getCurrentConnectionId: () => getCurrentConnectionId() || undefined, + effectiveModel: currentModel, + translatedBody: translatedBody as Record, + toolNameMap, + requestToolIdentityMap, + reasoningCacheScope, + clientHeaders: clientRawRequest?.headers ?? null, + isClaudeCodeCompatible, + log, + }, + nextSourceBody, + expectedConn + ) + ); + }, + logReceipt: (receipt) => reqLogger.logToolLoopReceipt(receipt), + }); + if (loopApply.kind === "error") { + return await finalizeToolLoopError({ + loop: loopApply.loop, + model, + provider, + connectionId, + providerRequest: loopApply.loop.finalProviderRequest || finalBody || translatedBody, + persistFailureUsage, + persistAttemptLogs, + trackPendingRequest, + }); + } + // `legResult` is declared as the full NonStreamingProviderLegResult union. The + // `kind === "error"` guard above narrows it to the ok variant, but the conditional + // reassignment below widens it back to the declared type, so every field read past + // this point lost the narrowing — 13 TS2339 diagnostics under + // tsconfig.typecheck-api.json, which pulls chatCore.ts in through the route while + // tsconfig.typecheck-core.json does not. Pin the ok variant in its own binding: + // `loopApply.leg` is already `NonStreamingProviderLegResult & { kind: "ok" }`, + // so no cast is involved. + let okLeg: NonStreamingProviderLegResult & { kind: "ok" } = legResult; + if (loopApply.kind === "ok") { + toolLoopRan = true; + toolLoopUsage = loopApply.usage; + okLeg = loopApply.leg; + } - if (legResult.kind === "error") { - const err = legResult.result; - const captured = providerRequestCapture.latest?.() ?? null; - finalBody = captured?.body ?? finalBody ?? translatedBody; - if (captured) { - reqLogger.logTargetRequest(captured.url, captured.headers, captured.body); + if (okLeg.upstreamResponse) { + providerResponse = okLeg.upstreamResponse; + providerHeaders = normalizeHeaders(okLeg.upstreamResponse.headers); + } else { + providerResponse = new Response(null, { + status: 200, + headers: okLeg.headers, + }); + providerHeaders = normalizeHeaders(okLeg.headers); } - reqLogger.logError(new Error(err.error || "Provider request failed"), finalBody); - const isNetworkThrow = Boolean(err.originalError); - if (err.response && !isNetworkThrow) { - reqLogger.logProviderResponse( - err.status, - err.response.statusText || "Error", - err.response.headers, - err.response - ); + finalBody = providerRequestCapture.body(okLeg.providerRequest || translatedBody); + const capturedOk = providerRequestCapture.latest?.(); + reqLogger.logTargetRequest( + okLeg.requestUrl || capturedOk?.url || "", + okLeg.requestHeaders || capturedOk?.headers || {}, + capturedOk?.body ?? finalBody + ); + const responseBody = okLeg.providerBody; + const responsePayloadFormat = okLeg.responsePayloadFormat; + const looksLikeSSE = okLeg.looksLikeSSE; + let translatedResponse = okLeg.response; + const memoryExtractionResponse = okLeg.responseForMemoryExtraction; + reqLogger.logProviderResponse( + 200, + "OK", + providerResponse.headers, + looksLikeSSE + ? { _streamed: true, _format: "sse-json", summary: responseBody } + : responseBody + ); + effectiveServiceTier = resolveReportedServiceTier(responseBody) ?? effectiveServiceTier; + if (onRequestSuccess) { + await onRequestSuccess(); } + const successConnectionId = getCurrentConnectionId(); + await maybeSyncClaudeExtraUsageState({ + provider, + connectionId: successConnectionId, + providerSpecificData: credentials?.providerSpecificData, + log, + }); + const usage = toolLoopUsage ?? extractUsageFromResponse(responseBody, provider); + const cacheUsageLogMeta = buildCacheUsageLogMeta(usage); + if (usage && typeof usage === "object") { + attachCompressionUsageReceiptAfterAnalytics(usage as Record, "provider"); + if (provider === "gemini") { + const promptTokens = + typeof (usage as Record).prompt_tokens === "number" + ? ((usage as Record).prompt_tokens as number) + : 0; + if (promptTokens > 0) incrementTokenUsage(model, promptTokens); + } + } + recordContextEditingTelemetryHook({ + contextEditingEnabled, + provider, + responseBody, + skillRequestId, + log, + }); appendRequestLog({ model, provider, - connectionId, - status: `FAILED ${err.status}`, + connectionId: successConnectionId, + tokens: usage, + status: "200 OK", }).catch(() => {}); - persistAttemptLogs({ - status: err.status, - error: err.error || "Provider request failed", - providerRequest: finalBody || translatedBody, - providerResponse: isNetworkThrow ? undefined : err.response, - clientResponse: buildErrorBody(err.status, err.error || "Provider request failed"), - cacheSource: "upstream", - }); - persistFailureUsage( - err.status, - err.errorCode || `upstream_${err.status}` - ); - trackPendingRequest(model, provider, connectionId, false); - return err; - } - - pipelineRecovered = true; - const expectedConn = managedLease - ? String(getCurrentConnectionId() || connectionId || "") || undefined - : undefined; - const loopApply = await applyServerOwnedToolLoopIfNeeded({ - enabled: isServerOwnedToolLoopEnabled(), - stream, - isResponsesEndpoint, - sourceFormat, - initialLeg: legResult, - sourceBody: (body || {}) as Record, - skillsModelId: getSkillsModelIdForFormat(sourceFormat), - executionContext: { - apiKeyId: memoryOwnerId || "local", - sessionId: pipelineSessionId, - requestId: skillRequestId, - requestIdentity: derivePostInjectionRequestIdentity({ - apiKeyId: memoryOwnerId || "local", - headers: clientRawRequest?.headers ?? null, - skillRequestId, - postInjectionBody: (body || {}) as Record, - }), - builtinToolNames: injectionResult.builtinToolNames, - injectedCustomSkillNames: injectionResult.injectedCustomSkillNames, - customSkillExecutionEnabled: - Boolean(memoryOwnerId) && memorySettings?.skillsEnabled === true, - executionFenceEnabled: true, + recordNonStreamingUsageStats(usage, { + traceEnabled, provider, - model: effectiveModel, - }, - abortSignal: clientRawRequest?.signal, - expectedConnectionId: expectedConn, - followUpLeg: async (nextSourceBody) => { - translatedBody = translateRequest( - sourceFormat, - targetFormat, - model, - { ...nextSourceBody }, - false, - credentials, - provider, - reqLogger, + connectionId: successConnectionId, + model, + startTime, + apiKeyInfo, + effectiveServiceTier, + isCombo, + comboStrategy, + endpoint: endpointPath, + }); + + // #12150 P1b surface 3 (fix round 1): a video-bridge-observed request's + // request- AND response-derived text both carry the full transcript (the + // flattened description on the request side, the model's own reply on + // the response side) — neither may populate durable Memory. See + // runMemoryExtractionGate for the shared gate + extraction wiring, unit + // tested directly in tests/unit/video-bridge-memory-suppression.test.ts. + runMemoryExtractionGate({ + memoryOwnerId, + memorySettings, + videoBridgeObserved, + pipelineSessionId, + requestBody: body as Record, + responseBody: memoryExtractionResponse as Record | null, + extractFacts, + log, + }); + + const customSkillExecutionEnabled = + Boolean(memoryOwnerId) && memorySettings?.skillsEnabled === true; + const builtinToolNames = [ + webSearchFallbackPlan.toolName, + webFetchFallbackPlan.toolName, + ...(memoryOwnerId && memorySettings?.enabled ? MEMORY_BUILTIN_TOOL_NAMES : []), + ].filter((name): name is string => Boolean(name)); + if (!toolLoopRan && (customSkillExecutionEnabled || builtinToolNames.length > 0)) { + const skillSessionId = pipelineSessionId; + + translatedResponse = await handleToolCallExecution( + translatedResponse, + getSkillsModelIdForFormat(sourceFormat), { - normalizeToolCallId: getModelNormalizeToolCallId(provider || "", model || "", sourceFormat), - preserveDeveloperRole: getModelPreserveOpenAIDeveloperRole( - provider || "", - model || "", - sourceFormat - ), - preserveCacheControl, - signatureNamespace: connectionId, - copilotClient: copilotCompatibleReasoning, - reasoningCacheScope, + apiKeyId: memoryOwnerId || "local", + sessionId: skillSessionId, + requestId: skillRequestId, + builtinToolNames, + customSkillExecutionEnabled, + provider, + model: effectiveModel, } ); - return runNonStreamingProviderLeg( - followUpLegInput( - { - executeProviderRequest: (modelToCall, allowDedup) => - executeProviderRequest(modelToCall, allowDedup), - runProviderExecution: runNonStreamingPipeline, - setRequestWireState: ({ translatedBody: nextBody, effectiveModel: nextModel }) => { - translatedBody = nextBody as typeof translatedBody; - currentModel = nextModel; - triedModels.add(nextModel); - }, - sourceFormat, - targetFormat, - clientResponseFormat, - provider, - model: effectiveModel, - connectionId: String(getCurrentConnectionId() || connectionId || ""), - getCurrentConnectionId: () => getCurrentConnectionId() || undefined, - effectiveModel: currentModel, - translatedBody: translatedBody as Record, - toolNameMap, - requestToolIdentityMap, - reasoningCacheScope, - clientHeaders: clientRawRequest?.headers ?? null, - isClaudeCodeCompatible, - log, - }, - nextSourceBody, - expectedConn - ) - ); - }, - logReceipt: (receipt) => reqLogger.logToolLoopReceipt(receipt), - }); - if (loopApply.kind === "error") { - return await finalizeToolLoopError({ - loop: loopApply.loop, + } + + const guardrailContext = buildPostCallGuardrailContext({ + apiKeyInfo, + body, + clientRawRequest, + log, model, provider, - connectionId, - providerRequest: loopApply.loop.finalProviderRequest || finalBody || translatedBody, - persistFailureUsage, - persistAttemptLogs, - trackPendingRequest, + responsePayloadFormat, + clientResponseFormat, }); - } - if (loopApply.kind === "ok") { - toolLoopRan = true; - toolLoopUsage = loopApply.usage; - legResult = loopApply.leg; - } - - if (legResult.upstreamResponse) { - providerResponse = legResult.upstreamResponse; - providerHeaders = normalizeHeaders(legResult.upstreamResponse.headers); - } else { - providerResponse = new Response(null, { - status: 200, - headers: legResult.headers, - }); - providerHeaders = normalizeHeaders(legResult.headers); - } - finalBody = providerRequestCapture.body(legResult.providerRequest || translatedBody); - const capturedOk = providerRequestCapture.latest?.(); - reqLogger.logTargetRequest( - legResult.requestUrl || capturedOk?.url || "", - legResult.requestHeaders || capturedOk?.headers || {}, - capturedOk?.body ?? finalBody - ); - const responseBody = legResult.providerBody; - const responsePayloadFormat = legResult.responsePayloadFormat; - const looksLikeSSE = legResult.looksLikeSSE; - let translatedResponse = legResult.response; - const memoryExtractionResponse = legResult.responseForMemoryExtraction; - reqLogger.logProviderResponse( - 200, - "OK", - providerResponse.headers, - looksLikeSSE - ? { _streamed: true, _format: "sse-json", summary: responseBody } - : responseBody - ); - effectiveServiceTier = resolveReportedServiceTier(responseBody) ?? effectiveServiceTier; - if (onRequestSuccess) { - await onRequestSuccess(); - } - const successConnectionId = getCurrentConnectionId(); - await maybeSyncClaudeExtraUsageState({ - provider, - connectionId: successConnectionId, - providerSpecificData: credentials?.providerSpecificData, - log, - }); - const usage = toolLoopUsage ?? extractUsageFromResponse(responseBody, provider); - const cacheUsageLogMeta = buildCacheUsageLogMeta(usage); - if (usage && typeof usage === "object") { - attachCompressionUsageReceiptAfterAnalytics(usage as Record, "provider"); - if (provider === "gemini") { - const promptTokens = - typeof (usage as Record).prompt_tokens === "number" - ? ((usage as Record).prompt_tokens as number) - : 0; - if (promptTokens > 0) incrementTokenUsage(model, promptTokens); - } - } - recordContextEditingTelemetryHook({ - contextEditingEnabled, - provider, - responseBody, - skillRequestId, - log, - }); - appendRequestLog({ - model, - provider, - connectionId: successConnectionId, - tokens: usage, - status: "200 OK", - }).catch(() => {}); - recordNonStreamingUsageStats(usage, { - traceEnabled, - provider, - connectionId: successConnectionId, - model, - startTime, - apiKeyInfo, - effectiveServiceTier, - isCombo, - comboStrategy, - endpoint: endpointPath, - }); - - // #12150 P1b surface 3 (fix round 1): a video-bridge-observed request's - // request- AND response-derived text both carry the full transcript (the - // flattened description on the request side, the model's own reply on - // the response side) — neither may populate durable Memory. See - // runMemoryExtractionGate for the shared gate + extraction wiring, unit - // tested directly in tests/unit/video-bridge-memory-suppression.test.ts. - runMemoryExtractionGate({ - memoryOwnerId, - memorySettings, - videoBridgeObserved, - pipelineSessionId, - requestBody: body as Record, - responseBody: memoryExtractionResponse as Record | null, - extractFacts, - log, - }); - - const customSkillExecutionEnabled = - Boolean(memoryOwnerId) && memorySettings?.skillsEnabled === true; - const builtinToolNames = [ - webSearchFallbackPlan.toolName, - webFetchFallbackPlan.toolName, - ...(memoryOwnerId && memorySettings?.enabled ? MEMORY_BUILTIN_TOOL_NAMES : []), - ].filter((name): name is string => Boolean(name)); - if (!toolLoopRan && (customSkillExecutionEnabled || builtinToolNames.length > 0)) { - const skillSessionId = pipelineSessionId; - - translatedResponse = await handleToolCallExecution( + const postCallGuardrails = await guardrailRegistry.runPostCallHooks( translatedResponse, - getSkillsModelIdForFormat(sourceFormat), - { - apiKeyId: memoryOwnerId || "local", - sessionId: skillSessionId, - requestId: skillRequestId, - builtinToolNames, - customSkillExecutionEnabled, - provider, - model: effectiveModel, - } + guardrailContext ); - } + translatedResponse = postCallGuardrails.response; - const guardrailContext = buildPostCallGuardrailContext({ - apiKeyInfo, - body, - clientRawRequest, - log, - model, - provider, - responsePayloadFormat, - clientResponseFormat, - }); - const postCallGuardrails = await guardrailRegistry.runPostCallHooks( - translatedResponse, - guardrailContext - ); - translatedResponse = postCallGuardrails.response; + const responseUsage = isJsonRecord(usage) + ? usage + : isJsonRecord(translatedResponse.usage) + ? translatedResponse.usage + : null; + const costUsage = normalizeUsage(responseUsage); + const estimatedCost = costUsage + ? await calculateCost(provider, model, costUsage, { serviceTier: effectiveServiceTier }) + : 0; - const responseUsage = isJsonRecord(usage) - ? usage - : isJsonRecord(translatedResponse.usage) - ? translatedResponse.usage - : null; - const costUsage = normalizeUsage(responseUsage); - const estimatedCost = costUsage - ? await calculateCost(provider, model, costUsage, { serviceTier: effectiveServiceTier }) - : 0; + if (postCallGuardrails.blocked) { + const guardrailMessage = postCallGuardrails.message || "Response blocked by guardrail"; + persistAttemptLogs({ + status: HTTP_STATUS.BAD_REQUEST, + tokens: usage, + responseBody, + providerRequest: finalBody || translatedBody, + providerResponse: looksLikeSSE + ? { + _streamed: true, + _format: "sse-json", + summary: responseBody, + } + : responseBody, + clientResponse: buildErrorBody(HTTP_STATUS.BAD_REQUEST, guardrailMessage), + claudeCacheMeta: claudePromptCacheLogMeta, + claudeCacheUsageMeta: cacheUsageLogMeta, + cacheSource: "upstream", + }); + if (apiKeyInfo?.id && estimatedCost > 0) { + recordCost(apiKeyInfo.id, estimatedCost); + } + log?.warn?.( + "GUARDRAIL", + `Response blocked by ${postCallGuardrails.guardrail || "guardrail"}: ${guardrailMessage}` + ); + finalizePendingScope(pendingScope, { + providerResponse: responseBody, + clientResponse: translatedResponse, + }); + return createErrorResult(HTTP_STATUS.BAD_REQUEST, guardrailMessage); + } - if (postCallGuardrails.blocked) { - const guardrailMessage = postCallGuardrails.message || "Response blocked by guardrail"; + // Validate the *translated* response actually carries client-usable output. + // isEmptyContentResponse (above) runs on the raw responseBody before translation; + // this check runs after translation + sanitization + tool-call execution to catch + // cases where a provider returns a structurally valid raw body that translates into + // choices:[] or output:[] with no usable content (Responses API shape included). + const malformedTranslatedReason = detectMalformedNonStream(translatedResponse); + if (malformedTranslatedReason) { + const totalLatency = Date.now() - startTime; + const rawBytes = (() => { + try { + return JSON.stringify(responseBody || {}).length; + } catch { + return -1; + } + })(); + reportMalformed200({ + mode: "nonstream", + provider, + model, + connectionId, + reason: malformedTranslatedReason, + recvBytes: rawBytes, + recvLines: -1, + emitted: -1, + events: {}, + ttftMs: totalLatency, + elapsedMs: totalLatency, + }); + appendRequestLog({ + model, + provider, + connectionId, + status: `FAILED ${HTTP_STATUS.BAD_GATEWAY}`, + }).catch(() => {}); + const malformed = describeMalformedNonStream(translatedResponse, malformedTranslatedReason); + const malformedMessage = `[${provider}/${model}] ${malformed.message}`; + const malformedClientBody = buildErrorBody( + HTTP_STATUS.BAD_GATEWAY, + malformedMessage, + undefined, + { code: malformed.code, type: malformed.type } + ); + persistAttemptLogs({ + status: HTTP_STATUS.BAD_GATEWAY, + tokens: usage, + responseBody, + providerRequest: finalBody || translatedBody, + providerResponse: looksLikeSSE + ? { _streamed: true, _format: "sse-json", summary: responseBody } + : responseBody, + clientResponse: malformedClientBody, + claudeCacheMeta: claudePromptCacheLogMeta, + claudeCacheUsageMeta: cacheUsageLogMeta, + cacheSource: "upstream", + }); + persistFailureUsage(HTTP_STATUS.BAD_GATEWAY, "malformed_translated_response"); + trackPendingRequest(model, provider, pendingConnId, false); + // Routing event (feedback foundation) — record the malformed outcome so + // the quality tracker de-prioritizes this model over time. + void emitRoutingEvent( + createRoutingEvent({ + requestId: traceId || pendingRequestId || "unknown", + provider: provider || "unknown", + model: model || "unknown", + strategy: isCombo ? (comboStrategy ?? "combo") : "direct", + latencyMs: Date.now() - startTime, + ttftMs: null, + inputTokens: null, + outputTokens: null, + cost: null, + retries: 0, + fallbackUsed: false, // combo-level fallback tracked by decisionTrace + outcome: "malformed", + status: HTTP_STATUS.BAD_GATEWAY, + finishReason: routingFinishReason(translatedResponse), + connectionId: credentials?.connectionId ?? null, + }) + ); + return createErrorResult( + HTTP_STATUS.BAD_GATEWAY, + malformedMessage, + null, + malformed.code, + malformed.type + ); + } + + // ── Phase 9.1: Cache store (non-streaming, temp=0) ── + storeSemanticCacheResponse({ + enabled: semanticCacheEnabled, + body: bodyForCacheWrite, + headers: clientRawRequest?.headers, + translatedResponse, + model, + apiKeyId: apiKeyInfo?.id ?? undefined, + usage, + log, + }); + + // ── Phase 9.2: Save for idempotency ── + // Reuse the key resolved by checkIdempotencyCache() above (single derivation per + // request). (#3821-review LEDGER-6) + saveIdempotency(idempotencyKey, translatedResponse, 200); + reqLogger.logConvertedResponse(translatedResponse); persistAttemptLogs({ - status: HTTP_STATUS.BAD_REQUEST, + status: 200, tokens: usage, responseBody, providerRequest: finalBody || translatedBody, @@ -5197,7 +5352,7 @@ export async function handleChatCore({ summary: responseBody, } : responseBody, - clientResponse: buildErrorBody(HTTP_STATUS.BAD_REQUEST, guardrailMessage), + clientResponse: translatedResponse, claudeCacheMeta: claudePromptCacheLogMeta, claudeCacheUsageMeta: cacheUsageLogMeta, cacheSource: "upstream", @@ -5205,76 +5360,61 @@ export async function handleChatCore({ if (apiKeyInfo?.id && estimatedCost > 0) { recordCost(apiKeyInfo.id, estimatedCost); } - log?.warn?.( - "GUARDRAIL", - `Response blocked by ${postCallGuardrails.guardrail || "guardrail"}: ${guardrailMessage}` - ); + + // === Quota Share POST-hook (B/F7) — fire-and-forget, fail-open === + await scheduleQuotaShareConsumption({ + apiKeyId: apiKeyInfo?.id, + connectionId: credentials?.connectionId, + provider, + model, + usage, + estimatedCost, + log, + }); + // === /Quota Share POST-hook === + + // ── Gamification event (fire-and-forget) ── + await emitRequestGamificationEvent({ apiKeyId: apiKeyInfo?.id, model, provider }); + finalizePendingScope(pendingScope, { providerResponse: responseBody, clientResponse: translatedResponse, }); - return createErrorResult(HTTP_STATUS.BAD_REQUEST, guardrailMessage); - } + const responseHeaders = buildNonStreamingResponseHeaders({ + provider, + model, + startTime, + responseUsage, + estimatedCost, + requestId: skillRequestId, + compressionResponseMeta, + comboStrategy, + }); + // #6426: align response body `model` with the `X-OmniRoute-Model` header + // (both must be the resolved backend model). Some upstreams (notably legacy + // /v1/completions text-completion path) return a body `model` field that + // differs from the resolved backend id we advertised in the header, leaving + // strict clients unable to reconcile the two. Rewrite body.model to `model` + // FIRST, then let #1311 echo override it when the opt-in setting is on. + if (typeof model === "string" && model) echoModelInObject(translatedResponse, model); + // #1311: echo the requested alias/combo name in the non-streaming response model. + if (echoModel) echoModelInObject(translatedResponse, echoModel); - // Validate the *translated* response actually carries client-usable output. - // isEmptyContentResponse (above) runs on the raw responseBody before translation; - // this check runs after translation + sanitization + tool-call execution to catch - // cases where a provider returns a structurally valid raw body that translates into - // choices:[] or output:[] with no usable content (Responses API shape included). - const malformedTranslatedReason = detectMalformedNonStream(translatedResponse); - if (malformedTranslatedReason) { - const totalLatency = Date.now() - startTime; - const rawBytes = (() => { - try { - return JSON.stringify(responseBody || {}).length; - } catch { - return -1; - } - })(); - reportMalformed200({ - mode: "nonstream", - provider, - model, - connectionId, - reason: malformedTranslatedReason, - recvBytes: rawBytes, - recvLines: -1, - emitted: -1, - events: {}, - ttftMs: totalLatency, - elapsedMs: totalLatency, - }); - appendRequestLog({ + // ── Plugin onResponse hook (fire-and-forget) ── + // #8395: the streaming branch below already calls this; the non-streaming + // (stream:false) branch returned without it, so onResponse never fired for + // non-streaming requests at all. + await runPluginOnResponseHook({ + requestId: traceId, + body, model, provider, - connectionId, - status: `FAILED ${HTTP_STATUS.BAD_GATEWAY}`, - }).catch(() => {}); - const malformed = describeMalformedNonStream(translatedResponse, malformedTranslatedReason); - const malformedMessage = `[${provider}/${model}] ${malformed.message}`; - const malformedClientBody = buildErrorBody( - HTTP_STATUS.BAD_GATEWAY, - malformedMessage, - undefined, - { code: malformed.code, type: malformed.type } - ); - persistAttemptLogs({ - status: HTTP_STATUS.BAD_GATEWAY, - tokens: usage, - responseBody, - providerRequest: finalBody || translatedBody, - providerResponse: looksLikeSSE - ? { _streamed: true, _format: "sse-json", summary: responseBody } - : responseBody, - clientResponse: malformedClientBody, - claudeCacheMeta: claudePromptCacheLogMeta, - claudeCacheUsageMeta: cacheUsageLogMeta, - cacheSource: "upstream", + apiKeyInfo, + headers: clientRawRequest?.headers, + response: { status: 200, data: translatedResponse }, }); - persistFailureUsage(HTTP_STATUS.BAD_GATEWAY, "malformed_translated_response"); - trackPendingRequest(model, provider, pendingConnId, false); - // Routing event (feedback foundation) — record the malformed outcome so - // the quality tracker de-prioritizes this model over time. + + // Routing event (feedback foundation) — fire-and-forget, cheap. void emitRoutingEvent( createRoutingEvent({ requestId: traceId || pendingRequestId || "unknown", @@ -5283,158 +5423,38 @@ export async function handleChatCore({ strategy: isCombo ? (comboStrategy ?? "combo") : "direct", latencyMs: Date.now() - startTime, ttftMs: null, - inputTokens: null, - outputTokens: null, - cost: null, + inputTokens: + usage && typeof usage === "object" + ? (() => { + const promptTokens = (usage as Record).prompt_tokens; + return typeof promptTokens === "number" && Number.isFinite(promptTokens) + ? promptTokens + : null; + })() + : null, + outputTokens: + usage && typeof usage === "object" + ? (() => { + const completionTokens = (usage as Record).completion_tokens; + return typeof completionTokens === "number" && Number.isFinite(completionTokens) + ? completionTokens + : null; + })() + : null, + cost: Number.isFinite(estimatedCost) ? estimatedCost : null, retries: 0, fallbackUsed: false, // combo-level fallback tracked by decisionTrace - outcome: "malformed", - status: HTTP_STATUS.BAD_GATEWAY, + outcome: "success", + status: 200, finishReason: routingFinishReason(translatedResponse), connectionId: credentials?.connectionId ?? null, }) ); - return createErrorResult( - HTTP_STATUS.BAD_GATEWAY, - malformedMessage, - null, - malformed.code, - malformed.type - ); - } - // ── Phase 9.1: Cache store (non-streaming, temp=0) ── - storeSemanticCacheResponse({ - enabled: semanticCacheEnabled, - body: bodyForCacheWrite, - headers: clientRawRequest?.headers, - translatedResponse, - model, - apiKeyId: apiKeyInfo?.id ?? undefined, - usage, - log, - }); - - // ── Phase 9.2: Save for idempotency ── - // Reuse the key resolved by checkIdempotencyCache() above (single derivation per - // request). (#3821-review LEDGER-6) - saveIdempotency(idempotencyKey, translatedResponse, 200); - reqLogger.logConvertedResponse(translatedResponse); - persistAttemptLogs({ - status: 200, - tokens: usage, - responseBody, - providerRequest: finalBody || translatedBody, - providerResponse: looksLikeSSE - ? { - _streamed: true, - _format: "sse-json", - summary: responseBody, - } - : responseBody, - clientResponse: translatedResponse, - claudeCacheMeta: claudePromptCacheLogMeta, - claudeCacheUsageMeta: cacheUsageLogMeta, - cacheSource: "upstream", - }); - if (apiKeyInfo?.id && estimatedCost > 0) { - recordCost(apiKeyInfo.id, estimatedCost); - } - - // === Quota Share POST-hook (B/F7) — fire-and-forget, fail-open === - await scheduleQuotaShareConsumption({ - apiKeyId: apiKeyInfo?.id, - connectionId: credentials?.connectionId, - provider, - model, - usage, - estimatedCost, - log, - }); - // === /Quota Share POST-hook === - - // ── Gamification event (fire-and-forget) ── - await emitRequestGamificationEvent({ apiKeyId: apiKeyInfo?.id, model, provider }); - - finalizePendingScope(pendingScope, { - providerResponse: responseBody, - clientResponse: translatedResponse, - }); - const responseHeaders = buildNonStreamingResponseHeaders({ - provider, - model, - startTime, - responseUsage, - estimatedCost, - requestId: skillRequestId, - compressionResponseMeta, - comboStrategy, - }); - // #6426: align response body `model` with the `X-OmniRoute-Model` header - // (both must be the resolved backend model). Some upstreams (notably legacy - // /v1/completions text-completion path) return a body `model` field that - // differs from the resolved backend id we advertised in the header, leaving - // strict clients unable to reconcile the two. Rewrite body.model to `model` - // FIRST, then let #1311 echo override it when the opt-in setting is on. - if (typeof model === "string" && model) echoModelInObject(translatedResponse, model); - // #1311: echo the requested alias/combo name in the non-streaming response model. - if (echoModel) echoModelInObject(translatedResponse, echoModel); - - // ── Plugin onResponse hook (fire-and-forget) ── - // #8395: the streaming branch below already calls this; the non-streaming - // (stream:false) branch returned without it, so onResponse never fired for - // non-streaming requests at all. - await runPluginOnResponseHook({ - requestId: traceId, - body, - model, - provider, - apiKeyInfo, - headers: clientRawRequest?.headers, - response: { status: 200, data: translatedResponse }, - }); - - // Routing event (feedback foundation) — fire-and-forget, cheap. - void emitRoutingEvent( - createRoutingEvent({ - requestId: traceId || pendingRequestId || "unknown", - provider: provider || "unknown", - model: model || "unknown", - strategy: isCombo ? (comboStrategy ?? "combo") : "direct", - latencyMs: Date.now() - startTime, - ttftMs: null, - inputTokens: - usage && typeof usage === "object" - ? (() => { - const promptTokens = (usage as Record).prompt_tokens; - return typeof promptTokens === "number" && Number.isFinite(promptTokens) - ? promptTokens - : null; - })() - : null, - outputTokens: - usage && typeof usage === "object" - ? (() => { - const completionTokens = (usage as Record).completion_tokens; - return typeof completionTokens === "number" && Number.isFinite(completionTokens) - ? completionTokens - : null; - })() - : null, - cost: Number.isFinite(estimatedCost) ? estimatedCost : null, - retries: 0, - fallbackUsed: false, // combo-level fallback tracked by decisionTrace - outcome: "success", - status: 200, - finishReason: routingFinishReason(translatedResponse), - connectionId: credentials?.connectionId ?? null, - }) - ); - - return { - success: true, - response: buildNonStreamingJsonResponse(translatedResponse, responseHeaders), - }; + return { + success: true, + response: buildNonStreamingJsonResponse(translatedResponse, responseHeaders), + }; } catch (error) { trackPendingRequest(model, provider, connectionId, false); if (isManagedLeaseFenceError(error)) return managedLeaseFenceErrorResult(error);