Feat : Auto restart after crash

This commit is contained in:
decolua
2026-03-14 09:37:29 +07:00
parent 549223c8cf
commit adae2605bf
31 changed files with 340 additions and 189 deletions

View File

@@ -169,7 +169,9 @@ export class AntigravityExecutor extends BaseExecutor {
let lastError = null;
let lastStatus = 0;
const MAX_AUTO_RETRIES = 3;
const MAX_RETRY_AFTER_RETRIES = 3;
const retryAttemptsByUrl = {}; // Track retry attempts per URL
const retryAfterAttemptsByUrl = {}; // Track Retry-After retries per URL
for (let urlIndex = 0; urlIndex < fallbackCount; urlIndex++) {
const url = this.buildUrl(model, stream, urlIndex);
@@ -177,10 +179,13 @@ export class AntigravityExecutor extends BaseExecutor {
const sessionId = transformedBody.request?.sessionId;
const headers = this.buildHeaders(credentials, stream, sessionId);
// Initialize retry counter for this URL
// Initialize retry counters for this URL
if (!retryAttemptsByUrl[urlIndex]) {
retryAttemptsByUrl[urlIndex] = 0;
}
if (!retryAfterAttemptsByUrl[urlIndex]) {
retryAfterAttemptsByUrl[urlIndex] = 0;
}
try {
const response = await proxyAwareFetch(url, {
@@ -206,8 +211,9 @@ export class AntigravityExecutor extends BaseExecutor {
}
}
if (retryMs && retryMs <= MAX_RETRY_AFTER_MS) {
log?.debug?.("RETRY", `${response.status} with Retry-After: ${Math.ceil(retryMs / 1000)}s, waiting...`);
if (retryMs && retryMs <= MAX_RETRY_AFTER_MS && retryAfterAttemptsByUrl[urlIndex] < MAX_RETRY_AFTER_RETRIES) {
retryAfterAttemptsByUrl[urlIndex]++;
log?.debug?.("RETRY", `${response.status} with Retry-After: ${Math.ceil(retryMs / 1000)}s, waiting... (${retryAfterAttemptsByUrl[urlIndex]}/${MAX_RETRY_AFTER_RETRIES})`);
await new Promise(resolve => setTimeout(resolve, retryMs));
urlIndex--;
continue;

View File

@@ -155,13 +155,30 @@ export class CursorExecutor extends BaseExecutor {
throw new Error("http2 module not available");
}
const HTTP2_TIMEOUT_MS = 60000; // 60s max — prevent hung sessions
return new Promise((resolve, reject) => {
const urlObj = new URL(url);
const client = http2.connect(`https://${urlObj.host}`);
const chunks = [];
let responseHeaders = {};
let settled = false;
client.on("error", reject);
// Ensure client is always closed on settle
const finish = (fn) => (...args) => {
if (settled) return;
settled = true;
clearTimeout(hangTimeout);
client.close();
fn(...args);
};
// Hard timeout: close session if server never responds
const hangTimeout = setTimeout(finish(() => {
reject(new Error("HTTP/2 request timed out"));
}), HTTP2_TIMEOUT_MS);
client.on("error", finish(reject));
const req = client.request({
":method": "POST",
@@ -173,25 +190,18 @@ export class CursorExecutor extends BaseExecutor {
req.on("response", (hdrs) => { responseHeaders = hdrs; });
req.on("data", (chunk) => { chunks.push(chunk); });
req.on("end", () => {
client.close();
req.on("end", finish(() => {
resolve({
status: responseHeaders[":status"],
headers: responseHeaders,
body: Buffer.concat(chunks)
});
});
req.on("error", (err) => {
client.close();
reject(err);
});
}));
req.on("error", finish(reject));
if (signal) {
signal.addEventListener("abort", () => {
req.close();
client.close();
reject(new Error("Request aborted"));
});
const onAbort = finish(() => reject(new Error("Request aborted")));
signal.addEventListener("abort", onAbort, { once: true });
}
req.write(body);

View File

@@ -198,6 +198,9 @@ export class GithubExecutor extends BaseExecutor {
}
});
if (!response.body) {
return { response: new Response("", { status: response.status, headers: response.headers }), url, headers, transformedBody };
}
const convertedStream = response.body.pipeThrough(transformStream);
return {

View File

@@ -345,6 +345,9 @@ export class KiroExecutor extends BaseExecutor {
});
// Pipe response body through transform stream
if (!response.body) {
return new Response("data: [DONE]\n\n", { status: response.status, headers: { "Content-Type": "text/event-stream" } });
}
const transformedStream = response.body.pipeThrough(transformStream);
return new Response(transformedStream, {

View File

@@ -56,6 +56,10 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
log?.debug?.("FORMAT", `${sourceFormat} → ${targetFormat} | stream=${stream}`);
let translatedBody = translateRequest(sourceFormat, targetFormat, model, body, stream, credentials, provider, reqLogger);
if (!translatedBody) {
trackPendingRequest(model, provider, connectionId, false, true);
return createErrorResult(HTTP_STATUS.BAD_REQUEST, `Failed to translate request for ${sourceFormat} → ${targetFormat}`);
}
const toolNameMap = translatedBody._toolNameMap;
delete translatedBody._toolNameMap;
translatedBody.model = model;
@@ -137,17 +141,23 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
// Handle 401/403 - try token refresh
if (providerResponse.status === HTTP_STATUS.UNAUTHORIZED || providerResponse.status === HTTP_STATUS.FORBIDDEN) {
const newCredentials = await refreshWithRetry(() => executor.refreshCredentials(credentials, log), 3, log);
if (newCredentials?.accessToken || newCredentials?.copilotToken) {
log?.info?.("TOKEN", `${provider.toUpperCase()} | refreshed`);
Object.assign(credentials, newCredentials);
if (onCredentialsRefreshed) await onCredentialsRefreshed(newCredentials);
try {
const retryResult = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions });
if (retryResult.response.ok) { providerResponse = retryResult.response; providerUrl = retryResult.url; }
} catch { log?.warn?.("TOKEN", `${provider.toUpperCase()} | retry after refresh failed`); }
} else {
log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`);
try {
const newCredentials = await refreshWithRetry(() => executor.refreshCredentials(credentials, log), 3, log);
if (newCredentials?.accessToken || newCredentials?.copilotToken) {
log?.info?.("TOKEN", `${provider.toUpperCase()} | refreshed`);
Object.assign(credentials, newCredentials);
if (onCredentialsRefreshed) {
try { await onCredentialsRefreshed(newCredentials); } catch (e) { log?.warn?.("TOKEN", `onCredentialsRefreshed failed: ${e.message}`); }
}
try {
const retryResult = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions });
if (retryResult.response.ok) { providerResponse = retryResult.response; providerUrl = retryResult.url; }
} catch { log?.warn?.("TOKEN", `${provider.toUpperCase()} | retry after refresh failed`); }
} else {
log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`);
}
} catch (e) {
log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh threw: ${e.message}`);
}
}
@@ -182,12 +192,14 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred
// Provider forced streaming but client wants JSON
if (!clientRequestedStreaming && providerRequiresStreaming) {
const result = await handleForcedSSEToJson({ ...sharedCtx, providerResponse, sourceFormat, trackDone, appendLog });
if (result) return result;
if (result) { streamController.handleComplete(); return result; }
}
// True non-streaming response
if (!stream) {
return handleNonStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat, reqLogger, trackDone, appendLog });
const result = await handleNonStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat, reqLogger, trackDone, appendLog });
streamController.handleComplete();
return result;
}
// Streaming response

View File

@@ -79,7 +79,8 @@ export function translateNonStreamingResponse(responseBody, targetFormat, source
for (const block of responseBody.content) {
if (block.type === "text") {
// Strip markdown code block markers (e.g. kimi wraps JSON in ```json...```)
const text = block.text.replace(/^\s*```\s*json\s*\n?/i, "").replace(/\n?\s*```\s*$/i, "");
const raw = block.text ?? "";
const text = raw.replace(/^\s*```\s*json\s*\n?/i, "").replace(/\n?\s*```\s*$/i, "");
textContent += text;
} else if (block.type === "thinking") thinkingContent += block.thinking || "";
else if (block.type === "tool_use") {

View File

@@ -76,6 +76,8 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r
ttft: ttftAt ? ttftAt - requestStartTime : Date.now() - requestStartTime,
total: Date.now() - requestStartTime
};
const safeContent = contentObj?.content || "[Empty streaming response]";
const safeThinking = contentObj?.thinking || null;
saveRequestDetail(buildRequestDetail({
provider, model, connectionId,
@@ -83,8 +85,8 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r
tokens: usage || { prompt_tokens: 0, completion_tokens: 0 },
request: extractRequestConfig(body, stream),
providerRequest: finalBody || translatedBody || null,
providerResponse: contentObj.content || "[Empty streaming response]",
response: { content: contentObj.content || "[Empty streaming response]", thinking: contentObj.thinking || null, type: "streaming" },
providerResponse: safeContent,
response: { content: safeContent, thinking: safeThinking, type: "streaming" },
status: "success"
}, { id: streamDetailId })).catch(err => {
console.error("[RequestDetail] Failed to update streaming content:", err.message);

View File

@@ -38,8 +38,15 @@ export async function handleComboChat({ body, models, handleSingleModel, log })
const modelStr = models[i];
log.info("COMBO", `Trying model ${i + 1}/${models.length}: ${modelStr}`);
const result = await handleSingleModel(body, modelStr);
let result;
try {
result = await handleSingleModel(body, modelStr);
} catch (e) {
lastError = `${modelStr}: ${e.message}`;
log.warn("COMBO", `Model threw exception, trying next`, { model: modelStr, error: e.message });
continue;
}
// Success or client error - return response
if (result.ok || result.status < 500) {
return result;

View File

@@ -254,7 +254,9 @@ async function onboardUser(accessToken, tierID, externalSignal) {
console.warn(`[ProjectId] onboardUser failed after ${MAX_ATTEMPTS} attempts: ${error.message}`);
return null;
}
throw error;
// Continue to next attempt instead of throwing (which would skip remaining retries)
console.warn(`[ProjectId] onboardUser attempt ${attempt} failed: ${error.message}, retrying...`);
await new Promise(resolve => setTimeout(resolve, 2000));
} finally {
clearTimeout(timeoutId);
externalSignal?.removeEventListener("abort", forwardAbort);

View File

@@ -255,7 +255,7 @@ export function buildProviderHeaders(provider, credentials, stream = true, body
}
break;
case "github":
case "github": {
// GitHub Copilot requires special headers to mimic VSCode
// Prioritize copilotToken from providerSpecificData, fallback to accessToken
const githubToken = credentials.copilotToken || credentials.accessToken;
@@ -279,6 +279,7 @@ export function buildProviderHeaders(provider, credentials, stream = true, body
headers["X-Initiator"] = "user";
headers["Accept"] = "application/json";
break;
}
case "codex":
case "qwen":

View File

@@ -69,83 +69,67 @@ export async function refreshAccessToken(provider, refreshToken, credentials, lo
* Specialized refresh for Claude OAuth tokens
*/
export async function refreshClaudeOAuthToken(refreshToken, log) {
const response = await fetch(OAUTH_ENDPOINTS.anthropic.token, {
method: "POST",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
},
body: JSON.stringify({
grant_type: "refresh_token",
refresh_token: refreshToken,
client_id: PROVIDERS.claude.clientId,
}),
});
if (!response.ok) {
const errorText = await response.text();
log?.error?.("TOKEN_REFRESH", "Failed to refresh Claude OAuth token", {
status: response.status,
error: errorText,
try {
const response = await fetch(OAUTH_ENDPOINTS.anthropic.token, {
method: "POST",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
},
body: JSON.stringify({
grant_type: "refresh_token",
refresh_token: refreshToken,
client_id: PROVIDERS.claude.clientId,
}),
});
if (!response.ok) {
const errorText = await response.text();
log?.error?.("TOKEN_REFRESH", "Failed to refresh Claude OAuth token", { status: response.status, error: errorText });
return null;
}
const tokens = await response.json();
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Claude OAuth token", { hasNewAccessToken: !!tokens.access_token, expiresIn: tokens.expires_in });
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
} catch (error) {
log?.error?.("TOKEN_REFRESH", `Network error refreshing Claude token: ${error.message}`);
return null;
}
const tokens = await response.json();
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Claude OAuth token", {
hasNewAccessToken: !!tokens.access_token,
hasNewRefreshToken: !!tokens.refresh_token,
expiresIn: tokens.expires_in,
});
return {
accessToken: tokens.access_token,
refreshToken: tokens.refresh_token || refreshToken,
expiresIn: tokens.expires_in,
};
}
/**
* Specialized refresh for Google providers (Gemini, Antigravity)
*/
export async function refreshGoogleToken(refreshToken, clientId, clientSecret, log) {
const response = await fetch(OAUTH_ENDPOINTS.google.token, {
method: "POST",
headers: {
"Content-Type": "application/x-www-form-urlencoded",
Accept: "application/json",
},
body: new URLSearchParams({
grant_type: "refresh_token",
refresh_token: refreshToken,
client_id: clientId,
client_secret: clientSecret,
}),
});
if (!response.ok) {
const errorText = await response.text();
log?.error?.("TOKEN_REFRESH", "Failed to refresh Google token", {
status: response.status,
error: errorText,
try {
const response = await fetch(OAUTH_ENDPOINTS.google.token, {
method: "POST",
headers: {
"Content-Type": "application/x-www-form-urlencoded",
Accept: "application/json",
},
body: new URLSearchParams({
grant_type: "refresh_token",
refresh_token: refreshToken,
client_id: clientId,
client_secret: clientSecret,
}),
});
if (!response.ok) {
const errorText = await response.text();
log?.error?.("TOKEN_REFRESH", "Failed to refresh Google token", { status: response.status, error: errorText });
return null;
}
const tokens = await response.json();
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Google token", { hasNewAccessToken: !!tokens.access_token, expiresIn: tokens.expires_in });
return { accessToken: tokens.access_token, refreshToken: tokens.refresh_token || refreshToken, expiresIn: tokens.expires_in };
} catch (error) {
log?.error?.("TOKEN_REFRESH", `Network error refreshing Google token: ${error.message}`);
return null;
}
const tokens = await response.json();
log?.info?.("TOKEN_REFRESH", "Successfully refreshed Google token", {
hasNewAccessToken: !!tokens.access_token,
hasNewRefreshToken: !!tokens.refresh_token,
expiresIn: tokens.expires_in,
});
return {
accessToken: tokens.access_token,
refreshToken: tokens.refresh_token || refreshToken,
expiresIn: tokens.expires_in,
};
}
/**
@@ -206,6 +190,7 @@ export async function refreshQwenToken(refreshToken, log) {
* Specialized refresh for Codex (OpenAI) OAuth tokens
*/
export async function refreshCodexToken(refreshToken, log) {
try {
const response = await fetch(OAUTH_ENDPOINTS.openai.token, {
method: "POST",
headers: {
@@ -242,6 +227,10 @@ export async function refreshCodexToken(refreshToken, log) {
refreshToken: tokens.refresh_token || refreshToken,
expiresIn: tokens.expires_in,
};
} catch (error) {
log?.error?.("TOKEN_REFRESH", `Network error refreshing Codex token: ${error.message}`);
return null;
}
}
/**

View File

@@ -213,30 +213,34 @@ async function getGeminiUsage(accessToken) {
*/
async function getAntigravityUsage(accessToken, providerSpecificData) {
try {
// First get project ID from subscription info
const projectId = await getAntigravityProjectId(accessToken);
// Fetch subscription info once — reuse for both projectId and plan
const subscriptionInfo = await getAntigravitySubscriptionInfo(accessToken);
const projectId = subscriptionInfo?.cloudaicompanionProject || null;
// Fetch quota data with timeout
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 10000); // 10s timeout
const response = await fetch(ANTIGRAVITY_CONFIG.quotaApiUrl, {
method: "POST",
headers: {
"Authorization": `Bearer ${accessToken}`,
"User-Agent": ANTIGRAVITY_CONFIG.userAgent,
"Content-Type": "application/json",
"X-Client-Name": "antigravity",
"X-Client-Version": "1.107.0",
"x-request-source": "local", // MITM bypass
},
body: JSON.stringify({
...(projectId ? { project: projectId } : {})
}),
signal: controller.signal,
});
clearTimeout(timeoutId);
let response;
try {
response = await fetch(ANTIGRAVITY_CONFIG.quotaApiUrl, {
method: "POST",
headers: {
"Authorization": `Bearer ${accessToken}`,
"User-Agent": ANTIGRAVITY_CONFIG.userAgent,
"Content-Type": "application/json",
"X-Client-Name": "antigravity",
"X-Client-Version": "1.107.0",
"x-request-source": "local", // MITM bypass
},
body: JSON.stringify({
...(projectId ? { project: projectId } : {})
}),
signal: controller.signal,
});
} finally {
clearTimeout(timeoutId);
}
if (response.status === 403) {
return {
@@ -302,9 +306,6 @@ async function getAntigravityUsage(accessToken, providerSpecificData) {
}
}
// Get subscription info for plan type
const subscriptionInfo = await getAntigravitySubscriptionInfo(accessToken);
return {
plan: subscriptionInfo?.currentTier?.name || "Unknown",
quotas,
@@ -332,10 +333,9 @@ async function getAntigravityProjectId(accessToken) {
* Get Antigravity subscription info
*/
async function getAntigravitySubscriptionInfo(accessToken) {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 10000); // 10s timeout
try {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 10000); // 10s timeout
const response = await fetch(ANTIGRAVITY_CONFIG.loadProjectApiUrl, {
method: "POST",
headers: {
@@ -347,15 +347,14 @@ async function getAntigravitySubscriptionInfo(accessToken) {
body: JSON.stringify({ metadata: CLIENT_METADATA, mode: 1 }),
signal: controller.signal,
});
clearTimeout(timeoutId);
if (!response.ok) return null;
return await response.json();
} catch (error) {
console.error("[Antigravity Subscription] Error:", error.message);
return null;
} finally {
clearTimeout(timeoutId);
}
}

View File

@@ -60,6 +60,7 @@ const FIELD = {
MSG_ID: 13,
MSG_TOOL_RESULTS: 18,
MSG_IS_AGENTIC: 29,
MSG_SERVER_BUBBLE_ID: 32,
MSG_UNIFIED_MODE: 47,
MSG_SUPPORTED_TOOLS: 51,
@@ -78,10 +79,19 @@ const FIELD = {
CLIENT_RESULT_TOOL_CALL_ID: 35,
CLIENT_RESULT_MODEL_CALL_ID: 48,
CLIENT_RESULT_TOOL_INDEX: 49,
// Aliases used by encodeClientSideToolV2Result
CV2R_TOOL: 1,
CV2R_MCP_RESULT: 28,
CV2R_CALL_ID: 35,
CV2R_MODEL_CALL_ID: 48,
CV2R_TOOL_INDEX: 49,
// MCPResult (nested inside ClientSideToolV2Result.mcp_result)
MCP_RESULT_SELECTED_TOOL: 1,
MCP_RESULT_RESULT: 2,
// Aliases used by encodeMcpResult
MCPR_SELECTED_TOOL: 1,
MCPR_RESULT: 2,
// ClientSideToolV2Call (nested inside ToolResult.tool_call)
CLIENT_CALL_TOOL: 1,
@@ -91,6 +101,14 @@ const FIELD = {
CLIENT_CALL_RAW_ARGS: 10,
CLIENT_CALL_TOOL_INDEX: 48,
CLIENT_CALL_MODEL_CALL_ID: 49,
// Aliases used by encodeClientSideToolV2Call
CV2C_TOOL: 1,
CV2C_MCP_PARAMS: 27,
CV2C_CALL_ID: 3,
CV2C_NAME: 9,
CV2C_RAW_ARGS: 10,
CV2C_TOOL_INDEX: 48,
CV2C_MODEL_CALL_ID: 49,
// Model
MODEL_NAME: 1,

View File

@@ -49,7 +49,7 @@ export function transformToOllama(response, model) {
const formattedCalls = toolCallsArr.map(tc => ({
function: {
name: tc.function.name,
arguments: JSON.parse(tc.function.arguments || "{}")
arguments: (() => { try { return JSON.parse(tc.function.arguments || "{}"); } catch { return {}; } })()
}
}));
const ollama = JSON.stringify({
@@ -75,6 +75,9 @@ export function transformToOllama(response, model) {
}
});
if (!response.body) {
return new Response("", { status: response.status, headers: { "Content-Type": "application/x-ndjson" } });
}
return new Response(response.body.pipeThrough(transform), {
headers: { "Content-Type": "application/x-ndjson", "Access-Control-Allow-Origin": "*" }
});

View File

@@ -5,8 +5,8 @@ const isCloud = typeof caches !== "undefined" && typeof caches === "object";
const originalFetch = globalThis.fetch;
const proxyDispatchers = new Map();
// DNS cache: { [hostname]: { ip, expiry } }
const DNS_CACHE = {};
// DNS cache — use Map to avoid prototype pollution via malformed hostnames
const DNS_CACHE = new Map();
const MITM_BYPASS_HOSTS = ["cloudcode-pa.googleapis.com", "daily-cloudcode-pa.googleapis.com", "googleapis.com"];
const MITM_BYPASS_HEADER = "x-request-source";
const MITM_BYPASS_VALUE = "local";
@@ -24,7 +24,7 @@ function normalizeString(value) {
* Resolve real IP using Google DNS (bypass system DNS)
*/
async function resolveRealIP(hostname) {
const cached = DNS_CACHE[hostname];
const cached = DNS_CACHE.get(hostname);
if (cached && Date.now() < cached.expiry) return cached.ip;
try {
@@ -34,7 +34,7 @@ async function resolveRealIP(hostname) {
resolver.setServers(GOOGLE_DNS_SERVERS);
const resolve4 = promisify(resolver.resolve4.bind(resolver));
const addresses = await resolve4(hostname);
DNS_CACHE[hostname] = { ip: addresses[0], expiry: Date.now() + MEMORY_CONFIG.dnsCacheTtlMs };
DNS_CACHE.set(hostname, { ip: addresses[0], expiry: Date.now() + MEMORY_CONFIG.dnsCacheTtlMs });
return addresses[0];
} catch (error) {
console.warn(`[ProxyFetch] DNS resolve failed for ${hostname}:`, error.message);
@@ -53,23 +53,27 @@ function shouldBypassMitmDns(url, options) {
headers[MITM_BYPASS_HEADER.charAt(0).toUpperCase() + MITM_BYPASS_HEADER.slice(1)] === MITM_BYPASS_VALUE;
if (!hasLocalMarker) {
// Debug: log when bypass is not triggered
const hostname = new URL(url).hostname;
if (MITM_BYPASS_HOSTS.some(host => hostname.includes(host))) {
console.warn(`[ProxyFetch] MITM bypass NOT triggered for ${hostname} - missing header`);
}
try {
const hostname = new URL(url).hostname;
if (MITM_BYPASS_HOSTS.some(host => hostname.includes(host))) {
console.warn(`[ProxyFetch] MITM bypass NOT triggered for ${hostname} - missing header`);
}
} catch { /* invalid URL — skip debug log */ }
return false;
}
const hostname = new URL(url).hostname;
return MITM_BYPASS_HOSTS.some(host => hostname.includes(host));
try {
const hostname = new URL(url).hostname;
return MITM_BYPASS_HOSTS.some(host => hostname.includes(host));
} catch { return false; }
}
function shouldBypassByNoProxy(targetUrl, noProxyValue) {
const noProxy = normalizeString(noProxyValue);
if (!noProxy) return false;
const hostname = new URL(targetUrl).hostname.toLowerCase();
let hostname;
try { hostname = new URL(targetUrl).hostname.toLowerCase(); } catch { return false; }
const patterns = noProxy.split(",").map((p) => p.trim().toLowerCase()).filter(Boolean);
return patterns.some((pattern) => {
@@ -86,7 +90,8 @@ function getEnvProxyUrl(targetUrl) {
const noProxy = process.env.NO_PROXY || process.env.no_proxy;
if (shouldBypassByNoProxy(targetUrl, noProxy)) return null;
const protocol = new URL(targetUrl).protocol;
let protocol;
try { protocol = new URL(targetUrl).protocol; } catch { return null; }
if (protocol === "https:") {
return process.env.HTTPS_PROXY || process.env.https_proxy ||
@@ -152,7 +157,7 @@ async function getDispatcher(proxyUrl) {
async function createBypassRequest(parsedUrl, realIP, options) {
const https = await import("https");
const net = await import("net");
const { Readable } = require("stream");
const { Readable } = await import("stream");
return new Promise((resolve, reject) => {
const socket = new net.Socket();

View File

@@ -44,7 +44,7 @@ async function createLogSession(sourceFormat, targetFormat, model) {
}
const timestamp = formatTimestamp();
const safeModel = model.replace(/[/:]/g, "-");
const safeModel = (model || "unknown").replace(/[/:]/g, "-");
const folderName = `${sourceFormat}_${targetFormat}_${safeModel}_${timestamp}`;
const sessionPath = path.join(LOGS_DIR, folderName);

View File

@@ -52,6 +52,13 @@ export function deriveSessionId(connectionId) {
return existing.sessionId;
}
// Evict oldest entry if store exceeds max size (safety cap between cleanup cycles)
const MAX_SESSIONS = 1000;
if (runtimeSessionStore.size >= MAX_SESSIONS) {
const oldest = runtimeSessionStore.keys().next().value;
runtimeSessionStore.delete(oldest);
}
const sessionId = generateBinaryStyleId();
runtimeSessionStore.set(connectionId, { sessionId, lastUsed: Date.now() });
return sessionId;

View File

@@ -6,7 +6,7 @@ import { parseSSELine, hasValuableContent, fixInvalidId, formatSSE } from "./str
export { COLORS, formatSSE };
const sharedDecoder = new TextDecoder();
// sharedEncoder is stateless — safe to share across streams
const sharedEncoder = new TextEncoder();
/**
@@ -49,6 +49,9 @@ export function createSSEStream(options = {}) {
let buffer = "";
let usage = null;
// Per-stream decoder with stream:true to correctly handle multi-byte chars split across chunks
const decoder = new TextDecoder("utf-8", { fatal: false });
const state = mode === STREAM_MODE.TRANSLATE ? { ...initState(sourceFormat), provider, toolNameMap, model } : null;
let totalContentLength = 0;
@@ -61,7 +64,7 @@ export function createSSEStream(options = {}) {
if (!ttftAt) {
ttftAt = Date.now();
}
const text = sharedDecoder.decode(chunk, { stream: true });
const text = decoder.decode(chunk, { stream: true });
buffer += text;
reqLogger?.appendProviderChunk?.(text);
@@ -253,7 +256,7 @@ export function createSSEStream(options = {}) {
flush(controller) {
trackPendingRequest(model, provider, connectionId, false);
try {
const remaining = sharedDecoder.decode();
const remaining = decoder.decode();
if (remaining) buffer += remaining;
if (mode === STREAM_MODE.PASSTHROUGH) {

View File

@@ -107,6 +107,9 @@ export function createDisconnectAwareStream(transformStream, streamController) {
controller.enqueue(value);
} catch (error) {
streamController.handleError(error);
// Cleanup reader/writer to avoid orphaned streams
reader.cancel().catch(() => {});
writer.abort().catch(() => {});
controller.error(error);
}
},