|
@@ -153,61 +153,90 @@ export class ChatBridgeService {
|
|
|
`tools=${Array.isArray(chatReq.tools) ? chatReq.tools.length : 0} messages=${messages.length} 输入≈${inputText.length}字 → ${upstream}`,
|
|
`tools=${Array.isArray(chatReq.tools) ? chatReq.tools.length : 0} messages=${messages.length} 输入≈${inputText.length}字 → ${upstream}`,
|
|
|
);
|
|
);
|
|
|
|
|
|
|
|
|
|
+ /**
|
|
|
|
|
+ * Codex 拿不到响应就会挂断重来,而自部署的多是单槽推理:
|
|
|
|
|
+ * 不取消上游的话,被放弃的那次生成会继续占着槽,下一次请求排在它后面 —— 越重试越慢。
|
|
|
|
|
+ */
|
|
|
|
|
+ const controller = new AbortController();
|
|
|
|
|
+ let clientGone = false;
|
|
|
|
|
+ res.once('close', () => {
|
|
|
|
|
+ if (res.writableEnded) return;
|
|
|
|
|
+ clientGone = true;
|
|
|
|
|
+ controller.abort();
|
|
|
|
|
+ this.#diag(`chat-bridge: #${id} Codex 提前断开(${secs()}),已取消上游请求`);
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 流式请求先把 SSE 开起来再等上游:单槽自部署服务光 prefill 就要几分钟,
|
|
|
|
|
+ * 那期间 Codex 一个字节都收不到就会掐断重发(重发又把同一个 prompt 重新排一遍队)。
|
|
|
|
|
+ */
|
|
|
|
|
+ const translator = chatReq.stream ? new ResponsesSseTranslator(responsesReq, (line) => this.#diag(line)) : null;
|
|
|
|
|
+ let events = 0;
|
|
|
|
|
+ const write = (items: BridgeSseEvent[]): void => {
|
|
|
|
|
+ events += items.length;
|
|
|
|
|
+ writeSseEvents(res, items);
|
|
|
|
|
+ };
|
|
|
|
|
+ if (translator) {
|
|
|
|
|
+ res.writeHead(200, {
|
|
|
|
|
+ 'content-type': 'text/event-stream',
|
|
|
|
|
+ 'cache-control': 'no-cache',
|
|
|
|
|
+ connection: 'keep-alive',
|
|
|
|
|
+ });
|
|
|
|
|
+ write(translator.begin());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
let upstreamRes: Response;
|
|
let upstreamRes: Response;
|
|
|
try {
|
|
try {
|
|
|
upstreamRes = await fetch(upstream, {
|
|
upstreamRes = await fetch(upstream, {
|
|
|
method: 'POST',
|
|
method: 'POST',
|
|
|
headers: this.#forwardHeaders(req),
|
|
headers: this.#forwardHeaders(req),
|
|
|
body: JSON.stringify(chatReq),
|
|
body: JSON.stringify(chatReq),
|
|
|
|
|
+ signal: controller.signal,
|
|
|
});
|
|
});
|
|
|
} catch (error) {
|
|
} catch (error) {
|
|
|
|
|
+ if (clientGone) return;
|
|
|
const message = error instanceof Error ? error.message : String(error);
|
|
const message = error instanceof Error ? error.message : String(error);
|
|
|
this.#diag(`chat-bridge: #${id} 上游不可达(${secs()}):${message}`);
|
|
this.#diag(`chat-bridge: #${id} 上游不可达(${secs()}):${message}`);
|
|
|
|
|
+ if (translator) {
|
|
|
|
|
+ write(translator.fail(`桥接上游不可达:${message}`));
|
|
|
|
|
+ res.end();
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
this.#failJson(res, 502, `桥接上游不可达:${message}`);
|
|
this.#failJson(res, 502, `桥接上游不可达:${message}`);
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
if (!upstreamRes.ok) {
|
|
if (!upstreamRes.ok) {
|
|
|
const text = await upstreamRes.text().catch(() => '');
|
|
const text = await upstreamRes.text().catch(() => '');
|
|
|
|
|
+ if (clientGone) return;
|
|
|
|
|
+ const detail = `桥接上游返回 HTTP ${upstreamRes.status}:${text.slice(0, 500)}`;
|
|
|
this.#diag(`chat-bridge: #${id} 上游 HTTP ${upstreamRes.status}(${secs()}):${text.slice(0, 300)}`);
|
|
this.#diag(`chat-bridge: #${id} 上游 HTTP ${upstreamRes.status}(${secs()}):${text.slice(0, 300)}`);
|
|
|
- this.#failJson(res, upstreamRes.status, `桥接上游返回 HTTP ${upstreamRes.status}:${text.slice(0, 500)}`);
|
|
|
|
|
|
|
+ if (translator) {
|
|
|
|
|
+ write(translator.fail(detail));
|
|
|
|
|
+ res.end();
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ this.#failJson(res, upstreamRes.status, detail);
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- if (!chatReq.stream) {
|
|
|
|
|
|
|
+ if (!translator) {
|
|
|
try {
|
|
try {
|
|
|
const chatJson = (await upstreamRes.json()) as Parameters<typeof chatResponseToResponses>[0];
|
|
const chatJson = (await upstreamRes.json()) as Parameters<typeof chatResponseToResponses>[0];
|
|
|
|
|
+ if (clientGone) return;
|
|
|
res.writeHead(200, { 'content-type': 'application/json' });
|
|
res.writeHead(200, { 'content-type': 'application/json' });
|
|
|
res.end(JSON.stringify(chatResponseToResponses(chatJson, responsesReq)));
|
|
res.end(JSON.stringify(chatResponseToResponses(chatJson, responsesReq)));
|
|
|
this.#diag(`chat-bridge: #${id} 非流式完成(${secs()})`);
|
|
this.#diag(`chat-bridge: #${id} 非流式完成(${secs()})`);
|
|
|
} catch (error) {
|
|
} catch (error) {
|
|
|
|
|
+ if (clientGone) return;
|
|
|
this.#diag(`chat-bridge: #${id} 非流式解析失败(${secs()})`);
|
|
this.#diag(`chat-bridge: #${id} 非流式解析失败(${secs()})`);
|
|
|
this.#failJson(res, 502, `桥接解析上游响应失败:${error instanceof Error ? error.message : String(error)}`);
|
|
this.#failJson(res, 502, `桥接解析上游响应失败:${error instanceof Error ? error.message : String(error)}`);
|
|
|
}
|
|
}
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- // 流式:边收边翻边发
|
|
|
|
|
- res.writeHead(200, {
|
|
|
|
|
- 'content-type': 'text/event-stream',
|
|
|
|
|
- 'cache-control': 'no-cache',
|
|
|
|
|
- connection: 'keep-alive',
|
|
|
|
|
- });
|
|
|
|
|
- const translator = new ResponsesSseTranslator(responsesReq, (line) => this.#diag(line));
|
|
|
|
|
|
|
+ // 流式:边收边翻边发(SSE 已在请求到达时就开好,translator 此时必定存在)
|
|
|
let bytes = 0;
|
|
let bytes = 0;
|
|
|
- let events = 0;
|
|
|
|
|
- let aborted = false;
|
|
|
|
|
- const write = (items: BridgeSseEvent[]): void => {
|
|
|
|
|
- events += items.length;
|
|
|
|
|
- writeSseEvents(res, items);
|
|
|
|
|
- };
|
|
|
|
|
- // Codex 提前挂断是这次排查的关键信号:没有这行就是上游自己停了
|
|
|
|
|
- res.once('close', () => {
|
|
|
|
|
- if (!res.writableEnded) {
|
|
|
|
|
- aborted = true;
|
|
|
|
|
- this.#diag(`chat-bridge: #${id} Codex 提前断开(${secs()},已发 ${events} 个事件)`);
|
|
|
|
|
- }
|
|
|
|
|
- });
|
|
|
|
|
try {
|
|
try {
|
|
|
if (!upstreamRes.body) throw new Error('上游响应没有 body');
|
|
if (!upstreamRes.body) throw new Error('上游响应没有 body');
|
|
|
const reader = upstreamRes.body.getReader();
|
|
const reader = upstreamRes.body.getReader();
|
|
@@ -222,6 +251,8 @@ export class ChatBridgeService {
|
|
|
write(translator.finish());
|
|
write(translator.finish());
|
|
|
this.#diag(`chat-bridge: #${id} 转发完成(${secs()},上游 ${bytes}B → ${events} 个事件)`);
|
|
this.#diag(`chat-bridge: #${id} 转发完成(${secs()},上游 ${bytes}B → ${events} 个事件)`);
|
|
|
} catch (error) {
|
|
} catch (error) {
|
|
|
|
|
+ // 客户端已走 = 我们自己主动 abort,不该再当成上游故障去报告警
|
|
|
|
|
+ if (clientGone) return;
|
|
|
const message = error instanceof Error ? error.message : String(error);
|
|
const message = error instanceof Error ? error.message : String(error);
|
|
|
this.#diag(`chat-bridge: #${id} 流式转发中断(${secs()}):${message}`);
|
|
this.#diag(`chat-bridge: #${id} 流式转发中断(${secs()}):${message}`);
|
|
|
write(translator.fail(message));
|
|
write(translator.fail(message));
|