| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179 |
- 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<string, unknown> | 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<void>((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<void>((resolve) => upstream.close(() => resolve()));
- });
- function makeBridge(): Promise<ChatBridgeService> {
- 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<string, unknown>;
- expect(body).toMatchObject({
- object: 'response',
- status: 'completed',
- usage: { input_tokens: 3, output_tokens: 4, total_tokens: 7 },
- });
- const output = body.output as Array<Record<string, unknown>>;
- 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<string, unknown> };
- });
- 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<string, unknown>;
- 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();
- });
- });
|