import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'node:http'; import { chatResponseToResponses, responsesRequestToChat, ResponsesSseTranslator, type BridgeSseEvent, type DiagnosticFn, type ResponsesCreateParams, } from './chatBridgeTranslate'; /** * 内置 Responses→Chat 桥接:给只会 /chat/completions 的端点(老 Ollama / vLLM / Xinference) * 套一层本地翻译代理。Codex 连这里的 /v1/responses,桥接翻译后转发上游。 * * 安全边界:只绑 127.0.0.1 随机端口;不持有任何密钥(incoming Authorization 原样透传, * key 由 Codex 侧 env_key 注入);只接受 POST /v1/responses,其余一律 404。 * 消息格式转换全部在 chatBridgeTranslate.ts,这里只管 HTTP 与生命周期。 */ export interface ChatBridgeOptions { /** 上游 chat 端点 API 根(形如 http://host:11434/v1) */ upstreamBaseUrl: string; /** 额外转发头(来自模型配置的 headersJson) */ headers?: Record | null; onDiagnostic?: DiagnosticFn; } export interface ChatBridgeInfo { url: string; port: number; upstreamBaseUrl: string; } function readBody(req: IncomingMessage): Promise { return new Promise((resolve, reject) => { const chunks: Buffer[] = []; req.on('data', (chunk: Buffer) => chunks.push(chunk)); req.on('end', () => resolve(Buffer.concat(chunks).toString('utf8'))); req.on('error', reject); }); } function writeSseEvents(res: ServerResponse, events: BridgeSseEvent[]): void { for (const item of events) { res.write(`event: ${item.event}\ndata: ${JSON.stringify(item.data)}\n\n`); } } export class ChatBridgeService { #server: Server | null = null; #info: ChatBridgeInfo | null = null; #options: ChatBridgeOptions | null = null; /** 诊断里区分同一轮里的多次上游请求 */ #seq = 0; /** 装配层(index.ts)注入的诊断出口;start 时的 onDiagnostic 优先 */ #diagnostic: DiagnosticFn | null = null; get info(): ChatBridgeInfo | null { return this.#info; } setDiagnostic(fn: DiagnosticFn | null): void { this.#diagnostic = fn; } #diag(line: string): void { (this.#options?.onDiagnostic ?? this.#diagnostic)?.(line); } async start(options: ChatBridgeOptions): Promise { await this.stop(); this.#options = options; const server = createServer((req, res) => { void this.#handle(req, res); }); this.#server = server; await new Promise((resolve, reject) => { server.once('error', reject); server.listen(0, '127.0.0.1', () => resolve()); }); const address = server.address(); const port = typeof address === 'object' && address ? address.port : 0; this.#info = { url: `http://127.0.0.1:${port}/v1`, port, upstreamBaseUrl: options.upstreamBaseUrl }; this.#diag(`chat-bridge: 已启动 ${this.#info.url} → ${options.upstreamBaseUrl}`); return this.#info; } async stop(): Promise { const server = this.#server; this.#server = null; this.#info = null; this.#options = null; if (!server) return; await new Promise((resolve) => server.close(() => resolve())); } #forwardHeaders(req: IncomingMessage): Record { const headers: Record = { 'content-type': 'application/json', // 上游 gzip 会把 SSE 行切碎,显式要求不压缩 'accept-encoding': 'identity', }; const authorization = req.headers.authorization; if (authorization) headers.authorization = authorization; for (const [key, value] of Object.entries(this.#options?.headers ?? {})) { if (value) headers[key.toLowerCase()] = value; } return headers; } #failJson(res: ServerResponse, status: number, message: string): void { if (res.headersSent) { res.end(); return; } res.writeHead(status, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: { message, type: 'bridge_error' } })); } async #handle(req: IncomingMessage, res: ServerResponse): Promise { if (req.method !== 'POST' || req.url !== '/v1/responses') { this.#failJson(res, 404, 'not found'); return; } let responsesReq: ResponsesCreateParams; try { responsesReq = JSON.parse(await readBody(req)) as ResponsesCreateParams; } catch { this.#failJson(res, 400, '请求体不是合法 JSON'); return; } const options = this.#options; if (!options) { this.#failJson(res, 503, '桥接未配置上游'); return; } const chatReq = responsesRequestToChat(responsesReq, (line) => this.#diag(line)); const upstream = `${options.upstreamBaseUrl.replace(/\/+$/u, '')}/chat/completions`; const id = ++this.#seq; const t0 = Date.now(); const secs = (at: number = Date.now()): string => `${((at - t0) / 1000).toFixed(1)}s`; const inputText = typeof responsesReq.input === 'string' ? responsesReq.input : JSON.stringify(responsesReq.input ?? ''); const messages = Array.isArray(chatReq.messages) ? chatReq.messages : []; this.#diag( `chat-bridge: #${id} 起 round model=${responsesReq.model ?? '?'} stream=${chatReq.stream ? 'Y' : 'N'} ` + `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; try { upstreamRes = await fetch(upstream, { method: 'POST', headers: this.#forwardHeaders(req), body: JSON.stringify(chatReq), signal: controller.signal, }); } catch (error) { if (clientGone) return; const message = error instanceof Error ? error.message : String(error); this.#diag(`chat-bridge: #${id} 上游不可达(${secs()}):${message}`); if (translator) { write(translator.fail(`桥接上游不可达:${message}`)); res.end(); return; } this.#failJson(res, 502, `桥接上游不可达:${message}`); return; } if (!upstreamRes.ok) { 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)}`); if (translator) { write(translator.fail(detail)); res.end(); return; } this.#failJson(res, upstreamRes.status, detail); return; } if (!translator) { try { const chatJson = (await upstreamRes.json()) as Parameters[0]; if (clientGone) return; res.writeHead(200, { 'content-type': 'application/json' }); res.end(JSON.stringify(chatResponseToResponses(chatJson, responsesReq))); this.#diag(`chat-bridge: #${id} 非流式完成(${secs()})`); } catch (error) { if (clientGone) return; this.#diag(`chat-bridge: #${id} 非流式解析失败(${secs()})`); this.#failJson(res, 502, `桥接解析上游响应失败:${error instanceof Error ? error.message : String(error)}`); } return; } // 流式:边收边翻边发(SSE 已在请求到达时就开好,translator 此时必定存在) let bytes = 0; try { if (!upstreamRes.body) throw new Error('上游响应没有 body'); const reader = upstreamRes.body.getReader(); const decoder = new TextDecoder(); for (;;) { const { done, value } = await reader.read(); if (done) break; bytes += value?.byteLength ?? 0; write(translator.push(decoder.decode(value, { stream: true }))); } write(translator.push(decoder.decode())); write(translator.finish()); this.#diag(`chat-bridge: #${id} 转发完成(${secs()},上游 ${bytes}B → ${events} 个事件)`); } catch (error) { // 客户端已走 = 我们自己主动 abort,不该再当成上游故障去报告警 if (clientGone) return; const message = error instanceof Error ? error.message : String(error); this.#diag(`chat-bridge: #${id} 流式转发中断(${secs()}):${message}`); write(translator.fail(message)); } res.end(); } } /** 模块级单例:providerService 应用/清除模型时启停,装配层退出时兜底 stop */ export const chatBridge = new ChatBridgeService();