diff --git a/src/client/stream.ts b/src/client/stream.ts index 75027506..be8bf4ad 100644 --- a/src/client/stream.ts +++ b/src/client/stream.ts @@ -10,11 +10,10 @@ export async function* parseSSE(response: Response): AsyncGenerator = {}; - const processLine = (rawLine: string): ServerSentEvent | undefined => { - const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine; - + const processLine = (line: string): ServerSentEvent | undefined => { if (line === '') { const completed = event.data !== undefined ? { data: event.data, event: event.event, id: event.id } @@ -51,9 +50,18 @@ export async function* parseSSE(response: Response): AsyncGenerator { expect(events).toEqual([{ event: 'message', data: 'hello' }]); }); + for (const newline of ['\n', '\r\n', '\r']) { + it(`preserves events at every byte boundary with ${JSON.stringify(newline)} line endings`, async () => { + const bytes = new TextEncoder().encode([ + ': keepalive', 'id: 7', 'event: message', 'data: 你好', + 'data: world', '', 'data: second', '', '', + ].join(newline)); + + for (let boundary = 0; boundary <= bytes.length; boundary++) { + const body = new ReadableStream({ + start(controller) { + controller.enqueue(bytes.slice(0, boundary)); + controller.enqueue(new Uint8Array()); + controller.enqueue(bytes.slice(boundary)); + controller.close(); + }, + }); + + expect(await collectEvents(new Response(body))).toEqual([ + { id: '7', event: 'message', data: '你好\nworld' }, + { data: 'second' }, + ]); + } + }); + } + + it('handles mixed line endings in byte-by-byte chunks', async () => { + const bytes = new TextEncoder().encode( + 'event: message\rdata: 你好\r\ndata: world\n\rdata: second\r\n\r\n', + ); + const body = new ReadableStream({ + start(controller) { + for (const byte of bytes) controller.enqueue(new Uint8Array([byte])); + controller.close(); + }, + }); + + expect(await collectEvents(new Response(body))).toEqual([ + { event: 'message', data: '你好\nworld' }, + { data: 'second' }, + ]); + }); + + it('dispatches a CR-terminated event before the stream closes', async () => { + let controller!: ReadableStreamDefaultController; + const body = new ReadableStream({ + start(streamController) { controller = streamController; }, + }); + const response = new Response(body); + const iterator = parseSSE(response); + const next = iterator.next(); + let timeout: ReturnType | undefined; + controller.enqueue(new TextEncoder().encode('data: hello\r\r')); + + try { + expect(await Promise.race([ + next, + new Promise((_resolve, reject) => { + timeout = setTimeout(() => reject(new Error('event was not dispatched')), 1000); + }), + ])).toEqual({ value: { data: 'hello' }, done: false }); + } finally { + clearTimeout(timeout); + controller.close(); + await next; + await iterator.return(undefined); + } + expect(response.body?.locked).toBe(false); + }); + // ------------------------------------------------------------------------- // Helpers // -------------------------------------------------------------------------