| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280 |
- /**
- * 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<string>(['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<ChatRouterInfo> {
- 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<void>((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<void> {
- const server = this.#server;
- this.#server = null;
- this.#upstreamBaseUrl = null;
- if (!server) return;
- await new Promise<void>((resolve) => {
- server.close(() => resolve());
- });
- this.#diagnose('路由层已停止');
- }
- async #handle(req: http.IncomingMessage, res: http.ServerResponse): Promise<void> {
- 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<string, string> = {
- '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<string, unknown>;
- 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();
|