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
10 changes: 10 additions & 0 deletions .changeset/seroval-chunk-reader.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
---
"@solidjs/start": patch
---

Make reading server function streams safer and faster.

- A malformed chunk after the first one is now logged instead of causing an unhandled promise rejection. On hosts that do not catch unhandled rejections, one bad request could stop the server process.
- The stream is cancelled when a chunk cannot be read or parsed.
- Chunk headers are now checked strictly, and a server rejects chunks over 64MB that a client sends.
- Large payloads that arrive in many small pieces are read in linear time. A 16MB chunk read in 16KB pieces took about 2 seconds and now takes about 30ms.
207 changes: 191 additions & 16 deletions packages/start/src/fns/serialization.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -198,11 +198,36 @@ describe("custom seroval plugins", () => {
});
});

function streamOf(pieces: (string | Uint8Array)[]) {
const encoder = new TextEncoder();
return new ReadableStream<Uint8Array>({
start(controller) {
for (const piece of pieces) {
controller.enqueue(typeof piece === "string" ? encoder.encode(piece) : piece);
}
controller.close();
},
});
}

function frame(data: string) {
const size = new TextEncoder().encode(data).length;
return `;0x${size.toString(16).padStart(8, "0")};${data}`;
}

/** Splits bytes into pieces of `size`, the way a network read can. */
function split(text: string, size: number) {
const bytes = new TextEncoder().encode(text);
const pieces: Uint8Array[] = [];
for (let i = 0; i < bytes.length; i += size) {
pieces.push(bytes.subarray(i, i + size));
}
return pieces;
}

