feat(open-sse): early-peek stream error detection + fix UTF-8 loss in CommandCode peek

- 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.
This commit is contained in:
2026-08-04 23:33:12 +07:00
parent 008e0ef311
commit e438a03f96
6 changed files with 1301 additions and 583 deletions

View File

@@ -27,7 +27,7 @@ export class CommandCodeExecutor extends BaseExecutor {
super("commandcode", PROVIDERS.commandcode); super("commandcode", PROVIDERS.commandcode);
} }
transformRequest(model, body, stream, credentials) { transformRequest(_model, body, _stream, _credentials) {
body.stream = true; body.stream = true;
return body; return body;
} }
@@ -131,13 +131,17 @@ export async function peekForUpstreamError(
) { ) {
const reader = originalResponse.body.getReader(); const reader = originalResponse.body.getReader();
const decoder = new TextDecoder(); const decoder = new TextDecoder();
const encoder = new TextEncoder();
const abortController = new AbortController(); const abortController = new AbortController();
const forwardAbort = () => abortController.abort(signal?.reason); const forwardAbort = () => abortController.abort(signal?.reason);
if (signal?.aborted) abortController.abort(signal?.reason); if (signal?.aborted) abortController.abort(signal?.reason);
else if (signal) else if (signal)
signal.addEventListener("abort", forwardAbort, { once: true }); 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 peeked = "";
let errorEvent = null; let errorEvent = null;
let committed = false; let committed = false;
@@ -167,6 +171,7 @@ export async function peekForUpstreamError(
Math.max(deadline - Date.now(), 1), Math.max(deadline - Date.now(), 1),
); );
if (done) break; if (done) break;
rawChunks.push(value);
peeked += decoder.decode(value, { stream: true }); peeked += decoder.decode(value, { stream: true });
const lines = peeked.split("\n"); const lines = peeked.split("\n");
// The last segment may be a partial line — only parse complete ones. // 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. // 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 (signal) signal.removeEventListener("abort", forwardAbort);
if (errorEvent) { if (errorEvent) {
@@ -210,7 +211,9 @@ export async function peekForUpstreamError(
start(controller) { start(controller) {
(async () => { (async () => {
try { 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) { while (true) {
const { done, value } = await reader.read(); const { done, value } = await reader.read();
if (done) break; if (done) break;

View File

@@ -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 { translateRequest } from "../translator/index.js";
import { stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js"; import { stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js";
import { FORMATS } from "../translator/formats.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 { createStreamController } from "../utils/streamHandler.js";
import { refreshWithRetry } from "../services/tokenRefresh.js"; import { refreshWithRetry } from "../services/tokenRefresh.js";
import { createRequestLogger } from "../utils/requestLogger.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 { 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 { HTTP_STATUS, TOKEN_SAVER_HEADER } from "../config/runtimeConfig.js";
import { handleBypassRequest } from "../utils/bypassHandler.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 { getExecutor } from "../executors/index.js";
import { supportsGrokCliReasoningEffort } from "../config/grokCli.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 { handleForcedSSEToJson } from "./chatCore/sseToJsonHandler.js";
import { handleNonStreamingResponse } from "./chatCore/nonStreamingHandler.js"; import { handleNonStreamingResponse } from "./chatCore/nonStreamingHandler.js";
import { handleStreamingResponse, buildOnStreamComplete } from "./chatCore/streamingHandler.js"; import {
import { detectClientTool, isNativePassthrough } from "../utils/clientDetector.js"; 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 { dedupeTools } from "../utils/toolDeduper.js";
import { injectCaveman } from "../rtk/caveman.js"; import { injectCaveman } from "../rtk/caveman.js";
import { injectPonytail } from "../rtk/ponytail.js"; import { injectPonytail } from "../rtk/ponytail.js";
import { compressMessages, formatRtkLog } from "../rtk/index.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 { compressWithPxpipe } from "../rtk/pxpipe.js";
import { getCapabilitiesForModel } from "../providers/capabilities.js"; import { getCapabilitiesForModel } from "../providers/capabilities.js";
import { stripUnsupportedModalities } from "../translator/concerns/modality.js"; import { stripUnsupportedModalities } from "../translator/concerns/modality.js";
@@ -38,31 +71,76 @@ import { resolveSessionId } from "../utils/sessionManager.js";
* @param {object} options.credentials - Provider credentials * @param {object} options.credentials - Provider credentials
* @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses") * @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 }) { 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 { provider, model } = modelInfo;
const requestStartTime = Date.now(); const requestStartTime = Date.now();
// Stable per-session color so all lines of one CLI conversation share a tag // Stable per-session color so all lines of one CLI conversation share a tag
const sessionSeed = (() => { const sessionSeed = (() => {
try { try {
return resolveSessionId({ headers: clientRawRequest?.headers, body, connectionId, scope: provider }); return resolveSessionId({
headers: clientRawRequest?.headers,
body,
connectionId,
scope: provider,
});
} catch { } catch {
return connectionId || ""; return connectionId || "";
} }
})(); })();
const reqTag = log?.tagForSession ? log.tagForSession(sessionSeed) : (log?.nextTag ? log.nextTag() : ""); 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) // Check for bypass patterns (warmup, skip, cc naming)
const bypassResponse = handleBypassRequest(body, model, userAgent, ccFilterNaming); const bypassResponse = handleBypassRequest(
body,
model,
userAgent,
ccFilterNaming,
);
if (bypassResponse) return bypassResponse; if (bypassResponse) return bypassResponse;
const alias = PROVIDER_ID_TO_ALIAS[provider] || provider; const alias = PROVIDER_ID_TO_ALIAS[provider] || provider;
const modelTargetFormat = getModelTargetFormat(alias, model); const modelTargetFormat = getModelTargetFormat(alias, model);
// Multi-endpoint providers: pick transport matching sourceFormat → zero translation // Multi-endpoint providers: pick transport matching sourceFormat → zero translation
const runtimeTransport = resolveTransport(provider, sourceFormat); const runtimeTransport = resolveTransport(provider, sourceFormat);
const targetFormat = modelTargetFormat || runtimeTransport?.format || getTargetFormat(provider); const targetFormat =
if (runtimeTransport && credentials) credentials.runtimeTransport = runtimeTransport; modelTargetFormat || runtimeTransport?.format || getTargetFormat(provider);
if (runtimeTransport && credentials)
credentials.runtimeTransport = runtimeTransport;
const stripList = getModelStrip(alias, model); const stripList = getModelStrip(alias, model);
const upstreamModel = getModelUpstreamId(alias, model); const upstreamModel = getModelUpstreamId(alias, model);
@@ -80,14 +158,22 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
} }
} }
const clientRequestedStreaming = body.stream === true || sourceFormat === FORMATS.ANTIGRAVITY || sourceFormat === FORMATS.GEMINI || sourceFormat === FORMATS.GEMINI_CLI; const clientRequestedStreaming =
body.stream === true ||
sourceFormat === FORMATS.ANTIGRAVITY ||
sourceFormat === FORMATS.GEMINI ||
sourceFormat === FORMATS.GEMINI_CLI;
const providerRequiresStreaming = PROVIDERS[provider]?.forceStream === true; const providerRequiresStreaming = PROVIDERS[provider]?.forceStream === true;
let stream = providerRequiresStreaming ? true : (body.stream !== false); let stream = providerRequiresStreaming ? true : body.stream !== false;
// Image generation models require non-streaming (Google v1internal:generateContent) // Image generation models require non-streaming (Google v1internal:generateContent)
const modelType = getModelType(alias, model); const modelType = getModelType(alias, model);
const isImageGenModel = modelType === "imageGen" || /image|imagen|image-generation/i.test(model); const isImageGenModel =
if (isImageGenModel && (provider === "antigravity" || provider === "gemini-cli")) { modelType === "imageGen" || /image|imagen|image-generation/i.test(model);
if (
isImageGenModel &&
(provider === "antigravity" || provider === "gemini-cli")
) {
stream = false; stream = false;
} }
@@ -102,14 +188,31 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
const acceptHeader = clientRawRequest?.headers?.accept || ""; const acceptHeader = clientRawRequest?.headers?.accept || "";
const clientPrefersJson = acceptHeader.includes("application/json"); const clientPrefersJson = acceptHeader.includes("application/json");
const clientPrefersSSE = acceptHeader.includes("text/event-stream"); const clientPrefersSSE = acceptHeader.includes("text/event-stream");
if (clientPrefersJson && !clientPrefersSSE && body.stream !== true && !providerRequiresStreaming) { if (
clientPrefersJson &&
!clientPrefersSSE &&
body.stream !== true &&
!providerRequiresStreaming
) {
stream = false; stream = false;
} }
const reqLogger = await createRequestLogger(sourceFormat, targetFormat, model); const reqLogger = await createRequestLogger(
if (clientRawRequest) reqLogger.logClientRawRequest(clientRawRequest.endpoint, clientRawRequest.body, clientRawRequest.headers); sourceFormat,
targetFormat,
model,
);
if (clientRawRequest)
reqLogger.logClientRawRequest(
clientRawRequest.endpoint,
clientRawRequest.body,
clientRawRequest.headers,
);
reqLogger.logRawRequest(body); reqLogger.logRawRequest(body);
log?.debug?.("FORMAT", `${sourceFormat} → ${targetFormat} | stream=${stream}`); log?.debug?.(
"FORMAT",
`${sourceFormat} → ${targetFormat} | stream=${stream}`,
);
// Native passthrough: CLI tool and provider are the same ecosystem // Native passthrough: CLI tool and provider are the same ecosystem
// Skip all translation/normalization — only model and Bearer are swapped // Skip all translation/normalization — only model and Bearer are swapped
@@ -123,27 +226,57 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
if (!passthrough) { if (!passthrough) {
const caps = getCapabilitiesForModel(provider, model); const caps = getCapabilitiesForModel(provider, model);
if (stripUnsupportedModalities(body, sourceFormat, caps)) { if (stripUnsupportedModalities(body, sourceFormat, caps)) {
log?.debug?.("MODALITY", `stripped unsupported media for ${provider}/${model}`); log?.debug?.(
"MODALITY",
`stripped unsupported media for ${provider}/${model}`,
);
} }
// Convert remote image URLs to base64 for targets that can't fetch URLs. // Convert remote image URLs to base64 for targets that can't fetch URLs.
try { try {
const n = await prefetchRemoteImages(body, sourceFormat, targetFormat, { signal: undefined }); const n = await prefetchRemoteImages(body, sourceFormat, targetFormat, {
if (n > 0) log?.debug?.("MODALITY", `prefetched ${n} remote image(s) for ${targetFormat}`); signal: undefined,
} catch (e) { log?.warn?.("MODALITY", `image prefetch failed: ${e.message}`); } });
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 translatedBody;
let toolNameMap; let toolNameMap;
if (passthrough) { if (passthrough) {
log?.debug?.("PASSTHROUGH", `${clientTool} → ${provider} | native lossless`); log?.debug?.(
"PASSTHROUGH",
`${clientTool} → ${provider} | native lossless`,
);
translatedBody = { ...body, model: stripThinkingSuffix(upstreamModel) }; translatedBody = { ...body, model: stripThinkingSuffix(upstreamModel) };
// Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects // Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects
if (clientTool === "claude") normalizeClaudePassthrough(translatedBody, translatedBody.model); if (clientTool === "claude")
normalizeClaudePassthrough(translatedBody, translatedBody.model);
} else { } else {
translatedBody = translateRequest(sourceFormat, targetFormat, upstreamModel, body, stream, credentials, provider, reqLogger, stripList, connectionId, clientTool); translatedBody = translateRequest(
sourceFormat,
targetFormat,
upstreamModel,
body,
stream,
credentials,
provider,
reqLogger,
stripList,
connectionId,
clientTool,
);
if (!translatedBody) { if (!translatedBody) {
trackPendingRequest(model, provider, connectionId, false, true); trackPendingRequest(model, provider, connectionId, false, true);
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Failed to translate request for ${sourceFormat} → ${targetFormat}`); return createErrorResult(
HTTP_STATUS.BAD_REQUEST,
`Failed to translate request for ${sourceFormat} → ${targetFormat}`,
);
} }
toolNameMap = translatedBody._toolNameMap; toolNameMap = translatedBody._toolNameMap;
delete translatedBody._toolNameMap; delete translatedBody._toolNameMap;
@@ -155,7 +288,10 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
const { tools: deduped, stripped } = dedupeTools(translatedBody.tools); const { tools: deduped, stripped } = dedupeTools(translatedBody.tools);
if (stripped.length > 0) { if (stripped.length > 0) {
translatedBody.tools = deduped; translatedBody.tools = deduped;
log?.debug?.("TOOLDEDUP", `stripped ${stripped.length}: ${stripped.slice(0, 3).join(", ")}${stripped.length > 3 ? "..." : ""}`); log?.debug?.(
"TOOLDEDUP",
`stripped ${stripped.length}: ${stripped.slice(0, 3).join(", ")}${stripped.length > 3 ? "..." : ""}`,
);
} }
} }
@@ -166,12 +302,26 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
// Request line: one correlated summary (fmt + thinking + counts + account) // Request line: one correlated summary (fmt + thinking + counts + account)
if (log?.line) { if (log?.line) {
const clientModel = clientRawRequest?.body?.model || `${provider}/${model}`; 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 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 toolN = translatedBody.tools?.length || body.tools?.length || 0;
const fmtStr = passthrough ? `FMT: ${sourceFormat} (passthrough)` : `FMT: ${sourceFormat}→${targetFormat}`; const fmtStr = passthrough
const showThinking = provider !== "grok-cli" || supportsGrokCliReasoningEffort(model); ? `FMT: ${sourceFormat} (passthrough)`
const think = showThinking ? log.fmtThink?.(extractThinking(translatedBody)) : null; : `FMT: ${sourceFormat}→${targetFormat}`;
const acc = credentials?.connectionName || credentials?.connectionId?.slice(0, 8) || "-"; 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 = [ const parts = [
`POST ${clientModel} → ${provider}/${model}`, `POST ${clientModel} → ${provider}/${model}`,
fmtStr, fmtStr,
@@ -186,29 +336,52 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
// TTS models don't support tool messages/function calling // TTS models don't support tool messages/function calling
if (getModelType(alias, model) === "tts" && translatedBody.messages) { if (getModelType(alias, model) === "tts" && translatedBody.messages) {
translatedBody.messages = translatedBody.messages.filter(msg => msg.role !== "tool"); translatedBody.messages = translatedBody.messages.filter(
(msg) => msg.role !== "tool",
);
delete translatedBody.tools; delete translatedBody.tools;
} }
// Per-request opt-out: client can bypass all token savers via header // Per-request opt-out: client can bypass all token savers via header
const tokenSaverEnabled = clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off"; const tokenSaverEnabled =
clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off";
// RTK: compress tool_result content // RTK: compress tool_result content
const rtkStats = compressMessages(translatedBody, tokenSaverEnabled && rtkEnabled); const rtkStats = compressMessages(
translatedBody,
tokenSaverEnabled && rtkEnabled,
);
const rtkLine = formatRtkLog(rtkStats); const rtkLine = formatRtkLog(rtkStats);
if (rtkLine) console.log(rtkLine); if (rtkLine) console.log(rtkLine);
// Headroom: optional external proxy compression; fail open if proxy is absent. // Headroom: optional external proxy compression; fail open if proxy is absent.
const headroomDiagnostics = {}; const headroomDiagnostics = {};
const headroomStats = await compressWithHeadroom(translatedBody, { enabled: tokenSaverEnabled && headroomEnabled, url: headroomUrl, model: upstreamModel, format: finalFormat, compressUserMessages: headroomCompressUserMessages, diagnostics: headroomDiagnostics }); const headroomStats = await compressWithHeadroom(translatedBody, {
enabled: tokenSaverEnabled && headroomEnabled,
url: headroomUrl,
model: upstreamModel,
format: finalFormat,
compressUserMessages: headroomCompressUserMessages,
diagnostics: headroomDiagnostics,
});
const headroomLine = formatHeadroomLog(headroomStats); const headroomLine = formatHeadroomLog(headroomStats);
const headroomSizeLine = formatHeadroomSizeLog(headroomDiagnostics); const headroomSizeLine = formatHeadroomSizeLog(headroomDiagnostics);
if (headroomLine) { if (headroomLine) {
log?.info?.("HEADROOM", `${headroomLine}${headroomSizeLine ? ` | ${headroomSizeLine}` : ""}`); log?.info?.(
"HEADROOM",
`${headroomLine}${headroomSizeLine ? ` | ${headroomSizeLine}` : ""}`,
);
if (isHeadroomPhantomSavings(headroomStats, headroomDiagnostics)) { if (isHeadroomPhantomSavings(headroomStats, headroomDiagnostics)) {
log?.warn?.("HEADROOM", `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${formatHeadroomSizeLog(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})` : ""}`); } 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. // Token-saver flags accumulator for the single "⚙" log line below.
const xf = []; const xf = [];
@@ -229,23 +402,42 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
let pxpipeSummary = null; let pxpipeSummary = null;
if (pxpipeEnabled) { if (pxpipeEnabled) {
const pxpipeResult = await compressWithPxpipe(translatedBody, { const pxpipeResult = await compressWithPxpipe(translatedBody, {
enabled: true, format: finalFormat, model: upstreamModel, enabled: true,
minChars: pxpipeMinChars, timeoutMs: pxpipeTimeoutMs, transform: pxpipeTransform, format: finalFormat,
model: upstreamModel,
minChars: pxpipeMinChars,
timeoutMs: pxpipeTimeoutMs,
transform: pxpipeTransform,
}); });
pxpipeSummary = pxpipeResult.summary; pxpipeSummary = pxpipeResult.summary;
if (pxpipeResult.body) translatedBody = pxpipeResult.body; if (pxpipeResult.body) translatedBody = pxpipeResult.body;
if (pxpipeSummary?.applied) xf.push(`PXPIPE:${pxpipeSummary.imageCount}img`); if (pxpipeSummary?.applied)
try { onPxpipeEvent?.({ provider, model, ...pxpipeSummary }); } catch { /* stats must not break requests */ } 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); const executor = getExecutor(provider);
trackPendingRequest(model, provider, connectionId, true); trackPendingRequest(model, provider, connectionId, true);
appendRequestLog({ model, provider, connectionId, status: "PENDING" }).catch(() => { }); appendRequestLog({ model, provider, connectionId, status: "PENDING" }).catch(
() => {},
);
const msgCount = translatedBody.messages?.length || translatedBody.input?.length || translatedBody.contents?.length || translatedBody.request?.contents?.length || 0; const msgCount =
log?.debug?.("REQUEST", `${provider.toUpperCase()} | ${model} | ${msgCount} msgs`); 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({ const streamController = createStreamController({
onDisconnect: (reason) => { onDisconnect: (reason) => {
@@ -253,21 +445,35 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
if (onDisconnect) onDisconnect(reason); if (onDisconnect) onDisconnect(reason);
}, },
onError: () => trackPendingRequest(model, provider, connectionId, false), onError: () => trackPendingRequest(model, provider, connectionId, false),
log, provider, model, reqTag log,
provider,
model,
reqTag,
}); });
const proxyOptions = { const proxyOptions = {
connectionProxyEnabled: credentials?.providerSpecificData?.connectionProxyEnabled === true, connectionProxyEnabled:
connectionProxyUrl: credentials?.providerSpecificData?.connectionProxyUrl || "", credentials?.providerSpecificData?.connectionProxyEnabled === true,
connectionNoProxy: credentials?.providerSpecificData?.connectionNoProxy || "", connectionProxyUrl:
credentials?.providerSpecificData?.connectionProxyUrl || "",
connectionNoProxy:
credentials?.providerSpecificData?.connectionNoProxy || "",
vercelRelayUrl: credentials?.providerSpecificData?.vercelRelayUrl || "", vercelRelayUrl: credentials?.providerSpecificData?.vercelRelayUrl || "",
}; };
if (proxyOptions.vercelRelayUrl) { if (proxyOptions.vercelRelayUrl) {
const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown"; const connectionName =
const poolId = credentials?.providerSpecificData?.connectionProxyPoolId || "none"; credentials?.connectionName || credentials?.connectionId || "unknown";
log?.info?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | vercel-relay=${proxyOptions.vercelRelayUrl}`); const poolId =
} else if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionProxyUrl) { 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; let maskedProxyUrl = proxyOptions.connectionProxyUrl;
try { try {
const parsed = new URL(proxyOptions.connectionProxyUrl); const parsed = new URL(proxyOptions.connectionProxyUrl);
@@ -279,14 +485,23 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
// Keep raw if URL parsing fails // Keep raw if URL parsing fails
} }
const poolId = credentials?.providerSpecificData?.connectionProxyPoolId || "none"; const poolId =
const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown"; credentials?.providerSpecificData?.connectionProxyPoolId || "none";
log?.info?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | url=${maskedProxyUrl}`); const connectionName =
credentials?.connectionName || credentials?.connectionId || "unknown";
log?.info?.(
"PROXY",
`${provider.toUpperCase()} | ${model} | conn=${connectionName} | pool=${poolId} | url=${maskedProxyUrl}`,
);
} }
if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionNoProxy) { if (proxyOptions.connectionProxyEnabled && proxyOptions.connectionNoProxy) {
const connectionName = credentials?.connectionName || credentials?.connectionId || "unknown"; const connectionName =
log?.debug?.("PROXY", `${provider.toUpperCase()} | ${model} | conn=${connectionName} | no_proxy=${proxyOptions.connectionNoProxy}`); credentials?.connectionName || credentials?.connectionId || "unknown";
log?.debug?.(
"PROXY",
`${provider.toUpperCase()} | ${model} | conn=${connectionName} | no_proxy=${proxyOptions.connectionNoProxy}`,
);
} }
// Execute request // Execute request
@@ -295,7 +510,15 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
// exception: it is decoded by the executor into OpenAI-compatible output. // exception: it is decoded by the executor into OpenAI-compatible output.
let providerResponseFormat = targetFormat; let providerResponseFormat = targetFormat;
try { try {
const result = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions }); const result = await executor.execute({
model,
body: translatedBody,
stream,
credentials,
signal: streamController.signal,
log,
proxyOptions,
});
providerResponse = result.response; providerResponse = result.response;
providerUrl = result.url; providerUrl = result.url;
providerHeaders = result.headers; providerHeaders = result.headers;
@@ -304,111 +527,254 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody);
} catch (error) { } catch (error) {
trackPendingRequest(model, provider, connectionId, false, true); trackPendingRequest(model, provider, connectionId, false, true);
appendRequestLog({ model, provider, connectionId, status: `FAILED ${error.name === "AbortError" ? 499 : HTTP_STATUS.BAD_GATEWAY}` }).catch(() => { }); appendRequestLog({
saveRequestDetail(buildRequestDetail({ model,
provider, model, connectionId, 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 }, latency: { ttft: 0, total: Date.now() - requestStartTime },
tokens: { prompt_tokens: 0, completion_tokens: 0 }, tokens: { prompt_tokens: 0, completion_tokens: 0 },
request: extractRequestConfig(body, stream), request: extractRequestConfig(body, stream),
providerRequest: translatedBody || null, providerRequest: translatedBody || null,
response: { error: error.message || String(error), status: error.name === "AbortError" ? 499 : 502, thinking: null }, response: {
error: error.message || String(error),
status: error.name === "AbortError" ? 499 : 502,
thinking: null,
},
pxpipe: pxpipeSummary, pxpipe: pxpipeSummary,
status: "error" status: "error",
})).catch(() => { }); }),
).catch(() => {});
if (error.name === "AbortError") { if (error.name === "AbortError") {
streamController.handleError(error); streamController.handleError(error);
return createErrorResult(499, "Request aborted"); return createErrorResult(499, "Request aborted");
} }
const errMsg = formatProviderError(error, provider, model, HTTP_STATUS.BAD_GATEWAY); const errMsg = formatProviderError(
error,
provider,
model,
HTTP_STATUS.BAD_GATEWAY,
);
if (log?.errorLine) { if (log?.errorLine) {
log.errorLine(reqTag, "✗", `ERROR 502 · ${provider}/${model} · ${Date.now() - requestStartTime}ms\n ${errMsg}${error.stack ? `\n ${error.stack}` : ""}`); 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); return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg);
} }
// 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 },
);
}
// Handle 401/403 - try token refresh (skip for noAuth providers) // Handle 401/403 - try token refresh (skip for noAuth providers)
if (!executor.noAuth && (providerResponse.status === HTTP_STATUS.UNAUTHORIZED || providerResponse.status === HTTP_STATUS.FORBIDDEN)) { if (
!executor.noAuth &&
(providerResponse.status === HTTP_STATUS.UNAUTHORIZED ||
providerResponse.status === HTTP_STATUS.FORBIDDEN)
) {
try { try {
// Mutate credentials after each successful refresh: rotating refresh_token // Mutate credentials after each successful refresh: rotating refresh_token
// providers (xAI/grok-cli) issue a new RT on every refresh; without this, // providers (xAI/grok-cli) issue a new RT on every refresh; without this,
// refreshWithRetry's 2nd/3rd attempt reuses the already-consumed RT → // refreshWithRetry's 2nd/3rd attempt reuses the already-consumed RT →
// invalid_grant → auth_failed retryable=false. // invalid_grant → auth_failed retryable=false.
const newCredentials = await refreshWithRetry(async () => { const newCredentials = await refreshWithRetry(
async () => {
const result = await executor.refreshCredentials(credentials, log); const result = await executor.refreshCredentials(credentials, log);
if (result?.refreshToken && result.refreshToken !== credentials.refreshToken) { if (
if (result.accessToken) credentials.accessToken = result.accessToken; result?.refreshToken &&
result.refreshToken !== credentials.refreshToken
) {
if (result.accessToken)
credentials.accessToken = result.accessToken;
credentials.refreshToken = result.refreshToken; credentials.refreshToken = result.refreshToken;
} }
return result; return result;
}, 3, log); },
3,
log,
);
if (newCredentials?.accessToken || newCredentials?.copilotToken) { if (newCredentials?.accessToken || newCredentials?.copilotToken) {
if (log?.line) log.line(reqTag, "🔑", `TOKEN REFRESHED · ${provider}/${model}`); if (log?.line)
log.line(reqTag, "🔑", `TOKEN REFRESHED · ${provider}/${model}`);
Object.assign(credentials, newCredentials); Object.assign(credentials, newCredentials);
if (onCredentialsRefreshed) { if (onCredentialsRefreshed) {
try { await onCredentialsRefreshed(newCredentials); } catch (e) { log?.warn?.("TOKEN", `onCredentialsRefreshed failed: ${e.message}`); } try {
await onCredentialsRefreshed(newCredentials);
} catch (e) {
log?.warn?.("TOKEN", `onCredentialsRefreshed failed: ${e.message}`);
}
} }
try { try {
const retryResult = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions }); const retryResult = await executor.execute({
model,
body: translatedBody,
stream,
credentials,
signal: streamController.signal,
log,
proxyOptions,
});
if (retryResult.response.ok) { if (retryResult.response.ok) {
providerResponse = retryResult.response; providerResponse = retryResult.response;
providerUrl = retryResult.url; providerUrl = retryResult.url;
providerResponseFormat = retryResult.responseFormat || targetFormat; providerResponseFormat = retryResult.responseFormat || targetFormat;
} }
} catch { log?.warn?.("TOKEN", `${provider.toUpperCase()} | retry after refresh failed`); } } catch {
log?.warn?.(
"TOKEN",
`${provider.toUpperCase()} | retry after refresh failed`,
);
}
} else { } else {
log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`); log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`);
} }
} catch (e) { } catch (e) {
log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh threw: ${e.message}`); log?.warn?.(
"TOKEN",
`${provider.toUpperCase()} | refresh threw: ${e.message}`,
);
} }
} }
// Provider returned error // Provider returned error
if (!providerResponse.ok) { if (!providerResponse.ok) {
trackPendingRequest(model, provider, connectionId, false, true); trackPendingRequest(model, provider, connectionId, false, true);
const { statusCode, message, resetsAtMs } = await parseUpstreamError(providerResponse, executor); const { statusCode, message, resetsAtMs } = await parseUpstreamError(
appendRequestLog({ model, provider, connectionId, status: `FAILED ${statusCode}` }).catch(() => { }); providerResponse,
saveRequestDetail(buildRequestDetail({ executor,
provider, model, connectionId, );
appendRequestLog({
model,
provider,
connectionId,
status: `FAILED ${statusCode}`,
}).catch(() => {});
saveRequestDetail(
buildRequestDetail({
provider,
model,
connectionId,
latency: { ttft: 0, total: Date.now() - requestStartTime }, latency: { ttft: 0, total: Date.now() - requestStartTime },
tokens: { prompt_tokens: 0, completion_tokens: 0 }, tokens: { prompt_tokens: 0, completion_tokens: 0 },
request: extractRequestConfig(body, stream), request: extractRequestConfig(body, stream),
providerRequest: finalBody || translatedBody || null, providerRequest: finalBody || translatedBody || null,
response: { error: message, status: statusCode, thinking: null }, response: { error: message, status: statusCode, thinking: null },
pxpipe: pxpipeSummary, pxpipe: pxpipeSummary,
status: "error" status: "error",
})).catch(() => { }); }),
).catch(() => {});
const errMsg = formatProviderError(new Error(message), provider, model, statusCode); const errMsg = formatProviderError(
new Error(message),
provider,
model,
statusCode,
);
if (log?.errorLine) { if (log?.errorLine) {
const urlStr = providerUrl ? `\n URL: ${providerUrl}` : ""; const urlStr = providerUrl ? `\n URL: ${providerUrl}` : "";
log.errorLine(reqTag, "✗", `ERROR ${statusCode} · ${provider}/${model} · ${Date.now() - requestStartTime}ms${urlStr}\n ${errMsg}`); log.errorLine(
reqTag,
"✗",
`ERROR ${statusCode} · ${provider}/${model} · ${Date.now() - requestStartTime}ms${urlStr}\n ${errMsg}`,
);
} }
reqLogger.logError(new Error(message), finalBody || translatedBody); reqLogger.logError(new Error(message), finalBody || translatedBody);
return createErrorResult(statusCode, errMsg, resetsAtMs); return createErrorResult(statusCode, errMsg, resetsAtMs);
} }
const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, pxpipe: pxpipeSummary, reqTag, log }; const sharedCtx = {
const appendLog = (extra) => appendRequestLog({ model, provider, connectionId, ...extra }).catch(() => { }); provider,
const trackDone = () => trackPendingRequest(model, provider, connectionId, false); 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);
// Provider forced streaming but client wants JSON // Provider forced streaming but client wants JSON
if (!clientRequestedStreaming && providerRequiresStreaming) { if (!clientRequestedStreaming && providerRequiresStreaming) {
const result = await handleForcedSSEToJson({ ...sharedCtx, providerResponse, sourceFormat, trackDone, appendLog }); const result = await handleForcedSSEToJson({
if (result) { streamController.handleComplete(); return result; } ...sharedCtx,
providerResponse,
sourceFormat,
trackDone,
appendLog,
});
if (result) {
streamController.handleComplete();
return result;
}
} }
// True non-streaming response // True non-streaming response
if (!stream) { if (!stream) {
const result = await handleNonStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat: providerResponseFormat, reqLogger, toolNameMap, trackDone, appendLog }); const result = await handleNonStreamingResponse({
...sharedCtx,
providerResponse,
sourceFormat,
targetFormat: providerResponseFormat,
reqLogger,
toolNameMap,
trackDone,
appendLog,
});
streamController.handleComplete(); streamController.handleComplete();
return result; return result;
} }
// Streaming response // Streaming response
const { onStreamComplete, streamDetailId } = buildOnStreamComplete({ ...sharedCtx }); const { onStreamComplete, streamDetailId } = buildOnStreamComplete({
return handleStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat: providerResponseFormat, userAgent, reqLogger, toolNameMap, streamController, onStreamComplete, streamDetailId }); ...sharedCtx,
});
return handleStreamingResponse({
...sharedCtx,
providerResponse,
sourceFormat,
targetFormat: providerResponseFormat,
userAgent,
reqLogger,
toolNameMap,
streamController,
onStreamComplete,
streamDetailId,
});
} }
export function isTokenExpiringSoon(expiresAt, bufferMs = 5 * 60 * 1000) { export function isTokenExpiringSoon(expiresAt, bufferMs = 5 * 60 * 1000) {

View File

@@ -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,
});
}

View File

@@ -20,7 +20,10 @@ import { handleBypassRequest } from "open-sse/utils/bypassHandler.js";
import { HTTP_STATUS } from "open-sse/config/runtimeConfig.js"; import { HTTP_STATUS } from "open-sse/config/runtimeConfig.js";
import { detectFormatByEndpoint } from "open-sse/translator/formats.js"; import { detectFormatByEndpoint } from "open-sse/translator/formats.js";
import * as log from "../utils/logger.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"; import { getProjectIdForConnection } from "open-sse/services/projectId.js";
/** /**
@@ -43,7 +46,7 @@ export async function handleChat(request, clientRawRequest = null) {
clientRawRequest = { clientRawRequest = {
endpoint: url.pathname, endpoint: url.pathname,
body, body,
headers: Object.fromEntries(request.headers.entries()) headers: Object.fromEntries(request.headers.entries()),
}; };
} }
cacheClaudeHeaders(clientRawRequest.headers); cacheClaudeHeaders(clientRawRequest.headers);
@@ -83,7 +86,12 @@ export async function handleChat(request, clientRawRequest = null) {
// Bypass naming/warmup requests before combo rotation to avoid wasting rotation slots // Bypass naming/warmup requests before combo rotation to avoid wasting rotation slots
const userAgent = request?.headers?.get("user-agent") || ""; const userAgent = request?.headers?.get("user-agent") || "";
const bypassResponse = handleBypassRequest(body, modelStr, userAgent, !!settings.ccFilterNaming); const bypassResponse = handleBypassRequest(
body,
modelStr,
userAgent,
!!settings.ccFilterNaming,
);
if (bypassResponse) return bypassResponse.response || bypassResponse; if (bypassResponse) return bypassResponse.response || bypassResponse;
// Check if model is a combo (has multiple models with fallback) // Check if model is a combo (has multiple models with fallback)
@@ -92,17 +100,22 @@ export async function handleChat(request, clientRawRequest = null) {
// Check for combo-specific strategy first, fallback to global // Check for combo-specific strategy first, fallback to global
const comboStrategies = settings.comboStrategies || {}; const comboStrategies = settings.comboStrategies || {};
const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy; const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy;
const comboStrategy = comboSpecificStrategy || settings.comboStrategy || "fallback"; const comboStrategy =
comboSpecificStrategy || settings.comboStrategy || "fallback";
if (comboStrategy === "fusion") { if (comboStrategy === "fusion") {
log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`); log.info(
"CHAT",
`Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`,
);
return handleFusionChat({ return handleFusionChat({
body, body,
models: comboModels, models: comboModels,
handleSingleModel: (b, m, isPanel) => { handleSingleModel: (b, m, isPanel) => {
let cleanRawReq = clientRawRequest; let cleanRawReq = clientRawRequest;
if (isPanel && clientRawRequest) { if (isPanel && clientRawRequest) {
const { tools, tool_choice, ...cleanBody } = clientRawRequest.body || {}; const { tools, tool_choice, ...cleanBody } =
clientRawRequest.body || {};
cleanRawReq = { ...clientRawRequest, body: cleanBody }; cleanRawReq = { ...clientRawRequest, body: cleanBody };
} }
return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); return handleSingleModelChat(b, m, cleanRawReq, request, apiKey);
@@ -115,26 +128,42 @@ export async function handleChat(request, clientRawRequest = null) {
} }
const comboStickyLimit = settings.comboStickyRoundRobinLimit; const comboStickyLimit = settings.comboStickyRoundRobinLimit;
log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`); log.info(
"CHAT",
`Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`,
);
return handleComboChat({ return handleComboChat({
body, body,
models: comboModels, models: comboModels,
handleSingleModel: (b, m) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey), handleSingleModel: (b, m) =>
handleSingleModelChat(b, m, clientRawRequest, request, apiKey),
log, log,
comboName: modelStr, comboName: modelStr,
comboStrategy, comboStrategy,
comboStickyLimit comboStickyLimit,
}); });
} }
// Single model request // Single model request
return handleSingleModelChat(body, modelStr, clientRawRequest, request, apiKey); return handleSingleModelChat(
body,
modelStr,
clientRawRequest,
request,
apiKey,
);
} }
/** /**
* Handle single model chat request * 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,
) {
const modelInfo = await getModelInfo(modelStr); const modelInfo = await getModelInfo(modelStr);
// If provider is null, this might be a combo name - check and handle // If provider is null, this might be a combo name - check and handle
@@ -145,17 +174,22 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
// Check for combo-specific strategy first, fallback to global // Check for combo-specific strategy first, fallback to global
const comboStrategies = chatSettings.comboStrategies || {}; const comboStrategies = chatSettings.comboStrategies || {};
const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy; const comboSpecificStrategy = comboStrategies[modelStr]?.fallbackStrategy;
const comboStrategy = comboSpecificStrategy || chatSettings.comboStrategy || "fallback"; const comboStrategy =
comboSpecificStrategy || chatSettings.comboStrategy || "fallback";
if (comboStrategy === "fusion") { if (comboStrategy === "fusion") {
log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`); log.info(
"CHAT",
`Combo "${modelStr}" with ${comboModels.length} models (strategy: fusion)`,
);
return handleFusionChat({ return handleFusionChat({
body, body,
models: comboModels, models: comboModels,
handleSingleModel: (b, m, isPanel) => { handleSingleModel: (b, m, isPanel) => {
let cleanRawReq = clientRawRequest; let cleanRawReq = clientRawRequest;
if (isPanel && clientRawRequest) { if (isPanel && clientRawRequest) {
const { tools, tool_choice, ...cleanBody } = clientRawRequest.body || {}; const { tools, tool_choice, ...cleanBody } =
clientRawRequest.body || {};
cleanRawReq = { ...clientRawRequest, body: cleanBody }; cleanRawReq = { ...clientRawRequest, body: cleanBody };
} }
return handleSingleModelChat(b, m, cleanRawReq, request, apiKey); return handleSingleModelChat(b, m, cleanRawReq, request, apiKey);
@@ -168,15 +202,19 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
} }
const comboStickyLimit = chatSettings.comboStickyRoundRobinLimit; const comboStickyLimit = chatSettings.comboStickyRoundRobinLimit;
log.info("CHAT", `Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`); log.info(
"CHAT",
`Combo "${modelStr}" with ${comboModels.length} models (strategy: ${comboStrategy}, sticky: ${comboStickyLimit})`,
);
return handleComboChat({ return handleComboChat({
body, body,
models: comboModels, models: comboModels,
handleSingleModel: (b, m) => handleSingleModelChat(b, m, clientRawRequest, request, apiKey), handleSingleModel: (b, m) =>
handleSingleModelChat(b, m, clientRawRequest, request, apiKey),
log, log,
comboName: modelStr, comboName: modelStr,
comboStrategy, comboStrategy,
comboStickyLimit comboStickyLimit,
}); });
} }
log.warn("CHAT", "Invalid model format", { model: modelStr }); log.warn("CHAT", "Invalid model format", { model: modelStr });
@@ -190,7 +228,8 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
// Extract userAgent from request // Extract userAgent from request
const userAgent = request?.headers?.get("user-agent") || ""; const userAgent = request?.headers?.get("user-agent") || "";
// Optional pin to a specific connection (dashboard test / client override) // Optional pin to a specific connection (dashboard test / client override)
const preferredConnectionId = request?.headers?.get("x-connection-id") || null; const preferredConnectionId =
request?.headers?.get("x-connection-id") || null;
// Try with available accounts (fallback on errors unless pinned) // Try with available accounts (fallback on errors unless pinned)
const excludeConnectionIds = new Set(); const excludeConnectionIds = new Set();
@@ -198,40 +237,75 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
let lastStatus = null; let lastStatus = null;
while (true) { while (true) {
const credentials = await getProviderCredentials(provider, excludeConnectionIds, model, { preferredConnectionId }); const credentials = await getProviderCredentials(
provider,
excludeConnectionIds,
model,
{ preferredConnectionId },
);
// All accounts unavailable // All accounts unavailable
if (!credentials || credentials.allRateLimited) { if (!credentials || credentials.allRateLimited) {
if (credentials?.allRateLimited) { if (credentials?.allRateLimited) {
const errorMsg = lastError || credentials.lastError || "Unavailable"; const errorMsg = lastError || credentials.lastError || "Unavailable";
const status = lastStatus || Number(credentials.lastErrorCode) || HTTP_STATUS.SERVICE_UNAVAILABLE; const status =
log.warn("CHAT", `[${provider}/${model}] ${errorMsg} (${credentials.retryAfterHuman})`); lastStatus ||
return unavailableResponse(status, `[${provider}/${model}] ${errorMsg}`, credentials.retryAfter, credentials.retryAfterHuman); 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) { if (excludeConnectionIds.size === 0) {
log.warn("AUTH", `No active credentials for provider: ${provider}`); log.warn("AUTH", `No active credentials for provider: ${provider}`);
return errorResponse(HTTP_STATUS.NOT_FOUND, `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 }); log.warn("CHAT", "No more accounts available", { provider });
return errorResponse(lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE, lastError || "All accounts unavailable"); return errorResponse(
lastStatus || HTTP_STATUS.SERVICE_UNAVAILABLE,
lastError || "All accounts unavailable",
);
} }
// Account selection shown in the unified "▶" line (acc:...) // Account selection shown in the unified "▶" line (acc:...)
const refreshedCredentials = await checkAndRefreshToken(provider, credentials); const refreshedCredentials = await checkAndRefreshToken(
provider,
credentials,
);
// Ensure real project ID is available for providers that need it (P0 fix: cold miss) // Ensure real project ID is available for providers that need it (P0 fix: cold miss)
if ((provider === "antigravity" || provider === "gemini-cli") && !refreshedCredentials.projectId) { if (
const pid = await getProjectIdForConnection(credentials.connectionId, refreshedCredentials.accessToken, provider); (provider === "antigravity" || provider === "gemini-cli") &&
!refreshedCredentials.projectId
) {
const pid = await getProjectIdForConnection(
credentials.connectionId,
refreshedCredentials.accessToken,
provider,
);
if (pid) { if (pid) {
refreshedCredentials.projectId = pid; refreshedCredentials.projectId = pid;
// Persist to DB in background so subsequent requests have it immediately // Persist to DB in background so subsequent requests have it immediately
updateProviderCredentials(credentials.connectionId, { projectId: pid }).catch(() => { }); updateProviderCredentials(credentials.connectionId, {
projectId: pid,
}).catch(() => {});
} }
} }
// Use shared chatCore // Use shared chatCore
const chatSettings = await getSettings(); const chatSettings = await getSettings();
const providerThinking = (chatSettings.providerThinking || {})[provider] || null; const providerThinking =
(chatSettings.providerThinking || {})[provider] || null;
const result = await handleChatCore({ const result = await handleChatCore({
body: { ...body, model: `${provider}/${model}` }, body: { ...body, model: `${provider}/${model}` },
modelInfo: { provider, model }, modelInfo: { provider, model },
@@ -254,35 +328,53 @@ async function handleSingleModelChat(body, modelStr, clientRawRequest = null, re
pxpipeMinChars: chatSettings.pxpipeMinChars, pxpipeMinChars: chatSettings.pxpipeMinChars,
pxpipeTimeoutMs: chatSettings.pxpipeTimeoutMs, pxpipeTimeoutMs: chatSettings.pxpipeTimeoutMs,
// Lazily warms the in-process module on first use; null when not installed (fail-open) // Lazily warms the in-process module on first use; null when not installed (fail-open)
pxpipeTransform: chatSettings.pxpipeEnabled ? await getPxpipeTransform() : null, pxpipeTransform: chatSettings.pxpipeEnabled
? await getPxpipeTransform()
: null,
onPxpipeEvent: appendPxpipeEvent, onPxpipeEvent: appendPxpipeEvent,
providerThinking, providerThinking,
streamErrorPatterns: chatSettings.streamErrorPatterns || {},
// Detect source format by endpoint + body // Detect source format by endpoint + body
sourceFormatOverride: request?.url ? detectFormatByEndpoint(new URL(request.url).pathname, body) : null, sourceFormatOverride: request?.url
? detectFormatByEndpoint(new URL(request.url).pathname, body)
: null,
onCredentialsRefreshed: async (newCreds) => { onCredentialsRefreshed: async (newCreds) => {
await updateProviderCredentials(credentials.connectionId, { await updateProviderCredentials(credentials.connectionId, {
...newCreds, ...newCreds,
existingProviderSpecificData: credentials.providerSpecificData, existingProviderSpecificData: credentials.providerSpecificData,
testStatus: "active" testStatus: "active",
}); });
}, },
onRequestSuccess: async () => { onRequestSuccess: async () => {
await clearAccountError(credentials.connectionId, credentials, model); 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) // 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); const { shouldFallback } = await markAccountUnavailable(
credentials.connectionId,
result.status,
result.error,
provider,
model,
result.resetsAtMs,
);
if (shouldFallback) { if (shouldFallback) {
// When a connection is explicitly pinned, never rotate to another account. // When a connection is explicitly pinned, never rotate to another account.
if (preferredConnectionId) { if (preferredConnectionId) {
log.warn("AUTH", `Pinned account ${credentials.connectionName} unavailable (${result.status}), no fallback`); log.warn(
"AUTH",
`Pinned account ${credentials.connectionName} unavailable (${result.status}), no fallback`,
);
return result.response; return result.response;
} }
log.warn("FALLBACK", `⇄ ACC:${credentials.connectionName} UNAVAILABLE (${result.status}) → NEXT ACCOUNT`); log.warn(
"FALLBACK",
`⇄ ACC:${credentials.connectionName} UNAVAILABLE (${result.status}) → NEXT ACCOUNT`,
);
excludeConnectionIds.add(credentials.connectionId); excludeConnectionIds.add(credentials.connectionId);
lastError = result.error; lastError = result.error;
lastStatus = result.status; lastStatus = result.status;

View File

@@ -98,4 +98,26 @@ describe("commandcode executor — early-error peek", () => {
expect(res.status).toBe(200); expect(res.status).toBe(200);
await res.body.cancel(); 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");
});
}); });

View File

@@ -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');
});
});