/** * ChatRouterService:本地 HTTP 路由层外壳。 * * codex(wire_api="responses")把 POST {url}/v1/responses 发到本服务, * 这里翻译成 POST {upstream}/chat/completions 转发给真实上游,并把 * 上游的 Chat SSE 流式翻译回 Responses SSE。翻译本体在 translate.ts(纯函数)。 * * 生命周期:applyProvider 探测到 chat-only 端点时 start,换模型/清除/退出时 stop。 * 监听 127.0.0.1 随机端口,不持有任何密钥——Authorization 与模型记录配置的自定义头 * 由 codex 请求带入后按白名单透传。 */ import http from 'node:http'; import type { AddressInfo } from 'node:net'; import { ChatSseTranslator, chatResponseToResponses, responsesToChatRequest, type ResponsesCreateParams, } from './translate'; export interface ChatRouterStartOptions { /** 真实上游 API 根,形如 https://api.deepseek.com/v1(去尾斜杠后拼 /chat/completions) */ upstreamBaseUrl: string; /** 额外透传给上游的请求头名(来自模型记录 headersJson 的键;authorization 恒透传) */ forwardHeaderNames?: readonly string[]; onDiagnostic?: (message: string) => void; } export interface ChatRouterInfo { url: string; port: number; upstreamBaseUrl: string; } const MAX_ERROR_BODY_BYTES = 4_096; export class ChatRouterService { #server: http.Server | null = null; #upstreamBaseUrl: string | null = null; #forwardHeaderNames = new Set(['authorization']); #onDiagnostic: ((message: string) => void) | null = null; #diagnosticSink: ((message: string) => void) | null = null; #requestSeq = 0; /** 装配层注入的常驻诊断出口;start 传入的 onDiagnostic 优先 */ setDiagnosticSink(sink: ((message: string) => void) | null): void { this.#diagnosticSink = sink; } get running(): boolean { return this.#server !== null; } get url(): string | null { return this.#upstreamBaseUrl ? this.#routerUrl() : null; } get upstreamBaseUrl(): string | null { return this.#upstreamBaseUrl; } info(): { running: boolean; upstream: string } | null { if (!this.#server || !this.#upstreamBaseUrl) return null; return { running: true, upstream: this.#upstreamBaseUrl }; } #routerUrl(): string { const address = this.#server?.address() as AddressInfo | null; return `http://127.0.0.1:${address?.port ?? 0}/v1`; } #diagnose(message: string): void { this.#onDiagnostic?.(message); this.#diagnosticSink?.(message); } /** 启动(已在跑则先停旧的再起新的,端口随机)。失败时抛错且不留半开服务 */ async start(options: ChatRouterStartOptions): Promise { await this.stop(); const upstream = options.upstreamBaseUrl.replace(/\/+$/u, ''); if (!/^https?:\/\//u.test(upstream)) throw new Error(`路由层上游地址无效:${options.upstreamBaseUrl}`); this.#upstreamBaseUrl = upstream; this.#onDiagnostic = options.onDiagnostic ?? null; this.#forwardHeaderNames = new Set(['authorization']); for (const name of options.forwardHeaderNames ?? []) { const normalized = name.trim().toLowerCase(); if (normalized) this.#forwardHeaderNames.add(normalized); } const server = http.createServer((req, res) => { void this.#handle(req, res); }); this.#server = server; await new Promise((resolve, reject) => { const onListening = () => { server.off('error', onError); resolve(); }; const onError = (error: Error) => { server.off('listening', onListening); this.#server = null; reject(error); }; server.once('listening', onListening); server.once('error', onError); server.listen(0, '127.0.0.1'); }); const info: ChatRouterInfo = { url: this.#routerUrl(), port: (this.#server!.address() as AddressInfo).port, upstreamBaseUrl: upstream }; this.#diagnose(`路由层已启动:${info.url} → ${upstream}/chat/completions`); return info; } async stop(): Promise { const server = this.#server; this.#server = null; this.#upstreamBaseUrl = null; if (!server) return; await new Promise((resolve) => { server.close(() => resolve()); }); this.#diagnose('路由层已停止'); } async #handle(req: http.IncomingMessage, res: http.ServerResponse): Promise { const seq = ++this.#requestSeq; const startedAt = Date.now(); const path = (req.url ?? '').replace(/\?.*$/u, ''); if (req.method !== 'POST' || (path !== '/v1/responses' && path !== '/responses')) { res.writeHead(404, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: { message: `路由层只接受 POST /v1/responses,收到 ${req.method} ${path}`, type: 'router_error' } })); return; } const bodyChunks: Buffer[] = []; for await (const chunk of req) bodyChunks.push(chunk as Buffer); let responsesRequest: ResponsesCreateParams; try { responsesRequest = JSON.parse(Buffer.concat(bodyChunks).toString('utf8')) as ResponsesCreateParams; } catch { res.writeHead(400, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: { message: '请求体不是合法 JSON', type: 'router_error' } })); return; } if (!responsesRequest || typeof responsesRequest !== 'object') { res.writeHead(400, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: { message: '请求体必须为 JSON 对象', type: 'router_error' } })); return; } const chatRequest = responsesToChatRequest(responsesRequest, (message) => { this.#diagnose(`路由 round #${seq}: ${message}`); }); const upstreamUrl = `${this.#upstreamBaseUrl}/chat/completions`; // 防单槽占死:客户端提前挂断 → 立刻取消上游 const controller = new AbortController(); let responseFinished = false; let clientAborted = false; res.once('close', () => { if (!responseFinished) { clientAborted = true; controller.abort(); } }); const forwardHeaders: Record = { 'content-type': 'application/json', // 防上游 gzip 把 SSE 行切碎 'accept-encoding': 'identity', accept: chatRequest.stream ? 'text/event-stream' : 'application/json', }; const incoming = req.headers; for (const name of this.#forwardHeaderNames) { const value = incoming[name]; if (typeof value === 'string' && value) forwardHeaders[name] = value; } const started = Date.now(); let upstreamBytes = 0; let upstreamStatus: number | null = null; let outcome = 'ok'; // 流式开始后出错时复用同一 translator:fail 帧的 response.id 与已发事件一致且幂等 let translator: ChatSseTranslator | null = null; try { // 流式:先把 SSE 头与 created/in_progress 发出去,再等上游(可能慢 prefill 数分钟) if (chatRequest.stream) { res.writeHead(200, { 'content-type': 'text/event-stream', 'cache-control': 'no-cache', connection: 'keep-alive', 'x-accel-buffering': 'no', }); translator = new ChatSseTranslator((message) => { this.#diagnose(`路由 round #${seq}: ${message}`); }); for (const frame of translator.begin()) res.write(frame); } const upstreamRes = await fetch(upstreamUrl, { method: 'POST', headers: forwardHeaders, body: JSON.stringify(chatRequest), signal: controller.signal, }); upstreamStatus = upstreamRes.status; if (!upstreamRes.ok) { const errorBody = (await upstreamRes.text()).slice(0, MAX_ERROR_BODY_BYTES); outcome = `upstream_${upstreamRes.status}`; this.#diagnose(`路由 round #${seq}: 上游 ${upstreamRes.status} —— ${errorBody}`); if (chatRequest.stream && translator) { // SSE 已开:以 response.failed 收尾,codex 会走错误映射并按策略重试 for (const frame of translator.fail(`上游 HTTP ${upstreamRes.status}: ${errorBody || '无响应体'}`)) res.write(frame); responseFinished = true; res.end(); return; } res.writeHead(upstreamRes.status, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: { message: `上游 HTTP ${upstreamRes.status}: ${errorBody}`, type: 'upstream_error' } })); return; } if (!chatRequest.stream) { const payload = await upstreamRes.json() as Record; responseFinished = true; res.writeHead(200, { 'content-type': 'application/json' }); res.end(JSON.stringify(chatResponseToResponses(payload, `resp_router_${seq}`))); return; } const reader = upstreamRes.body?.getReader() ?? null; if (!reader) { for (const frame of translator!.fail('上游未返回可读响应体')) res.write(frame); responseFinished = true; res.end(); return; } const decoder = new TextDecoder(); for (;;) { const { done, value } = await reader.read(); if (done) break; upstreamBytes += value.byteLength; for (const frame of translator!.push(decoder.decode(value, { stream: true }))) res.write(frame); } for (const frame of translator!.finish()) res.write(frame); responseFinished = true; res.end(); } catch (error) { const message = error instanceof Error ? error.message : String(error); if (clientAborted) { outcome = 'client_abort'; this.#diagnose(`路由 round #${seq}: 客户端提前挂断,已取消上游请求`); return; } outcome = 'fetch_error'; this.#diagnose(`路由 round #${seq}: 上游请求失败 —— ${message}`); if (!res.headersSent) { res.writeHead(502, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: { message: `路由层访问上游失败:${message}`, type: 'router_error' } })); return; } if (translator && !responseFinished) { for (const frame of translator.fail(`上游请求失败:${message}`)) res.write(frame); responseFinished = true; res.end(); } } finally { this.#diagnose( `路由 round #${seq} 结束:状态=${outcome} 上游=${upstreamStatus ?? '-'} ` + `字节=${upstreamBytes} 耗时=${Date.now() - started}ms`, ); } } } /** 应用级单例:路由层随当前应用的 provider 启停 */ export const chatRouter = new ChatRouterService();