merge: integrate origin/master (v0.5.50) into gitea/new_feature
- Resolve conflicts in chatCore handlers: keep apiKey/streamErrorPatterns from the details-filters feature, adopt origin's stripContinuityFields, customToolNames, cache-inclusive usage accounting, and Responses-API SSE→JSON conversion - Adopt origin's provider usage handlers (codebuddy-intl, qoder creds) and modality detection (audio/video inputs) - Keep requestDetails apiKey column (schema v2) + masked key persistence Co-authored-by: CommandCodeBot <noreply@commandcode.ai>
This commit is contained in:
@@ -4,7 +4,7 @@ import {
|
||||
resolveTransport,
|
||||
} from "../services/provider.js";
|
||||
import { translateRequest } from "../translator/index.js";
|
||||
import { stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js";
|
||||
import { applyThinking, extractThinking, stripThinkingSuffix } from "../translator/concerns/thinkingUnified.js";
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { normalizeClaudePassthrough } from "../translator/formats/claude.js";
|
||||
import { createStreamController } from "../utils/streamHandler.js";
|
||||
@@ -61,7 +61,6 @@ 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";
|
||||
|
||||
/**
|
||||
@@ -71,6 +70,26 @@ import { resolveSessionId } from "../utils/sessionManager.js";
|
||||
* @param {object} options.credentials - Provider credentials
|
||||
* @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses")
|
||||
*/
|
||||
/**
|
||||
* 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;
|
||||
}
|
||||
|
||||
export async function handleChatCore({
|
||||
body,
|
||||
modelInfo,
|
||||
@@ -138,7 +157,7 @@ export async function handleChatCore({
|
||||
// Multi-endpoint providers: pick transport matching sourceFormat → zero translation
|
||||
const runtimeTransport = resolveTransport(provider, sourceFormat);
|
||||
const targetFormat =
|
||||
modelTargetFormat || runtimeTransport?.format || getTargetFormat(provider);
|
||||
modelTargetFormat || runtimeTransport?.format || getTargetFormat(provider, credentials);
|
||||
if (runtimeTransport && credentials)
|
||||
credentials.runtimeTransport = runtimeTransport;
|
||||
const stripList = getModelStrip(alias, model);
|
||||
@@ -248,12 +267,25 @@ export async function handleChatCore({
|
||||
|
||||
let translatedBody;
|
||||
let toolNameMap;
|
||||
let customToolNames;
|
||||
if (passthrough) {
|
||||
log?.debug?.(
|
||||
"PASSTHROUGH",
|
||||
`${clientTool} → ${provider} | native lossless`,
|
||||
);
|
||||
translatedBody = { ...body, model: stripThinkingSuffix(upstreamModel) };
|
||||
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;
|
||||
}
|
||||
}
|
||||
// Normalize newer Cowork/CC beta shapes (adaptive thinking, mid-conversation system) the API rejects
|
||||
if (clientTool === "claude")
|
||||
normalizeClaudePassthrough(translatedBody, translatedBody.model);
|
||||
@@ -280,7 +312,10 @@ export async function handleChatCore({
|
||||
}
|
||||
toolNameMap = translatedBody._toolNameMap;
|
||||
delete translatedBody._toolNameMap;
|
||||
customToolNames = translatedBody._customToolNames;
|
||||
delete translatedBody._customToolNames;
|
||||
translatedBody.model = stripThinkingSuffix(upstreamModel);
|
||||
stripContinuityFields(translatedBody);
|
||||
}
|
||||
|
||||
// Dedupe duplicate built-in tools when equivalent MCP tools are present (Claude clients only).
|
||||
@@ -736,6 +771,8 @@ export async function handleChatCore({
|
||||
...sharedCtx,
|
||||
providerResponse,
|
||||
sourceFormat,
|
||||
targetFormat: providerResponseFormat,
|
||||
customToolNames,
|
||||
trackDone,
|
||||
appendLog,
|
||||
});
|
||||
@@ -754,6 +791,7 @@ export async function handleChatCore({
|
||||
targetFormat: providerResponseFormat,
|
||||
reqLogger,
|
||||
toolNameMap,
|
||||
customToolNames,
|
||||
trackDone,
|
||||
appendLog,
|
||||
});
|
||||
@@ -773,6 +811,7 @@ export async function handleChatCore({
|
||||
userAgent,
|
||||
reqLogger,
|
||||
toolNameMap,
|
||||
customToolNames,
|
||||
streamController,
|
||||
onStreamComplete,
|
||||
streamDetailId,
|
||||
|
||||
@@ -9,6 +9,7 @@ import { parseSSEToOpenAIResponse } from "./sseToJsonHandler.js";
|
||||
import { buildRequestDetail, extractRequestConfig, extractUsageFromResponse, saveUsageStats, formatDoneLine } from "./requestDetail.js";
|
||||
import { appendRequestLog, saveRequestDetail } from "@/lib/usageDb.js";
|
||||
import { decloakToolNames } from "../../utils/claudeCloaking.js";
|
||||
import { ROLE, RESPONSES_ITEM } from "../../translator/schema/index.js";
|
||||
|
||||
function parseToolArguments(value) {
|
||||
if (!value) return {};
|
||||
@@ -60,11 +61,93 @@ function openAICompletionToClaudeMessage(responseBody) {
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert an OpenAI Chat Completions non-streaming response body into the
|
||||
* OpenAI Responses API shape. Used when a Responses-format client (e.g. Codex)
|
||||
* is routed to a Chat Completions upstream and `stream:false` — the streaming
|
||||
* path already emits Responses events, but the JSON path returned a raw
|
||||
* `chat.completion` body, so tool_calls were invisible to Responses clients.
|
||||
*/
|
||||
function extractCustomToolInput(argumentsValue) {
|
||||
const argumentsText = typeof argumentsValue === "string" ? argumentsValue : JSON.stringify(argumentsValue || {});
|
||||
try {
|
||||
const parsed = JSON.parse(argumentsText);
|
||||
if (parsed && typeof parsed === "object" && typeof parsed.input === "string") return parsed.input;
|
||||
} catch { /* raw freeform input */ }
|
||||
return argumentsText;
|
||||
}
|
||||
|
||||
function openAICompletionToResponses(responseBody, customToolNames = null) {
|
||||
const choice = responseBody?.choices?.[0];
|
||||
if (!choice) return responseBody;
|
||||
|
||||
const message = choice.message || {};
|
||||
const output = [];
|
||||
|
||||
// Reasoning → a reasoning item (summary text), mirroring the streaming path.
|
||||
const reasoning = message.reasoning_content || message.reasoning;
|
||||
if (typeof reasoning === "string" && reasoning.length > 0) {
|
||||
output.push({
|
||||
type: RESPONSES_ITEM.REASONING,
|
||||
summary: [{ type: RESPONSES_ITEM.SUMMARY_TEXT, text: reasoning }],
|
||||
});
|
||||
}
|
||||
|
||||
// Assistant text → a message item with output_text content.
|
||||
const text = typeof message.content === "string" ? message.content : "";
|
||||
if (text.length > 0) {
|
||||
output.push({
|
||||
type: RESPONSES_ITEM.MESSAGE,
|
||||
role: ROLE.ASSISTANT,
|
||||
content: [{ type: RESPONSES_ITEM.OUTPUT_TEXT, text, annotations: [] }],
|
||||
});
|
||||
}
|
||||
|
||||
// tool_calls → function_call/custom_tool_call items (Responses-native tool shape).
|
||||
for (const tc of message.tool_calls || []) {
|
||||
const fn = tc.function || {};
|
||||
const custom = customToolNames?.has(fn.name);
|
||||
output.push({
|
||||
type: custom ? RESPONSES_ITEM.CUSTOM_TOOL_CALL : RESPONSES_ITEM.FUNCTION_CALL,
|
||||
id: `${custom ? "ctc" : "fc"}_${tc.id || ""}`,
|
||||
call_id: tc.id || "",
|
||||
name: fn.name || "",
|
||||
...(custom
|
||||
? { input: extractCustomToolInput(fn.arguments) }
|
||||
: { arguments: typeof fn.arguments === "string" ? fn.arguments : JSON.stringify(fn.arguments || {}) }),
|
||||
});
|
||||
}
|
||||
|
||||
const usage = responseBody.usage || {};
|
||||
const status = choice.finish_reason === "tool_calls" ? "completed" : (choice.finish_reason === "stop" ? "completed" : (choice.finish_reason || "completed"));
|
||||
|
||||
return {
|
||||
id: `resp_${responseBody.id || ""}`.replace(/^resp_chatcmpl-/, "resp_"),
|
||||
object: "response",
|
||||
created_at: responseBody.created || Math.floor(Date.now() / 1000),
|
||||
model: responseBody.model || "unknown",
|
||||
status,
|
||||
background: false,
|
||||
error: null,
|
||||
output,
|
||||
usage: {
|
||||
input_tokens: usage.prompt_tokens || usage.input_tokens || 0,
|
||||
output_tokens: usage.completion_tokens || usage.output_tokens || 0,
|
||||
total_tokens: usage.total_tokens || (usage.prompt_tokens || 0) + (usage.completion_tokens || 0),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Translate non-streaming response body from provider format → OpenAI format.
|
||||
*/
|
||||
export function translateNonStreamingResponse(responseBody, targetFormat, sourceFormat) {
|
||||
export function translateNonStreamingResponse(responseBody, targetFormat, sourceFormat, customToolNames = null) {
|
||||
if (targetFormat === sourceFormat) return responseBody;
|
||||
// Provider responded in OpenAI Chat Completions shape but the client speaks
|
||||
// Responses API — convert so tool_calls/text surface as Responses `output`.
|
||||
if (targetFormat === FORMATS.OPENAI && sourceFormat === FORMATS.OPENAI_RESPONSES) {
|
||||
return openAICompletionToResponses(responseBody, customToolNames);
|
||||
}
|
||||
if (targetFormat === FORMATS.OPENAI && sourceFormat === FORMATS.CLAUDE) {
|
||||
return openAICompletionToClaudeMessage(responseBody);
|
||||
}
|
||||
@@ -198,7 +281,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, pxpipe, reqTag, log }) {
|
||||
export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, customToolNames, trackDone, appendLog, pxpipe, reqTag, log }) {
|
||||
trackDone();
|
||||
const contentType = providerResponse.headers.get("content-type") || "";
|
||||
let responseBody;
|
||||
@@ -239,9 +322,12 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
if (log?.line) log.line(reqTag, "📊", formatDoneLine({ usage, latency: { total: Date.now() - requestStartTime } }));
|
||||
|
||||
const translatedResponse = needsTranslation(targetFormat, sourceFormat)
|
||||
? translateNonStreamingResponse(responseBody, targetFormat, sourceFormat)
|
||||
? translateNonStreamingResponse(responseBody, targetFormat, sourceFormat, customToolNames)
|
||||
: responseBody;
|
||||
const isClaudeMessageResponse = sourceFormat === FORMATS.CLAUDE && translatedResponse?.type === "message";
|
||||
// Responses-format translation produces a `object:"response"` body with no
|
||||
// `choices`; skip the Chat-Completions-specific post-processing below for it.
|
||||
const isResponsesResponse = sourceFormat === FORMATS.OPENAI_RESPONSES && translatedResponse?.object === "response";
|
||||
|
||||
// Fix finish_reason for tool_calls: some providers return non-standard values (e.g. "other")
|
||||
if (translatedResponse?.choices?.[0]) {
|
||||
@@ -254,13 +340,13 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
}
|
||||
|
||||
// Ensure OpenAI-required fields
|
||||
if (!isClaudeMessageResponse) {
|
||||
if (!isClaudeMessageResponse && !isResponsesResponse) {
|
||||
if (!translatedResponse.object) translatedResponse.object = "chat.completion";
|
||||
if (!translatedResponse.created) translatedResponse.created = Math.floor(Date.now() / 1000);
|
||||
}
|
||||
|
||||
// Strip Azure-specific fields
|
||||
if (!isClaudeMessageResponse) {
|
||||
if (!isClaudeMessageResponse && !isResponsesResponse) {
|
||||
delete translatedResponse.prompt_filter_results;
|
||||
if (translatedResponse?.choices) {
|
||||
for (const choice of translatedResponse.choices) delete choice.content_filter_results;
|
||||
@@ -274,7 +360,7 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
// Strip reasoning_content only when content is non-empty.
|
||||
// When content is empty (e.g. thinking models that used all tokens for reasoning),
|
||||
// reasoning_content is the only useful output and must be preserved.
|
||||
if (!isClaudeMessageResponse && translatedResponse?.choices) {
|
||||
if (!isClaudeMessageResponse && !isResponsesResponse && translatedResponse?.choices) {
|
||||
for (const choice of translatedResponse.choices) {
|
||||
if (choice?.message?.reasoning_content && choice.message.content) {
|
||||
delete choice.message.reasoning_content;
|
||||
|
||||
@@ -4,12 +4,8 @@ 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,
|
||||
formatDoneLine,
|
||||
} from "./requestDetail.js";
|
||||
import { buildRequestDetail, extractRequestConfig, saveUsageStats, formatDoneLine } from "./requestDetail.js";
|
||||
import { ROLE, RESPONSES_ITEM } from "../../translator/schema/index.js";
|
||||
|
||||
// Responses-API providers (e.g. codex) may emit SSE without content-type + use Responses output shape
|
||||
const isResponsesProvider = (p) =>
|
||||
@@ -41,6 +37,76 @@ function pickAssistantMessageForChatCompletion(output) {
|
||||
return { msgItem: last, textContent: textFromResponsesMessageItem(last) };
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert an OpenAI Chat Completions JSON body into the Responses API shape.
|
||||
* Inlined here (not imported from nonStreamingHandler.js) to avoid a circular
|
||||
* import. Mirrors openAICompletionToResponses in nonStreamingHandler.js.
|
||||
*/
|
||||
function extractCustomToolInput(argumentsValue) {
|
||||
const argumentsText = typeof argumentsValue === "string" ? argumentsValue : JSON.stringify(argumentsValue || {});
|
||||
try {
|
||||
const parsed = JSON.parse(argumentsText);
|
||||
if (parsed && typeof parsed === "object" && typeof parsed.input === "string") return parsed.input;
|
||||
} catch { /* raw freeform input */ }
|
||||
return argumentsText;
|
||||
}
|
||||
|
||||
function chatCompletionToResponses(responseBody, customToolNames = null) {
|
||||
const choice = responseBody?.choices?.[0];
|
||||
if (!choice) return responseBody;
|
||||
|
||||
const message = choice.message || {};
|
||||
const output = [];
|
||||
|
||||
const reasoning = message.reasoning_content || message.reasoning;
|
||||
if (typeof reasoning === "string" && reasoning.length > 0) {
|
||||
output.push({
|
||||
type: RESPONSES_ITEM.REASONING,
|
||||
summary: [{ type: RESPONSES_ITEM.SUMMARY_TEXT, text: reasoning }],
|
||||
});
|
||||
}
|
||||
|
||||
const text = typeof message.content === "string" ? message.content : "";
|
||||
if (text.length > 0) {
|
||||
output.push({
|
||||
type: RESPONSES_ITEM.MESSAGE,
|
||||
role: ROLE.ASSISTANT,
|
||||
content: [{ type: RESPONSES_ITEM.OUTPUT_TEXT, text, annotations: [] }],
|
||||
});
|
||||
}
|
||||
|
||||
for (const tc of message.tool_calls || []) {
|
||||
const fn = tc.function || {};
|
||||
const custom = customToolNames?.has(fn.name);
|
||||
output.push({
|
||||
type: custom ? RESPONSES_ITEM.CUSTOM_TOOL_CALL : RESPONSES_ITEM.FUNCTION_CALL,
|
||||
id: `${custom ? "ctc" : "fc"}_${tc.id || ""}`,
|
||||
call_id: tc.id || "",
|
||||
name: fn.name || "",
|
||||
...(custom
|
||||
? { input: extractCustomToolInput(fn.arguments) }
|
||||
: { arguments: typeof fn.arguments === "string" ? fn.arguments : JSON.stringify(fn.arguments || {}) }),
|
||||
});
|
||||
}
|
||||
|
||||
const usage = responseBody.usage || {};
|
||||
return {
|
||||
id: `resp_${responseBody.id || ""}`.replace(/^resp_chatcmpl-/, "resp_"),
|
||||
object: "response",
|
||||
created_at: responseBody.created || Math.floor(Date.now() / 1000),
|
||||
model: responseBody.model || "unknown",
|
||||
status: "completed",
|
||||
background: false,
|
||||
error: null,
|
||||
output,
|
||||
usage: {
|
||||
input_tokens: usage.prompt_tokens || usage.input_tokens || 0,
|
||||
output_tokens: usage.completion_tokens || usage.output_tokens || 0,
|
||||
total_tokens: usage.total_tokens || (usage.prompt_tokens || 0) + (usage.completion_tokens || 0),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse OpenAI-style SSE text into a single chat completion JSON.
|
||||
* Used when provider forces streaming but client wants non-streaming.
|
||||
@@ -136,6 +202,7 @@ export function parseSSEToOpenAIResponse(rawSSE, fallbackModel) {
|
||||
export async function handleForcedSSEToJson({
|
||||
providerResponse,
|
||||
sourceFormat,
|
||||
targetFormat,
|
||||
provider,
|
||||
model,
|
||||
body,
|
||||
@@ -147,6 +214,7 @@ export async function handleForcedSSEToJson({
|
||||
apiKey,
|
||||
clientRawRequest,
|
||||
onRequestSuccess,
|
||||
customToolNames,
|
||||
trackDone,
|
||||
appendLog,
|
||||
reqTag,
|
||||
@@ -170,8 +238,11 @@ export async function handleForcedSSEToJson({
|
||||
};
|
||||
|
||||
// Codex/Responses API SSE path
|
||||
// Branch on the UPSTREAM format (targetFormat = format we spoke to the provider in),
|
||||
// not the client format: a Responses-API client behind a chat-native forced-streaming
|
||||
// provider still receives chat SSE chunks, which must go through the standard path.
|
||||
const isCodexResponsesApi =
|
||||
isResponsesProvider(provider) || sourceFormat === FORMATS.OPENAI_RESPONSES;
|
||||
isResponsesProvider(provider) || targetFormat === FORMATS.OPENAI_RESPONSES;
|
||||
if (isCodexResponsesApi) {
|
||||
try {
|
||||
const jsonResponse = await convertResponsesStreamToJson(
|
||||
@@ -200,6 +271,11 @@ export async function handleForcedSSEToJson({
|
||||
}),
|
||||
);
|
||||
|
||||
// Same cache-inclusive total for the recorded detail, so the DB and the
|
||||
// client-facing usage can never disagree.
|
||||
const inTokensForLog = (usage.input_tokens || 0)
|
||||
+ (usage.cache_read_input_tokens || usage.cached_tokens || 0)
|
||||
+ (usage.cache_creation_input_tokens || 0);
|
||||
const { msgItem, textContent } = pickAssistantMessageForChatCompletion(
|
||||
jsonResponse.output,
|
||||
);
|
||||
@@ -209,9 +285,10 @@ export async function handleForcedSSEToJson({
|
||||
buildRequestDetail(
|
||||
{
|
||||
...ctx,
|
||||
apiKey,
|
||||
latency: { ttft: totalLatency, total: totalLatency },
|
||||
tokens: {
|
||||
prompt_tokens: usage.input_tokens || 0,
|
||||
prompt_tokens: inTokensForLog,
|
||||
completion_tokens: usage.output_tokens || 0,
|
||||
},
|
||||
response: {
|
||||
@@ -238,9 +315,22 @@ export async function handleForcedSSEToJson({
|
||||
};
|
||||
}
|
||||
|
||||
// Build client-format response
|
||||
const inTokens = usage.input_tokens || 0;
|
||||
// Build client-format response.
|
||||
// input_tokens EXCLUDES cached tokens on cache-capable upstreams, so summing
|
||||
// only input+output under-reports prompt_tokens — measured: 2012 reported
|
||||
// where the real prompt was ~5344 with 5332 served from cache. Fold the cache
|
||||
// counters in, and keep them visible in prompt_tokens_details so a client can
|
||||
// tell a cache hit from a small prompt.
|
||||
const cacheRead = usage.cache_read_input_tokens || usage.cached_tokens || 0;
|
||||
const cacheCreate = usage.cache_creation_input_tokens || 0;
|
||||
const inTokens = (usage.input_tokens || 0) + cacheRead + cacheCreate;
|
||||
const outTokens = usage.output_tokens || 0;
|
||||
const cacheDetails = (cacheRead > 0 || cacheCreate > 0)
|
||||
? {
|
||||
prompt_tokens_details: {
|
||||
...(cacheRead > 0 ? { cached_tokens: cacheRead } : {}),
|
||||
...(cacheCreate > 0 ? { cache_creation_tokens: cacheCreate } : {}) } }
|
||||
: {};
|
||||
let finalResp;
|
||||
|
||||
// Extract tool calls from Responses API output (function_call items)
|
||||
@@ -309,6 +399,7 @@ export async function handleForcedSSEToJson({
|
||||
prompt_tokens: inTokens,
|
||||
completion_tokens: outTokens,
|
||||
total_tokens: inTokens + outTokens,
|
||||
...cacheDetails,
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -384,23 +475,14 @@ export async function handleForcedSSEToJson({
|
||||
}),
|
||||
);
|
||||
|
||||
const totalLatency = Date.now() - requestStartTime;
|
||||
saveRequestDetail(
|
||||
buildRequestDetail(
|
||||
{
|
||||
...ctx,
|
||||
latency: { ttft: totalLatency, total: totalLatency },
|
||||
tokens: usage,
|
||||
response: {
|
||||
content: parsed.choices?.[0]?.message?.content || null,
|
||||
thinking: parsed.choices?.[0]?.message?.reasoning_content || null,
|
||||
finish_reason: parsed.choices?.[0]?.finish_reason || "unknown",
|
||||
},
|
||||
status: "success",
|
||||
},
|
||||
{ endpoint: clientRawRequest?.endpoint || null },
|
||||
),
|
||||
).catch(() => {});
|
||||
// Re-attach usage explicitly. This handler already HAS the correct usage — it is
|
||||
// the same object written to the usage DB, and for a cached Claude request that DB
|
||||
// row reads cache_read_input_tokens: 11022 — yet the client was observed receiving
|
||||
// no usage field at all (verified 2026-08-04 with a fingerprinted payload matched
|
||||
// on both sides). Whatever drops it between assembly and serialisation, the client
|
||||
// must not be left unable to account for its own token spend: a caller cannot tell
|
||||
// a 90%-cached request from a cheap one without this.
|
||||
if (usage && Object.keys(usage).length > 0) parsed.usage = usage;
|
||||
|
||||
// Strip reasoning_content only when content is non-empty.
|
||||
// When content is empty (e.g. thinking models that used all tokens for reasoning),
|
||||
@@ -414,9 +496,19 @@ export async function handleForcedSSEToJson({
|
||||
}
|
||||
}
|
||||
|
||||
// A Responses-format client (e.g. Codex) forced this provider to stream,
|
||||
// but wants JSON back. parseSSEToOpenAIResponse yields a Chat Completions
|
||||
// body; convert it to the Responses `output` shape so tool_calls are not
|
||||
// lost on the non-streaming return path. Inlined (not imported from
|
||||
// nonStreamingHandler.js) to avoid a circular import: nonStreamingHandler
|
||||
// already imports parseSSEToOpenAIResponse from this module.
|
||||
const finalBody = sourceFormat === FORMATS.OPENAI_RESPONSES
|
||||
? chatCompletionToResponses(parsed, customToolNames)
|
||||
: parsed;
|
||||
|
||||
return {
|
||||
success: true,
|
||||
response: new Response(JSON.stringify(parsed), {
|
||||
response: new Response(JSON.stringify(finalBody), {
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"Access-Control-Allow-Origin": "*",
|
||||
|
||||
@@ -38,6 +38,7 @@ function buildTransformStream({
|
||||
userAgent,
|
||||
reqLogger,
|
||||
toolNameMap,
|
||||
customToolNames,
|
||||
model,
|
||||
connectionId,
|
||||
body,
|
||||
@@ -68,6 +69,7 @@ function buildTransformStream({
|
||||
body,
|
||||
onStreamComplete,
|
||||
apiKey,
|
||||
customToolNames,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -83,6 +85,7 @@ function buildTransformStream({
|
||||
body,
|
||||
onStreamComplete,
|
||||
apiKey,
|
||||
customToolNames,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -118,6 +121,7 @@ export async function handleStreamingResponse({
|
||||
onRequestSuccess,
|
||||
reqLogger,
|
||||
toolNameMap,
|
||||
customToolNames,
|
||||
streamController,
|
||||
onStreamComplete,
|
||||
streamDetailId,
|
||||
@@ -200,6 +204,7 @@ export async function handleStreamingResponse({
|
||||
userAgent,
|
||||
reqLogger,
|
||||
toolNameMap,
|
||||
customToolNames,
|
||||
model,
|
||||
connectionId,
|
||||
body,
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
import createOpenAIEmbeddingAdapter from "./openai.js";
|
||||
import gemini from "./gemini.js";
|
||||
import openaiCompatNode from "./openaiCompatNode.js";
|
||||
import selfhostedEmbedding from "./selfhostedEmbedding.js";
|
||||
|
||||
const OPENAI_COMPAT_PROVIDERS = [
|
||||
"openai", "openrouter", "mistral", "voyage-ai", "fireworks",
|
||||
@@ -13,6 +14,12 @@ const ADAPTERS = {
|
||||
...Object.fromEntries(OPENAI_COMPAT_PROVIDERS.map((id) => [id, createOpenAIEmbeddingAdapter(id)])),
|
||||
gemini,
|
||||
google_ai_studio: gemini,
|
||||
// Self-hosted reads creds.providerSpecificData.baseUrl (one provider, many
|
||||
// servers) — but via its OWN adapter, not openaiCompatNode: that one falls back
|
||||
// to api.openai.com when no baseUrl is set, which under a provider called
|
||||
// "Self-hosted Embedding" means silently shipping the input and API key to
|
||||
// OpenAI. selfhostedEmbedding refuses instead.
|
||||
"selfhosted-embedding": selfhostedEmbedding,
|
||||
};
|
||||
|
||||
export function getEmbeddingAdapter(provider) {
|
||||
|
||||
46
open-sse/handlers/embeddingProviders/selfhostedEmbedding.js
Normal file
46
open-sse/handlers/embeddingProviders/selfhostedEmbedding.js
Normal file
@@ -0,0 +1,46 @@
|
||||
// Self-hosted embeddings — like openaiCompatNode, but the baseUrl is REQUIRED.
|
||||
//
|
||||
// openaiCompatNode falls back to https://api.openai.com/v1 when a connection
|
||||
// carries no providerSpecificData.baseUrl. For a custom NODE that default is
|
||||
// defensible: the node was created by pointing at some OpenAI-compatible URL, and
|
||||
// OpenAI is the archetype. For a provider whose entire purpose is "my own
|
||||
// server", it is actively harmful — a connection saved without a baseUrl sends
|
||||
// the INPUT TEXT and the API KEY to OpenAI, silently, under a provider named
|
||||
// "Self-hosted Embedding".
|
||||
//
|
||||
// Observed exactly that with a placeholder connection (2026-08-04):
|
||||
//
|
||||
// [selfhosted-embedding/embedding] [401]: Incorrect API key provided: abc.
|
||||
// You can find your API key at https://platform.openai.com/account/api-keys.
|
||||
//
|
||||
// The key "abc" was typed as a throwaway for a LOCAL server and left the network.
|
||||
// A self-hosted provider must never have a cloud fallback, so this one refuses
|
||||
// instead: no baseUrl means a configuration error, reported as such.
|
||||
import createOpenAIEmbeddingAdapter from "./openai.js";
|
||||
|
||||
const baseAdapter = createOpenAIEmbeddingAdapter("openai");
|
||||
|
||||
export class MissingBaseUrlError extends Error {
|
||||
constructor() {
|
||||
super(
|
||||
"Self-hosted Embedding needs an endpoint: set this connection's baseUrl to " +
|
||||
"the OpenAI base URL of your server, e.g. http://host:8080/v1 (note the /v1 — " +
|
||||
"\"/embeddings\" is appended to it). Refusing to fall back to api.openai.com, " +
|
||||
"which would send your input and API key to OpenAI."
|
||||
);
|
||||
this.name = "MissingBaseUrlError";
|
||||
this.isConfigError = true;
|
||||
}
|
||||
}
|
||||
|
||||
export default {
|
||||
...baseAdapter,
|
||||
buildUrl: (_model, creds) => {
|
||||
const rawBaseUrl = creds?.providerSpecificData?.baseUrl;
|
||||
if (!rawBaseUrl || !String(rawBaseUrl).trim()) throw new MissingBaseUrlError();
|
||||
// Accept either the OpenAI base or a full embeddings URL, so a value pasted
|
||||
// from a curl example works as well as one typed from the help text.
|
||||
const baseUrl = String(rawBaseUrl).trim().replace(/\/$/, "").replace(/\/embeddings$/, "");
|
||||
return `${baseUrl}/embeddings`;
|
||||
},
|
||||
};
|
||||
@@ -1,5 +1,5 @@
|
||||
import { createErrorResult, parseUpstreamError, formatProviderError } from "../utils/error.js";
|
||||
import { HTTP_STATUS } from "../config/runtimeConfig.js";
|
||||
import { HTTP_STATUS, FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js";
|
||||
import { getExecutor } from "../executors/index.js";
|
||||
import { refreshWithRetry } from "../services/tokenRefresh.js";
|
||||
import { getEmbeddingAdapter } from "./embeddingProviders/index.js";
|
||||
@@ -38,13 +38,24 @@ export async function handleEmbeddingsCore({
|
||||
}
|
||||
|
||||
const ctx = { input };
|
||||
const url = adapter.buildUrl(model, credentials, ctx);
|
||||
const headers = adapter.buildHeaders(credentials, ctx);
|
||||
const requestBody = adapter.buildBody(model, {
|
||||
input,
|
||||
encoding_format: body.encoding_format || "float",
|
||||
dimensions: body.dimensions,
|
||||
});
|
||||
// buildUrl/buildHeaders/buildBody were called bare. An adapter that rejects a
|
||||
// misconfigured connection — selfhosted-embedding throws when no baseUrl is set
|
||||
// rather than silently falling back to api.openai.com — would have escaped this
|
||||
// function uncaught, surfacing as a 500 or a request that never settles. A
|
||||
// configuration mistake is a 400 with the reason in it.
|
||||
let url, headers, requestBody;
|
||||
try {
|
||||
url = adapter.buildUrl(model, credentials, ctx);
|
||||
headers = adapter.buildHeaders(credentials, ctx);
|
||||
requestBody = adapter.buildBody(model, {
|
||||
input,
|
||||
encoding_format: body.encoding_format || "float",
|
||||
dimensions: body.dimensions,
|
||||
});
|
||||
} catch (error) {
|
||||
log?.debug?.("EMBEDDINGS", `Request build failed: ${error.message}`);
|
||||
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `[${provider}/${model}] ${error.message}`);
|
||||
}
|
||||
|
||||
log?.debug?.("EMBEDDINGS", `${provider.toUpperCase()} | ${model} | input_type=${Array.isArray(input) ? `array[${input.length}]` : "string"}`);
|
||||
|
||||
@@ -54,6 +65,9 @@ export async function handleEmbeddingsCore({
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify(requestBody),
|
||||
...(typeof AbortSignal?.timeout === "function"
|
||||
? { signal: AbortSignal.timeout(FETCH_CONNECT_TIMEOUT_MS) }
|
||||
: {}),
|
||||
});
|
||||
} catch (error) {
|
||||
const errMsg = formatProviderError(error, provider, model, HTTP_STATUS.BAD_GATEWAY);
|
||||
|
||||
@@ -170,9 +170,17 @@ export async function handleSttCore({ provider, model, formData, credentials, st
|
||||
const file = formData.get("file");
|
||||
if (!file) return createErrorResult(HTTP_STATUS.BAD_REQUEST, "Missing required field: file");
|
||||
|
||||
const cfg = sttConfig;
|
||||
let cfg = sttConfig;
|
||||
if (!cfg) return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Provider '${provider}' does not support STT`);
|
||||
|
||||
// Per-connection endpoint override. Registry entries carry a fixed baseUrl,
|
||||
// which is right for a named cloud service but useless for a self-hosted one
|
||||
// whose address only the operator knows. Opt-in: absent unless the connection
|
||||
// sets it, so cloud providers are untouched. Mirrors the custom embedding
|
||||
// providers, which already resolve baseUrl the same way.
|
||||
const overrideUrl = credentials?.providerSpecificData?.baseUrl;
|
||||
if (overrideUrl) cfg = { ...cfg, baseUrl: String(overrideUrl).replace(/\/+$/, "") };
|
||||
|
||||
const token = cfg.authType === "none" ? null : (credentials?.apiKey || credentials?.accessToken);
|
||||
if (cfg.authType !== "none" && !token) {
|
||||
return createErrorResult(HTTP_STATUS.UNAUTHORIZED, `No credentials for STT provider: ${provider}`);
|
||||
|
||||
@@ -48,16 +48,16 @@ function createTtsResponse(base64Audio, format, responseFormat) {
|
||||
*
|
||||
* @returns {Promise<{success, response, status?, error?}>}
|
||||
*/
|
||||
export async function handleTtsCore({ provider, model, input, credentials, responseFormat = "mp3", language }) {
|
||||
export async function handleTtsCore({ provider, model, input, credentials, responseFormat = "mp3", language, style }) {
|
||||
if (!input?.trim()) {
|
||||
return createErrorResult(HTTP_STATUS.BAD_REQUEST, "Missing required field: input");
|
||||
}
|
||||
|
||||
try {
|
||||
// Special-case adapters (google-tts, edge-tts, local-device, elevenlabs, openai, openrouter, gemini)
|
||||
// Special-case adapters (google-tts, edge-tts, local-device, elevenlabs, openai, openrouter, gemini, xiaomi-mimo)
|
||||
const adapter = getTtsAdapter(provider);
|
||||
if (adapter) {
|
||||
const result = await adapter.synthesize(input.trim(), model, credentials, responseFormat, { language });
|
||||
const result = await adapter.synthesize(input.trim(), model, credentials, responseFormat, { language, style });
|
||||
// Adapter may return a full {success, response} (legacy) or {base64, format}
|
||||
if (result.success !== undefined) return result;
|
||||
return createTtsResponse(result.base64, result.format, responseFormat);
|
||||
|
||||
@@ -6,6 +6,8 @@ import elevenlabs, { fetchElevenLabsVoices } from "./elevenlabs.js";
|
||||
import openai from "./openai.js";
|
||||
import openrouter from "./openrouter.js";
|
||||
import gemini, { fetchGeminiVoices } from "./gemini.js";
|
||||
import xiaomiMimo from "./xiaomi-mimo.js";
|
||||
import selfhostedTts from "./selfhostedTts.js";
|
||||
import { FORMAT_HANDLERS } from "./genericFormats.js";
|
||||
import { parseModelVoice } from "./_base.js";
|
||||
|
||||
@@ -18,6 +20,8 @@ const SPECIAL_ADAPTERS = {
|
||||
openai,
|
||||
openrouter,
|
||||
gemini,
|
||||
"xiaomi-mimo": xiaomiMimo,
|
||||
"selfhosted-tts": selfhostedTts,
|
||||
};
|
||||
|
||||
export function getTtsAdapter(provider) {
|
||||
|
||||
69
open-sse/handlers/ttsProviders/selfhostedTts.js
Normal file
69
open-sse/handlers/ttsProviders/selfhostedTts.js
Normal file
@@ -0,0 +1,69 @@
|
||||
// Self-hosted OpenAI-compatible TTS — POST {baseUrl}/v1/audio/speech.
|
||||
//
|
||||
// A SPECIAL_ADAPTER rather than a genericFormats handler on purpose: the generic
|
||||
// dispatcher resolves baseUrl from the static registry entry
|
||||
// (`synthesizeViaConfig` reads `cfg.baseUrl`) and never looks at the connection,
|
||||
// which is exactly the limitation this provider exists to lift.
|
||||
import { Buffer } from "node:buffer";
|
||||
|
||||
const DEFAULT_BASE_URL = "http://localhost:8880";
|
||||
const DEFAULT_MODEL = "kokoro";
|
||||
const DEFAULT_VOICE = "af_heart";
|
||||
|
||||
export default {
|
||||
async synthesize(text, model, credentials, responseFormat = "mp3") {
|
||||
// Accept either providerSpecificData.baseUrl (how the custom embedding and
|
||||
// STT providers carry it) or a bare credentials.baseUrl (how the OpenAI TTS
|
||||
// adapter does), so a connection configured either way works.
|
||||
const raw = credentials?.providerSpecificData?.baseUrl || credentials?.baseUrl || DEFAULT_BASE_URL;
|
||||
// Tolerate a baseUrl given as the full endpoint or with a trailing /v1 —
|
||||
// both are natural things to paste, and silently double-appending the path
|
||||
// would 404 with nothing pointing at the cause.
|
||||
const base = String(raw)
|
||||
.replace(/\/+$/, "")
|
||||
.replace(/\/v1\/audio\/speech$/, "")
|
||||
.replace(/\/v1$/, "");
|
||||
|
||||
// The provider prefix is already stripped by getModelInfo, so `model` here is
|
||||
// "kokoro" or "kokoro/af_heart" — NOT "selfhosted-tts/...".
|
||||
//
|
||||
// A bare value is the MODEL, not the voice. The OpenAI adapter reads a bare
|
||||
// value as a voice, which is right for a service whose model is fixed
|
||||
// ("tts-1") and whose voice varies — but wrong here, where the model is the
|
||||
// variable part. Treating it as a voice sent voice="kokoro" upstream and
|
||||
// Kokoro answered 400, so `selfhosted-tts/kokoro` — the obvious way to
|
||||
// address this provider — was the one form that did not work (verified
|
||||
// against a live Kokoro through 9router, 2026-08-03).
|
||||
let ttsModel = DEFAULT_MODEL;
|
||||
let voice = DEFAULT_VOICE;
|
||||
if (model) {
|
||||
const parts = String(model).split("/").filter(Boolean);
|
||||
if (parts.length >= 2) {
|
||||
ttsModel = parts[0];
|
||||
voice = parts.slice(1).join("/");
|
||||
} else if (parts.length === 1) {
|
||||
ttsModel = parts[0];
|
||||
}
|
||||
}
|
||||
|
||||
const res = await fetch(`${base}/v1/audio/speech`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
...(credentials?.apiKey ? { Authorization: `Bearer ${credentials.apiKey}` } : {}),
|
||||
},
|
||||
body: JSON.stringify({
|
||||
model: ttsModel,
|
||||
voice,
|
||||
input: text,
|
||||
response_format: responseFormat,
|
||||
}),
|
||||
});
|
||||
if (!res.ok) {
|
||||
const err = await res.json().catch(() => ({}));
|
||||
throw new Error(err?.error?.message || `Self-hosted TTS failed: ${res.status}`);
|
||||
}
|
||||
const buf = await res.arrayBuffer();
|
||||
return { base64: Buffer.from(buf).toString("base64"), format: responseFormat };
|
||||
},
|
||||
};
|
||||
65
open-sse/handlers/ttsProviders/xiaomi-mimo.js
Normal file
65
open-sse/handlers/ttsProviders/xiaomi-mimo.js
Normal file
@@ -0,0 +1,65 @@
|
||||
// Xiaomi MiMo TTS — via OpenAI-compatible chat completions (non-streaming).
|
||||
// Docs: https://mimo.mi.com/docs/zh-CN/quick-start/usage-guide/audio/speech-synthesis-v2.5
|
||||
// Message contract: target text in `role: assistant` content, style/voice
|
||||
// instructions in `role: user` content. Voice is selected via the top-level
|
||||
// `audio.voice` field (NOT embedded in the model name).
|
||||
import { parseModelVoice } from "./_base.js";
|
||||
|
||||
const DEFAULT_MODEL = "mimo-v2.5-tts";
|
||||
const DEFAULT_VOICE = "mimo_default";
|
||||
|
||||
export default {
|
||||
synthesize(text, model, credentials, responseFormat, { style, language } = {}) {
|
||||
if (!credentials?.apiKey) throw new Error("xiaomi-mimo API key required");
|
||||
return synthesizeMiMo(text, model, credentials.apiKey, style, language);
|
||||
},
|
||||
};
|
||||
|
||||
export async function synthesizeMiMo(text, model, apiKey, style, language) {
|
||||
const { modelId, voiceId } = parseModelVoice(model, DEFAULT_MODEL, DEFAULT_VOICE, [DEFAULT_MODEL]);
|
||||
|
||||
// Language and style are soft instructions → prepend as a role:user message.
|
||||
// MiMo auto-detects the spoken language of the text; the hint only nudges it
|
||||
// (e.g. "Speak in English.") and is independent of the chosen voice.
|
||||
const instructions = [];
|
||||
if (language) instructions.push(`Speak in ${language}.`);
|
||||
if (style) instructions.push(style);
|
||||
|
||||
const messages = [{ role: "assistant", content: text }];
|
||||
if (instructions.length) messages.unshift({ role: "user", content: instructions.join(" ") });
|
||||
|
||||
const res = await fetch("https://api.xiaomimimo.com/v1/chat/completions", {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": `Bearer ${apiKey}`,
|
||||
},
|
||||
body: JSON.stringify({
|
||||
model: modelId,
|
||||
stream: false,
|
||||
messages,
|
||||
audio: {
|
||||
format: "wav",
|
||||
voice: voiceId || DEFAULT_VOICE,
|
||||
},
|
||||
}),
|
||||
});
|
||||
|
||||
const rawText = await res.text();
|
||||
let data = {};
|
||||
if (rawText) {
|
||||
try { data = JSON.parse(rawText); } catch { data = {}; }
|
||||
}
|
||||
|
||||
if (!res.ok) {
|
||||
throw new Error(data?.error?.message || rawText || `MiMo TTS error (${res.status})`);
|
||||
}
|
||||
|
||||
const audio = data?.choices?.[0]?.message?.audio?.data;
|
||||
if (!audio) throw new Error(data?.error?.message || "MiMo TTS returned no audio");
|
||||
|
||||
return {
|
||||
base64: audio,
|
||||
format: data?.choices?.[0]?.message?.audio?.format || "wav",
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user