merge origin/master into gitea/new_feature
Bring local branch up to v0.5.35 while keeping xAI image/edit, SuperGrok quota tracking, per-provider timeouts, and pinned model-test actions.
This commit is contained in:
@@ -1,18 +1,19 @@
|
||||
import { detectFormat, getTargetFormat, resolveTransport } from "../services/provider.js";
|
||||
import { translateRequest } from "../translator/index.js";
|
||||
import { stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js";
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { normalizeClaudePassthrough } from "../translator/formats/claude.js";
|
||||
import { COLORS } from "../utils/stream.js";
|
||||
import { createStreamController } from "../utils/streamHandler.js";
|
||||
import { refreshWithRetry } from "../services/tokenRefresh.js";
|
||||
import { createRequestLogger } from "../utils/requestLogger.js";
|
||||
import { getModelTargetFormat, getModelStrip, getModelUpstreamId, getModelType, PROVIDER_ID_TO_ALIAS } from "../config/providerModels.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { createErrorResult, parseUpstreamError, formatProviderError } from "../utils/error.js";
|
||||
import { HTTP_STATUS } from "../config/runtimeConfig.js";
|
||||
import { HTTP_STATUS, TOKEN_SAVER_HEADER } from "../config/runtimeConfig.js";
|
||||
import { handleBypassRequest } from "../utils/bypassHandler.js";
|
||||
import { trackPendingRequest, appendRequestLog, saveRequestDetail } from "@/lib/usageDb.js";
|
||||
import { getExecutor } from "../executors/index.js";
|
||||
import { supportsGrokCliReasoningEffort } from "../config/grokCli.js";
|
||||
import { buildRequestDetail, extractRequestConfig } from "./chatCore/requestDetail.js";
|
||||
import { handleForcedSSEToJson } from "./chatCore/sseToJsonHandler.js";
|
||||
import { handleNonStreamingResponse } from "./chatCore/nonStreamingHandler.js";
|
||||
@@ -23,9 +24,12 @@ import { injectCaveman } from "../rtk/caveman.js";
|
||||
import { injectPonytail } from "../rtk/ponytail.js";
|
||||
import { compressMessages, formatRtkLog } from "../rtk/index.js";
|
||||
import { compressWithHeadroom, formatHeadroomLog, formatHeadroomSizeLog, isHeadroomPhantomSavings } from "../rtk/headroom.js";
|
||||
import { compressWithPxpipe } from "../rtk/pxpipe.js";
|
||||
import { getCapabilitiesForModel } from "../providers/capabilities.js";
|
||||
import { stripUnsupportedModalities } from "../translator/concerns/modality.js";
|
||||
import { prefetchRemoteImages } from "../translator/concerns/prefetch.js";
|
||||
import { extractThinking } from "../translator/concerns/thinkingUnified.js";
|
||||
import { resolveSessionId } from "../utils/sessionManager.js";
|
||||
|
||||
/**
|
||||
* Core chat handler - shared between SSE and Worker
|
||||
@@ -34,9 +38,18 @@ import { prefetchRemoteImages } from "../translator/concerns/prefetch.js";
|
||||
* @param {object} options.credentials - Provider credentials
|
||||
* @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses")
|
||||
*/
|
||||
export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, headroomEnabled, headroomUrl, headroomCompressUserMessages, cavemanEnabled, cavemanLevel, ponytailEnabled, ponytailLevel, 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 }) {
|
||||
const { provider, model } = modelInfo;
|
||||
const requestStartTime = Date.now();
|
||||
// Stable per-session color so all lines of one CLI conversation share a tag
|
||||
const sessionSeed = (() => {
|
||||
try {
|
||||
return resolveSessionId({ headers: clientRawRequest?.headers, body, connectionId, scope: provider });
|
||||
} catch {
|
||||
return connectionId || "";
|
||||
}
|
||||
})();
|
||||
const reqTag = log?.tagForSession ? log.tagForSession(sessionSeed) : (log?.nextTag ? log.nextTag() : "");
|
||||
|
||||
const sourceFormat = sourceFormatOverride || detectFormat(body);
|
||||
|
||||
@@ -123,9 +136,9 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
let toolNameMap;
|
||||
if (passthrough) {
|
||||
log?.debug?.("PASSTHROUGH", `${clientTool} → ${provider} | native lossless`);
|
||||
translatedBody = { ...body, model: upstreamModel };
|
||||
translatedBody = { ...body, model: stripThinkingSuffix(upstreamModel) };
|
||||
// Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects
|
||||
if (clientTool === "claude") normalizeClaudePassthrough(translatedBody, upstreamModel);
|
||||
if (clientTool === "claude") normalizeClaudePassthrough(translatedBody, translatedBody.model);
|
||||
} else {
|
||||
translatedBody = translateRequest(sourceFormat, targetFormat, upstreamModel, body, stream, credentials, provider, reqLogger, stripList, connectionId, clientTool);
|
||||
if (!translatedBody) {
|
||||
@@ -134,7 +147,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
}
|
||||
toolNameMap = translatedBody._toolNameMap;
|
||||
delete translatedBody._toolNameMap;
|
||||
translatedBody.model = upstreamModel;
|
||||
translatedBody.model = stripThinkingSuffix(upstreamModel);
|
||||
}
|
||||
|
||||
// Dedupe duplicate built-in tools when equivalent MCP tools are present (Claude clients only).
|
||||
@@ -150,41 +163,83 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
// Covers both passthrough (source shape) and translated (target shape) flows
|
||||
const finalFormat = passthrough ? sourceFormat : targetFormat;
|
||||
|
||||
// Request line: one correlated summary (fmt + thinking + counts + account)
|
||||
if (log?.line) {
|
||||
const clientModel = clientRawRequest?.body?.model || `${provider}/${model}`;
|
||||
const msgN = translatedBody.messages?.length || translatedBody.input?.length || translatedBody.contents?.length || body.messages?.length || body.input?.length || 0;
|
||||
const toolN = translatedBody.tools?.length || body.tools?.length || 0;
|
||||
const fmtStr = passthrough ? `FMT: ${sourceFormat} (passthrough)` : `FMT: ${sourceFormat}→${targetFormat}`;
|
||||
const showThinking = provider !== "grok-cli" || supportsGrokCliReasoningEffort(model);
|
||||
const think = showThinking ? log.fmtThink?.(extractThinking(translatedBody)) : null;
|
||||
const acc = credentials?.connectionName || credentials?.connectionId?.slice(0, 8) || "-";
|
||||
const parts = [
|
||||
`POST ${clientModel} → ${provider}/${model}`,
|
||||
fmtStr,
|
||||
stream ? "STREAM" : "JSON",
|
||||
`${msgN} MSG`,
|
||||
];
|
||||
if (toolN) parts.push(`${toolN} TOOL`);
|
||||
if (think) parts.push(`THINK:${think}`);
|
||||
parts.push(`ACC:${acc}`);
|
||||
log.line(reqTag, "▶", parts.join(" · "));
|
||||
}
|
||||
|
||||
// TTS models don't support tool messages/function calling
|
||||
if (getModelType(alias, model) === "tts" && translatedBody.messages) {
|
||||
translatedBody.messages = translatedBody.messages.filter(msg => msg.role !== "tool");
|
||||
delete translatedBody.tools;
|
||||
}
|
||||
|
||||
// Per-request opt-out: client can bypass all token savers via header
|
||||
const tokenSaverEnabled = clientRawRequest?.headers?.[TOKEN_SAVER_HEADER]?.toLowerCase() !== "off";
|
||||
|
||||
// RTK: compress tool_result content
|
||||
const rtkStats = compressMessages(translatedBody, rtkEnabled);
|
||||
const rtkStats = compressMessages(translatedBody, tokenSaverEnabled && rtkEnabled);
|
||||
const rtkLine = formatRtkLog(rtkStats);
|
||||
if (rtkLine) console.log(rtkLine);
|
||||
|
||||
// Headroom: optional external proxy compression; fail open if proxy is absent.
|
||||
const headroomDiagnostics = {};
|
||||
const headroomStats = await compressWithHeadroom(translatedBody, { enabled: 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 headroomSizeLine = formatHeadroomSizeLog(headroomDiagnostics);
|
||||
if (headroomLine) {
|
||||
log?.info?.("HEADROOM", `${headroomLine}${headroomSizeLine ? ` | ${headroomSizeLine}` : ""}`);
|
||||
if (isHeadroomPhantomSavings(headroomStats, headroomDiagnostics)) {
|
||||
log?.warn?.("HEADROOM", `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${headroomSizeLine}`);
|
||||
log?.warn?.("HEADROOM", `reported token delta, but outbound JSON shrank <5%; provider may bill near-original payload | ${formatHeadroomSizeLog(headroomDiagnostics)}`);
|
||||
}
|
||||
} else if (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.
|
||||
const xf = [];
|
||||
|
||||
// Caveman: inject terse-style system prompt
|
||||
if (cavemanEnabled && cavemanLevel) {
|
||||
if (tokenSaverEnabled && cavemanEnabled && cavemanLevel) {
|
||||
injectCaveman(translatedBody, finalFormat, cavemanLevel);
|
||||
log?.debug?.("CAVEMAN", `${cavemanLevel} | ${finalFormat}`);
|
||||
xf.push(`CAVEMAN:${cavemanLevel}`);
|
||||
}
|
||||
|
||||
// Ponytail: inject lazy-senior-dev system prompt
|
||||
if (ponytailEnabled && ponytailLevel) {
|
||||
if (tokenSaverEnabled && ponytailEnabled && ponytailLevel) {
|
||||
injectPonytail(translatedBody, finalFormat, ponytailLevel);
|
||||
log?.debug?.("PONYTAIL", `${ponytailLevel} | ${finalFormat}`);
|
||||
xf.push(`PONYTAIL:${ponytailLevel}`);
|
||||
}
|
||||
|
||||
// PXPIPE: image bulky context (Claude-format bodies only), last saver before dispatch
|
||||
let pxpipeSummary = null;
|
||||
if (pxpipeEnabled) {
|
||||
const pxpipeResult = await compressWithPxpipe(translatedBody, {
|
||||
enabled: true, format: finalFormat, model: upstreamModel,
|
||||
minChars: pxpipeMinChars, timeoutMs: pxpipeTimeoutMs, transform: pxpipeTransform,
|
||||
});
|
||||
pxpipeSummary = pxpipeResult.summary;
|
||||
if (pxpipeResult.body) translatedBody = pxpipeResult.body;
|
||||
if (pxpipeSummary?.applied) xf.push(`PXPIPE:${pxpipeSummary.imageCount}img`);
|
||||
try { onPxpipeEvent?.({ provider, model, ...pxpipeSummary }); } catch { /* stats must not break requests */ }
|
||||
}
|
||||
|
||||
if (xf.length && log?.line) log.line(reqTag, "⚙", xf.join(" · "));
|
||||
|
||||
const executor = getExecutor(provider);
|
||||
trackPendingRequest(model, provider, connectionId, true);
|
||||
appendRequestLog({ model, provider, connectionId, status: "PENDING" }).catch(() => { });
|
||||
@@ -198,7 +253,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
if (onDisconnect) onDisconnect(reason);
|
||||
},
|
||||
onError: () => trackPendingRequest(model, provider, connectionId, false),
|
||||
log, provider, model
|
||||
log, provider, model, reqTag
|
||||
});
|
||||
|
||||
const proxyOptions = {
|
||||
@@ -253,6 +308,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
request: extractRequestConfig(body, stream),
|
||||
providerRequest: translatedBody || null,
|
||||
response: { error: error.message || String(error), status: error.name === "AbortError" ? 499 : 502, thinking: null },
|
||||
pxpipe: pxpipeSummary,
|
||||
status: "error"
|
||||
})).catch(() => { });
|
||||
|
||||
@@ -261,7 +317,9 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
return createErrorResult(499, "Request aborted");
|
||||
}
|
||||
const errMsg = formatProviderError(error, provider, model, HTTP_STATUS.BAD_GATEWAY);
|
||||
console.log(`${COLORS.red}[ERROR] ${errMsg}${COLORS.reset}`);
|
||||
if (log?.errorLine) {
|
||||
log.errorLine(reqTag, "✗", `ERROR 502 · ${provider}/${model} · ${Date.now() - requestStartTime}ms\n ${errMsg}${error.stack ? `\n ${error.stack}` : ""}`);
|
||||
}
|
||||
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg);
|
||||
}
|
||||
|
||||
@@ -270,7 +328,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
try {
|
||||
const newCredentials = await refreshWithRetry(() => executor.refreshCredentials(credentials, log), 3, log);
|
||||
if (newCredentials?.accessToken || newCredentials?.copilotToken) {
|
||||
log?.info?.("TOKEN", `${provider.toUpperCase()} | refreshed`);
|
||||
if (log?.line) log.line(reqTag, "🔑", `TOKEN REFRESHED · ${provider}/${model}`);
|
||||
Object.assign(credentials, newCredentials);
|
||||
if (onCredentialsRefreshed) {
|
||||
try { await onCredentialsRefreshed(newCredentials); } catch (e) { log?.warn?.("TOKEN", `onCredentialsRefreshed failed: ${e.message}`); }
|
||||
@@ -299,16 +357,20 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
request: extractRequestConfig(body, stream),
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
response: { error: message, status: statusCode, thinking: null },
|
||||
pxpipe: pxpipeSummary,
|
||||
status: "error"
|
||||
})).catch(() => { });
|
||||
|
||||
const errMsg = formatProviderError(new Error(message), provider, model, statusCode);
|
||||
console.log(`${COLORS.red}[ERROR] ${errMsg}${COLORS.reset}`);
|
||||
if (log?.errorLine) {
|
||||
const urlStr = providerUrl ? `\n URL: ${providerUrl}` : "";
|
||||
log.errorLine(reqTag, "✗", `ERROR ${statusCode} · ${provider}/${model} · ${Date.now() - requestStartTime}ms${urlStr}\n ${errMsg}`);
|
||||
}
|
||||
reqLogger.logError(new Error(message), finalBody || translatedBody);
|
||||
return createErrorResult(statusCode, errMsg, resetsAtMs);
|
||||
}
|
||||
|
||||
const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess };
|
||||
const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, pxpipe: pxpipeSummary, reqTag, log };
|
||||
const appendLog = (extra) => appendRequestLog({ model, provider, connectionId, ...extra }).catch(() => { });
|
||||
const trackDone = () => trackPendingRequest(model, provider, connectionId, false);
|
||||
|
||||
|
||||
@@ -6,7 +6,7 @@ import { addBufferToUsage, filterUsageForFormat } from "../../utils/usageTrackin
|
||||
import { createErrorResult } from "../../utils/error.js";
|
||||
import { HTTP_STATUS } from "../../config/runtimeConfig.js";
|
||||
import { parseSSEToOpenAIResponse } from "./sseToJsonHandler.js";
|
||||
import { buildRequestDetail, extractRequestConfig, extractUsageFromResponse, saveUsageStats } from "./requestDetail.js";
|
||||
import { buildRequestDetail, extractRequestConfig, extractUsageFromResponse, saveUsageStats, formatDoneLine } from "./requestDetail.js";
|
||||
import { appendRequestLog, saveRequestDetail } from "@/lib/usageDb.js";
|
||||
import { decloakToolNames } from "../../utils/claudeCloaking.js";
|
||||
|
||||
@@ -198,7 +198,7 @@ export function translateNonStreamingResponse(responseBody, targetFormat, source
|
||||
/**
|
||||
* Handle non-streaming response from provider.
|
||||
*/
|
||||
export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, trackDone, appendLog }) {
|
||||
export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, trackDone, appendLog, pxpipe, reqTag, log }) {
|
||||
trackDone();
|
||||
const contentType = providerResponse.headers.get("content-type") || "";
|
||||
let responseBody;
|
||||
@@ -235,7 +235,8 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
|
||||
const usage = extractUsageFromResponse(responseBody);
|
||||
appendLog({ tokens: usage, status: "200 OK" });
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint });
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint, silent: true });
|
||||
if (log?.line) log.line(reqTag, "📊", formatDoneLine({ usage, latency: { total: Date.now() - requestStartTime } }));
|
||||
|
||||
const translatedResponse = needsTranslation(targetFormat, sourceFormat)
|
||||
? translateNonStreamingResponse(responseBody, targetFormat, sourceFormat)
|
||||
@@ -296,6 +297,7 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
thinking: translatedResponse?.choices?.[0]?.message?.reasoning_content || translatedResponse?.reasoning_content || null,
|
||||
finish_reason: translatedResponse?.choices?.[0]?.finish_reason || "unknown"
|
||||
},
|
||||
pxpipe,
|
||||
status: "success"
|
||||
}, { endpoint: clientRawRequest?.endpoint || null })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to save:", err.message);
|
||||
|
||||
@@ -69,12 +69,31 @@ export function buildRequestDetail(base, overrides = {}) {
|
||||
providerRequest: base.providerRequest || null,
|
||||
providerResponse: base.providerResponse || null,
|
||||
response: base.response || {},
|
||||
pxpipe: base.pxpipe || undefined,
|
||||
status: base.status || "success",
|
||||
...overrides
|
||||
};
|
||||
}
|
||||
|
||||
export function saveUsageStats({ provider, model, tokens, connectionId, apiKey, endpoint, label = "USAGE" }) {
|
||||
// Build the "done" summary: duration, ttft, in/out tokens with cache breakdown
|
||||
export function formatDoneLine({ usage, latency }) {
|
||||
const u = usage || {};
|
||||
const inTok = u.prompt_tokens ?? u.input_tokens ?? 0;
|
||||
const outTok = u.completion_tokens ?? u.output_tokens ?? 0;
|
||||
const cacheRead = u.cache_read_input_tokens ?? u.cached_tokens ?? u.prompt_tokens_details?.cached_tokens ?? 0;
|
||||
const cacheCreate = u.cache_creation_input_tokens ?? 0;
|
||||
let inStr = `IN ${inTok}`;
|
||||
if (cacheRead || cacheCreate) {
|
||||
const parts = [];
|
||||
if (cacheRead) parts.push(`↻${cacheRead}`);
|
||||
if (cacheCreate) parts.push(`+${cacheCreate}`);
|
||||
inStr += ` (CACHE ${parts.join(" ")})`;
|
||||
}
|
||||
const ttftStr = latency?.ttft ? ` · TTFT ${latency.ttft}ms` : "";
|
||||
return `DONE ${latency?.total ?? 0}ms${ttftStr} · ${inStr} · OUT ${outTok}`;
|
||||
}
|
||||
|
||||
export function saveUsageStats({ provider, model, tokens, connectionId, apiKey, endpoint, label = "USAGE", silent = false }) {
|
||||
if (!tokens || typeof tokens !== "object") return;
|
||||
|
||||
const inTokens = tokens.input_tokens ?? tokens.prompt_tokens ?? 0;
|
||||
@@ -82,9 +101,11 @@ export function saveUsageStats({ provider, model, tokens, connectionId, apiKey,
|
||||
|
||||
if (inTokens === 0 && outTokens === 0) return;
|
||||
|
||||
const time = new Date().toLocaleTimeString("en-US", { hour12: false, hour: "2-digit", minute: "2-digit", second: "2-digit" });
|
||||
const accountSuffix = connectionId ? ` | account=${connectionId.slice(0, 8)}...` : "";
|
||||
console.log(`${COLORS.green}[${time}] 📊 [${label}] ${provider.toUpperCase()} | in=${inTokens} | out=${outTokens}${accountSuffix}${COLORS.reset}`);
|
||||
if (!silent) {
|
||||
const time = new Date().toLocaleTimeString("en-US", { hour12: false, hour: "2-digit", minute: "2-digit", second: "2-digit" });
|
||||
const accountSuffix = connectionId ? ` | account=${connectionId.slice(0, 8)}...` : "";
|
||||
console.log(`${COLORS.green}[${time}] 📊 [${label}] ${provider.toUpperCase()} | in=${inTokens} | out=${outTokens}${accountSuffix}${COLORS.reset}`);
|
||||
}
|
||||
|
||||
// Canonicalize to one storage convention (prompt_tokens cache-inclusive) so
|
||||
// cached/cache-creation tokens survive to cost calc + stats. See canonicalizeUsage.
|
||||
|
||||
@@ -3,7 +3,7 @@ import { createErrorResult } from "../../utils/error.js";
|
||||
import { HTTP_STATUS } from "../../config/runtimeConfig.js";
|
||||
import { FORMATS } from "../../translator/formats.js";
|
||||
import { PROVIDERS } from "../../config/providers.js";
|
||||
import { buildRequestDetail, extractRequestConfig, saveUsageStats } from "./requestDetail.js";
|
||||
import { buildRequestDetail, extractRequestConfig, saveUsageStats, formatDoneLine } from "./requestDetail.js";
|
||||
|
||||
// Responses-API providers (e.g. codex) may emit SSE without content-type + use Responses output shape
|
||||
const isResponsesProvider = (p) => PROVIDERS[p]?.format === FORMATS.OPENAI_RESPONSES;
|
||||
@@ -102,7 +102,7 @@ export function parseSSEToOpenAIResponse(rawSSE, fallbackModel) {
|
||||
* Handle case: provider forced streaming but client wants JSON.
|
||||
* Supports both Codex/Responses API SSE and standard Chat Completions SSE.
|
||||
*/
|
||||
export async function handleForcedSSEToJson({ providerResponse, sourceFormat, provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, trackDone, appendLog }) {
|
||||
export async function handleForcedSSEToJson({ providerResponse, sourceFormat, provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, trackDone, appendLog, reqTag, log }) {
|
||||
const contentType = providerResponse.headers.get("content-type") || "";
|
||||
const isSSE = contentType.includes("text/event-stream") || (contentType === "" && isResponsesProvider(provider));
|
||||
if (!isSSE) return null; // not handled here
|
||||
@@ -124,7 +124,8 @@ export async function handleForcedSSEToJson({ providerResponse, sourceFormat, pr
|
||||
|
||||
const usage = jsonResponse.usage || {};
|
||||
appendLog({ tokens: usage, status: "200 OK" });
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint });
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint, silent: true });
|
||||
if (log?.line) log.line(reqTag, "📊", formatDoneLine({ usage, latency: { total: Date.now() - requestStartTime } }));
|
||||
|
||||
const { msgItem, textContent } = pickAssistantMessageForChatCompletion(jsonResponse.output);
|
||||
const totalLatency = Date.now() - requestStartTime;
|
||||
@@ -200,7 +201,8 @@ export async function handleForcedSSEToJson({ providerResponse, sourceFormat, pr
|
||||
|
||||
const usage = parsed.usage || {};
|
||||
appendLog({ tokens: usage, status: "200 OK" });
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint });
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint, silent: true });
|
||||
if (log?.line) log.line(reqTag, "📊", formatDoneLine({ usage, latency: { total: Date.now() - requestStartTime } }));
|
||||
|
||||
const totalLatency = Date.now() - requestStartTime;
|
||||
saveRequestDetail(buildRequestDetail({
|
||||
|
||||
@@ -5,7 +5,7 @@ import { pipeWithDisconnect } from "../../utils/streamHandler.js";
|
||||
import { PROVIDERS } from "../../config/providers.js";
|
||||
import { STREAM_STALL_TIMEOUT_MS } from "../../config/runtimeConfig.js";
|
||||
import { buildAbortedResponsesTerminalBytes } from "../../utils/responsesStreamHelpers.js";
|
||||
import { buildRequestDetail, extractRequestConfig, saveUsageStats } from "./requestDetail.js";
|
||||
import { buildRequestDetail, extractRequestConfig, saveUsageStats, formatDoneLine } from "./requestDetail.js";
|
||||
import { saveRequestDetail } from "@/lib/usageDb.js";
|
||||
import { SSE_HEADERS_CORS as SSE_HEADERS } from "../../utils/sseConstants.js";
|
||||
|
||||
@@ -43,7 +43,7 @@ function buildTransformStream({ provider, sourceFormat, targetFormat, userAgent,
|
||||
/**
|
||||
* 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, streamController, onStreamComplete, streamDetailId }) {
|
||||
export async function handleStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, userAgent, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, streamController, onStreamComplete, streamDetailId, pxpipe, reqTag, log }) {
|
||||
if (onRequestSuccess) {
|
||||
Promise.resolve()
|
||||
.then(onRequestSuccess)
|
||||
@@ -67,7 +67,8 @@ export async function handleStreamingResponse({ providerResponse, provider, mode
|
||||
const shortMsg = sanitizedTitle
|
||||
|| (bodyText.length < 200 ? bodyText.replace(/<[^>]*>/g, '').trim().slice(0, 160) : `Upstream returned non-SSE response (${upstreamContentType})`);
|
||||
const status = providerResponse.status || 502;
|
||||
console.warn(`[STREAM] ${provider} | ${model} | blocked pipe: ${shortMsg} [${status}]`);
|
||||
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,
|
||||
@@ -94,6 +95,7 @@ export async function handleStreamingResponse({ providerResponse, provider, mode
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
providerResponse: "[Streaming - raw response not captured]",
|
||||
response: { content: "[Streaming in progress...]", thinking: null, type: "streaming" },
|
||||
pxpipe,
|
||||
status: "success"
|
||||
}, { id: streamDetailId })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to save streaming request:", err.message);
|
||||
@@ -108,7 +110,7 @@ export async function handleStreamingResponse({ providerResponse, provider, mode
|
||||
/**
|
||||
* Build onStreamComplete callback for streaming usage tracking.
|
||||
*/
|
||||
export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest }) {
|
||||
export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest, pxpipe, reqTag, log }) {
|
||||
const streamDetailId = `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`;
|
||||
|
||||
const onStreamComplete = (contentObj, usage, ttftAt) => {
|
||||
@@ -127,12 +129,15 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
providerResponse: safeContent,
|
||||
response: { content: safeContent, thinking: safeThinking, type: "streaming" },
|
||||
pxpipe,
|
||||
status: "success"
|
||||
}, { id: streamDetailId })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to update streaming content:", err.message);
|
||||
});
|
||||
|
||||
saveUsageStats({ provider, model, tokens: usage, connectionId, apiKey, endpoint: clientRawRequest?.endpoint, label: "STREAM USAGE" });
|
||||
// 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 };
|
||||
|
||||
@@ -273,6 +273,53 @@ const CHAT_SEARCH_CONFIG = {
|
||||
const tokens = data?.usage?.total_tokens || 0;
|
||||
return { text, citations, tokens };
|
||||
}
|
||||
},
|
||||
|
||||
"perplexity-agent": {
|
||||
endpoint: () => searchEndpoint("perplexity-agent"),
|
||||
buildBody: (query, model) => ({
|
||||
model,
|
||||
input: query,
|
||||
tools: [{ type: "web_search" }]
|
||||
}),
|
||||
buildHeaders: (token) => ({
|
||||
"Content-Type": "application/json",
|
||||
Authorization: `Bearer ${token}`
|
||||
}),
|
||||
extractAnswer: (data) => {
|
||||
const output = Array.isArray(data?.output) ? data.output : [];
|
||||
let text = "";
|
||||
const citations = [];
|
||||
for (const item of output) {
|
||||
const parts = Array.isArray(item?.content) ? item.content : [];
|
||||
for (const p of parts) {
|
||||
if (typeof p?.text === "string") text += p.text;
|
||||
const anns = Array.isArray(p?.annotations) ? p.annotations : [];
|
||||
for (const a of anns) {
|
||||
const c = normalizeCitation(a?.url ? a : a?.url_citation);
|
||||
if (c) citations.push(c);
|
||||
}
|
||||
}
|
||||
const results = Array.isArray(item?.results) ? item.results : [];
|
||||
for (const r of results) {
|
||||
const url = r?.url || r?.link;
|
||||
if (!url) continue;
|
||||
citations.push({
|
||||
url,
|
||||
title: r?.title || "",
|
||||
snippet: r?.snippet || ""
|
||||
});
|
||||
}
|
||||
}
|
||||
if (!citations.length && Array.isArray(data?.citations)) {
|
||||
for (const c of data.citations) {
|
||||
const n = normalizeCitation(c);
|
||||
if (n) citations.push(n);
|
||||
}
|
||||
}
|
||||
const tokens = data?.usage?.total_tokens || 0;
|
||||
return { text, citations, tokens };
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
166
open-sse/handlers/videoCore.js
Normal file
166
open-sse/handlers/videoCore.js
Normal file
@@ -0,0 +1,166 @@
|
||||
import { createErrorResult } from "../utils/error.js";
|
||||
import { HTTP_STATUS } from "../config/runtimeConfig.js";
|
||||
import { refreshTokenByProvider } from "../services/tokenRefresh.js";
|
||||
import { PROVIDER_MEDIA } from "../providers/index.js";
|
||||
|
||||
// Upstream fetch deadline for video job submission/polling (the job itself is
|
||||
// async upstream — this only bounds the HTTP round-trip, not video rendering).
|
||||
const VIDEO_FETCH_TIMEOUT_MS = Number(process.env.VIDEO_FETCH_TIMEOUT_MS || 120000);
|
||||
|
||||
// POST /videos/* creates a billable upstream job. A network error after the
|
||||
// request left the socket may still have created the job, so creation is NEVER
|
||||
// auto-retried (the only re-send is the auth retry after a 401/403 refresh,
|
||||
// which upstream rejects before job creation).
|
||||
export const VIDEO_ACTIONS = new Set(["generations", "edits", "extensions"]);
|
||||
|
||||
export function getVideoConfig(provider) {
|
||||
return PROVIDER_MEDIA[provider]?.videoConfig || null;
|
||||
}
|
||||
|
||||
/** Strip bearer tokens / obvious secrets from text destined for clients or logs. */
|
||||
export function sanitizeSecrets(text, credentials = null) {
|
||||
if (!text) return text;
|
||||
let out = String(text).replace(/Bearer\s+[A-Za-z0-9._~+/=-]{8,}/gi, "Bearer [redacted]");
|
||||
for (const key of ["accessToken", "refreshToken", "apiKey"]) {
|
||||
const secret = credentials?.[key];
|
||||
if (typeof secret === "string" && secret.length >= 8) {
|
||||
out = out.split(secret).join("[redacted]");
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
function buildUpstreamUrl(config, action, requestId) {
|
||||
const base = config.baseUrl.replace(/\/$/, "");
|
||||
return requestId ? `${base}/${encodeURIComponent(requestId)}` : `${base}/${action}`;
|
||||
}
|
||||
|
||||
function buildHeaders({ token, contentType, idempotencyKey }) {
|
||||
const headers = { Accept: "application/json" };
|
||||
if (token) headers.Authorization = `Bearer ${token}`;
|
||||
if (contentType) headers["Content-Type"] = contentType;
|
||||
if (idempotencyKey) headers["Idempotency-Key"] = idempotencyKey;
|
||||
return headers;
|
||||
}
|
||||
|
||||
function combineSignals(signal, timeoutMs) {
|
||||
const timeoutSignal = typeof AbortSignal?.timeout === "function" ? AbortSignal.timeout(timeoutMs) : null;
|
||||
if (signal && timeoutSignal && typeof AbortSignal.any === "function") {
|
||||
return AbortSignal.any([signal, timeoutSignal]);
|
||||
}
|
||||
return signal || timeoutSignal || undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Transparent proxy for async video jobs (xAI Grok Imagine shape).
|
||||
*
|
||||
* - Forwards the raw body byte-for-byte (JSON or multipart) — no reshaping.
|
||||
* - Passes upstream JSON (request_id, status, video.url, error) back verbatim.
|
||||
* - 401/403 with a refresh token: refresh ONCE, retry ONCE. No other retry.
|
||||
* - Upstream error text is sanitized before it reaches the client.
|
||||
*
|
||||
* @param {object} options
|
||||
* @param {string} options.provider - Provider id (must have registry videoConfig)
|
||||
* @param {"generations"|"edits"|"extensions"|null} options.action - Creation action (POST)
|
||||
* @param {string|null} [options.requestId] - Poll target (GET /videos/{id})
|
||||
* @param {Buffer|string|null} [options.rawBody] - Exact body to forward
|
||||
* @param {string|null} [options.contentType] - Original Content-Type header
|
||||
* @param {string|null} [options.idempotencyKey] - Forwarded Idempotency-Key
|
||||
* @param {object} options.credentials - { accessToken?, apiKey?, refreshToken?, authType? }
|
||||
* @param {AbortSignal} [options.signal] - Client cancellation signal
|
||||
* @param {number} [options.timeoutMs]
|
||||
* @param {object} [options.log]
|
||||
* @param {function} [options.onCredentialsRefreshed]
|
||||
* @returns {Promise<{ success: boolean, response: Response, status?: number, error?: string }>}
|
||||
*/
|
||||
export async function handleVideoProxyCore({
|
||||
provider,
|
||||
action = null,
|
||||
requestId = null,
|
||||
rawBody = null,
|
||||
contentType = null,
|
||||
idempotencyKey = null,
|
||||
credentials,
|
||||
signal,
|
||||
timeoutMs = VIDEO_FETCH_TIMEOUT_MS,
|
||||
log,
|
||||
onCredentialsRefreshed,
|
||||
}) {
|
||||
const config = getVideoConfig(provider);
|
||||
if (!config) {
|
||||
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Provider '${provider}' does not support video generation`);
|
||||
}
|
||||
if (!requestId && !VIDEO_ACTIONS.has(action)) {
|
||||
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Unknown video action: ${action}`);
|
||||
}
|
||||
|
||||
const method = requestId ? "GET" : "POST";
|
||||
const url = buildUpstreamUrl(config, action, requestId);
|
||||
const fetchSignal = combineSignals(signal, timeoutMs);
|
||||
|
||||
const doFetch = (token) =>
|
||||
fetch(url, {
|
||||
method,
|
||||
headers: buildHeaders({ token, contentType: method === "POST" ? contentType : null, idempotencyKey: method === "POST" ? idempotencyKey : null }),
|
||||
body: method === "POST" ? rawBody : undefined,
|
||||
signal: fetchSignal,
|
||||
});
|
||||
|
||||
let upstream;
|
||||
try {
|
||||
upstream = await doFetch(credentials?.accessToken || credentials?.apiKey);
|
||||
} catch (error) {
|
||||
if (error?.name === "AbortError" || error?.name === "TimeoutError") {
|
||||
return createErrorResult(HTTP_STATUS.REQUEST_TIMEOUT, `[${provider}] video ${method} aborted: ${error.message}`);
|
||||
}
|
||||
// Never re-send a creation POST on network error — the job may already exist upstream.
|
||||
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, sanitizeSecrets(`[${provider}] video upstream fetch failed: ${error.message}`, credentials));
|
||||
}
|
||||
|
||||
// 401/403 → refresh once → retry once (OAuth accounts only; API keys can't refresh)
|
||||
if (
|
||||
(upstream.status === HTTP_STATUS.UNAUTHORIZED || upstream.status === HTTP_STATUS.FORBIDDEN) &&
|
||||
credentials?.refreshToken
|
||||
) {
|
||||
let refreshed = null;
|
||||
try {
|
||||
refreshed = await refreshTokenByProvider(provider, credentials, log);
|
||||
} catch (error) {
|
||||
log?.warn?.("TOKEN", `${provider} | video refresh error: ${sanitizeSecrets(error.message, credentials)}`);
|
||||
}
|
||||
if (refreshed?.accessToken) {
|
||||
log?.info?.("TOKEN", `${provider.toUpperCase()} | refreshed for video ${method}`);
|
||||
Object.assign(credentials, refreshed);
|
||||
if (onCredentialsRefreshed) await onCredentialsRefreshed(refreshed);
|
||||
try {
|
||||
await upstream.body?.cancel?.();
|
||||
} catch { /* noop */ }
|
||||
try {
|
||||
upstream = await doFetch(credentials.accessToken || credentials.apiKey);
|
||||
} catch (error) {
|
||||
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, sanitizeSecrets(`[${provider}] video retry after refresh failed: ${error.message}`, credentials));
|
||||
}
|
||||
} else {
|
||||
log?.warn?.("TOKEN", `${provider.toUpperCase()} | video refresh failed — account needs re-auth`);
|
||||
}
|
||||
}
|
||||
|
||||
const bodyText = await upstream.text().catch(() => "");
|
||||
|
||||
if (!upstream.ok) {
|
||||
const message = sanitizeSecrets(bodyText || `HTTP ${upstream.status}`, credentials);
|
||||
return createErrorResult(upstream.status, `[${provider}] ${message.slice(0, 2000)}`);
|
||||
}
|
||||
|
||||
// Success: pass the upstream JSON through untouched (request_id / status / video.url).
|
||||
return {
|
||||
success: true,
|
||||
response: new Response(bodyText, {
|
||||
status: upstream.status,
|
||||
headers: {
|
||||
"Content-Type": upstream.headers.get("content-type") || "application/json",
|
||||
"Access-Control-Allow-Origin": "*",
|
||||
},
|
||||
}),
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user