From b2ba0101aed3e38a7743f42205a50b85a26d91ab Mon Sep 17 00:00:00 2001 From: "Alexis H. Munsayac" Date: Wed, 7 Oct 2026 04:28:54 +0800 Subject: [PATCH 1/3] fix(web): end cancelled chunk streams cleanly and validate chunk headers - `cancel()` now ends a pending `next()` as done even partway through a frame. It used to throw "Malformed server function stream." whenever the cancel landed mid-frame. Frames reported a superseded response with that message instead of the supersession reason. - A chunk header must be exactly `;0x` plus 8 hex digits plus `;`. The old `parseInt` read accepted wrong delimiters and trailing junk such as `;0x5zzzzzzz;`. - Payloads are decoded with a fatal UTF-8 decoder, so invalid bytes fail as a malformed stream instead of becoming U+FFFD. - `drain()` and a failed first frame in `deserializeStream` cancel the body instead of leaving it locked and unread. - The reader drops an oversized store once it holds no unread bytes. One large frame used to keep its doubled allocation until the stream ended. - `createChunk` refuses a payload whose length does not fit the 8-digit header. One encoder and one decoder are shared at module level. The wire format is unchanged. Co-Authored-By: Claude Opus 5.5 --- .changeset/chunk-reader-hardening.md | 9 + packages/web/server-functions/src/shared.ts | 141 ++++++++++--- .../runtime/chunk-reader-hardening.spec.ts | 187 ++++++++++++++++++ 3 files changed, 315 insertions(+), 22 deletions(-) create mode 100644 .changeset/chunk-reader-hardening.md create mode 100644 packages/web/test/runtime/chunk-reader-hardening.spec.ts diff --git a/.changeset/chunk-reader-hardening.md b/.changeset/chunk-reader-hardening.md new file mode 100644 index 000000000..b3c00fe19 --- /dev/null +++ b/.changeset/chunk-reader-hardening.md @@ -0,0 +1,9 @@ +--- +"@solidjs/web": patch +--- + +Cancelling a server function or frames stream partway through a chunk now ends it cleanly instead of failing as a malformed stream. + +- Chunk headers must be exactly `;0x` plus 8 hex digits plus `;`, and payloads must be valid UTF-8. +- A stream that fails to decode is cancelled instead of being left unread. +- The reader frees the memory from a large chunk once it is read. diff --git a/packages/web/server-functions/src/shared.ts b/packages/web/server-functions/src/shared.ts index 75dbdcf65..fcb29e78f 100644 --- a/packages/web/server-functions/src/shared.ts +++ b/packages/web/server-functions/src/shared.ts @@ -1195,15 +1195,42 @@ export function createChunk(data: string): Uint8Array; // streams frame their chunks identically — see frame-transport.js) so there // is exactly one framing implementation. export function createChunk(data) { - const encoder = new TextEncoder(); - const encodeData = encoder.encode(data); + const encodeData = CHUNK_ENCODER.encode(data); const bytes = encodeData.length; + // The header carries exactly 8 hex digits. A longer length would widen it + // past 12 bytes and every reader would misframe the rest of the stream. + if (bytes > MAX_CHUNK_BYTES) { + throw new RangeError("Server function chunk is too large to frame."); + } const chunk = new Uint8Array(12 + bytes); - chunk.set(encoder.encode(`;0x${bytes.toString(16).padStart(8, "0")};`)); // 32-bit + chunk.set(CHUNK_ENCODER.encode(`;0x${bytes.toString(16).padStart(8, "0")};`)); // 32-bit chunk.set(encodeData, 12); return chunk; } +const CHUNK_ENCODER = /* @__PURE__ */ new TextEncoder(); +// Fatal, so a payload that is not valid UTF-8 is refused instead of being +// decoded with U+FFFD substituted. `createChunk` always writes valid UTF-8, +// so only a corrupted or hand-built stream can trip it. Stateless use only: +// every call decodes one whole payload, so one instance serves every reader. +const CHUNK_DECODER = /* @__PURE__ */ new TextDecoder("utf-8", { fatal: true }); +const MAX_CHUNK_BYTES = 0xffffffff; +// A store this large is released once the reader holds no unread bytes, +// unless the frames it is serving are big enough to need it again. +const RETAINED_STORE_BYTES = 64 * 1024; + +function malformedStream() { + return new Error("Malformed server function stream."); +} + +/** The value of one ASCII hex digit, or -1. */ +function hexDigit(byte) { + if (byte >= 0x30 && byte <= 0x39) return byte - 0x30; + if (byte >= 0x61 && byte <= 0x66) return byte - 0x57; + if (byte >= 0x41 && byte <= 0x46) return byte - 0x37; + return -1; +} + export class ChunkReader { constructor(stream) { this.reader = stream.getReader(); @@ -1212,6 +1239,9 @@ export class ChunkReader { this.store = new Uint8Array(0); this.buffer = this.store; this.done = false; + // Set by `cancel()`. A cancelled read ends cleanly even when it stops + // partway through a frame. + this.cancelled = false; } async readChunk() { @@ -1253,6 +1283,10 @@ export class ChunkReader { } async next() { + // A cancelled read is over. Frames still buffered are not delivered: + // the caller asked to stop, and frames cancels a superseded response + // precisely so its later chunks are not applied. + if (this.cancelled) return { done: true, value: undefined }; // A network read boundary can land anywhere — inside the 12-byte header // just as easily as inside a payload — so buffer until the whole header // is present before parsing it. Parsing a truncated header used to @@ -1260,41 +1294,82 @@ export class ChunkReader { // frame, which no localhost test ever produces. while (this.buffer.length < 12) { if (this.done) { - if (this.buffer.length === 0) return { done: true, value: undefined }; - throw new Error("Malformed server function stream."); + // A partial frame is truncation, unless `cancel()` is why the body + // ended. Then it is the clean end the caller asked for. + if (this.buffer.length === 0 || this.cancelled) return { done: true, value: undefined }; + throw malformedStream(); } await this.readChunk(); } - // `;0x00000000;` — the hex length names how many payload bytes to wait for - const decoder = new TextDecoder(); - const bytes = Number.parseInt(decoder.decode(this.buffer.subarray(1, 11)), 16); - if (Number.isNaN(bytes)) { - throw new Error("Malformed server function stream."); + // `;0x00000000;`, exactly: the delimiters, the `0x`, then 8 hex digits + // naming how many payload bytes to wait for. `parseInt` used to accept + // any 10 bytes it could read a number from, so a header with the wrong + // delimiters or trailing junk (`;0x5zzzzzzz;`) decoded as a real frame. + const header = this.buffer; + if (header[0] !== 0x3b || header[1] !== 0x30 || header[2] !== 0x78 || header[11] !== 0x3b) { + throw malformedStream(); + } + let bytes = 0; + for (let i = 3; i < 11; i++) { + const digit = hexDigit(header[i]); + if (digit < 0) throw malformedStream(); + bytes = bytes * 16 + digit; } while (bytes > this.buffer.length - 12) { if (this.done) { - throw new Error("Malformed server function stream."); + if (this.cancelled) return { done: true, value: undefined }; + throw malformedStream(); } await this.readChunk(); } - const partial = decoder.decode(this.buffer.subarray(12, 12 + bytes)); + let partial; + try { + partial = CHUNK_DECODER.decode(this.buffer.subarray(12, 12 + bytes)); + } catch { + throw malformedStream(); + } this.buffer = this.buffer.subarray(12 + bytes); + this.releaseStore(bytes); return { done: false, value: partial }; } + /** + * Drops an oversized store once nothing in it is unread. Without this, the + * doubled allocation from one large frame was kept until the stream ended. + * For a live source or a frames connection, that can be as long as the page + * is open. The store is kept while frames of the size just read still need + * it, so a steady stream of large frames does not reallocate per frame. + */ + releaseStore(frameBytes) { + if (this.buffer.length !== 0) return; + const keep = Math.max(RETAINED_STORE_BYTES, (12 + frameBytes) * 4); + if (this.store.length <= keep) return; + this.store = new Uint8Array(0); + this.buffer = this.store; + } + async drain(interpret) { - while (true) { - const result = await this.next(); - if (result.done) { - break; + try { + while (true) { + const result = await this.next(); + if (result.done) { + break; + } + interpret(result.value); } - interpret(result.value); + } catch (error) { + // Nothing will read the rest of the body, so release it instead of + // leaving it locked and unread until it is collected. + this.cancel(error).catch(() => {}); + throw error; } } /** End the read: the body is cancelled through the lock this reader - * holds, and a pending `next()` resolves done. */ + * holds, and a pending `next()` resolves done, even partway through a + * frame. */ cancel(reason) { + this.cancelled = true; return this.reader.cancel(reason); } } @@ -1323,7 +1398,7 @@ export function createEventChunk(data, id) { // end of its line, and the payload's lines were split above it. let event = id !== undefined && !id.includes("\0") ? `id: ${id}\n` : ""; for (const line of data.split(/\r\n|\r|\n/)) event += `data: ${line}\n`; - return new TextEncoder().encode(event + "\n"); + return CHUNK_ENCODER.encode(event + "\n"); } /** @@ -1625,14 +1700,31 @@ export async function deserializeStream(source, codecOptions, wire) { } const reader = wire && isEventStream(source) ? wire.open(source.body) : new ChunkReader(source.body); - const result = await reader.next(); + // When the first frame fails, nothing will read the rest of the body, so + // it is cancelled rather than left locked and unread. Once the drain below + // is running, the same cancel also ends it. + const abandon = error => { + try { + const cancelled = reader.cancel(error); + if (cancelled && typeof cancelled.then === "function") cancelled.then(undefined, () => {}); + } catch {} + }; + let result; + try { + result = await reader.next(); + } catch (error) { + abandon(error); + throw error; + } if (!result.done) { // An error trailer as the FIRST frame is the whole answer: encoding // failed before any value was delivered, and the failure is the result // (#3117). Thrown here so the caller sees a failed call, never a void // success. if (result.value.startsWith(ERROR_TRAILER_PREFIX)) { - throw errorFromTrailer(result.value); + const error = errorFromTrailer(result.value); + abandon(error); + throw error; } // The codec's decode half loads here — when a Serialized body has // actually arrived — so a client whose responses all ride the JSON fast @@ -1688,7 +1780,12 @@ export async function deserializeStream(source, codecOptions, wire) { error => end(error) ); - return interpretChunk(result.value); + try { + return interpretChunk(result.value); + } catch (error) { + abandon(error); + throw error; + } } return undefined; } /** diff --git a/packages/web/test/runtime/chunk-reader-hardening.spec.ts b/packages/web/test/runtime/chunk-reader-hardening.spec.ts new file mode 100644 index 000000000..d567dc04b --- /dev/null +++ b/packages/web/test/runtime/chunk-reader-hardening.spec.ts @@ -0,0 +1,187 @@ +/** + * ChunkReader cancellation, header validation, and cleanup. + * + * - `cancel()` ends the read cleanly even partway through a frame. Frames + * relies on this to report a superseded response as superseded rather + * than as a malformed stream. + * - A header must be exactly `;0x` plus 8 hex digits plus `;`. + * - A failed read releases the body instead of leaving it locked. + * - An oversized store is released once nothing in it is unread. + */ +import { describe, expect, it } from "vitest"; +import { + ChunkReader, + createChunk, + deserializeStream +} from "../../server-functions/src/shared.js"; + +const encoder = new TextEncoder(); + +function concat(chunks: Uint8Array[]) { + const total = chunks.reduce((size, chunk) => size + chunk.length, 0); + const bytes = new Uint8Array(total); + let offset = 0; + for (const chunk of chunks) { + bytes.set(chunk, offset); + offset += chunk.length; + } + return bytes; +} + +/** + * Delivers `pieces` one per read. With `hold`, the stream then stays open + * instead of closing, like a connection waiting on the server. Records + * whether it was cancelled. + */ +function source(pieces: Uint8Array[], { hold = false } = {}) { + let index = 0; + let release: (() => void) | undefined; + const state = { cancelled: false }; + const stream = new ReadableStream( + { + pull(controller) { + if (index < pieces.length) return controller.enqueue(pieces[index++]); + if (hold) return new Promise(resolve => (release = resolve)); + controller.close(); + }, + cancel() { + state.cancelled = true; + release?.(); + } + }, + { highWaterMark: 0 } + ); + return { stream, state }; +} + +const tick = () => new Promise(resolve => setTimeout(resolve, 0)); + +/** The reader's backing allocation. Internal, so it is not on the declared type. */ +const storeOf = (reader: InstanceType): Uint8Array => + (reader as unknown as { store: Uint8Array }).store; + +describe("ChunkReader cancel()", () => { + const frame = createChunk("hello"); + + for (const [label, buffered] of [ + ["nothing buffered", new Uint8Array(0)], + ["partway through the header", frame.subarray(0, 5)], + ["partway through the payload", frame.subarray(0, 14)] + ] as const) { + it(`resolves a pending next() as done with ${label}`, async () => { + const { stream } = source(buffered.length ? [buffered] : [], { hold: true }); + const reader = new ChunkReader(stream); + const pending = reader.next(); + await tick(); + await reader.cancel(new Error("superseded")); + await expect(pending).resolves.toEqual({ done: true, value: undefined }); + }); + } + + it("does not deliver frames that were already buffered", async () => { + const { stream } = source([concat([createChunk("a"), createChunk("b")])], { hold: true }); + const reader = new ChunkReader(stream); + expect((await reader.next()).value).toBe("a"); + await reader.cancel(undefined); + expect(await reader.next()).toEqual({ done: true, value: undefined }); + }); + + it("still refuses a body that ends partway through a frame without a cancel", async () => { + const { stream } = source([frame.subarray(0, 14)]); + await expect(new ChunkReader(stream).next()).rejects.toThrow( + "Malformed server function stream." + ); + }); +}); + +describe("ChunkReader header validation", () => { + const payload = encoder.encode("hello"); + + it("accepts the canonical header, in either hex case", async () => { + for (const header of [";0x00000005;", ";0x0000000A;", ";0x0000000a;"]) { + const body = header === ";0x00000005;" ? payload : encoder.encode("helloworld"); + const reader = new ChunkReader(source([concat([encoder.encode(header), body])]).stream); + expect((await reader.next()).done).toBe(false); + } + }); + + for (const header of [ + ";0x5zzzzzzz;", // trailing junk after a valid digit + "X0x00000005X", // wrong delimiters + ";0000000005;", // no 0x prefix + ";0x0000005 ;", // whitespace inside the digits + ";0X00000005;", // uppercase X + ";-0x0000005;" // a sign + ]) { + it(`refuses ${JSON.stringify(header)}`, async () => { + const reader = new ChunkReader(source([concat([encoder.encode(header), payload])]).stream); + await expect(reader.next()).rejects.toThrow("Malformed server function stream."); + }); + } + + it("refuses a payload that is not valid UTF-8", async () => { + const invalid = Uint8Array.of(0x22, 0xff, 0xfe, 0x22); + const header = encoder.encode(";0x00000004;"); + const reader = new ChunkReader(source([concat([header, invalid])]).stream); + await expect(reader.next()).rejects.toThrow("Malformed server function stream."); + }); + + it("still decodes multi-byte characters split across reads", async () => { + const value = '{"text":"héllo 😀 世界"}'; + const bytes = createChunk(value); + for (let cut = 0; cut <= bytes.length; cut++) { + const reader = new ChunkReader(source([bytes.subarray(0, cut), bytes.subarray(cut)]).stream); + expect((await reader.next()).value).toBe(value); + } + }); +}); + +describe("ChunkReader cleanup", () => { + it("cancels the body when the drain fails", async () => { + const bytes = concat([createChunk("1"), createChunk("not json"), createChunk("3")]); + const { stream, state } = source([bytes], { hold: true }); + await expect( + new ChunkReader(stream).drain((value: string) => JSON.parse(value)) + ).rejects.toThrow(SyntaxError); + expect(state.cancelled).toBe(true); + }); + + it("cancels the body when deserializeStream's first frame is malformed", async () => { + const { stream, state } = source([encoder.encode(";0x5zzzzzzz;hello")], { hold: true }); + await expect(deserializeStream(new Response(stream))).rejects.toThrow( + "Malformed server function stream." + ); + expect(state.cancelled).toBe(true); + }); + + it("cancels the body when deserializeStream's first value fails to decode", async () => { + const { stream, state } = source([createChunk("not json")], { hold: true }); + await expect(deserializeStream(new Response(stream))).rejects.toThrow(SyntaxError); + expect(state.cancelled).toBe(true); + }); + + it("releases an oversized store once nothing in it is unread", async () => { + const large = createChunk("x".repeat(1 << 20)); + const pieces: Uint8Array[] = []; + for (let offset = 0; offset < large.length; offset += 65_536) { + pieces.push(large.subarray(offset, offset + 65_536)); + } + pieces.push(createChunk("small")); + const reader = new ChunkReader(source(pieces, { hold: true }).stream); + await reader.next(); + await reader.next(); + expect(storeOf(reader).length).toBeLessThanOrEqual(64 * 1024); + reader.cancel(undefined); + }); + + it("keeps the store while frames of the same size keep arriving", async () => { + const frame = () => createChunk("y".repeat(200_000)); + const reader = new ChunkReader(source([frame(), frame(), frame()], { hold: true }).stream); + await reader.next(); + const store = storeOf(reader); + await reader.next(); + await reader.next(); + expect(storeOf(reader)).toBe(store); + reader.cancel(undefined); + }); +}); From 2b84aae3ce4bfa2aa8d5317bfede531f0f6b220c Mon Sep 17 00:00:00 2001 From: "Alexis H. Munsayac" Date: Wed, 7 Oct 2026 08:49:44 +0800 Subject: [PATCH 2/3] fix(web): release a large chunk's store right away and recheck cancel after each read - `releaseStore` now shrinks a store over 64 KiB to its unread bytes after every frame. The old rule compared the store with four times the frame just read. A store that grew to fit one frame is at most about twice that frame, so it was never released right after it. A connection that went idle after a large frame kept the allocation until the stream ended. - A steady stream of frames over 64 KiB now regrows its store for each one. Frames under 64 KiB never shrink it, so the #3154 steady state for small frames is unchanged. - `next()` checks `cancelled` after every read, not only on entry. A read that had already resolved with data when `cancel()` ran used to finish its frame and deliver it. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/web/server-functions/src/shared.ts | 50 +++++++++------ .../runtime/chunk-reader-hardening.spec.ts | 62 ++++++++++++++----- 2 files changed, 79 insertions(+), 33 deletions(-) diff --git a/packages/web/server-functions/src/shared.ts b/packages/web/server-functions/src/shared.ts index fcb29e78f..d71ced29a 100644 --- a/packages/web/server-functions/src/shared.ts +++ b/packages/web/server-functions/src/shared.ts @@ -1215,8 +1215,8 @@ const CHUNK_ENCODER = /* @__PURE__ */ new TextEncoder(); // every call decodes one whole payload, so one instance serves every reader. const CHUNK_DECODER = /* @__PURE__ */ new TextDecoder("utf-8", { fatal: true }); const MAX_CHUNK_BYTES = 0xffffffff; -// A store this large is released once the reader holds no unread bytes, -// unless the frames it is serving are big enough to need it again. +// A store over this size is shrunk to its unread bytes after each frame (see +// `releaseStore`). A store at or under it is kept as is. const RETAINED_STORE_BYTES = 64 * 1024; function malformedStream() { @@ -1285,7 +1285,9 @@ export class ChunkReader { async next() { // A cancelled read is over. Frames still buffered are not delivered: // the caller asked to stop, and frames cancels a superseded response - // precisely so its later chunks are not applied. + // precisely so its later chunks are not applied. The check repeats after + // every read below, because `cancel()` can run while a read is pending, + // or after one has already resolved with data. if (this.cancelled) return { done: true, value: undefined }; // A network read boundary can land anywhere — inside the 12-byte header // just as easily as inside a payload — so buffer until the whole header @@ -1294,12 +1296,11 @@ export class ChunkReader { // frame, which no localhost test ever produces. while (this.buffer.length < 12) { if (this.done) { - // A partial frame is truncation, unless `cancel()` is why the body - // ended. Then it is the clean end the caller asked for. - if (this.buffer.length === 0 || this.cancelled) return { done: true, value: undefined }; + if (this.buffer.length === 0) return { done: true, value: undefined }; throw malformedStream(); } await this.readChunk(); + if (this.cancelled) return { done: true, value: undefined }; } // `;0x00000000;`, exactly: the delimiters, the `0x`, then 8 hex digits // naming how many payload bytes to wait for. `parseInt` used to accept @@ -1317,10 +1318,10 @@ export class ChunkReader { } while (bytes > this.buffer.length - 12) { if (this.done) { - if (this.cancelled) return { done: true, value: undefined }; throw malformedStream(); } await this.readChunk(); + if (this.cancelled) return { done: true, value: undefined }; } let partial; try { @@ -1329,23 +1330,34 @@ export class ChunkReader { throw malformedStream(); } this.buffer = this.buffer.subarray(12 + bytes); - this.releaseStore(bytes); + this.releaseStore(); return { done: false, value: partial }; } /** - * Drops an oversized store once nothing in it is unread. Without this, the - * doubled allocation from one large frame was kept until the stream ended. - * For a live source or a frames connection, that can be as long as the page - * is open. The store is kept while frames of the size just read still need - * it, so a steady stream of large frames does not reallocate per frame. + * Shrinks an oversized store to the bytes still unread, after each frame. + * Without this, the doubled allocation from one large frame was kept until + * the stream ended. For a live source or a frames connection, that can be + * as long as the page is open. + * + * The rule looks only at what is left to read, so a connection that goes + * idle right after a large frame releases it at once. The cost is that a + * steady stream of frames over 64 KiB regrows its store for each one. That + * is a few doubling allocations per frame, small next to decoding the + * frame. Frames under 64 KiB never trigger it, so the #3154 steady state + * for small frames is unchanged. */ - releaseStore(frameBytes) { - if (this.buffer.length !== 0) return; - const keep = Math.max(RETAINED_STORE_BYTES, (12 + frameBytes) * 4); - if (this.store.length <= keep) return; - this.store = new Uint8Array(0); - this.buffer = this.store; + releaseStore() { + if (this.store.length <= RETAINED_STORE_BYTES) return; + const unread = this.buffer.length; + // Bytes of the next frame that arrived with this one. Copying up to 64 + // KiB of them into a fresh store is still far cheaper than keeping the + // large one. + if (unread > RETAINED_STORE_BYTES) return; + const kept = new Uint8Array(unread); + kept.set(this.buffer); + this.store = kept; + this.buffer = kept; } async drain(interpret) { diff --git a/packages/web/test/runtime/chunk-reader-hardening.spec.ts b/packages/web/test/runtime/chunk-reader-hardening.spec.ts index d567dc04b..b92944d7a 100644 --- a/packages/web/test/runtime/chunk-reader-hardening.spec.ts +++ b/packages/web/test/runtime/chunk-reader-hardening.spec.ts @@ -86,6 +86,23 @@ describe("ChunkReader cancel()", () => { expect(await reader.next()).toEqual({ done: true, value: undefined }); }); + it("does not deliver a frame whose last read resolved just before the cancel", async () => { + // The pending read resolves with the rest of the frame, and `cancel()` + // runs before `next()` resumes. The frame must still not be delivered. + let controller!: ReadableStreamDefaultController; + const stream = new ReadableStream( + { start: c => void (controller = c) }, + { highWaterMark: 0 } + ); + const reader = new ChunkReader(stream); + controller.enqueue(frame.subarray(0, 5)); + const pending = reader.next(); + await tick(); + controller.enqueue(frame.subarray(5)); + reader.cancel(undefined); + await expect(pending).resolves.toEqual({ done: true, value: undefined }); + }); + it("still refuses a body that ends partway through a frame without a cancel", async () => { const { stream } = source([frame.subarray(0, 14)]); await expect(new ChunkReader(stream).next()).rejects.toThrow( @@ -160,27 +177,44 @@ describe("ChunkReader cleanup", () => { expect(state.cancelled).toBe(true); }); - it("releases an oversized store once nothing in it is unread", async () => { - const large = createChunk("x".repeat(1 << 20)); + /** Splits `bytes` into 64 KiB reads. */ + const reads = (bytes: Uint8Array) => { const pieces: Uint8Array[] = []; - for (let offset = 0; offset < large.length; offset += 65_536) { - pieces.push(large.subarray(offset, offset + 65_536)); + for (let offset = 0; offset < bytes.length; offset += 65_536) { + pieces.push(bytes.subarray(offset, offset + 65_536)); } - pieces.push(createChunk("small")); - const reader = new ChunkReader(source(pieces, { hold: true }).stream); - await reader.next(); - await reader.next(); - expect(storeOf(reader).length).toBeLessThanOrEqual(64 * 1024); + return pieces; + }; + + it("releases a large frame's store as soon as the stream goes idle", async () => { + // No frame follows: the connection waits, as a live source does. + const reader = new ChunkReader(source(reads(createChunk("x".repeat(8 << 20))), { hold: true }).stream); + expect((await reader.next()).value).toHaveLength(8 << 20); + expect(storeOf(reader).length).toBe(0); reader.cancel(undefined); }); - it("keeps the store while frames of the same size keep arriving", async () => { - const frame = () => createChunk("y".repeat(200_000)); - const reader = new ChunkReader(source([frame(), frame(), frame()], { hold: true }).stream); + it("keeps the next frame's bytes when it shrinks the store", async () => { + // The last read carries the end of the large frame and the start of the + // next one, so the shrink must carry those bytes over. + const large = createChunk("x".repeat(1 << 20)); + const next = createChunk("after"); + const bytes = concat([large, next.subarray(0, 8)]); + const reader = new ChunkReader( + source([...reads(bytes), next.subarray(8)], { hold: true }).stream + ); + expect((await reader.next()).value).toHaveLength(1 << 20); + expect(storeOf(reader).length).toBe(8); + expect((await reader.next()).value).toBe("after"); + reader.cancel(undefined); + }); + + it("keeps a store at or under 64 KiB across frames", async () => { + const frames = Array.from({ length: 20 }, (_, i) => createChunk(`frame-${i}-` + "y".repeat(300))); + const reader = new ChunkReader(source([concat(frames)], { hold: true }).stream); await reader.next(); const store = storeOf(reader); - await reader.next(); - await reader.next(); + for (let i = 1; i < frames.length; i++) await reader.next(); expect(storeOf(reader)).toBe(store); reader.cancel(undefined); }); From f6fecadfa147b1dcf44eae7c35ad64a7f65a0e2c Mon Sep 17 00:00:00 2001 From: Ryan Carniato Date: Wed, 7 Oct 2026 01:41:47 -0700 Subject: [PATCH 3/3] fix(web): trim the hardened chunk reader; raise the frames/page size caps (#3846) No behavior change. Header check is one regex plus parseInt (hexDigit and the byte compares go), deserializeStream releases the body from one try/catch, the store release is inlined into next(), and readChunk answers whether the reader was cancelled so next() returns one shared done result. Frames client +849 -> +518 B minified against next. Size-Exception for the remaining cost, accepted by the maintainer 2026-10-07 (server-function/frames reader correctness): frames eager client consumer 13.79 -> 13.98 KB, page base 44.89 -> 45.12 KB, page live 48.60 -> 48.82 KB, each at measured + 10 B with its minified recorded. Co-authored-by: Cursor --- packages/web/server-functions/src/shared.ts | 224 ++++++++------------ scripts/size/floor-caps.json | 8 +- scripts/size/scenarios.js | 32 ++- 3 files changed, 127 insertions(+), 137 deletions(-) diff --git a/packages/web/server-functions/src/shared.ts b/packages/web/server-functions/src/shared.ts index d71ced29a..7b7eecf48 100644 --- a/packages/web/server-functions/src/shared.ts +++ b/packages/web/server-functions/src/shared.ts @@ -1216,21 +1216,15 @@ const CHUNK_ENCODER = /* @__PURE__ */ new TextEncoder(); const CHUNK_DECODER = /* @__PURE__ */ new TextDecoder("utf-8", { fatal: true }); const MAX_CHUNK_BYTES = 0xffffffff; // A store over this size is shrunk to its unread bytes after each frame (see -// `releaseStore`). A store at or under it is kept as is. +// the end of `ChunkReader.next`). A store at or under it is kept as is. const RETAINED_STORE_BYTES = 64 * 1024; +const DONE = { done: true, value: undefined }; + function malformedStream() { return new Error("Malformed server function stream."); } -/** The value of one ASCII hex digit, or -1. */ -function hexDigit(byte) { - if (byte >= 0x30 && byte <= 0x39) return byte - 0x30; - if (byte >= 0x61 && byte <= 0x66) return byte - 0x57; - if (byte >= 0x41 && byte <= 0x46) return byte - 0x37; - return -1; -} - export class ChunkReader { constructor(stream) { this.reader = stream.getReader(); @@ -1244,11 +1238,12 @@ export class ChunkReader { this.cancelled = false; } + /** Reads once into the buffer; answers whether the reader was cancelled. */ async readChunk() { const chunk = await this.reader.read(); if (chunk.done) { this.done = true; - return; + return this.cancelled; } // Amortized growth (#3154). Reallocating the whole buffer per network // read made one frame O(reads²): a 1 MiB argument payload delivered at @@ -1280,6 +1275,7 @@ export class ChunkReader { this.store = grown; this.buffer = grown.subarray(0, needed); } + return this.cancelled; } async next() { @@ -1288,7 +1284,7 @@ export class ChunkReader { // precisely so its later chunks are not applied. The check repeats after // every read below, because `cancel()` can run while a read is pending, // or after one has already resolved with data. - if (this.cancelled) return { done: true, value: undefined }; + if (this.cancelled) return DONE; // A network read boundary can land anywhere — inside the 12-byte header // just as easily as inside a payload — so buffer until the whole header // is present before parsing it. Parsing a truncated header used to @@ -1296,32 +1292,23 @@ export class ChunkReader { // frame, which no localhost test ever produces. while (this.buffer.length < 12) { if (this.done) { - if (this.buffer.length === 0) return { done: true, value: undefined }; + if (this.buffer.length === 0) return DONE; throw malformedStream(); } - await this.readChunk(); - if (this.cancelled) return { done: true, value: undefined }; + if (await this.readChunk()) return DONE; } // `;0x00000000;`, exactly: the delimiters, the `0x`, then 8 hex digits // naming how many payload bytes to wait for. `parseInt` used to accept // any 10 bytes it could read a number from, so a header with the wrong // delimiters or trailing junk (`;0x5zzzzzzz;`) decoded as a real frame. - const header = this.buffer; - if (header[0] !== 0x3b || header[1] !== 0x30 || header[2] !== 0x78 || header[11] !== 0x3b) { - throw malformedStream(); - } - let bytes = 0; - for (let i = 3; i < 11; i++) { - const digit = hexDigit(header[i]); - if (digit < 0) throw malformedStream(); - bytes = bytes * 16 + digit; - } + const header = String.fromCharCode(...this.buffer.subarray(0, 12)); + if (!/^;0x[\dA-Fa-f]{8};$/.test(header)) throw malformedStream(); + const bytes = parseInt(header.slice(3, 11), 16); while (bytes > this.buffer.length - 12) { if (this.done) { throw malformedStream(); } - await this.readChunk(); - if (this.cancelled) return { done: true, value: undefined }; + if (await this.readChunk()) return DONE; } let partial; try { @@ -1330,36 +1317,25 @@ export class ChunkReader { throw malformedStream(); } this.buffer = this.buffer.subarray(12 + bytes); - this.releaseStore(); + // Shrink an oversized store to the bytes still unread, after each frame. + // Without this, the doubled allocation from one large frame was kept + // until the stream ended. For a live source or a frames connection, that + // can be as long as the page is open. + // + // The rule looks only at what is left to read, so a connection that goes + // idle right after a large frame releases it at once. The cost is that a + // steady stream of frames over 64 KiB regrows its store for each one. + // That is a few doubling allocations per frame, small next to decoding + // the frame. Frames under 64 KiB never trigger it, so the #3154 steady + // state for small frames is unchanged. Unread bytes (the next frame's, + // arriving with this one) are copied over only up to 64 KiB, which is + // still far cheaper than keeping the large store. + if (this.store.length > RETAINED_STORE_BYTES && this.buffer.length <= RETAINED_STORE_BYTES) { + this.store = this.buffer = this.buffer.slice(); + } return { done: false, value: partial }; } - /** - * Shrinks an oversized store to the bytes still unread, after each frame. - * Without this, the doubled allocation from one large frame was kept until - * the stream ended. For a live source or a frames connection, that can be - * as long as the page is open. - * - * The rule looks only at what is left to read, so a connection that goes - * idle right after a large frame releases it at once. The cost is that a - * steady stream of frames over 64 KiB regrows its store for each one. That - * is a few doubling allocations per frame, small next to decoding the - * frame. Frames under 64 KiB never trigger it, so the #3154 steady state - * for small frames is unchanged. - */ - releaseStore() { - if (this.store.length <= RETAINED_STORE_BYTES) return; - const unread = this.buffer.length; - // Bytes of the next frame that arrived with this one. Copying up to 64 - // KiB of them into a fresh store is still far cheaper than keeping the - // large one. - if (unread > RETAINED_STORE_BYTES) return; - const kept = new Uint8Array(unread); - kept.set(this.buffer); - this.store = kept; - this.buffer = kept; - } - async drain(interpret) { try { while (true) { @@ -1715,91 +1691,77 @@ export async function deserializeStream(source, codecOptions, wire) { // When the first frame fails, nothing will read the rest of the body, so // it is cancelled rather than left locked and unread. Once the drain below // is running, the same cancel also ends it. - const abandon = error => { - try { - const cancelled = reader.cancel(error); - if (cancelled && typeof cancelled.then === "function") cancelled.then(undefined, () => {}); - } catch {} - }; - let result; try { - result = await reader.next(); - } catch (error) { - abandon(error); - throw error; - } - if (!result.done) { - // An error trailer as the FIRST frame is the whole answer: encoding - // failed before any value was delivered, and the failure is the result - // (#3117). Thrown here so the caller sees a failed call, never a void - // success. - if (result.value.startsWith(ERROR_TRAILER_PREFIX)) { - const error = errorFromTrailer(result.value); - abandon(error); - throw error; - } - // The codec's decode half loads here — when a Serialized body has - // actually arrived — so a client whose responses all ride the JSON fast - // path never pays for it (see the loading notes at the top). The - // decode-only module: reading a payload never needs the encode half - // (that loads separately, when rich arguments serialize). - const { createJSONDeserializer } = await import("../../serialization/src/serializer-decode.js"); - // Cross-references between chunks resolve through state inside the - // deserializer, so one instance handles the whole stream. - const deserializeChunk = createJSONDeserializer(codecOptions); - - function interpretChunk(chunk) { - // Mid-stream, the trailer means a LATER value's encoding failed: the - // resolved head keeps its data, and the throw below rejects the drain, - // whose abort sweep fails every still-pending async value with the - // carried reason instead of a generic truncation error (#3117). - if (chunk.startsWith(ERROR_TRAILER_PREFIX)) { - throw errorFromTrailer(chunk); + const result = await reader.next(); + if (!result.done) { + // An error trailer as the FIRST frame is the whole answer: encoding + // failed before any value was delivered, and the failure is the result + // (#3117). Thrown here so the caller sees a failed call, never a void + // success. + if (result.value.startsWith(ERROR_TRAILER_PREFIX)) { + throw errorFromTrailer(result.value); + } + // The codec's decode half loads here — when a Serialized body has + // actually arrived — so a client whose responses all ride the JSON fast + // path never pays for it (see the loading notes at the top). The + // decode-only module: reading a payload never needs the encode half + // (that loads separately, when rich arguments serialize). + const { createJSONDeserializer } = + await import("../../serialization/src/serializer-decode.js"); + // Cross-references between chunks resolve through state inside the + // deserializer, so one instance handles the whole stream. + const deserializeChunk = createJSONDeserializer(codecOptions); + + function interpretChunk(chunk) { + // Mid-stream, the trailer means a LATER value's encoding failed: the + // resolved head keeps its data, and the throw below rejects the drain, + // whose abort sweep fails every still-pending async value with the + // carried reason instead of a generic truncation error (#3117). + if (chunk.startsWith(ERROR_TRAILER_PREFIX)) { + throw errorFromTrailer(chunk); + } + return deserializeChunk(JSON.parse(chunk)); } - return deserializeChunk(JSON.parse(chunk)); - } - // Failure wiring for the drain: a network drop or malformed frame must - // fail every value still waiting on later chunks — otherwise their - // promises hang forever and open streams never terminate (and the drain - // rejection itself goes unhandled). Normal completion runs the same - // sweep: on a well-formed stream every value has already settled and the - // sweep no-ops, while a truncation that lands exactly on a frame - // boundary — indistinguishable from completion — leaves stranded values - // that can never settle once the body is done. - // - // A live loop (`wire`) is told instead of swept: the body's end is the - // connection's lifetime signal — a death when deferreds are still open, - // a completion when none are — and the loop decides what happens to the - // open ones. A death it will reconnect from leaves them pending (the - // re-yielded answer supersedes them; throwing into them would surface - // the death the loop exists to erase). An iteration ending for good - // settles them by how it ended: `sweep` fails them (ended by error), - // `close` completes the streams and leaves promises pending (ended by - // the consumer or by completion) — see the loop's emitClosed. - const connection = wire && wire.connection; - let endConnection; - if (connection) connection.ended = new Promise(resolve => (endConnection = resolve)); - const end = error => { - const sweep = () => deserializeChunk.abort(error); - if (connection) { - const close = () => deserializeChunk.close(); - endConnection({ open: deserializeChunk.open(), error, sweep, close }); - } else sweep(); - }; - reader.drain(interpretChunk).then( - () => end(new Error("Server function stream ended unexpectedly.")), - error => end(error) - ); + // Failure wiring for the drain: a network drop or malformed frame must + // fail every value still waiting on later chunks — otherwise their + // promises hang forever and open streams never terminate (and the drain + // rejection itself goes unhandled). Normal completion runs the same + // sweep: on a well-formed stream every value has already settled and the + // sweep no-ops, while a truncation that lands exactly on a frame + // boundary — indistinguishable from completion — leaves stranded values + // that can never settle once the body is done. + // + // A live loop (`wire`) is told instead of swept: the body's end is the + // connection's lifetime signal — a death when deferreds are still open, + // a completion when none are — and the loop decides what happens to the + // open ones. A death it will reconnect from leaves them pending (the + // re-yielded answer supersedes them; throwing into them would surface + // the death the loop exists to erase). An iteration ending for good + // settles them by how it ended: `sweep` fails them (ended by error), + // `close` completes the streams and leaves promises pending (ended by + // the consumer or by completion) — see the loop's emitClosed. + const connection = wire && wire.connection; + let endConnection; + if (connection) connection.ended = new Promise(resolve => (endConnection = resolve)); + const end = error => { + const sweep = () => deserializeChunk.abort(error); + if (connection) { + const close = () => deserializeChunk.close(); + endConnection({ open: deserializeChunk.open(), error, sweep, close }); + } else sweep(); + }; + reader.drain(interpretChunk).then( + () => end(new Error("Server function stream ended unexpectedly.")), + error => end(error) + ); - try { return interpretChunk(result.value); - } catch (error) { - abandon(error); - throw error; } + } catch (error) { + reader.cancel(error).catch(() => {}); + throw error; } - return undefined; } /** * `deserializeStream` for an already-buffered string. * diff --git a/scripts/size/floor-caps.json b/scripts/size/floor-caps.json index f2a19921e..f195a5ac3 100644 --- a/scripts/size/floor-caps.json +++ b/scripts/size/floor-caps.json @@ -12,12 +12,12 @@ "minified": 52626 }, "page: base server components (hydrating + dynamic + frames + sf reference)": { - "cap": "44.89 KB", - "minified": 145757 + "cap": "45.12 KB", + "minified": 146295 }, "page: live server components (base + live/GET + action + isPending/latest)": { - "cap": "48.60 KB", - "minified": 157720 + "cap": "48.82 KB", + "minified": 158257 }, "server: floor (getRequestEvent + isServer)": { "cap": "1.34 KB", diff --git a/scripts/size/scenarios.js b/scripts/size/scenarios.js index 69cbbb3c2..37987d0bb 100644 --- a/scripts/size/scenarios.js +++ b/scripts/size/scenarios.js @@ -3382,8 +3382,22 @@ module.exports = [ // 10 B; recorded minified 43,414 B. Accepted by the maintainer // (2026-10-06, "pay the cost for correctness"). The cap is frozen again // at 13.79 KB. - limit: "13.79 KB", - capMinified: 43414, + // Size-Exception (#3846, 2026-10-07): 13.79 -> 13.98 KB, measured at + // 13,969 B against `next` @ 721eb0676's 13,787 (+182 B; 179 B over the + // cap; +518 B minified, 43,414 -> 43,932; sf shared slice +597, frames + // client -79). The server-function/frames `ChunkReader` hardening: + // cancel ends a read mid-frame and is rechecked after every read, strict + // `;0x` + 8 hex digit headers, a fatal UTF-8 decoder, the body cancelled + // on a failed drain or first frame, the store released after a large + // frame, and `createChunk`'s 4 GiB `RangeError`. Measured after a shave + // (regex header check, one `try/catch` in `deserializeStream`, inlined + // store release, shared done result) that took the PR from +849 B + // minified to +518. Cap set at measured + 10 B rounded up to 0.01 KB; + // recorded minified 43,932 B. Accepted by the maintainer (2026-10-07: + // server-function/frames reader correctness). The cap is frozen again at + // 13.98 KB. + limit: "13.98 KB", + capMinified: 43932, alias: framesAlias, external: framesExternal }, @@ -3555,6 +3569,13 @@ module.exports = [ // at or below measured + 10 B; recorded minified 145,757 B. Accepted by // the maintainer (2026-10-06, "pay the cost for correctness"). The cap is // frozen again at 44.89 KB. + // Size-Exception (#3846, 2026-10-07): 44.89 -> 45.12 KB + // (floor-caps.json), measured at 45,102 B against `next` @ 721eb0676's + // 44,968 (+134 B; 212 B over the cap; +519 B minified, 145,757 recorded + // -> 146,295) — the `ChunkReader` hardening (the frames note). Cap set + // at measured + 10 B rounded up to 0.01 KB; recorded minified 146,295 B. + // Accepted by the maintainer (2026-10-07: server-function/frames reader + // correctness). The cap is frozen again at 45.12 KB. limit: floorCaps["page: base server components (hydrating + dynamic + frames + sf reference)"], capMinified: floorMinified["page: base server components (hydrating + dynamic + frames + sf reference)"], @@ -3676,6 +3697,13 @@ module.exports = [ // at or below measured + 10 B; recorded minified 157,720 B. Accepted by // the maintainer (2026-10-06, "pay the cost for correctness"). The cap is // frozen again at 48.60 KB. + // Size-Exception (#3846, 2026-10-07): 48.60 -> 48.82 KB + // (floor-caps.json), measured at 48,803 B against `next` @ 721eb0676's + // 48,573 (+230 B; 203 B over the cap; +519 B minified, 157,720 recorded + // -> 158,257) — the same bytes as the base page. Cap set at measured + + // 10 B rounded up to 0.01 KB; recorded minified 158,257 B. Accepted by + // the maintainer (2026-10-07: server-function/frames reader + // correctness). The cap is frozen again at 48.82 KB. limit: floorCaps["page: live server components (base + live/GET + action + isPending/latest)"], capMinified: floorMinified["page: live server components (base + live/GET + action + isPending/latest)"],