import { FORMATS } from "../../translator/formats.js"; import { needsTranslation } from "../../translator/index.js"; import { createSSETransformStreamWithLogger, createPassthroughStreamWithLogger, } from "../../utils/stream.js"; import { pipeWithDisconnect } from "../../utils/streamHandler.js"; import { PROVIDERS } from "../../config/providers.js"; import { HTTP_STATUS, STREAM_STALL_TIMEOUT_MS } from "../../config/runtimeConfig.js"; import { buildAbortedResponsesTerminalBytes } from "../../utils/responsesStreamHelpers.js"; import { buildStreamErrorBytes } from "../../utils/streamHelpers.js"; import { buildRequestDetail, extractRequestConfig, saveUsageStats, formatDoneLine, tokensForDetail, shouldPersistRequestDetail, } from "./requestDetail.js"; import { streamStatusForContent } from "../../utils/streamErrorPatterns.js"; import { saveRequestDetail } from "@/lib/usageDb.js"; import { SSE_HEADERS_CORS as SSE_HEADERS } from "../../utils/sseConstants.js"; import { upstreamResponseHeaders } from "../../utils/upstreamHeaders.js"; // Codex returns Responses API SSE → which client format to translate INTO, by request sourceFormat. // Gemini-family all map to ANTIGRAVITY decoder; unknown sources fall back to OPENAI. const CODEX_SOURCE_TO_TARGET = { [FORMATS.OPENAI_RESPONSES]: FORMATS.OPENAI_RESPONSES, [FORMATS.CLAUDE]: FORMATS.CLAUDE, [FORMATS.ANTIGRAVITY]: FORMATS.ANTIGRAVITY, [FORMATS.GEMINI]: FORMATS.ANTIGRAVITY, [FORMATS.GEMINI_CLI]: FORMATS.ANTIGRAVITY, }; /** * Determine which SSE transform stream to use based on provider/format. */ function buildTransformStream({ provider, sourceFormat, targetFormat, userAgent, reqLogger, toolNameMap, customToolNames, model, connectionId, body, onStreamComplete, apiKey, credentials }) { const isDroidCLI = userAgent?.toLowerCase().includes("droid") || userAgent?.toLowerCase().includes("codex-cli"); // Responses-API providers (e.g. codex) emit Responses SSE → translate into client format const isResponsesProvider = PROVIDERS[provider]?.format === FORMATS.OPENAI_RESPONSES; const needsCodexTranslation = isResponsesProvider && targetFormat === FORMATS.OPENAI_RESPONSES && !isDroidCLI; if (needsCodexTranslation) { const codexTarget = CODEX_SOURCE_TO_TARGET[sourceFormat] || FORMATS.OPENAI; return createSSETransformStreamWithLogger(FORMATS.OPENAI_RESPONSES, codexTarget, provider, reqLogger, toolNameMap, model, connectionId, body, onStreamComplete, apiKey, customToolNames, credentials); } if (needsTranslation(targetFormat, sourceFormat)) { return createSSETransformStreamWithLogger(targetFormat, sourceFormat, provider, reqLogger, toolNameMap, model, connectionId, body, onStreamComplete, apiKey, customToolNames, credentials); } return createPassthroughStreamWithLogger( provider, reqLogger, model, connectionId, body, onStreamComplete, apiKey, ); } /** * Handle streaming response — pipe provider SSE through transform stream to client. */ export async function handleStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, userAgent, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, customToolNames, streamController, onStreamComplete, pxpipe, reqTag, log, credentials }) { if (onRequestSuccess) { Promise.resolve() .then(onRequestSuccess) .catch(err => { console.error("[ChatCore] onRequestSuccess failed:", err?.message || err); }); } // When upstream returns HTML/text instead of SSE (e.g. Cloudflare 5xx error // page), piping it through the SSE transform stream causes Next.js // "failed to pipe response" and crashes the chat router. Read the body, // pull a short human-readable message from the , sanitize it, and // return a clean JSON error instead. The message is stripped of HTML tags // and clamped so untrusted upstream text never reaches the client verbatim // (the UI may render error.message as HTML). const upstreamContentType = ( providerResponse.headers.get("content-type") || "" ).toLowerCase(); if ( upstreamContentType && !upstreamContentType.includes("text/event-stream") && !upstreamContentType.includes("application/json") ) { const bodyText = await providerResponse.text().catch(() => ""); const titleMatch = bodyText.match(/<title>([^<]+)<\/title>/i); const sanitizedTitle = (titleMatch?.[1] || "") .replace(/<[^>]*>/g, "") .replace(/[\r\n]+/g, " ") .trim() .slice(0, 160); const shortMsg = sanitizedTitle || (bodyText.length < 200 ? bodyText .replace(/<[^>]*>/g, "") .trim() .slice(0, 160) : `Upstream returned non-SSE response (${upstreamContentType})`); const status = providerResponse.status || 502; if (log?.errorLine) log.errorLine( reqTag, "✗", `BLOCKED ${status} · ${provider}/${model} · non-SSE (${upstreamContentType})\n ${shortMsg}`, ); else console.warn( `[STREAM] ${provider} | ${model} | blocked pipe: ${shortMsg} [${status}]`, ); streamController?.handleError?.(new Error(`upstream non-SSE: ${status}`)); return { success: false, response: new Response( JSON.stringify({ error: { message: `[${status}]: ${shortMsg}` } }), { status, headers: { "Content-Type": "application/json", "Access-Control-Allow-Origin": "*", }, }, ), }; } const transformStream = buildTransformStream({ provider, sourceFormat, targetFormat, userAgent, reqLogger, toolNameMap, customToolNames, model, connectionId, body, onStreamComplete, apiKey, credentials }); // Terminal bytes when the stream aborts after HTTP 200 was already sent, so the // client sees a real error instead of a silently truncated stream. // Responses passthrough keeps its own response.failed shape; every other client // format gets the OpenAI error frame + [DONE], or `event: error` for Claude. const isResponsesPassthrough = sourceFormat === FORMATS.OPENAI_RESPONSES && targetFormat === FORMATS.OPENAI_RESPONSES; const onAbortTerminal = isResponsesPassthrough ? buildAbortedResponsesTerminalBytes : (message) => buildStreamErrorBytes( HTTP_STATUS.GATEWAY_TIMEOUT, message, sourceFormat, ); const stallTimeoutMs = PROVIDERS[provider]?.stallTimeoutMs || STREAM_STALL_TIMEOUT_MS; const transformedBody = pipeWithDisconnect( providerResponse, transformStream, streamController, onAbortTerminal, stallTimeoutMs, ); return { success: true, response: new Response(transformedBody, { headers: { ...SSE_HEADERS, ...upstreamResponseHeaders(providerResponse.headers), }, }), }; } /** * Build onStreamComplete callback for streaming usage tracking. */ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest, pxpipe, reqTag, log, streamErrorPatterns, persistUsage = "all", }) { const streamDetailId = `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`; const onStreamComplete = (contentObj, usage, ttftAt) => { const latency = { ttft: ttftAt ? ttftAt - requestStartTime : Date.now() - requestStartTime, total: Date.now() - requestStartTime, }; const safeContent = contentObj?.content || "[Empty streaming response]"; const safeThinking = contentObj?.thinking || null; const rawProviderText = typeof contentObj?.rawProviderText === "string" ? contentObj.rawProviderText : ""; if (shouldPersistRequestDetail(persistUsage, "success")) { saveRequestDetail( buildRequestDetail( { provider, model, connectionId, apiKey, latency, tokens: tokensForDetail(usage), request: extractRequestConfig(body, stream), providerRequest: finalBody || translatedBody || null, providerResponse: rawProviderText || safeContent, response: { content: safeContent, thinking: safeThinking, type: "streaming", }, pxpipe, status: streamStatusForContent( streamErrorPatterns?.[provider], safeContent, ), }, { id: streamDetailId }, ), ).catch((err) => { console.error( "[RequestDetail] Failed to update streaming content:", err.message, ); }); } // Persist stream usage to DB (no console line; the "📊 done" line below is authoritative) saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint, label: "STREAM USAGE", silent: true, }); if (log?.line) log.line(reqTag, "📊", formatDoneLine({ usage, latency })); }; return { onStreamComplete, streamDetailId }; }