From e438a03f96de8eea2110a66b5f48a223cc90626a Mon Sep 17 00:00:00 2001 From: luulam Date: Tue, 4 Aug 2026 23:33:12 +0700 Subject: [PATCH] feat(open-sse): early-peek stream error detection + fix UTF-8 loss in CommandCode peek MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - new open-sse/utils/streamErrorPeek.js: bounded peek of the first bytes of a 200 stream; configured pattern match → 502 so account/combo fallback can run before any byte reaches the client (streaming included). Re-emits RAW bytes so split multi-byte UTF-8 sequences survive the peek (never re-encode decoded text — TextDecoder flush corrupts a lone leading byte to U+FFFD). - chatCore: run the peek after executor.execute when the provider has streamErrorPatterns configured; chat.js passes the settings through. - commandcode executor: same raw-bytes fix in peekForUpstreamError + regression test that fails against the old flush-based re-encode. --- open-sse/executors/commandcode.js | 17 +- open-sse/handlers/chatCore.js | 1046 +++++++++++++++-------- open-sse/utils/streamErrorPeek.js | 136 +++ src/sse/handlers/chat.js | 564 +++++++----- tests/unit/commandcode-executor.test.js | 22 + tests/unit/stream-error-peek.test.js | 99 +++ 6 files changed, 1301 insertions(+), 583 deletions(-) create mode 100644 open-sse/utils/streamErrorPeek.js create mode 100644 tests/unit/stream-error-peek.test.js diff --git a/open-sse/executors/commandcode.js b/open-sse/executors/commandcode.js index e047ad57..36a97a3d 100644 --- a/open-sse/executors/commandcode.js +++ b/open-sse/executors/commandcode.js @@ -27,7 +27,7 @@ export class CommandCodeExecutor extends BaseExecutor { super("commandcode", PROVIDERS.commandcode); } - transformRequest(model, body, stream, credentials) { + transformRequest(_model, body, _stream, _credentials) { body.stream = true; return body; } @@ -131,13 +131,17 @@ export async function peekForUpstreamError( ) { const reader = originalResponse.body.getReader(); const decoder = new TextDecoder(); - const encoder = new TextEncoder(); const abortController = new AbortController(); const forwardAbort = () => abortController.abort(signal?.reason); if (signal?.aborted) abortController.abort(signal?.reason); else if (signal) signal.addEventListener("abort", forwardAbort, { once: true }); + // Raw bytes for lossless re-emission; decoded text is only used for line + // parsing / error detection. Never re-encode decoded text: TextDecoder + // holds a split multi-byte char internally and flush() would replace it + // with U+FFFD, corrupting the stream. + const rawChunks = []; let peeked = ""; let errorEvent = null; let committed = false; @@ -167,6 +171,7 @@ export async function peekForUpstreamError( Math.max(deadline - Date.now(), 1), ); if (done) break; + rawChunks.push(value); peeked += decoder.decode(value, { stream: true }); const lines = peeked.split("\n"); // The last segment may be a partial line — only parse complete ones. @@ -188,10 +193,6 @@ export async function peekForUpstreamError( // the downstream stream pipeline (stall detection, abort handling) takes over. } - // Flush any partial multi-byte UTF-8 sequence held by the decoder so the - // re-encoded peeked bytes round-trip losslessly. - peeked += decoder.decode(); - if (signal) signal.removeEventListener("abort", forwardAbort); if (errorEvent) { @@ -210,7 +211,9 @@ export async function peekForUpstreamError( start(controller) { (async () => { try { - if (peeked) controller.enqueue(encoder.encode(peeked)); + // Re-emit RAW bytes (never re-encoded decoded text) so split + // multi-byte UTF-8 sequences survive the peek untouched. + for (const c of rawChunks) controller.enqueue(c); while (true) { const { done, value } = await reader.read(); if (done) break; diff --git a/open-sse/handlers/chatCore.js b/open-sse/handlers/chatCore.js index 4f91e020..38cf1490 100644 --- a/open-sse/handlers/chatCore.js +++ b/open-sse/handlers/chatCore.js @@ -1,4 +1,8 @@ -import { detectFormat, getTargetFormat, resolveTransport } from "../services/provider.js"; +import { + detectFormat, + getTargetFormat, + resolveTransport, +} from "../services/provider.js"; import { translateRequest } from "../translator/index.js"; import { stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js"; import { FORMATS } from "../translator/formats.js"; @@ -6,24 +10,53 @@ import { normalizeClaudePassthrough } from "../translator/formats/claude.js"; import { createStreamController } from "../utils/streamHandler.js"; 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 { + 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 { + 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"; +import { + trackPendingRequest, + appendRequestLog, + saveRequestDetail, +} from "@/lib/usageDb.js"; import { getExecutor } from "../executors/index.js"; import { supportsGrokCliReasoningEffort } from "../config/grokCli.js"; -import { buildRequestDetail, extractRequestConfig } from "./chatCore/requestDetail.js"; +import { + buildRequestDetail, + extractRequestConfig, +} from "./chatCore/requestDetail.js"; import { handleForcedSSEToJson } from "./chatCore/sseToJsonHandler.js"; import { handleNonStreamingResponse } from "./chatCore/nonStreamingHandler.js"; -import { handleStreamingResponse, buildOnStreamComplete } from "./chatCore/streamingHandler.js"; -import { detectClientTool, isNativePassthrough } from "../utils/clientDetector.js"; +import { + handleStreamingResponse, + buildOnStreamComplete, +} from "./chatCore/streamingHandler.js"; +import { maybeRejectEarlyStreamError } from "../utils/streamErrorPeek.js"; +import { + detectClientTool, + isNativePassthrough, +} from "../utils/clientDetector.js"; import { dedupeTools } from "../utils/toolDeduper.js"; import { injectCaveman } from "../rtk/caveman.js"; import { injectPonytail } from "../rtk/ponytail.js"; import { compressMessages, formatRtkLog } from "../rtk/index.js"; -import { compressWithHeadroom, formatHeadroomLog, formatHeadroomSizeLog, isHeadroomPhantomSavings } from "../rtk/headroom.js"; +import { + compressWithHeadroom, + formatHeadroomLog, + formatHeadroomSizeLog, + isHeadroomPhantomSavings, +} from "../rtk/headroom.js"; import { compressWithPxpipe } from "../rtk/pxpipe.js"; import { getCapabilitiesForModel } from "../providers/capabilities.js"; import { stripUnsupportedModalities } from "../translator/concerns/modality.js"; @@ -38,380 +71,713 @@ import { resolveSessionId } from "../utils/sessionManager.js"; * @param {object} options.credentials - Provider credentials * @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses") */ -export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, headroomEnabled, headroomUrl, headroomCompressUserMessages, cavemanEnabled, cavemanLevel, ponytailEnabled, ponytailLevel, pxpipeEnabled, pxpipeMinChars, pxpipeTimeoutMs, pxpipeTransform, onPxpipeEvent, sourceFormatOverride, providerThinking }) { - const { provider, model } = modelInfo; - const requestStartTime = Date.now(); - // Stable per-session color so all lines of one CLI conversation share a tag - const sessionSeed = (() => { - try { - return resolveSessionId({ headers: clientRawRequest?.headers, body, connectionId, scope: provider }); - } catch { - return connectionId || ""; - } - })(); - const reqTag = log?.tagForSession ? log.tagForSession(sessionSeed) : (log?.nextTag ? log.nextTag() : ""); +export async function handleChatCore({ + body, + modelInfo, + credentials, + log, + onCredentialsRefreshed, + onRequestSuccess, + onDisconnect, + clientRawRequest, + connectionId, + userAgent, + apiKey, + ccFilterNaming, + rtkEnabled, + headroomEnabled, + headroomUrl, + headroomCompressUserMessages, + cavemanEnabled, + cavemanLevel, + ponytailEnabled, + ponytailLevel, + pxpipeEnabled, + pxpipeMinChars, + pxpipeTimeoutMs, + pxpipeTransform, + onPxpipeEvent, + sourceFormatOverride, + providerThinking, + streamErrorPatterns, +}) { + const { provider, model } = modelInfo; + const requestStartTime = Date.now(); + // Stable per-session color so all lines of one CLI conversation share a tag + const sessionSeed = (() => { + try { + return resolveSessionId({ + headers: clientRawRequest?.headers, + body, + connectionId, + scope: provider, + }); + } catch { + return connectionId || ""; + } + })(); + const reqTag = log?.tagForSession + ? log.tagForSession(sessionSeed) + : log?.nextTag + ? log.nextTag() + : ""; - const sourceFormat = sourceFormatOverride || detectFormat(body); + const sourceFormat = sourceFormatOverride || detectFormat(body); - // Check for bypass patterns (warmup, skip, cc naming) - const bypassResponse = handleBypassRequest(body, model, userAgent, ccFilterNaming); - if (bypassResponse) return bypassResponse; + // Check for bypass patterns (warmup, skip, cc naming) + const bypassResponse = handleBypassRequest( + body, + model, + userAgent, + ccFilterNaming, + ); + if (bypassResponse) return bypassResponse; - const alias = PROVIDER_ID_TO_ALIAS[provider] || provider; - const modelTargetFormat = getModelTargetFormat(alias, model); - // Multi-endpoint providers: pick transport matching sourceFormat → zero translation - const runtimeTransport = resolveTransport(provider, sourceFormat); - const targetFormat = modelTargetFormat || runtimeTransport?.format || getTargetFormat(provider); - if (runtimeTransport && credentials) credentials.runtimeTransport = runtimeTransport; - const stripList = getModelStrip(alias, model); - const upstreamModel = getModelUpstreamId(alias, model); + const alias = PROVIDER_ID_TO_ALIAS[provider] || provider; + const modelTargetFormat = getModelTargetFormat(alias, model); + // Multi-endpoint providers: pick transport matching sourceFormat → zero translation + const runtimeTransport = resolveTransport(provider, sourceFormat); + const targetFormat = + modelTargetFormat || runtimeTransport?.format || getTargetFormat(provider); + if (runtimeTransport && credentials) + credentials.runtimeTransport = runtimeTransport; + const stripList = getModelStrip(alias, model); + const upstreamModel = getModelUpstreamId(alias, model); - // Inject provider-level thinking config override (only if client hasn't set) - // on/off → extended type (body.thinking), none/low/medium/high → effort type (body.reasoning_effort) - if (providerThinking?.mode && providerThinking.mode !== "auto") { - const mode = providerThinking.mode; - if (mode === "on" && !body.thinking) { - console.log("Injecting provider-level thinking config override: on"); - body = { ...body, thinking: { type: "enabled", budget_tokens: 10000 } }; - } else if (mode === "off" && !body.thinking) { - body = { ...body, thinking: { type: "disabled" } }; - } else if (!body.reasoning_effort) { - body = { ...body, reasoning_effort: mode }; - } - } + // Inject provider-level thinking config override (only if client hasn't set) + // on/off → extended type (body.thinking), none/low/medium/high → effort type (body.reasoning_effort) + if (providerThinking?.mode && providerThinking.mode !== "auto") { + const mode = providerThinking.mode; + if (mode === "on" && !body.thinking) { + console.log("Injecting provider-level thinking config override: on"); + body = { ...body, thinking: { type: "enabled", budget_tokens: 10000 } }; + } else if (mode === "off" && !body.thinking) { + body = { ...body, thinking: { type: "disabled" } }; + } else if (!body.reasoning_effort) { + body = { ...body, reasoning_effort: mode }; + } + } - const clientRequestedStreaming = body.stream === true || sourceFormat === FORMATS.ANTIGRAVITY || sourceFormat === FORMATS.GEMINI || sourceFormat === FORMATS.GEMINI_CLI; - const providerRequiresStreaming = PROVIDERS[provider]?.forceStream === true; - let stream = providerRequiresStreaming ? true : (body.stream !== false); + const clientRequestedStreaming = + body.stream === true || + sourceFormat === FORMATS.ANTIGRAVITY || + sourceFormat === FORMATS.GEMINI || + sourceFormat === FORMATS.GEMINI_CLI; + const providerRequiresStreaming = PROVIDERS[provider]?.forceStream === true; + let stream = providerRequiresStreaming ? true : body.stream !== false; - // Image generation models require non-streaming (Google v1internal:generateContent) - const modelType = getModelType(alias, model); - const isImageGenModel = modelType === "imageGen" || /image|imagen|image-generation/i.test(model); - if (isImageGenModel && (provider === "antigravity" || provider === "gemini-cli")) { - stream = false; - } + // Image generation models require non-streaming (Google v1internal:generateContent) + const modelType = getModelType(alias, model); + const isImageGenModel = + modelType === "imageGen" || /image|imagen|image-generation/i.test(model); + if ( + isImageGenModel && + (provider === "antigravity" || provider === "gemini-cli") + ) { + stream = false; + } - // DeepSeek-TUI: interactive TUI panel sends stream:true and needs SSE. - // Non-interactive mode (-p flag) sends without stream and can't parse SSE. - // Only force non-streaming when client didn't explicitly request it. - const detectedTool = detectClientTool(clientRawRequest?.headers || {}, body); - if (detectedTool === "deepseek-tui" && body.stream !== true) stream = false; + // DeepSeek-TUI: interactive TUI panel sends stream:true and needs SSE. + // Non-interactive mode (-p flag) sends without stream and can't parse SSE. + // Only force non-streaming when client didn't explicitly request it. + const detectedTool = detectClientTool(clientRawRequest?.headers || {}, body); + if (detectedTool === "deepseek-tui" && body.stream !== true) stream = false; - // Check client Accept header preference for non-streaming requests - // This fixes AI SDK compatibility where clients send Accept: application/json - const acceptHeader = clientRawRequest?.headers?.accept || ""; - const clientPrefersJson = acceptHeader.includes("application/json"); - const clientPrefersSSE = acceptHeader.includes("text/event-stream"); - if (clientPrefersJson && !clientPrefersSSE && body.stream !== true && !providerRequiresStreaming) { - stream = false; - } + // Check client Accept header preference for non-streaming requests + // This fixes AI SDK compatibility where clients send Accept: application/json + const acceptHeader = clientRawRequest?.headers?.accept || ""; + const clientPrefersJson = acceptHeader.includes("application/json"); + const clientPrefersSSE = acceptHeader.includes("text/event-stream"); + if ( + clientPrefersJson && + !clientPrefersSSE && + body.stream !== true && + !providerRequiresStreaming + ) { + stream = false; + } - const reqLogger = await createRequestLogger(sourceFormat, targetFormat, model); - if (clientRawRequest) reqLogger.logClientRawRequest(clientRawRequest.endpoint, clientRawRequest.body, clientRawRequest.headers); - reqLogger.logRawRequest(body); - log?.debug?.("FORMAT", `${sourceFormat} → ${targetFormat} | stream=${stream}`); + const reqLogger = await createRequestLogger( + sourceFormat, + targetFormat, + model, + ); + if (clientRawRequest) + reqLogger.logClientRawRequest( + clientRawRequest.endpoint, + clientRawRequest.body, + clientRawRequest.headers, + ); + reqLogger.logRawRequest(body); + log?.debug?.( + "FORMAT", + `${sourceFormat} → ${targetFormat} | stream=${stream}`, + ); - // Native passthrough: CLI tool and provider are the same ecosystem - // Skip all translation/normalization — only model and Bearer are swapped - const clientTool = detectClientTool(clientRawRequest?.headers || {}, body); - const passthrough = isNativePassthrough(clientTool, provider); + // Native passthrough: CLI tool and provider are the same ecosystem + // Skip all translation/normalization — only model and Bearer are swapped + const clientTool = detectClientTool(clientRawRequest?.headers || {}, body); + const passthrough = isNativePassthrough(clientTool, provider); - // Expose raw client headers to translators/executors for session-id resolution - if (credentials) credentials.rawHeaders = clientRawRequest?.headers || {}; + // Expose raw client headers to translators/executors for session-id resolution + if (credentials) credentials.rawHeaders = clientRawRequest?.headers || {}; - // Auto-strip media blocks the model can't read (vision/audio/pdf) before translation. - if (!passthrough) { - const caps = getCapabilitiesForModel(provider, model); - if (stripUnsupportedModalities(body, sourceFormat, caps)) { - log?.debug?.("MODALITY", `stripped unsupported media for ${provider}/${model}`); - } - // Convert remote image URLs to base64 for targets that can't fetch URLs. - try { - const n = await prefetchRemoteImages(body, sourceFormat, targetFormat, { signal: undefined }); - if (n > 0) log?.debug?.("MODALITY", `prefetched ${n} remote image(s) for ${targetFormat}`); - } catch (e) { log?.warn?.("MODALITY", `image prefetch failed: ${e.message}`); } - } + // Auto-strip media blocks the model can't read (vision/audio/pdf) before translation. + if (!passthrough) { + const caps = getCapabilitiesForModel(provider, model); + if (stripUnsupportedModalities(body, sourceFormat, caps)) { + log?.debug?.( + "MODALITY", + `stripped unsupported media for ${provider}/${model}`, + ); + } + // Convert remote image URLs to base64 for targets that can't fetch URLs. + try { + const n = await prefetchRemoteImages(body, sourceFormat, targetFormat, { + signal: undefined, + }); + if (n > 0) + log?.debug?.( + "MODALITY", + `prefetched ${n} remote image(s) for ${targetFormat}`, + ); + } catch (e) { + log?.warn?.("MODALITY", `image prefetch failed: ${e.message}`); + } + } - let translatedBody; - let toolNameMap; - if (passthrough) { - log?.debug?.("PASSTHROUGH", `${clientTool} → ${provider} | native lossless`); - translatedBody = { ...body, model: stripThinkingSuffix(upstreamModel) }; - // Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects - if (clientTool === "claude") normalizeClaudePassthrough(translatedBody, translatedBody.model); - } else { - translatedBody = translateRequest(sourceFormat, targetFormat, upstreamModel, body, stream, credentials, provider, reqLogger, stripList, connectionId, clientTool); - if (!translatedBody) { - trackPendingRequest(model, provider, connectionId, false, true); - return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Failed to translate request for ${sourceFormat} → ${targetFormat}`); - } - toolNameMap = translatedBody._toolNameMap; - delete translatedBody._toolNameMap; - translatedBody.model = stripThinkingSuffix(upstreamModel); - } + let translatedBody; + let toolNameMap; + if (passthrough) { + log?.debug?.( + "PASSTHROUGH", + `${clientTool} → ${provider} | native lossless`, + ); + translatedBody = { ...body, model: stripThinkingSuffix(upstreamModel) }; + // Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects + if (clientTool === "claude") + normalizeClaudePassthrough(translatedBody, translatedBody.model); + } else { + translatedBody = translateRequest( + sourceFormat, + targetFormat, + upstreamModel, + body, + stream, + credentials, + provider, + reqLogger, + stripList, + connectionId, + clientTool, + ); + if (!translatedBody) { + trackPendingRequest(model, provider, connectionId, false, true); + return createErrorResult( + HTTP_STATUS.BAD_REQUEST, + `Failed to translate request for ${sourceFormat} → ${targetFormat}`, + ); + } + toolNameMap = translatedBody._toolNameMap; + delete translatedBody._toolNameMap; + translatedBody.model = stripThinkingSuffix(upstreamModel); + } - // Dedupe duplicate built-in tools when equivalent MCP tools are present (Claude clients only). - if (clientTool === "claude" && Array.isArray(translatedBody.tools)) { - const { tools: deduped, stripped } = dedupeTools(translatedBody.tools); - if (stripped.length > 0) { - translatedBody.tools = deduped; - log?.debug?.("TOOLDEDUP", `stripped ${stripped.length}: ${stripped.slice(0, 3).join(", ")}${stripped.length > 3 ? "..." : ""}`); - } - } + // Dedupe duplicate built-in tools when equivalent MCP tools are present (Claude clients only). + if (clientTool === "claude" && Array.isArray(translatedBody.tools)) { + const { tools: deduped, stripped } = dedupeTools(translatedBody.tools); + if (stripped.length > 0) { + translatedBody.tools = deduped; + log?.debug?.( + "TOOLDEDUP", + `stripped ${stripped.length}: ${stripped.slice(0, 3).join(", ")}${stripped.length > 3 ? "..." : ""}`, + ); + } + } - // Token savers: applied at the final body just before dispatch - // Covers both passthrough (source shape) and translated (target shape) flows - const finalFormat = passthrough ? sourceFormat : targetFormat; + // Token savers: applied at the final body just before dispatch + // Covers both passthrough (source shape) and translated (target shape) flows + const finalFormat = passthrough ? sourceFormat : targetFormat; - // Request line: one correlated summary (fmt + thinking + counts + account) - if (log?.line) { - const clientModel = clientRawRequest?.body?.model || `${provider}/${model}`; - const msgN = translatedBody.messages?.length || translatedBody.input?.length || translatedBody.contents?.length || body.messages?.length || body.input?.length || 0; - const toolN = translatedBody.tools?.length || body.tools?.length || 0; - const fmtStr = passthrough ? `FMT: ${sourceFormat} (passthrough)` : `FMT: ${sourceFormat}→${targetFormat}`; - const showThinking = provider !== "grok-cli" || supportsGrokCliReasoningEffort(model); - const think = showThinking ? log.fmtThink?.(extractThinking(translatedBody)) : null; - const acc = credentials?.connectionName || credentials?.connectionId?.slice(0, 8) || "-"; - const parts = [ - `POST ${clientModel} → ${provider}/${model}`, - fmtStr, - stream ? "STREAM" : "JSON", - `${msgN} MSG`, - ]; - if (toolN) parts.push(`${toolN} TOOL`); - if (think) parts.push(`THINK:${think}`); - parts.push(`ACC:${acc}`); - log.line(reqTag, "▶", parts.join(" · ")); - } + // Request line: one correlated summary (fmt + thinking + counts + account) + if (log?.line) { + const clientModel = clientRawRequest?.body?.model || `${provider}/${model}`; + const msgN = + translatedBody.messages?.length || + translatedBody.input?.length || + translatedBody.contents?.length || + body.messages?.length || + body.input?.length || + 0; + const toolN = translatedBody.tools?.length || body.tools?.length || 0; + const fmtStr = passthrough + ? `FMT: ${sourceFormat} (passthrough)` + : `FMT: ${sourceFormat}→${targetFormat}`; + const showThinking = + provider !== "grok-cli" || supportsGrokCliReasoningEffort(model); + const think = showThinking + ? log.fmtThink?.(extractThinking(translatedBody)) + : null; + const acc = + credentials?.connectionName || + credentials?.connectionId?.slice(0, 8) || + "-"; + const parts = [ + `POST ${clientModel} → ${provider}/${model}`, + fmtStr, + stream ? "STREAM" : "JSON", + `${msgN} MSG`, + ]; + if (toolN) parts.push(`${toolN} TOOL`); + if (think) parts.push(`THINK:${think}`); + parts.push(`ACC:${acc}`); + log.line(reqTag, "▶", parts.join(" · ")); + } - // TTS models don't support tool messages/function calling - if (getModelType(alias, model) === "tts" && translatedBody.messages) { - translatedBody.messages = translatedBody.messages.filter(msg => msg.role !== "tool"); - delete translatedBody.tools; - } + // TTS models don't support tool messages/function calling + if (getModelType(alias, model) === "tts" && translatedBody.messages) { + translatedBody.messages = translatedBody.messages.filter( + (msg) => msg.role !== "tool", + ); + delete translatedBody.tools; + } - // Per-request opt-out: client can bypass all token savers via header - const tokenSaverEnabled = clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off"; + // Per-request opt-out: client can bypass all token savers via header + const tokenSaverEnabled = + clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off"; - // RTK: compress tool_result content - const rtkStats = compressMessages(translatedBody, tokenSaverEnabled && rtkEnabled); - const rtkLine = formatRtkLog(rtkStats); - if (rtkLine) console.log(rtkLine); + // RTK: compress tool_result content + const rtkStats = compressMessages( + translatedBody, + tokenSaverEnabled && rtkEnabled, + ); + const rtkLine = formatRtkLog(rtkStats); + if (rtkLine) console.log(rtkLine); - // Headroom: optional external proxy compression; fail open if proxy is absent. - const headroomDiagnostics = {}; - const headroomStats = await compressWithHeadroom(translatedBody, { enabled: tokenSaverEnabled && headroomEnabled, url: headroomUrl, model: upstreamModel, format: finalFormat, compressUserMessages: headroomCompressUserMessages, diagnostics: headroomDiagnostics }); - const headroomLine = formatHeadroomLog(headroomStats); - const headroomSizeLine = formatHeadroomSizeLog(headroomDiagnostics); - if (headroomLine) { - log?.info?.("HEADROOM", `${headroomLine}${headroomSizeLine ? ` | ${headroomSizeLine}` : ""}`); - if (isHeadroomPhantomSavings(headroomStats, headroomDiagnostics)) { - log?.warn?.("HEADROOM", `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${formatHeadroomSizeLog(headroomDiagnostics)}`); - } - } else if (tokenSaverEnabled && headroomEnabled) log?.warn?.("HEADROOM", `skipped: ${headroomDiagnostics.reason || "compression unavailable"}${headroomDiagnostics.endpoint ? ` (${headroomDiagnostics.endpoint})` : ""}`); + // Headroom: optional external proxy compression; fail open if proxy is absent. + const headroomDiagnostics = {}; + const headroomStats = await compressWithHeadroom(translatedBody, { + enabled: tokenSaverEnabled && headroomEnabled, + url: headroomUrl, + model: upstreamModel, + format: finalFormat, + compressUserMessages: headroomCompressUserMessages, + diagnostics: headroomDiagnostics, + }); + const headroomLine = formatHeadroomLog(headroomStats); + const headroomSizeLine = formatHeadroomSizeLog(headroomDiagnostics); + if (headroomLine) { + log?.info?.( + "HEADROOM", + `${headroomLine}${headroomSizeLine ? ` | ${headroomSizeLine}` : ""}`, + ); + if (isHeadroomPhantomSavings(headroomStats, headroomDiagnostics)) { + log?.warn?.( + "HEADROOM", + `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${formatHeadroomSizeLog(headroomDiagnostics)}`, + ); + } + } else if (tokenSaverEnabled && headroomEnabled) + log?.warn?.( + "HEADROOM", + `skipped: ${headroomDiagnostics.reason || "compression unavailable"}${headroomDiagnostics.endpoint ? ` (${headroomDiagnostics.endpoint})` : ""}`, + ); - // Token-saver flags accumulator for the single "⚙" log line below. - const xf = []; + // Token-saver flags accumulator for the single "⚙" log line below. + const xf = []; - // Caveman: inject terse-style system prompt - if (tokenSaverEnabled && cavemanEnabled && cavemanLevel) { - injectCaveman(translatedBody, finalFormat, cavemanLevel); - xf.push(`CAVEMAN:${cavemanLevel}`); - } + // Caveman: inject terse-style system prompt + if (tokenSaverEnabled && cavemanEnabled && cavemanLevel) { + injectCaveman(translatedBody, finalFormat, cavemanLevel); + xf.push(`CAVEMAN:${cavemanLevel}`); + } - // Ponytail: inject lazy-senior-dev system prompt - if (tokenSaverEnabled && ponytailEnabled && ponytailLevel) { - injectPonytail(translatedBody, finalFormat, ponytailLevel); - xf.push(`PONYTAIL:${ponytailLevel}`); - } + // Ponytail: inject lazy-senior-dev system prompt + if (tokenSaverEnabled && ponytailEnabled && ponytailLevel) { + injectPonytail(translatedBody, finalFormat, ponytailLevel); + xf.push(`PONYTAIL:${ponytailLevel}`); + } - // PXPIPE: image bulky context (Claude-format bodies only), last saver before dispatch - let pxpipeSummary = null; - if (pxpipeEnabled) { - const pxpipeResult = await compressWithPxpipe(translatedBody, { - enabled: true, format: finalFormat, model: upstreamModel, - minChars: pxpipeMinChars, timeoutMs: pxpipeTimeoutMs, transform: pxpipeTransform, - }); - pxpipeSummary = pxpipeResult.summary; - if (pxpipeResult.body) translatedBody = pxpipeResult.body; - if (pxpipeSummary?.applied) xf.push(`PXPIPE:${pxpipeSummary.imageCount}img`); - try { onPxpipeEvent?.({ provider, model, ...pxpipeSummary }); } catch { /* stats must not break requests */ } - } + // PXPIPE: image bulky context (Claude-format bodies only), last saver before dispatch + let pxpipeSummary = null; + if (pxpipeEnabled) { + const pxpipeResult = await compressWithPxpipe(translatedBody, { + enabled: true, + format: finalFormat, + model: upstreamModel, + minChars: pxpipeMinChars, + timeoutMs: pxpipeTimeoutMs, + transform: pxpipeTransform, + }); + pxpipeSummary = pxpipeResult.summary; + if (pxpipeResult.body) translatedBody = pxpipeResult.body; + if (pxpipeSummary?.applied) + xf.push(`PXPIPE:${pxpipeSummary.imageCount}img`); + try { + onPxpipeEvent?.({ provider, model, ...pxpipeSummary }); + } catch { + /* stats must not break requests */ + } + } - if (xf.length && log?.line) log.line(reqTag, "⚙", xf.join(" · ")); + if (xf.length && log?.line) log.line(reqTag, "⚙", xf.join(" · ")); - const executor = getExecutor(provider); - trackPendingRequest(model, provider, connectionId, true); - appendRequestLog({ model, provider, connectionId, status: "PENDING" }).catch(() => { }); + const executor = getExecutor(provider); + trackPendingRequest(model, provider, connectionId, true); + appendRequestLog({ model, provider, connectionId, status: "PENDING" }).catch( + () => {}, + ); - const msgCount = translatedBody.messages?.length || translatedBody.input?.length || translatedBody.contents?.length || translatedBody.request?.contents?.length || 0; - log?.debug?.("REQUEST", `${provider.toUpperCase()} | ${model} | ${msgCount} msgs`); + const msgCount = + translatedBody.messages?.length || + translatedBody.input?.length || + translatedBody.contents?.length || + translatedBody.request?.contents?.length || + 0; + log?.debug?.( + "REQUEST", + `${provider.toUpperCase()} | ${model} | ${msgCount} msgs`, + ); - const streamController = createStreamController({ - onDisconnect: (reason) => { - trackPendingRequest(model, provider, connectionId, false); - if (onDisconnect) onDisconnect(reason); - }, - onError: () => trackPendingRequest(model, provider, connectionId, false), - log, provider, model, reqTag - }); + const streamController = createStreamController({ + onDisconnect: (reason) => { + trackPendingRequest(model, provider, connectionId, false); + if (onDisconnect) onDisconnect(reason); + }, + onError: () => trackPendingRequest(model, provider, connectionId, false), + log, + provider, + model, + reqTag, + }); - const proxyOptions = { - connectionProxyEnabled: credentials?.providerSpecificData?.connectionProxyEnabled === true, - connectionProxyUrl: credentials?.providerSpecificData?.connectionProxyUrl || "", - connectionNoProxy: credentials?.providerSpecificData?.connectionNoProxy || "", - vercelRelayUrl: credentials?.providerSpecificData?.vercelRelayUrl || "", - }; + const proxyOptions = { + connectionProxyEnabled: + credentials?.providerSpecificData?.connectionProxyEnabled === true, + connectionProxyUrl: + credentials?.providerSpecificData?.connectionProxyUrl || "", + connectionNoProxy: + credentials?.providerSpecificData?.connectionNoProxy || "", + vercelRelayUrl: credentials?.providerSpecificData?.vercelRelayUrl || "", + }; - if (proxyOptions.vercelRelayUrl) { - const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown"; - const poolId = credentials?.providerSpecificData?.connectionProxyPoolId || "none"; - log?.info?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | vercel-relay=${proxyOptions.vercelRelayUrl}`); - } else if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionProxyUrl) { - let maskedProxyUrl = proxyOptions.connectionProxyUrl; - try { - const parsed = new URL(proxyOptions.connectionProxyUrl); - const host = parsed.hostname || ""; - const port = parsed.port ? `:${parsed.port}` : ""; - const protocol = parsed.protocol || "http:"; - maskedProxyUrl = `${protocol}//${host}${port}`; - } catch { - // Keep raw if URL parsing fails - } + if (proxyOptions.vercelRelayUrl) { + const connectionName = + credentials?.connectionName || credentials?.connectionId || "unknown"; + const poolId = + credentials?.providerSpecificData?.connectionProxyPoolId || "none"; + log?.info?.( + "PROXY", + `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | vercel-relay=${proxyOptions.vercelRelayUrl}`, + ); + } else if ( + proxyOptions.connectionProxyEnabled && + proxyOptions.connectionProxyUrl + ) { + let maskedProxyUrl = proxyOptions.connectionProxyUrl; + try { + const parsed = new URL(proxyOptions.connectionProxyUrl); + const host = parsed.hostname || ""; + const port = parsed.port ? `:${parsed.port}` : ""; + const protocol = parsed.protocol || "http:"; + maskedProxyUrl = `${protocol}//${host}${port}`; + } catch { + // Keep raw if URL parsing fails + } - const poolId = credentials?.providerSpecificData?.connectionProxyPoolId || "none"; - const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown"; - log?.info?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | url=${maskedProxyUrl}`); - } + const poolId = + credentials?.providerSpecificData?.connectionProxyPoolId || "none"; + const connectionName = + credentials?.connectionName || credentials?.connectionId || "unknown"; + log?.info?.( + "PROXY", + `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | url=${maskedProxyUrl}`, + ); + } - if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionNoProxy) { - const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown"; - log?.debug?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | no_proxy=${proxyOptions.connectionNoProxy}`); - } + if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionNoProxy) { + const connectionName = + credentials?.connectionName || credentials?.connectionId || "unknown"; + log?.debug?.( + "PROXY", + `${provider.toUpperCase()} | ${model} | conn=${connectionName} | no_proxy=${proxyOptions.connectionNoProxy}`, + ); + } - // Execute request - let providerResponse, providerUrl, providerHeaders, finalBody; - // Most executors return their registry format. Cursor AgentService is an - // exception: it is decoded by the executor into OpenAI-compatible output. - let providerResponseFormat = targetFormat; - try { - const result = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions }); - providerResponse = result.response; - providerUrl = result.url; - providerHeaders = result.headers; - finalBody = result.transformedBody; - providerResponseFormat = result.responseFormat || targetFormat; - reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - } catch (error) { - trackPendingRequest(model, provider, connectionId, false, true); - appendRequestLog({ model, provider, connectionId, status: `FAILED ${error.name === "AbortError" ? 499 : HTTP_STATUS.BAD_GATEWAY}` }).catch(() => { }); - saveRequestDetail(buildRequestDetail({ - provider, model, connectionId, - latency: { ttft: 0, total: Date.now() - requestStartTime }, - tokens: { prompt_tokens: 0, completion_tokens: 0 }, - request: extractRequestConfig(body, stream), - providerRequest: translatedBody || null, - response: { error: error.message || String(error), status: error.name === "AbortError" ? 499 : 502, thinking: null }, - pxpipe: pxpipeSummary, - status: "error" - })).catch(() => { }); + // Execute request + let providerResponse, providerUrl, providerHeaders, finalBody; + // Most executors return their registry format. Cursor AgentService is an + // exception: it is decoded by the executor into OpenAI-compatible output. + let providerResponseFormat = targetFormat; + try { + const result = await executor.execute({ + model, + body: translatedBody, + stream, + credentials, + signal: streamController.signal, + log, + proxyOptions, + }); + providerResponse = result.response; + providerUrl = result.url; + providerHeaders = result.headers; + finalBody = result.transformedBody; + providerResponseFormat = result.responseFormat || targetFormat; + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); + } catch (error) { + trackPendingRequest(model, provider, connectionId, false, true); + appendRequestLog({ + model, + provider, + connectionId, + status: `FAILED ${error.name === "AbortError" ? 499 : HTTP_STATUS.BAD_GATEWAY}`, + }).catch(() => {}); + saveRequestDetail( + buildRequestDetail({ + provider, + model, + connectionId, + latency: { ttft: 0, total: Date.now() - requestStartTime }, + tokens: { prompt_tokens: 0, completion_tokens: 0 }, + request: extractRequestConfig(body, stream), + providerRequest: translatedBody || null, + response: { + error: error.message || String(error), + status: error.name === "AbortError" ? 499 : 502, + thinking: null, + }, + pxpipe: pxpipeSummary, + status: "error", + }), + ).catch(() => {}); - if (error.name === "AbortError") { - streamController.handleError(error); - return createErrorResult(499, "Request aborted"); - } - const errMsg = formatProviderError(error, provider, model, HTTP_STATUS.BAD_GATEWAY); - if (log?.errorLine) { - log.errorLine(reqTag, "✗", `ERROR 502 · ${provider}/${model} · ${Date.now() - requestStartTime}ms\n ${errMsg}${error.stack ? `\n ${error.stack}` : ""}`); - } - return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg); - } + if (error.name === "AbortError") { + streamController.handleError(error); + return createErrorResult(499, "Request aborted"); + } + const errMsg = formatProviderError( + error, + provider, + model, + HTTP_STATUS.BAD_GATEWAY, + ); + if (log?.errorLine) { + log.errorLine( + reqTag, + "✗", + `ERROR 502 · ${provider}/${model} · ${Date.now() - requestStartTime}ms\n ${errMsg}${error.stack ? `\n ${error.stack}` : ""}`, + ); + } + return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg); + } - // Handle 401/403 - try token refresh (skip for noAuth providers) - if (!executor.noAuth && (providerResponse.status === HTTP_STATUS.UNAUTHORIZED || providerResponse.status === HTTP_STATUS.FORBIDDEN)) { - try { - // Mutate credentials after each successful refresh: rotating refresh_token - // providers (xAI/grok-cli) issue a new RT on every refresh; without this, - // refreshWithRetry's 2nd/3rd attempt reuses the already-consumed RT → - // invalid_grant → auth_failed retryable=false. - const newCredentials = await refreshWithRetry(async () => { - const result = await executor.refreshCredentials(credentials, log); - if (result?.refreshToken && result.refreshToken !== credentials.refreshToken) { - if (result.accessToken) credentials.accessToken = result.accessToken; - credentials.refreshToken = result.refreshToken; - } - return result; - }, 3, log); - if (newCredentials?.accessToken || newCredentials?.copilotToken) { - if (log?.line) log.line(reqTag, "🔑", `TOKEN REFRESHED · ${provider}/${model}`); - Object.assign(credentials, newCredentials); - if (onCredentialsRefreshed) { - try { await onCredentialsRefreshed(newCredentials); } catch (e) { log?.warn?.("TOKEN", `onCredentialsRefreshed failed: ${e.message}`); } - } - try { - const retryResult = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions }); - if (retryResult.response.ok) { - providerResponse = retryResult.response; - providerUrl = retryResult.url; - providerResponseFormat = retryResult.responseFormat || targetFormat; - } - } catch { log?.warn?.("TOKEN", `${provider.toUpperCase()} | retry after refresh failed`); } - } else { - log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`); - } - } catch (e) { - log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh threw: ${e.message}`); - } - } + // Config-driven in-stream error detection: peek the first bytes of a 200 + // stream; if a configured pattern matches, fail fast with 502 so account + // /combo fallback can still run before any byte reaches the client. + const configuredPatterns = streamErrorPatterns?.[provider]; + if ( + providerResponse?.ok && + Array.isArray(configuredPatterns) && + configuredPatterns.length > 0 + ) { + providerResponse = await maybeRejectEarlyStreamError( + providerResponse, + configuredPatterns, + { signal: streamController.signal }, + ); + } - // Provider returned error - if (!providerResponse.ok) { - trackPendingRequest(model, provider, connectionId, false, true); - const { statusCode, message, resetsAtMs } = await parseUpstreamError(providerResponse, executor); - appendRequestLog({ model, provider, connectionId, status: `FAILED ${statusCode}` }).catch(() => { }); - saveRequestDetail(buildRequestDetail({ - provider, model, connectionId, - latency: { ttft: 0, total: Date.now() - requestStartTime }, - tokens: { prompt_tokens: 0, completion_tokens: 0 }, - request: extractRequestConfig(body, stream), - providerRequest: finalBody || translatedBody || null, - response: { error: message, status: statusCode, thinking: null }, - pxpipe: pxpipeSummary, - status: "error" - })).catch(() => { }); + // Handle 401/403 - try token refresh (skip for noAuth providers) + if ( + !executor.noAuth && + (providerResponse.status === HTTP_STATUS.UNAUTHORIZED || + providerResponse.status === HTTP_STATUS.FORBIDDEN) + ) { + try { + // Mutate credentials after each successful refresh: rotating refresh_token + // providers (xAI/grok-cli) issue a new RT on every refresh; without this, + // refreshWithRetry's 2nd/3rd attempt reuses the already-consumed RT → + // invalid_grant → auth_failed retryable=false. + const newCredentials = await refreshWithRetry( + async () => { + const result = await executor.refreshCredentials(credentials, log); + if ( + result?.refreshToken && + result.refreshToken !== credentials.refreshToken + ) { + if (result.accessToken) + credentials.accessToken = result.accessToken; + credentials.refreshToken = result.refreshToken; + } + return result; + }, + 3, + log, + ); + if (newCredentials?.accessToken || newCredentials?.copilotToken) { + if (log?.line) + log.line(reqTag, "🔑", `TOKEN REFRESHED · ${provider}/${model}`); + Object.assign(credentials, newCredentials); + if (onCredentialsRefreshed) { + try { + await onCredentialsRefreshed(newCredentials); + } catch (e) { + log?.warn?.("TOKEN", `onCredentialsRefreshed failed: ${e.message}`); + } + } + try { + const retryResult = await executor.execute({ + model, + body: translatedBody, + stream, + credentials, + signal: streamController.signal, + log, + proxyOptions, + }); + if (retryResult.response.ok) { + providerResponse = retryResult.response; + providerUrl = retryResult.url; + providerResponseFormat = retryResult.responseFormat || targetFormat; + } + } catch { + log?.warn?.( + "TOKEN", + `${provider.toUpperCase()} | retry after refresh failed`, + ); + } + } else { + log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`); + } + } catch (e) { + log?.warn?.( + "TOKEN", + `${provider.toUpperCase()} | refresh threw: ${e.message}`, + ); + } + } - const errMsg = formatProviderError(new Error(message), provider, model, statusCode); - if (log?.errorLine) { - const urlStr = providerUrl ? `\n URL: ${providerUrl}` : ""; - 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); - } + // Provider returned error + if (!providerResponse.ok) { + trackPendingRequest(model, provider, connectionId, false, true); + const { statusCode, message, resetsAtMs } = await parseUpstreamError( + providerResponse, + executor, + ); + appendRequestLog({ + model, + provider, + connectionId, + status: `FAILED ${statusCode}`, + }).catch(() => {}); + saveRequestDetail( + buildRequestDetail({ + provider, + model, + connectionId, + latency: { ttft: 0, total: Date.now() - requestStartTime }, + tokens: { prompt_tokens: 0, completion_tokens: 0 }, + request: extractRequestConfig(body, stream), + providerRequest: finalBody || translatedBody || null, + response: { error: message, status: statusCode, thinking: null }, + pxpipe: pxpipeSummary, + status: "error", + }), + ).catch(() => {}); - const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, pxpipe: pxpipeSummary, reqTag, log }; - const appendLog = (extra) => appendRequestLog({ model, provider, connectionId, ...extra }).catch(() => { }); - const trackDone = () => trackPendingRequest(model, provider, connectionId, false); + const errMsg = formatProviderError( + new Error(message), + provider, + model, + statusCode, + ); + if (log?.errorLine) { + const urlStr = providerUrl ? `\n URL: ${providerUrl}` : ""; + 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); + } - // Provider forced streaming but client wants JSON - if (!clientRequestedStreaming && providerRequiresStreaming) { - const result = await handleForcedSSEToJson({ ...sharedCtx, providerResponse, sourceFormat, trackDone, appendLog }); - if (result) { streamController.handleComplete(); return result; } - } + const sharedCtx = { + provider, + model, + body, + stream, + translatedBody, + finalBody, + requestStartTime, + connectionId, + apiKey, + clientRawRequest, + onRequestSuccess, + pxpipe: pxpipeSummary, + reqTag, + log, + streamErrorPatterns, + }; + const appendLog = (extra) => + appendRequestLog({ model, provider, connectionId, ...extra }).catch( + () => {}, + ); + const trackDone = () => + trackPendingRequest(model, provider, connectionId, false); - // True non-streaming response - if (!stream) { - const result = await handleNonStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat: providerResponseFormat, reqLogger, toolNameMap, trackDone, appendLog }); - streamController.handleComplete(); - return result; - } + // Provider forced streaming but client wants JSON + if (!clientRequestedStreaming && providerRequiresStreaming) { + const result = await handleForcedSSEToJson({ + ...sharedCtx, + providerResponse, + sourceFormat, + trackDone, + appendLog, + }); + if (result) { + streamController.handleComplete(); + return result; + } + } - // Streaming response - const { onStreamComplete, streamDetailId } = buildOnStreamComplete({ ...sharedCtx }); - return handleStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat: providerResponseFormat, userAgent, reqLogger, toolNameMap, streamController, onStreamComplete, streamDetailId }); + // True non-streaming response + if (!stream) { + const result = await handleNonStreamingResponse({ + ...sharedCtx, + providerResponse, + sourceFormat, + targetFormat: providerResponseFormat, + reqLogger, + toolNameMap, + trackDone, + appendLog, + }); + streamController.handleComplete(); + return result; + } + + // Streaming response + const { onStreamComplete, streamDetailId } = buildOnStreamComplete({ + ...sharedCtx, + }); + return handleStreamingResponse({ + ...sharedCtx, + providerResponse, + sourceFormat, + targetFormat: providerResponseFormat, + userAgent, + reqLogger, + toolNameMap, + streamController, + onStreamComplete, + streamDetailId, + }); } export function isTokenExpiringSoon(expiresAt, bufferMs = 5 * 60 * 1000) { - if (!expiresAt) return false; - return new Date(expiresAt).getTime() - Date.now() < bufferMs; + if (!expiresAt) return false; + return new Date(expiresAt).getTime() - Date.now() < bufferMs; } diff --git a/open-sse/utils/streamErrorPeek.js b/open-sse/utils/streamErrorPeek.js new file mode 100644 index 00000000..3e180733 --- /dev/null +++ b/open-sse/utils/streamErrorPeek.js @@ -0,0 +1,136 @@ +import { matchStreamErrorPatterns } from "./streamErrorPatterns.js"; +import { HTTP_STATUS } from "../config/runtimeConfig.js"; + +const DEFAULT_TIMEOUT_MS = (() => { + const raw = process.env.STREAM_ERROR_PEEK_TIMEOUT_MS; + const n = raw ? parseInt(raw, 10) : NaN; + return Number.isFinite(n) && n > 0 ? n : 3000; +})(); + +const DEFAULT_MAX_BYTES = (() => { + const raw = process.env.STREAM_ERROR_PEEK_MAX_BYTES; + const n = raw ? parseInt(raw, 10) : NaN; + return Number.isFinite(n) && n > 0 ? n : 8192; +})(); + +function makeAbortError(reason) { + const error = new Error(reason?.message || reason || "Request aborted"); + error.name = "AbortError"; + return error; +} + +/** + * Read the first bytes of a 200 response and reject it with a 502 when a + * configured stream-error pattern matches — BEFORE any byte reaches the + * client, so account/model fallback can still kick in for streaming too. + * On no-match/timeout/abort the response is re-emitted (buffered bytes + + * rest of stream) unchanged in spirit. Fail-open: never throws. + */ +export async function maybeRejectEarlyStreamError( + response, + patterns, + { + signal = null, + timeoutMs = DEFAULT_TIMEOUT_MS, + maxBytes = DEFAULT_MAX_BYTES, + } = {}, +) { + if (!Array.isArray(patterns) || patterns.length === 0) return response; + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + const abortController = new AbortController(); + const forwardAbort = () => abortController.abort(signal?.reason); + if (signal?.aborted) abortController.abort(signal?.reason); + else if (signal) + signal.addEventListener("abort", forwardAbort, { once: true }); + + // Raw bytes for lossless re-emission; decoded text ONLY for pattern matching. + // Never re-encode decoded text: TextDecoder holds a split multi-byte char + // internally and flush() would replace it with U+FFFD, corrupting the stream. + const rawChunks = []; + let peekedText = ""; + let total = 0; + + const readWithTimeout = (ms) => { + if (abortController.signal.aborted) + return Promise.reject(makeAbortError(abortController.signal.reason)); + const timeoutPromise = new Promise((_, reject) => { + const t = setTimeout( + () => reject(new Error("stream error peek timeout")), + ms, + ); + t.unref?.(); + }); + const abortPromise = new Promise((_, reject) => { + abortController.signal.addEventListener( + "abort", + () => reject(makeAbortError(abortController.signal.reason)), + { once: true }, + ); + }); + return Promise.race([reader.read(), timeoutPromise, abortPromise]); + }; + + try { + const deadline = Date.now() + timeoutMs; + while (total < maxBytes && Date.now() < deadline) { + const { done, value } = await readWithTimeout( + Math.max(deadline - Date.now(), 1), + ); + if (done) break; + rawChunks.push(value); + peekedText += decoder.decode(value, { stream: true }); + total += value.byteLength; + const matched = matchStreamErrorPatterns(patterns, peekedText); + if (matched) { + await reader.cancel("stream error pattern matched").catch(() => {}); + return new Response( + JSON.stringify({ + error: { + message: `Stream error pattern matched: ${matched}`, + type: "upstream_error", + }, + }), + { + status: HTTP_STATUS.BAD_GATEWAY, + statusText: String(matched).slice(0, 200), + headers: { "Content-Type": "application/json" }, + }, + ); + } + } + } catch { + // timeout / abort / read failure → commit; downstream stall/abort handling takes over. + } + + if (signal) signal.removeEventListener("abort", forwardAbort); + + const remaining = new ReadableStream({ + start(controller) { + (async () => { + try { + // Re-emit RAW bytes (never re-encoded decoded text) so split + // multi-byte UTF-8 sequences survive the peek untouched. + for (const c of rawChunks) controller.enqueue(c); + while (true) { + const { done, value } = await reader.read(); + if (done) break; + controller.enqueue(value); + } + controller.close(); + } catch (err) { + controller.error(err); + } + })(); + }, + cancel() { + reader.cancel("stream cancelled during peek commit").catch(() => {}); + }, + }); + + return new Response(remaining, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }); +} diff --git a/src/sse/handlers/chat.js b/src/sse/handlers/chat.js index 383bdec9..8d56a0eb 100644 --- a/src/sse/handlers/chat.js +++ b/src/sse/handlers/chat.js @@ -1,11 +1,11 @@ import "open-sse/index.js"; import { - getProviderCredentials, - markAccountUnavailable, - clearAccountError, - extractApiKey, - isValidApiKey, + getProviderCredentials, + markAccountUnavailable, + clearAccountError, + extractApiKey, + isValidApiKey, } from "../services/auth.js"; import { cacheClaudeHeaders } from "open-sse/utils/claudeHeaderCache.js"; import { getSettings } from "@/lib/localDb"; @@ -20,7 +20,10 @@ import { handleBypassRequest } from "open-sse/utils/bypassHandler.js"; import { HTTP_STATUS } from "open-sse/config/runtimeConfig.js"; import { detectFormatByEndpoint } from "open-sse/translator/formats.js"; import * as log from "../utils/logger.js"; -import { updateProviderCredentials, checkAndRefreshToken } from "../services/tokenRefresh.js"; +import { + updateProviderCredentials, + checkAndRefreshToken, +} from "../services/tokenRefresh.js"; import { getProjectIdForConnection } from "open-sse/services/projectId.js"; /** @@ -29,266 +32,355 @@ import { getProjectIdForConnection } from "open-sse/services/projectId.js"; * Format detection and translation handled by translator */ export async function handleChat(request, clientRawRequest = null) { - let body; - try { - body = await request.json(); - } catch { - log.warn("CHAT", "Invalid JSON body"); - return errorResponse(HTTP_STATUS.BAD_REQUEST, "Invalid JSON body"); - } + let body; + try { + body = await request.json(); + } catch { + log.warn("CHAT", "Invalid JSON body"); + return errorResponse(HTTP_STATUS.BAD_REQUEST, "Invalid JSON body"); + } - // Build clientRawRequest for logging (if not provided) - if (!clientRawRequest) { - const url = new URL(request.url); - clientRawRequest = { - endpoint: url.pathname, - body, - headers: Object.fromEntries(request.headers.entries()) - }; - } - cacheClaudeHeaders(clientRawRequest.headers); + // Build clientRawRequest for logging (if not provided) + if (!clientRawRequest) { + const url = new URL(request.url); + clientRawRequest = { + endpoint: url.pathname, + body, + headers: Object.fromEntries(request.headers.entries()), + }; + } + cacheClaudeHeaders(clientRawRequest.headers); - const modelStr = body.model; + const modelStr = body.model; - // Request summary is emitted as the unified "▶" line in chatCore (has fmt/thinking/account) + // Request summary is emitted as the unified "▶" line in chatCore (has fmt/thinking/account) - // Log API key (masked) - const authHeader = request.headers.get("Authorization"); - const apiKey = extractApiKey(request); - if (authHeader && apiKey) { - const masked = log.maskKey(apiKey); - log.debug("AUTH", `API Key: ${masked}`); - } else { - log.debug("AUTH", "No API key provided (local mode)"); - } + // Log API key (masked) + const authHeader = request.headers.get("Authorization"); + const apiKey = extractApiKey(request); + if (authHeader && apiKey) { + const masked = log.maskKey(apiKey); + log.debug("AUTH", `API Key: ${masked}`); + } else { + log.debug("AUTH", "No API key provided (local mode)"); + } - // Enforce API key if enabled in settings - const settings = await getSettings(); - if (settings.requireApiKey) { - if (!apiKey) { - log.warn("AUTH", "Missing API key (requireApiKey=true)"); - return errorResponse(HTTP_STATUS.UNAUTHORIZED, "Missing API key"); - } - const valid = await isValidApiKey(apiKey); - if (!valid) { - log.warn("AUTH", "Invalid API key (requireApiKey=true)"); - return errorResponse(HTTP_STATUS.UNAUTHORIZED, "Invalid API key"); - } - } + // Enforce API key if enabled in settings + const settings = await getSettings(); + if (settings.requireApiKey) { + if (!apiKey) { + log.warn("AUTH", "Missing API key (requireApiKey=true)"); + return errorResponse(HTTP_STATUS.UNAUTHORIZED, "Missing API key"); + } + const valid = await isValidApiKey(apiKey); + if (!valid) { + log.warn("AUTH", "Invalid API key (requireApiKey=true)"); + return errorResponse(HTTP_STATUS.UNAUTHORIZED, "Invalid API key"); + } + } - if (!modelStr) { - log.warn("CHAT", "Missing model"); - return errorResponse(HTTP_STATUS.BAD_REQUEST, "Missing model"); - } + if (!modelStr) { + log.warn("CHAT", "Missing model"); + return errorResponse(HTTP_STATUS.BAD_REQUEST, "Missing model"); + } - // Bypass naming/warmup requests before combo rotation to avoid wasting rotation slots - const userAgent = request?.headers?.get("user-agent") || ""; - const bypassResponse = handleBypassRequest(body, modelStr, userAgent, !!settings.ccFilterNaming); - if (bypassResponse) return bypassResponse.response || bypassResponse; + // Bypass naming/warmup requests before combo rotation to avoid wasting rotation slots + const userAgent = request?.headers?.get("user-agent") || ""; + const bypassResponse = handleBypassRequest( + body, + modelStr, + userAgent, + !!settings.ccFilterNaming, + ); + if (bypassResponse) return bypassResponse.response || bypassResponse; - // Check if model is a combo (has multiple models with fallback) - const comboModels = await getComboModels(modelStr); - if (comboModels) { - // Check for combo-specific strategy first, fallback to global - const comboStrategies = settings.comboStrategies || {}; - const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy; - const comboStrategy = comboSpecificStrategy || settings.comboStrategy || "fallback"; + // Check if model is a combo (has multiple models with fallback) + const comboModels = await getComboModels(modelStr); + if (comboModels) { + // Check for combo-specific strategy first, fallback to global + const comboStrategies = settings.comboStrategies || {}; + const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy; + const comboStrategy = + comboSpecificStrategy || settings.comboStrategy || "fallback"; - if (comboStrategy === "fusion") { - log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`); - return handleFusionChat({ - body, - models: comboModels, - handleSingleModel: (b, m, isPanel) => { - let cleanRawReq = clientRawRequest; - if (isPanel && clientRawRequest) { - const { tools, tool_choice, ...cleanBody } = clientRawRequest.body || {}; - cleanRawReq = { ...clientRawRequest, body: cleanBody }; - } - return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); - }, - log, - comboName: modelStr, - judgeModel: comboStrategies[modelStr]?.judgeModel, - tuning: comboStrategies[modelStr]?.fusionTuning, - }); - } + if (comboStrategy === "fusion") { + log.info( + "CHAT", + `Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`, + ); + return handleFusionChat({ + body, + models: comboModels, + handleSingleModel: (b, m, isPanel) => { + let cleanRawReq = clientRawRequest; + if (isPanel && clientRawRequest) { + const { tools, tool_choice, ...cleanBody } = + clientRawRequest.body || {}; + cleanRawReq = { ...clientRawRequest, body: cleanBody }; + } + return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); + }, + log, + comboName: modelStr, + judgeModel: comboStrategies[modelStr]?.judgeModel, + tuning: comboStrategies[modelStr]?.fusionTuning, + }); + } - const comboStickyLimit = settings.comboStickyRoundRobinLimit; - log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`); - return handleComboChat({ - body, - models: comboModels, - handleSingleModel: (b, m) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey), - log, - comboName: modelStr, - comboStrategy, - comboStickyLimit - }); - } + const comboStickyLimit = settings.comboStickyRoundRobinLimit; + log.info( + "CHAT", + `Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`, + ); + return handleComboChat({ + body, + models: comboModels, + handleSingleModel: (b, m) => + handleSingleModelChat(b, m, clientRawRequest, request, apiKey), + log, + comboName: modelStr, + comboStrategy, + comboStickyLimit, + }); + } - // Single model request - return handleSingleModelChat(body, modelStr, clientRawRequest, request, apiKey); + // Single model request + return handleSingleModelChat( + body, + modelStr, + clientRawRequest, + request, + apiKey, + ); } /** * Handle single model chat request */ -async function handleSingleModelChat(body, modelStr, clientRawRequest = null, request = null, apiKey = null) { - const modelInfo = await getModelInfo(modelStr); +async function handleSingleModelChat( + body, + modelStr, + clientRawRequest = null, + request = null, + apiKey = null, +) { + const modelInfo = await getModelInfo(modelStr); - // If provider is null, this might be a combo name - check and handle - if (!modelInfo.provider) { - const comboModels = await getComboModels(modelStr); - if (comboModels) { - const chatSettings = await getSettings(); - // Check for combo-specific strategy first, fallback to global - const comboStrategies = chatSettings.comboStrategies || {}; - const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy; - const comboStrategy = comboSpecificStrategy || chatSettings.comboStrategy || "fallback"; + // If provider is null, this might be a combo name - check and handle + if (!modelInfo.provider) { + const comboModels = await getComboModels(modelStr); + if (comboModels) { + const chatSettings = await getSettings(); + // Check for combo-specific strategy first, fallback to global + const comboStrategies = chatSettings.comboStrategies || {}; + const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy; + const comboStrategy = + comboSpecificStrategy || chatSettings.comboStrategy || "fallback"; - if (comboStrategy === "fusion") { - log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`); - return handleFusionChat({ - body, - models: comboModels, - handleSingleModel: (b, m, isPanel) => { - let cleanRawReq = clientRawRequest; - if (isPanel && clientRawRequest) { - const { tools, tool_choice, ...cleanBody } = clientRawRequest.body || {}; - cleanRawReq = { ...clientRawRequest, body: cleanBody }; - } - return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); - }, - log, - comboName: modelStr, - judgeModel: comboStrategies[modelStr]?.judgeModel, - tuning: comboStrategies[modelStr]?.fusionTuning, - }); - } + if (comboStrategy === "fusion") { + log.info( + "CHAT", + `Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`, + ); + return handleFusionChat({ + body, + models: comboModels, + handleSingleModel: (b, m, isPanel) => { + let cleanRawReq = clientRawRequest; + if (isPanel && clientRawRequest) { + const { tools, tool_choice, ...cleanBody } = + clientRawRequest.body || {}; + cleanRawReq = { ...clientRawRequest, body: cleanBody }; + } + return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); + }, + log, + comboName: modelStr, + judgeModel: comboStrategies[modelStr]?.judgeModel, + tuning: comboStrategies[modelStr]?.fusionTuning, + }); + } - const comboStickyLimit = chatSettings.comboStickyRoundRobinLimit; - log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`); - return handleComboChat({ - body, - models: comboModels, - handleSingleModel: (b, m) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey), - log, - comboName: modelStr, - comboStrategy, - comboStickyLimit - }); - } - log.warn("CHAT", "Invalid model format", { model: modelStr }); - return errorResponse(HTTP_STATUS.BAD_REQUEST, "Invalid model format"); - } + const comboStickyLimit = chatSettings.comboStickyRoundRobinLimit; + log.info( + "CHAT", + `Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`, + ); + return handleComboChat({ + body, + models: comboModels, + handleSingleModel: (b, m) => + handleSingleModelChat(b, m, clientRawRequest, request, apiKey), + log, + comboName: modelStr, + comboStrategy, + comboStickyLimit, + }); + } + log.warn("CHAT", "Invalid model format", { model: modelStr }); + return errorResponse(HTTP_STATUS.BAD_REQUEST, "Invalid model format"); + } - const { provider, model } = modelInfo; + const { provider, model } = modelInfo; - // Routing shown in the unified "▶" line (client model → provider/model) + // Routing shown in the unified "▶" line (client model → provider/model) - // Extract userAgent from request - const userAgent = request?.headers?.get("user-agent") || ""; - // Optional pin to a specific connection (dashboard test / client override) - const preferredConnectionId = request?.headers?.get("x-connection-id") || null; + // Extract userAgent from request + const userAgent = request?.headers?.get("user-agent") || ""; + // Optional pin to a specific connection (dashboard test / client override) + const preferredConnectionId = + request?.headers?.get("x-connection-id") || null; - // Try with available accounts (fallback on errors unless pinned) - const excludeConnectionIds = new Set(); - let lastError = null; - let lastStatus = null; + // Try with available accounts (fallback on errors unless pinned) + const excludeConnectionIds = new Set(); + let lastError = null; + let lastStatus = null; - while (true) { - const credentials = await getProviderCredentials(provider, excludeConnectionIds, model, { preferredConnectionId }); + while (true) { + const credentials = await getProviderCredentials( + provider, + excludeConnectionIds, + model, + { preferredConnectionId }, + ); - // 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; - log.warn("CHAT", `[${provider}/${model}] ${errorMsg} (${credentials.retryAfterHuman})`); - return unavailableResponse(status, `[${provider}/${model}] ${errorMsg}`, credentials.retryAfter, credentials.retryAfterHuman); - } - if (excludeConnectionIds.size === 0) { - log.warn("AUTH", `No active credentials for provider: ${provider}`); - return errorResponse(HTTP_STATUS.NOT_FOUND, `No active credentials for provider: ${provider}`); - } - log.warn("CHAT", "No more accounts available", { provider }); - return errorResponse(lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE, lastError || "All accounts unavailable"); - } + // 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; + log.warn( + "CHAT", + `[${provider}/${model}] ${errorMsg} (${credentials.retryAfterHuman})`, + ); + return unavailableResponse( + status, + `[${provider}/${model}] ${errorMsg}`, + credentials.retryAfter, + credentials.retryAfterHuman, + ); + } + if (excludeConnectionIds.size === 0) { + log.warn("AUTH", `No active credentials for provider: ${provider}`); + return errorResponse( + HTTP_STATUS.NOT_FOUND, + `No active credentials for provider: ${provider}`, + ); + } + log.warn("CHAT", "No more accounts available", { provider }); + return errorResponse( + lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE, + lastError || "All accounts unavailable", + ); + } - // Account selection shown in the unified "▶" line (acc:...) - const refreshedCredentials = await checkAndRefreshToken(provider, credentials); + // Account selection shown in the unified "▶" line (acc:...) + const refreshedCredentials = await checkAndRefreshToken( + provider, + credentials, + ); - // Ensure real project ID is available for providers that need it (P0 fix: cold miss) - if ((provider === "antigravity" || provider === "gemini-cli") && !refreshedCredentials.projectId) { - const pid = await getProjectIdForConnection(credentials.connectionId, refreshedCredentials.accessToken, provider); - if (pid) { - refreshedCredentials.projectId = pid; - // Persist to DB in background so subsequent requests have it immediately - updateProviderCredentials(credentials.connectionId, { projectId: pid }).catch(() => { }); - } - } + // Ensure real project ID is available for providers that need it (P0 fix: cold miss) + if ( + (provider === "antigravity" || provider === "gemini-cli") && + !refreshedCredentials.projectId + ) { + const pid = await getProjectIdForConnection( + credentials.connectionId, + refreshedCredentials.accessToken, + provider, + ); + if (pid) { + refreshedCredentials.projectId = pid; + // Persist to DB in background so subsequent requests have it immediately + updateProviderCredentials(credentials.connectionId, { + projectId: pid, + }).catch(() => {}); + } + } - // Use shared chatCore - const chatSettings = await getSettings(); - const providerThinking = (chatSettings.providerThinking || {})[provider] || null; - const result = await handleChatCore({ - body: { ...body, model: `${provider}/${model}` }, - modelInfo: { provider, model }, - credentials: refreshedCredentials, - log, - clientRawRequest, - connectionId: credentials.connectionId, - userAgent, - apiKey, - ccFilterNaming: !!chatSettings.ccFilterNaming, - rtkEnabled: !!chatSettings.rtkEnabled, - headroomEnabled: !!chatSettings.headroomEnabled, - headroomUrl: chatSettings.headroomUrl || DEFAULT_HEADROOM_URL, - headroomCompressUserMessages: !!chatSettings.headroomCompressUserMessages, - cavemanEnabled: !!chatSettings.cavemanEnabled, - cavemanLevel: chatSettings.cavemanLevel || "full", - ponytailEnabled: !!chatSettings.ponytailEnabled, - ponytailLevel: chatSettings.ponytailLevel || "full", - pxpipeEnabled: !!chatSettings.pxpipeEnabled, - pxpipeMinChars: chatSettings.pxpipeMinChars, - pxpipeTimeoutMs: chatSettings.pxpipeTimeoutMs, - // Lazily warms the in-process module on first use; null when not installed (fail-open) - pxpipeTransform: chatSettings.pxpipeEnabled ? await getPxpipeTransform() : null, - onPxpipeEvent: appendPxpipeEvent, - providerThinking, - // Detect source format by endpoint + body - sourceFormatOverride: request?.url ? detectFormatByEndpoint(new URL(request.url).pathname, body) : null, - onCredentialsRefreshed: async (newCreds) => { - await updateProviderCredentials(credentials.connectionId, { - ...newCreds, - existingProviderSpecificData: credentials.providerSpecificData, - testStatus: "active" - }); - }, - onRequestSuccess: async () => { - await clearAccountError(credentials.connectionId, credentials, model); - } - }); + // Use shared chatCore + const chatSettings = await getSettings(); + const providerThinking = + (chatSettings.providerThinking || {})[provider] || null; + const result = await handleChatCore({ + body: { ...body, model: `${provider}/${model}` }, + modelInfo: { provider, model }, + credentials: refreshedCredentials, + log, + clientRawRequest, + connectionId: credentials.connectionId, + userAgent, + apiKey, + ccFilterNaming: !!chatSettings.ccFilterNaming, + rtkEnabled: !!chatSettings.rtkEnabled, + headroomEnabled: !!chatSettings.headroomEnabled, + headroomUrl: chatSettings.headroomUrl || DEFAULT_HEADROOM_URL, + headroomCompressUserMessages: !!chatSettings.headroomCompressUserMessages, + cavemanEnabled: !!chatSettings.cavemanEnabled, + cavemanLevel: chatSettings.cavemanLevel || "full", + ponytailEnabled: !!chatSettings.ponytailEnabled, + ponytailLevel: chatSettings.ponytailLevel || "full", + pxpipeEnabled: !!chatSettings.pxpipeEnabled, + pxpipeMinChars: chatSettings.pxpipeMinChars, + pxpipeTimeoutMs: chatSettings.pxpipeTimeoutMs, + // Lazily warms the in-process module on first use; null when not installed (fail-open) + pxpipeTransform: chatSettings.pxpipeEnabled + ? await getPxpipeTransform() + : null, + onPxpipeEvent: appendPxpipeEvent, + providerThinking, + streamErrorPatterns: chatSettings.streamErrorPatterns || {}, + // Detect source format by endpoint + body + sourceFormatOverride: request?.url + ? detectFormatByEndpoint(new URL(request.url).pathname, body) + : null, + onCredentialsRefreshed: async (newCreds) => { + await updateProviderCredentials(credentials.connectionId, { + ...newCreds, + existingProviderSpecificData: credentials.providerSpecificData, + testStatus: "active", + }); + }, + onRequestSuccess: async () => { + await clearAccountError(credentials.connectionId, credentials, model); + }, + }); - if (result.success) return result.response; + if (result.success) return 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); + // 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, + ); - if (shouldFallback) { - // When a connection is explicitly pinned, never rotate to another account. - if (preferredConnectionId) { - log.warn("AUTH", `Pinned account ${credentials.connectionName} unavailable (${result.status}), no fallback`); - return result.response; - } - log.warn("FALLBACK", `⇄ ACC:${credentials.connectionName} UNAVAILABLE (${result.status}) → NEXT ACCOUNT`); - excludeConnectionIds.add(credentials.connectionId); - lastError = result.error; - lastStatus = result.status; - continue; - } + if (shouldFallback) { + // When a connection is explicitly pinned, never rotate to another account. + if (preferredConnectionId) { + log.warn( + "AUTH", + `Pinned account ${credentials.connectionName} unavailable (${result.status}), no fallback`, + ); + return result.response; + } + log.warn( + "FALLBACK", + `⇄ ACC:${credentials.connectionName} UNAVAILABLE (${result.status}) → NEXT ACCOUNT`, + ); + excludeConnectionIds.add(credentials.connectionId); + lastError = result.error; + lastStatus = result.status; + continue; + } - return result.response; - } + return result.response; + } } diff --git a/tests/unit/commandcode-executor.test.js b/tests/unit/commandcode-executor.test.js index 1ad7a376..63cfd489 100644 --- a/tests/unit/commandcode-executor.test.js +++ b/tests/unit/commandcode-executor.test.js @@ -98,4 +98,26 @@ describe("commandcode executor — early-error peek", () => { expect(res.status).toBe(200); await res.body.cancel(); }); + + it("re-emits raw bytes so a multi-byte char split across the peek boundary survives", async () => { + // chunk1 ends mid-é (0xC3); the peek commits on the first complete line + // (text-delta "hi") while the decoder still holds 0xC3. Re-emission must + // use RAW bytes — re-encoding decoded text would replace 0xC3 with U+FFFD. + const first = Buffer.from( + '{"type":"text-delta","text":"hi"}\n{"type":"text-delta","text":"caf', + ); + const chunk1 = new Uint8Array([...first, 0xc3]); + const rest = new Uint8Array([0xa9, 0x22, 0x7d, 0x0a]); // é"}\n + const body = new ReadableStream({ + start(c) { + c.enqueue(chunk1); + c.enqueue(rest); + c.close(); + }, + }); + const res = await peekForUpstreamError(new Response(body, { status: 200 }), "m"); + const text = await res.text(); + expect(text).toContain("café"); + expect(text).not.toContain("\uFFFD"); + }); }); diff --git a/tests/unit/stream-error-peek.test.js b/tests/unit/stream-error-peek.test.js new file mode 100644 index 00000000..3a5df51c --- /dev/null +++ b/tests/unit/stream-error-peek.test.js @@ -0,0 +1,99 @@ +import { describe, it, expect } from "vitest"; +import { maybeRejectEarlyStreamError } from "../../open-sse/utils/streamErrorPeek.js"; + +const encoder = new TextEncoder(); +const sseResp = (lines) => + new Response( + new ReadableStream({ + start(c) { + for (const l of lines) c.enqueue(encoder.encode(l)); + c.close(); + }, + }), + { status: 200, headers: { "content-type": "text/event-stream" } }, + ); + +describe("maybeRejectEarlyStreamError", () => { + it("returns 502 when a pattern matches early stream text", async () => { + const res = await maybeRejectEarlyStreamError( + sseResp([ + '{"type":"start"}\n{"type":"error","error":{"type":"server_error","message":"Network connection lost."}}\n', + ]), + ["server_error"], + ); + expect(res.status).toBe(502); + const body = await res.json(); + expect(body.error.message).toContain("server_error"); + }); + + it("passes the stream through unchanged when nothing matches", async () => { + const res = await maybeRejectEarlyStreamError( + sseResp([ + 'data: {"choices":[{"delta":{"content":"hello"},"finish_reason":null}]}\n\ndata: [DONE]\n\n', + ]), + ["server_error"], + ); + expect(res.status).toBe(200); + const text = await res.text(); + expect(text).toContain("hello"); + expect(text).toContain("[DONE]"); + }); + + it("commits on timeout without hanging", async () => { + const stalled = new Response(new ReadableStream({ start() {} }), { + status: 200, + }); + const res = await maybeRejectEarlyStreamError(stalled, ["x"], { + timeoutMs: 50, + }); + expect(res.status).toBe(200); + await res.body.cancel(); + }); + + it("commits on abort without hanging", async () => { + const ctrl = new AbortController(); + const stalled = new Response(new ReadableStream({ start() {} }), { + status: 200, + }); + setTimeout(() => ctrl.abort(new Error("gone")), 10); + const res = await maybeRejectEarlyStreamError(stalled, ["x"], { + signal: ctrl.signal, + timeoutMs: 2000, + }); + expect(res.status).toBe(200); + await res.body.cancel(); + }); + + it("commits (passthrough) when patterns are empty", async () => { + const res = await maybeRejectEarlyStreamError( + sseResp(["data: hi\n\n"]), + [], + ); + expect(res.status).toBe(200); + expect(await res.text()).toBe("data: hi\n\n"); + }); + + it("multi-byte UTF-8 split across peek boundary round-trips losslessly", async () => { + // "café" split mid-é: the peek consumes bytes up to and including 0xC3 + // (the first half of é); the re-emitted stream must contain the RAW bytes + // (never a re-encoded decoded string — TextDecoder flush would replace the + // lone 0xC3 with U+FFFD and corrupt the output). + const bytes = [ + 0x64, 0x61, 0x74, 0x61, 0x3a, 0x20, 0x22, 0x63, 0x61, 0x66, 0xc3, + ]; + const rest = new Uint8Array([0xa9, 0x22, 0x0a, 0x0a]); + const body = new ReadableStream({ + start(c) { + c.enqueue(new Uint8Array(bytes)); + c.enqueue(rest); + c.close(); + }, + }); + const res = await maybeRejectEarlyStreamError( + new Response(body, { status: 200 }), + ["nomatch"], + { maxBytes: 11 }, + ); + expect(await res.text()).toBe('data: "café"\n\n'); + }); +});