Sign In

@tanstack/start-client-core

Package Overview
Dependencies
Maintainers
3
Versions
407
Alerts
File Explorer

Advanced tools

Socket logo

Install Socket

Detect and block malicious and high-risk dependencies

Install

@tanstack/start-client-core - npm Package Compare versions

Comparing version
1.170.19
to
1.170.20
+20
-10
dist/esm/client-rpc/frame-decoder.js

@@ -88,3 +88,14 @@ import { FrameType } from "../constants.js";

const bufferList = [];
let bufferHead = 0;
let totalLength = 0;
function advanceBufferHead() {
bufferList[bufferHead++] = EMPTY_BUFFER;
if (bufferHead === bufferList.length) {
bufferList.length = 0;
bufferHead = 0;
} else if (bufferHead >= 32) {
bufferList.splice(0, bufferHead);
bufferHead = 0;
}
}
/**

@@ -96,3 +107,3 @@ * Reads header bytes from buffer chunks without flattening.

if (totalLength < 9) return null;
const first = bufferList[0];
const first = bufferList[bufferHead];
if (first.length >= 9) return {

@@ -106,3 +117,3 @@ type: first[0],

let remaining = 9;
for (let i = 0; i < bufferList.length && remaining > 0; i++) {
for (let i = bufferHead; i < bufferList.length && remaining > 0; i++) {
const chunk = bufferList[i];

@@ -125,7 +136,7 @@ const toCopy = Math.min(chunk.length, remaining);

if (count === 0) return EMPTY_BUFFER;
const first = bufferList[0];
const first = bufferList[bufferHead];
if (first && first.length >= count) {
const result = first.subarray(0, count);
if (first.length === count) bufferList.shift();
else bufferList[0] = first.subarray(count);
if (first.length === count) advanceBufferHead();
else bufferList[bufferHead] = first.subarray(count);
totalLength -= count;

@@ -137,5 +148,4 @@ return result;

let remaining = count;
while (remaining > 0 && bufferList.length > 0) {
const chunk = bufferList[0];
if (!chunk) break;
while (remaining > 0 && bufferHead < bufferList.length) {
const chunk = bufferList[bufferHead];
const toCopy = Math.min(chunk.length, remaining);

@@ -145,4 +155,4 @@ result.set(chunk.subarray(0, toCopy), offset);

remaining -= toCopy;
if (toCopy === chunk.length) bufferList.shift();
else bufferList[0] = chunk.subarray(toCopy);
if (toCopy === chunk.length) advanceBufferHead();
else bufferList[bufferHead] = chunk.subarray(toCopy);
}

@@ -149,0 +159,0 @@ totalLength -= count;

@@ -1,1 +0,1 @@

{"version":3,"file":"frame-decoder.js","names":[],"sources":["../../../src/client-rpc/frame-decoder.ts"],"sourcesContent":["/**\n * Client-side frame decoder for multiplexed responses.\n *\n * Decodes binary frame protocol and reconstructs:\n * - JSON stream (NDJSON lines for seroval)\n * - Raw streams (binary data as ReadableStream<Uint8Array>)\n */\n\nimport { FRAME_HEADER_SIZE, FrameType } from '../constants'\n\n/** Cached TextDecoder for frame decoding */\nconst textDecoder = new TextDecoder()\n\n/** Shared empty buffer for empty buffer case - avoids allocation */\nconst EMPTY_BUFFER = new Uint8Array(0)\n\n/** Hardening limits to prevent memory/CPU DoS */\nconst MAX_FRAME_PAYLOAD_SIZE = 16 * 1024 * 1024 // 16MiB\nconst MAX_BUFFERED_BYTES = 32 * 1024 * 1024 // 32MiB\nconst MAX_STREAMS = 1024\nconst MAX_FRAMES = 100_000 // Limit total frames to prevent CPU DoS\n\n/**\n * Result of frame decoding.\n */\nexport interface FrameDecoderResult {\n /** Gets or creates a raw stream by ID (for use by deserialize plugin) */\n getStream: (id: number) => ReadableStream<Uint8Array>\n /** Stream of JSON strings (NDJSON lines) */\n chunks: ReadableStream<string>\n}\n\n/**\n * Creates a frame decoder that processes a multiplexed response stream.\n *\n * @param input The raw response body stream\n * @returns Decoded JSON stream and stream getter function\n */\nexport function createFrameDecoder(\n input: ReadableStream<Uint8Array>,\n): FrameDecoderResult {\n const streamControllers = new Map<\n number,\n ReadableStreamDefaultController<Uint8Array>\n >()\n const streams = new Map<number, ReadableStream<Uint8Array>>()\n const cancelledStreamIds = new Set<number>()\n\n let cancelled = false as boolean\n let inputReader: ReadableStreamReader<Uint8Array> | null = null\n let frameCount = 0\n\n let jsonController!: ReadableStreamDefaultController<string>\n const jsonChunks = new ReadableStream<string>({\n start(controller) {\n jsonController = controller\n },\n cancel() {\n cancelled = true\n try {\n inputReader?.cancel()\n } catch {\n // Ignore\n }\n\n streamControllers.forEach((ctrl) => {\n try {\n ctrl.error(new Error('Framed response cancelled'))\n } catch {\n // Ignore\n }\n })\n streamControllers.clear()\n streams.clear()\n cancelledStreamIds.clear()\n },\n })\n\n /**\n * Gets or creates a stream for a given stream ID.\n * Called by deserialize plugin when it encounters a RawStream reference.\n */\n function getOrCreateStream(id: number): ReadableStream<Uint8Array> {\n const existing = streams.get(id)\n if (existing) {\n return existing\n }\n\n // If we already received an END/ERROR for this streamId, returning a fresh stream\n // would hang consumers. Return an already-closed stream instead.\n if (cancelledStreamIds.has(id)) {\n return new ReadableStream<Uint8Array>({\n start(controller) {\n controller.close()\n },\n })\n }\n\n if (streams.size >= MAX_STREAMS) {\n throw new Error(\n `Too many raw streams in framed response (max ${MAX_STREAMS})`,\n )\n }\n\n const stream = new ReadableStream<Uint8Array>({\n start(ctrl) {\n streamControllers.set(id, ctrl)\n },\n cancel() {\n cancelledStreamIds.add(id)\n streamControllers.delete(id)\n streams.delete(id)\n },\n })\n streams.set(id, stream)\n return stream\n }\n\n /**\n * Ensures stream exists and returns its controller for enqueuing data.\n * Used for CHUNK frames where we need to ensure stream is created.\n */\n function ensureController(\n id: number,\n ): ReadableStreamDefaultController<Uint8Array> | undefined {\n getOrCreateStream(id)\n return streamControllers.get(id)\n }\n\n // Process frames asynchronously\n ;(async () => {\n const reader = input.getReader()\n inputReader = reader\n\n const bufferList: Array<Uint8Array> = []\n let totalLength = 0\n\n /**\n * Reads header bytes from buffer chunks without flattening.\n * Returns header data or null if not enough bytes available.\n */\n function readHeader(): {\n type: number\n streamId: number\n length: number\n } | null {\n if (totalLength < FRAME_HEADER_SIZE) return null\n\n const first = bufferList[0]!\n\n // Fast path: header fits entirely in first chunk (common case)\n if (first.length >= FRAME_HEADER_SIZE) {\n const type = first[0]!\n const streamId =\n ((first[1]! << 24) |\n (first[2]! << 16) |\n (first[3]! << 8) |\n first[4]!) >>>\n 0\n const length =\n ((first[5]! << 24) |\n (first[6]! << 16) |\n (first[7]! << 8) |\n first[8]!) >>>\n 0\n return { type, streamId, length }\n }\n\n // Slow path: header spans multiple chunks - flatten header bytes only\n const headerBytes = new Uint8Array(FRAME_HEADER_SIZE)\n let offset = 0\n let remaining = FRAME_HEADER_SIZE\n for (let i = 0; i < bufferList.length && remaining > 0; i++) {\n const chunk = bufferList[i]!\n const toCopy = Math.min(chunk.length, remaining)\n headerBytes.set(chunk.subarray(0, toCopy), offset)\n offset += toCopy\n remaining -= toCopy\n }\n\n const type = headerBytes[0]!\n const streamId =\n ((headerBytes[1]! << 24) |\n (headerBytes[2]! << 16) |\n (headerBytes[3]! << 8) |\n headerBytes[4]!) >>>\n 0\n const length =\n ((headerBytes[5]! << 24) |\n (headerBytes[6]! << 16) |\n (headerBytes[7]! << 8) |\n headerBytes[8]!) >>>\n 0\n\n return { type, streamId, length }\n }\n\n /**\n * Flattens buffer list into single Uint8Array and removes from list.\n */\n function extractFlattened(count: number): Uint8Array {\n if (count === 0) return EMPTY_BUFFER\n\n // Fast path: the requested bytes are fully contained in the first buffered\n // chunk (the common case — most frames arrive within a single network\n // read). Return a subarray view instead of allocating a new buffer and\n // copying `count` bytes. The view shares the chunk's backing ArrayBuffer,\n // which is safe because buffered chunks are never mutated in place after\n // being read from the network.\n const first = bufferList[0]\n if (first && first.length >= count) {\n const result = first.subarray(0, count)\n if (first.length === count) {\n bufferList.shift()\n } else {\n bufferList[0] = first.subarray(count)\n }\n totalLength -= count\n return result\n }\n\n // Slow path: the requested bytes span multiple chunks — flatten by copying.\n const result = new Uint8Array(count)\n let offset = 0\n let remaining = count\n\n while (remaining > 0 && bufferList.length > 0) {\n const chunk = bufferList[0]\n if (!chunk) break\n const toCopy = Math.min(chunk.length, remaining)\n result.set(chunk.subarray(0, toCopy), offset)\n\n offset += toCopy\n remaining -= toCopy\n\n if (toCopy === chunk.length) {\n bufferList.shift()\n } else {\n bufferList[0] = chunk.subarray(toCopy)\n }\n }\n\n totalLength -= count\n return result\n }\n\n try {\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n while (true) {\n const { done, value } = await reader.read()\n if (cancelled) break\n if (done) break\n\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (!value) continue\n\n // Append incoming chunk to buffer list\n if (totalLength + value.length > MAX_BUFFERED_BYTES) {\n throw new Error(\n `Framed response buffer exceeded ${MAX_BUFFERED_BYTES} bytes`,\n )\n }\n bufferList.push(value)\n totalLength += value.length\n\n // Parse complete frames from buffer\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n while (true) {\n const header = readHeader()\n if (!header) break // Not enough bytes for header\n\n const { type, streamId, length } = header\n\n if (\n type !== FrameType.JSON &&\n type !== FrameType.CHUNK &&\n type !== FrameType.END &&\n type !== FrameType.ERROR\n ) {\n throw new Error(`Unknown frame type: ${type}`)\n }\n\n // Enforce stream id conventions: JSON uses streamId 0, raw streams use non-zero ids\n if (type === FrameType.JSON) {\n if (streamId !== 0) {\n throw new Error('Invalid JSON frame streamId (expected 0)')\n }\n } else {\n if (streamId === 0) {\n throw new Error('Invalid raw frame streamId (expected non-zero)')\n }\n }\n\n if (length > MAX_FRAME_PAYLOAD_SIZE) {\n throw new Error(\n `Frame payload too large: ${length} bytes (max ${MAX_FRAME_PAYLOAD_SIZE})`,\n )\n }\n\n const frameSize = FRAME_HEADER_SIZE + length\n if (totalLength < frameSize) break // Wait for more data\n\n if (++frameCount > MAX_FRAMES) {\n throw new Error(\n `Too many frames in framed response (max ${MAX_FRAMES})`,\n )\n }\n\n // Extract and consume header bytes\n extractFlattened(FRAME_HEADER_SIZE)\n\n // Extract payload\n const payload = extractFlattened(length)\n\n // Process frame by type\n switch (type) {\n case FrameType.JSON: {\n try {\n jsonController.enqueue(textDecoder.decode(payload))\n } catch {\n // JSON stream may be cancelled/closed\n }\n break\n }\n\n case FrameType.CHUNK: {\n const ctrl = ensureController(streamId)\n if (ctrl) {\n ctrl.enqueue(payload)\n }\n break\n }\n\n case FrameType.END: {\n const ctrl = ensureController(streamId)\n cancelledStreamIds.add(streamId)\n if (ctrl) {\n try {\n ctrl.close()\n } catch {\n // Already closed\n }\n streamControllers.delete(streamId)\n }\n break\n }\n\n case FrameType.ERROR: {\n const ctrl = ensureController(streamId)\n cancelledStreamIds.add(streamId)\n if (ctrl) {\n const message = textDecoder.decode(payload)\n ctrl.error(new Error(message))\n streamControllers.delete(streamId)\n }\n break\n }\n }\n }\n }\n\n if (totalLength !== 0) {\n throw new Error('Incomplete frame at end of framed response')\n }\n\n // Close JSON stream when done\n try {\n jsonController.close()\n } catch {\n // JSON stream may be cancelled/closed\n }\n\n // Close any remaining streams (shouldn't happen in normal operation)\n streamControllers.forEach((ctrl) => {\n try {\n ctrl.close()\n } catch {\n // Already closed\n }\n })\n streamControllers.clear()\n } catch (error) {\n // Error reading - propagate to all streams\n try {\n jsonController.error(error)\n } catch {\n // Already errored/closed\n }\n streamControllers.forEach((ctrl) => {\n try {\n ctrl.error(error)\n } catch {\n // Already errored/closed\n }\n })\n streamControllers.clear()\n } finally {\n try {\n reader.releaseLock()\n } catch {\n // Ignore\n }\n inputReader = null\n }\n })()\n\n return { getStream: getOrCreateStream, chunks: jsonChunks }\n}\n"],"mappings":";;;;;;;;;;AAWA,IAAM,cAAc,IAAI,YAAY;;AAGpC,IAAM,eAAe,IAAI,WAAW,CAAC;;AAGrC,IAAM,yBAAyB,KAAK,OAAO;AAC3C,IAAM,qBAAqB,KAAK,OAAO;AACvC,IAAM,cAAc;AACpB,IAAM,aAAa;;;;;;;AAkBnB,SAAgB,mBACd,OACoB;CACpB,MAAM,oCAAoB,IAAI,IAG5B;CACF,MAAM,0BAAU,IAAI,IAAwC;CAC5D,MAAM,qCAAqB,IAAI,IAAY;CAE3C,IAAI,YAAY;CAChB,IAAI,cAAuD;CAC3D,IAAI,aAAa;CAEjB,IAAI;CACJ,MAAM,aAAa,IAAI,eAAuB;EAC5C,MAAM,YAAY;GAChB,iBAAiB;EACnB;EACA,SAAS;GACP,YAAY;GACZ,IAAI;IACF,aAAa,OAAO;GACtB,QAAQ,CAER;GAEA,kBAAkB,SAAS,SAAS;IAClC,IAAI;KACF,KAAK,sBAAM,IAAI,MAAM,2BAA2B,CAAC;IACnD,QAAQ,CAER;GACF,CAAC;GACD,kBAAkB,MAAM;GACxB,QAAQ,MAAM;GACd,mBAAmB,MAAM;EAC3B;CACF,CAAC;;;;;CAMD,SAAS,kBAAkB,IAAwC;EACjE,MAAM,WAAW,QAAQ,IAAI,EAAE;EAC/B,IAAI,UACF,OAAO;EAKT,IAAI,mBAAmB,IAAI,EAAE,GAC3B,OAAO,IAAI,eAA2B,EACpC,MAAM,YAAY;GAChB,WAAW,MAAM;EACnB,EACF,CAAC;EAGH,IAAI,QAAQ,QAAQ,aAClB,MAAM,IAAI,MACR,gDAAgD,YAAY,EAC9D;EAGF,MAAM,SAAS,IAAI,eAA2B;GAC5C,MAAM,MAAM;IACV,kBAAkB,IAAI,IAAI,IAAI;GAChC;GACA,SAAS;IACP,mBAAmB,IAAI,EAAE;IACzB,kBAAkB,OAAO,EAAE;IAC3B,QAAQ,OAAO,EAAE;GACnB;EACF,CAAC;EACD,QAAQ,IAAI,IAAI,MAAM;EACtB,OAAO;CACT;;;;;CAMA,SAAS,iBACP,IACyD;EACzD,kBAAkB,EAAE;EACpB,OAAO,kBAAkB,IAAI,EAAE;CACjC;CAGC,CAAC,YAAY;EACZ,MAAM,SAAS,MAAM,UAAU;EAC/B,cAAc;EAEd,MAAM,aAAgC,CAAC;EACvC,IAAI,cAAc;;;;;EAMlB,SAAS,aAIA;GACP,IAAI,cAAA,GAAiC,OAAO;GAE5C,MAAM,QAAQ,WAAW;GAGzB,IAAI,MAAM,UAAA,GAcR,OAAO;IAAE,MAbI,MAAM;IAaJ,WAXX,MAAM,MAAO,KACZ,MAAM,MAAO,KACb,MAAM,MAAO,IACd,MAAM,QACR;IAOuB,SALrB,MAAM,MAAO,KACZ,MAAM,MAAO,KACb,MAAM,MAAO,IACd,MAAM,QACR;GAC8B;GAIlC,MAAM,cAAc,IAAI,WAAA,CAA4B;GACpD,IAAI,SAAS;GACb,IAAI,YAAA;GACJ,KAAK,IAAI,IAAI,GAAG,IAAI,WAAW,UAAU,YAAY,GAAG,KAAK;IAC3D,MAAM,QAAQ,WAAW;IACzB,MAAM,SAAS,KAAK,IAAI,MAAM,QAAQ,SAAS;IAC/C,YAAY,IAAI,MAAM,SAAS,GAAG,MAAM,GAAG,MAAM;IACjD,UAAU;IACV,aAAa;GACf;GAgBA,OAAO;IAAE,MAdI,YAAY;IAcV,WAZX,YAAY,MAAO,KAClB,YAAY,MAAO,KACnB,YAAY,MAAO,IACpB,YAAY,QACd;IAQuB,SANrB,YAAY,MAAO,KAClB,YAAY,MAAO,KACnB,YAAY,MAAO,IACpB,YAAY,QACd;GAE8B;EAClC;;;;EAKA,SAAS,iBAAiB,OAA2B;GACnD,IAAI,UAAU,GAAG,OAAO;GAQxB,MAAM,QAAQ,WAAW;GACzB,IAAI,SAAS,MAAM,UAAU,OAAO;IAClC,MAAM,SAAS,MAAM,SAAS,GAAG,KAAK;IACtC,IAAI,MAAM,WAAW,OACnB,WAAW,MAAM;SAEjB,WAAW,KAAK,MAAM,SAAS,KAAK;IAEtC,eAAe;IACf,OAAO;GACT;GAGA,MAAM,SAAS,IAAI,WAAW,KAAK;GACnC,IAAI,SAAS;GACb,IAAI,YAAY;GAEhB,OAAO,YAAY,KAAK,WAAW,SAAS,GAAG;IAC7C,MAAM,QAAQ,WAAW;IACzB,IAAI,CAAC,OAAO;IACZ,MAAM,SAAS,KAAK,IAAI,MAAM,QAAQ,SAAS;IAC/C,OAAO,IAAI,MAAM,SAAS,GAAG,MAAM,GAAG,MAAM;IAE5C,UAAU;IACV,aAAa;IAEb,IAAI,WAAW,MAAM,QACnB,WAAW,MAAM;SAEjB,WAAW,KAAK,MAAM,SAAS,MAAM;GAEzC;GAEA,eAAe;GACf,OAAO;EACT;EAEA,IAAI;GAEF,OAAO,MAAM;IACX,MAAM,EAAE,MAAM,UAAU,MAAM,OAAO,KAAK;IAC1C,IAAI,WAAW;IACf,IAAI,MAAM;IAGV,IAAI,CAAC,OAAO;IAGZ,IAAI,cAAc,MAAM,SAAS,oBAC/B,MAAM,IAAI,MACR,mCAAmC,mBAAmB,OACxD;IAEF,WAAW,KAAK,KAAK;IACrB,eAAe,MAAM;IAIrB,OAAO,MAAM;KACX,MAAM,SAAS,WAAW;KAC1B,IAAI,CAAC,QAAQ;KAEb,MAAM,EAAE,MAAM,UAAU,WAAW;KAEnC,IACE,SAAS,UAAU,QACnB,SAAS,UAAU,SACnB,SAAS,UAAU,OACnB,SAAS,UAAU,OAEnB,MAAM,IAAI,MAAM,uBAAuB,MAAM;KAI/C,IAAI,SAAS,UAAU;UACjB,aAAa,GACf,MAAM,IAAI,MAAM,0CAA0C;KAAA,OAG5D,IAAI,aAAa,GACf,MAAM,IAAI,MAAM,gDAAgD;KAIpE,IAAI,SAAS,wBACX,MAAM,IAAI,MACR,4BAA4B,OAAO,cAAc,uBAAuB,EAC1E;KAGF,MAAM,YAAA,IAAgC;KACtC,IAAI,cAAc,WAAW;KAE7B,IAAI,EAAE,aAAa,YACjB,MAAM,IAAI,MACR,2CAA2C,WAAW,EACxD;KAIF,iBAAA,CAAkC;KAGlC,MAAM,UAAU,iBAAiB,MAAM;KAGvC,QAAQ,MAAR;MACE,KAAK,UAAU;OACb,IAAI;QACF,eAAe,QAAQ,YAAY,OAAO,OAAO,CAAC;OACpD,QAAQ,CAER;OACA;MAGF,KAAK,UAAU,OAAO;OACpB,MAAM,OAAO,iBAAiB,QAAQ;OACtC,IAAI,MACF,KAAK,QAAQ,OAAO;OAEtB;MACF;MAEA,KAAK,UAAU,KAAK;OAClB,MAAM,OAAO,iBAAiB,QAAQ;OACtC,mBAAmB,IAAI,QAAQ;OAC/B,IAAI,MAAM;QACR,IAAI;SACF,KAAK,MAAM;QACb,QAAQ,CAER;QACA,kBAAkB,OAAO,QAAQ;OACnC;OACA;MACF;MAEA,KAAK,UAAU,OAAO;OACpB,MAAM,OAAO,iBAAiB,QAAQ;OACtC,mBAAmB,IAAI,QAAQ;OAC/B,IAAI,MAAM;QACR,MAAM,UAAU,YAAY,OAAO,OAAO;QAC1C,KAAK,MAAM,IAAI,MAAM,OAAO,CAAC;QAC7B,kBAAkB,OAAO,QAAQ;OACnC;OACA;MACF;KACF;IACF;GACF;GAEA,IAAI,gBAAgB,GAClB,MAAM,IAAI,MAAM,4CAA4C;GAI9D,IAAI;IACF,eAAe,MAAM;GACvB,QAAQ,CAER;GAGA,kBAAkB,SAAS,SAAS;IAClC,IAAI;KACF,KAAK,MAAM;IACb,QAAQ,CAER;GACF,CAAC;GACD,kBAAkB,MAAM;EAC1B,SAAS,OAAO;GAEd,IAAI;IACF,eAAe,MAAM,KAAK;GAC5B,QAAQ,CAER;GACA,kBAAkB,SAAS,SAAS;IAClC,IAAI;KACF,KAAK,MAAM,KAAK;IAClB,QAAQ,CAER;GACF,CAAC;GACD,kBAAkB,MAAM;EAC1B,UAAU;GACR,IAAI;IACF,OAAO,YAAY;GACrB,QAAQ,CAER;GACA,cAAc;EAChB;CACF,GAAG;CAEH,OAAO;EAAE,WAAW;EAAmB,QAAQ;CAAW;AAC5D"}
{"version":3,"file":"frame-decoder.js","names":[],"sources":["../../../src/client-rpc/frame-decoder.ts"],"sourcesContent":["/**\n * Client-side frame decoder for multiplexed responses.\n *\n * Decodes binary frame protocol and reconstructs:\n * - JSON stream (NDJSON lines for seroval)\n * - Raw streams (binary data as ReadableStream<Uint8Array>)\n */\n\nimport { FRAME_HEADER_SIZE, FrameType } from '../constants'\n\n/** Cached TextDecoder for frame decoding */\nconst textDecoder = new TextDecoder()\n\n/** Shared empty buffer for empty buffer case - avoids allocation */\nconst EMPTY_BUFFER = new Uint8Array(0)\n\n/** Hardening limits to prevent memory/CPU DoS */\nconst MAX_FRAME_PAYLOAD_SIZE = 16 * 1024 * 1024 // 16MiB\nconst MAX_BUFFERED_BYTES = 32 * 1024 * 1024 // 32MiB\nconst MAX_STREAMS = 1024\nconst MAX_FRAMES = 100_000 // Limit total frames to prevent CPU DoS\n\n/**\n * Result of frame decoding.\n */\nexport interface FrameDecoderResult {\n /** Gets or creates a raw stream by ID (for use by deserialize plugin) */\n getStream: (id: number) => ReadableStream<Uint8Array>\n /** Stream of JSON strings (NDJSON lines) */\n chunks: ReadableStream<string>\n}\n\n/**\n * Creates a frame decoder that processes a multiplexed response stream.\n *\n * @param input The raw response body stream\n * @returns Decoded JSON stream and stream getter function\n */\nexport function createFrameDecoder(\n input: ReadableStream<Uint8Array>,\n): FrameDecoderResult {\n const streamControllers = new Map<\n number,\n ReadableStreamDefaultController<Uint8Array>\n >()\n const streams = new Map<number, ReadableStream<Uint8Array>>()\n const cancelledStreamIds = new Set<number>()\n\n let cancelled = false as boolean\n let inputReader: ReadableStreamReader<Uint8Array> | null = null\n let frameCount = 0\n\n let jsonController!: ReadableStreamDefaultController<string>\n const jsonChunks = new ReadableStream<string>({\n start(controller) {\n jsonController = controller\n },\n cancel() {\n cancelled = true\n try {\n inputReader?.cancel()\n } catch {\n // Ignore\n }\n\n streamControllers.forEach((ctrl) => {\n try {\n ctrl.error(new Error('Framed response cancelled'))\n } catch {\n // Ignore\n }\n })\n streamControllers.clear()\n streams.clear()\n cancelledStreamIds.clear()\n },\n })\n\n /**\n * Gets or creates a stream for a given stream ID.\n * Called by deserialize plugin when it encounters a RawStream reference.\n */\n function getOrCreateStream(id: number): ReadableStream<Uint8Array> {\n const existing = streams.get(id)\n if (existing) {\n return existing\n }\n\n // If we already received an END/ERROR for this streamId, returning a fresh stream\n // would hang consumers. Return an already-closed stream instead.\n if (cancelledStreamIds.has(id)) {\n return new ReadableStream<Uint8Array>({\n start(controller) {\n controller.close()\n },\n })\n }\n\n if (streams.size >= MAX_STREAMS) {\n throw new Error(\n `Too many raw streams in framed response (max ${MAX_STREAMS})`,\n )\n }\n\n const stream = new ReadableStream<Uint8Array>({\n start(ctrl) {\n streamControllers.set(id, ctrl)\n },\n cancel() {\n cancelledStreamIds.add(id)\n streamControllers.delete(id)\n streams.delete(id)\n },\n })\n streams.set(id, stream)\n return stream\n }\n\n /**\n * Ensures stream exists and returns its controller for enqueuing data.\n * Used for CHUNK frames where we need to ensure stream is created.\n */\n function ensureController(\n id: number,\n ): ReadableStreamDefaultController<Uint8Array> | undefined {\n getOrCreateStream(id)\n return streamControllers.get(id)\n }\n\n // Process frames asynchronously\n ;(async () => {\n const reader = input.getReader()\n inputReader = reader\n\n const bufferList: Array<Uint8Array> = []\n // Index of the first un-consumed chunk in bufferList. Advancing this\n // pointer is O(1); using bufferList.shift() to drop a consumed chunk is\n // O(n) and degrades to O(n^2) when a single large frame is assembled from\n // many small chunks (e.g. a big RawStream payload split across reads).\n let bufferHead = 0\n let totalLength = 0\n\n function advanceBufferHead(): void {\n bufferList[bufferHead++] = EMPTY_BUFFER\n\n // Reset drained buffers immediately and compact long-lived buffers in batches.\n if (bufferHead === bufferList.length) {\n bufferList.length = 0\n bufferHead = 0\n } else if (bufferHead >= 32) {\n bufferList.splice(0, bufferHead)\n bufferHead = 0\n }\n }\n\n /**\n * Reads header bytes from buffer chunks without flattening.\n * Returns header data or null if not enough bytes available.\n */\n function readHeader(): {\n type: number\n streamId: number\n length: number\n } | null {\n if (totalLength < FRAME_HEADER_SIZE) return null\n\n const first = bufferList[bufferHead]!\n\n // Fast path: header fits entirely in first chunk (common case)\n if (first.length >= FRAME_HEADER_SIZE) {\n const type = first[0]!\n const streamId =\n ((first[1]! << 24) |\n (first[2]! << 16) |\n (first[3]! << 8) |\n first[4]!) >>>\n 0\n const length =\n ((first[5]! << 24) |\n (first[6]! << 16) |\n (first[7]! << 8) |\n first[8]!) >>>\n 0\n return { type, streamId, length }\n }\n\n // Slow path: header spans multiple chunks - flatten header bytes only\n const headerBytes = new Uint8Array(FRAME_HEADER_SIZE)\n let offset = 0\n let remaining = FRAME_HEADER_SIZE\n for (let i = bufferHead; i < bufferList.length && remaining > 0; i++) {\n const chunk = bufferList[i]!\n const toCopy = Math.min(chunk.length, remaining)\n headerBytes.set(chunk.subarray(0, toCopy), offset)\n offset += toCopy\n remaining -= toCopy\n }\n\n const type = headerBytes[0]!\n const streamId =\n ((headerBytes[1]! << 24) |\n (headerBytes[2]! << 16) |\n (headerBytes[3]! << 8) |\n headerBytes[4]!) >>>\n 0\n const length =\n ((headerBytes[5]! << 24) |\n (headerBytes[6]! << 16) |\n (headerBytes[7]! << 8) |\n headerBytes[8]!) >>>\n 0\n\n return { type, streamId, length }\n }\n\n /**\n * Flattens buffer list into single Uint8Array and removes from list.\n */\n function extractFlattened(count: number): Uint8Array {\n if (count === 0) return EMPTY_BUFFER\n\n // Fast path: the requested bytes are fully contained in the first buffered\n // chunk (the common case — most frames arrive within a single network\n // read). Return a subarray view instead of allocating a new buffer and\n // copying `count` bytes. The view shares the chunk's backing ArrayBuffer,\n // which is safe because buffered chunks are never mutated in place after\n // being read from the network.\n const first = bufferList[bufferHead]\n if (first && first.length >= count) {\n const result = first.subarray(0, count)\n if (first.length === count) {\n advanceBufferHead()\n } else {\n bufferList[bufferHead] = first.subarray(count)\n }\n totalLength -= count\n return result\n }\n\n // Slow path: the requested bytes span multiple chunks — flatten by copying.\n const result = new Uint8Array(count)\n let offset = 0\n let remaining = count\n\n while (remaining > 0 && bufferHead < bufferList.length) {\n const chunk = bufferList[bufferHead]!\n const toCopy = Math.min(chunk.length, remaining)\n result.set(chunk.subarray(0, toCopy), offset)\n\n offset += toCopy\n remaining -= toCopy\n\n if (toCopy === chunk.length) {\n advanceBufferHead()\n } else {\n bufferList[bufferHead] = chunk.subarray(toCopy)\n }\n }\n\n totalLength -= count\n return result\n }\n\n try {\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n while (true) {\n const { done, value } = await reader.read()\n if (cancelled) break\n if (done) break\n\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (!value) continue\n\n // Append incoming chunk to buffer list\n if (totalLength + value.length > MAX_BUFFERED_BYTES) {\n throw new Error(\n `Framed response buffer exceeded ${MAX_BUFFERED_BYTES} bytes`,\n )\n }\n bufferList.push(value)\n totalLength += value.length\n\n // Parse complete frames from buffer\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n while (true) {\n const header = readHeader()\n if (!header) break // Not enough bytes for header\n\n const { type, streamId, length } = header\n\n if (\n type !== FrameType.JSON &&\n type !== FrameType.CHUNK &&\n type !== FrameType.END &&\n type !== FrameType.ERROR\n ) {\n throw new Error(`Unknown frame type: ${type}`)\n }\n\n // Enforce stream id conventions: JSON uses streamId 0, raw streams use non-zero ids\n if (type === FrameType.JSON) {\n if (streamId !== 0) {\n throw new Error('Invalid JSON frame streamId (expected 0)')\n }\n } else {\n if (streamId === 0) {\n throw new Error('Invalid raw frame streamId (expected non-zero)')\n }\n }\n\n if (length > MAX_FRAME_PAYLOAD_SIZE) {\n throw new Error(\n `Frame payload too large: ${length} bytes (max ${MAX_FRAME_PAYLOAD_SIZE})`,\n )\n }\n\n const frameSize = FRAME_HEADER_SIZE + length\n if (totalLength < frameSize) break // Wait for more data\n\n if (++frameCount > MAX_FRAMES) {\n throw new Error(\n `Too many frames in framed response (max ${MAX_FRAMES})`,\n )\n }\n\n // Extract and consume header bytes\n extractFlattened(FRAME_HEADER_SIZE)\n\n // Extract payload\n const payload = extractFlattened(length)\n\n // Process frame by type\n switch (type) {\n case FrameType.JSON: {\n try {\n jsonController.enqueue(textDecoder.decode(payload))\n } catch {\n // JSON stream may be cancelled/closed\n }\n break\n }\n\n case FrameType.CHUNK: {\n const ctrl = ensureController(streamId)\n if (ctrl) {\n ctrl.enqueue(payload)\n }\n break\n }\n\n case FrameType.END: {\n const ctrl = ensureController(streamId)\n cancelledStreamIds.add(streamId)\n if (ctrl) {\n try {\n ctrl.close()\n } catch {\n // Already closed\n }\n streamControllers.delete(streamId)\n }\n break\n }\n\n case FrameType.ERROR: {\n const ctrl = ensureController(streamId)\n cancelledStreamIds.add(streamId)\n if (ctrl) {\n const message = textDecoder.decode(payload)\n ctrl.error(new Error(message))\n streamControllers.delete(streamId)\n }\n break\n }\n }\n }\n }\n\n if (totalLength !== 0) {\n throw new Error('Incomplete frame at end of framed response')\n }\n\n // Close JSON stream when done\n try {\n jsonController.close()\n } catch {\n // JSON stream may be cancelled/closed\n }\n\n // Close any remaining streams (shouldn't happen in normal operation)\n streamControllers.forEach((ctrl) => {\n try {\n ctrl.close()\n } catch {\n // Already closed\n }\n })\n streamControllers.clear()\n } catch (error) {\n // Error reading - propagate to all streams\n try {\n jsonController.error(error)\n } catch {\n // Already errored/closed\n }\n streamControllers.forEach((ctrl) => {\n try {\n ctrl.error(error)\n } catch {\n // Already errored/closed\n }\n })\n streamControllers.clear()\n } finally {\n try {\n reader.releaseLock()\n } catch {\n // Ignore\n }\n inputReader = null\n }\n })()\n\n return { getStream: getOrCreateStream, chunks: jsonChunks }\n}\n"],"mappings":";;;;;;;;;;AAWA,IAAM,cAAc,IAAI,YAAY;;AAGpC,IAAM,eAAe,IAAI,WAAW,CAAC;;AAGrC,IAAM,yBAAyB,KAAK,OAAO;AAC3C,IAAM,qBAAqB,KAAK,OAAO;AACvC,IAAM,cAAc;AACpB,IAAM,aAAa;;;;;;;AAkBnB,SAAgB,mBACd,OACoB;CACpB,MAAM,oCAAoB,IAAI,IAG5B;CACF,MAAM,0BAAU,IAAI,IAAwC;CAC5D,MAAM,qCAAqB,IAAI,IAAY;CAE3C,IAAI,YAAY;CAChB,IAAI,cAAuD;CAC3D,IAAI,aAAa;CAEjB,IAAI;CACJ,MAAM,aAAa,IAAI,eAAuB;EAC5C,MAAM,YAAY;GAChB,iBAAiB;EACnB;EACA,SAAS;GACP,YAAY;GACZ,IAAI;IACF,aAAa,OAAO;GACtB,QAAQ,CAER;GAEA,kBAAkB,SAAS,SAAS;IAClC,IAAI;KACF,KAAK,sBAAM,IAAI,MAAM,2BAA2B,CAAC;IACnD,QAAQ,CAER;GACF,CAAC;GACD,kBAAkB,MAAM;GACxB,QAAQ,MAAM;GACd,mBAAmB,MAAM;EAC3B;CACF,CAAC;;;;;CAMD,SAAS,kBAAkB,IAAwC;EACjE,MAAM,WAAW,QAAQ,IAAI,EAAE;EAC/B,IAAI,UACF,OAAO;EAKT,IAAI,mBAAmB,IAAI,EAAE,GAC3B,OAAO,IAAI,eAA2B,EACpC,MAAM,YAAY;GAChB,WAAW,MAAM;EACnB,EACF,CAAC;EAGH,IAAI,QAAQ,QAAQ,aAClB,MAAM,IAAI,MACR,gDAAgD,YAAY,EAC9D;EAGF,MAAM,SAAS,IAAI,eAA2B;GAC5C,MAAM,MAAM;IACV,kBAAkB,IAAI,IAAI,IAAI;GAChC;GACA,SAAS;IACP,mBAAmB,IAAI,EAAE;IACzB,kBAAkB,OAAO,EAAE;IAC3B,QAAQ,OAAO,EAAE;GACnB;EACF,CAAC;EACD,QAAQ,IAAI,IAAI,MAAM;EACtB,OAAO;CACT;;;;;CAMA,SAAS,iBACP,IACyD;EACzD,kBAAkB,EAAE;EACpB,OAAO,kBAAkB,IAAI,EAAE;CACjC;CAGC,CAAC,YAAY;EACZ,MAAM,SAAS,MAAM,UAAU;EAC/B,cAAc;EAEd,MAAM,aAAgC,CAAC;EAKvC,IAAI,aAAa;EACjB,IAAI,cAAc;EAElB,SAAS,oBAA0B;GACjC,WAAW,gBAAgB;GAG3B,IAAI,eAAe,WAAW,QAAQ;IACpC,WAAW,SAAS;IACpB,aAAa;GACf,OAAO,IAAI,cAAc,IAAI;IAC3B,WAAW,OAAO,GAAG,UAAU;IAC/B,aAAa;GACf;EACF;;;;;EAMA,SAAS,aAIA;GACP,IAAI,cAAA,GAAiC,OAAO;GAE5C,MAAM,QAAQ,WAAW;GAGzB,IAAI,MAAM,UAAA,GAcR,OAAO;IAAE,MAbI,MAAM;IAaJ,WAXX,MAAM,MAAO,KACZ,MAAM,MAAO,KACb,MAAM,MAAO,IACd,MAAM,QACR;IAOuB,SALrB,MAAM,MAAO,KACZ,MAAM,MAAO,KACb,MAAM,MAAO,IACd,MAAM,QACR;GAC8B;GAIlC,MAAM,cAAc,IAAI,WAAA,CAA4B;GACpD,IAAI,SAAS;GACb,IAAI,YAAA;GACJ,KAAK,IAAI,IAAI,YAAY,IAAI,WAAW,UAAU,YAAY,GAAG,KAAK;IACpE,MAAM,QAAQ,WAAW;IACzB,MAAM,SAAS,KAAK,IAAI,MAAM,QAAQ,SAAS;IAC/C,YAAY,IAAI,MAAM,SAAS,GAAG,MAAM,GAAG,MAAM;IACjD,UAAU;IACV,aAAa;GACf;GAgBA,OAAO;IAAE,MAdI,YAAY;IAcV,WAZX,YAAY,MAAO,KAClB,YAAY,MAAO,KACnB,YAAY,MAAO,IACpB,YAAY,QACd;IAQuB,SANrB,YAAY,MAAO,KAClB,YAAY,MAAO,KACnB,YAAY,MAAO,IACpB,YAAY,QACd;GAE8B;EAClC;;;;EAKA,SAAS,iBAAiB,OAA2B;GACnD,IAAI,UAAU,GAAG,OAAO;GAQxB,MAAM,QAAQ,WAAW;GACzB,IAAI,SAAS,MAAM,UAAU,OAAO;IAClC,MAAM,SAAS,MAAM,SAAS,GAAG,KAAK;IACtC,IAAI,MAAM,WAAW,OACnB,kBAAkB;SAElB,WAAW,cAAc,MAAM,SAAS,KAAK;IAE/C,eAAe;IACf,OAAO;GACT;GAGA,MAAM,SAAS,IAAI,WAAW,KAAK;GACnC,IAAI,SAAS;GACb,IAAI,YAAY;GAEhB,OAAO,YAAY,KAAK,aAAa,WAAW,QAAQ;IACtD,MAAM,QAAQ,WAAW;IACzB,MAAM,SAAS,KAAK,IAAI,MAAM,QAAQ,SAAS;IAC/C,OAAO,IAAI,MAAM,SAAS,GAAG,MAAM,GAAG,MAAM;IAE5C,UAAU;IACV,aAAa;IAEb,IAAI,WAAW,MAAM,QACnB,kBAAkB;SAElB,WAAW,cAAc,MAAM,SAAS,MAAM;GAElD;GAEA,eAAe;GACf,OAAO;EACT;EAEA,IAAI;GAEF,OAAO,MAAM;IACX,MAAM,EAAE,MAAM,UAAU,MAAM,OAAO,KAAK;IAC1C,IAAI,WAAW;IACf,IAAI,MAAM;IAGV,IAAI,CAAC,OAAO;IAGZ,IAAI,cAAc,MAAM,SAAS,oBAC/B,MAAM,IAAI,MACR,mCAAmC,mBAAmB,OACxD;IAEF,WAAW,KAAK,KAAK;IACrB,eAAe,MAAM;IAIrB,OAAO,MAAM;KACX,MAAM,SAAS,WAAW;KAC1B,IAAI,CAAC,QAAQ;KAEb,MAAM,EAAE,MAAM,UAAU,WAAW;KAEnC,IACE,SAAS,UAAU,QACnB,SAAS,UAAU,SACnB,SAAS,UAAU,OACnB,SAAS,UAAU,OAEnB,MAAM,IAAI,MAAM,uBAAuB,MAAM;KAI/C,IAAI,SAAS,UAAU;UACjB,aAAa,GACf,MAAM,IAAI,MAAM,0CAA0C;KAAA,OAG5D,IAAI,aAAa,GACf,MAAM,IAAI,MAAM,gDAAgD;KAIpE,IAAI,SAAS,wBACX,MAAM,IAAI,MACR,4BAA4B,OAAO,cAAc,uBAAuB,EAC1E;KAGF,MAAM,YAAA,IAAgC;KACtC,IAAI,cAAc,WAAW;KAE7B,IAAI,EAAE,aAAa,YACjB,MAAM,IAAI,MACR,2CAA2C,WAAW,EACxD;KAIF,iBAAA,CAAkC;KAGlC,MAAM,UAAU,iBAAiB,MAAM;KAGvC,QAAQ,MAAR;MACE,KAAK,UAAU;OACb,IAAI;QACF,eAAe,QAAQ,YAAY,OAAO,OAAO,CAAC;OACpD,QAAQ,CAER;OACA;MAGF,KAAK,UAAU,OAAO;OACpB,MAAM,OAAO,iBAAiB,QAAQ;OACtC,IAAI,MACF,KAAK,QAAQ,OAAO;OAEtB;MACF;MAEA,KAAK,UAAU,KAAK;OAClB,MAAM,OAAO,iBAAiB,QAAQ;OACtC,mBAAmB,IAAI,QAAQ;OAC/B,IAAI,MAAM;QACR,IAAI;SACF,KAAK,MAAM;QACb,QAAQ,CAER;QACA,kBAAkB,OAAO,QAAQ;OACnC;OACA;MACF;MAEA,KAAK,UAAU,OAAO;OACpB,MAAM,OAAO,iBAAiB,QAAQ;OACtC,mBAAmB,IAAI,QAAQ;OAC/B,IAAI,MAAM;QACR,MAAM,UAAU,YAAY,OAAO,OAAO;QAC1C,KAAK,MAAM,IAAI,MAAM,OAAO,CAAC;QAC7B,kBAAkB,OAAO,QAAQ;OACnC;OACA;MACF;KACF;IACF;GACF;GAEA,IAAI,gBAAgB,GAClB,MAAM,IAAI,MAAM,4CAA4C;GAI9D,IAAI;IACF,eAAe,MAAM;GACvB,QAAQ,CAER;GAGA,kBAAkB,SAAS,SAAS;IAClC,IAAI;KACF,KAAK,MAAM;IACb,QAAQ,CAER;GACF,CAAC;GACD,kBAAkB,MAAM;EAC1B,SAAS,OAAO;GAEd,IAAI;IACF,eAAe,MAAM,KAAK;GAC5B,QAAQ,CAER;GACA,kBAAkB,SAAS,SAAS;IAClC,IAAI;KACF,KAAK,MAAM,KAAK;IAClB,QAAQ,CAER;GACF,CAAC;GACD,kBAAkB,MAAM;EAC1B,UAAU;GACR,IAAI;IACF,OAAO,YAAY;GACrB,QAAQ,CAER;GACA,cAAc;EAChB;CACF,GAAG;CAEH,OAAO;EAAE,WAAW;EAAmB,QAAQ;CAAW;AAC5D"}
{
"name": "@tanstack/start-client-core",
"version": "1.170.19",
"version": "1.170.20",
"description": "Modern and scalable routing for React applications",

@@ -94,5 +94,5 @@ "author": "Tanner Linsley",

"seroval": "^1.6.2",
"@tanstack/router-core": "1.171.19",
"@tanstack/router-core": "1.171.20",
"@tanstack/start-fn-stubs": "1.162.0",
"@tanstack/start-storage-context": "1.167.21"
"@tanstack/start-storage-context": "1.167.22"
},

