| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204 |
- import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'node:http';
- import {
- chatResponseToResponses,
- responsesRequestToChat,
- ResponsesSseTranslator,
- type BridgeSseEvent,
- type DiagnosticFn,
- type ResponsesCreateParams,
- } from './chatBridgeTranslate';
- /**
- * 内置 Responses→Chat 桥接:给只会 /chat/completions 的端点(老 Ollama / vLLM / Xinference)
- * 套一层本地翻译代理。Codex 连这里的 /v1/responses,桥接翻译后转发上游。
- *
- * 安全边界:只绑 127.0.0.1 随机端口;不持有任何密钥(incoming Authorization 原样透传,
- * key 由 Codex 侧 env_key 注入);只接受 POST /v1/responses,其余一律 404。
- * 消息格式转换全部在 chatBridgeTranslate.ts,这里只管 HTTP 与生命周期。
- */
- export interface ChatBridgeOptions {
- /** 上游 chat 端点 API 根(形如 http://host:11434/v1) */
- upstreamBaseUrl: string;
- /** 额外转发头(来自模型配置的 headersJson) */
- headers?: Record<string, string> | null;
- onDiagnostic?: DiagnosticFn;
- }
- export interface ChatBridgeInfo {
- url: string;
- port: number;
- upstreamBaseUrl: string;
- }
- function readBody(req: IncomingMessage): Promise<string> {
- return new Promise((resolve, reject) => {
- const chunks: Buffer[] = [];
- req.on('data', (chunk: Buffer) => chunks.push(chunk));
- req.on('end', () => resolve(Buffer.concat(chunks).toString('utf8')));
- req.on('error', reject);
- });
- }
- function writeSseEvents(res: ServerResponse, events: BridgeSseEvent[]): void {
- for (const item of events) {
- res.write(`event: ${item.event}\ndata: ${JSON.stringify(item.data)}\n\n`);
- }
- }
- export class ChatBridgeService {
- #server: Server | null = null;
- #info: ChatBridgeInfo | null = null;
- #options: ChatBridgeOptions | null = null;
- /** 装配层(index.ts)注入的诊断出口;start 时的 onDiagnostic 优先 */
- #diagnostic: DiagnosticFn | null = null;
- get info(): ChatBridgeInfo | null {
- return this.#info;
- }
- setDiagnostic(fn: DiagnosticFn | null): void {
- this.#diagnostic = fn;
- }
- #diag(line: string): void {
- (this.#options?.onDiagnostic ?? this.#diagnostic)?.(line);
- }
- async start(options: ChatBridgeOptions): Promise<ChatBridgeInfo> {
- await this.stop();
- this.#options = options;
- const server = createServer((req, res) => {
- void this.#handle(req, res);
- });
- this.#server = server;
- await new Promise<void>((resolve, reject) => {
- server.once('error', reject);
- server.listen(0, '127.0.0.1', () => resolve());
- });
- const address = server.address();
- const port = typeof address === 'object' && address ? address.port : 0;
- this.#info = { url: `http://127.0.0.1:${port}/v1`, port, upstreamBaseUrl: options.upstreamBaseUrl };
- this.#diag(`chat-bridge: 已启动 ${this.#info.url} → ${options.upstreamBaseUrl}`);
- return this.#info;
- }
- async stop(): Promise<void> {
- const server = this.#server;
- this.#server = null;
- this.#info = null;
- this.#options = null;
- if (!server) return;
- await new Promise<void>((resolve) => server.close(() => resolve()));
- }
- #forwardHeaders(req: IncomingMessage): Record<string, string> {
- const headers: Record<string, string> = {
- 'content-type': 'application/json',
- // 上游 gzip 会把 SSE 行切碎,显式要求不压缩
- 'accept-encoding': 'identity',
- };
- const authorization = req.headers.authorization;
- if (authorization) headers.authorization = authorization;
- for (const [key, value] of Object.entries(this.#options?.headers ?? {})) {
- if (value) headers[key.toLowerCase()] = value;
- }
- return headers;
- }
- #failJson(res: ServerResponse, status: number, message: string): void {
- if (res.headersSent) {
- res.end();
- return;
- }
- res.writeHead(status, { 'content-type': 'application/json' });
- res.end(JSON.stringify({ error: { message, type: 'bridge_error' } }));
- }
- async #handle(req: IncomingMessage, res: ServerResponse): Promise<void> {
- if (req.method !== 'POST' || req.url !== '/v1/responses') {
- this.#failJson(res, 404, 'not found');
- return;
- }
- let responsesReq: ResponsesCreateParams;
- try {
- responsesReq = JSON.parse(await readBody(req)) as ResponsesCreateParams;
- } catch {
- this.#failJson(res, 400, '请求体不是合法 JSON');
- return;
- }
- const options = this.#options;
- if (!options) {
- this.#failJson(res, 503, '桥接未配置上游');
- return;
- }
- const chatReq = responsesRequestToChat(responsesReq, (line) => this.#diag(line));
- const upstream = `${options.upstreamBaseUrl.replace(/\/+$/u, '')}/chat/completions`;
- let upstreamRes: Response;
- try {
- upstreamRes = await fetch(upstream, {
- method: 'POST',
- headers: this.#forwardHeaders(req),
- body: JSON.stringify(chatReq),
- });
- } catch (error) {
- const message = error instanceof Error ? error.message : String(error);
- this.#diag(`chat-bridge: 上游不可达:${message}`);
- this.#failJson(res, 502, `桥接上游不可达:${message}`);
- return;
- }
- if (!upstreamRes.ok) {
- const text = await upstreamRes.text().catch(() => '');
- this.#diag(`chat-bridge: 上游 HTTP ${upstreamRes.status}:${text.slice(0, 300)}`);
- this.#failJson(res, upstreamRes.status, `桥接上游返回 HTTP ${upstreamRes.status}:${text.slice(0, 500)}`);
- return;
- }
- if (!chatReq.stream) {
- try {
- const chatJson = (await upstreamRes.json()) as Parameters<typeof chatResponseToResponses>[0];
- res.writeHead(200, { 'content-type': 'application/json' });
- res.end(JSON.stringify(chatResponseToResponses(chatJson, responsesReq)));
- } catch (error) {
- this.#failJson(res, 502, `桥接解析上游响应失败:${error instanceof Error ? error.message : String(error)}`);
- }
- return;
- }
- // 流式:边收边翻边发
- res.writeHead(200, {
- 'content-type': 'text/event-stream',
- 'cache-control': 'no-cache',
- connection: 'keep-alive',
- });
- const translator = new ResponsesSseTranslator(responsesReq, (line) => this.#diag(line));
- try {
- if (!upstreamRes.body) throw new Error('上游响应没有 body');
- const reader = upstreamRes.body.getReader();
- const decoder = new TextDecoder();
- for (;;) {
- const { done, value } = await reader.read();
- if (done) break;
- writeSseEvents(res, translator.push(decoder.decode(value, { stream: true })));
- }
- writeSseEvents(res, translator.push(decoder.decode()));
- writeSseEvents(res, translator.finish());
- } catch (error) {
- const message = error instanceof Error ? error.message : String(error);
- this.#diag(`chat-bridge: 流式转发中断:${message}`);
- writeSseEvents(res, translator.fail(message));
- }
- res.end();
- }
- }
- /** 模块级单例:providerService 应用/清除模型时启停,装配层退出时兜底 stop */
- export const chatBridge = new ChatBridgeService();
|