routerService.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280
  1. /**
  2. * ChatRouterService:本地 HTTP 路由层外壳。
  3. *
  4. * codex(wire_api="responses")把 POST {url}/v1/responses 发到本服务,
  5. * 这里翻译成 POST {upstream}/chat/completions 转发给真实上游,并把
  6. * 上游的 Chat SSE 流式翻译回 Responses SSE。翻译本体在 translate.ts(纯函数)。
  7. *
  8. * 生命周期:applyProvider 探测到 chat-only 端点时 start,换模型/清除/退出时 stop。
  9. * 监听 127.0.0.1 随机端口,不持有任何密钥——Authorization 与模型记录配置的自定义头
  10. * 由 codex 请求带入后按白名单透传。
  11. */
  12. import http from 'node:http';
  13. import type { AddressInfo } from 'node:net';
  14. import {
  15. ChatSseTranslator,
  16. chatResponseToResponses,
  17. responsesToChatRequest,
  18. type ResponsesCreateParams,
  19. } from './translate';
  20. export interface ChatRouterStartOptions {
  21. /** 真实上游 API 根,形如 https://api.deepseek.com/v1(去尾斜杠后拼 /chat/completions) */
  22. upstreamBaseUrl: string;
  23. /** 额外透传给上游的请求头名(来自模型记录 headersJson 的键;authorization 恒透传) */
  24. forwardHeaderNames?: readonly string[];
  25. onDiagnostic?: (message: string) => void;
  26. }
  27. export interface ChatRouterInfo {
  28. url: string;
  29. port: number;
  30. upstreamBaseUrl: string;
  31. }
  32. const MAX_ERROR_BODY_BYTES = 4_096;
  33. export class ChatRouterService {
  34. #server: http.Server | null = null;
  35. #upstreamBaseUrl: string | null = null;
  36. #forwardHeaderNames = new Set<string>(['authorization']);
  37. #onDiagnostic: ((message: string) => void) | null = null;
  38. #diagnosticSink: ((message: string) => void) | null = null;
  39. #requestSeq = 0;
  40. /** 装配层注入的常驻诊断出口;start 传入的 onDiagnostic 优先 */
  41. setDiagnosticSink(sink: ((message: string) => void) | null): void {
  42. this.#diagnosticSink = sink;
  43. }
  44. get running(): boolean {
  45. return this.#server !== null;
  46. }
  47. get url(): string | null {
  48. return this.#upstreamBaseUrl ? this.#routerUrl() : null;
  49. }
  50. get upstreamBaseUrl(): string | null {
  51. return this.#upstreamBaseUrl;
  52. }
  53. info(): { running: boolean; upstream: string } | null {
  54. if (!this.#server || !this.#upstreamBaseUrl) return null;
  55. return { running: true, upstream: this.#upstreamBaseUrl };
  56. }
  57. #routerUrl(): string {
  58. const address = this.#server?.address() as AddressInfo | null;
  59. return `http://127.0.0.1:${address?.port ?? 0}/v1`;
  60. }
  61. #diagnose(message: string): void {
  62. this.#onDiagnostic?.(message);
  63. this.#diagnosticSink?.(message);
  64. }
  65. /** 启动(已在跑则先停旧的再起新的,端口随机)。失败时抛错且不留半开服务 */
  66. async start(options: ChatRouterStartOptions): Promise<ChatRouterInfo> {
  67. await this.stop();
  68. const upstream = options.upstreamBaseUrl.replace(/\/+$/u, '');
  69. if (!/^https?:\/\//u.test(upstream)) throw new Error(`路由层上游地址无效:${options.upstreamBaseUrl}`);
  70. this.#upstreamBaseUrl = upstream;
  71. this.#onDiagnostic = options.onDiagnostic ?? null;
  72. this.#forwardHeaderNames = new Set(['authorization']);
  73. for (const name of options.forwardHeaderNames ?? []) {
  74. const normalized = name.trim().toLowerCase();
  75. if (normalized) this.#forwardHeaderNames.add(normalized);
  76. }
  77. const server = http.createServer((req, res) => {
  78. void this.#handle(req, res);
  79. });
  80. this.#server = server;
  81. await new Promise<void>((resolve, reject) => {
  82. const onListening = () => {
  83. server.off('error', onError);
  84. resolve();
  85. };
  86. const onError = (error: Error) => {
  87. server.off('listening', onListening);
  88. this.#server = null;
  89. reject(error);
  90. };
  91. server.once('listening', onListening);
  92. server.once('error', onError);
  93. server.listen(0, '127.0.0.1');
  94. });
  95. const info: ChatRouterInfo = { url: this.#routerUrl(), port: (this.#server!.address() as AddressInfo).port, upstreamBaseUrl: upstream };
  96. this.#diagnose(`路由层已启动:${info.url} → ${upstream}/chat/completions`);
  97. return info;
  98. }
  99. async stop(): Promise<void> {
  100. const server = this.#server;
  101. this.#server = null;
  102. this.#upstreamBaseUrl = null;
  103. if (!server) return;
  104. await new Promise<void>((resolve) => {
  105. server.close(() => resolve());
  106. });
  107. this.#diagnose('路由层已停止');
  108. }
  109. async #handle(req: http.IncomingMessage, res: http.ServerResponse): Promise<void> {
  110. const seq = ++this.#requestSeq;
  111. const startedAt = Date.now();
  112. const path = (req.url ?? '').replace(/\?.*$/u, '');
  113. if (req.method !== 'POST' || (path !== '/v1/responses' && path !== '/responses')) {
  114. res.writeHead(404, { 'content-type': 'application/json' });
  115. res.end(JSON.stringify({ error: { message: `路由层只接受 POST /v1/responses,收到 ${req.method} ${path}`, type: 'router_error' } }));
  116. return;
  117. }
  118. const bodyChunks: Buffer[] = [];
  119. for await (const chunk of req) bodyChunks.push(chunk as Buffer);
  120. let responsesRequest: ResponsesCreateParams;
  121. try {
  122. responsesRequest = JSON.parse(Buffer.concat(bodyChunks).toString('utf8')) as ResponsesCreateParams;
  123. } catch {
  124. res.writeHead(400, { 'content-type': 'application/json' });
  125. res.end(JSON.stringify({ error: { message: '请求体不是合法 JSON', type: 'router_error' } }));
  126. return;
  127. }
  128. if (!responsesRequest || typeof responsesRequest !== 'object') {
  129. res.writeHead(400, { 'content-type': 'application/json' });
  130. res.end(JSON.stringify({ error: { message: '请求体必须为 JSON 对象', type: 'router_error' } }));
  131. return;
  132. }
  133. const chatRequest = responsesToChatRequest(responsesRequest, (message) => {
  134. this.#diagnose(`路由 round #${seq}: ${message}`);
  135. });
  136. const upstreamUrl = `${this.#upstreamBaseUrl}/chat/completions`;
  137. // 防单槽占死:客户端提前挂断 → 立刻取消上游
  138. const controller = new AbortController();
  139. let responseFinished = false;
  140. let clientAborted = false;
  141. res.once('close', () => {
  142. if (!responseFinished) {
  143. clientAborted = true;
  144. controller.abort();
  145. }
  146. });
  147. const forwardHeaders: Record<string, string> = {
  148. 'content-type': 'application/json',
  149. // 防上游 gzip 把 SSE 行切碎
  150. 'accept-encoding': 'identity',
  151. accept: chatRequest.stream ? 'text/event-stream' : 'application/json',
  152. };
  153. const incoming = req.headers;
  154. for (const name of this.#forwardHeaderNames) {
  155. const value = incoming[name];
  156. if (typeof value === 'string' && value) forwardHeaders[name] = value;
  157. }
  158. const started = Date.now();
  159. let upstreamBytes = 0;
  160. let upstreamStatus: number | null = null;
  161. let outcome = 'ok';
  162. // 流式开始后出错时复用同一 translator:fail 帧的 response.id 与已发事件一致且幂等
  163. let translator: ChatSseTranslator | null = null;
  164. try {
  165. // 流式:先把 SSE 头与 created/in_progress 发出去,再等上游(可能慢 prefill 数分钟)
  166. if (chatRequest.stream) {
  167. res.writeHead(200, {
  168. 'content-type': 'text/event-stream',
  169. 'cache-control': 'no-cache',
  170. connection: 'keep-alive',
  171. 'x-accel-buffering': 'no',
  172. });
  173. translator = new ChatSseTranslator((message) => {
  174. this.#diagnose(`路由 round #${seq}: ${message}`);
  175. });
  176. for (const frame of translator.begin()) res.write(frame);
  177. }
  178. const upstreamRes = await fetch(upstreamUrl, {
  179. method: 'POST',
  180. headers: forwardHeaders,
  181. body: JSON.stringify(chatRequest),
  182. signal: controller.signal,
  183. });
  184. upstreamStatus = upstreamRes.status;
  185. if (!upstreamRes.ok) {
  186. const errorBody = (await upstreamRes.text()).slice(0, MAX_ERROR_BODY_BYTES);
  187. outcome = `upstream_${upstreamRes.status}`;
  188. this.#diagnose(`路由 round #${seq}: 上游 ${upstreamRes.status} —— ${errorBody}`);
  189. if (chatRequest.stream && translator) {
  190. // SSE 已开:以 response.failed 收尾,codex 会走错误映射并按策略重试
  191. for (const frame of translator.fail(`上游 HTTP ${upstreamRes.status}: ${errorBody || '无响应体'}`)) res.write(frame);
  192. responseFinished = true;
  193. res.end();
  194. return;
  195. }
  196. res.writeHead(upstreamRes.status, { 'content-type': 'application/json' });
  197. res.end(JSON.stringify({ error: { message: `上游 HTTP ${upstreamRes.status}: ${errorBody}`, type: 'upstream_error' } }));
  198. return;
  199. }
  200. if (!chatRequest.stream) {
  201. const payload = await upstreamRes.json() as Record<string, unknown>;
  202. responseFinished = true;
  203. res.writeHead(200, { 'content-type': 'application/json' });
  204. res.end(JSON.stringify(chatResponseToResponses(payload, `resp_router_${seq}`)));
  205. return;
  206. }
  207. const reader = upstreamRes.body?.getReader() ?? null;
  208. if (!reader) {
  209. for (const frame of translator!.fail('上游未返回可读响应体')) res.write(frame);
  210. responseFinished = true;
  211. res.end();
  212. return;
  213. }
  214. const decoder = new TextDecoder();
  215. for (;;) {
  216. const { done, value } = await reader.read();
  217. if (done) break;
  218. upstreamBytes += value.byteLength;
  219. for (const frame of translator!.push(decoder.decode(value, { stream: true }))) res.write(frame);
  220. }
  221. for (const frame of translator!.finish()) res.write(frame);
  222. responseFinished = true;
  223. res.end();
  224. } catch (error) {
  225. const message = error instanceof Error ? error.message : String(error);
  226. if (clientAborted) {
  227. outcome = 'client_abort';
  228. this.#diagnose(`路由 round #${seq}: 客户端提前挂断,已取消上游请求`);
  229. return;
  230. }
  231. outcome = 'fetch_error';
  232. this.#diagnose(`路由 round #${seq}: 上游请求失败 —— ${message}`);
  233. if (!res.headersSent) {
  234. res.writeHead(502, { 'content-type': 'application/json' });
  235. res.end(JSON.stringify({ error: { message: `路由层访问上游失败:${message}`, type: 'router_error' } }));
  236. return;
  237. }
  238. if (translator && !responseFinished) {
  239. for (const frame of translator.fail(`上游请求失败:${message}`)) res.write(frame);
  240. responseFinished = true;
  241. res.end();
  242. }
  243. } finally {
  244. this.#diagnose(
  245. `路由 round #${seq} 结束:状态=${outcome} 上游=${upstreamStatus ?? '-'} ` +
  246. `字节=${upstreamBytes} 耗时=${Date.now() - started}ms`,
  247. );
  248. }
  249. }
  250. }
  251. /** 应用级单例:路由层随当前应用的 provider 启停 */
  252. export const chatRouter = new ChatRouterService();