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
5 changes: 5 additions & 0 deletions .changeset/calm-streams-flow.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@sveltejs/kit': patch
---

perf: parse large streamed frames in linear time
47 changes: 27 additions & 20 deletions packages/kit/src/runtime/client/stream.js
Original file line number Diff line number Diff line change
Expand Up @@ -8,32 +8,39 @@
*/
export async function* read_stream(reader, delimiter, options) {
let done = false;
let buffer = '';
/** @type {string[]} */
let parts = [];
let rest = '';
const decoder = new TextDecoder(undefined, options);
const carry = delimiter.length - 1;

while (true) {
let split = buffer.indexOf(delimiter);
while (split !== -1) {
yield buffer.slice(0, split);
buffer = buffer.slice(split + delimiter.length);
split = buffer.indexOf(delimiter);
}
while (!done) {
const chunk = await reader.read();
done = chunk.done;
let text = rest;
if (chunk.value) text += decoder.decode(chunk.value, { stream: true });
if (done) text += decoder.decode();

if (done) {
if (buffer) {
yield buffer;
let start = 0;
let split = text.indexOf(delimiter);
while (split !== -1) {
if (parts.length > 0) {
parts.push(text.slice(start, split));
yield parts.join('');
parts = [];
} else {
yield text.slice(start, split);
}
return;
start = split + delimiter.length;
split = text.indexOf(delimiter, start);
}

const chunk = await reader.read();
done = chunk.done;
if (chunk.value) {
buffer += decoder.decode(chunk.value, { stream: true });
}
const keep = Math.max(start, text.length - carry);
if (keep > start) parts.push(text.slice(start, keep));
rest = text.slice(keep);
}

if (done) {
buffer += decoder.decode();
}
if (parts.length > 0 || rest) {
yield parts.join('') + rest;
}
}
71 changes: 71 additions & 0 deletions packages/kit/src/runtime/client/stream.spec.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
import { expect, test } from 'vitest';
import { read_stream } from './stream.js';

const encoder = new TextEncoder();

/**
* @param {Uint8Array[]} chunks
* @param {string} delimiter
* @param {TextDecoderOptions} [options]
* @returns {Promise<string[]>}
*/
async function collect(chunks, delimiter, options) {
const values = [];
const reader = new ReadableStream({
start(controller) {
for (const chunk of chunks) {
controller.enqueue(chunk);
}
controller.close();
}
}).getReader();

for await (const value of read_stream(reader, delimiter, options)) {
values.push(value);
}

return values;
}

test('parses a delimiter split across chunks', async () => {
const chunks = [encoder.encode('a\n'), encoder.encode('\nb')];

await expect(collect(chunks, '\n\n')).resolves.toEqual(['a', 'b']);
});

test('parses UTF-8 code points split across chunks', async () => {
const bytes = encoder.encode('a\u00e9\n');
const chunks = Array.from(bytes, (byte) => new Uint8Array([byte]));

await expect(collect(chunks, '\n')).resolves.toEqual(['a\u00e9']);
});

test('parses a frame much larger than a chunk', async () => {
const size = 100 * 1024;
const chunks = Array.from({ length: size / 1024 }, () => encoder.encode('x'.repeat(1024)));
chunks.push(encoder.encode('\n'));

const values = await collect(chunks, '\n');
expect(values).toHaveLength(1);
expect(values[0]).toHaveLength(size);
expect(values[0]?.slice(0, 5)).toBe('xxxxx');
expect(values[0]?.slice(-5)).toBe('xxxxx');
});

test('parses short frames after a long frame', async () => {
const prefix = 'x'.repeat(100 * 1024);
const chunks = [encoder.encode(prefix), encoder.encode('end\nshort1\nshort2\n')];

await expect(collect(chunks, '\n')).resolves.toEqual([`${prefix}end`, 'short1', 'short2']);
});

test('parses a trailing frame without yielding trailing empty frames', async () => {
await expect(collect([encoder.encode('trailing')], '\n')).resolves.toEqual(['trailing']);
await expect(collect([encoder.encode('complete\n')], '\n')).resolves.toEqual(['complete']);
});

test('rejects malformed UTF-8', async () => {
const chunks = [encoder.encode('valid'), new Uint8Array([0xff]), encoder.encode('\n')];

await expect(collect(chunks, '\n', { fatal: true })).rejects.toThrow(TypeError);
});
Loading