Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
327 changes: 324 additions & 3 deletions cli/builtin-proxy.js
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,15 @@ function createBuiltinProxyRuntimeController(deps = {}) {
return false;
}

function shouldFallbackFromUpstreamResponsesFailure(error) {
const text = String(error || '').trim();
if (!text) return false;
if (/timeout/i.test(text)) return true;
if (/socket hang up/i.test(text)) return true;
if (/ECONNRESET/i.test(text)) return true;
return false;
}

function proxyRequestJson(targetUrl, options = {}) {
const parsed = new URL(targetUrl);
const transport = parsed.protocol === 'https:' ? https : http;
Expand Down Expand Up @@ -628,6 +637,291 @@ function createBuiltinProxyRuntimeController(deps = {}) {
writeSse(res, 'done', '[DONE]');
}

function appendChatStreamToolCall(target, toolCall) {
if (!toolCall || typeof toolCall !== 'object') return;
const index = Number.isFinite(toolCall.index) ? toolCall.index : target.length;
if (!target[index]) {
target[index] = {
id: '',
type: 'function',
function: { name: '', arguments: '' }
};
}
const current = target[index];
if (typeof toolCall.id === 'string' && toolCall.id) current.id = toolCall.id;
if (typeof toolCall.type === 'string' && toolCall.type) current.type = toolCall.type;
const fn = toolCall.function && typeof toolCall.function === 'object' ? toolCall.function : null;
if (fn) {
if (typeof fn.name === 'string' && fn.name) current.function.name = fn.name;
if (typeof fn.arguments === 'string') current.function.arguments += fn.arguments;
}
}

function writeChatCompletionChunkAsResponsesSse(state, chunk) {
if (!chunk || typeof chunk !== 'object') return;
if (typeof chunk.model === 'string' && chunk.model) {
state.model = chunk.model;
}
const choices = Array.isArray(chunk.choices) ? chunk.choices : [];
for (const choice of choices) {
const delta = choice && choice.delta && typeof choice.delta === 'object' ? choice.delta : null;
if (!delta) continue;

if (typeof delta.content === 'string' && delta.content) {
if (!state.messageItem) {
state.messageItem = {
id: `msg_${crypto.randomBytes(8).toString('hex')}`,
type: 'message',
role: 'assistant',
content: [{ type: 'output_text', text: '' }]
};
state.output.push(state.messageItem);
writeSse(state.res, 'response.output_item.added', {
type: 'response.output_item.added',
output_index: state.output.length - 1,
item: state.messageItem
});
}
state.messageText += delta.content;
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,
sequence_number: state.nextSeq()
});
}

if (Array.isArray(delta.tool_calls)) {
for (const toolCall of delta.tool_calls) {
appendChatStreamToolCall(state.toolCalls, toolCall);
}
}
}
}

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

if (state.messageItem) {
const outputIndex = state.output.indexOf(state.messageItem);
writeSse(state.res, 'response.output_text.done', {
type: 'response.output_text.done',
item_id: state.messageItem.id,
output_index: outputIndex,
content_index: 0,
text: state.messageText,
sequence_number: state.nextSeq()
});
writeSse(state.res, 'response.output_item.done', {
type: 'response.output_item.done',
output_index: outputIndex,
item: state.messageItem,
sequence_number: state.nextSeq()
});
}

for (const toolCall of state.toolCalls) {
if (!toolCall) continue;
const item = {
type: 'function_call',
call_id: toolCall.id || `call_${crypto.randomBytes(8).toString('hex')}`,
name: toolCall.function && typeof toolCall.function.name === 'string' ? toolCall.function.name : '',
arguments: toolCall.function && typeof toolCall.function.arguments === 'string' ? toolCall.function.arguments : ''
};
const outputIndex = state.output.length;
state.output.push(item);
writeSse(state.res, 'response.output_item.added', {
type: 'response.output_item.added',
output_index: outputIndex,
item
});
writeSse(state.res, 'response.output_item.done', {
type: 'response.output_item.done',
output_index: outputIndex,
item,
sequence_number: state.nextSeq()
});
}

