From 61f630c60d42734280b5b2bc8387584e457ec00c Mon Sep 17 00:00:00 2001 From: Peter Dave Hello <3691490+PeterDaveHello@users.noreply.github.com> Date: Fri, 11 Sep 2026 23:09:23 +0800 Subject: [PATCH] Decode SSE chunks incrementally --- src/utils/eventsource-parser.mjs | 9 +-- tests/unit/utils/eventsource-parser.test.mjs | 72 +++++++++++++++++++- 2 files changed, 74 insertions(+), 7 deletions(-) diff --git a/src/utils/eventsource-parser.mjs b/src/utils/eventsource-parser.mjs index 31e019e92..84d8f6dfd 100644 --- a/src/utils/eventsource-parser.mjs +++ b/src/utils/eventsource-parser.mjs @@ -2,7 +2,7 @@ function createParser(onParse) { let isFirstChunk - let bytes + let decoder let buffer let startingPosition let startingFieldLength @@ -18,7 +18,7 @@ function createParser(onParse) { } function reset() { isFirstChunk = true - bytes = [] + decoder = new TextDecoder() buffer = '' startingPosition = 0 startingFieldLength = -1 @@ -29,8 +29,7 @@ function createParser(onParse) { } function feed(chunk) { - bytes = bytes.concat(Array.from(chunk)) - buffer = new TextDecoder().decode(new Uint8Array(bytes)) + buffer += decoder.decode(chunk, { stream: true }) if (isFirstChunk && hasBom(buffer)) { buffer = buffer.slice(BOM.length) } @@ -70,10 +69,8 @@ function createParser(onParse) { position += lineLength + 1 } if (position === length) { - bytes = [] buffer = '' } else if (position > 0) { - bytes = bytes.slice(new TextEncoder().encode(buffer.slice(0, position)).length) buffer = buffer.slice(position) } } diff --git a/tests/unit/utils/eventsource-parser.test.mjs b/tests/unit/utils/eventsource-parser.test.mjs index debe86ee0..6ca24ccfb 100644 --- a/tests/unit/utils/eventsource-parser.test.mjs +++ b/tests/unit/utils/eventsource-parser.test.mjs @@ -115,7 +115,22 @@ test('createParser preserves split CRLF state across empty chunks', () => { test('createParser produces the same events at every single chunk boundary', () => { const stream = toBytes('data: alpha\r\n\ndata: beta\ndata: gamma\r\r') - const expected = parseChunks(stream) + const expected = [ + { + type: 'event', + id: undefined, + event: undefined, + data: 'alpha', + extra: undefined, + }, + { + type: 'event', + id: undefined, + event: undefined, + data: 'beta\ngamma', + extra: undefined, + }, + ] for (let split = 0; split <= stream.length; ++split) { const actual = parseChunks(stream.slice(0, split), stream.slice(split)) @@ -123,6 +138,61 @@ test('createParser produces the same events at every single chunk boundary', () } }) +test('createParser preserves UTF-8 data at every byte boundary', () => { + const stream = toBytes('data: 台灣🙂 café\n\n') + const expected = [ + { + type: 'event', + id: undefined, + event: undefined, + data: '台灣🙂 café', + extra: undefined, + }, + ] + + for (let split = 0; split <= stream.length; ++split) { + const actual = parseChunks(stream.slice(0, split), stream.slice(split)) + assert.deepEqual(actual, expected, `split at byte ${split}`) + } +}) + +test('createParser preserves pending data after a leading UTF-8 BOM', () => { + const parsed = parseChunks(toBytes('\uFEFFdata: a\n\ndata:'), toBytes(' b\n\n')) + + assert.deepEqual( + parsed.map((event) => event.data), + ['a', 'b'], + ) +}) + +test('createParser preserves pending data after invalid UTF-8 replacement', () => { + const firstChunk = new Uint8Array([ + ...toBytes('data: '), + 0xff, + ...toBytes('\n\ndata:'), + ]) + const parsed = parseChunks(firstChunk, toBytes(' b\n\n')) + + assert.deepEqual( + parsed.map((event) => event.data), + ['�', 'b'], + ) +}) + +test('createParser reset discards pending decoder bytes', () => { + const parsed = [] + const parser = createParser((event) => parsed.push(event)) + + parser.feed(toBytes('🙂').slice(0, 2)) + parser.reset() + parser.feed(toBytes('data: clean\n\n')) + + assert.deepEqual( + parsed.map((event) => event.data), + ['clean'], + ) +}) + test('createParser handles \\r only line endings', () => { const parsed = [] const parser = createParser((event) => parsed.push(event))