update: sync local changes with latest features
This commit is contained in:
@@ -28,7 +28,7 @@ import { compressMessages, formatRtkLog } from "../rtk/index.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, cavemanEnabled, cavemanLevel, sourceFormatOverride, providerThinking }) {
|
||||
export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, onMidStreamError, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, cavemanEnabled, cavemanLevel, sourceFormatOverride, providerThinking }) {
|
||||
const { provider, model } = modelInfo;
|
||||
const requestStartTime = Date.now();
|
||||
|
||||
@@ -185,13 +185,14 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
}
|
||||
|
||||
// Execute request
|
||||
let providerResponse, providerUrl, providerHeaders, finalBody;
|
||||
let providerResponse, providerUrl, providerHeaders, finalBody, midStreamError;
|
||||
try {
|
||||
const result = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions });
|
||||
providerResponse = result.response;
|
||||
providerUrl = result.url;
|
||||
providerHeaders = result.headers;
|
||||
finalBody = result.transformedBody;
|
||||
midStreamError = result.midStreamError;
|
||||
reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody);
|
||||
} catch (error) {
|
||||
trackPendingRequest(model, provider, connectionId, false, true);
|
||||
@@ -258,7 +259,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
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, midStreamError, onMidStreamError };
|
||||
const appendLog = (extra) => appendRequestLog({ model, provider, connectionId, ...extra }).catch(() => { });
|
||||
const trackDone = () => trackPendingRequest(model, provider, connectionId, false);
|
||||
|
||||
@@ -275,8 +276,8 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
|
||||
return result;
|
||||
}
|
||||
|
||||
// Streaming response
|
||||
const { onStreamComplete } = buildOnStreamComplete({ ...sharedCtx });
|
||||
// Streaming response (midStreamError + onMidStreamError already in sharedCtx)
|
||||
const { onStreamComplete } = buildOnStreamComplete(sharedCtx);
|
||||
return handleStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat, userAgent, reqLogger, toolNameMap, streamController, onStreamComplete });
|
||||
}
|
||||
|
||||
|
||||
@@ -132,7 +132,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, midStreamError, onMidStreamError }) {
|
||||
trackDone();
|
||||
const contentType = providerResponse.headers.get("content-type") || "";
|
||||
let responseBody;
|
||||
@@ -144,6 +144,21 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m
|
||||
appendLog({ status: `FAILED ${HTTP_STATUS.BAD_GATEWAY}` });
|
||||
return createErrorResult(HTTP_STATUS.BAD_GATEWAY, "Invalid SSE response for non-streaming request");
|
||||
}
|
||||
|
||||
// Check for mid-stream errors detected during SSE parsing (e.g. Qoder quota exceeded).
|
||||
// The upstream returned HTTP 200 but the SSE envelope contained a non-200 statusCodeValue.
|
||||
// Without this check, the error message would be silently served as normal content with HTTP 200.
|
||||
if (midStreamError?.error) {
|
||||
const err = midStreamError.error;
|
||||
if (typeof onMidStreamError === "function") {
|
||||
onMidStreamError(err).catch(e => {
|
||||
console.error("[MidStreamError] Failed to apply cooldown:", e.message);
|
||||
});
|
||||
}
|
||||
appendLog({ status: `FAILED ${err.status || 429}` });
|
||||
return createErrorResult(err.status || 429, err.message || "Upstream error detected in SSE stream");
|
||||
}
|
||||
|
||||
responseBody = parsed;
|
||||
} else {
|
||||
try {
|
||||
|
||||
@@ -75,8 +75,10 @@ export function handleStreamingResponse({ providerResponse, provider, model, sou
|
||||
|
||||
/**
|
||||
* 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, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest }) {
|
||||
export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest, midStreamError, onMidStreamError }) {
|
||||
const streamDetailId = `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`;
|
||||
|
||||
const onStreamComplete = (contentObj, usage, ttftAt) => {
|
||||
@@ -87,6 +89,15 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r
|
||||
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,
|
||||
latency,
|
||||
@@ -95,7 +106,7 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r
|
||||
providerRequest: finalBody || translatedBody || null,
|
||||
providerResponse: safeContent,
|
||||
response: { content: safeContent, thinking: safeThinking, type: "streaming" },
|
||||
status: "success"
|
||||
status: midStreamError?.error ? "error" : "success"
|
||||
}, { id: streamDetailId })).catch(err => {
|
||||
console.error("[RequestDetail] Failed to update streaming content:", err.message);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user