const response = ensureResponseMetadata({
id: state.responseId,
model: state.model,
created_at: state.createdAt,
status: 'completed',
output: state.output
});
writeSse(state.res, 'response.completed', { type: 'response.completed', response });
writeSse(state.res, 'done', '[DONE]');
state.res.end();
}

function streamChatCompletionsAsResponsesSse(targetUrl, options = {}) {
const parsed = new URL(targetUrl);
const transport = parsed.protocol === 'https:' ? https : http;
const bodyText = options.body ? JSON.stringify(options.body) : '';
const headers = {
'Accept': 'text/event-stream',
...(options.body ? { 'Content-Type': 'application/json' } : {}),
...(options.headers || {})
};
if (options.body) {
headers['Content-Length'] = Buffer.byteLength(bodyText, 'utf-8');
}
const timeoutMs = Number.isFinite(options.timeoutMs)
? Math.max(1000, Number(options.timeoutMs))
: 30000;
const res = options.res;
const model = typeof options.model === 'string' ? options.model : '';

return new Promise((resolve) => {
let settled = false;
const finish = (value) => {
if (settled) return;
settled = true;
resolve(value);
};
const req = transport.request({
protocol: parsed.protocol,
hostname: parsed.hostname,
port: parsed.port || (parsed.protocol === 'https:' ? 443 : 80),
method: options.method || 'POST',
path: `${parsed.pathname}${parsed.search}`,
headers,
agent: parsed.protocol === 'https:' ? HTTPS_KEEP_ALIVE_AGENT : HTTP_KEEP_ALIVE_AGENT
}, (upstreamRes) => {
const status = upstreamRes.statusCode || 0;
const chunks = [];
const contentType = String(upstreamRes.headers && upstreamRes.headers['content-type'] || '');

if (status === 404 || status === 405) {
upstreamRes.on('data', (chunk) => chunk && chunks.push(chunk));
upstreamRes.on('end', () => finish({ retry: true, status, bodyText: chunks.length ? Buffer.concat(chunks).toString('utf-8') : '' }));
return;
}

if (status >= 400) {
upstreamRes.on('data', (chunk) => chunk && chunks.push(chunk));
upstreamRes.on('end', () => finish({ ok: false, status, bodyText: chunks.length ? Buffer.concat(chunks).toString('utf-8') : '' }));
return;
}

res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no'
});

if (!/text\/event-stream/i.test(contentType)) {
upstreamRes.on('data', (chunk) => chunk && chunks.push(chunk));
upstreamRes.on('end', () => {
const text = chunks.length ? Buffer.concat(chunks).toString('utf-8') : '';
const parsedJson = parseJsonOrError(text);
if (parsedJson.error) {
writeSse(res, 'response.failed', { type: 'response.failed', error: `invalid upstream response: ${parsedJson.error}` });
writeSse(res, 'done', '[DONE]');
res.end();
finish({ ok: true });
return;
}
sendResponsesSse(res, buildResponsesPayloadFromChatCompletion(parsedJson.value, model));
res.end();
finish({ ok: true });
});
return;
}

let sequence = 0;
const state = {
res,
responseId: `resp_${crypto.randomBytes(10).toString('hex')}`,
model,
createdAt: Math.floor(Date.now() / 1000),
output: [],
messageItem: null,
messageText: '',
toolCalls: [],
finished: false,
nextSeq: () => {
sequence += 1;
return sequence;
}
};
writeSse(res, 'response.created', {
type: 'response.created',
response: {
id: state.responseId,
model: state.model,
created_at: state.createdAt
}
});

let buffer = '';
const handleEventBlock = (block) => {
const dataLines = String(block || '')
.split(/\r?\n/)
.filter((line) => line.startsWith('data:'))
.map((line) => line.slice(5).trimStart());
if (dataLines.length === 0) return;
const data = dataLines.join('\n').trim();
if (!data) return;
if (data === '[DONE]') {
finishChatStreamResponsesSse(state);
finish({ ok: true });
return;
}
const parsedChunk = parseJsonOrError(data);
if (!parsedChunk.error) {
writeChatCompletionChunkAsResponsesSse(state, parsedChunk.value);
}
};

upstreamRes.on('data', (chunk) => {
buffer += chunk.toString('utf-8');
let boundary = buffer.search(/\r?\n\r?\n/);
while (boundary >= 0) {
const block = buffer.slice(0, boundary);
const match = buffer.slice(boundary).match(/^\r?\n\r?\n/);
buffer = buffer.slice(boundary + (match ? match[0].length : 2));
handleEventBlock(block);
boundary = buffer.search(/\r?\n\r?\n/);
}
});
upstreamRes.on('end', () => {
if (buffer.trim()) handleEventBlock(buffer);
finishChatStreamResponsesSse(state);
finish({ ok: true });
});
});
req.setTimeout(timeoutMs, () => {
try { req.destroy(new Error('timeout')); } catch (_) {}
finish({ ok: false, error: 'timeout' });
});
req.on('error', (err) => finish({ ok: false, error: err && err.message ? err.message : 'request failed' }));
if (bodyText) req.write(bodyText);
req.end();
});
}