/** Frames one serialized node the way `serializeToJSONStream` does. */
function frame(node: unknown) {
const data = JSON.stringify(node);
const size = new TextEncoder().encode(data).length.toString(16).padStart(8, "0");
return `;0x${size};${data}`;
function frameNode(node: unknown) {
return frame(JSON.stringify(node));
}

/** An argument list whose only item is a promise that a later frame was meant to settle. */
Expand Down Expand Up @@ -238,13 +263,30 @@ async function collectUnhandledRejections(run: () => Promise<void>) {
return reasons;
}

/** Reads one whole frame off a serialized stream; its header and data can arrive as separate pieces. */
async function readFirstFrame(stream: ReadableStream<Uint8Array>) {
const reader = stream.getReader();
let bytes = new Uint8Array(0);
let end = Infinity;
while (bytes.length < end) {
const { done, value } = await reader.read();
if (done) break;
const joined = new Uint8Array(bytes.length + value.length);
joined.set(bytes);
joined.set(value, bytes.length);
bytes = joined;
if (end === Infinity && bytes.length >= 12) {
end = 12 + Number.parseInt(new TextDecoder().decode(bytes.subarray(3, 11)), 16);
}
}
await reader.cancel();
return new TextDecoder().decode(bytes.subarray(0, end));
}

/** The first frame of a JSON stream that never completes on its own. */
async function firstFrame(value: unknown) {
const { serializeToJSONStream } = await loadSerialization(true);
const reader = serializeToJSONStream(value).getReader();
const { value: chunk } = await reader.read();
await reader.cancel();
return new TextDecoder().decode(chunk);
return await readFirstFrame(serializeToJSONStream(value));
}

describe("values waiting on a later frame", () => {
Expand All @@ -254,12 +296,13 @@ describe("values waiting on a later frame", () => {

afterEach(() => {
vi.unstubAllEnvs();
vi.restoreAllMocks();
});

it("rejects a promise that is still pending when the body ends", async () => {
const { deserializeJSONStream } = await loadSerialization(true);

const [arg] = (await deserializeJSONStream(new Response(frame(PENDING_PROMISE_ARGS)))) as [
const [arg] = (await deserializeJSONStream(new Response(frameNode(PENDING_PROMISE_ARGS)))) as [
Promise<unknown>,
];

Expand All @@ -283,25 +326,30 @@ describe("values waiting on a later frame", () => {

it("rejects pending values with the failure when a later frame is malformed", async () => {
const { deserializeJSONStream } = await loadSerialization(true);
const consoleError = vi.spyOn(console, "error").mockImplementation(() => {});
let outcome: Awaited<ReturnType<typeof settleWithin>> | undefined;

const unhandled = await collectUnhandledRejections(async () => {
const [arg] = (await deserializeJSONStream(
new Response(frame(PENDING_PROMISE_ARGS) + ";0xZZZZZZZZ;junk"),
new Response(frameNode(PENDING_PROMISE_ARGS) + ";0xZZZZZZZZ;junk"),
)) as [Promise<unknown>];
outcome = await settleWithin(arg);
});

expect(unhandled).toEqual([]);
expect(outcome?.status).toBe("rejected");
expect((outcome as { reason: Error }).reason.message).toBe("Malformed server function stream.");
expect(consoleError).toHaveBeenCalledWith(
expect.stringContaining("server function stream"),
(outcome as { reason: Error }).reason,
);
});

it("does not report a pending promise nobody awaits once the body ends", async () => {
const { deserializeJSONStream } = await loadSerialization(true);

const unhandled = await collectUnhandledRejections(async () => {
await deserializeJSONStream(new Response(frame(PENDING_PROMISE_ARGS)));
await deserializeJSONStream(new Response(frameNode(PENDING_PROMISE_ARGS)));
});

expect(unhandled).toEqual([]);
Expand All @@ -313,7 +361,7 @@ describe("values waiting on a later frame", () => {
let arg: Promise<unknown> | undefined;

const unhandled = await collectUnhandledRejections(async () => {
[arg] = (await deserializeJSONStream(new Response(frame(rejectedArgs)))) as [
[arg] = (await deserializeJSONStream(new Response(frameNode(rejectedArgs)))) as [
Promise<unknown>,
];
});
Expand Down Expand Up @@ -362,16 +410,14 @@ describe("values waiting on a later JS frame", () => {

afterEach(() => {
vi.unstubAllEnvs();
vi.restoreAllMocks();
delete (globalThis as any).self;
delete (globalThis as any).$R;
});

async function firstJSFrame(id: string, value: unknown) {
const { serializeToJSStream } = await loadSerialization(true);
const reader = serializeToJSStream(id, value).getReader();
const { value: chunk } = await reader.read();
await reader.cancel();
return new TextDecoder().decode(chunk);
return await readFirstFrame(serializeToJSStream(id, value));
}

it("rejects a promise that is still pending when the body ends", async () => {
Expand Down Expand Up @@ -406,6 +452,7 @@ describe("values waiting on a later JS frame", () => {
it("rejects pending values with the failure when a later frame is malformed", async () => {
const body = await firstJSFrame("server-fn:1", [new Promise(() => {})]);
const { deserializeJSStream } = await loadSerialization(true);
const consoleError = vi.spyOn(console, "error").mockImplementation(() => {});
let outcome: Awaited<ReturnType<typeof settleWithin>> | undefined;

const unhandled = await collectUnhandledRejections(async () => {
Expand All @@ -419,6 +466,134 @@ describe("values waiting on a later JS frame", () => {
expect(unhandled).toEqual([]);
expect(outcome?.status).toBe("rejected");
expect((outcome as { reason: Error }).reason.message).toBe("Malformed server function stream.");
expect(consoleError).toHaveBeenCalledWith(
expect.stringContaining("server function stream"),
(outcome as { reason: Error }).reason,
);
expect((globalThis as any).$R["server-fn:1"]).toBeUndefined();
});

it("releases the scope and cancels the body when the first frame fails", async () => {
const body = await firstJSFrame("server-fn:3", [new Promise(() => {})]);
const source = body.slice(12) + ';throw new Error("first frame failed")';
const { deserializeJSStream } = await loadSerialization(true);
let cancelled = false;
const stream = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode(frame(source)));
},
cancel() {
cancelled = true;
},
});

await expect(deserializeJSStream("server-fn:3", new Response(stream))).rejects.toThrow(
"first frame failed",
);
expect(cancelled).toBe(true);
expect((globalThis as any).$R["server-fn:3"]).toBeUndefined();
});
});

describe("SerovalChunkReader", () => {
beforeEach(() => {
vi.resetModules();
});

afterEach(() => {
vi.unstubAllEnvs();
vi.restoreAllMocks();
});

it("writes a fixed-size header before each chunk", async () => {
const { serializeToJSONString } = await loadSerialization(true);

const payload = await serializeToJSONString(1);

expect(payload).toMatch(/^;0x[0-9a-f]{8};/);
expect(payload).toBe(frame(payload.slice(12)));
});

it("reads chunks split at any byte, including inside a character", async () => {
const { SerovalChunkReader } = await loadSerialization(true);
const body = frame("héllo wörld ✓") + frame("") + frame("second");

const reader = new SerovalChunkReader(streamOf(split(body, 1)));
const chunks: string[] = [];
await reader.drain(chunk => chunks.push(chunk));

expect(chunks).toEqual(["héllo wörld ✓", "", "second"]);
});

it("reads a large chunk in small pieces in linear time", async () => {
const { SerovalChunkReader } = await loadSerialization(true);
const data = "x".repeat(16 * 1024 * 1024);

const start = performance.now();
const result = await new SerovalChunkReader(streamOf(split(frame(data), 16 * 1024))).next();

expect(result.value).toHaveLength(data.length);
// Copying the whole buffer on every piece took about 2 seconds here.
expect(performance.now() - start).toBeLessThan(500);
});

it.each([
["a missing delimiter", "X0x00000003Yabc"],
["a non-hex size", ";0x0000zz03;abc"],
["a missing 0x prefix", ";0000000003;abc"],
["a truncated header", ";0x0000"],
["truncated data", ";0x000000ff;abc"],
])("rejects %s", async (_, body) => {
const { SerovalChunkReader } = await loadSerialization(true);

await expect(new SerovalChunkReader(streamOf([body])).next()).rejects.toThrow(
"Malformed server function stream.",
);
});

it("rejects a chunk over the size limit before buffering it", async () => {
const { SerovalChunkReader } = await loadSerialization(true);
const reader = new SerovalChunkReader(streamOf([";0xffffffff;"]), { maxChunkSize: 1024 });

await expect(reader.next()).rejects.toThrow(/larger than the limit/);
});

it("reports a bad later chunk instead of leaving an unhandled rejection", async () => {
const { serializeToJSONString, deserializeJSONStream } = await loadSerialization(true);
const consoleError = vi.spyOn(console, "error").mockImplementation(() => {});
const unhandled: unknown[] = [];
const onUnhandled = (error: unknown) => unhandled.push(error);
process.on("unhandledRejection", onUnhandled);

try {
const first = await serializeToJSONString([1, 2]);
const value = await deserializeJSONStream(new Response(first + frame("{nope")));
await new Promise(resolve => setTimeout(resolve, 20));

expect(value).toEqual([1, 2]);
expect(unhandled).toEqual([]);
expect(consoleError).toHaveBeenCalledWith(
expect.stringContaining("server function stream"),
expect.any(SyntaxError),
);
} finally {
process.off("unhandledRejection", onUnhandled);
}
});

it("cancels the body when the first chunk cannot be parsed", async () => {
const { deserializeJSONStream } = await loadSerialization(true);
let cancelled = false;
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new TextEncoder().encode(frame("{nope")));
},
cancel() {
cancelled = true;
},
});

await expect(deserializeJSONStream(new Response(body))).rejects.toThrow(SyntaxError);
expect(cancelled).toBe(true);
});
});
Loading
Loading