import { createServer, type Server } from 'node:http'; import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { ChatBridgeService } from './chatBridgeService'; /** mock 上游:记录最近一次 /chat/completions 请求,按配置返回非流式或 SSE 流式 */ let upstream: Server; let upstreamPort = 0; let lastUpstream: { url: string | null; auth: string | null; extra: string | null; body: Record | null }; let upstreamMode: 'json' | 'sse' | 'error' = 'json'; const SSE_BODY = [ `data: ${JSON.stringify({ choices: [{ delta: { content: '你好' }, finish_reason: null }] })}`, '', `data: ${JSON.stringify({ choices: [{ delta: { content: '世界' }, finish_reason: null }] })}`, '', `data: ${JSON.stringify({ choices: [{ delta: {}, finish_reason: 'stop' }], usage: { prompt_tokens: 5, completion_tokens: 2, total_tokens: 7 } })}`, '', 'data: [DONE]', '', '', ].join('\n'); beforeEach(async () => { upstreamMode = 'json'; lastUpstream = { url: null, auth: null, extra: null, body: null }; upstream = createServer((req, res) => { const chunks: Buffer[] = []; req.on('data', (chunk: Buffer) => chunks.push(chunk)); req.on('end', () => { lastUpstream = { url: req.url ?? null, auth: req.headers.authorization ?? null, extra: (req.headers['x-extra'] as string) ?? null, body: JSON.parse(Buffer.concat(chunks).toString('utf8') || 'null'), }; if (upstreamMode === 'error') { res.writeHead(500, { 'content-type': 'application/json' }); res.end(JSON.stringify({ error: 'boom' })); return; } if (upstreamMode === 'sse') { res.writeHead(200, { 'content-type': 'text/event-stream' }); res.end(SSE_BODY); return; } res.writeHead(200, { 'content-type': 'application/json' }); res.end( JSON.stringify({ id: 'chatcmpl-1', created: 1700000000, model: 'qwen', choices: [{ finish_reason: 'stop', message: { role: 'assistant', content: '非流式回复' } }], usage: { prompt_tokens: 3, completion_tokens: 4, total_tokens: 7 }, }), ); }); }); await new Promise((resolve) => upstream.listen(0, '127.0.0.1', resolve)); const address = upstream.address(); upstreamPort = typeof address === 'object' && address ? address.port : 0; }); afterEach(async () => { await new Promise((resolve) => upstream.close(() => resolve())); }); function makeBridge(): Promise { return Promise.resolve(new ChatBridgeService()); } describe('ChatBridgeService', () => { it('非流式:请求翻译后转发上游,响应翻译回 Responses 形状,头原样透传', async () => { const bridge = await makeBridge(); const info = await bridge.start({ upstreamBaseUrl: `http://127.0.0.1:${upstreamPort}/v1`, headers: { 'X-Extra': 'yes' }, }); try { const response = await fetch(`${info.url}/responses`, { method: 'POST', headers: { 'content-type': 'application/json', authorization: 'Bearer sk-live' }, body: JSON.stringify({ model: 'qwen', instructions: '你是助手', input: '你好', max_output_tokens: 64, }), }); expect(response.status).toBe(200); // 上游收到的是 chat 形状 expect(lastUpstream.url).toBe('/v1/chat/completions'); expect(lastUpstream.auth).toBe('Bearer sk-live'); expect(lastUpstream.extra).toBe('yes'); expect(lastUpstream.body).toMatchObject({ model: 'qwen', max_tokens: 64, messages: [ { role: 'developer', content: '你是助手' }, { role: 'user', content: '你好' }, ], }); // 下游拿到的是 responses 形状 const body = (await response.json()) as Record; expect(body).toMatchObject({ object: 'response', status: 'completed', usage: { input_tokens: 3, output_tokens: 4, total_tokens: 7 }, }); const output = body.output as Array>; expect((output[0].content as Array<{ text: string }>)[0].text).toBe('非流式回复'); } finally { await bridge.stop(); } }); it('流式:上游 SSE 逐事件翻译下发,completed 带 usage', async () => { upstreamMode = 'sse'; const bridge = await makeBridge(); const info = await bridge.start({ upstreamBaseUrl: `http://127.0.0.1:${upstreamPort}/v1` }); try { const response = await fetch(`${info.url}/responses`, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ model: 'qwen', input: '你好', stream: true }), }); expect(response.status).toBe(200); expect(response.headers.get('content-type')).toContain('text/event-stream'); const text = await response.text(); const events = text .split('\n\n') .filter((block) => block.trim()) .map((block) => { const eventMatch = block.match(/^event: (.+)$/mu); const dataMatch = block.match(/^data: (.+)$/mu); return { event: eventMatch?.[1], data: JSON.parse(dataMatch?.[1] ?? '{}') as Record }; }); const types = events.map((item) => item.event); expect(types[0]).toBe('response.created'); expect(types.at(-1)).toBe('response.completed'); // 请求侧补了 include_usage expect(lastUpstream.body?.stream_options).toEqual({ include_usage: true }); const deltas = events .filter((item) => item.event === 'response.output_text.delta') .map((item) => item.data.delta); expect(deltas.join('')).toBe('你好世界'); const completed = events.at(-1)?.data.response as Record; expect(completed.usage).toEqual({ input_tokens: 5, output_tokens: 2, total_tokens: 7 }); } finally { await bridge.stop(); } }); it('上游非 2xx:状态码与错误体透回', async () => { upstreamMode = 'error'; const bridge = await makeBridge(); const info = await bridge.start({ upstreamBaseUrl: `http://127.0.0.1:${upstreamPort}/v1` }); try { const response = await fetch(`${info.url}/responses`, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ model: 'qwen', input: 'x' }), }); expect(response.status).toBe(500); const body = (await response.json()) as { error: { message: string } }; expect(body.error.message).toContain('500'); } finally { await bridge.stop(); } }); it('只接 POST /v1/responses;stop 后端口释放', async () => { const bridge = await makeBridge(); const info = await bridge.start({ upstreamBaseUrl: `http://127.0.0.1:${upstreamPort}/v1` }); const wrong = await fetch(`${info.url}/models`); expect(wrong.status).toBe(404); await bridge.stop(); await expect(fetch(`${info.url}/responses`, { method: 'POST' })).rejects.toThrow(); }); }); /** * Codex 等不到就会挂断重来,而自部署的多是单槽推理: * 桥接不取消上游,被放弃的那次生成就继续占着槽,下一次请求排在它后面 —— 越重试越慢。 */ describe('Codex 提前挂断时取消上游', () => { it('上游请求在挂断后立即中止', async () => { let diedAt = 0; const neverEnds = createServer((req, res) => { req.resume(); req.on('end', () => { res.writeHead(200, { 'content-type': 'text/event-stream' }); res.write(`data: ${JSON.stringify({ choices: [{ delta: { content: '第一段' }, finish_reason: null }] })}\n\n`); const timer = setInterval(() => res.write(': tick\n\n'), 50); const stop = (): void => { clearInterval(timer); if (!diedAt) diedAt = Date.now(); }; res.on('close', stop); req.on('aborted', stop); }); }); await new Promise((resolve) => neverEnds.listen(0, '127.0.0.1', resolve)); const port = (neverEnds.address() as { port: number }).port; const bridge = new ChatBridgeService(); try { const info = await bridge.start({ upstreamBaseUrl: `http://127.0.0.1:${port}/v1` }); const controller = new AbortController(); const response = await fetch(`${info.url}/responses`, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ model: 'm', input: '讲个长故事', stream: true }), signal: controller.signal, }); // 收到第一个事件再挂断,确保桥接确实已经在上游吞吐数据 await response.body?.getReader().read(); const abortAt = Date.now(); controller.abort(); const deadline = abortAt + 2_000; while (!diedAt && Date.now() < deadline) { await new Promise((resolve) => setTimeout(resolve, 25)); } expect(diedAt).toBeGreaterThan(0); expect(diedAt - abortAt).toBeLessThan(1_000); } finally { await bridge.stop(); await new Promise((resolve) => neverEnds.close(() => resolve())); } }, 20_000); }); /** * 单槽自部署服务光 prefill 就要几分钟才吐首个字节。桥接要是等上游响应头到手才给 Codex 写东西, * Codex 面对的就是一条一个字节都没有的死连接 —— 它会掐了重发,于是越重试越慢。 * 流式请求一到就先把 SSE 建起来,上游慢不影响客户端侧有字节可收。 */ describe('上游首字节很慢时先把流建起来', () => { it('立刻回 response.created,不等上游', async () => { const slowUpstream = createServer((req, res) => { req.resume(); req.on('end', () => { setTimeout(() => { res.writeHead(200, { 'content-type': 'text/event-stream' }); res.end(SSE_BODY); }, 6_000); }); }); await new Promise((resolve) => slowUpstream.listen(0, '127.0.0.1', resolve)); const port = (slowUpstream.address() as { port: number }).port; const bridge = new ChatBridgeService(); const controller = new AbortController(); try { const info = await bridge.start({ upstreamBaseUrl: `http://127.0.0.1:${port}/v1` }); const startedAt = Date.now(); const response = await fetch(`${info.url}/responses`, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ model: 'm', input: '讲个故事', stream: true }), signal: controller.signal, }); const reader = response.body?.getReader(); if (!reader) throw new Error('没有响应体'); const first = await reader.read(); const firstAt = Date.now() - startedAt; const firstText = new TextDecoder().decode(first.value ?? new Uint8Array()); expect(firstText).toContain('event: response.created'); expect(firstAt).toBeLessThan(1_500); // 流没断:上游 6 秒后才回,最终回复照样要翻给客户端 let all = firstText; for (;;) { const { done, value } = await reader.read(); if (done) break; all += new TextDecoder().decode(value); } expect(all).toContain('你好'); expect(all).toContain('event: response.completed'); } finally { controller.abort(); await bridge.stop(); await new Promise((resolve) => slowUpstream.close(() => resolve())); } }, 30_000); });