chatBridgeService.ts 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234
  1. import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'node:http';
  2. import {
  3. chatResponseToResponses,
  4. responsesRequestToChat,
  5. ResponsesSseTranslator,
  6. type BridgeSseEvent,
  7. type DiagnosticFn,
  8. type ResponsesCreateParams,
  9. } from './chatBridgeTranslate';
  10. /**
  11. * 内置 Responses→Chat 桥接:给只会 /chat/completions 的端点(老 Ollama / vLLM / Xinference)
  12. * 套一层本地翻译代理。Codex 连这里的 /v1/responses,桥接翻译后转发上游。
  13. *
  14. * 安全边界:只绑 127.0.0.1 随机端口;不持有任何密钥(incoming Authorization 原样透传,
  15. * key 由 Codex 侧 env_key 注入);只接受 POST /v1/responses,其余一律 404。
  16. * 消息格式转换全部在 chatBridgeTranslate.ts,这里只管 HTTP 与生命周期。
  17. */
  18. export interface ChatBridgeOptions {
  19. /** 上游 chat 端点 API 根(形如 http://host:11434/v1) */
  20. upstreamBaseUrl: string;
  21. /** 额外转发头(来自模型配置的 headersJson) */
  22. headers?: Record<string, string> | null;
  23. onDiagnostic?: DiagnosticFn;
  24. }
  25. export interface ChatBridgeInfo {
  26. url: string;
  27. port: number;
  28. upstreamBaseUrl: string;
  29. }
  30. function readBody(req: IncomingMessage): Promise<string> {
  31. return new Promise((resolve, reject) => {
  32. const chunks: Buffer[] = [];
  33. req.on('data', (chunk: Buffer) => chunks.push(chunk));
  34. req.on('end', () => resolve(Buffer.concat(chunks).toString('utf8')));
  35. req.on('error', reject);
  36. });
  37. }
  38. function writeSseEvents(res: ServerResponse, events: BridgeSseEvent[]): void {
  39. for (const item of events) {
  40. res.write(`event: ${item.event}\ndata: ${JSON.stringify(item.data)}\n\n`);
  41. }
  42. }
  43. export class ChatBridgeService {
  44. #server: Server | null = null;
  45. #info: ChatBridgeInfo | null = null;
  46. #options: ChatBridgeOptions | null = null;
  47. /** 诊断里区分同一轮里的多次上游请求 */
  48. #seq = 0;
  49. /** 装配层(index.ts)注入的诊断出口;start 时的 onDiagnostic 优先 */
  50. #diagnostic: DiagnosticFn | null = null;
  51. get info(): ChatBridgeInfo | null {
  52. return this.#info;
  53. }
  54. setDiagnostic(fn: DiagnosticFn | null): void {
  55. this.#diagnostic = fn;
  56. }
  57. #diag(line: string): void {
  58. (this.#options?.onDiagnostic ?? this.#diagnostic)?.(line);
  59. }
  60. async start(options: ChatBridgeOptions): Promise<ChatBridgeInfo> {
  61. await this.stop();
  62. this.#options = options;
  63. const server = createServer((req, res) => {
  64. void this.#handle(req, res);
  65. });
  66. this.#server = server;
  67. await new Promise<void>((resolve, reject) => {
  68. server.once('error', reject);
  69. server.listen(0, '127.0.0.1', () => resolve());
  70. });
  71. const address = server.address();
  72. const port = typeof address === 'object' && address ? address.port : 0;
  73. this.#info = { url: `http://127.0.0.1:${port}/v1`, port, upstreamBaseUrl: options.upstreamBaseUrl };
  74. this.#diag(`chat-bridge: 已启动 ${this.#info.url} → ${options.upstreamBaseUrl}`);
  75. return this.#info;
  76. }
  77. async stop(): Promise<void> {
  78. const server = this.#server;
  79. this.#server = null;
  80. this.#info = null;
  81. this.#options = null;
  82. if (!server) return;
  83. await new Promise<void>((resolve) => server.close(() => resolve()));
  84. }
  85. #forwardHeaders(req: IncomingMessage): Record<string, string> {
  86. const headers: Record<string, string> = {
  87. 'content-type': 'application/json',
  88. // 上游 gzip 会把 SSE 行切碎,显式要求不压缩
  89. 'accept-encoding': 'identity',
  90. };
  91. const authorization = req.headers.authorization;
  92. if (authorization) headers.authorization = authorization;
  93. for (const [key, value] of Object.entries(this.#options?.headers ?? {})) {
  94. if (value) headers[key.toLowerCase()] = value;
  95. }
  96. return headers;
  97. }
  98. #failJson(res: ServerResponse, status: number, message: string): void {
  99. if (res.headersSent) {
  100. res.end();
  101. return;
  102. }
  103. res.writeHead(status, { 'content-type': 'application/json' });
  104. res.end(JSON.stringify({ error: { message, type: 'bridge_error' } }));
  105. }
  106. async #handle(req: IncomingMessage, res: ServerResponse): Promise<void> {
  107. if (req.method !== 'POST' || req.url !== '/v1/responses') {
  108. this.#failJson(res, 404, 'not found');
  109. return;
  110. }
  111. let responsesReq: ResponsesCreateParams;
  112. try {
  113. responsesReq = JSON.parse(await readBody(req)) as ResponsesCreateParams;
  114. } catch {
  115. this.#failJson(res, 400, '请求体不是合法 JSON');
  116. return;
  117. }
  118. const options = this.#options;
  119. if (!options) {
  120. this.#failJson(res, 503, '桥接未配置上游');
  121. return;
  122. }
  123. const chatReq = responsesRequestToChat(responsesReq, (line) => this.#diag(line));
  124. const upstream = `${options.upstreamBaseUrl.replace(/\/+$/u, '')}/chat/completions`;
  125. const id = ++this.#seq;
  126. const t0 = Date.now();
  127. const secs = (at: number = Date.now()): string => `${((at - t0) / 1000).toFixed(1)}s`;
  128. const inputText = typeof responsesReq.input === 'string' ? responsesReq.input : JSON.stringify(responsesReq.input ?? '');
  129. const messages = Array.isArray(chatReq.messages) ? chatReq.messages : [];
  130. this.#diag(
  131. `chat-bridge: #${id} 起 round model=${responsesReq.model ?? '?'} stream=${chatReq.stream ? 'Y' : 'N'} ` +
  132. `tools=${Array.isArray(chatReq.tools) ? chatReq.tools.length : 0} messages=${messages.length} 输入≈${inputText.length}字 → ${upstream}`,
  133. );
  134. let upstreamRes: Response;
  135. try {
  136. upstreamRes = await fetch(upstream, {
  137. method: 'POST',
  138. headers: this.#forwardHeaders(req),
  139. body: JSON.stringify(chatReq),
  140. });
  141. } catch (error) {
  142. const message = error instanceof Error ? error.message : String(error);
  143. this.#diag(`chat-bridge: #${id} 上游不可达(${secs()}):${message}`);
  144. this.#failJson(res, 502, `桥接上游不可达:${message}`);
  145. return;
  146. }
  147. if (!upstreamRes.ok) {
  148. const text = await upstreamRes.text().catch(() => '');
  149. this.#diag(`chat-bridge: #${id} 上游 HTTP ${upstreamRes.status}(${secs()}):${text.slice(0, 300)}`);
  150. this.#failJson(res, upstreamRes.status, `桥接上游返回 HTTP ${upstreamRes.status}:${text.slice(0, 500)}`);
  151. return;
  152. }
  153. if (!chatReq.stream) {
  154. try {
  155. const chatJson = (await upstreamRes.json()) as Parameters<typeof chatResponseToResponses>[0];
  156. res.writeHead(200, { 'content-type': 'application/json' });
  157. res.end(JSON.stringify(chatResponseToResponses(chatJson, responsesReq)));
  158. this.#diag(`chat-bridge: #${id} 非流式完成(${secs()})`);
  159. } catch (error) {
  160. this.#diag(`chat-bridge: #${id} 非流式解析失败(${secs()})`);
  161. this.#failJson(res, 502, `桥接解析上游响应失败:${error instanceof Error ? error.message : String(error)}`);
  162. }
  163. return;
  164. }
  165. // 流式:边收边翻边发
  166. res.writeHead(200, {
  167. 'content-type': 'text/event-stream',
  168. 'cache-control': 'no-cache',
  169. connection: 'keep-alive',
  170. });
  171. const translator = new ResponsesSseTranslator(responsesReq, (line) => this.#diag(line));
  172. let bytes = 0;
  173. let events = 0;
  174. let aborted = false;
  175. const write = (items: BridgeSseEvent[]): void => {
  176. events += items.length;
  177. writeSseEvents(res, items);
  178. };
  179. // Codex 提前挂断是这次排查的关键信号:没有这行就是上游自己停了
  180. res.once('close', () => {
  181. if (!res.writableEnded) {
  182. aborted = true;
  183. this.#diag(`chat-bridge: #${id} Codex 提前断开(${secs()},已发 ${events} 个事件)`);
  184. }
  185. });
  186. try {
  187. if (!upstreamRes.body) throw new Error('上游响应没有 body');
  188. const reader = upstreamRes.body.getReader();
  189. const decoder = new TextDecoder();
  190. for (;;) {
  191. const { done, value } = await reader.read();
  192. if (done) break;
  193. bytes += value?.byteLength ?? 0;
  194. write(translator.push(decoder.decode(value, { stream: true })));
  195. }
  196. write(translator.push(decoder.decode()));
  197. write(translator.finish());
  198. this.#diag(`chat-bridge: #${id} 转发完成(${secs()},上游 ${bytes}B → ${events} 个事件)`);
  199. } catch (error) {
  200. const message = error instanceof Error ? error.message : String(error);
  201. this.#diag(`chat-bridge: #${id} 流式转发中断(${secs()}):${message}`);
  202. write(translator.fail(message));
  203. }
  204. res.end();
  205. }
  206. }
  207. /** 模块级单例:providerService 应用/清除模型时启停,装配层退出时兜底 stop */
  208. export const chatBridge = new ChatBridgeService();