import assert from 'node:assert/strict'; import test from 'node:test'; import { ReadableStream } from 'node:stream/web'; import { classifySessionProviderErrorMessage, consumeSessionEventsUntilFinish, eventMatchesRequest, parseSessionStreamEvent, } from './session-reply-wait.mjs'; test('parseSessionStreamEvent parses SSE data payload', () => { const event = parseSessionStreamEvent('id: 1\ndata: {"type":"Finish","token_state":null}\n'); assert.equal(event?.type, 'Finish'); }); test('eventMatchesRequest scopes by request id when present', () => { assert.equal(eventMatchesRequest({ request_id: 'req-1' }, 'req-1'), true); assert.equal(eventMatchesRequest({ request_id: 'req-2' }, 'req-1'), false); assert.equal(eventMatchesRequest({ type: 'Finish' }, 'req-1'), true); }); test('consumeSessionEventsUntilFinish resolves on Finish', async () => { const frames = [ 'data: {"type":"Message","request_id":"req-1","message":{"role":"assistant"}}\n\n', 'data: {"type":"Finish","request_id":"req-1","token_state":{"totalTokens":12}}\n\n', ]; const stream = new ReadableStream({ start(controller) { for (const frame of frames) controller.enqueue(new TextEncoder().encode(frame)); controller.close(); }, }); const result = await consumeSessionEventsUntilFinish(stream, { requestId: 'req-1', timeoutMs: 5000 }); assert.equal(result.finishEvent.type, 'Finish'); assert.equal(result.tokenState.totalTokens, 12); }); test('consumeSessionEventsUntilFinish ignores a stale unscoped Finish before request activity', async () => { const frames = [ 'data: {"type":"Finish","token_state":{"totalTokens":0}}\n\n', 'data: {"type":"ActiveRequests","request_id":"req-fresh","request_ids":["req-fresh"]}\n\n', 'data: {"type":"Finish","request_id":"req-fresh","token_state":{"totalTokens":0}}\n\n', 'data: {"type":"Message","request_id":"req-fresh","message":{"role":"assistant"}}\n\n', 'data: {"type":"Finish","token_state":{"totalTokens":12}}\n\n', ]; const stream = new ReadableStream({ start(controller) { for (const frame of frames) controller.enqueue(new TextEncoder().encode(frame)); controller.close(); }, }); const result = await consumeSessionEventsUntilFinish(stream, { requestId: 'req-fresh', timeoutMs: 5000, requireRequestActivityBeforeUnscopedFinish: true, }); assert.equal(result.tokenState.totalTokens, 12); }); test('consumeSessionEventsUntilFinish fails fast when an empty Finish has no later activity', async () => { const stream = new ReadableStream({ start(controller) { controller.enqueue(new TextEncoder().encode( 'data: {"type":"Finish","request_id":"req-empty","token_state":{"totalTokens":0}}\n\n', )); }, }); await assert.rejects( consumeSessionEventsUntilFinish(stream, { requestId: 'req-empty', timeoutMs: 5000, emptyFinishGraceMs: 5, requireRequestActivityBeforeUnscopedFinish: true, }), (error) => { assert.equal(error.code, 'SESSION_EMPTY_FINISH'); return true; }, ); }); test('consumeSessionEventsUntilFinish records a successful raster generate_image result', async () => { const toolResult = JSON.stringify({ ok: true, jobId: 'job-1', source: { mimeType: 'image/webp' }, asset: { id: 'asset-1', publicUrl: '/MindSpace/user/public/images/hero.webp', workspaceRelativePath: 'public/images/hero.webp', }, }); const frames = [ `data: ${JSON.stringify({ type: 'Message', request_id: 'req-1', message: { content: [{ type: 'toolRequest', id: 'call-1', toolCall: { value: { name: 'sandbox-fs__generate_image', arguments: { purpose: 'hero' } } }, }], }, })}\n\n`, `data: ${JSON.stringify({ type: 'Message', request_id: 'req-1', message: { content: [{ type: 'toolResponse', id: 'call-1', toolResult: { status: 'success', value: { content: [{ type: 'text', text: toolResult }], isError: false } }, }], }, })}\n\n`, 'data: {"type":"Finish","request_id":"req-1"}\n\n', ]; const stream = new ReadableStream({ start(controller) { for (const frame of frames) controller.enqueue(new TextEncoder().encode(frame)); controller.close(); }, }); const result = await consumeSessionEventsUntilFinish(stream, { requestId: 'req-1', timeoutMs: 5000 }); assert.equal(result.toolEvidence.generateImage.called, true); assert.equal(result.toolEvidence.generateImage.succeeded, true); assert.equal(result.toolEvidence.generateImage.jobId, 'job-1'); }); test('consumeSessionEventsUntilFinish surfaces reasoning_content protocol errors before Finish', async () => { const frames = [ `data: ${JSON.stringify({ type: 'Message', request_id: 'req-1', message: { role: 'assistant', content: [{ type: 'text', text: 'Ran into this error: Bad request (400): The reasoning_content in the thinking mode must be passed back to the API.\n\nPlease retry if you think this is a transient or recoverable error.', }], }, })}\n\n`, 'data: {"type":"Finish","request_id":"req-1"}\n\n', ]; const stream = new ReadableStream({ start(controller) { for (const frame of frames) controller.enqueue(new TextEncoder().encode(frame)); controller.close(); }, }); await assert.rejects( consumeSessionEventsUntilFinish(stream, { requestId: 'req-1', timeoutMs: 5000 }), (error) => { assert.equal(error.code, 'SESSION_REASONING_CONTENT_POISONED'); assert.match(error.message, /reasoning_content/); return true; }, ); }); test('consumeSessionEventsUntilFinish surfaces provider tool history errors before Finish', async () => { const frames = [ `data: ${JSON.stringify({ type: 'Message', request_id: 'req-1', message: { role: 'assistant', content: [{ type: 'text', text: "Ran into this error: Request failed: Bad request (400): An assistant message with 'tool_calls' must be followed by tool messages responding to each 'tool_call_id'. (insufficient tool messages following tool_calls message).\n\nPlease retry if you think this is a transient or recoverable error.", }], }, })}\n\n`, 'data: {"type":"Finish","request_id":"req-1"}\n\n', ]; const stream = new ReadableStream({ start(controller) { for (const frame of frames) controller.enqueue(new TextEncoder().encode(frame)); controller.close(); }, }); await assert.rejects( consumeSessionEventsUntilFinish(stream, { requestId: 'req-1', timeoutMs: 5000 }), (error) => { assert.equal(error.code, 'SESSION_TOOL_HISTORY_POISONED'); assert.match(error.message, /tool_calls/); return true; }, ); }); test('consumeSessionEventsUntilFinish classifies unsupported visual history as degradable', async () => { const message = 'Request failed: Bad request (400): Failed to deserialize the JSON body into the target type: messages[8]: unknown variant `image_url`, expected `text`'; assert.equal( classifySessionProviderErrorMessage(message), 'SESSION_VISUAL_CONTEXT_UNSUPPORTED', ); const frames = [ `data: ${JSON.stringify({ type: 'Message', request_id: 'req-visual', message: { role: 'assistant', content: [{ type: 'text', text: `Ran into this error: ${message}\n\nPlease retry if you think this is a transient or recoverable error.`, }], }, })}\n\n`, 'data: {"type":"Finish","request_id":"req-visual"}\n\n', ]; const stream = new ReadableStream({ start(controller) { for (const frame of frames) controller.enqueue(new TextEncoder().encode(frame)); controller.close(); }, }); await assert.rejects( consumeSessionEventsUntilFinish(stream, { requestId: 'req-visual', timeoutMs: 5000, }), (error) => { assert.equal(error.code, 'SESSION_VISUAL_CONTEXT_UNSUPPORTED'); assert.equal(error.retryable, false); return true; }, ); });