import { describe, expect, it } from 'vitest' import { AnthropicStreamError, parseAnthropicSSE } from '../../../src/main/utils/sse-parser' import { readNdjsonLines } from '../../../src/main/utils/ndjson-reader' const encoder = new TextEncoder() function chunked(parts: string[]): ReadableStream { return new ReadableStream({ start(controller) { for (const part of parts) controller.enqueue(encoder.encode(part)) controller.close() }, }) } async function drain(stream: ReadableStream): Promise { const out: string[] = [] for await (const token of parseAnthropicSSE(stream)) out.push(token) return out } const deltaLine = (text: string): string => `data: ${JSON.stringify({ type: 'content_block_delta', delta: { type: 'text_delta', text } })}\n` describe('parseAnthropicSSE', () => { it('yields text deltas split across chunk boundaries and stops at message_stop', async () => { const body = `event: content_block_delta\n${deltaLine('Hel')}\n${deltaLine('lo')}\ndata: {"type":"message_stop"}\n\n` const parts = [body.slice(0, 17), body.slice(17, 60), body.slice(60)] expect(await drain(chunked(parts))).toEqual(['Hel', 'lo']) }) it('accepts [DONE], data: without a space, and a final line without newline', async () => { expect(await drain(chunked([deltaLine('a'), 'data: [DONE]\n']))).toEqual(['a']) expect(await drain(chunked([ `data:${JSON.stringify({ type: 'content_block_delta', delta: { type: 'text_delta', text: 'b' } })}\n`, 'data: {"type":"message_stop"}', ]))).toEqual(['b']) }) it('skips malformed JSON lines and ping events', async () => { expect(await drain(chunked([ 'data: {not json}\n', 'event: ping\ndata: {"type":"ping"}\n', deltaLine('x'), 'data: {"type":"message_stop"}\n', ]))).toEqual(['x']) }) it('throws a provider_error on an error event', async () => { const stream = chunked([ deltaLine('partial'), 'event: error\ndata: {"type":"error","error":{"type":"overloaded_error","message":"Overloaded"}}\n\n', ]) const tokens: string[] = [] const error = await (async () => { try { for await (const token of parseAnthropicSSE(stream)) tokens.push(token) return null } catch (err) { return err } })() expect(tokens).toEqual(['partial']) expect(error).toBeInstanceOf(AnthropicStreamError) expect((error as AnthropicStreamError).kind).toBe('provider_error') expect((error as AnthropicStreamError).providerErrorType).toBe('overloaded_error') }) it('throws incomplete when the stream ends before message_stop', async () => { await expect(drain(chunked([deltaLine('cut')]))).rejects.toMatchObject({ name: 'AnthropicStreamError', kind: 'incomplete', }) await expect(drain(chunked([]))).rejects.toMatchObject({ kind: 'incomplete' }) }) }) describe('readNdjsonLines', () => { it('reassembles lines across chunks, skips blanks and yields a trailing line', async () => { const reader = chunked(['{"a":', '1}\n\n \n{"b"', ':2}\n{"c":3}']).getReader() const lines: string[] = [] for await (const line of readNdjsonLines(reader)) lines.push(line) expect(lines).toEqual(['{"a":1}', '{"b":2}', '{"c":3}']) }) it('decodes multi-byte characters split across chunks', async () => { const bytes = encoder.encode('{"t":"한글"}\n') const reader = new ReadableStream({ start(controller) { controller.enqueue(bytes.slice(0, 8)) controller.enqueue(bytes.slice(8)) controller.close() }, }).getReader() const lines: string[] = [] for await (const line of readNdjsonLines(reader)) lines.push(line) expect(lines).toEqual(['{"t":"한글"}']) }) })