- Qoder: handle mid-stream errors by returning proper error Response instead of embedding in stream - Qoder: add refreshCredentials() to validate token via quota endpoint before requests - chatCore: validate and refresh provider tokens before sending chat requests - chatCore: bail early with 401 on unrecoverable token refresh errors - Usage stats: track apiKey, comboName, fallbackHistory in request details - Dashboard: improve Combos, Endpoint, Provider, Usage, and RequestDetails pages - API keys route: upsert logic with provider_type support - DB repos: usageRepo query improvements, requestDetailsRepo pagination, apiKeysRepo updates
118 lines
6.6 KiB
JavaScript
118 lines
6.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 { STREAM_STALL_TIMEOUT_MS } from "../../config/runtimeConfig.js";
|
|
import { buildAbortedResponsesTerminalBytes } from "../../utils/responsesStreamHelpers.js";
|
|
import { buildRequestDetail, extractRequestConfig } from "./requestDetail.js";
|
|
import { saveRequestDetail } from "@/lib/usageDb.js";
|
|
import { SSE_HEADERS_CORS as SSE_HEADERS } from "../../utils/sseConstants.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, model, connectionId, body, onStreamComplete, apiKey, endpoint, comboName }) {
|
|
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, endpoint, comboName);
|
|
}
|
|
|
|
if (needsTranslation(targetFormat, sourceFormat)) {
|
|
return createSSETransformStreamWithLogger(targetFormat, sourceFormat, provider, reqLogger, toolNameMap, model, connectionId, body, onStreamComplete, apiKey, endpoint, comboName);
|
|
}
|
|
|
|
return createPassthroughStreamWithLogger(provider, reqLogger, model, connectionId, body, onStreamComplete, apiKey, endpoint, comboName);
|
|
}
|
|
|
|
/**
|
|
* Handle streaming response — pipe provider SSE through transform stream to client.
|
|
*/
|
|
export function handleStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, userAgent, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, comboName, fallbackHistory, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, streamController, onStreamComplete }) {
|
|
if (onRequestSuccess) onRequestSuccess();
|
|
|
|
const transformStream = buildTransformStream({ provider, sourceFormat, targetFormat, userAgent, reqLogger, toolNameMap, model, connectionId, body, onStreamComplete, apiKey, endpoint: clientRawRequest?.endpoint, comboName });
|
|
|
|
// Responses passthrough: synthesize response.failed + [DONE] if the stream aborts/stalls before a terminal event
|
|
const isResponsesPassthrough = sourceFormat === FORMATS.OPENAI_RESPONSES && targetFormat === FORMATS.OPENAI_RESPONSES;
|
|
const onAbortTerminal = isResponsesPassthrough ? buildAbortedResponsesTerminalBytes : null;
|
|
const stallTimeoutMs = PROVIDERS[provider]?.stallTimeoutMs || STREAM_STALL_TIMEOUT_MS;
|
|
const transformedBody = pipeWithDisconnect(providerResponse, transformStream, streamController, onAbortTerminal, stallTimeoutMs);
|
|
|
|
const streamDetailId = `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`;
|
|
saveRequestDetail(buildRequestDetail({
|
|
provider, model, connectionId, apiKey, comboName, fallbackHistory,
|
|
latency: { ttft: 0, total: Date.now() - requestStartTime },
|
|
tokens: { prompt_tokens: 0, completion_tokens: 0 },
|
|
request: extractRequestConfig(body, stream),
|
|
providerRequest: finalBody || translatedBody || null,
|
|
providerResponse: "[Streaming - raw response not captured]",
|
|
response: { content: "[Streaming in progress...]", thinking: null, type: "streaming" },
|
|
status: "success"
|
|
}, { id: streamDetailId })).catch(err => {
|
|
console.error("[RequestDetail] Failed to save streaming request:", err.message);
|
|
});
|
|
|
|
return {
|
|
success: true,
|
|
response: new Response(transformedBody, { headers: SSE_HEADERS })
|
|
};
|
|
}
|
|
|
|
/**
|
|
* Build onStreamComplete callback for streaming usage tracking.
|
|
* @param {object} options.midStreamError - Shared state object from executor (filled during streaming)
|
|
* @param {function} options.onMidStreamError - Callback to invoke when mid-stream error is detected
|
|
*/
|
|
export function buildOnStreamComplete({ provider, model, connectionId, apiKey, comboName, fallbackHistory, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest, midStreamError, onMidStreamError }) {
|
|
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;
|
|
|
|
// Check if a mid-stream error was detected during streaming (e.g. Qoder quota exceeded)
|
|
// The stream already returned HTTP 200, so the pre-stream error path in chatCore didn't fire.
|
|
// Apply cooldown here so the account gets locked for subsequent requests.
|
|
if (midStreamError?.error && typeof onMidStreamError === "function") {
|
|
onMidStreamError(midStreamError.error).catch(err => {
|
|
console.error("[MidStreamError] Failed to apply cooldown:", err.message);
|
|
});
|
|
}
|
|
|
|
saveRequestDetail(buildRequestDetail({
|
|
provider, model, connectionId, apiKey, comboName, fallbackHistory,
|
|
latency,
|
|
tokens: usage || { prompt_tokens: 0, completion_tokens: 0 },
|
|
request: extractRequestConfig(body, stream),
|
|
providerRequest: finalBody || translatedBody || null,
|
|
providerResponse: safeContent,
|
|
response: { content: safeContent, thinking: safeThinking, type: "streaming" },
|
|
status: midStreamError?.error ? "error" : "success"
|
|
}, { id: streamDetailId })).catch(err => {
|
|
console.error("[RequestDetail] Failed to update streaming content:", err.message);
|
|
});
|
|
};
|
|
|
|
return { onStreamComplete, streamDetailId };
|
|
}
|