Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 13 additions & 5 deletions src/client/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,10 @@ export async function* parseSSE(response: Response): AsyncGenerator<ServerSentEv

const decoder = new TextDecoder();
let buffer = '';
let skipLeadingLF = false;
let event: Partial<ServerSentEvent> = {};

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 }
Expand Down Expand Up @@ -51,9 +50,18 @@ export async function* parseSSE(response: Response): AsyncGenerator<ServerSentEv
const { done, value } = await reader.read();
if (done) break;

buffer += decoder.decode(value, { stream: true });
let chunk = decoder.decode(value, { stream: true });
if (chunk.length === 0) continue;

// A CR completes its line immediately; swallow a paired LF in the next chunk.
if (skipLeadingLF) {
if (chunk.startsWith('\n')) chunk = chunk.slice(1);
skipLeadingLF = false;
}
buffer += chunk;
skipLeadingLF = buffer.endsWith('\r');

const lines = buffer.split('\n');
const lines = buffer.split(/\r\n|\r|\n/);
buffer = lines.pop() || '';

for (const line of lines) {
Expand Down
69 changes: 69 additions & 0 deletions test/client/stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,75 @@ describe('parseSSE', () => {
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<Uint8Array>({
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<Uint8Array>({
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<Uint8Array>;
const body = new ReadableStream<Uint8Array>({
start(streamController) { controller = streamController; },
});
const response = new Response(body);
const iterator = parseSSE(response);
const next = iterator.next();
let timeout: ReturnType<typeof setTimeout> | undefined;
controller.enqueue(new TextEncoder().encode('data: hello\r\r'));

try {
expect(await Promise.race([
next,
new Promise<never>((_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
// -------------------------------------------------------------------------
Expand Down