import assert from 'node:assert/strict'; import { setTimeout as delay } from 'node:timers/promises'; import { fileURLToPath } from 'node:url'; import { test } from 'node:test'; import { Chat } from '@ai-sdk/react'; import { DefaultChatTransport } from 'ai'; import { createJiti } from 'jiti'; const jiti = createJiti(import.meta.url, { alias: { '@': fileURLToPath(new URL('../src/', import.meta.url)) }, jsx: true, }); async function waitFor(predicate) { const deadline = Date.now() + 3000; while (!predicate()) { assert.ok(Date.now() < deadline, 'Expected streaming update before timeout'); await delay(5); } } function textOf(chat) { return chat.messages.at(-1)?.parts .filter((part) => part.type === 'text') .map((part) => part.text) .join(''); } async function setup(t) { const environment = { OPENROUTER_API_KEY: process.env.OPENROUTER_API_KEY, FIRECRAWL_API_KEY: process.env.FIRECRAWL_API_KEY, OPENROUTER_MODEL: process.env.OPENROUTER_MODEL, }; process.env.OPENROUTER_API_KEY = 'streaming-test-key'; delete process.env.FIRECRAWL_API_KEY; t.after(() => { for (const [key, value] of Object.entries(environment)) { if (value === undefined) delete process.env[key]; else process.env[key] = value; } }); let controller; let upstreamSignal; let response; let upstreamRequests = 0; let upstreamModel; let closed = false; const encoder = new TextEncoder(); const { default: api } = await jiti.import('../src/server/app.ts'); const POST = (request) => api.fetch(request); // Exercise the real route, OpenRouter adapter, and client transport. Only the // external provider is replaced, so no credentials or network are required. t.mock.method(globalThis, 'fetch', async (url, init) => { assert.equal(String(url), 'https://openrouter.ai/api/v1/chat/completions'); assert.equal(JSON.parse(init.body).stream, true); upstreamModel = JSON.parse(init.body).model; upstreamRequests += 1; upstreamSignal = init.signal; return new Response(new ReadableStream({ start(value) { controller = value; upstreamSignal?.addEventListener('abort', () => { if (!closed) { closed = true; controller.error(new DOMException('Aborted', 'AbortError')); } }, { once: true }); }, cancel() { closed = true; }, }), { headers: { 'Content-Type': 'text/event-stream' } }); }); const chat = new Chat({ transport: new DefaultChatTransport({ api: 'http://localhost/api/chat', fetch: async (url, init) => { response = await POST(new Request(url, init)); return response; }, }), }); t.after(async () => { await chat.stop(); if (controller && !closed) { closed = true; controller.close(); } }); function event(data) { controller.enqueue(encoder.encode(`data: ${JSON.stringify(data)}\n\n`)); } function chunk(delta, finishReason = null) { event({ id: 'test-generation', object: 'chat.completion.chunk', created: 1, model: 'test-model', choices: [{ index: 0, delta, finish_reason: finishReason }], }); } return { chat, get response() { return response; }, get upstreamRequests() { return upstreamRequests; }, get upstreamModel() { return upstreamModel; }, get upstreamSignal() { return upstreamSignal; }, ready: () => waitFor(() => Boolean(controller)), chunk, finish() { chunk({}, 'stop'); controller.enqueue(encoder.encode('data: [DONE]\n\n')); closed = true; controller.close(); }, fail() { event({ error: { code: 502, message: 'Provider disconnected' } }); closed = true; controller.close(); }, }; } test('streams reasoning and text to the client before the provider finishes', async (t) => { const stream = await setup(t); let completed = false; const sending = stream.chat.sendMessage({ text: 'Hello' }).then(() => { completed = true; }); await stream.ready(); assert.match(stream.response.headers.get('content-type'), /text\/event-stream/); assert.match(stream.response.headers.get('cache-control'), /no-transform/); assert.equal(stream.response.headers.get('x-accel-buffering'), 'no'); stream.chunk({ reasoning: 'Thinking through the answer.' }); await waitFor(() => stream.chat.messages.at(-1)?.parts.some( (part) => part.type === 'reasoning' && part.text === 'Thinking through the answer.' )); stream.chunk({ content: 'Hello' }); await waitFor(() => textOf(stream.chat) === 'Hello'); assert.equal(stream.chat.status, 'streaming'); assert.equal(completed, false); stream.chunk({ content: ' world' }); await waitFor(() => textOf(stream.chat) === 'Hello world'); assert.equal(completed, false); stream.finish(); await sending; assert.equal(stream.chat.status, 'ready'); assert.equal(textOf(stream.chat), 'Hello world'); assert.equal(stream.upstreamRequests, 1); }); for (const withPartialText of [false, true]) { test(`stop aborts the upstream request ${withPartialText ? 'after a text chunk' : 'before the first token'}`, async (t) => { const stream = await setup(t); const errors = t.mock.method(console, 'error', () => {}); const sending = stream.chat.sendMessage({ text: 'Write a long response' }); await stream.ready(); if (withPartialText) { stream.chunk({ content: 'Partial answer' }); await waitFor(() => textOf(stream.chat) === 'Partial answer'); } await stream.chat.stop(); await sending; assert.equal(stream.upstreamSignal.aborted, true); assert.equal(stream.chat.status, 'ready'); assert.equal(stream.chat.error, undefined); if (withPartialText) assert.equal(textOf(stream.chat), 'Partial answer'); assert.equal(stream.upstreamRequests, 1); assert.equal(errors.mock.callCount(), 0); }); } test('a mid-stream provider error preserves partial text and reaches the client', async (t) => { const stream = await setup(t); t.mock.method(console, 'error', () => {}); const sending = stream.chat.sendMessage({ text: 'Hello' }); await stream.ready(); stream.chunk({ content: 'Partial answer' }); await waitFor(() => textOf(stream.chat) === 'Partial answer'); stream.fail(); await sending; assert.equal(stream.chat.status, 'error'); assert.equal(stream.chat.error.message, 'The response was interrupted. Please try again.'); assert.equal(textOf(stream.chat), 'Partial answer'); }); test('missing credentials return an error without contacting the provider', async (t) => { const stream = await setup(t); delete process.env.OPENROUTER_API_KEY; await stream.chat.sendMessage({ text: 'Hello' }); assert.equal(stream.response.status, 503); assert.equal(stream.chat.status, 'error'); assert.equal(stream.upstreamRequests, 0); }); test('changing the selected model routes each message to that OpenRouter model', async (t) => { const stream = await setup(t); for (const [index, modelId] of ['anthropic/claude-opus-5', 'google/gemini-3.8-flash'].entries()) { const sending = stream.chat.sendMessage({ text: `Message ${index + 1}` }, { body: { modelId } }); await waitFor(() => stream.upstreamRequests === index + 1); assert.equal(stream.upstreamModel, modelId); stream.chunk({ content: 'Response from the selected model.' }); stream.finish(); await sending; assert.equal(stream.chat.status, 'ready'); } }); test('requests without a selection use the configured catalog model', async (t) => { const stream = await setup(t); process.env.OPENROUTER_MODEL = 'deepseek/deepseek-v4-flash-0731'; const sending = stream.chat.sendMessage({ text: 'Hello' }); await stream.ready(); assert.equal(stream.upstreamModel, 'deepseek/deepseek-v4-flash-0731'); stream.finish(); await sending; }); for (const modelId of ['unknown/model', null, { id: 'anthropic/claude-opus-5' }]) { test(`rejects an invalid model selection (${JSON.stringify(modelId)}) without contacting OpenRouter`, async (t) => { const stream = await setup(t); await stream.chat.sendMessage({ text: 'Hello' }, { body: { modelId } }); assert.equal(stream.response.status, 400); assert.equal(stream.chat.status, 'error'); assert.equal(stream.upstreamRequests, 0); }); }