chatBridgeService.ts 7.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204
  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. /** 装配层(index.ts)注入的诊断出口;start 时的 onDiagnostic 优先 */
  48. #diagnostic: DiagnosticFn | null = null;
  49. get info(): ChatBridgeInfo | null {
  50. return this.#info;
  51. }
  52. setDiagnostic(fn: DiagnosticFn | null): void {
  53. this.#diagnostic = fn;
  54. }
  55. #diag(line: string): void {
  56. (this.#options?.onDiagnostic ?? this.#diagnostic)?.(line);
  57. }
  58. async start(options: ChatBridgeOptions): Promise<ChatBridgeInfo> {
  59. await this.stop();
  60. this.#options = options;
  61. const server = createServer((req, res) => {
  62. void this.#handle(req, res);
  63. });
  64. this.#server = server;
  65. await new Promise<void>((resolve, reject) => {
  66. server.once('error', reject);
  67. server.listen(0, '127.0.0.1', () => resolve());
  68. });
  69. const address = server.address();
  70. const port = typeof address === 'object' && address ? address.port : 0;
  71. this.#info = { url: `http://127.0.0.1:${port}/v1`, port, upstreamBaseUrl: options.upstreamBaseUrl };
  72. this.#diag(`chat-bridge: 已启动 ${this.#info.url} → ${options.upstreamBaseUrl}`);
  73. return this.#info;
  74. }
  75. async stop(): Promise<void> {
  76. const server = this.#server;
  77. this.#server = null;
  78. this.#info = null;
  79. this.#options = null;
  80. if (!server) return;
  81. await new Promise<void>((resolve) => server.close(() => resolve()));
  82. }
  83. #forwardHeaders(req: IncomingMessage): Record<string, string> {
  84. const headers: Record<string, string> = {
  85. 'content-type': 'application/json',
  86. // 上游 gzip 会把 SSE 行切碎,显式要求不压缩
  87. 'accept-encoding': 'identity',
  88. };
  89. const authorization = req.headers.authorization;
  90. if (authorization) headers.authorization = authorization;
  91. for (const [key, value] of Object.entries(this.#options?.headers ?? {})) {
  92. if (value) headers[key.toLowerCase()] = value;
  93. }
  94. return headers;
  95. }
  96. #failJson(res: ServerResponse, status: number, message: string): void {
  97. if (res.headersSent) {
  98. res.end();
  99. return;
  100. }
  101. res.writeHead(status, { 'content-type': 'application/json' });
  102. res.end(JSON.stringify({ error: { message, type: 'bridge_error' } }));
  103. }
  104. async #handle(req: IncomingMessage, res: ServerResponse): Promise<void> {
  105. if (req.method !== 'POST' || req.url !== '/v1/responses') {
  106. this.#failJson(res, 404, 'not found');
  107. return;
  108. }
  109. let responsesReq: ResponsesCreateParams;
  110. try {
  111. responsesReq = JSON.parse(await readBody(req)) as ResponsesCreateParams;
  112. } catch {
  113. this.#failJson(res, 400, '请求体不是合法 JSON');
  114. return;
  115. }
  116. const options = this.#options;
  117. if (!options) {
  118. this.#failJson(res, 503, '桥接未配置上游');
  119. return;
  120. }
  121. const chatReq = responsesRequestToChat(responsesReq, (line) => this.#diag(line));
  122. const upstream = `${options.upstreamBaseUrl.replace(/\/+$/u, '')}/chat/completions`;
  123. let upstreamRes: Response;
  124. try {
  125. upstreamRes = await fetch(upstream, {
  126. method: 'POST',
  127. headers: this.#forwardHeaders(req),
  128. body: JSON.stringify(chatReq),
  129. });
  130. } catch (error) {
  131. const message = error instanceof Error ? error.message : String(error);
  132. this.#diag(`chat-bridge: 上游不可达:${message}`);
  133. this.#failJson(res, 502, `桥接上游不可达:${message}`);
  134. return;
  135. }
  136. if (!upstreamRes.ok) {
  137. const text = await upstreamRes.text().catch(() => '');
  138. this.#diag(`chat-bridge: 上游 HTTP ${upstreamRes.status}:${text.slice(0, 300)}`);
  139. this.#failJson(res, upstreamRes.status, `桥接上游返回 HTTP ${upstreamRes.status}:${text.slice(0, 500)}`);
  140. return;
  141. }
  142. if (!chatReq.stream) {
  143. try {
  144. const chatJson = (await upstreamRes.json()) as Parameters<typeof chatResponseToResponses>[0];
  145. res.writeHead(200, { 'content-type': 'application/json' });
  146. res.end(JSON.stringify(chatResponseToResponses(chatJson, responsesReq)));
  147. } catch (error) {
  148. this.#failJson(res, 502, `桥接解析上游响应失败:${error instanceof Error ? error.message : String(error)}`);
  149. }
  150. return;
  151. }
  152. // 流式:边收边翻边发
  153. res.writeHead(200, {
  154. 'content-type': 'text/event-stream',
  155. 'cache-control': 'no-cache',
  156. connection: 'keep-alive',
  157. });
  158. const translator = new ResponsesSseTranslator(responsesReq, (line) => this.#diag(line));
  159. try {
  160. if (!upstreamRes.body) throw new Error('上游响应没有 body');
  161. const reader = upstreamRes.body.getReader();
  162. const decoder = new TextDecoder();
  163. for (;;) {
  164. const { done, value } = await reader.read();
  165. if (done) break;
  166. writeSseEvents(res, translator.push(decoder.decode(value, { stream: true })));
  167. }
  168. writeSseEvents(res, translator.push(decoder.decode()));
  169. writeSseEvents(res, translator.finish());
  170. } catch (error) {
  171. const message = error instanceof Error ? error.message : String(error);
  172. this.#diag(`chat-bridge: 流式转发中断:${message}`);
  173. writeSseEvents(res, translator.fail(message));
  174. }
  175. res.end();
  176. }
  177. }
  178. /** 模块级单例:providerService 应用/清除模型时启停,装配层退出时兜底 stop */
  179. export const chatBridge = new ChatBridgeService();