Resolve conflicts: - streamingHandler.js: adopt upstreamResponseHeaders while keeping 0-token detail row avoidance - capabilities.js: preserve user-asserted caps and globalThis slots without local caching of catalogSource - AddCustomModelModal.js & providers/[id]/page.js: wire STT transport marker with custom model edits/assertions - models/custom/route.js & aliasRepo.js: persist custom model transport and invalidate user caps - usageRepo.js: key byApiKey live stats by full API key and keep tail in maskApiKey - UsageStats.js: lazy load charts dynamically
252 lines
8.6 KiB
JavaScript
252 lines
8.6 KiB
JavaScript
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 <title>, 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 };
|
|
}
|