Skip to content
Merged
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
9 changes: 9 additions & 0 deletions .changeset/chunk-reader-hardening.md
Original file line number Diff line number Diff line change
@@ -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.
241 changes: 156 additions & 85 deletions packages/web/server-functions/src/shared.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1195,15 +1195,36 @@ 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 over this size is shrunk to its unread bytes after each frame (see
// 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.");
}

export class ChunkReader {
constructor(stream) {
this.reader = stream.getReader();
Expand All @@ -1212,13 +1233,17 @@ 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;
}

/** 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
Expand Down Expand Up @@ -1250,51 +1275,89 @@ export class ChunkReader {
this.store = grown;
this.buffer = grown.subarray(0, needed);
}
return this.cancelled;
}

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. 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;
// 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
// mis-read the length (or throw) whenever a proxy/TLS record split a
// 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.");
if (this.buffer.length === 0) return DONE;
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.");
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 = 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 new Error("Malformed server function stream.");
throw malformedStream();
}
await this.readChunk();
if (await this.readChunk()) return DONE;
}
let partial;
try {
partial = CHUNK_DECODER.decode(this.buffer.subarray(12, 12 + bytes));
} catch {
throw malformedStream();
}
const partial = decoder.decode(this.buffer.subarray(12, 12 + bytes));
this.buffer = this.buffer.subarray(12 + bytes);
// 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 };
}

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);
}
}
Expand Down Expand Up @@ -1323,7 +1386,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");
}

/**
Expand Down Expand Up @@ -1625,72 +1688,80 @@ 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();
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);
// 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.
try {
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)
);

return interpretChunk(result.value);
return interpretChunk(result.value);
}
} catch (error) {
reader.cancel(error).catch(() => {});
throw error;
}
return undefined;
} /**
* `deserializeStream` for an already-buffered string.
*
Expand Down
Loading
Loading