Skip to content
Merged
Show file tree
Hide file tree
Changes from 5 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
52 changes: 44 additions & 8 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 @@ -1266,7 +1302,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 +1311,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 +1361,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 +1373,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 +1393,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 +1403,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 +1454,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 +1464,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
101 changes: 101 additions & 0 deletions tests/unit/builtin-proxy-responses-shim.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -554,3 +554,104 @@ test('builtin-proxy /v1/responses maps Responses tool items through chat fallbac
await closeServer(upstream);
}
});

test('builtin-proxy /v1/responses stream=true emits response.failed when upstream stream aborts mid-flight', async () => {
const sockets = new Set();
const upstream = http.createServer((req, res) => {
if (req.url === '/v1/responses' && req.method === 'POST') {
res.writeHead(404, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ error: 'responses endpoint unavailable' }));
return;
}
if (req.url === '/v1/chat/completions' && req.method === 'POST') {
const chunks = [];
req.on('data', (chunk) => chunks.push(chunk));
req.on('end', () => {
res.writeHead(200, { 'Content-Type': 'text/event-stream; charset=utf-8' });
res.write('data: {"id":"chatcmpl_partial","model":"gpt-test","choices":[{"delta":{"role":"assistant"}}]}\n\n');
res.write('data: {"id":"chatcmpl_partial","model":"gpt-test","choices":[{"delta":{"content":"partial"}}]}\n\n');
setTimeout(() => {
try { req.socket.destroy(); } catch (_) {}
}, 30);
});
return;
}
res.writeHead(404, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ error: 'not found' }));
});
upstream.on('connection', (socket) => {
sockets.add(socket);
socket.on('close', () => sockets.delete(socket));
});
const { port: upstreamPort } = await listen(upstream);
let proxyRuntime = null;

try {
proxyRuntime = await startTestProxy(upstreamPort);
const proxyPort = proxyRuntime.server.address().port;
const sse = await requestText(`http://127.0.0.1:${proxyPort}/v1/responses`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: { model: 'gpt-test', input: 'ping', stream: true }
});
assert.equal(sse.status, 200);
assert.match(sse.headers['content-type'], /text\/event-stream/i);
assert.match(sse.text, /event: response\.created/);
assert.match(sse.text, /"delta":"partial"/);
assert.match(sse.text, /event: response\.failed/);
assert.match(sse.text, /data: \[DONE\]/);
} finally {
if (proxyRuntime) {
await closeServer(proxyRuntime.server, proxyRuntime.connections);
}
await closeServer(upstream, sockets);
}
});

test('builtin-proxy /v1/responses retries upstream after a transient connection reset', async () => {
let connectionCount = 0;
const upstream = http.createServer((req, res) => {
if (req.url === '/v1/responses' && req.method === 'POST') {
res.writeHead(404, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ error: 'responses endpoint unavailable' }));
return;
}
if (req.url === '/v1/chat/completions' && req.method === 'POST') {
connectionCount += 1;
if (connectionCount === 1) {
try { req.socket.destroy(); } catch (_) {}
return;
}
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({
id: 'chatcmpl_after_reset',
model: 'gpt-test',
choices: [{ message: { role: 'assistant', content: 'recovered' } }]
}));
return;
}
res.writeHead(404, { 'Content-Type': 'application/json' });
res.end(JSON.stringify({ error: 'not found' }));
});
const { port: upstreamPort } = await listen(upstream);
let proxyRuntime = null;

try {
proxyRuntime = await startTestProxy(upstreamPort);
const proxyPort = proxyRuntime.server.address().port;
const resp = await requestText(`http://127.0.0.1:${proxyPort}/v1/responses`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: { model: 'gpt-test', input: 'ping', stream: false }
});
assert.equal(resp.status, 200);
assert.ok(connectionCount >= 2, 'transient reset should be retried');
const parsed = JSON.parse(resp.text);
assert.equal(parsed.output[0].content[0].text, 'recovered');
} finally {
if (proxyRuntime) {
await closeServer(proxyRuntime.server, proxyRuntime.connections);
}
await closeServer(upstream);
}
});
Loading
Loading