diff --git a/open-sse/config/errorConfig.js b/open-sse/config/errorConfig.js index 71491a4d2c..f03477b943 100644 --- a/open-sse/config/errorConfig.js +++ b/open-sse/config/errorConfig.js @@ -28,6 +28,37 @@ export const DEFAULT_ERROR_MESSAGES = { 504: "Gateway timeout" }; +export const CODEX_REQUEST_SCHEMA_ERROR_CODES = new Set([ + "unknown_parameter", + "unsupported_value", +]); + +export const CODEX_REQUEST_SCHEMA_MESSAGE_PATTERN = /\b(?:unknown[_ ]parameter|unsupported[_ ]value)\b/i; +export const CODEX_REQUEST_SCHEMA_PARAM_ROOTS = new Set([ + "input", + "instructions", + "tools", + "tool_choice", + "parallel_tool_calls", + "stream", + "store", + "reasoning", + "service_tier", + "include", + "prompt_cache_key", + "client_metadata", + "text", +]); +export const CODEX_ITEM_ID_PARAM_PATTERN = /^input\[\d+\]\.id$/; +export const CODEX_ITEM_ID_MESSAGE_PATTERN = /expected an id that begins with ["'`]\w+["'`]/i; + +export const REQUEST_SCHEMA_CLASSIFICATION = Object.freeze({ + category: "request_schema", + accountFallback: false, + cooldownMs: 0, + comboScope: "provider", +}); + // Exponential backoff config for rate limits export const BACKOFF_CONFIG = { base: 2000, diff --git a/open-sse/executors/codex.js b/open-sse/executors/codex.js index 43c36f5350..426016b34a 100644 --- a/open-sse/executors/codex.js +++ b/open-sse/executors/codex.js @@ -5,7 +5,7 @@ import { refreshProviderCredentials, shouldRefreshCredentials, } from "../services/oauthCredentialManager.js"; -import { normalizeResponsesInput } from "../translator/formats/responsesApi.js"; +import { normalizeResponsesInput, normalizeStatelessResponseInput } from "../translator/formats/responsesApi.js"; import { fetchImageAsBase64 } from "../translator/concerns/image.js"; import { getModelUpstreamId } from "../config/providerModels.js"; import { DEFAULT_RETRY_CONFIG, HTTP_STATUS, resolveRetryEntry } from "../config/runtimeConfig.js"; @@ -24,9 +24,6 @@ const CODEX_SSE_USER_OUTPUT_PATTERNS = [ const CODEX_SSE_PEEK_BYTES = 256 * 1024; const CODEX_MODEL_CAPACITY_MESSAGE = "Selected model is at capacity. Please try a different model."; -// Server-generated item id prefixes that Codex /responses cannot resolve when store=false -const SERVER_ID_PATTERN = /^(rs|fc|resp|msg)_/; - // Hosted tool types that Codex/OpenAI Responses executes server-side const CODEX_HOSTED_TOOL_TYPES = new Set([ "image_generation", "web_search", "web_search_preview", "file_search", @@ -54,21 +51,6 @@ function convertSystemToDeveloperRole(body) { } } -// Strip invalid or stored item IDs before sending a store=false request. -function stripStoredItemReferences(body) { - if (!Array.isArray(body.input)) return; - body.input = body.input.filter((item) => { - if (typeof item === "string" && SERVER_ID_PATTERN.test(item)) return false; - if (item && typeof item === "object" && !Array.isArray(item)) { - if (item.type === "item_reference") return false; - // function_call.id is optional input metadata; call_id carries tool-result correlation. - if (item.type === "function_call") delete item.id; - if (typeof item.id === "string" && SERVER_ID_PATTERN.test(item.id)) delete item.id; - } - return true; - }); -} - // Flatten Chat-Completions tool shape into Responses flat format + filter unsupported tools function normalizeCodexTools(body) { if (!Array.isArray(body.tools)) return; @@ -403,8 +385,8 @@ export class CodexExecutor extends BaseExecutor { // Keep system prompts in body.input as role=developer so they stay in the cacheable prefix convertSystemToDeveloperRole(body); - // Strip invalid function-call IDs and stored references that Codex cannot resolve with store=false - stripStoredItemReferences(body); + // Strip optional call item IDs and stored references that store=false cannot resolve. + body.input = normalizeStatelessResponseInput(body.input); // Flatten function tools + drop unsupported types normalizeCodexTools(body); diff --git a/open-sse/handlers/chatCore.js b/open-sse/handlers/chatCore.js index 4f91e020a0..9b32d2b237 100644 --- a/open-sse/handlers/chatCore.js +++ b/open-sse/handlers/chatCore.js @@ -8,7 +8,7 @@ import { refreshWithRetry } from "../services/tokenRefresh.js"; import { createRequestLogger } from "../utils/requestLogger.js"; import { getModelTargetFormat, getModelStrip, getModelUpstreamId, getModelType, PROVIDER_ID_TO_ALIAS } from "../config/providerModels.js"; import { PROVIDERS } from "../config/providers.js"; -import { createErrorResult, parseUpstreamError, formatProviderError } from "../utils/error.js"; +import { cloneUpstreamErrorResponse, createErrorResult, parseUpstreamError, formatProviderError } from "../utils/error.js"; import { HTTP_STATUS, TOKEN_SAVER_HEADER } from "../config/runtimeConfig.js"; import { handleBypassRequest } from "../utils/bypassHandler.js"; import { trackPendingRequest, appendRequestLog, saveRequestDetail } from "@/lib/usageDb.js"; @@ -367,6 +367,9 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred // Provider returned error if (!providerResponse.ok) { trackPendingRequest(model, provider, connectionId, false, true); + const upstreamResponse = provider === "codex" && providerResponse.status === HTTP_STATUS.BAD_REQUEST + ? cloneUpstreamErrorResponse(providerResponse) + : null; const { statusCode, message, resetsAtMs } = await parseUpstreamError(providerResponse, executor); appendRequestLog({ model, provider, connectionId, status: `FAILED ${statusCode}` }).catch(() => { }); saveRequestDetail(buildRequestDetail({ @@ -386,7 +389,8 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred log.errorLine(reqTag, "✗", `ERROR ${statusCode} · ${provider}/${model} · ${Date.now() - requestStartTime}ms${urlStr}\n ${errMsg}`); } reqLogger.logError(new Error(message), finalBody || translatedBody); - return createErrorResult(statusCode, errMsg, resetsAtMs); + const errorResult = createErrorResult(statusCode, errMsg, resetsAtMs); + return upstreamResponse ? { ...errorResult, upstreamResponse } : errorResult; } const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, pxpipe: pxpipeSummary, reqTag, log }; diff --git a/open-sse/services/accountFallback.js b/open-sse/services/accountFallback.js index 8d280da412..f8611bb45a 100644 --- a/open-sse/services/accountFallback.js +++ b/open-sse/services/accountFallback.js @@ -1,4 +1,115 @@ -import { ERROR_RULES, BACKOFF_CONFIG, TRANSIENT_COOLDOWN_MS } from "../config/errorConfig.js"; +import { + ERROR_RULES, + BACKOFF_CONFIG, + TRANSIENT_COOLDOWN_MS, + CODEX_REQUEST_SCHEMA_ERROR_CODES, + CODEX_REQUEST_SCHEMA_MESSAGE_PATTERN, + CODEX_REQUEST_SCHEMA_PARAM_ROOTS, + CODEX_ITEM_ID_PARAM_PATTERN, + CODEX_ITEM_ID_MESSAGE_PATTERN, + REQUEST_SCHEMA_CLASSIFICATION, +} from "../config/errorConfig.js"; + +function parseJsonErrorText(value) { + if (typeof value !== "string") return null; + const text = value.trim().replace(/^\[\d+\]:\s*/, ""); + const candidates = [text]; + const firstBrace = text.indexOf("{"); + const lastBrace = text.lastIndexOf("}"); + if (firstBrace >= 0 && lastBrace > firstBrace) candidates.push(text.slice(firstBrace, lastBrace + 1)); + for (const candidate of candidates) { + try { return JSON.parse(candidate); } catch { /* try the next shape */ } + } + return null; +} + +function hasErrorMetadata(value) { + return Boolean(value?.type || value?.code || value?.param); +} + +function isGenericBadRequestWrapper(value) { + return String(value?.type || "").toLowerCase() === "invalid_request_error" + && String(value?.code || "").toLowerCase() === "bad_request" + && typeof value?.message === "string"; +} + +function normalizeErrorPayload(value, depth = 0) { + if (depth > 6) return { message: "" }; + if (typeof value === "string") { + const parsed = parseJsonErrorText(value); + return parsed ? normalizeErrorPayload(parsed, depth + 1) : { message: value }; + } + if (!value || typeof value !== "object" || Array.isArray(value)) { + return { message: String(value || "") }; + } + if (String(value.type || "").toLowerCase() === "error" + && value.error && typeof value.error === "object" && !Array.isArray(value.error)) { + return normalizeErrorPayload(value.error, depth + 1); + } + if (hasErrorMetadata(value) && !isGenericBadRequestWrapper(value)) return value; + if (isGenericBadRequestWrapper(value)) { + const parsed = parseJsonErrorText(value.message); + return parsed ? normalizeErrorPayload(parsed, depth + 1) : { message: value.message }; + } + if (value.error && typeof value.error === "object" && !Array.isArray(value.error)) { + return normalizeErrorPayload(value.error, depth + 1); + } + if (typeof value.error === "string") return normalizeErrorPayload(value.error, depth + 1); + if (typeof value.message === "string") { + const parsed = parseJsonErrorText(value.message); + if (parsed) return normalizeErrorPayload(parsed, depth + 1); + } + return value; +} + +function getSchemaParamRoot(param, message) { + const direct = String(param || "").match(/^([a-z_]\w*)/i)?.[1]; + if (direct) return direct.toLowerCase(); + const embedded = String(message || "").match( + /\b(?:unknown[_ ]parameter\s*:\s*|unsupported[_ ]value\s+(?:for|at)\s+)["'`]?([a-z_]\w*)/i + )?.[1]; + return embedded?.toLowerCase() || null; +} + +export function isCodexRequestSchemaError(provider, status, errorValue = "") { + if (provider !== "codex" || Number(status) !== 400) return false; + + const error = normalizeErrorPayload(errorValue); + const type = String(error?.type || "").toLowerCase(); + const code = String(error?.code || "").toLowerCase(); + const param = String(error?.param || ""); + const message = String(error?.message || (typeof error?.error === "string" ? error.error : "")); + + if (code === "invalid_prompt" || type === "invalid_prompt") return false; + + const itemIdParam = CODEX_ITEM_ID_PARAM_PATTERN.test(param) + || /input\[\d+\]\.id/i.test(message); + const itemIdMetadata = (!type || type === "invalid_request_error") + && (!code || code === "invalid_value"); + if (itemIdMetadata && itemIdParam && CODEX_ITEM_ID_MESSAGE_PATTERN.test(message)) return true; + + const schemaCode = CODEX_REQUEST_SCHEMA_ERROR_CODES.has(code) + ? code + : (CODEX_REQUEST_SCHEMA_ERROR_CODES.has(type) ? type : null); + const metadataAllowsMessageOnly = !code && (!type || type === "invalid_request_error"); + const schemaField = CODEX_REQUEST_SCHEMA_PARAM_ROOTS.has(getSchemaParamRoot(param, message)); + if (schemaCode === "unknown_parameter" || schemaCode === "unsupported_value") return schemaField; + return metadataAllowsMessageOnly && schemaField && CODEX_REQUEST_SCHEMA_MESSAGE_PATTERN.test(message); +} + +export function classifyProviderError(provider, status, errorText, backoffLevel = 0) { + if (isCodexRequestSchemaError(provider, status, errorText)) { + return { ...REQUEST_SCHEMA_CLASSIFICATION }; + } + const { shouldFallback, cooldownMs, newBackoffLevel } = checkFallbackError(status, errorText, backoffLevel); + return { + category: "provider_error", + accountFallback: shouldFallback, + cooldownMs, + comboScope: "model", + ...(newBackoffLevel === undefined ? {} : { newBackoffLevel }), + }; +} /** * Calculate exponential backoff cooldown for rate limits (429) @@ -118,10 +229,19 @@ export function getModelLockKey(model) { * Reads flat field `modelLock_${model}` (or `modelLock___all` when model=null). */ export function isModelLockActive(connection, model) { - const key = getModelLockKey(model); - const expiry = connection[key] || connection[MODEL_LOCK_ALL]; - if (!expiry) return false; - return new Date(expiry).getTime() > Date.now(); + return Boolean(getModelLockUntil(connection, model)); +} + +/** Return the latest active lock that applies to the requested model. */ +export function getModelLockUntil(connection, model) { + if (!connection) return null; + const now = Date.now(); + const expiries = [connection[getModelLockKey(model)], connection[MODEL_LOCK_ALL]] + .filter(Boolean) + .map(value => new Date(value).getTime()) + .filter(value => Number.isFinite(value) && value > now); + if (expiries.length === 0) return null; + return new Date(Math.max(...expiries)).toISOString(); } /** diff --git a/open-sse/services/combo.js b/open-sse/services/combo.js index 9216ab2fcf..9e18be6264 100644 --- a/open-sse/services/combo.js +++ b/open-sse/services/combo.js @@ -2,7 +2,7 @@ * Shared combo (model combo) handling with fallback support */ -import { checkFallbackError, formatRetryAfter } from "./accountFallback.js"; +import { classifyProviderError, formatRetryAfter } from "./accountFallback.js"; import { unavailableResponse } from "../utils/error.js"; import { getCapabilitiesForModel } from "../providers/capabilities.js"; import { extractTextContent } from "../translator/formats/gemini.js"; @@ -224,9 +224,11 @@ export function getComboModelsFromData(modelStr, combosData) { * @param {string} [options.comboName] - Name of the combo (for round-robin tracking) * @param {string} [options.comboStrategy] - Strategy: "fallback" or "round-robin" * @param {number|string} [options.comboStickyLimit=1] - Requests per combo model before switching + * @param {Set} [options.blockedProviders] - Provider exclusions shared by nested combos + * @param {Map} [options.providerTails] - Per-provider queues shared by nested fusion combos * @returns {Promise} */ -export async function handleComboChat({ body, models, handleSingleModel, log, comboName, comboStrategy, comboStickyLimit = 1, autoSwitch = true }) { +export async function handleComboChat({ body, models, handleSingleModel, log, comboName, comboStrategy, comboStickyLimit = 1, autoSwitch = true, resolveModelProvider = null, blockedProviders = new Set(), providerTails = null }) { // Apply rotation strategy if enabled let rotatedModels = getRotatedModels(models, comboName, comboStrategy, comboStickyLimit); @@ -245,13 +247,37 @@ export async function handleComboChat({ body, models, handleSingleModel, log, co let lastError = null; let earliestRetryAfter = null; let lastStatus = null; + let requestSchemaResponse = null; for (let i = 0; i < rotatedModels.length; i++) { const modelStr = rotatedModels[i]; - log.info("COMBO", `Trying model ${i + 1}/${rotatedModels.length}: ${modelStr}`); try { - const result = await handleSingleModel(body, modelStr); + const provider = resolveModelProvider ? await resolveModelProvider(modelStr) : null; + if (provider && blockedProviders.has(provider)) { + log.info("COMBO", `Skipping model ${modelStr}: provider ${provider} rejected request schema`); + continue; + } + log.info("COMBO", `Trying model ${i + 1}/${rotatedModels.length}: ${modelStr}`); + let blockedBefore; + let result; + if (provider && providerTails) { + const previous = providerTails.get(provider) || Promise.resolve(); + const call = previous.catch(() => {}).then(() => { + if (blockedProviders.has(provider)) return null; + blockedBefore = new Set(blockedProviders); + return handleSingleModel(body, modelStr, blockedProviders, providerTails); + }); + providerTails.set(provider, call); + result = await call; + if (!result) { + log.info("COMBO", `Skipping model ${modelStr}: provider ${provider} rejected request schema`); + continue; + } + } else { + blockedBefore = new Set(blockedProviders); + result = await handleSingleModel(body, modelStr, blockedProviders, providerTails); + } // Success (2xx) - return response if (result.ok) { @@ -262,8 +288,10 @@ export async function handleComboChat({ body, models, handleSingleModel, log, co // Extract error info from response let errorText = result.statusText || ""; let retryAfter = null; + let errorPayload = null; try { const errorBody = await result.clone().json(); + errorPayload = errorBody; errorText = errorBody?.error?.message || errorBody?.error || errorBody?.message || errorText; retryAfter = errorBody?.retryAfter || null; } catch { @@ -280,10 +308,17 @@ export async function handleComboChat({ body, models, handleSingleModel, log, co try { errorText = JSON.stringify(errorText); } catch { errorText = String(errorText); } } - // Check if should fallback to next model - const { shouldFallback, cooldownMs } = checkFallbackError(result.status, errorText); + const nestedSchemaProvider = provider ? null : [...blockedProviders].find((item) => !blockedBefore.has(item)); + const errorProvider = provider || nestedSchemaProvider; + const classification = classifyProviderError(errorProvider, result.status, errorPayload || errorText); + if (classification.comboScope === "provider") { + if (!requestSchemaResponse) requestSchemaResponse = result; + if (errorProvider) blockedProviders.add(errorProvider); + log.warn("COMBO", `Provider ${errorProvider || "unknown"} rejected the request schema`, { status: result.status }); + continue; + } - if (!shouldFallback) { + if (!classification.accountFallback) { log.warn("COMBO", `Model ${modelStr} failed (no fallback)`, { status: result.status }); return result; } @@ -291,10 +326,10 @@ export async function handleComboChat({ body, models, handleSingleModel, log, co // For transient errors (503/502/504), wait for cooldown before falling through // so a briefly-overloaded provider gets a chance to recover rather than being // skipped immediately (fixes: combo falls through on transient 503) - if (cooldownMs && cooldownMs > 0 && cooldownMs <= 5000 && + if (classification.cooldownMs && classification.cooldownMs > 0 && classification.cooldownMs <= 5000 && (result.status === 503 || result.status === 502 || result.status === 504)) { - log.info("COMBO", `Model ${modelStr} transient ${result.status}, waiting ${cooldownMs}ms before next`); - await new Promise(r => setTimeout(r, cooldownMs)); + log.info("COMBO", `Model ${modelStr} transient ${result.status}, waiting ${classification.cooldownMs}ms before next`); + await new Promise(r => setTimeout(r, classification.cooldownMs)); } // Fallback to next model @@ -309,6 +344,8 @@ export async function handleComboChat({ body, models, handleSingleModel, log, co } } + if (requestSchemaResponse) return requestSchemaResponse; + // All models failed // Use 503 (Service Unavailable) rather than 406 (Not Acceptable) — 406 implies // the request itself is invalid, but here the providers are simply unavailable @@ -491,10 +528,24 @@ function collectPanel(calls, { minPanel, stragglerGraceMs, panelHardTimeoutMs }) * @param {string} [options.comboName] - Combo name (logging) * @param {string} [options.judgeModel] - Judge model; falls back to panel[0] * @param {Object} [options.tuning] - Override FUSION_DEFAULTS (minPanel, grace, timeout) + * @param {Function} [options.resolveModelProvider] - Resolve a model string to its provider + * @param {Set} [options.blockedProviders] - Provider exclusions shared by nested combos + * @param {Map} [options.providerTails] - Per-provider queues shared by nested fusion combos * @returns {Promise} */ -export async function handleFusionChat({ body, models, handleSingleModel, log, comboName, judgeModel, tuning }) { +export async function handleFusionChat({ body, models, handleSingleModel, log, comboName, judgeModel, tuning, resolveModelProvider = null, blockedProviders = new Set(), providerTails = new Map() }) { const panel = Array.isArray(models) ? models.filter(Boolean) : []; + const resolveProvider = async (model) => { + if (!resolveModelProvider) return null; + try { return await resolveModelProvider(model); } catch { return null; } + }; + const enqueueProviderCall = (provider, unknownKey, task) => { + const queueKey = provider || unknownKey; + const previous = providerTails.get(queueKey) || Promise.resolve(); + const call = previous.catch(() => {}).then(task); + providerTails.set(queueKey, call); + return call; + }; if (panel.length === 0) { return new Response( JSON.stringify({ error: { message: "Fusion combo has no models" } }), @@ -504,7 +555,16 @@ export async function handleFusionChat({ body, models, handleSingleModel, log, c // A single-model fusion has nothing to fuse — just answer directly. if (panel.length === 1) { - return handleSingleModel(body, panel[0]); + const provider = await resolveProvider(panel[0]); + return enqueueProviderCall(provider, Symbol("unknown-provider"), () => { + if (provider && blockedProviders.has(provider)) { + return new Response( + JSON.stringify({ error: { message: `Provider ${provider} rejected the request schema` } }), + { status: 503, headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": "*" } } + ); + } + return handleSingleModel(body, panel[0], undefined, blockedProviders, providerTails); + }); } const cfg = { ...FUSION_DEFAULTS, ...(tuning || {}) }; @@ -524,8 +584,59 @@ export async function handleFusionChat({ body, models, handleSingleModel, log, c } const t0 = Date.now(); - const calls = panel.map((m) => withTimeout(handleSingleModel(panelBody, m, true), cfg.panelHardTimeoutMs)); + let requestSchemaResponse = null; + let acceptingPanelCalls = true; + const providers = resolveModelProvider + ? await Promise.all(panel.map(resolveProvider)) + : panel.map(() => null); + const unknownProviderQueue = Symbol("unknown-provider"); + const recordSchemaFailure = async (result, provider, blockedBefore) => { + let errorPayload = null; + let errorText = result?.statusText || ""; + try { + errorPayload = await result.clone().json(); + errorText = errorPayload?.error?.message || errorPayload?.error || errorPayload?.message || errorText; + } catch { /* non-JSON provider error */ } + const nestedSchemaProvider = provider ? null : [...blockedProviders].find((item) => !blockedBefore.has(item)); + const errorProvider = provider || nestedSchemaProvider; + const classification = classifyProviderError(errorProvider, result?.status, errorPayload || errorText); + if (classification.comboScope !== "provider") return false; + if (!requestSchemaResponse) requestSchemaResponse = result; + if (errorProvider) blockedProviders.add(errorProvider); + return true; + }; + const calls = panel.map((model, index) => { + const provider = providers[index]; + const unknownKey = resolveModelProvider ? unknownProviderQueue : Symbol(`model-${index}`); + return enqueueProviderCall(provider, unknownKey, async () => { + if (!acceptingPanelCalls) return { __dropped: true }; + if (provider && blockedProviders.has(provider)) return { __schemaSkipped: true }; + + const blockedBefore = new Set(blockedProviders); + const result = await withTimeout( + handleSingleModel(panelBody, model, true, blockedProviders, providerTails), + cfg.panelHardTimeoutMs + ); + if (acceptingPanelCalls && !result?.ok && !result?.__timeout && !result?.__error) { + await recordSchemaFailure(result, provider, blockedBefore); + } + return result; + }); + }); const settled = await collectPanel(calls, { ...cfg, minPanel }); + acceptingPanelCalls = false; + const busyPanelProviders = new Set(); + const busyPanelModels = new Set(); + let unknownProviderQueueBusy = false; + for (let i = 0; i < providers.length; i++) { + if ((i in settled) && !settled[i]?.__timeout) continue; + busyPanelModels.add(panel[i]); + if (providers[i]) busyPanelProviders.add(providers[i]); + else if (resolveModelProvider) unknownProviderQueueBusy = true; + } + const isPanelProviderBusy = (model, provider) => provider + ? busyPanelProviders.has(provider) + : busyPanelModels.has(model) || unknownProviderQueueBusy; log.info("FUSION", `fan-out collected in ${Date.now() - t0}ms`); // 2. Collect successful answers. @@ -534,6 +645,7 @@ export async function handleFusionChat({ body, models, handleSingleModel, log, c const res = settled[i]; const model = panel[i]; if (!res) { log.warn("FUSION", `Panel ${model} dropped (straggler/timeout)`); continue; } + if (res.__dropped || res.__schemaSkipped) { log.warn("FUSION", `Panel ${model} skipped`); continue; } if (res.__timeout) { log.warn("FUSION", `Panel ${model} timed out`); continue; } if (res.__error) { log.warn("FUSION", `Panel ${model} threw`, { error: res.__error?.message || String(res.__error) }); continue; } if (!res.ok) { log.warn("FUSION", `Panel ${model} failed`, { status: res.status }); continue; } @@ -541,7 +653,7 @@ export async function handleFusionChat({ body, models, handleSingleModel, log, c const json = await res.clone().json(); const text = extractPanelText(json); if (text) { - answers.push({ model, text }); + answers.push({ model, text, response: res }); log.info("FUSION", `Panel ${model} ok (${text.length} chars)`); } else { log.warn("FUSION", `Panel ${model} returned empty content`); @@ -553,6 +665,7 @@ export async function handleFusionChat({ body, models, handleSingleModel, log, c // 3. Degrade gracefully when the panel is too thin to fuse. if (answers.length === 0) { + if (requestSchemaResponse) return requestSchemaResponse; log.warn("FUSION", "All panel models failed"); return new Response( JSON.stringify({ error: { message: "All fusion panel models failed" } }), @@ -561,11 +674,55 @@ export async function handleFusionChat({ body, models, handleSingleModel, log, c } if (answers.length === 1) { log.info("FUSION", `Only ${answers[0].model} succeeded — answering directly (no fusion)`); - return handleSingleModel(body, answers[0].model); + const answerProvider = await resolveProvider(answers[0].model); + if ( + (answerProvider && blockedProviders.has(answerProvider)) + || isPanelProviderBusy(answers[0].model, answerProvider) + ) { + return answers[0].response; + } + const answerUnknownKey = resolveModelProvider ? unknownProviderQueue : Symbol("unknown-provider"); + const result = await enqueueProviderCall(answerProvider, answerUnknownKey, () => { + if (answerProvider && blockedProviders.has(answerProvider)) return null; + return handleSingleModel(body, answers[0].model, undefined, blockedProviders, providerTails); + }); + return result || requestSchemaResponse || answers[0].response; } // 4. Judge analyzes + writes one final answer (streams to client if requested). const judgeBody = appendUserTurn(body, buildJudgePrompt(answers)); - log.info("FUSION", `Judging ${answers.length} answers with ${judge}`); - return handleSingleModel(judgeBody, judge); + const judgeCandidates = [judge, ...answers.map(({ model }) => model)]; + const attemptedProviders = new Set(); + const attemptedModels = new Set(); + let judgeAttempted = false; + let busyJudgeCandidate = false; + for (const candidate of judgeCandidates) { + if (attemptedModels.has(candidate)) continue; + attemptedModels.add(candidate); + const provider = await resolveProvider(candidate); + if (provider && ( + blockedProviders.has(provider) + || attemptedProviders.has(provider) + )) continue; + if (isPanelProviderBusy(candidate, provider)) { + busyJudgeCandidate = true; + continue; + } + if (provider) attemptedProviders.add(provider); + + log.info("FUSION", `Judging ${answers.length} answers with ${candidate}`); + judgeAttempted = true; + const blockedBefore = new Set(blockedProviders); + const judgeUnknownKey = resolveModelProvider ? unknownProviderQueue : Symbol("unknown-provider"); + const result = await enqueueProviderCall(provider, judgeUnknownKey, () => { + if (provider && blockedProviders.has(provider)) return null; + return handleSingleModel(judgeBody, candidate, undefined, blockedProviders, providerTails); + }); + if (!result) continue; + if (result.ok) return result; + const isSchemaFailure = await recordSchemaFailure(result, provider, blockedBefore); + if (!isSchemaFailure) return result; + } + if (!judgeAttempted && busyJudgeCandidate) return answers[0].response; + return requestSchemaResponse || answers[0].response; } diff --git a/open-sse/translator/formats/responsesApi.js b/open-sse/translator/formats/responsesApi.js index c41ee470db..6b5e4c67f3 100644 --- a/open-sse/translator/formats/responsesApi.js +++ b/open-sse/translator/formats/responsesApi.js @@ -1,5 +1,32 @@ import { ROLE, OPENAI_BLOCK, RESPONSES_ITEM } from "../schema/index.js"; +const STORED_ITEM_REFERENCE_PATTERN = /^(?:at|msg|amsg|rs|lsh|fc|tsc|fco|ctc|ctco|tso|ws|ig|cmp|resp)_/; +const STATELESS_CALL_ITEM_TYPES = new Set([ + RESPONSES_ITEM.FUNCTION_CALL, + RESPONSES_ITEM.FUNCTION_CALL_OUTPUT, + RESPONSES_ITEM.CUSTOM_TOOL_CALL, + RESPONSES_ITEM.CUSTOM_TOOL_CALL_OUTPUT, +]); + +/** + * Remove stored references and optional call item IDs from a stateless Responses replay. + * call_id remains the correlation key; IDs on every other item type are preserved. + */ +export function normalizeStatelessResponseInput(input) { + if (!Array.isArray(input)) return input; + + return input.flatMap((item) => { + if (typeof item === "string" && STORED_ITEM_REFERENCE_PATTERN.test(item)) return []; + if (!item || typeof item !== "object" || Array.isArray(item)) return [item]; + if (item.type === RESPONSES_ITEM.ITEM_REFERENCE) return []; + if (!STATELESS_CALL_ITEM_TYPES.has(item.type) || !Object.hasOwn(item, "id")) return [item]; + + const normalizedItem = { ...item }; + delete normalizedItem.id; + return [normalizedItem]; + }); +} + /** * Normalize Responses API input to array format. * Accepts string or array, returns array of message items. diff --git a/open-sse/translator/schema/blocks.js b/open-sse/translator/schema/blocks.js index 61c2564593..81cbc0af70 100644 --- a/open-sse/translator/schema/blocks.js +++ b/open-sse/translator/schema/blocks.js @@ -27,6 +27,9 @@ export const RESPONSES_ITEM = { MESSAGE: "message", FUNCTION_CALL: "function_call", FUNCTION_CALL_OUTPUT: "function_call_output", + CUSTOM_TOOL_CALL: "custom_tool_call", + CUSTOM_TOOL_CALL_OUTPUT: "custom_tool_call_output", + ITEM_REFERENCE: "item_reference", REASONING: "reasoning", OUTPUT_TEXT: "output_text", INPUT_TEXT: "input_text", diff --git a/open-sse/utils/error.js b/open-sse/utils/error.js index 315723e303..a144710e69 100644 --- a/open-sse/utils/error.js +++ b/open-sse/utils/error.js @@ -37,6 +37,21 @@ export function errorResponse(statusCode, message) { }); } +export function cloneUpstreamErrorResponse(response) { + const cloned = response.clone(); + const headers = new Headers(response.headers); + headers.delete("content-encoding"); + headers.delete("content-length"); + headers.delete("transfer-encoding"); + headers.delete("set-cookie"); + headers.set("Access-Control-Allow-Origin", "*"); + return new Response(cloned.body, { + status: response.status, + statusText: response.statusText, + headers, + }); +} + /** * Write error to SSE stream (for streaming) * @param {WritableStreamDefaultWriter} writer - Stream writer diff --git a/src/sse/handlers/chat.js b/src/sse/handlers/chat.js index af2914a451..c93e062b17 100644 --- a/src/sse/handlers/chat.js +++ b/src/sse/handlers/chat.js @@ -1,4 +1,5 @@ import "open-sse/index.js"; +import { classifyProviderError } from "open-sse/services/accountFallback.js"; import { getProviderCredentials, @@ -99,18 +100,19 @@ export async function handleChat(request, clientRawRequest = null) { return handleFusionChat({ body, models: comboModels, - handleSingleModel: (b, m, isPanel) => { + handleSingleModel: (b, m, isPanel, blockedProviders, providerTails) => { let cleanRawReq = clientRawRequest; if (isPanel && clientRawRequest) { const { tools, tool_choice, ...cleanBody } = clientRawRequest.body || {}; cleanRawReq = { ...clientRawRequest, body: cleanBody }; } - return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); + return handleSingleModelChat(b, m, cleanRawReq, request, apiKey, blockedProviders, providerTails); }, log, comboName: modelStr, judgeModel: comboStrategies[modelStr]?.judgeModel, tuning: comboStrategies[modelStr]?.fusionTuning, + resolveModelProvider: async (candidate) => (await getModelInfo(candidate)).provider, }); } @@ -119,11 +121,12 @@ export async function handleChat(request, clientRawRequest = null) { return handleComboChat({ body, models: comboModels, - handleSingleModel: (b, m) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey), + handleSingleModel: (b, m, blockedProviders) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey, blockedProviders), log, comboName: modelStr, comboStrategy, - comboStickyLimit + comboStickyLimit, + resolveModelProvider: async (candidate) => (await getModelInfo(candidate)).provider, }); } @@ -134,7 +137,7 @@ export async function handleChat(request, clientRawRequest = null) { /** * Handle single model chat request */ -async function handleSingleModelChat(body, modelStr, clientRawRequest = null, request = null, apiKey = null) { +async function handleSingleModelChat(body, modelStr, clientRawRequest = null, request = null, apiKey = null, blockedProviders = null, providerTails = null) { const modelInfo = await getModelInfo(modelStr); // If provider is null, this might be a combo name - check and handle @@ -152,18 +155,21 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re return handleFusionChat({ body, models: comboModels, - handleSingleModel: (b, m, isPanel) => { + handleSingleModel: (b, m, isPanel, sharedBlockedProviders, sharedProviderTails) => { let cleanRawReq = clientRawRequest; if (isPanel && clientRawRequest) { const { tools, tool_choice, ...cleanBody } = clientRawRequest.body || {}; cleanRawReq = { ...clientRawRequest, body: cleanBody }; } - return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); + return handleSingleModelChat(b, m, cleanRawReq, request, apiKey, sharedBlockedProviders, sharedProviderTails); }, log, comboName: modelStr, judgeModel: comboStrategies[modelStr]?.judgeModel, tuning: comboStrategies[modelStr]?.fusionTuning, + resolveModelProvider: async (candidate) => (await getModelInfo(candidate)).provider, + ...(blockedProviders ? { blockedProviders } : {}), + ...(providerTails ? { providerTails } : {}), }); } @@ -172,11 +178,14 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re return handleComboChat({ body, models: comboModels, - handleSingleModel: (b, m) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey), + handleSingleModel: (b, m, sharedBlockedProviders) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey, sharedBlockedProviders, providerTails), log, comboName: modelStr, comboStrategy, - comboStickyLimit + comboStickyLimit, + resolveModelProvider: async (candidate) => (await getModelInfo(candidate)).provider, + ...(blockedProviders ? { blockedProviders } : {}), + ...(providerTails ? { providerTails } : {}), }); } log.warn("CHAT", "Invalid model format", { model: modelStr }); @@ -201,8 +210,8 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re // All accounts unavailable if (!credentials || credentials.allRateLimited) { if (credentials?.allRateLimited) { - const errorMsg = lastError || credentials.lastError || "Unavailable"; - const status = lastStatus || Number(credentials.lastErrorCode) || HTTP_STATUS.SERVICE_UNAVAILABLE; + const errorMsg = lastError || "Temporarily unavailable"; + const status = lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE; log.warn("CHAT", `[${provider}/${model}] ${errorMsg} (${credentials.retryAfterHuman})`); return unavailableResponse(status, `[${provider}/${model}] ${errorMsg}`, credentials.retryAfter, credentials.retryAfterHuman); } @@ -271,6 +280,13 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re if (result.success) return result.response; + const classification = classifyProviderError(provider, result.status, result.error); + if (classification.category === "request_schema") { + blockedProviders?.add(provider); + log.warn("REQUEST", `Non-retryable Codex request schema error (${result.status})`, { provider }); + return result.upstreamResponse || result.response; + } + // Mark account unavailable (auto-calculates cooldown with exponential backoff, or precise resetsAtMs) const { shouldFallback } = await markAccountUnavailable(credentials.connectionId, result.status, result.error, provider, model, result.resetsAtMs); diff --git a/src/sse/services/auth.js b/src/sse/services/auth.js index 36fd6c4962..95a9f13428 100644 --- a/src/sse/services/auth.js +++ b/src/sse/services/auth.js @@ -1,6 +1,6 @@ import { getProviderConnections, validateApiKey, updateProviderConnection, getSettings, getProxyPools } from "@/lib/localDb"; import { resolveConnectionProxyConfig, pickProxyPoolId } from "@/lib/network/connectionProxy"; -import { formatRetryAfter, checkFallbackError, isModelLockActive, buildModelLockUpdate, getEarliestModelLockUntil } from "open-sse/services/accountFallback.js"; +import { formatRetryAfter, classifyProviderError, isModelLockActive, buildModelLockUpdate, getModelLockUntil } from "open-sse/services/accountFallback.js"; import { MAX_RATE_LIMIT_COOLDOWN_MS } from "open-sse/config/errorConfig.js"; import { resolveProviderId, FREE_PROVIDERS } from "@/shared/constants/providers.js"; import * as log from "../utils/logger.js"; @@ -79,25 +79,21 @@ export async function getProviderCredentials(provider, excludeConnectionIds = nu const excluded = excludeSet.has(c.id); const locked = isModelLockActive(c, model); if (excluded || locked) { - const lockUntil = getEarliestModelLockUntil(c); + const lockUntil = getModelLockUntil(c, model); log.debug("AUTH", ` → ${c.id?.slice(0, 8)} | ${excluded ? "excluded" : ""} ${locked ? `modelLocked(${model}) until ${lockUntil}` : ""}`); } }); if (availableConnections.length === 0) { // Find earliest lock expiry across all connections for retry timing - const lockedConns = connections.filter(c => isModelLockActive(c, model)); - const expiries = lockedConns.map(c => getEarliestModelLockUntil(c)).filter(Boolean); + const expiries = connections.map(c => getModelLockUntil(c, model)).filter(Boolean); const earliest = expiries.sort()[0] || null; if (earliest) { - const earliestConn = lockedConns[0]; - log.warn("AUTH", `${provider} | all ${connections.length} accounts locked for ${model || "all"} (${formatRetryAfter(earliest)}) | lastError=${earliestConn?.lastError?.slice(0, 50)}`); + log.warn("AUTH", `${provider} | all ${connections.length} accounts locked for ${model || "all"} (${formatRetryAfter(earliest)})`); return { allRateLimited: true, retryAfter: earliest, retryAfterHuman: formatRetryAfter(earliest), - lastError: earliestConn?.lastError || null, - lastErrorCode: earliestConn?.errorCode || null }; } log.warn("AUTH", `${provider} | all ${connections.length} accounts unavailable`); @@ -209,6 +205,10 @@ export async function getProviderCredentials(provider, excludeConnectionIds = nu */ export async function markAccountUnavailable(connectionId, status, errorText, provider = null, model = null, resetsAtMs = null) { if (!connectionId || connectionId === "noauth") return { shouldFallback: false, cooldownMs: 0 }; + const preliminaryClassification = classifyProviderError(provider, status, errorText); + if (preliminaryClassification.category === "request_schema") { + return { shouldFallback: false, cooldownMs: 0 }; + } const connections = await getProviderConnections({ provider }); const conn = connections.find(c => c.id === connectionId); const backoffLevel = conn?.backoffLevel || 0; @@ -220,7 +220,12 @@ export async function markAccountUnavailable(connectionId, status, errorText, pr cooldownMs = Math.min(resetsAtMs - Date.now(), MAX_RATE_LIMIT_COOLDOWN_MS); newBackoffLevel = 0; } else { - ({ shouldFallback, cooldownMs, newBackoffLevel } = checkFallbackError(status, errorText, backoffLevel)); + const classification = backoffLevel === 0 + ? preliminaryClassification + : classifyProviderError(provider, status, errorText, backoffLevel); + shouldFallback = classification.accountFallback; + cooldownMs = classification.cooldownMs; + newBackoffLevel = classification.newBackoffLevel; } if (!shouldFallback) return { shouldFallback: false, cooldownMs: 0 }; diff --git a/tests/unit/auth-model-lock-isolation.test.js b/tests/unit/auth-model-lock-isolation.test.js new file mode 100644 index 0000000000..ca5df26e4f --- /dev/null +++ b/tests/unit/auth-model-lock-isolation.test.js @@ -0,0 +1,132 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + getProviderConnections: vi.fn(), + updateProviderConnection: vi.fn(), + getSettings: vi.fn(), +})); + +vi.mock("@/lib/localDb", () => ({ + getProviderConnections: mocks.getProviderConnections, + updateProviderConnection: mocks.updateProviderConnection, + getSettings: mocks.getSettings, + getProxyPools: vi.fn().mockResolvedValue([]), + validateApiKey: vi.fn(), +})); +vi.mock("@/lib/network/connectionProxy", () => ({ + resolveConnectionProxyConfig: vi.fn().mockResolvedValue({ + connectionProxyEnabled: false, + connectionProxyUrl: "", + connectionNoProxy: "", + proxyPoolId: null, + vercelRelayUrl: "", + }), + pickProxyPoolId: vi.fn(), +})); +vi.mock("@/shared/constants/providers.js", () => ({ + resolveProviderId: vi.fn(provider => provider), + FREE_PROVIDERS: {}, +})); +vi.mock("@/sse/utils/logger.js", () => ({ info: vi.fn(), debug: vi.fn(), warn: vi.fn() })); + +import { getProviderCredentials, markAccountUnavailable } from "../../src/sse/services/auth.js"; +import { getModelLockUntil, isModelLockActive } from "../../open-sse/services/accountFallback.js"; + +const NOW = new Date("2026-07-31T03:00:00.000Z"); +const at = seconds => new Date(NOW.getTime() + seconds * 1000).toISOString(); + +function connection(id, fields = {}) { + return { + id, + provider: "codex", + isActive: true, + accessToken: `TOKEN_${id}`, + displayName: id, + providerSpecificData: {}, + ...fields, + }; +} + +describe("model lock isolation", () => { + beforeEach(() => { + vi.useFakeTimers(); + vi.setSystemTime(NOW); + vi.clearAllMocks(); + mocks.getSettings.mockResolvedValue({ fallbackStrategy: "fill-first" }); + mocks.updateProviderConnection.mockResolvedValue(undefined); + }); + + afterEach(() => vi.useRealTimers()); + + it("uses only the requested model and global locks for retry timing", async () => { + const connections = [ + connection("account-a", { + modelLock_gpt: at(60), + modelLock___all: at(120), + modelLock_other: at(10), + lastError: "old item_probe_a", + }), + connection("account-b", { + modelLock_gpt: at(90), + modelLock_other: at(300), + lastError: "old item_probe_b", + }), + ]; + mocks.getProviderConnections.mockResolvedValue(connections); + + expect(getModelLockUntil(connections[0], "gpt")).toBe(at(120)); + expect(getModelLockUntil(connections[1], "gpt")).toBe(at(90)); + + const result = await getProviderCredentials("codex", null, "gpt"); + expect(result).toEqual({ + allRateLimited: true, + retryAfter: at(90), + retryAfterHuman: "reset after 1m 30s", + }); + }); + + it("does not let an expired model lock hide an active global lock", () => { + const account = connection("account-a", { + modelLock_gpt: at(-10), + modelLock___all: at(30), + }); + expect(isModelLockActive(account, "gpt")).toBe(true); + expect(getModelLockUntil(account, "gpt")).toBe(at(30)); + }); + + it("ignores another model's lock", async () => { + mocks.getProviderConnections.mockResolvedValue([ + connection("account-a", { modelLock_other: at(300), lastError: "old item_probe_other" }), + ]); + const result = await getProviderCredentials("codex", null, "gpt"); + expect(result.connectionId).toBe("account-a"); + }); + + it("rejects schema errors before DB reads even with a future reset timestamp", async () => { + const error = `[400]: ${JSON.stringify({ error: { + type: "invalid_request_error", + code: "invalid_value", + param: "input[58].id", + message: "Expected an ID that begins with 'ctc' for input[58].id", + } })}`; + const result = await markAccountUnavailable( + "account-a", 400, error, "codex", "gpt", NOW.getTime() + 30000, + ); + + expect(result).toEqual({ shouldFallback: false, cooldownMs: 0 }); + expect(mocks.getProviderConnections).not.toHaveBeenCalled(); + expect(mocks.updateProviderConnection).not.toHaveBeenCalled(); + }); + + it("keeps rate-limit locking unchanged", async () => { + mocks.getProviderConnections.mockResolvedValue([connection("account-a", { backoffLevel: 0 })]); + const result = await markAccountUnavailable("account-a", 429, "rate limit", "codex", "gpt"); + + expect(result).toEqual({ shouldFallback: true, cooldownMs: 2000 }); + expect(mocks.updateProviderConnection).toHaveBeenCalledWith("account-a", expect.objectContaining({ + modelLock_gpt: at(2), + errorCode: 429, + backoffLevel: 1, + })); + }); +}); diff --git a/tests/unit/chat-codex-schema-isolation.test.js b/tests/unit/chat-codex-schema-isolation.test.js new file mode 100644 index 0000000000..16349d7f9e --- /dev/null +++ b/tests/unit/chat-codex-schema-isolation.test.js @@ -0,0 +1,359 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + getProviderCredentials: vi.fn(), + markAccountUnavailable: vi.fn(), + clearAccountError: vi.fn(), + getSettings: vi.fn(), + getModelInfo: vi.fn(), + getComboModels: vi.fn(), + handleChatCore: vi.fn(), + checkAndRefreshToken: vi.fn(), +})); + +vi.mock("open-sse/index.js", () => ({})); +vi.mock("@/sse/services/auth.js", () => ({ + getProviderCredentials: mocks.getProviderCredentials, + markAccountUnavailable: mocks.markAccountUnavailable, + clearAccountError: mocks.clearAccountError, + extractApiKey: vi.fn(() => null), + isValidApiKey: vi.fn(), +})); +vi.mock("@/lib/localDb", () => ({ getSettings: mocks.getSettings })); +vi.mock("@/sse/services/model.js", () => ({ + getModelInfo: mocks.getModelInfo, + getComboModels: mocks.getComboModels, +})); +vi.mock("open-sse/handlers/chatCore.js", () => ({ handleChatCore: mocks.handleChatCore })); +vi.mock("@/sse/services/tokenRefresh.js", () => ({ + checkAndRefreshToken: mocks.checkAndRefreshToken, + updateProviderCredentials: vi.fn(), +})); +vi.mock("@/lib/headroom/detect", () => ({ DEFAULT_HEADROOM_URL: "http://localhost:8787" })); +vi.mock("@/lib/pxpipe/loader.js", () => ({ getTransform: vi.fn(() => null) })); +vi.mock("@/lib/pxpipe/events.js", () => ({ appendPxpipeEvent: vi.fn() })); +vi.mock("open-sse/utils/bypassHandler.js", () => ({ handleBypassRequest: vi.fn(() => null) })); +vi.mock("open-sse/services/projectId.js", () => ({ getProjectIdForConnection: vi.fn() })); +vi.mock("@/sse/utils/logger.js", () => ({ + info: vi.fn(), debug: vi.fn(), warn: vi.fn(), error: vi.fn(), maskKey: vi.fn(() => "masked"), +})); + +import { handleChat } from "../../src/sse/handlers/chat.js"; + +const PROBE_ID = "item_8e297850f5942c40d91db6c2"; +const SCHEMA_ERROR = `[codex/gpt-5.6-sol] [400]: ${JSON.stringify({ error: { + type: "invalid_request_error", + code: "invalid_value", + param: "input[58].id", + message: `Invalid 'input[58].id': '${PROBE_ID}'. Expected an ID that begins with 'ctc'.`, +} })} (reset after 19s)`; + +function account(id) { + return { + connectionId: id, + connectionName: id, + accessToken: "TOKEN", + providerSpecificData: {}, + _connection: { id }, + }; +} + +function request(body = {}) { + return new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ + model: "cx/gpt-5.6-sol", + input: [{ type: "message", role: "user", content: "hello" }], + ...body, + }), + }); +} + +describe("Codex schema 400 account isolation", () => { + beforeEach(() => { + vi.clearAllMocks(); + mocks.getSettings.mockResolvedValue({ requireApiKey: false }); + mocks.getModelInfo.mockResolvedValue({ provider: "codex", model: "gpt-5.6-sol" }); + mocks.getComboModels.mockResolvedValue(null); + mocks.getProviderCredentials.mockResolvedValue(account("codex-account-1")); + mocks.checkAndRefreshToken.mockImplementation(async (_provider, credentials) => credentials); + mocks.markAccountUnavailable.mockResolvedValue({ shouldFallback: true, cooldownMs: 30000 }); + }); + + it("returns the original 400 before account writes or rotation", async () => { + const originalResponse = new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }); + mocks.handleChatCore.mockResolvedValue({ + success: false, + status: 400, + error: SCHEMA_ERROR, + response: originalResponse, + resetsAtMs: Date.now() + 30000, + }); + + const response = await handleChat(request()); + + expect(response).toBe(originalResponse); + expect(mocks.getProviderCredentials).toHaveBeenCalledOnce(); + expect(mocks.handleChatCore).toHaveBeenCalledOnce(); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + expect(response.headers.get("Retry-After")).toBeNull(); + }); + + it("does not leak a schema error into the next valid request", async () => { + const schemaResponse = new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }); + const successResponse = new Response(JSON.stringify({ id: "resp_valid", output: [] }), { status: 200 }); + mocks.handleChatCore + .mockResolvedValueOnce({ + success: false, + status: 400, + error: SCHEMA_ERROR, + response: new Response("synthetic", { status: 400 }), + upstreamResponse: schemaResponse, + }) + .mockResolvedValueOnce({ success: true, response: successResponse }); + + const firstResponse = await handleChat(request()); + const secondResponse = await handleChat(request({ input: [{ type: "message", role: "user", content: "valid" }] })); + const secondBody = await secondResponse.text(); + + expect(firstResponse).toBe(schemaResponse); + expect(secondResponse).toBe(successResponse); + expect(secondBody).not.toContain(PROBE_ID); + expect(mocks.handleChatCore).toHaveBeenCalledTimes(2); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); + + it("keeps 429 account fallback unchanged", async () => { + const success = new Response("ok", { status: 200 }); + mocks.getProviderCredentials + .mockReset() + .mockResolvedValueOnce(account("codex-account-1")) + .mockResolvedValueOnce(account("codex-account-2")); + mocks.handleChatCore + .mockResolvedValueOnce({ success: false, status: 429, error: "rate limit", response: new Response("rate limit", { status: 429 }) }) + .mockResolvedValueOnce({ success: true, response: success }); + + const response = await handleChat(request()); + + expect(response).toBe(success); + expect(mocks.getProviderCredentials).toHaveBeenCalledTimes(2); + expect(mocks.markAccountUnavailable).toHaveBeenCalledOnce(); + }); + + it("does not retry Codex through a nested combo", async () => { + const schemaResponse = new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }); + mocks.getComboModels.mockImplementation(async (model) => { + if (model === "outer") return ["inner", "cx/gpt-5.6-sol"]; + if (model === "inner") return ["cx/gpt-5.6-sol", "cx/gpt-5.6-codex"]; + return null; + }); + mocks.getModelInfo.mockImplementation(async (model) => { + if (model === "outer" || model === "inner") return { provider: null, model }; + return { provider: "codex", model: model.split("/").at(-1) }; + }); + mocks.handleChatCore.mockResolvedValue({ + success: false, + status: 400, + error: SCHEMA_ERROR, + response: new Response("synthetic", { status: 400 }), + upstreamResponse: schemaResponse, + }); + + const response = await handleChat(request({ model: "outer" })); + + expect(response).toBe(schemaResponse); + expect(mocks.handleChatCore).toHaveBeenCalledOnce(); + expect(mocks.getProviderCredentials).toHaveBeenCalledOnce(); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); + + it("shares provider blocks with a nested fusion combo", async () => { + const schemaResponse = new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }); + mocks.getSettings.mockResolvedValue({ + requireApiKey: false, + comboStrategies: { innerFusion: { fallbackStrategy: "fusion" } }, + }); + mocks.getComboModels.mockImplementation(async (model) => { + if (model === "outer") return ["cx/gpt-5.6-sol", "innerFusion"]; + if (model === "innerFusion") return ["cx/gpt-5.6-codex", "other/model"]; + return null; + }); + mocks.getModelInfo.mockImplementation(async (model) => { + if (model === "outer" || model === "innerFusion") return { provider: null, model }; + const [prefix, resolvedModel] = model.split("/"); + return { provider: prefix === "cx" ? "codex" : "other", model: resolvedModel }; + }); + mocks.handleChatCore.mockImplementation(async ({ modelInfo }) => { + if (modelInfo.provider === "codex") { + return { + success: false, + status: 400, + error: SCHEMA_ERROR, + response: new Response("synthetic", { status: 400 }), + upstreamResponse: schemaResponse, + }; + } + return { + success: true, + response: new Response(JSON.stringify({ + output: [{ type: "message", content: [{ type: "output_text", text: "ok" }] }], + }), { + status: 200, + headers: { "Content-Type": "application/json" }, + }), + }; + }); + + const response = await handleChat(request({ model: "outer" })); + const calledProviders = mocks.handleChatCore.mock.calls.map(([options]) => options.modelInfo.provider); + + expect(response.ok).toBe(true); + expect(calledProviders.filter((provider) => provider === "codex")).toHaveLength(1); + expect(calledProviders.filter((provider) => provider === "other")).toHaveLength(2); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); + + it("skips a blocked provider in a single-model fusion", async () => { + const schemaResponse = new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }); + mocks.getSettings.mockResolvedValue({ + requireApiKey: false, + comboStrategies: { singleFusion: { fallbackStrategy: "fusion" } }, + }); + mocks.getComboModels.mockImplementation(async (model) => { + if (model === "outer") return ["cx/gpt-5.6-sol", "singleFusion"]; + if (model === "singleFusion") return ["cx/gpt-5.6-codex"]; + return null; + }); + mocks.getModelInfo.mockImplementation(async (model) => { + if (model === "outer" || model === "singleFusion") return { provider: null, model }; + return { provider: "codex", model: model.split("/").at(-1) }; + }); + mocks.handleChatCore.mockResolvedValue({ + success: false, + status: 400, + error: SCHEMA_ERROR, + response: new Response("synthetic", { status: 400 }), + upstreamResponse: schemaResponse, + }); + + const response = await handleChat(request({ model: "outer" })); + + expect(response).toBe(schemaResponse); + expect(mocks.handleChatCore).toHaveBeenCalledOnce(); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); + + it("shares provider queues across nested fusion combos", async () => { + mocks.getSettings.mockResolvedValue({ + requireApiKey: false, + comboStrategies: { + outerFusion: { fallbackStrategy: "fusion" }, + innerFusion: { fallbackStrategy: "fusion" }, + }, + }); + mocks.getComboModels.mockImplementation(async (model) => { + if (model === "outerFusion") return ["cx/gpt-5.6-sol", "innerFusion", "other/a"]; + if (model === "innerFusion") return ["cx/gpt-5.6-codex", "other/b"]; + return null; + }); + mocks.getModelInfo.mockImplementation(async (model) => { + if (model === "outerFusion" || model === "innerFusion") return { provider: null, model }; + const [prefix, resolvedModel] = model.split("/"); + return { provider: prefix === "cx" ? "codex" : "other", model: resolvedModel }; + }); + mocks.handleChatCore.mockImplementation(async ({ modelInfo }) => { + if (modelInfo.provider === "codex") { + return { + success: false, + status: 400, + error: SCHEMA_ERROR, + response: new Response("synthetic", { status: 400 }), + upstreamResponse: new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }), + }; + } + return { + success: true, + response: new Response(JSON.stringify({ + output: [{ type: "message", content: [{ type: "output_text", text: `ok-${modelInfo.model}` }] }], + }), { status: 200, headers: { "Content-Type": "application/json" } }), + }; + }); + + const response = await handleChat(request({ model: "outerFusion" })); + const calledProviders = mocks.handleChatCore.mock.calls.map(([options]) => options.modelInfo.provider); + + expect(response.ok).toBe(true); + expect(calledProviders.filter((provider) => provider === "codex")).toHaveLength(1); + expect(calledProviders.filter((provider) => provider === "other").length).toBeGreaterThanOrEqual(2); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); + + it("shares provider queues from fusion into a nested fallback combo", async () => { + let activeCodexCalls = 0; + let maxActiveCodexCalls = 0; + mocks.getSettings.mockResolvedValue({ + requireApiKey: false, + comboStrategies: { outerFusion: { fallbackStrategy: "fusion" } }, + }); + mocks.getComboModels.mockImplementation(async (model) => { + if (model === "outerFusion") return ["cx/gpt-5.6-sol", "innerFallback", "other/a"]; + if (model === "innerFallback") return ["cx/gpt-5.6-codex", "other/b"]; + return null; + }); + mocks.getModelInfo.mockImplementation(async (model) => { + if (model === "outerFusion" || model === "innerFallback") return { provider: null, model }; + const [prefix, resolvedModel] = model.split("/"); + return { provider: prefix === "cx" ? "codex" : "other", model: resolvedModel }; + }); + mocks.handleChatCore.mockImplementation(async ({ modelInfo }) => { + if (modelInfo.provider === "codex") { + activeCodexCalls += 1; + maxActiveCodexCalls = Math.max(maxActiveCodexCalls, activeCodexCalls); + await new Promise((resolve) => setTimeout(resolve, 20)); + activeCodexCalls -= 1; + return { + success: false, + status: 400, + error: SCHEMA_ERROR, + response: new Response("synthetic", { status: 400 }), + upstreamResponse: new Response(JSON.stringify({ error: { message: SCHEMA_ERROR } }), { status: 400 }), + }; + } + return { + success: true, + response: new Response(JSON.stringify({ + output: [{ type: "message", content: [{ type: "output_text", text: `ok-${modelInfo.model}` }] }], + }), { status: 200, headers: { "Content-Type": "application/json" } }), + }; + }); + + const response = await handleChat(request({ model: "outerFusion" })); + const calledProviders = mocks.handleChatCore.mock.calls.map(([options]) => options.modelInfo.provider); + + expect(response.ok).toBe(true); + expect(maxActiveCodexCalls).toBe(1); + expect(calledProviders.filter((provider) => provider === "codex")).toHaveLength(1); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); + + it("does not replay a stored lastError when all accounts were already locked", async () => { + const retryAfter = new Date(Date.now() + 45000).toISOString(); + mocks.getProviderCredentials.mockResolvedValue({ + allRateLimited: true, + retryAfter, + retryAfterHuman: "reset after 45s", + lastError: `old ${PROBE_ID}`, + lastErrorCode: 400, + }); + + const response = await handleChat(request()); + const text = await response.text(); + + expect(response.status).toBe(503); + expect(text).toContain("Temporarily unavailable"); + expect(text).not.toContain(PROBE_ID); + expect(mocks.handleChatCore).not.toHaveBeenCalled(); + expect(mocks.markAccountUnavailable).not.toHaveBeenCalled(); + }); +}); diff --git a/tests/unit/chat-core-upstream-error.test.js b/tests/unit/chat-core-upstream-error.test.js new file mode 100644 index 0000000000..efdee7bde6 --- /dev/null +++ b/tests/unit/chat-core-upstream-error.test.js @@ -0,0 +1,124 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const { executeMock } = vi.hoisted(() => ({ + executeMock: vi.fn(), +})); + +vi.mock("../../open-sse/executors/index.js", () => ({ + getExecutor: () => ({ + noAuth: true, + execute: executeMock, + }), +})); + +vi.mock("../../open-sse/utils/requestLogger.js", () => ({ + createRequestLogger: async () => ({ + logClientRawRequest: vi.fn(), + logRawRequest: vi.fn(), + logTargetRequest: vi.fn(), + logError: vi.fn(), + }), +})); + +vi.mock("@/lib/usageDb.js", () => ({ + trackPendingRequest: vi.fn(), + appendRequestLog: vi.fn(async () => {}), + saveRequestDetail: vi.fn(async () => {}), +})); + +const { handleChatCore } = await import("../../open-sse/handlers/chatCore.js"); + +describe("handleChatCore upstream errors", () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it("preserves the original provider response before parsing its body", async () => { + const upstreamBody = { + error: { + type: "invalid_request_error", + code: "invalid_value", + param: "input[434].id", + message: "Invalid 'input[434].id': 'item_probe'. Expected an ID that begins with 'ctc'.", + }, + }; + executeMock.mockResolvedValue({ + response: new Response(JSON.stringify(upstreamBody), { + status: 400, + statusText: "Bad Request", + headers: { + "content-type": "application/json", + "content-encoding": "gzip", + "content-length": "999", + "transfer-encoding": "chunked", + "set-cookie": "upstream-session=secret", + "x-request-id": "req_schema_probe", + }, + }), + url: "https://chatgpt.com/backend-api/codex/responses", + headers: {}, + transformedBody: null, + }); + + const result = await handleChatCore({ + body: { + model: "gpt-5.6-sol", + stream: false, + input: [{ type: "message", role: "user", content: "hello" }], + }, + modelInfo: { provider: "codex", model: "gpt-5.6-sol" }, + credentials: { accessToken: "TOKEN", providerSpecificData: {} }, + connectionId: "codex-account-1", + sourceFormatOverride: "openai-responses", + clientRawRequest: { + endpoint: "/v1/responses", + body: {}, + headers: { accept: "application/json" }, + }, + log: { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn() }, + }); + + expect(result.success).toBe(false); + expect(result.status).toBe(400); + expect(result.upstreamResponse.status).toBe(400); + expect(result.upstreamResponse.statusText).toBe("Bad Request"); + expect(result.upstreamResponse.headers.get("x-request-id")).toBe("req_schema_probe"); + expect(result.upstreamResponse.headers.get("Access-Control-Allow-Origin")).toBe("*"); + expect(result.upstreamResponse.headers.get("content-encoding")).toBeNull(); + expect(result.upstreamResponse.headers.get("content-length")).toBeNull(); + expect(result.upstreamResponse.headers.get("transfer-encoding")).toBeNull(); + expect(result.upstreamResponse.headers.get("set-cookie")).toBeNull(); + await expect(result.upstreamResponse.json()).resolves.toEqual(upstreamBody); + await expect(result.response.json()).resolves.not.toEqual(upstreamBody); + }); + + it("keeps non-schema error results unchanged", async () => { + executeMock.mockResolvedValue({ + response: new Response(JSON.stringify({ error: { message: "rate limit" } }), { status: 429 }), + url: "https://chatgpt.com/backend-api/codex/responses", + headers: {}, + transformedBody: null, + }); + + const result = await handleChatCore({ + body: { + model: "gpt-5.6-sol", + stream: false, + input: [{ type: "message", role: "user", content: "hello" }], + }, + modelInfo: { provider: "codex", model: "gpt-5.6-sol" }, + credentials: { accessToken: "TOKEN", providerSpecificData: {} }, + connectionId: "codex-account-1", + sourceFormatOverride: "openai-responses", + clientRawRequest: { + endpoint: "/v1/responses", + body: {}, + headers: { accept: "application/json" }, + }, + log: { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn() }, + }); + + expect(result.status).toBe(429); + expect(result).not.toHaveProperty("upstreamResponse"); + }); +}); diff --git a/tests/unit/codex-function-call-item-id.test.js b/tests/unit/codex-function-call-item-id.test.js index a3df61dd0a..b27b29b932 100644 --- a/tests/unit/codex-function-call-item-id.test.js +++ b/tests/unit/codex-function-call-item-id.test.js @@ -11,96 +11,117 @@ function transformInput(input) { }; executor.transformRequest("gpt-5.6-sol", body, true, { - connectionId: "test-codex-function-call-item-id", + connectionId: "test-codex-stateless-item-id", providerSpecificData: {}, }); return body.input; } -describe("CodexExecutor function-call item ids", () => { - it("strips replayed item ids without changing the tool correlation id", () => { - const input = transformInput([ - { - type: "function_call", - id: "item_a80a215e158de93e3e66cd2c", - call_id: "call_shell_1", - name: "shell", - arguments: "{}", - }, - { - type: "function_call_output", - id: "item_output_1", - call_id: "call_shell_1", - output: "done", - }, - ]); +const TARGET_ITEMS = { + function_call: { + legalId: "fc_valid_1", + payload: { call_id: "call_function", name: "shell", arguments: "{\"cmd\":\"pwd\"}", status: "completed" }, + }, + function_call_output: { + legalId: "fco_valid_1", + payload: { call_id: "call_function", output: "done", status: "completed" }, + }, + custom_tool_call: { + legalId: "ctc_valid_1", + payload: { call_id: "call_custom", name: "codex_app", input: "PAYLOAD", status: "completed" }, + }, + custom_tool_call_output: { + legalId: "ctco_valid_1", + payload: { call_id: "call_custom", output: "RESULT", status: "completed" }, + }, +}; - expect(input[0]).toEqual({ - type: "function_call", - call_id: "call_shell_1", - name: "shell", - arguments: "{}", - }); - expect(input[1].id).toBe("item_output_1"); - expect(input[1].call_id).toBe("call_shell_1"); - }); +describe("CodexExecutor stateless item IDs", () => { + it.each(Object.entries(TARGET_ITEMS))("removes every optional %s id and preserves its payload", (type, fixture) => { + const ids = ["item_replayed_1", fixture.legalId, 42, null]; + const source = [ + ...ids.map((id, sequence) => ({ type, id, ...fixture.payload, sequence })), + { type, ...fixture.payload, sequence: ids.length }, + ]; - it("cleans an invalid function-call id at the reported history index", () => { - const history = Array.from({ length: 59 }, (_, index) => ({ - type: "message", - id: `item_message_${index}`, - role: "user", - content: [{ type: "input_text", text: `step ${index}` }], - })); - history.push({ - type: "function_call", - id: "item_6a5f72cd0d444378d96b2841", - call_id: "call_reported_59", - name: "local_tool", - arguments: "{}", + const input = transformInput(source); + + expect(input).toHaveLength(source.length); + input.forEach((item, sequence) => { + expect(item).toEqual({ type, ...fixture.payload, sequence }); }); + }); - const input = transformInput(history); + it("preserves IDs and payloads on non-call items", () => { + const source = [ + { type: "message", id: "msg_history_1", role: "assistant", content: [{ type: "output_text", text: "continue" }] }, + { type: "message", id: "item_message_1", role: "user", content: [{ type: "input_text", text: "again" }] }, + { type: "reasoning", id: "rs_history_1", encrypted_content: "ENCRYPTED_REASONING" }, + { type: "reasoning", id: "item_reasoning_1", summary: [{ type: "summary_text", text: "summary" }] }, + { type: "future_response_item", id: "item_future_1", payload: "PAYLOAD" }, + { id: "item_implicit_message_1", role: "user", content: "hello" }, + ]; - expect(input[59].id).toBeUndefined(); - expect(input[59].call_id).toBe("call_reported_59"); + expect(transformInput(source)).toEqual(source); }); - it("removes every function-call item id when store is disabled", () => { + it("removes stored item references while preserving ordinary input", () => { + const storedReferences = [ + "at_stored", "msg_stored", "amsg_stored", "rs_stored", "lsh_stored", + "fc_stored", "tsc_stored", "fco_stored", "ctc_stored", "ctco_stored", + "tso_stored", "ws_stored", "ig_stored", "cmp_stored", "resp_stored", + ]; + const input = transformInput([ - { - type: "function_call", - id: "fc_valid_1", - call_id: "call_valid_1", - name: "shell", - arguments: "{}", - }, - { - type: "function_call", - id: 42, - call_id: "call_numeric_1", - name: "shell", - arguments: "{}", - }, + ...storedReferences, + { type: "item_reference", id: "item_stored" }, + "ordinary text", + { type: "message", id: "msg_kept", role: "user", content: "continue" }, ]); - expect(input[0].id).toBeUndefined(); - expect(input[0].call_id).toBe("call_valid_1"); - expect(input[1].id).toBeUndefined(); - expect(input[1].call_id).toBe("call_numeric_1"); + expect(input).toEqual([ + "ordinary text", + { type: "message", id: "msg_kept", role: "user", content: "continue" }, + ]); }); - it("does not strip item ids from non-function-call input items", () => { + it("preserves function and custom call/output pairing through call_id", () => { const input = transformInput([ - { - type: "message", - id: "item_message_1", - role: "user", - content: [{ type: "input_text", text: "continue" }], - }, + { type: "function_call", id: "item_fc", call_id: "call_function", name: "shell", arguments: "{}" }, + { type: "function_call_output", id: "item_fco", call_id: "call_function", output: "done" }, + { type: "custom_tool_call", id: "item_ctc", call_id: "call_custom", name: "codex_app", input: "PAYLOAD" }, + { type: "custom_tool_call_output", id: "item_ctco", call_id: "call_custom", output: "RESULT" }, + ]); + + expect(input).toEqual([ + { type: "function_call", call_id: "call_function", name: "shell", arguments: "{}" }, + { type: "function_call_output", call_id: "call_function", output: "done" }, + { type: "custom_tool_call", call_id: "call_custom", name: "codex_app", input: "PAYLOAD" }, + { type: "custom_tool_call_output", call_id: "call_custom", output: "RESULT" }, ]); + }); + + it.each([ + [58, "function_call", "function_call_output", { name: "shell", arguments: "{}" }], + [434, "custom_tool_call", "custom_tool_call_output", { name: "codex_app", input: "PAYLOAD" }], + ])("normalizes a call/output pair after %i replayed history items", (targetIndex, callType, outputType, callPayload) => { + const callId = `call_reported_${targetIndex}`; + const history = Array.from({ length: targetIndex }, (_, index) => ({ + type: "message", + id: `msg_history_${index}`, + role: "user", + content: [{ type: "input_text", text: `step ${index}` }], + })); + history.push( + { type: callType, id: `item_call_${targetIndex}`, call_id: callId, ...callPayload }, + { type: outputType, id: `item_output_${targetIndex}`, call_id: callId, output: "RESULT" }, + ); + + const input = transformInput(history); - expect(input[0].id).toBe("item_message_1"); + expect(input[targetIndex - 1].id).toBe(`msg_history_${targetIndex - 1}`); + expect(input[targetIndex]).toEqual({ type: callType, call_id: callId, ...callPayload }); + expect(input[targetIndex + 1]).toEqual({ type: outputType, call_id: callId, output: "RESULT" }); }); }); diff --git a/tests/unit/codex-schema-fallback.test.js b/tests/unit/codex-schema-fallback.test.js new file mode 100644 index 0000000000..9a60a3f604 --- /dev/null +++ b/tests/unit/codex-schema-fallback.test.js @@ -0,0 +1,124 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +import { classifyProviderError, isCodexRequestSchemaError } from "../../open-sse/services/accountFallback.js"; +import { handleComboChat } from "../../open-sse/services/combo.js"; + +const ITEM_ID_ERROR = { + type: "invalid_request_error", + code: "invalid_value", + param: "input[434].id", + message: "Invalid 'input[434].id': 'item_probe_434'. Expected an ID that begins with 'ctc'.", +}; + +const log = { info: vi.fn(), warn: vi.fn() }; + +function errorResponse(status, error) { + return new Response(JSON.stringify({ error }), { + status, + headers: { "Content-Type": "application/json" }, + }); +} + +function wrappedSchemaResponse() { + return errorResponse(400, { + type: "invalid_request_error", + code: "bad_request", + message: `[400]: ${JSON.stringify({ error: ITEM_ID_ERROR })}`, + }); +} + +describe("Codex request schema classification", () => { + beforeEach(() => vi.clearAllMocks()); + + it.each([ + ["structured item ID", { error: ITEM_ID_ERROR }], + ["raw JSON", JSON.stringify({ error: ITEM_ID_ERROR })], + ["status-wrapped JSON", `[400]: ${JSON.stringify({ error: ITEM_ID_ERROR })}`], + ["outer bad_request wrapper", { error: { type: "invalid_request_error", code: "bad_request", message: `[400]: ${JSON.stringify({ error: ITEM_ID_ERROR })}` } }], + ["message-only item ID", `[400]: ${ITEM_ID_ERROR.message}`], + ["unknown parameter", { error: { type: "invalid_request_error", code: "unknown_parameter", param: "input[2].namespace", message: "Unknown parameter" } }], + ["message-only unknown parameter", "[400]: Unknown parameter: 'input[2].namespace'."], + ["unsupported value", { error: { type: "invalid_request_error", code: "unsupported_value", param: "tool_choice", message: "Unsupported value" } }], + ["message-only unsupported value", "[400]: Unsupported value for 'service_tier': 'BAD'."], + ])("classifies %s", (_name, value) => { + expect(classifyProviderError("codex", 400, value)).toEqual({ + category: "request_schema", + accountFallback: false, + cooldownMs: 0, + comboScope: "provider", + }); + }); + + it.each([ + ["other provider", "openai", 400, { error: ITEM_ID_ERROR }], + ["unauthorized", "codex", 401, { error: ITEM_ID_ERROR }], + ["forbidden", "codex", 403, { error: ITEM_ID_ERROR }], + ["rate limit", "codex", 429, "rate limit"], + ["capacity", "codex", 400, "Selected model is at capacity"], + ["invalid prompt", "codex", 400, { error: { type: "invalid_request_error", code: "invalid_prompt", message: ITEM_ID_ERROR.message } }], + ["unsupported account model", "codex", 400, { error: { type: "invalid_request_error", code: "unsupported_value", param: "model", message: "This model is not supported for the current account" } }], + ["unrelated invalid value", "codex", 400, { error: { type: "invalid_request_error", code: "invalid_value", param: "reasoning.effort", message: "Invalid value" } }], + ])("does not classify %s", (_name, provider, status, value) => { + expect(isCodexRequestSchemaError(provider, status, value)).toBe(false); + }); +}); + +describe("Codex provider-scoped Combo fallback", () => { + beforeEach(() => vi.clearAllMocks()); + + it("skips remaining Codex models and continues another provider", async () => { + const calls = []; + const providers = { + "cx/gpt-5.6-sol": "codex", + "codex/gpt-5.5": "codex", + "openai/gpt-5.5": "openai", + }; + const response = await handleComboChat({ + body: {}, + models: Object.keys(providers), + handleSingleModel: vi.fn(async (_body, model) => { + calls.push(model); + return model.startsWith("openai/") ? new Response("ok", { status: 200 }) : wrappedSchemaResponse(); + }), + resolveModelProvider: async model => providers[model], + log, + autoSwitch: false, + }); + + expect(response.status).toBe(200); + expect(calls).toEqual(["cx/gpt-5.6-sol", "openai/gpt-5.5"]); + }); + + it("returns the first schema 400 after one all-Codex call", async () => { + const first = wrappedSchemaResponse(); + const handleSingleModel = vi.fn().mockResolvedValue(first); + const response = await handleComboChat({ + body: {}, + models: ["cx/gpt-5.6-sol", "codex/gpt-5.5"], + handleSingleModel, + resolveModelProvider: async () => "codex", + log, + autoSwitch: false, + }); + + expect(response).toBe(first); + expect(handleSingleModel).toHaveBeenCalledOnce(); + }); + + it.each([429, 502, 503, 504])("keeps normal fallback for HTTP %s", async status => { + const handleSingleModel = vi.fn() + .mockResolvedValueOnce(errorResponse(status, { message: status === 429 ? "rate limit" : "upstream unavailable" })) + .mockResolvedValueOnce(new Response("ok", { status: 200 })); + const response = await handleComboChat({ + body: {}, + models: ["cx/gpt-5.6-sol", "codex/gpt-5.5"], + handleSingleModel, + resolveModelProvider: async () => "codex", + log, + autoSwitch: false, + }); + + expect(response.status).toBe(200); + expect(handleSingleModel).toHaveBeenCalledTimes(2); + }); +}); diff --git a/tests/unit/combo-fusion.test.js b/tests/unit/combo-fusion.test.js index 52949cb00b..b0a24325fd 100644 --- a/tests/unit/combo-fusion.test.js +++ b/tests/unit/combo-fusion.test.js @@ -17,6 +17,21 @@ function errResponse(status = 500) { return make(); } +function schemaResponse() { + return new Response(JSON.stringify({ error: { + type: "invalid_request_error", + code: "invalid_value", + param: "input[434].id", + message: "Invalid 'input[434].id'. Expected an ID that begins with 'ctc'.", + } }), { status: 400, headers: { "Content-Type": "application/json" } }); +} + +function deferred() { + let resolve; + const promise = new Promise((done) => { resolve = done; }); + return { promise, resolve }; +} + describe("fusion combo", () => { it("answers directly with a single-model panel (nothing to fuse)", async () => { const handleSingleModel = vi.fn(async () => okResponse("solo")); @@ -112,6 +127,217 @@ describe("fusion combo", () => { expect(judgeText).not.toContain("slow"); }); + it("uses a non-busy provider for judging after quorum leaves a provider straggler", async () => { + const slowGate = deferred(); + const slowDone = deferred(); + const activeByProvider = new Map(); + const maxActiveByProvider = new Map(); + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + const provider = model.split("/")[0]; + const active = (activeByProvider.get(provider) || 0) + 1; + activeByProvider.set(provider, active); + maxActiveByProvider.set(provider, Math.max(maxActiveByProvider.get(provider) || 0, active)); + try { + if (model === "p/slow") { + await slowGate.promise; + slowDone.resolve(); + return okResponse("slow"); + } + if (!isPanel) return okResponse("FINAL"); + return okResponse(`fast-${model}`); + } finally { + activeByProvider.set(provider, activeByProvider.get(provider) - 1); + } + }); + + const fusion = handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["p/fast", "q/fast", "p/slow"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + let raced; + try { + raced = await Promise.race([ + fusion.then((response) => ({ response })), + new Promise((resolve) => setTimeout(() => resolve({ timedOut: true }), 250)), + ]); + } finally { + slowGate.resolve(); + } + const response = await fusion; + await slowDone.promise; + + expect(raced.timedOut).not.toBe(true); + expect(response).toBe(raced.response); + expect(response.ok).toBe(true); + expect(handleSingleModel.mock.calls.filter(([, , isPanel]) => isPanel === undefined).map(([, model]) => model)).toEqual(["q/fast"]); + expect(maxActiveByProvider.get("p")).toBe(1); + expect(handleSingleModel).toHaveBeenCalledTimes(4); + }); + + it("keeps a timed-out panel provider busy while its upstream work may still run", async () => { + const detachedGate = deferred(); + const detachedDone = deferred(); + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + if (model === "p/timed-out") { + detachedGate.promise.then(() => detachedDone.resolve()); + return { __timeout: true }; + } + if (!isPanel) return okResponse("FINAL"); + return okResponse(`fast-${model}`); + }); + + let response; + try { + response = await handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["p/fast", "q/fast", "p/timed-out"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + } finally { + detachedGate.resolve(); + } + await detachedDone.promise; + + expect(response.ok).toBe(true); + expect(handleSingleModel.mock.calls.filter(([, , isPanel]) => isPanel === undefined).map(([, model]) => model)).toEqual(["q/fast"]); + }); + + it("does not re-enter a busy unknown provider used by a nested combo", async () => { + const nestedGate = deferred(); + const nestedDone = deferred(); + let activeUnknown = 0; + let maxActiveUnknown = 0; + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + if (model === "inner-combo") { + activeUnknown++; + maxActiveUnknown = Math.max(maxActiveUnknown, activeUnknown); + try { + await nestedGate.promise; + nestedDone.resolve(); + return okResponse("nested answer"); + } finally { + activeUnknown--; + } + } + if (!isPanel) return okResponse("FINAL"); + return okResponse(`fast-${model}`); + }); + + const fusion = handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["inner-combo", "q/fast", "r/fast"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model === "inner-combo" ? null : model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + let raced; + try { + raced = await Promise.race([ + fusion.then((result) => ({ result })), + new Promise((resolve) => setTimeout(() => resolve({ timedOut: true }), 250)), + ]); + } finally { + nestedGate.resolve(); + } + await fusion; + await nestedDone.promise; + + expect(raced.timedOut).not.toBe(true); + expect(raced.result.ok).toBe(true); + expect(maxActiveUnknown).toBe(1); + expect(handleSingleModel.mock.calls.filter(([, , isPanel]) => isPanel === undefined).map(([, model]) => model)).toEqual(["q/fast"]); + }); + + it("returns completed panel output when every judge provider still has a straggler", async () => { + const slowGate = deferred(); + const slowDone = deferred(); + const firstResponse = okResponse("fast-p/first"); + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + if (model === "p/slow") { + await slowGate.promise; + slowDone.resolve(); + return okResponse("slow"); + } + if (model === "codex/schema") return schemaResponse(); + if (!isPanel) throw new Error("busy provider was reused as judge"); + if (model === "p/first") return firstResponse; + return okResponse(`fast-${model}`); + }); + + const fusion = handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["p/first", "p/second", "p/slow", "codex/schema"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + let raced; + try { + raced = await Promise.race([ + fusion.then((response) => ({ response })), + new Promise((resolve) => setTimeout(() => resolve({ timedOut: true }), 250)), + ]); + } finally { + slowGate.resolve(); + } + await fusion; + await slowDone.promise; + + expect(raced.timedOut).not.toBe(true); + expect(raced.response).toBe(firstResponse); + expect(handleSingleModel.mock.calls.filter(([, , isPanel]) => isPanel === undefined)).toHaveLength(0); + expect(handleSingleModel).toHaveBeenCalledTimes(4); + }); + + it("does not retry a lone answer through its busy provider", async () => { + const slowGate = deferred(); + const slowDone = deferred(); + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + if (model === "p/slow") { + await slowGate.promise; + slowDone.resolve(); + return okResponse("slow"); + } + if (model === "q/empty") return okResponse(""); + if (!isPanel) throw new Error("busy provider was retried"); + return okResponse("only answer"); + }); + + const fusion = handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["p/fast", "q/empty", "p/slow"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + let raced; + try { + raced = await Promise.race([ + fusion.then((response) => ({ response })), + new Promise((resolve) => setTimeout(() => resolve({ timedOut: true }), 250)), + ]); + } finally { + slowGate.resolve(); + } + await fusion; + await slowDone.promise; + + expect(raced.timedOut).not.toBe(true); + expect(raced.response.ok).toBe(true); + expect(handleSingleModel.mock.calls.filter(([, , isPanel]) => isPanel === undefined)).toHaveLength(0); + expect(handleSingleModel).toHaveBeenCalledTimes(3); + }); + it("returns the lone survivor directly when only one panel model succeeds", async () => { const handleSingleModel = vi.fn(async (_body, model) => { if (model === "p/ok") return okResponse("lone"); @@ -142,6 +368,98 @@ describe("fusion combo", () => { expect(res.status).toBe(503); }); + it("stops a provider after its first schema failure", async () => { + const firstSchemaResponse = schemaResponse(); + const handleSingleModel = vi.fn(async (_body, model) => { + if (model === "codex/a") return firstSchemaResponse; + return schemaResponse(); + }); + + const response = await handleFusionChat({ + body: { input: [{ role: "user", content: "Q" }] }, + models: ["codex/a", "codex/b", "codex/c"], + handleSingleModel, + log, + resolveModelProvider: async () => "codex", + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + + expect(response).toBe(firstSchemaResponse); + expect(handleSingleModel).toHaveBeenCalledOnce(); + expect(handleSingleModel).toHaveBeenCalledWith( + expect.any(Object), + "codex/a", + true, + expect.any(Set), + expect.any(Map) + ); + }); + + it("continues an unblocked provider after a Codex schema failure", async () => { + const handleSingleModel = vi.fn(async (_body, model) => { + if (model === "codex/a") return schemaResponse(); + if (model === "codex/b") throw new Error("blocked Codex model was called"); + return okResponse("other answer"); + }); + + const response = await handleFusionChat({ + body: { input: [{ role: "user", content: "Q" }] }, + models: ["codex/a", "codex/b", "other/c"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + + expect(response.ok).toBe(true); + expect(handleSingleModel.mock.calls.filter(([, model]) => model === "codex/a")).toHaveLength(1); + expect(handleSingleModel.mock.calls.filter(([, model]) => model === "codex/b")).toHaveLength(0); + expect(handleSingleModel.mock.calls.filter(([, model]) => model === "other/c")).toHaveLength(2); + }); + + it("does not reuse a blocked provider as the judge", async () => { + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + if (model === "codex/b") return schemaResponse(); + if (!isPanel) return okResponse("final answer"); + return okResponse(`answer from ${model}`); + }); + + await handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["codex/a", "codex/b", "other/c"], + handleSingleModel, + log, + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + + const judgeCall = handleSingleModel.mock.calls.find(([, , isPanel]) => isPanel === undefined); + expect(judgeCall[1]).toBe("other/c"); + expect(handleSingleModel.mock.calls.filter(([, model]) => model.startsWith("codex/"))).toHaveLength(2); + }); + + it("tries another provider when the judge returns a schema error", async () => { + const handleSingleModel = vi.fn(async (_body, model, isPanel) => { + if (model === "codex/judge") return schemaResponse(); + if (!isPanel) return okResponse("fallback final answer"); + return okResponse(`answer from ${model}`); + }); + + const response = await handleFusionChat({ + body: { messages: [{ role: "user", content: "Q" }] }, + models: ["other/a", "other/b"], + handleSingleModel, + log, + judgeModel: "codex/judge", + resolveModelProvider: async (model) => model.split("/")[0], + tuning: { minPanel: 2, stragglerGraceMs: 10, panelHardTimeoutMs: 1000 }, + }); + + expect(response.ok).toBe(true); + expect(handleSingleModel.mock.calls.filter(([, model]) => model === "codex/judge")).toHaveLength(1); + expect(handleSingleModel.mock.calls.at(-1)[1]).toBe("other/a"); + }); + it("flattens previous tool history and assistant tool_calls into prose for panel calls", async () => { const handleSingleModel = vi.fn(async () => okResponse("ans")); await handleFusionChat({