@@ -99,0 +99,0 @@ "devDependencies": {

@@ -136,4 +136,22 @@ /**

const bufferList: Array<Uint8Array> = []
// Index of the first un-consumed chunk in bufferList. Advancing this
// pointer is O(1); using bufferList.shift() to drop a consumed chunk is
// O(n) and degrades to O(n^2) when a single large frame is assembled from
// many small chunks (e.g. a big RawStream payload split across reads).
let bufferHead = 0
let totalLength = 0
function advanceBufferHead(): void {
bufferList[bufferHead++] = EMPTY_BUFFER
// Reset drained buffers immediately and compact long-lived buffers in batches.
if (bufferHead === bufferList.length) {
bufferList.length = 0
bufferHead = 0
} else if (bufferHead >= 32) {
bufferList.splice(0, bufferHead)
bufferHead = 0
}
}
/**

@@ -150,3 +168,3 @@ * Reads header bytes from buffer chunks without flattening.

const first = bufferList[0]!
const first = bufferList[bufferHead]!

@@ -175,3 +193,3 @@ // Fast path: header fits entirely in first chunk (common case)

let remaining = FRAME_HEADER_SIZE
for (let i = 0; i < bufferList.length && remaining > 0; i++) {
for (let i = bufferHead; i < bufferList.length && remaining > 0; i++) {
const chunk = bufferList[i]!

@@ -213,9 +231,9 @@ const toCopy = Math.min(chunk.length, remaining)

// being read from the network.
const first = bufferList[0]
const first = bufferList[bufferHead]
if (first && first.length >= count) {
const result = first.subarray(0, count)
if (first.length === count) {
bufferList.shift()
advanceBufferHead()
} else {
bufferList[0] = first.subarray(count)
bufferList[bufferHead] = first.subarray(count)
}

@@ -231,5 +249,4 @@ totalLength -= count

while (remaining > 0 && bufferList.length > 0) {
const chunk = bufferList[0]
if (!chunk) break
while (remaining > 0 && bufferHead < bufferList.length) {
const chunk = bufferList[bufferHead]!
const toCopy = Math.min(chunk.length, remaining)

@@ -242,5 +259,5 @@ result.set(chunk.subarray(0, toCopy), offset)

if (toCopy === chunk.length) {
bufferList.shift()
advanceBufferHead()
} else {
bufferList[0] = chunk.subarray(toCopy)
bufferList[bufferHead] = chunk.subarray(toCopy)
}

@@ -247,0 +264,0 @@ }