2026-06-20 10:09:50 +07:00
import { detectFormat , getTargetFormat , resolveTransport } from "../services/provider.js" ;
2026-02-27 11:15:12 +07:00
import { translateRequest } from "../translator/index.js" ;
2026-08-05 11:39:07 +07:00
import { applyThinking , extractThinking , stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js" ;
2026-01-05 10:37:09 +07:00
import { FORMATS } from "../translator/formats.js" ;
2026-08-14 16:08:30 +07:00
import { normalizeClaudePassthrough , anchorClaudeCache } from "../translator/formats/claude.js" ;
2026-02-27 11:15:12 +07:00
import { createStreamController } from "../utils/streamHandler.js" ;
2026-01-11 21:45:01 +07:00
import { refreshWithRetry } from "../services/tokenRefresh.js" ;
2026-01-05 10:37:09 +07:00
import { createRequestLogger } from "../utils/requestLogger.js" ;
2026-08-14 16:51:48 +07:00
import { getModelTargetFormat , getModelSupportedFormats , getModelStrip , getModelUpstreamId , getModelType , PROVIDER_ID_TO_ALIAS } from "../config/providerModels.js" ;
2026-06-13 22:05:01 +07:00
import { PROVIDERS } from "../config/providers.js" ;
2026-01-05 10:37:09 +07:00
import { createErrorResult , parseUpstreamError , formatProviderError } from "../utils/error.js" ;
2026-07-16 11:18:27 +07:00
import { HTTP_STATUS , TOKEN_SAVER_HEADER } from "../config/runtimeConfig.js" ;
2026-01-05 10:37:09 +07:00
import { handleBypassRequest } from "../utils/bypassHandler.js" ;
2026-02-27 11:15:12 +07:00
import { trackPendingRequest , appendRequestLog , saveRequestDetail } from "@/lib/usageDb.js" ;
2026-01-11 21:45:01 +07:00
import { getExecutor } from "../executors/index.js" ;
2026-07-16 15:28:25 +07:00
import { supportsGrokCliReasoningEffort } from "../config/grokCli.js" ;
2026-02-27 11:15:12 +07:00
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" ;
2026-04-04 23:47:39 +07:00
import { detectClientTool , isNativePassthrough } from "../utils/clientDetector.js" ;
2026-05-12 09:19:50 +07:00
import { dedupeTools } from "../utils/toolDeduper.js" ;
2026-04-30 18:00:38 +07:00
import { injectCaveman } from "../rtk/caveman.js" ;
2026-06-20 10:09:50 +07:00
import { injectPonytail } from "../rtk/ponytail.js" ;
2026-04-30 18:00:38 +07:00
import { compressMessages , formatRtkLog } from "../rtk/index.js" ;
2026-06-26 11:08:33 +07:00
import { compressWithHeadroom , formatHeadroomLog , formatHeadroomSizeLog , isHeadroomPhantomSavings } from "../rtk/headroom.js" ;
2026-07-10 18:01:20 +07:00
import { compressWithPxpipe } from "../rtk/pxpipe.js" ;
2026-06-15 18:18:04 +07:00
import { getCapabilitiesForModel } from "../providers/capabilities.js" ;
import { stripUnsupportedModalities } from "../translator/concerns/modality.js" ;
import { prefetchRemoteImages } from "../translator/concerns/prefetch.js" ;
2026-08-28 11:07:07 +07:00
import { defaultClaudeToolType } from "../translator/concerns/toolCall.js" ;
2026-07-10 18:01:20 +07:00
import { resolveSessionId } from "../utils/sessionManager.js" ;
2026-02-08 16:45:31 +07:00
2026-01-05 10:37:09 +07:00
/**
* Core chat handler - shared between SSE and Worker
* @param {object} options.body - Request body
* @param {object} options.modelInfo - { provider, model }
* @param {object} options.credentials - Provider credentials
2026-02-27 11:15:12 +07:00
* @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses")
2026-01-05 10:37:09 +07:00
*/
2026-08-05 13:23:03 +07:00
/**
* Remove translator-internal continuity fields from the outbound upstream
* body. The Responses→Chat request translator stashes reasoning
* `encrypted_content` on assistant messages so a later openai→responses
* round-trip can restore the store=false continuity blob; that stash must
* never reach an upstream provider. Chat-native proxies reject the unknown
* assistant-message field and answer every turn with a literal "400" body
* (observed with multi-turn Codex sessions via OpenAI-compatible nodes).
*/
export function stripContinuityFields ( body ) {
if ( ! body || ! Array . isArray ( body . messages )) return body ;
for ( const msg of body . messages ) {
if ( msg && typeof msg === "object" ) {
delete msg . encrypted_content ;
delete msg . reasoning_encrypted_content ;
}
}
return body ;
}
2026-08-28 16:30:25 +07:00
export async function handleChatCore ({ body , modelInfo , credentials , log , onCredentialsRefreshed , onRequestSuccess , onDisconnect , clientRawRequest , connectionId , userAgent , apiKey , ccFilterNaming , rtkEnabled , headroomEnabled , headroomUrl , headroomCompressUserMessages , headroomTimeoutMs , cavemanEnabled , cavemanLevel , ponytailEnabled , ponytailLevel , pxpipeEnabled , pxpipeMinChars , pxpipeTimeoutMs , pxpipeTransform , onPxpipeEvent , sourceFormatOverride , providerThinking }) {
2026-01-05 10:37:09 +07:00
const { provider , model } = modelInfo ;
2026-02-09 11:30:42 +08:00
const requestStartTime = Date . now ();
2026-07-10 18:01:20 +07:00
// 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 () : "" );
2026-01-05 10:37:09 +07:00
2026-02-27 11:15:12 +07:00
const sourceFormat = sourceFormatOverride || detectFormat ( body );
2026-01-05 10:37:09 +07:00
2026-03-13 09:41:40 +07:00
// Check for bypass patterns (warmup, skip, cc naming)
const bypassResponse = handleBypassRequest ( body , model , userAgent , ccFilterNaming );
2026-02-27 11:15:12 +07:00
if ( bypassResponse ) return bypassResponse ;
2026-01-05 10:37:09 +07:00
const alias = PROVIDER_ID_TO_ALIAS [ provider ] || provider ;
const modelTargetFormat = getModelTargetFormat ( alias , model );
2026-08-14 16:51:48 +07:00
// Multi-endpoint providers: pick transport matching sourceFormat → zero translation.
// Per-model guard: only use the transport when the model declares support for that
// sourceFormat — opencode-go models differ in endpoint support (kimi/glm only do
// /chat/completions), so without this guard a claude-format request would wrongly
// route kimi to /messages.
const modelSupportedFormats = getModelSupportedFormats ( alias , model );
2026-06-20 10:09:50 +07:00
const runtimeTransport = resolveTransport ( provider , sourceFormat );
2026-08-14 16:51:48 +07:00
// Per-model guard: when a model declares supportedFormats, only use the
// sourceFormat-matched transport if that format is declared (opencode-go models
// differ — kimi/glm only do /chat/completions). Undeclared models keep the
// upstream default (use the transport), preserving behavior for glm/deepseek/...
const useTransport = ( ! modelSupportedFormats || modelSupportedFormats . includes ( sourceFormat )) ? runtimeTransport : null ;
2026-08-28 16:33:29 +07:00
// A source-format-matched endpoint keeps the request lossless. Prefer it
// over a model-level targetFormat, which is only the fallback for clients
// whose wire format has no supported transport (for example MiniMax-M3:
// OpenAI clients should stay on /chat/completions; other clients can fall
// back to its declared Claude target).
const targetFormat = useTransport ? . format || modelTargetFormat || getTargetFormat ( provider , credentials );
2026-08-14 16:51:48 +07:00
if ( useTransport && credentials ) credentials . runtimeTransport = useTransport ;
2026-04-07 10:06:55 +07:00
const stripList = getModelStrip ( alias , model );
2026-06-06 10:36:03 +07:00
const upstreamModel = getModelUpstreamId ( alias , model );
2026-02-02 09:17:15 +07:00
2026-04-13 12:04:57 +07:00
// 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 };
}
}
2026-02-22 21:44:11 +07:00
const clientRequestedStreaming = body . stream === true || sourceFormat === FORMATS . ANTIGRAVITY || sourceFormat === FORMATS . GEMINI || sourceFormat === FORMATS . GEMINI_CLI ;
2026-06-13 22:05:01 +07:00
const providerRequiresStreaming = PROVIDERS [ provider ] ? . forceStream === true ;
2026-03-13 10:00:47 +07:00
let stream = providerRequiresStreaming ? true : ( body . stream !== false );
2026-06-21 17:50:21 +07:00
// 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 ;
}
2026-05-13 21:10:42 +05:30
// 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 ;
2026-03-13 10:00:47 +07:00
// 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" );
2026-06-26 10:33:29 +07:00
if ( clientPrefersJson && ! clientPrefersSSE && body . stream !== true && ! providerRequiresStreaming ) {
2026-03-13 10:00:47 +07:00
stream = false ;
}
2026-01-05 10:37:09 +07:00
const reqLogger = await createRequestLogger ( sourceFormat , targetFormat , model );
2026-02-27 11:15:12 +07:00
if ( clientRawRequest ) reqLogger . logClientRawRequest ( clientRawRequest . endpoint , clientRawRequest . body , clientRawRequest . headers );
2026-01-05 10:37:09 +07:00
reqLogger . logRawRequest ( body );
log ? . debug ? .( "FORMAT" , ` ${ sourceFormat } → ${ targetFormat } | stream= ${ stream } ` );
2026-04-04 23:47:39 +07:00
// 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 );
2026-06-15 11:38:43 +07:00
// Expose raw client headers to translators/executors for session-id resolution
if ( credentials ) credentials . rawHeaders = clientRawRequest ? . headers || {};
2026-06-15 18:18:04 +07:00
// 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 } ` ); }
}
2026-04-04 23:47:39 +07:00
let translatedBody ;
let toolNameMap ;
2026-08-05 13:23:03 +07:00
let customToolNames ;
2026-04-04 23:47:39 +07:00
if ( passthrough ) {
log ? . debug ? .( "PASSTHROUGH" , ` ${ clientTool } → ${ provider } | native lossless` );
2026-07-07 16:29:11 +07:00
translatedBody = { ... body , model : stripThinkingSuffix ( upstreamModel ) };
2026-08-05 11:39:07 +07:00
if ( provider === "codex" ) {
const suffixThinking = {};
applyThinking ( sourceFormat , upstreamModel , suffixThinking , provider );
if ( suffixThinking . reasoning_effort ) {
const reasoning = translatedBody . reasoning ;
translatedBody . reasoning = {
...( reasoning && typeof reasoning === "object" && ! Array . isArray ( reasoning ) ? reasoning : {}),
effort : suffixThinking . reasoning_effort ,
};
delete translatedBody . reasoning_effort ;
}
}
2026-06-08 15:37:01 +07:00
// Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects
2026-07-07 16:29:11 +07:00
if ( clientTool === "claude" ) normalizeClaudePassthrough ( translatedBody , translatedBody . model );
2026-04-04 23:47:39 +07:00
} else {
2026-06-06 10:36:03 +07:00
translatedBody = translateRequest ( sourceFormat , targetFormat , upstreamModel , body , stream , credentials , provider , reqLogger , stripList , connectionId , clientTool );
2026-04-04 23:47:39 +07:00
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 ;
2026-08-05 13:23:03 +07:00
customToolNames = translatedBody . _customToolNames ;
delete translatedBody . _customToolNames ;
2026-07-07 16:29:11 +07:00
translatedBody . model = stripThinkingSuffix ( upstreamModel );
2026-08-05 13:23:03 +07:00
stripContinuityFields ( translatedBody );
2026-03-14 09:37:29 +07:00
}
2026-01-05 10:37:09 +07:00
2026-05-12 09:19:50 +07:00
// 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 ? "..." : "" } ` );
}
}
2026-04-30 18:00:38 +07:00
// 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 ;
2026-07-10 18:01:20 +07:00
// 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 } ` ;
2026-07-16 15:28:25 +07:00
const showThinking = provider !== "grok-cli" || supportsGrokCliReasoningEffort ( model );
const think = showThinking ? log . fmtThink ? .( extractThinking ( translatedBody )) : null ;
2026-07-10 18:01:20 +07:00
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 ( " · " ));
}
2026-06-06 11:31:52 +07:00
// 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 ;
}
2026-08-28 11:07:07 +07:00
// Claude tool schema requires `type` to be explicitly set; strict gateways (e.g., MiniMax)
// reject legacy payloads that omit it with HTTP 400. Default to "custom" when missing.
if ( finalFormat === FORMATS . CLAUDE && Array . isArray ( translatedBody . tools )) {
translatedBody . tools = defaultClaudeToolType ( translatedBody . tools );
}
2026-07-16 11:18:27 +07:00
// Per-request opt-out: client can bypass all token savers via header
const tokenSaverEnabled = clientRawRequest ? . headers ? .[ TOKEN_SAVER_HEADER ] ? . toLowerCase () !== "off" ;
2026-04-30 18:00:38 +07:00
// RTK: compress tool_result content
2026-07-16 11:18:27 +07:00
const rtkStats = compressMessages ( translatedBody , tokenSaverEnabled && rtkEnabled );
2026-04-30 18:00:38 +07:00
const rtkLine = formatRtkLog ( rtkStats );
if ( rtkLine ) console . log ( rtkLine );
2026-06-20 10:09:50 +07:00
// Headroom: optional external proxy compression; fail open if proxy is absent.
2026-06-26 11:08:33 +07:00
const headroomDiagnostics = {};
2026-08-28 16:30:25 +07:00
const headroomStats = await compressWithHeadroom ( translatedBody , { enabled : tokenSaverEnabled && headroomEnabled , url : headroomUrl , model : upstreamModel , format : finalFormat , compressUserMessages : headroomCompressUserMessages , timeoutMs : headroomTimeoutMs , diagnostics : headroomDiagnostics });
2026-06-20 10:09:50 +07:00
const headroomLine = formatHeadroomLog ( headroomStats );
2026-06-26 11:08:33 +07:00
const headroomSizeLine = formatHeadroomSizeLog ( headroomDiagnostics );
if ( headroomLine ) {
log ? . info ? .( "HEADROOM" , ` ${ headroomLine }${ headroomSizeLine ? ` | ${ headroomSizeLine } ` : "" } ` );
if ( isHeadroomPhantomSavings ( headroomStats , headroomDiagnostics )) {
2026-07-10 18:01:20 +07:00
log ? . warn ? .( "HEADROOM" , `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${ formatHeadroomSizeLog ( headroomDiagnostics ) } ` );
2026-06-26 11:08:33 +07:00
}
2026-07-16 11:18:27 +07:00
} else if ( tokenSaverEnabled && headroomEnabled ) log ? . warn ? .( "HEADROOM" , `skipped: ${ headroomDiagnostics . reason || "compression unavailable" }${ headroomDiagnostics . endpoint ? ` ( ${ headroomDiagnostics . endpoint } )` : "" } ` );
2026-06-20 10:09:50 +07:00
2026-07-16 18:13:51 +07:00
// Token-saver flags accumulator for the single "⚙" log line below.
const xf = [];
2026-04-30 18:00:38 +07:00
// Caveman: inject terse-style system prompt
2026-07-16 11:18:27 +07:00
if ( tokenSaverEnabled && cavemanEnabled && cavemanLevel ) {
2026-04-30 18:00:38 +07:00
injectCaveman ( translatedBody , finalFormat , cavemanLevel );
2026-07-10 18:01:20 +07:00
xf . push ( `CAVEMAN: ${ cavemanLevel } ` );
2026-04-30 18:00:38 +07:00
}
2026-06-20 10:09:50 +07:00
// Ponytail: inject lazy-senior-dev system prompt
2026-07-16 11:18:27 +07:00
if ( tokenSaverEnabled && ponytailEnabled && ponytailLevel ) {
2026-06-20 10:09:50 +07:00
injectPonytail ( translatedBody , finalFormat , ponytailLevel );
2026-07-10 18:01:20 +07:00
xf . push ( `PONYTAIL: ${ ponytailLevel } ` );
2026-06-20 10:09:50 +07:00
}
2026-07-10 16:10:16 +07:00
// 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 ;
2026-07-10 18:01:20 +07:00
if ( pxpipeSummary ? . applied ) xf . push ( `PXPIPE: ${ pxpipeSummary . imageCount } img` );
2026-07-10 16:10:16 +07:00
try { onPxpipeEvent ? .({ provider , model , ... pxpipeSummary }); } catch { /* stats must not break requests */ }
}
2026-07-10 18:01:20 +07:00
if ( xf . length && log ? . line ) log . line ( reqTag , "⚙" , xf . join ( " · " ));
2026-08-14 16:08:30 +07:00
// Pin cache breakpoints to the final body — every saver above can reshape
// system/tools/messages, and a stale anchor costs a full prefix rewrite.
if ( passthrough && clientTool === "claude" ) anchorClaudeCache ( translatedBody );
2026-01-11 21:45:01 +07:00
const executor = getExecutor ( provider );
2026-01-06 20:41:53 +02:00
trackPendingRequest ( model , provider , connectionId , true );
2026-05-13 21:10:42 +05:30
appendRequestLog ({ model , provider , connectionId , status : "PENDING" }). catch (() => { });
2026-01-06 20:41:53 +02:00
2026-02-27 11:15:12 +07:00
const msgCount = translatedBody . messages ? . length || translatedBody . input ? . length || translatedBody . contents ? . length || translatedBody . request ? . contents ? . length || 0 ;
2026-01-05 10:37:09 +07:00
log ? . debug ? .( "REQUEST" , ` ${ provider . toUpperCase () } | ${ model } | ${ msgCount } msgs` );
2026-02-27 11:15:12 +07:00
const streamController = createStreamController ({
2026-02-15 12:51:37 +08:00
onDisconnect : ( reason ) => {
trackPendingRequest ( model , provider , connectionId , false );
if ( onDisconnect ) onDisconnect ( reason );
},
2026-02-27 11:15:12 +07:00
onError : () => trackPendingRequest ( model , provider , connectionId , false ),
2026-07-10 18:01:20 +07:00
log , provider , model , reqTag
2026-02-15 12:51:37 +08:00
});
2026-01-05 10:37:09 +07:00
2026-03-09 15:46:06 +07:00
const proxyOptions = {
connectionProxyEnabled : credentials ? . providerSpecificData ? . connectionProxyEnabled === true ,
connectionProxyUrl : credentials ? . providerSpecificData ? . connectionProxyUrl || "" ,
connectionNoProxy : credentials ? . providerSpecificData ? . connectionNoProxy || "" ,
2026-04-13 10:08:24 +07:00
vercelRelayUrl : credentials ? . providerSpecificData ? . vercelRelayUrl || "" ,
2026-03-09 15:46:06 +07:00
};
2026-04-13 10:08:24 +07:00
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 ) {
2026-03-09 15:46:06 +07:00
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 } ` );
}
if ( proxyOptions . connectionProxyEnabled && proxyOptions . connectionNoProxy ) {
const connectionName = credentials ? . connectionName || credentials ? . connectionId || "unknown" ;
log ? . debug ? .( "PROXY" , ` ${ provider . toUpperCase () } | ${ model } | conn= ${ connectionName } | no_proxy= ${ proxyOptions . connectionNoProxy } ` );
}
2026-02-27 11:15:12 +07:00
// Execute request
let providerResponse , providerUrl , providerHeaders , finalBody ;
2026-07-20 15:39:17 +07:00
// Most executors return their registry format. Cursor AgentService is an
// exception: it is decoded by the executor into OpenAI-compatible output.
let providerResponseFormat = targetFormat ;
2026-01-05 10:37:09 +07:00
try {
2026-09-05 21:07:51 +07:00
const result = await executor . execute ({
model ,
body : translatedBody ,
stream ,
credentials ,
providerSessionId : sessionSeed ,
clientTool ,
signal : streamController . signal ,
log ,
proxyOptions ,
});
2026-01-11 21:45:01 +07:00
providerResponse = result . response ;
providerUrl = result . url ;
providerHeaders = result . headers ;
finalBody = result . transformedBody ;
2026-07-20 15:39:17 +07:00
providerResponseFormat = result . responseFormat || targetFormat ;
2026-01-27 10:49:16 +07:00
reqLogger . logTargetRequest ( providerUrl , providerHeaders , finalBody );
2026-01-05 10:37:09 +07:00
} catch ( error ) {
2026-02-21 16:42:46 +07:00
trackPendingRequest ( model , provider , connectionId , false , true );
2026-05-13 21:10:42 +05:30
appendRequestLog ({ model , provider , connectionId , status : `FAILED ${ error . name === "AbortError" ? 499 : HTTP_STATUS . BAD_GATEWAY } ` }). catch (() => { });
2026-02-27 11:15:12 +07:00
saveRequestDetail ( buildRequestDetail ({
provider , model , connectionId ,
2026-02-09 11:30:42 +08:00
latency : { ttft : 0 , total : Date . now () - requestStartTime },
tokens : { prompt_tokens : 0 , completion_tokens : 0 },
request : extractRequestConfig ( body , stream ),
providerRequest : translatedBody || null ,
2026-02-27 11:15:12 +07:00
response : { error : error . message || String ( error ), status : error . name === "AbortError" ? 499 : 502 , thinking : null },
2026-07-10 16:10:16 +07:00
pxpipe : pxpipeSummary ,
2026-02-09 11:30:42 +08:00
status : "error"
2026-05-13 21:10:42 +05:30
})). catch (() => { });
2026-02-09 11:30:42 +08:00
2026-01-05 10:37:09 +07:00
if ( error . name === "AbortError" ) {
streamController . handleError ( error );
return createErrorResult ( 499 , "Request aborted" );
}
2026-02-07 11:17:06 +07:00
const errMsg = formatProviderError ( error , provider , model , HTTP_STATUS . BAD_GATEWAY );
2026-07-10 18:01:20 +07:00
if ( log ? . errorLine ) {
log . errorLine ( reqTag , "✗" , `ERROR 502 · ${ provider } / ${ model } · ${ Date . now () - requestStartTime } ms\n ${ errMsg }${ error . stack ? `\n ${ error . stack } ` : "" } ` );
}
2026-02-07 11:17:06 +07:00
return createErrorResult ( HTTP_STATUS . BAD_GATEWAY , errMsg );
2026-01-05 10:37:09 +07:00
}
2026-04-14 10:14:50 +07:00
// Handle 401/403 - try token refresh (skip for noAuth providers)
if ( ! executor . noAuth && ( providerResponse . status === HTTP_STATUS . UNAUTHORIZED || providerResponse . status === HTTP_STATUS . FORBIDDEN )) {
2026-03-14 09:37:29 +07:00
try {
2026-07-25 17:30:11 +07:00
// 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 );
2026-03-14 09:37:29 +07:00
if ( newCredentials ? . accessToken || newCredentials ? . copilotToken ) {
2026-07-10 18:01:20 +07:00
if ( log ? . line ) log . line ( reqTag , "🔑" , `TOKEN REFRESHED · ${ provider } / ${ model } ` );
2026-03-14 09:37:29 +07:00
Object . assign ( credentials , newCredentials );
if ( onCredentialsRefreshed ) {
try { await onCredentialsRefreshed ( newCredentials ); } catch ( e ) { log ? . warn ? .( "TOKEN" , `onCredentialsRefreshed failed: ${ e . message } ` ); }
}
try {
2026-09-05 21:07:51 +07:00
const retryResult = await executor . execute ({
model ,
body : translatedBody ,
stream ,
credentials ,
providerSessionId : sessionSeed ,
clientTool ,
signal : streamController . signal ,
log ,
proxyOptions ,
});
2026-07-20 15:39:17 +07:00
if ( retryResult . response . ok ) {
providerResponse = retryResult . response ;
providerUrl = retryResult . url ;
providerResponseFormat = retryResult . responseFormat || targetFormat ;
}
2026-03-14 09:37:29 +07:00
} 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 } ` );
2026-01-05 10:37:09 +07:00
}
}
2026-02-27 11:15:12 +07:00
// Provider returned error
2026-01-05 10:37:09 +07:00
if ( ! providerResponse . ok ) {
2026-02-21 16:42:46 +07:00
trackPendingRequest ( model , provider , connectionId , false , true );
2026-04-24 11:36:16 +07:00
const { statusCode , message , resetsAtMs } = await parseUpstreamError ( providerResponse , executor );
2026-05-13 21:10:42 +05:30
appendRequestLog ({ model , provider , connectionId , status : `FAILED ${ statusCode } ` }). catch (() => { });
2026-02-27 11:15:12 +07:00
saveRequestDetail ( buildRequestDetail ({
provider , model , connectionId ,
2026-02-09 11:30:42 +08:00
latency : { ttft : 0 , total : Date . now () - requestStartTime },
tokens : { prompt_tokens : 0 , completion_tokens : 0 },
request : extractRequestConfig ( body , stream ),
providerRequest : finalBody || translatedBody || null ,
2026-02-27 11:15:12 +07:00
response : { error : message , status : statusCode , thinking : null },
2026-07-10 16:10:16 +07:00
pxpipe : pxpipeSummary ,
2026-02-09 11:30:42 +08:00
status : "error"
2026-05-13 21:10:42 +05:30
})). catch (() => { });
2026-02-09 11:30:42 +08:00
2026-01-05 10:57:45 +07:00
const errMsg = formatProviderError ( new Error ( message ), provider , model , statusCode );
2026-07-10 18:01:20 +07:00
if ( log ? . errorLine ) {
const urlStr = providerUrl ? `\n URL: ${ providerUrl } ` : "" ;
log . errorLine ( reqTag , "✗" , `ERROR ${ statusCode } · ${ provider } / ${ model } · ${ Date . now () - requestStartTime } ms ${ urlStr } \n ${ errMsg } ` );
}
2026-01-11 21:45:01 +07:00
reqLogger . logError ( new Error ( message ), finalBody || translatedBody );
2026-04-24 11:36:16 +07:00
return createErrorResult ( statusCode , errMsg , resetsAtMs );
2026-01-05 10:37:09 +07:00
}
2026-07-10 18:01:20 +07:00
const sharedCtx = { provider , model , body , stream , translatedBody , finalBody , requestStartTime , connectionId , apiKey , clientRawRequest , onRequestSuccess , pxpipe : pxpipeSummary , reqTag , log };
2026-05-13 21:10:42 +05:30
const appendLog = ( extra ) => appendRequestLog ({ model , provider , connectionId , ... extra }). catch (() => { });
2026-02-27 11:15:12 +07:00
const trackDone = () => trackPendingRequest ( model , provider , connectionId , false );
// Provider forced streaming but client wants JSON
2026-02-15 11:47:55 +07:00
if ( ! clientRequestedStreaming && providerRequiresStreaming ) {
2026-08-05 13:23:03 +07:00
const result = await handleForcedSSEToJson ({ ... sharedCtx , providerResponse , sourceFormat , targetFormat : providerResponseFormat , customToolNames , trackDone , appendLog });
2026-03-14 09:37:29 +07:00
if ( result ) { streamController . handleComplete (); return result ; }
2026-02-15 11:47:55 +07:00
}
2026-02-27 11:15:12 +07:00
// True non-streaming response
2026-01-05 10:37:09 +07:00
if ( ! stream ) {
2026-08-05 13:23:03 +07:00
const result = await handleNonStreamingResponse ({ ... sharedCtx , providerResponse , sourceFormat , targetFormat : providerResponseFormat , reqLogger , toolNameMap , customToolNames , trackDone , appendLog });
2026-03-14 09:37:29 +07:00
streamController . handleComplete ();
return result ;
2026-01-05 10:37:09 +07:00
}
// Streaming response
2026-07-03 15:07:20 +07:00
const { onStreamComplete , streamDetailId } = buildOnStreamComplete ({ ... sharedCtx });
2026-09-03 17:55:52 +07:00
return handleStreamingResponse ({ ... sharedCtx , providerResponse , sourceFormat , targetFormat : providerResponseFormat , userAgent , reqLogger , toolNameMap , customToolNames , streamController , onStreamComplete , streamDetailId , credentials });
2026-01-05 10:37:09 +07:00
}
export function isTokenExpiringSoon ( expiresAt , bufferMs = 5 * 60 * 1000 ) {
if ( ! expiresAt ) return false ;
2026-02-27 11:15:12 +07:00
return new Date ( expiresAt ). getTime () - Date . now () < bufferMs ;
2026-02-08 16:45:31 +07:00
}