async function streamChatCompletionsAsResponsesSseWithFallbackUrls(baseUrl, pathSuffix, options = {}) {
const urls = buildUpstreamUrlCandidates(baseUrl, pathSuffix);
if (urls.length === 0) {
return { ok: false, error: 'failed to build upstream URL' };
}
let lastResult = null;
for (const url of urls) {
const result = await streamChatCompletionsAsResponsesSse(url, options);
lastResult = result;
if (result && result.retry) continue;
return result;
}
return lastResult || { ok: false, error: 'failed to build upstream URL' };
}

function canListenPort(host, port) {
return new Promise((resolve) => {
const tester = net.createServer();
Expand Down Expand Up @@ -1044,14 +1338,41 @@ function createBuiltinProxyRuntimeController(deps = {}) {
}

if (!upstreamResponses.ok) {
res.writeHead(502, { 'Content-Type': 'application/json; charset=utf-8' });
res.end(JSON.stringify({ error: upstreamResponses.error || 'Upstream request failed' }));
return;
if (!shouldFallbackFromUpstreamResponsesFailure(upstreamResponses.error)) {
res.writeHead(502, { 'Content-Type': 'application/json; charset=utf-8' });
res.end(JSON.stringify({ error: upstreamResponses.error || 'Upstream request failed' }));
return;
}
// Some OpenAI-compatible gateways accept /responses but never complete it.
// Treat that as an unsupported Responses endpoint and try the chat fallback.
}

const model = typeof payload.model === 'string' ? payload.model : '';
const chatBody = buildChatCompletionsBodyFromResponsesPayload(payload);

if (wantsStream) {
const streamingChatBody = { ...chatBody, stream: true };
const streamed = await streamChatCompletionsAsResponsesSseWithFallbackUrls(upstream.baseUrl, 'chat/completions', {
method: 'POST',
headers: commonHeaders,
timeoutMs,
body: streamingChatBody,
res,
model
});
if (!streamed.ok) {
if (!res.headersSent) {
res.writeHead(streamed.status && streamed.status >= 400 ? streamed.status : 502, { 'Content-Type': 'application/json; charset=utf-8' });
res.end(streamed.bodyText || JSON.stringify({ error: streamed.error || 'proxy request failed' }));
} else if (!res.writableEnded) {
writeSse(res, 'response.failed', { type: 'response.failed', error: streamed.error || streamed.bodyText || 'proxy request failed' });
writeSse(res, 'done', '[DONE]');
res.end();
}
}
return;
}

const upstreamChat = await proxyRequestJsonWithFallbackUrls(upstream.baseUrl, 'chat/completions', {
method: 'POST',
headers: commonHeaders,
Expand Down
Loading
Loading