Skip to content
Merged
Show file tree
Hide file tree
Changes from 9 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 10 additions & 2 deletions cli.js
Original file line number Diff line number Diff line change
Expand Up @@ -304,8 +304,16 @@ const CLI_INSTALL_TARGETS = Object.freeze([
}
]);

const HTTP_KEEP_ALIVE_AGENT = new http.Agent({ keepAlive: true });
const HTTPS_KEEP_ALIVE_AGENT = new https.Agent({ keepAlive: true });
const HTTP_KEEP_ALIVE_AGENT = new http.Agent({
keepAlive: true,
keepAliveMsecs: 1000,
maxFreeSockets: 4
});
const HTTPS_KEEP_ALIVE_AGENT = new https.Agent({
keepAlive: true,
keepAliveMsecs: 1000,
maxFreeSockets: 4
});

const openaiBridgeHandler = createOpenaiBridgeHttpHandler({
settingsFile: OPENAI_BRIDGE_SETTINGS_FILE,
Expand Down
109 changes: 107 additions & 2 deletions cli/builtin-proxy.js
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,42 @@ function createBuiltinProxyRuntimeController(deps = {}) {
return false;
}

function isTransientNetworkError(error) {
const text = String(error || '').trim();
if (!text) return false;
if (/socket hang up/i.test(text)) return true;
if (/ECONNRESET|ECONNREFUSED|EPIPE|EPROTO|ETIMEDOUT/i.test(text)) return true;
if (/EAI_AGAIN/i.test(text)) return true;
if (/UND_ERR_SOCKET/i.test(text)) return true;
if (/disconnected before|secure tls|tls handshake/i.test(text)) return true;
return false;
}

const TRANSIENT_RETRY_DELAYS_MS = [200, 600];

async function retryTransientRequest(executor) {
let lastResult = null;
for (let attempt = 0; attempt <= TRANSIENT_RETRY_DELAYS_MS.length; attempt += 1) {
if (attempt > 0) {
const delay = TRANSIENT_RETRY_DELAYS_MS[attempt - 1];
// eslint-disable-next-line no-await-in-loop
await new Promise((r) => {
const t = setTimeout(r, delay);
if (typeof t.unref === 'function') t.unref();
});
}
// eslint-disable-next-line no-await-in-loop
const result = await executor(attempt);
lastResult = result;
if (!result) return result;
if (result.ok) return result;
if (result.retry) return result;
if (result.status && result.status > 0) return result;
if (!isTransientNetworkError(result.error)) return result;
}
return lastResult;
}

function proxyRequestJson(targetUrl, options = {}) {
const parsed = new URL(targetUrl);
const transport = parsed.protocol === 'https:' ? https : http;
Expand Down Expand Up @@ -206,7 +242,7 @@ function createBuiltinProxyRuntimeController(deps = {}) {
}
let lastResult = null;
for (let index = 0; index < urls.length; index += 1) {
const result = await proxyRequestJson(urls[index], options);
const result = await retryTransientRequest(() => proxyRequestJson(urls[index], options));
lastResult = result;
if (!result.ok) {
return result;
Expand Down Expand Up @@ -702,9 +738,34 @@ function createBuiltinProxyRuntimeController(deps = {}) {
}
}

function stopChatStreamHeartbeat(state) {
if (!state || !state.heartbeatTimer) return;
clearInterval(state.heartbeatTimer);
state.heartbeatTimer = null;
}

function startChatStreamHeartbeat(state) {
if (!state || state.heartbeatTimer) return;
const timer = setInterval(() => {
if (state.finished) {
stopChatStreamHeartbeat(state);
return;
}
const target = state.res;
if (!target || target.writableEnded || target.destroyed) {
stopChatStreamHeartbeat(state);
return;
}
try { target.write(': keepalive\n\n'); } catch (_) {}
}, 15000);
if (typeof timer.unref === 'function') timer.unref();
state.heartbeatTimer = timer;
}

function finishChatStreamResponsesSse(state) {
if (state.finished) return;
state.finished = true;
stopChatStreamHeartbeat(state);

if (state.messageItem) {
const outputIndex = state.output.indexOf(state.messageItem);
Expand Down Expand Up @@ -759,6 +820,22 @@ function createBuiltinProxyRuntimeController(deps = {}) {
state.res.end();
}

function failResponsesSseRaw(res, message) {
if (!res || res.writableEnded || res.destroyed) return;
try {
writeSse(res, 'response.failed', { type: 'response.failed', error: message || 'upstream stream failed' });
writeSse(res, 'done', '[DONE]');
res.end();
} catch (_) {}
}

function failChatStreamResponsesSse(state, message) {
if (!state || state.finished) return;
state.finished = true;
stopChatStreamHeartbeat(state);
failResponsesSseRaw(state.res, message);
}

function streamChatCompletionsAsResponsesSse(targetUrl, options = {}) {
const parsed = new URL(targetUrl);
const transport = parsed.protocol === 'https:' ? https : http;
Expand Down Expand Up @@ -796,6 +873,29 @@ function createBuiltinProxyRuntimeController(deps = {}) {
const status = upstreamRes.statusCode || 0;
const chunks = [];
const contentType = String(upstreamRes.headers && upstreamRes.headers['content-type'] || '');
let streamState = null;

const handleAbort = (reason) => {
if (settled) return;
if (streamState) {
failChatStreamResponsesSse(streamState, reason);
finish({ ok: true });
return;
}
if (res.headersSent) {
failResponsesSseRaw(res, reason);
finish({ ok: true });
return;
}
finish({
ok: false,
status,
error: reason,
bodyText: chunks.length ? Buffer.concat(chunks).toString('utf-8') : ''
});
};
upstreamRes.on('error', (err) => handleAbort(err && err.message ? err.message : 'upstream stream failed'));
upstreamRes.on('aborted', () => handleAbort('upstream stream aborted'));

if (status === 404 || status === 405) {
upstreamRes.on('data', (chunk) => chunk && chunks.push(chunk));
Expand Down Expand Up @@ -851,6 +951,11 @@ function createBuiltinProxyRuntimeController(deps = {}) {
return sequence;
}
};
streamState = state;
startChatStreamHeartbeat(state);
if (typeof res.on === 'function') {
res.on('close', () => stopChatStreamHeartbeat(state));
}
writeSse(res, 'response.created', {
type: 'response.created',
response: {
Expand Down Expand Up @@ -914,7 +1019,7 @@ function createBuiltinProxyRuntimeController(deps = {}) {
}
let lastResult = null;
for (const url of urls) {
const result = await streamChatCompletionsAsResponsesSse(url, options);
const result = await retryTransientRequest(() => streamChatCompletionsAsResponsesSse(url, options));
lastResult = result;
if (result && result.retry) continue;
return result;
Expand Down
63 changes: 53 additions & 10 deletions cli/openai-bridge.js
Original file line number Diff line number Diff line change
Expand Up @@ -716,6 +716,42 @@ function isLoopbackAddress(address) {
return value === '127.0.0.1' || value === '::1' || value === '::ffff:127.0.0.1';
}

function isTransientNetworkError(error) {
const text = String(error || '').trim();
if (!text) return false;
if (/socket hang up/i.test(text)) return true;
if (/ECONNRESET|ECONNREFUSED|EPIPE|EPROTO|ETIMEDOUT/i.test(text)) return true;
if (/EAI_AGAIN/i.test(text)) return true;
if (/UND_ERR_SOCKET/i.test(text)) return true;
if (/disconnected before|secure tls|tls handshake/i.test(text)) return true;
return false;
}

const TRANSIENT_RETRY_DELAYS_MS = [200, 600];

async function retryTransientRequest(executor) {
let lastResult = null;
for (let attempt = 0; attempt <= TRANSIENT_RETRY_DELAYS_MS.length; attempt += 1) {
if (attempt > 0) {
const delay = TRANSIENT_RETRY_DELAYS_MS[attempt - 1];
// eslint-disable-next-line no-await-in-loop
await new Promise((r) => {
const t = setTimeout(r, delay);
if (typeof t.unref === 'function') t.unref();
});
}
// eslint-disable-next-line no-await-in-loop
const result = await executor(attempt);
lastResult = result;
if (!result) return result;
if (result.ok) return result;
if (result.retry) return result;
if (result.status && result.status > 0) return result;
if (!isTransientNetworkError(result.error)) return result;
}
return lastResult;
}

function writeSse(res, eventName, dataObj) {
if (!res || res.writableEnded || res.destroyed) return;
if (eventName) {
Expand Down Expand Up @@ -758,7 +794,14 @@ function writeChatCompletionChunkAsResponsesSse(state, chunk) {
const delta = choice && choice.delta && typeof choice.delta === 'object' ? choice.delta : null;
if (!delta) continue;

const segments = [];
if (typeof delta.reasoning_content === 'string' && delta.reasoning_content) {
segments.push(delta.reasoning_content);
}
if (typeof delta.content === 'string' && delta.content) {
segments.push(delta.content);
}
for (const seg of segments) {
if (!state.messageItem) {
state.messageItem = {
id: `msg_${crypto.randomBytes(8).toString('hex')}`,
Expand All @@ -773,14 +816,14 @@ function writeChatCompletionChunkAsResponsesSse(state, chunk) {
item: state.messageItem
});
}
state.messageText += delta.content;
state.messageText += seg;
state.messageItem.content[0].text = state.messageText;
writeSse(state.res, 'response.output_text.delta', {
type: 'response.output_text.delta',
item_id: state.messageItem.id,
output_index: state.output.length - 1,
content_index: 0,
delta: delta.content,
delta: seg,
sequence_number: state.nextSeq()
});
}
Expand Down Expand Up @@ -1266,7 +1309,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
}

const url = joinApiUrl(upstream.baseUrl, 'models');
const result = await proxyRequestJson(url, {
const result = await retryTransientRequest(() => proxyRequestJson(url, {
method: 'GET',
headers: {
...(authHeader ? { Authorization: authHeader } : {}),
Expand All @@ -1275,7 +1318,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
maxBytes: maxUpstreamBytes,
httpAgent,
httpsAgent
});
}));
if (!result.ok) {
res.writeHead(502, { 'Content-Type': 'application/json; charset=utf-8' });
res.end(JSON.stringify({ error: `Upstream request failed: ${result.error}` }));
Expand Down Expand Up @@ -1325,7 +1368,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
}
const upstreamUrl = joinApiUrl(upstream.baseUrl, 'chat/completions');
const chatBody = { ...converted.chat, stream: true };
const streamed = await streamChatCompletionsAsResponsesSse(upstreamUrl, {
const streamed = await retryTransientRequest(() => streamChatCompletionsAsResponsesSse(upstreamUrl, {
method: 'POST',
body: chatBody,
headers: {
Expand All @@ -1337,7 +1380,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
httpsAgent,
res,
model: typeof chatBody.model === 'string' ? chatBody.model : ''
});
}));
if (!streamed.ok) {
if (res.writableEnded || res.destroyed) {
return;
Expand All @@ -1357,7 +1400,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
// Maxx-style behavior: prefer upstream /responses if supported.
// Fallback to /chat/completions conversion when upstream does not implement /responses (404/405).
const upstreamResponsesUrl = joinApiUrl(upstream.baseUrl, 'responses');
const upstreamResponsesResult = await proxyRequestJson(upstreamResponsesUrl, {
const upstreamResponsesResult = await retryTransientRequest(() => proxyRequestJson(upstreamResponsesUrl, {
method: 'POST',
body: toUpstreamNonStreamingResponsesPayload(responsesRequest),
headers: {
Expand All @@ -1367,7 +1410,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
maxBytes: maxUpstreamBytes,
httpAgent,
httpsAgent
});
}));

if (upstreamResponsesResult.ok && upstreamResponsesResult.status >= 200 && upstreamResponsesResult.status < 300) {
const upstreamJson = parseJsonOrError(upstreamResponsesResult.bodyText);
Expand Down Expand Up @@ -1418,7 +1461,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
}

const upstreamUrl = joinApiUrl(upstream.baseUrl, 'chat/completions');
const upstreamResult = await proxyRequestJson(upstreamUrl, {
const upstreamResult = await retryTransientRequest(() => proxyRequestJson(upstreamUrl, {
method: 'POST',
body: converted.chat,
headers: {
Expand All @@ -1428,7 +1471,7 @@ function createOpenaiBridgeHttpHandler(options = {}) {
maxBytes: maxUpstreamBytes,
httpAgent,
httpsAgent
});
}));
if (!upstreamResult.ok) {
res.writeHead(502, { 'Content-Type': 'application/json; charset=utf-8' });
res.end(JSON.stringify({ error: `Upstream request failed: ${upstreamResult.error}` }));
Expand Down
Loading
Loading