From ee3bf285877efb7c4ab47bfed2230007c25ca365 Mon Sep 17 00:00:00 2001 From: Nic Polumeyv <162764842+Nic-Polumeyv@users.noreply.github.com> Date: Thu, 23 Jul 2026 12:41:36 -0400 Subject: [PATCH] perf: parse large streamed frames in linear time --- .changeset/calm-streams-flow.md | 5 ++ packages/kit/src/runtime/client/stream.js | 47 ++++++------ .../kit/src/runtime/client/stream.spec.js | 71 +++++++++++++++++++ 3 files changed, 103 insertions(+), 20 deletions(-) create mode 100644 .changeset/calm-streams-flow.md create mode 100644 packages/kit/src/runtime/client/stream.spec.js diff --git a/.changeset/calm-streams-flow.md b/.changeset/calm-streams-flow.md new file mode 100644 index 000000000000..86eb2dfe9493 --- /dev/null +++ b/.changeset/calm-streams-flow.md @@ -0,0 +1,5 @@ +--- +'@sveltejs/kit': patch +--- + +perf: parse large streamed frames in linear time diff --git a/packages/kit/src/runtime/client/stream.js b/packages/kit/src/runtime/client/stream.js index 8256f0609632..da36c7e9fec7 100644 --- a/packages/kit/src/runtime/client/stream.js +++ b/packages/kit/src/runtime/client/stream.js @@ -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; } } diff --git a/packages/kit/src/runtime/client/stream.spec.js b/packages/kit/src/runtime/client/stream.spec.js new file mode 100644 index 000000000000..5e6d03c3a342 --- /dev/null +++ b/packages/kit/src/runtime/client/stream.spec.js @@ -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} + */ +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); +});