@tanstack/start-client-core
Advanced tools
@@ -13,5 +13,5 @@ /** | ||
| /** Gets or creates a raw stream by ID (for use by deserialize plugin) */ | ||
| getOrCreateStream: (id: number) => ReadableStream<Uint8Array>; | ||
| getStream: (id: number) => ReadableStream<Uint8Array>; | ||
| /** Stream of JSON strings (NDJSON lines) */ | ||
| jsonChunks: ReadableStream<string>; | ||
| chunks: ReadableStream<string>; | ||
| } | ||
@@ -18,0 +18,0 @@ /** |
@@ -232,4 +232,4 @@ import { FrameType } from "../constants.js"; | ||
| return { | ||
| getOrCreateStream, | ||
| jsonChunks | ||
| getStream: getOrCreateStream, | ||
| chunks: jsonChunks | ||
| }; | ||
@@ -236,0 +236,0 @@ } |
@@ -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 getOrCreateStream: (id: number) => ReadableStream<Uint8Array>\n /** Stream of JSON strings (NDJSON lines) */\n jsonChunks: 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 { getOrCreateStream, 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;EAAmB;CAAW;AACzC"} | ||
| {"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"} |
@@ -146,7 +146,7 @@ import { TSS_CONTENT_TYPE_FRAMED, TSS_FORMDATA_CONTEXT, validateFramedProtocolVersion } from "../constants.js"; | ||
| if (!response.body) throw new Error("No response body for framed response"); | ||
| const { getOrCreateStream, jsonChunks } = createFrameDecoder(response.body); | ||
| const plugins = [createRawStreamDeserializePlugin(getOrCreateStream), ...serovalPlugins || []]; | ||
| const { getStream, chunks } = createFrameDecoder(response.body); | ||
| const plugins = [createRawStreamDeserializePlugin(getStream), ...serovalPlugins || []]; | ||
| const refs = /* @__PURE__ */ new Map(); | ||
| result = await processFramedResponse({ | ||
| jsonStream: jsonChunks, | ||
| jsonStream: chunks, | ||
| onMessage: (msg) => fromCrossJSON(msg, { | ||
@@ -153,0 +153,0 @@ refs, |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"serverFnFetcher.js","names":[],"sources":["../../../src/client-rpc/serverFnFetcher.ts"],"sourcesContent":["import {\n createRawStreamDeserializePlugin,\n encode,\n invariant,\n isNotFound,\n parseRedirect,\n} from '@tanstack/router-core'\nimport { fromCrossJSON, toJSONAsync } from 'seroval'\nimport { getDefaultSerovalPlugins } from '../getDefaultSerovalPlugins'\nimport {\n TSS_CONTENT_TYPE_FRAMED,\n TSS_FORMDATA_CONTEXT,\n X_TSS_RAW_RESPONSE,\n X_TSS_SERIALIZED,\n validateFramedProtocolVersion,\n} from '../constants'\nimport { createFrameDecoder } from './frame-decoder'\nimport type { FunctionMiddlewareClientFnOptions } from '../createMiddleware'\nimport type { Plugin as SerovalPlugin } from 'seroval'\n\nlet serovalPlugins: Array<SerovalPlugin<any, any>> | null = null\n\n/**\n * Current async post-processing context for deserialization.\n *\n * Some deserializers need to perform async work after synchronous deserialization\n * (e.g., decoding RSC payloads, fetching remote data). This context allows them\n * to register promises that must complete before the deserialized value is used.\n *\n * This uses a synchronous execution context pattern:\n * - Each call to `fromCrossJSON` is synchronous\n * - Within that synchronous execution, all `fromSerializable` calls happen\n * - We set the context before `fromCrossJSON`, clear it after\n * - For streaming chunks, we set/clear context around each `onMessage` call\n *\n * Even with concurrent server function calls, each individual deserialization\n * is atomic (synchronous), so promises are correctly scoped to their call.\n */\nlet currentPostProcessContext: Array<Promise<unknown>> | null = null\n\n/**\n * Set the current post-processing context for async deserialization work.\n * Called before deserialization starts.\n *\n * @param ctx - Array to collect async work promises, or null to clear\n */\nexport function setPostProcessContext(\n ctx: Array<Promise<unknown>> | null,\n): void {\n currentPostProcessContext = ctx\n}\n\n/**\n * Get the current post-processing context.\n * Returns null if no deserialization is in progress.\n */\nexport function getPostProcessContext(): Array<Promise<unknown>> | null {\n return currentPostProcessContext\n}\n\n/**\n * Track an async post-processing promise in the current deserialization context.\n * Called by deserializers that need to perform async work after sync deserialization.\n *\n * If no context is active (e.g., on server), this is a no-op.\n *\n * @param promise - The async work promise to track\n */\nexport function trackPostProcessPromise(promise: Promise<unknown>): void {\n if (currentPostProcessContext) {\n currentPostProcessContext.push(promise)\n }\n}\n\n/**\n * Helper to await all post-processing promises.\n * Uses Promise.allSettled to ensure all promises complete even if some reject.\n */\nasync function awaitPostProcessPromises(\n promises: Array<Promise<unknown>>,\n): Promise<void> {\n if (promises.length > 0) {\n await Promise.allSettled(promises)\n }\n}\n\n/**\n * Checks if an object has at least one own enumerable property.\n * More efficient than Object.keys(obj).length > 0 as it short-circuits on first property.\n */\nconst hop = Object.prototype.hasOwnProperty\nfunction hasOwnProperties(obj: object): boolean {\n for (const _ in obj) {\n if (hop.call(obj, _)) {\n return true\n }\n }\n return false\n}\n// caller =>\n// serverFnFetcher =>\n// client =>\n// server =>\n// fn =>\n// seroval =>\n// client middleware =>\n// serverFnFetcher =>\n// caller\n\nexport async function serverFnFetcher(\n url: string,\n args: Array<any>,\n handler: (url: string, requestInit: RequestInit) => Promise<Response>,\n) {\n if (!serovalPlugins) {\n serovalPlugins = getDefaultSerovalPlugins()\n }\n const _first = args[0]\n\n const first = _first as FunctionMiddlewareClientFnOptions<any, any, any> & {\n headers?: HeadersInit\n }\n\n // Use custom fetch if provided, otherwise fall back to the passed handler (global fetch)\n const fetchImpl = first.fetch ?? handler\n\n const type = first.data instanceof FormData ? 'formData' : 'payload'\n\n // Arrange the headers\n const headers = first.headers ? new Headers(first.headers) : new Headers()\n headers.set('x-tsr-serverFn', 'true')\n\n if (type === 'payload') {\n headers.set(\n 'accept',\n `${TSS_CONTENT_TYPE_FRAMED}, application/x-ndjson, application/json`,\n )\n }\n\n // If the method is GET, we need to move the payload to the query string\n if (first.method === 'GET') {\n if (type === 'formData') {\n throw new Error('FormData is not supported with GET requests')\n }\n const serializedPayload = await serializePayload(first)\n if (serializedPayload !== undefined) {\n const encodedPayload = encode({\n payload: serializedPayload,\n })\n if (url.includes('?')) {\n url += `&${encodedPayload}`\n } else {\n url += `?${encodedPayload}`\n }\n }\n }\n\n let body = undefined\n if (first.method === 'POST') {\n body = await getFetchBody(first)\n if (typeof body === 'string') {\n headers.set('content-type', 'application/json')\n }\n }\n\n return await getResponse(async () =>\n fetchImpl(url, {\n method: first.method,\n headers,\n signal: first.signal,\n body,\n }),\n )\n}\n\nasync function serializePayload(\n opts: FunctionMiddlewareClientFnOptions<any, any, any>,\n): Promise<string | undefined> {\n let payloadAvailable = false\n const payloadToSerialize: any = {}\n if (opts.data !== undefined) {\n payloadAvailable = true\n payloadToSerialize['data'] = opts.data\n }\n\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (opts.context && hasOwnProperties(opts.context)) {\n payloadAvailable = true\n payloadToSerialize['context'] = opts.context\n }\n\n if (payloadAvailable) {\n return serialize(payloadToSerialize)\n }\n return undefined\n}\n\nasync function serialize(data: any) {\n return JSON.stringify(\n await Promise.resolve(toJSONAsync(data, { plugins: serovalPlugins! })),\n )\n}\n\nasync function getFetchBody(\n opts: FunctionMiddlewareClientFnOptions<any, any, any>,\n): Promise<FormData | string | undefined> {\n if (opts.data instanceof FormData) {\n let serializedContext = undefined\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (opts.context && hasOwnProperties(opts.context)) {\n serializedContext = await serialize(opts.context)\n }\n if (serializedContext !== undefined) {\n opts.data.set(TSS_FORMDATA_CONTEXT, serializedContext)\n }\n return opts.data\n }\n const serializedBody = await serializePayload(opts)\n if (serializedBody) {\n return serializedBody\n }\n return undefined\n}\n\n/**\n * Retrieves a response from a given function and manages potential errors\n * and special response types including redirects and not found errors.\n *\n * @param fn - The function to execute for obtaining the response.\n * @returns The processed response from the function.\n * @throws If the response is invalid or an error occurs during processing.\n */\nasync function getResponse(fn: () => Promise<Response>) {\n let response: Response\n try {\n response = await fn() // client => server => fn => server => client\n } catch (error) {\n if (error instanceof Response) {\n response = error\n } else {\n console.log(error)\n throw error\n }\n }\n\n if (response.headers.get(X_TSS_RAW_RESPONSE) === 'true') {\n return response\n }\n\n const contentType = response.headers.get('content-type')\n if (!contentType) {\n if (process.env.NODE_ENV !== 'production') {\n throw new Error(\n 'Invariant failed: expected content-type header to be set',\n )\n }\n\n invariant()\n }\n const serializedByStart = !!response.headers.get(X_TSS_SERIALIZED)\n\n // If the response is serialized by the start server, we need to process it\n // differently than a normal response.\n if (serializedByStart) {\n let result\n\n // If it's a framed response (contains RawStream), use frame decoder\n if (contentType.includes(TSS_CONTENT_TYPE_FRAMED)) {\n // Validate protocol version compatibility\n validateFramedProtocolVersion(contentType)\n\n if (!response.body) {\n throw new Error('No response body for framed response')\n }\n\n const { getOrCreateStream, jsonChunks } = createFrameDecoder(\n response.body,\n )\n\n // Create deserialize plugin that wires up the raw streams\n const rawStreamPlugin =\n createRawStreamDeserializePlugin(getOrCreateStream)\n const plugins = [rawStreamPlugin, ...(serovalPlugins || [])]\n\n const refs = new Map()\n result = await processFramedResponse({\n jsonStream: jsonChunks,\n onMessage: (msg: any) => fromCrossJSON(msg, { refs, plugins }),\n onError(msg, error) {\n console.error(msg, error)\n },\n })\n }\n // If it's a JSON response, it can be simpler\n else if (contentType.includes('application/json')) {\n const jsonPayload = await response.json()\n // Track async post-processing work for this deserialization\n const postProcessPromises: Array<Promise<unknown>> = []\n setPostProcessContext(postProcessPromises)\n try {\n result = fromCrossJSON(jsonPayload, { plugins: serovalPlugins! })\n } finally {\n setPostProcessContext(null)\n }\n // Await any async post-processing before returning\n await awaitPostProcessPromises(postProcessPromises)\n }\n\n if (!result) {\n if (process.env.NODE_ENV !== 'production') {\n throw new Error('Invariant failed: expected result to be resolved')\n }\n\n invariant()\n }\n if (result instanceof Error) {\n throw result\n }\n\n return result\n }\n\n // If it wasn't processed by the start serializer, check\n // if it's JSON\n if (contentType.includes('application/json')) {\n const jsonPayload = await response.json()\n const redirect = parseRedirect(jsonPayload)\n if (redirect) {\n throw redirect\n }\n if (isNotFound(jsonPayload)) {\n throw jsonPayload\n }\n return jsonPayload\n }\n\n // Otherwise, if it's not OK, throw the content\n if (!response.ok) {\n throw new Error(await response.text())\n }\n\n // Or return the response itself\n return response\n}\n\n/**\n * Processes a framed response where each JSON chunk is a complete JSON string\n * (already decoded by frame decoder).\n *\n * Uses per-chunk post-processing context to ensure async deserialization work\n * completes before the next chunk is processed. This prevents issues when\n * streaming values require async post-processing (e.g., RSC decoding).\n */\nasync function processFramedResponse({\n jsonStream,\n onMessage,\n onError,\n}: {\n jsonStream: ReadableStream<string>\n onMessage: (msg: any) => any\n onError?: (msg: string, error?: any) => void\n}) {\n const reader = jsonStream.getReader()\n\n // Read first JSON frame - this is the main result\n const { value: firstValue, done: firstDone } = await reader.read()\n if (firstDone || !firstValue) {\n throw new Error('Stream ended before first object')\n }\n\n // Each frame is a complete JSON string\n const firstObject = JSON.parse(firstValue)\n\n // Process remaining frames for streaming refs like RawStream.\n // Keep draining until the server closes the stream.\n // Each chunk gets its own post-processing context to properly scope async work.\n let drainCancelled = false as boolean\n const drain = (async () => {\n try {\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n while (true) {\n const { value, done } = await reader.read()\n if (done) break\n if (value) {\n try {\n // Set up post-processing context for this chunk\n const chunkPostProcessPromises: Array<Promise<unknown>> = []\n setPostProcessContext(chunkPostProcessPromises)\n try {\n onMessage(JSON.parse(value))\n } finally {\n setPostProcessContext(null)\n }\n // Await any async post-processing from this chunk before processing next.\n // This ensures values requiring async work are ready before their\n // containing Promise/Stream resolves/emits to consumers.\n await awaitPostProcessPromises(chunkPostProcessPromises)\n } catch (e) {\n onError?.(`Invalid JSON: ${value}`, e)\n }\n }\n }\n } catch (err) {\n if (!drainCancelled) {\n onError?.('Stream processing error:', err)\n }\n }\n })()\n\n // Process first object with its own post-processing context\n let result: any\n const initialPostProcessPromises: Array<Promise<unknown>> = []\n setPostProcessContext(initialPostProcessPromises)\n try {\n result = onMessage(firstObject)\n } catch (err) {\n setPostProcessContext(null)\n drainCancelled = true\n reader.cancel().catch(() => {})\n throw err\n }\n setPostProcessContext(null)\n\n // Await initial post-processing promises before returning result\n await awaitPostProcessPromises(initialPostProcessPromises)\n\n // If the initial decode fails async, stop draining to avoid holding\n // onto the response body and raw stream buffers unnecessarily.\n Promise.resolve(result).catch(() => {\n drainCancelled = true\n reader.cancel().catch(() => {})\n })\n\n // Detach reader once draining completes.\n drain.finally(() => {\n try {\n reader.releaseLock()\n } catch {\n // Ignore\n }\n })\n\n return result\n}\n"],"mappings":";;;;;;AAoBA,IAAI,iBAAwD;;;;;;;;;;;;;;;;;AAkB5D,IAAI,4BAA4D;;;;;;;AAQhE,SAAgB,sBACd,KACM;CACN,4BAA4B;AAC9B;;;;;;;;;AAkBA,SAAgB,wBAAwB,SAAiC;CACvE,IAAI,2BACF,0BAA0B,KAAK,OAAO;AAE1C;;;;;AAMA,eAAe,yBACb,UACe;CACf,IAAI,SAAS,SAAS,GACpB,MAAM,QAAQ,WAAW,QAAQ;AAErC;;;;;AAMA,IAAM,MAAM,OAAO,UAAU;AAC7B,SAAS,iBAAiB,KAAsB;CAC9C,KAAK,MAAM,KAAK,KACd,IAAI,IAAI,KAAK,KAAK,CAAC,GACjB,OAAO;CAGX,OAAO;AACT;AAWA,eAAsB,gBACpB,KACA,MACA,SACA;CACA,IAAI,CAAC,gBACH,iBAAiB,yBAAyB;CAI5C,MAAM,QAFS,KAAK;CAOpB,MAAM,YAAY,MAAM,SAAS;CAEjC,MAAM,OAAO,MAAM,gBAAgB,WAAW,aAAa;CAG3D,MAAM,UAAU,MAAM,UAAU,IAAI,QAAQ,MAAM,OAAO,IAAI,IAAI,QAAQ;CACzE,QAAQ,IAAI,kBAAkB,MAAM;CAEpC,IAAI,SAAS,WACX,QAAQ,IACN,UACA,GAAG,wBAAwB,yCAC7B;CAIF,IAAI,MAAM,WAAW,OAAO;EAC1B,IAAI,SAAS,YACX,MAAM,IAAI,MAAM,6CAA6C;EAE/D,MAAM,oBAAoB,MAAM,iBAAiB,KAAK;EACtD,IAAI,sBAAsB,KAAA,GAAW;GACnC,MAAM,iBAAiB,OAAO,EAC5B,SAAS,kBACX,CAAC;GACD,IAAI,IAAI,SAAS,GAAG,GAClB,OAAO,IAAI;QAEX,OAAO,IAAI;EAEf;CACF;CAEA,IAAI,OAAO,KAAA;CACX,IAAI,MAAM,WAAW,QAAQ;EAC3B,OAAO,MAAM,aAAa,KAAK;EAC/B,IAAI,OAAO,SAAS,UAClB,QAAQ,IAAI,gBAAgB,kBAAkB;CAElD;CAEA,OAAO,MAAM,YAAY,YACvB,UAAU,KAAK;EACb,QAAQ,MAAM;EACd;EACA,QAAQ,MAAM;EACd;CACF,CAAC,CACH;AACF;AAEA,eAAe,iBACb,MAC6B;CAC7B,IAAI,mBAAmB;CACvB,MAAM,qBAA0B,CAAC;CACjC,IAAI,KAAK,SAAS,KAAA,GAAW;EAC3B,mBAAmB;EACnB,mBAAmB,UAAU,KAAK;CACpC;CAGA,IAAI,KAAK,WAAW,iBAAiB,KAAK,OAAO,GAAG;EAClD,mBAAmB;EACnB,mBAAmB,aAAa,KAAK;CACvC;CAEA,IAAI,kBACF,OAAO,UAAU,kBAAkB;AAGvC;AAEA,eAAe,UAAU,MAAW;CAClC,OAAO,KAAK,UACV,MAAM,QAAQ,QAAQ,YAAY,MAAM,EAAE,SAAS,eAAgB,CAAC,CAAC,CACvE;AACF;AAEA,eAAe,aACb,MACwC;CACxC,IAAI,KAAK,gBAAgB,UAAU;EACjC,IAAI,oBAAoB,KAAA;EAExB,IAAI,KAAK,WAAW,iBAAiB,KAAK,OAAO,GAC/C,oBAAoB,MAAM,UAAU,KAAK,OAAO;EAElD,IAAI,sBAAsB,KAAA,GACxB,KAAK,KAAK,IAAI,sBAAsB,iBAAiB;EAEvD,OAAO,KAAK;CACd;CACA,MAAM,iBAAiB,MAAM,iBAAiB,IAAI;CAClD,IAAI,gBACF,OAAO;AAGX;;;;;;;;;AAUA,eAAe,YAAY,IAA6B;CACtD,IAAI;CACJ,IAAI;EACF,WAAW,MAAM,GAAG;CACtB,SAAS,OAAO;EACd,IAAI,iBAAiB,UACnB,WAAW;OACN;GACL,QAAQ,IAAI,KAAK;GACjB,MAAM;EACR;CACF;CAEA,IAAI,SAAS,QAAQ,IAAA,WAAsB,MAAM,QAC/C,OAAO;CAGT,MAAM,cAAc,SAAS,QAAQ,IAAI,cAAc;CACvD,IAAI,CAAC,aAAa;EAChB,IAAA,QAAA,IAAA,aAA6B,cAC3B,MAAM,IAAI,MACR,0DACF;EAGF,UAAU;CACZ;CAKA,IAAI,CAJuB,CAAC,SAAS,QAAQ,IAAA,kBAAoB,GAI1C;EACrB,IAAI;EAGJ,IAAI,YAAY,SAAA,0BAAgC,GAAG;GAEjD,8BAA8B,WAAW;GAEzC,IAAI,CAAC,SAAS,MACZ,MAAM,IAAI,MAAM,sCAAsC;GAGxD,MAAM,EAAE,mBAAmB,eAAe,mBACxC,SAAS,IACX;GAKA,MAAM,UAAU,CADd,iCAAiC,iBAClB,GAAiB,GAAI,kBAAkB,CAAC,CAAE;GAE3D,MAAM,uBAAO,IAAI,IAAI;GACrB,SAAS,MAAM,sBAAsB;IACnC,YAAY;IACZ,YAAY,QAAa,cAAc,KAAK;KAAE;KAAM;IAAQ,CAAC;IAC7D,QAAQ,KAAK,OAAO;KAClB,QAAQ,MAAM,KAAK,KAAK;IAC1B;GACF,CAAC;EACH,OAEK,IAAI,YAAY,SAAS,kBAAkB,GAAG;GACjD,MAAM,cAAc,MAAM,SAAS,KAAK;GAExC,MAAM,sBAA+C,CAAC;GACtD,sBAAsB,mBAAmB;GACzC,IAAI;IACF,SAAS,cAAc,aAAa,EAAE,SAAS,eAAgB,CAAC;GAClE,UAAU;IACR,sBAAsB,IAAI;GAC5B;GAEA,MAAM,yBAAyB,mBAAmB;EACpD;EAEA,IAAI,CAAC,QAAQ;GACX,IAAA,QAAA,IAAA,aAA6B,cAC3B,MAAM,IAAI,MAAM,kDAAkD;GAGpE,UAAU;EACZ;EACA,IAAI,kBAAkB,OACpB,MAAM;EAGR,OAAO;CACT;CAIA,IAAI,YAAY,SAAS,kBAAkB,GAAG;EAC5C,MAAM,cAAc,MAAM,SAAS,KAAK;EACxC,MAAM,WAAW,cAAc,WAAW;EAC1C,IAAI,UACF,MAAM;EAER,IAAI,WAAW,WAAW,GACxB,MAAM;EAER,OAAO;CACT;CAGA,IAAI,CAAC,SAAS,IACZ,MAAM,IAAI,MAAM,MAAM,SAAS,KAAK,CAAC;CAIvC,OAAO;AACT;;;;;;;;;AAUA,eAAe,sBAAsB,EACnC,YACA,WACA,WAKC;CACD,MAAM,SAAS,WAAW,UAAU;CAGpC,MAAM,EAAE,OAAO,YAAY,MAAM,cAAc,MAAM,OAAO,KAAK;CACjE,IAAI,aAAa,CAAC,YAChB,MAAM,IAAI,MAAM,kCAAkC;CAIpD,MAAM,cAAc,KAAK,MAAM,UAAU;CAKzC,IAAI,iBAAiB;CACrB,MAAM,SAAS,YAAY;EACzB,IAAI;GAEF,OAAO,MAAM;IACX,MAAM,EAAE,OAAO,SAAS,MAAM,OAAO,KAAK;IAC1C,IAAI,MAAM;IACV,IAAI,OACF,IAAI;KAEF,MAAM,2BAAoD,CAAC;KAC3D,sBAAsB,wBAAwB;KAC9C,IAAI;MACF,UAAU,KAAK,MAAM,KAAK,CAAC;KAC7B,UAAU;MACR,sBAAsB,IAAI;KAC5B;KAIA,MAAM,yBAAyB,wBAAwB;IACzD,SAAS,GAAG;KACV,UAAU,iBAAiB,SAAS,CAAC;IACvC;GAEJ;EACF,SAAS,KAAK;GACZ,IAAI,CAAC,gBACH,UAAU,4BAA4B,GAAG;EAE7C;CACF,GAAG;CAGH,IAAI;CACJ,MAAM,6BAAsD,CAAC;CAC7D,sBAAsB,0BAA0B;CAChD,IAAI;EACF,SAAS,UAAU,WAAW;CAChC,SAAS,KAAK;EACZ,sBAAsB,IAAI;EAC1B,iBAAiB;EACjB,OAAO,OAAO,EAAE,YAAY,CAAC,CAAC;EAC9B,MAAM;CACR;CACA,sBAAsB,IAAI;CAG1B,MAAM,yBAAyB,0BAA0B;CAIzD,QAAQ,QAAQ,MAAM,EAAE,YAAY;EAClC,iBAAiB;EACjB,OAAO,OAAO,EAAE,YAAY,CAAC,CAAC;CAChC,CAAC;CAGD,MAAM,cAAc;EAClB,IAAI;GACF,OAAO,YAAY;EACrB,QAAQ,CAER;CACF,CAAC;CAED,OAAO;AACT"} | ||
| {"version":3,"file":"serverFnFetcher.js","names":[],"sources":["../../../src/client-rpc/serverFnFetcher.ts"],"sourcesContent":["import {\n createRawStreamDeserializePlugin,\n encode,\n invariant,\n isNotFound,\n parseRedirect,\n} from '@tanstack/router-core'\nimport { fromCrossJSON, toJSONAsync } from 'seroval'\nimport { getDefaultSerovalPlugins } from '../getDefaultSerovalPlugins'\nimport {\n TSS_CONTENT_TYPE_FRAMED,\n TSS_FORMDATA_CONTEXT,\n X_TSS_RAW_RESPONSE,\n X_TSS_SERIALIZED,\n validateFramedProtocolVersion,\n} from '../constants'\nimport { createFrameDecoder } from './frame-decoder'\nimport type { FunctionMiddlewareClientFnOptions } from '../createMiddleware'\nimport type { Plugin as SerovalPlugin } from 'seroval'\n\nlet serovalPlugins: Array<SerovalPlugin<any, any>> | null = null\n\n/**\n * Current async post-processing context for deserialization.\n *\n * Some deserializers need to perform async work after synchronous deserialization\n * (e.g., decoding RSC payloads, fetching remote data). This context allows them\n * to register promises that must complete before the deserialized value is used.\n *\n * This uses a synchronous execution context pattern:\n * - Each call to `fromCrossJSON` is synchronous\n * - Within that synchronous execution, all `fromSerializable` calls happen\n * - We set the context before `fromCrossJSON`, clear it after\n * - For streaming chunks, we set/clear context around each `onMessage` call\n *\n * Even with concurrent server function calls, each individual deserialization\n * is atomic (synchronous), so promises are correctly scoped to their call.\n */\nlet currentPostProcessContext: Array<Promise<unknown>> | null = null\n\n/**\n * Set the current post-processing context for async deserialization work.\n * Called before deserialization starts.\n *\n * @param ctx - Array to collect async work promises, or null to clear\n */\nexport function setPostProcessContext(\n ctx: Array<Promise<unknown>> | null,\n): void {\n currentPostProcessContext = ctx\n}\n\n/**\n * Get the current post-processing context.\n * Returns null if no deserialization is in progress.\n */\nexport function getPostProcessContext(): Array<Promise<unknown>> | null {\n return currentPostProcessContext\n}\n\n/**\n * Track an async post-processing promise in the current deserialization context.\n * Called by deserializers that need to perform async work after sync deserialization.\n *\n * If no context is active (e.g., on server), this is a no-op.\n *\n * @param promise - The async work promise to track\n */\nexport function trackPostProcessPromise(promise: Promise<unknown>): void {\n if (currentPostProcessContext) {\n currentPostProcessContext.push(promise)\n }\n}\n\n/**\n * Helper to await all post-processing promises.\n * Uses Promise.allSettled to ensure all promises complete even if some reject.\n */\nasync function awaitPostProcessPromises(\n promises: Array<Promise<unknown>>,\n): Promise<void> {\n if (promises.length > 0) {\n await Promise.allSettled(promises)\n }\n}\n\n/**\n * Checks if an object has at least one own enumerable property.\n * More efficient than Object.keys(obj).length > 0 as it short-circuits on first property.\n */\nconst hop = Object.prototype.hasOwnProperty\nfunction hasOwnProperties(obj: object): boolean {\n for (const _ in obj) {\n if (hop.call(obj, _)) {\n return true\n }\n }\n return false\n}\n// caller =>\n// serverFnFetcher =>\n// client =>\n// server =>\n// fn =>\n// seroval =>\n// client middleware =>\n// serverFnFetcher =>\n// caller\n\nexport async function serverFnFetcher(\n url: string,\n args: Array<any>,\n handler: (url: string, requestInit: RequestInit) => Promise<Response>,\n) {\n if (!serovalPlugins) {\n serovalPlugins = getDefaultSerovalPlugins()\n }\n const _first = args[0]\n\n const first = _first as FunctionMiddlewareClientFnOptions<any, any, any> & {\n headers?: HeadersInit\n }\n\n // Use custom fetch if provided, otherwise fall back to the passed handler (global fetch)\n const fetchImpl = first.fetch ?? handler\n\n const type = first.data instanceof FormData ? 'formData' : 'payload'\n\n // Arrange the headers\n const headers = first.headers ? new Headers(first.headers) : new Headers()\n headers.set('x-tsr-serverFn', 'true')\n\n if (type === 'payload') {\n headers.set(\n 'accept',\n `${TSS_CONTENT_TYPE_FRAMED}, application/x-ndjson, application/json`,\n )\n }\n\n // If the method is GET, we need to move the payload to the query string\n if (first.method === 'GET') {\n if (type === 'formData') {\n throw new Error('FormData is not supported with GET requests')\n }\n const serializedPayload = await serializePayload(first)\n if (serializedPayload !== undefined) {\n const encodedPayload = encode({\n payload: serializedPayload,\n })\n if (url.includes('?')) {\n url += `&${encodedPayload}`\n } else {\n url += `?${encodedPayload}`\n }\n }\n }\n\n let body = undefined\n if (first.method === 'POST') {\n body = await getFetchBody(first)\n if (typeof body === 'string') {\n headers.set('content-type', 'application/json')\n }\n }\n\n return await getResponse(async () =>\n fetchImpl(url, {\n method: first.method,\n headers,\n signal: first.signal,\n body,\n }),\n )\n}\n\nasync function serializePayload(\n opts: FunctionMiddlewareClientFnOptions<any, any, any>,\n): Promise<string | undefined> {\n let payloadAvailable = false\n const payloadToSerialize: any = {}\n if (opts.data !== undefined) {\n payloadAvailable = true\n payloadToSerialize['data'] = opts.data\n }\n\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (opts.context && hasOwnProperties(opts.context)) {\n payloadAvailable = true\n payloadToSerialize['context'] = opts.context\n }\n\n if (payloadAvailable) {\n return serialize(payloadToSerialize)\n }\n return undefined\n}\n\nasync function serialize(data: any) {\n return JSON.stringify(\n await Promise.resolve(toJSONAsync(data, { plugins: serovalPlugins! })),\n )\n}\n\nasync function getFetchBody(\n opts: FunctionMiddlewareClientFnOptions<any, any, any>,\n): Promise<FormData | string | undefined> {\n if (opts.data instanceof FormData) {\n let serializedContext = undefined\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (opts.context && hasOwnProperties(opts.context)) {\n serializedContext = await serialize(opts.context)\n }\n if (serializedContext !== undefined) {\n opts.data.set(TSS_FORMDATA_CONTEXT, serializedContext)\n }\n return opts.data\n }\n const serializedBody = await serializePayload(opts)\n if (serializedBody) {\n return serializedBody\n }\n return undefined\n}\n\n/**\n * Retrieves a response from a given function and manages potential errors\n * and special response types including redirects and not found errors.\n *\n * @param fn - The function to execute for obtaining the response.\n * @returns The processed response from the function.\n * @throws If the response is invalid or an error occurs during processing.\n */\nasync function getResponse(fn: () => Promise<Response>) {\n let response: Response\n try {\n response = await fn() // client => server => fn => server => client\n } catch (error) {\n if (error instanceof Response) {\n response = error\n } else {\n console.log(error)\n throw error\n }\n }\n\n if (response.headers.get(X_TSS_RAW_RESPONSE) === 'true') {\n return response\n }\n\n const contentType = response.headers.get('content-type')\n if (!contentType) {\n if (process.env.NODE_ENV !== 'production') {\n throw new Error(\n 'Invariant failed: expected content-type header to be set',\n )\n }\n\n invariant()\n }\n const serializedByStart = !!response.headers.get(X_TSS_SERIALIZED)\n\n // If the response is serialized by the start server, we need to process it\n // differently than a normal response.\n if (serializedByStart) {\n let result\n\n // If it's a framed response (contains RawStream), use frame decoder\n if (contentType.includes(TSS_CONTENT_TYPE_FRAMED)) {\n // Validate protocol version compatibility\n validateFramedProtocolVersion(contentType)\n\n if (!response.body) {\n throw new Error('No response body for framed response')\n }\n\n const { getStream, chunks } = createFrameDecoder(response.body)\n\n // Create deserialize plugin that wires up the raw streams\n const rawStreamPlugin = createRawStreamDeserializePlugin(getStream)\n const plugins = [rawStreamPlugin, ...(serovalPlugins || [])]\n\n const refs = new Map()\n result = await processFramedResponse({\n jsonStream: chunks,\n onMessage: (msg: any) => fromCrossJSON(msg, { refs, plugins }),\n onError(msg, error) {\n console.error(msg, error)\n },\n })\n }\n // If it's a JSON response, it can be simpler\n else if (contentType.includes('application/json')) {\n const jsonPayload = await response.json()\n // Track async post-processing work for this deserialization\n const postProcessPromises: Array<Promise<unknown>> = []\n setPostProcessContext(postProcessPromises)\n try {\n result = fromCrossJSON(jsonPayload, { plugins: serovalPlugins! })\n } finally {\n setPostProcessContext(null)\n }\n // Await any async post-processing before returning\n await awaitPostProcessPromises(postProcessPromises)\n }\n\n if (!result) {\n if (process.env.NODE_ENV !== 'production') {\n throw new Error('Invariant failed: expected result to be resolved')\n }\n\n invariant()\n }\n if (result instanceof Error) {\n throw result\n }\n\n return result\n }\n\n // If it wasn't processed by the start serializer, check\n // if it's JSON\n if (contentType.includes('application/json')) {\n const jsonPayload = await response.json()\n const redirect = parseRedirect(jsonPayload)\n if (redirect) {\n throw redirect\n }\n if (isNotFound(jsonPayload)) {\n throw jsonPayload\n }\n return jsonPayload\n }\n\n // Otherwise, if it's not OK, throw the content\n if (!response.ok) {\n throw new Error(await response.text())\n }\n\n // Or return the response itself\n return response\n}\n\n/**\n * Processes a framed response where each JSON chunk is a complete JSON string\n * (already decoded by frame decoder).\n *\n * Uses per-chunk post-processing context to ensure async deserialization work\n * completes before the next chunk is processed. This prevents issues when\n * streaming values require async post-processing (e.g., RSC decoding).\n */\nasync function processFramedResponse({\n jsonStream,\n onMessage,\n onError,\n}: {\n jsonStream: ReadableStream<string>\n onMessage: (msg: any) => any\n onError?: (msg: string, error?: any) => void\n}) {\n const reader = jsonStream.getReader()\n\n // Read first JSON frame - this is the main result\n const { value: firstValue, done: firstDone } = await reader.read()\n if (firstDone || !firstValue) {\n throw new Error('Stream ended before first object')\n }\n\n // Each frame is a complete JSON string\n const firstObject = JSON.parse(firstValue)\n\n // Process remaining frames for streaming refs like RawStream.\n // Keep draining until the server closes the stream.\n // Each chunk gets its own post-processing context to properly scope async work.\n let drainCancelled = false as boolean\n const drain = (async () => {\n try {\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n while (true) {\n const { value, done } = await reader.read()\n if (done) break\n if (value) {\n try {\n // Set up post-processing context for this chunk\n const chunkPostProcessPromises: Array<Promise<unknown>> = []\n setPostProcessContext(chunkPostProcessPromises)\n try {\n onMessage(JSON.parse(value))\n } finally {\n setPostProcessContext(null)\n }\n // Await any async post-processing from this chunk before processing next.\n // This ensures values requiring async work are ready before their\n // containing Promise/Stream resolves/emits to consumers.\n await awaitPostProcessPromises(chunkPostProcessPromises)\n } catch (e) {\n onError?.(`Invalid JSON: ${value}`, e)\n }\n }\n }\n } catch (err) {\n if (!drainCancelled) {\n onError?.('Stream processing error:', err)\n }\n }\n })()\n\n // Process first object with its own post-processing context\n let result: any\n const initialPostProcessPromises: Array<Promise<unknown>> = []\n setPostProcessContext(initialPostProcessPromises)\n try {\n result = onMessage(firstObject)\n } catch (err) {\n setPostProcessContext(null)\n drainCancelled = true\n reader.cancel().catch(() => {})\n throw err\n }\n setPostProcessContext(null)\n\n // Await initial post-processing promises before returning result\n await awaitPostProcessPromises(initialPostProcessPromises)\n\n // If the initial decode fails async, stop draining to avoid holding\n // onto the response body and raw stream buffers unnecessarily.\n Promise.resolve(result).catch(() => {\n drainCancelled = true\n reader.cancel().catch(() => {})\n })\n\n // Detach reader once draining completes.\n drain.finally(() => {\n try {\n reader.releaseLock()\n } catch {\n // Ignore\n }\n })\n\n return result\n}\n"],"mappings":";;;;;;AAoBA,IAAI,iBAAwD;;;;;;;;;;;;;;;;;AAkB5D,IAAI,4BAA4D;;;;;;;AAQhE,SAAgB,sBACd,KACM;CACN,4BAA4B;AAC9B;;;;;;;;;AAkBA,SAAgB,wBAAwB,SAAiC;CACvE,IAAI,2BACF,0BAA0B,KAAK,OAAO;AAE1C;;;;;AAMA,eAAe,yBACb,UACe;CACf,IAAI,SAAS,SAAS,GACpB,MAAM,QAAQ,WAAW,QAAQ;AAErC;;;;;AAMA,IAAM,MAAM,OAAO,UAAU;AAC7B,SAAS,iBAAiB,KAAsB;CAC9C,KAAK,MAAM,KAAK,KACd,IAAI,IAAI,KAAK,KAAK,CAAC,GACjB,OAAO;CAGX,OAAO;AACT;AAWA,eAAsB,gBACpB,KACA,MACA,SACA;CACA,IAAI,CAAC,gBACH,iBAAiB,yBAAyB;CAI5C,MAAM,QAFS,KAAK;CAOpB,MAAM,YAAY,MAAM,SAAS;CAEjC,MAAM,OAAO,MAAM,gBAAgB,WAAW,aAAa;CAG3D,MAAM,UAAU,MAAM,UAAU,IAAI,QAAQ,MAAM,OAAO,IAAI,IAAI,QAAQ;CACzE,QAAQ,IAAI,kBAAkB,MAAM;CAEpC,IAAI,SAAS,WACX,QAAQ,IACN,UACA,GAAG,wBAAwB,yCAC7B;CAIF,IAAI,MAAM,WAAW,OAAO;EAC1B,IAAI,SAAS,YACX,MAAM,IAAI,MAAM,6CAA6C;EAE/D,MAAM,oBAAoB,MAAM,iBAAiB,KAAK;EACtD,IAAI,sBAAsB,KAAA,GAAW;GACnC,MAAM,iBAAiB,OAAO,EAC5B,SAAS,kBACX,CAAC;GACD,IAAI,IAAI,SAAS,GAAG,GAClB,OAAO,IAAI;QAEX,OAAO,IAAI;EAEf;CACF;CAEA,IAAI,OAAO,KAAA;CACX,IAAI,MAAM,WAAW,QAAQ;EAC3B,OAAO,MAAM,aAAa,KAAK;EAC/B,IAAI,OAAO,SAAS,UAClB,QAAQ,IAAI,gBAAgB,kBAAkB;CAElD;CAEA,OAAO,MAAM,YAAY,YACvB,UAAU,KAAK;EACb,QAAQ,MAAM;EACd;EACA,QAAQ,MAAM;EACd;CACF,CAAC,CACH;AACF;AAEA,eAAe,iBACb,MAC6B;CAC7B,IAAI,mBAAmB;CACvB,MAAM,qBAA0B,CAAC;CACjC,IAAI,KAAK,SAAS,KAAA,GAAW;EAC3B,mBAAmB;EACnB,mBAAmB,UAAU,KAAK;CACpC;CAGA,IAAI,KAAK,WAAW,iBAAiB,KAAK,OAAO,GAAG;EAClD,mBAAmB;EACnB,mBAAmB,aAAa,KAAK;CACvC;CAEA,IAAI,kBACF,OAAO,UAAU,kBAAkB;AAGvC;AAEA,eAAe,UAAU,MAAW;CAClC,OAAO,KAAK,UACV,MAAM,QAAQ,QAAQ,YAAY,MAAM,EAAE,SAAS,eAAgB,CAAC,CAAC,CACvE;AACF;AAEA,eAAe,aACb,MACwC;CACxC,IAAI,KAAK,gBAAgB,UAAU;EACjC,IAAI,oBAAoB,KAAA;EAExB,IAAI,KAAK,WAAW,iBAAiB,KAAK,OAAO,GAC/C,oBAAoB,MAAM,UAAU,KAAK,OAAO;EAElD,IAAI,sBAAsB,KAAA,GACxB,KAAK,KAAK,IAAI,sBAAsB,iBAAiB;EAEvD,OAAO,KAAK;CACd;CACA,MAAM,iBAAiB,MAAM,iBAAiB,IAAI;CAClD,IAAI,gBACF,OAAO;AAGX;;;;;;;;;AAUA,eAAe,YAAY,IAA6B;CACtD,IAAI;CACJ,IAAI;EACF,WAAW,MAAM,GAAG;CACtB,SAAS,OAAO;EACd,IAAI,iBAAiB,UACnB,WAAW;OACN;GACL,QAAQ,IAAI,KAAK;GACjB,MAAM;EACR;CACF;CAEA,IAAI,SAAS,QAAQ,IAAA,WAAsB,MAAM,QAC/C,OAAO;CAGT,MAAM,cAAc,SAAS,QAAQ,IAAI,cAAc;CACvD,IAAI,CAAC,aAAa;EAChB,IAAA,QAAA,IAAA,aAA6B,cAC3B,MAAM,IAAI,MACR,0DACF;EAGF,UAAU;CACZ;CAKA,IAAI,CAJuB,CAAC,SAAS,QAAQ,IAAA,kBAAoB,GAI1C;EACrB,IAAI;EAGJ,IAAI,YAAY,SAAA,0BAAgC,GAAG;GAEjD,8BAA8B,WAAW;GAEzC,IAAI,CAAC,SAAS,MACZ,MAAM,IAAI,MAAM,sCAAsC;GAGxD,MAAM,EAAE,WAAW,WAAW,mBAAmB,SAAS,IAAI;GAI9D,MAAM,UAAU,CADQ,iCAAiC,SACxC,GAAiB,GAAI,kBAAkB,CAAC,CAAE;GAE3D,MAAM,uBAAO,IAAI,IAAI;GACrB,SAAS,MAAM,sBAAsB;IACnC,YAAY;IACZ,YAAY,QAAa,cAAc,KAAK;KAAE;KAAM;IAAQ,CAAC;IAC7D,QAAQ,KAAK,OAAO;KAClB,QAAQ,MAAM,KAAK,KAAK;IAC1B;GACF,CAAC;EACH,OAEK,IAAI,YAAY,SAAS,kBAAkB,GAAG;GACjD,MAAM,cAAc,MAAM,SAAS,KAAK;GAExC,MAAM,sBAA+C,CAAC;GACtD,sBAAsB,mBAAmB;GACzC,IAAI;IACF,SAAS,cAAc,aAAa,EAAE,SAAS,eAAgB,CAAC;GAClE,UAAU;IACR,sBAAsB,IAAI;GAC5B;GAEA,MAAM,yBAAyB,mBAAmB;EACpD;EAEA,IAAI,CAAC,QAAQ;GACX,IAAA,QAAA,IAAA,aAA6B,cAC3B,MAAM,IAAI,MAAM,kDAAkD;GAGpE,UAAU;EACZ;EACA,IAAI,kBAAkB,OACpB,MAAM;EAGR,OAAO;CACT;CAIA,IAAI,YAAY,SAAS,kBAAkB,GAAG;EAC5C,MAAM,cAAc,MAAM,SAAS,KAAK;EACxC,MAAM,WAAW,cAAc,WAAW;EAC1C,IAAI,UACF,MAAM;EAER,IAAI,WAAW,WAAW,GACxB,MAAM;EAER,OAAO;CACT;CAGA,IAAI,CAAC,SAAS,IACZ,MAAM,IAAI,MAAM,MAAM,SAAS,KAAK,CAAC;CAIvC,OAAO;AACT;;;;;;;;;AAUA,eAAe,sBAAsB,EACnC,YACA,WACA,WAKC;CACD,MAAM,SAAS,WAAW,UAAU;CAGpC,MAAM,EAAE,OAAO,YAAY,MAAM,cAAc,MAAM,OAAO,KAAK;CACjE,IAAI,aAAa,CAAC,YAChB,MAAM,IAAI,MAAM,kCAAkC;CAIpD,MAAM,cAAc,KAAK,MAAM,UAAU;CAKzC,IAAI,iBAAiB;CACrB,MAAM,SAAS,YAAY;EACzB,IAAI;GAEF,OAAO,MAAM;IACX,MAAM,EAAE,OAAO,SAAS,MAAM,OAAO,KAAK;IAC1C,IAAI,MAAM;IACV,IAAI,OACF,IAAI;KAEF,MAAM,2BAAoD,CAAC;KAC3D,sBAAsB,wBAAwB;KAC9C,IAAI;MACF,UAAU,KAAK,MAAM,KAAK,CAAC;KAC7B,UAAU;MACR,sBAAsB,IAAI;KAC5B;KAIA,MAAM,yBAAyB,wBAAwB;IACzD,SAAS,GAAG;KACV,UAAU,iBAAiB,SAAS,CAAC;IACvC;GAEJ;EACF,SAAS,KAAK;GACZ,IAAI,CAAC,gBACH,UAAU,4BAA4B,GAAG;EAE7C;CACF,GAAG;CAGH,IAAI;CACJ,MAAM,6BAAsD,CAAC;CAC7D,sBAAsB,0BAA0B;CAChD,IAAI;EACF,SAAS,UAAU,WAAW;CAChC,SAAS,KAAK;EACZ,sBAAsB,IAAI;EAC1B,iBAAiB;EACjB,OAAO,OAAO,EAAE,YAAY,CAAC,CAAC;EAC9B,MAAM;CACR;CACA,sBAAsB,IAAI;CAG1B,MAAM,yBAAyB,0BAA0B;CAIzD,QAAQ,QAAQ,MAAM,EAAE,YAAY;EAClC,iBAAiB;EACjB,OAAO,OAAO,EAAE,YAAY,CAAC,CAAC;CAChC,CAAC;CAGD,MAAM,cAAc;EAClB,IAAI;GACF,OAAO,YAAY;EACrB,QAAQ,CAER;CACF,CAAC;CAED,OAAO;AACT"} |
@@ -75,11 +75,11 @@ import { hydrateIdAttribute, hydrateWhenAttribute } from "./constants.js"; | ||
| return new Promise((resolve) => { | ||
| const state = { disposed: false }; | ||
| const cleanupStrategyRef = { current: void 0 }; | ||
| let disposed = false; | ||
| let cleanupStrategy = void 0; | ||
| let cleanupHydrate = () => {}; | ||
| const finish = (reason) => { | ||
| if (state.disposed) return; | ||
| state.disposed = true; | ||
| if (disposed) return; | ||
| disposed = true; | ||
| options.signal.removeEventListener("abort", onAbort); | ||
| cleanupHydrate(); | ||
| runHydrationStrategyCleanup(cleanupStrategyRef.current)?.(); | ||
| runHydrationStrategyCleanup(cleanupStrategy)?.(); | ||
| resolve(reason); | ||
@@ -90,8 +90,7 @@ }; | ||
| cleanupHydrate = options.onHydrate(() => finish("hydrate")); | ||
| const cleanupStrategy = strategy._s?.({ | ||
| cleanupStrategy = strategy._s?.({ | ||
| element: options.element, | ||
| prefetch: () => finish("prefetch") | ||
| }); | ||
| cleanupStrategyRef.current = cleanupStrategy; | ||
| if (state.disposed) runHydrationStrategyCleanup(cleanupStrategy)?.(); | ||
| if (disposed) runHydrationStrategyCleanup(cleanupStrategy)?.(); | ||
| }); | ||
@@ -98,0 +97,0 @@ } |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"runtime.js","names":[],"sources":["../../../src/hydration/runtime.ts"],"sourcesContent":["import { hydrateIdAttribute, hydrateWhenAttribute } from './constants'\nimport type {\n HydrationPrefetchStrategy,\n HydrationPrefetchWaitReason,\n HydrationRuntimeGate,\n HydrationWhen,\n} from './types'\n\nconst hydrateIdSelector = `[${hydrateIdAttribute}]`\n\nexport type HydrationGateRecord = HydrationRuntimeGate & {\n id: string\n when: HydrationWhen\n promise: Promise<void>\n consumers: number\n resolveListeners: Set<() => void>\n}\n\nconst gateRegistry = /* @__PURE__ */ new Map<string, HydrationGateRecord>()\nconst resolvedGateIds = /* @__PURE__ */ new Set<string>()\nconst fallbackHtmlByGateId = /* @__PURE__ */ new Map<string, string>()\n\nexport function createResolvedGate(\n id: string,\n when: HydrationWhen,\n): HydrationGateRecord {\n return {\n id,\n when,\n promise: Promise.resolve(),\n resolve: () => {},\n resolved: true,\n consumers: 0,\n resolveListeners: new Set<() => void>(),\n }\n}\n\nexport function getOrCreateGate(\n id: string,\n when: HydrationWhen,\n): HydrationGateRecord {\n const existing = gateRegistry.get(id)\n if (existing?.when === when) {\n existing.consumers++\n return existing\n }\n\n let resolvePromise!: () => void\n const promise = new Promise<void>((resolve) => {\n resolvePromise = resolve\n })\n\n const gate: HydrationGateRecord = {\n id,\n promise,\n resolved: false,\n consumers: 1,\n when,\n resolveListeners: new Set(),\n resolve: () => {\n if (gate.resolved) return\n gate.resolved = true\n resolvePromise()\n gate.resolveListeners.forEach((listener) => listener())\n gate.resolveListeners.clear()\n },\n }\n\n gateRegistry.set(id, gate)\n if (when !== 'never' && resolvedGateIds.has(id)) {\n resolvedGateIds.delete(id)\n gate.resolve()\n }\n return gate\n}\n\nexport function releaseGate(gate: HydrationGateRecord) {\n resolvedGateIds.delete(gate.id)\n gate.consumers--\n if (gate.consumers > 0) return\n if (gateRegistry.get(gate.id) === gate) {\n gateRegistry.delete(gate.id)\n fallbackHtmlByGateId.delete(gate.id)\n gate.resolveListeners.clear()\n }\n}\n\nexport function onGateResolve(gate: HydrationGateRecord, listener: () => void) {\n if (gate.resolved) {\n listener()\n return () => {}\n }\n\n gate.resolveListeners.add(listener)\n return () => {\n gate.resolveListeners.delete(listener)\n }\n}\n\nexport function runHydrationStrategyCleanup(cleanup: void | (() => void)) {\n if (typeof cleanup === 'function') return cleanup\n return undefined\n}\n\nexport function waitForHydrationPrefetchStrategy(\n strategy: HydrationPrefetchStrategy,\n options: {\n element: Element | null\n signal: AbortSignal\n onHydrate: (listener: () => void) => () => void\n },\n): Promise<HydrationPrefetchWaitReason> {\n if (options.signal.aborted) {\n return Promise.resolve('abort')\n }\n\n return new Promise((resolve) => {\n const state = { disposed: false }\n const cleanupStrategyRef: { current: void | (() => void) } = {\n current: undefined,\n }\n let cleanupHydrate = () => {}\n\n const finish = (reason: HydrationPrefetchWaitReason) => {\n if (state.disposed) return\n state.disposed = true\n options.signal.removeEventListener('abort', onAbort)\n cleanupHydrate()\n runHydrationStrategyCleanup(cleanupStrategyRef.current)?.()\n resolve(reason)\n }\n\n const onAbort = () => finish('abort')\n\n options.signal.addEventListener('abort', onAbort, { once: true })\n cleanupHydrate = options.onHydrate(() => finish('hydrate'))\n const cleanupStrategy = strategy._s?.({\n element: options.element,\n prefetch: () => finish('prefetch'),\n })\n cleanupStrategyRef.current = cleanupStrategy\n if (state.disposed) {\n runHydrationStrategyCleanup(cleanupStrategy)?.()\n }\n })\n}\n\nexport function getMarkerGate(marker: Element) {\n const id = marker.getAttribute(hydrateIdAttribute)\n return id ? gateRegistry.get(id) : undefined\n}\n\nexport function resolveHydrationMarker(marker: Element) {\n const id = marker.getAttribute(hydrateIdAttribute)\n const when = marker.getAttribute(hydrateWhenAttribute)\n if (!id || !when || when === 'never') {\n return\n }\n\n const gate = gateRegistry.get(id)\n if (gate) {\n if (gate.when !== 'never') gate.resolve()\n return\n }\n\n resolvedGateIds.add(id)\n}\n\nexport function clearResolvedGateIdsInMarker(marker: Element) {\n const ownId = marker.getAttribute(hydrateIdAttribute)\n if (ownId) {\n resolvedGateIds.delete(ownId)\n }\n\n marker.querySelectorAll(hydrateIdSelector).forEach((childMarker) => {\n const childId = childMarker.getAttribute(hydrateIdAttribute)\n if (childId) {\n resolvedGateIds.delete(childId)\n }\n })\n}\n\nexport function saveFallbackHtml(id: string, element: Element) {\n if (!fallbackHtmlByGateId.has(id)) {\n fallbackHtmlByGateId.set(id, element.innerHTML)\n }\n}\n\nexport function getFallbackHtml(id: string) {\n return fallbackHtmlByGateId.get(id)\n}\n"],"mappings":";;AAQA,IAAM,oBAAoB,IAAI,mBAAmB;AAUjD,IAAM,+BAA+B,IAAI,IAAiC;AAC1E,IAAM,kCAAkC,IAAI,IAAY;AACxD,IAAM,uCAAuC,IAAI,IAAoB;AAErE,SAAgB,mBACd,IACA,MACqB;CACrB,OAAO;EACL;EACA;EACA,SAAS,QAAQ,QAAQ;EACzB,eAAe,CAAC;EAChB,UAAU;EACV,WAAW;EACX,kCAAkB,IAAI,IAAgB;CACxC;AACF;AAEA,SAAgB,gBACd,IACA,MACqB;CACrB,MAAM,WAAW,aAAa,IAAI,EAAE;CACpC,IAAI,UAAU,SAAS,MAAM;EAC3B,SAAS;EACT,OAAO;CACT;CAEA,IAAI;CAKJ,MAAM,OAA4B;EAChC;EACA,SAAA,IANkB,SAAe,YAAY;GAC7C,iBAAiB;EACnB,CAIE;EACA,UAAU;EACV,WAAW;EACX;EACA,kCAAkB,IAAI,IAAI;EAC1B,eAAe;GACb,IAAI,KAAK,UAAU;GACnB,KAAK,WAAW;GAChB,eAAe;GACf,KAAK,iBAAiB,SAAS,aAAa,SAAS,CAAC;GACtD,KAAK,iBAAiB,MAAM;EAC9B;CACF;CAEA,aAAa,IAAI,IAAI,IAAI;CACzB,IAAI,SAAS,WAAW,gBAAgB,IAAI,EAAE,GAAG;EAC/C,gBAAgB,OAAO,EAAE;EACzB,KAAK,QAAQ;CACf;CACA,OAAO;AACT;AAEA,SAAgB,YAAY,MAA2B;CACrD,gBAAgB,OAAO,KAAK,EAAE;CAC9B,KAAK;CACL,IAAI,KAAK,YAAY,GAAG;CACxB,IAAI,aAAa,IAAI,KAAK,EAAE,MAAM,MAAM;EACtC,aAAa,OAAO,KAAK,EAAE;EAC3B,qBAAqB,OAAO,KAAK,EAAE;EACnC,KAAK,iBAAiB,MAAM;CAC9B;AACF;AAEA,SAAgB,cAAc,MAA2B,UAAsB;CAC7E,IAAI,KAAK,UAAU;EACjB,SAAS;EACT,aAAa,CAAC;CAChB;CAEA,KAAK,iBAAiB,IAAI,QAAQ;CAClC,aAAa;EACX,KAAK,iBAAiB,OAAO,QAAQ;CACvC;AACF;AAEA,SAAgB,4BAA4B,SAA8B;CACxE,IAAI,OAAO,YAAY,YAAY,OAAO;AAE5C;AAEA,SAAgB,iCACd,UACA,SAKsC;CACtC,IAAI,QAAQ,OAAO,SACjB,OAAO,QAAQ,QAAQ,OAAO;CAGhC,OAAO,IAAI,SAAS,YAAY;EAC9B,MAAM,QAAQ,EAAE,UAAU,MAAM;EAChC,MAAM,qBAAuD,EAC3D,SAAS,KAAA,EACX;EACA,IAAI,uBAAuB,CAAC;EAE5B,MAAM,UAAU,WAAwC;GACtD,IAAI,MAAM,UAAU;GACpB,MAAM,WAAW;GACjB,QAAQ,OAAO,oBAAoB,SAAS,OAAO;GACnD,eAAe;GACf,4BAA4B,mBAAmB,OAAO,IAAI;GAC1D,QAAQ,MAAM;EAChB;EAEA,MAAM,gBAAgB,OAAO,OAAO;EAEpC,QAAQ,OAAO,iBAAiB,SAAS,SAAS,EAAE,MAAM,KAAK,CAAC;EAChE,iBAAiB,QAAQ,gBAAgB,OAAO,SAAS,CAAC;EAC1D,MAAM,kBAAkB,SAAS,KAAK;GACpC,SAAS,QAAQ;GACjB,gBAAgB,OAAO,UAAU;EACnC,CAAC;EACD,mBAAmB,UAAU;EAC7B,IAAI,MAAM,UACR,4BAA4B,eAAe,IAAI;CAEnD,CAAC;AACH;AAEA,SAAgB,cAAc,QAAiB;CAC7C,MAAM,KAAK,OAAO,aAAa,kBAAkB;CACjD,OAAO,KAAK,aAAa,IAAI,EAAE,IAAI,KAAA;AACrC;AAEA,SAAgB,uBAAuB,QAAiB;CACtD,MAAM,KAAK,OAAO,aAAa,kBAAkB;CACjD,MAAM,OAAO,OAAO,aAAa,oBAAoB;CACrD,IAAI,CAAC,MAAM,CAAC,QAAQ,SAAS,SAC3B;CAGF,MAAM,OAAO,aAAa,IAAI,EAAE;CAChC,IAAI,MAAM;EACR,IAAI,KAAK,SAAS,SAAS,KAAK,QAAQ;EACxC;CACF;CAEA,gBAAgB,IAAI,EAAE;AACxB;AAEA,SAAgB,6BAA6B,QAAiB;CAC5D,MAAM,QAAQ,OAAO,aAAa,kBAAkB;CACpD,IAAI,OACF,gBAAgB,OAAO,KAAK;CAG9B,OAAO,iBAAiB,iBAAiB,EAAE,SAAS,gBAAgB;EAClE,MAAM,UAAU,YAAY,aAAa,kBAAkB;EAC3D,IAAI,SACF,gBAAgB,OAAO,OAAO;CAElC,CAAC;AACH;AAEA,SAAgB,iBAAiB,IAAY,SAAkB;CAC7D,IAAI,CAAC,qBAAqB,IAAI,EAAE,GAC9B,qBAAqB,IAAI,IAAI,QAAQ,SAAS;AAElD;AAEA,SAAgB,gBAAgB,IAAY;CAC1C,OAAO,qBAAqB,IAAI,EAAE;AACpC"} | ||
| {"version":3,"file":"runtime.js","names":[],"sources":["../../../src/hydration/runtime.ts"],"sourcesContent":["import { hydrateIdAttribute, hydrateWhenAttribute } from './constants'\nimport type {\n HydrationPrefetchStrategy,\n HydrationPrefetchWaitReason,\n HydrationRuntimeGate,\n HydrationWhen,\n} from './types'\n\nconst hydrateIdSelector = `[${hydrateIdAttribute}]`\n\nexport type HydrationGateRecord = HydrationRuntimeGate & {\n id: string\n when: HydrationWhen\n promise: Promise<void>\n consumers: number\n resolveListeners: Set<() => void>\n}\n\nconst gateRegistry = /* @__PURE__ */ new Map<string, HydrationGateRecord>()\nconst resolvedGateIds = /* @__PURE__ */ new Set<string>()\nconst fallbackHtmlByGateId = /* @__PURE__ */ new Map<string, string>()\n\nexport function createResolvedGate(\n id: string,\n when: HydrationWhen,\n): HydrationGateRecord {\n return {\n id,\n when,\n promise: Promise.resolve(),\n resolve: () => {},\n resolved: true,\n consumers: 0,\n resolveListeners: new Set<() => void>(),\n }\n}\n\nexport function getOrCreateGate(\n id: string,\n when: HydrationWhen,\n): HydrationGateRecord {\n const existing = gateRegistry.get(id)\n if (existing?.when === when) {\n existing.consumers++\n return existing\n }\n\n let resolvePromise!: () => void\n const promise = new Promise<void>((resolve) => {\n resolvePromise = resolve\n })\n\n const gate: HydrationGateRecord = {\n id,\n promise,\n resolved: false,\n consumers: 1,\n when,\n resolveListeners: new Set(),\n resolve: () => {\n if (gate.resolved) return\n gate.resolved = true\n resolvePromise()\n gate.resolveListeners.forEach((listener) => listener())\n gate.resolveListeners.clear()\n },\n }\n\n gateRegistry.set(id, gate)\n if (when !== 'never' && resolvedGateIds.has(id)) {\n resolvedGateIds.delete(id)\n gate.resolve()\n }\n return gate\n}\n\nexport function releaseGate(gate: HydrationGateRecord) {\n resolvedGateIds.delete(gate.id)\n gate.consumers--\n if (gate.consumers > 0) return\n if (gateRegistry.get(gate.id) === gate) {\n gateRegistry.delete(gate.id)\n fallbackHtmlByGateId.delete(gate.id)\n gate.resolveListeners.clear()\n }\n}\n\nexport function onGateResolve(gate: HydrationGateRecord, listener: () => void) {\n if (gate.resolved) {\n listener()\n return () => {}\n }\n\n gate.resolveListeners.add(listener)\n return () => {\n gate.resolveListeners.delete(listener)\n }\n}\n\nexport function runHydrationStrategyCleanup(cleanup: void | (() => void)) {\n if (typeof cleanup === 'function') return cleanup\n return undefined\n}\n\nexport function waitForHydrationPrefetchStrategy(\n strategy: HydrationPrefetchStrategy,\n options: {\n element: Element | null\n signal: AbortSignal\n onHydrate: (listener: () => void) => () => void\n },\n): Promise<HydrationPrefetchWaitReason> {\n if (options.signal.aborted) {\n return Promise.resolve('abort')\n }\n\n return new Promise((resolve) => {\n let disposed = false\n // The strategy may finish synchronously before returning its cleanup.\n let cleanupStrategy: void | (() => void) = undefined\n let cleanupHydrate = () => {}\n\n const finish = (reason: HydrationPrefetchWaitReason) => {\n if (disposed) {\n return\n }\n disposed = true\n options.signal.removeEventListener('abort', onAbort)\n cleanupHydrate()\n runHydrationStrategyCleanup(cleanupStrategy)?.()\n resolve(reason)\n }\n\n const onAbort = () => finish('abort')\n\n options.signal.addEventListener('abort', onAbort, { once: true })\n cleanupHydrate = options.onHydrate(() => finish('hydrate'))\n cleanupStrategy = strategy._s?.({\n element: options.element,\n prefetch: () => finish('prefetch'),\n })\n // A synchronous finish must immediately run the cleanup just returned.\n // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition\n if (disposed) {\n runHydrationStrategyCleanup(cleanupStrategy)?.()\n }\n })\n}\n\nexport function getMarkerGate(marker: Element) {\n const id = marker.getAttribute(hydrateIdAttribute)\n return id ? gateRegistry.get(id) : undefined\n}\n\nexport function resolveHydrationMarker(marker: Element) {\n const id = marker.getAttribute(hydrateIdAttribute)\n const when = marker.getAttribute(hydrateWhenAttribute)\n if (!id || !when || when === 'never') {\n return\n }\n\n const gate = gateRegistry.get(id)\n if (gate) {\n if (gate.when !== 'never') gate.resolve()\n return\n }\n\n resolvedGateIds.add(id)\n}\n\nexport function clearResolvedGateIdsInMarker(marker: Element) {\n const ownId = marker.getAttribute(hydrateIdAttribute)\n if (ownId) {\n resolvedGateIds.delete(ownId)\n }\n\n marker.querySelectorAll(hydrateIdSelector).forEach((childMarker) => {\n const childId = childMarker.getAttribute(hydrateIdAttribute)\n if (childId) {\n resolvedGateIds.delete(childId)\n }\n })\n}\n\nexport function saveFallbackHtml(id: string, element: Element) {\n if (!fallbackHtmlByGateId.has(id)) {\n fallbackHtmlByGateId.set(id, element.innerHTML)\n }\n}\n\nexport function getFallbackHtml(id: string) {\n return fallbackHtmlByGateId.get(id)\n}\n"],"mappings":";;AAQA,IAAM,oBAAoB,IAAI,mBAAmB;AAUjD,IAAM,+BAA+B,IAAI,IAAiC;AAC1E,IAAM,kCAAkC,IAAI,IAAY;AACxD,IAAM,uCAAuC,IAAI,IAAoB;AAErE,SAAgB,mBACd,IACA,MACqB;CACrB,OAAO;EACL;EACA;EACA,SAAS,QAAQ,QAAQ;EACzB,eAAe,CAAC;EAChB,UAAU;EACV,WAAW;EACX,kCAAkB,IAAI,IAAgB;CACxC;AACF;AAEA,SAAgB,gBACd,IACA,MACqB;CACrB,MAAM,WAAW,aAAa,IAAI,EAAE;CACpC,IAAI,UAAU,SAAS,MAAM;EAC3B,SAAS;EACT,OAAO;CACT;CAEA,IAAI;CAKJ,MAAM,OAA4B;EAChC;EACA,SAAA,IANkB,SAAe,YAAY;GAC7C,iBAAiB;EACnB,CAIE;EACA,UAAU;EACV,WAAW;EACX;EACA,kCAAkB,IAAI,IAAI;EAC1B,eAAe;GACb,IAAI,KAAK,UAAU;GACnB,KAAK,WAAW;GAChB,eAAe;GACf,KAAK,iBAAiB,SAAS,aAAa,SAAS,CAAC;GACtD,KAAK,iBAAiB,MAAM;EAC9B;CACF;CAEA,aAAa,IAAI,IAAI,IAAI;CACzB,IAAI,SAAS,WAAW,gBAAgB,IAAI,EAAE,GAAG;EAC/C,gBAAgB,OAAO,EAAE;EACzB,KAAK,QAAQ;CACf;CACA,OAAO;AACT;AAEA,SAAgB,YAAY,MAA2B;CACrD,gBAAgB,OAAO,KAAK,EAAE;CAC9B,KAAK;CACL,IAAI,KAAK,YAAY,GAAG;CACxB,IAAI,aAAa,IAAI,KAAK,EAAE,MAAM,MAAM;EACtC,aAAa,OAAO,KAAK,EAAE;EAC3B,qBAAqB,OAAO,KAAK,EAAE;EACnC,KAAK,iBAAiB,MAAM;CAC9B;AACF;AAEA,SAAgB,cAAc,MAA2B,UAAsB;CAC7E,IAAI,KAAK,UAAU;EACjB,SAAS;EACT,aAAa,CAAC;CAChB;CAEA,KAAK,iBAAiB,IAAI,QAAQ;CAClC,aAAa;EACX,KAAK,iBAAiB,OAAO,QAAQ;CACvC;AACF;AAEA,SAAgB,4BAA4B,SAA8B;CACxE,IAAI,OAAO,YAAY,YAAY,OAAO;AAE5C;AAEA,SAAgB,iCACd,UACA,SAKsC;CACtC,IAAI,QAAQ,OAAO,SACjB,OAAO,QAAQ,QAAQ,OAAO;CAGhC,OAAO,IAAI,SAAS,YAAY;EAC9B,IAAI,WAAW;EAEf,IAAI,kBAAuC,KAAA;EAC3C,IAAI,uBAAuB,CAAC;EAE5B,MAAM,UAAU,WAAwC;GACtD,IAAI,UACF;GAEF,WAAW;GACX,QAAQ,OAAO,oBAAoB,SAAS,OAAO;GACnD,eAAe;GACf,4BAA4B,eAAe,IAAI;GAC/C,QAAQ,MAAM;EAChB;EAEA,MAAM,gBAAgB,OAAO,OAAO;EAEpC,QAAQ,OAAO,iBAAiB,SAAS,SAAS,EAAE,MAAM,KAAK,CAAC;EAChE,iBAAiB,QAAQ,gBAAgB,OAAO,SAAS,CAAC;EAC1D,kBAAkB,SAAS,KAAK;GAC9B,SAAS,QAAQ;GACjB,gBAAgB,OAAO,UAAU;EACnC,CAAC;EAGD,IAAI,UACF,4BAA4B,eAAe,IAAI;CAEnD,CAAC;AACH;AAEA,SAAgB,cAAc,QAAiB;CAC7C,MAAM,KAAK,OAAO,aAAa,kBAAkB;CACjD,OAAO,KAAK,aAAa,IAAI,EAAE,IAAI,KAAA;AACrC;AAEA,SAAgB,uBAAuB,QAAiB;CACtD,MAAM,KAAK,OAAO,aAAa,kBAAkB;CACjD,MAAM,OAAO,OAAO,aAAa,oBAAoB;CACrD,IAAI,CAAC,MAAM,CAAC,QAAQ,SAAS,SAC3B;CAGF,MAAM,OAAO,aAAa,IAAI,EAAE;CAChC,IAAI,MAAM;EACR,IAAI,KAAK,SAAS,SAAS,KAAK,QAAQ;EACxC;CACF;CAEA,gBAAgB,IAAI,EAAE;AACxB;AAEA,SAAgB,6BAA6B,QAAiB;CAC5D,MAAM,QAAQ,OAAO,aAAa,kBAAkB;CACpD,IAAI,OACF,gBAAgB,OAAO,KAAK;CAG9B,OAAO,iBAAiB,iBAAiB,EAAE,SAAS,gBAAgB;EAClE,MAAM,UAAU,YAAY,aAAa,kBAAkB;EAC3D,IAAI,SACF,gBAAgB,OAAO,OAAO;CAElC,CAAC;AACH;AAEA,SAAgB,iBAAiB,IAAY,SAAkB;CAC7D,IAAI,CAAC,qBAAqB,IAAI,EAAE,GAC9B,qBAAqB,IAAI,IAAI,QAAQ,SAAS;AAElD;AAEA,SAAgB,gBAAgB,IAAY;CAC1C,OAAO,qBAAqB,IAAI,EAAE;AACpC"} |
| //#region src/hydration/visible.ts | ||
| var visibleType = "visible"; | ||
| var observerRegistry = /* @__PURE__ */ new Map(); | ||
| function cleanupVisibleObserverEntry(observerEntry) { | ||
| if (observerEntry.elements.size > 0) return; | ||
| observerEntry.observer.disconnect(); | ||
| observerRegistry.delete(observerEntry.key); | ||
| function cleanupVisibleObserverEntry(key, observer, elements) { | ||
| if (elements.size > 0) return; | ||
| observer.disconnect(); | ||
| observerRegistry.delete(key); | ||
| } | ||
@@ -24,38 +24,36 @@ /* @__NO_SIDE_EFFECTS__ */ | ||
| if (!observerEntry) { | ||
| const entry = { | ||
| key, | ||
| elements: /* @__PURE__ */ new Map(), | ||
| observer: new IntersectionObserver((entries) => { | ||
| for (const intersectingEntry of entries) { | ||
| if (!intersectingEntry.isIntersecting) continue; | ||
| const callbacks = entry.elements.get(intersectingEntry.target); | ||
| if (!callbacks) continue; | ||
| callbacks.forEach((callback) => callback()); | ||
| entry.elements.delete(intersectingEntry.target); | ||
| entry.observer.unobserve(intersectingEntry.target); | ||
| cleanupVisibleObserverEntry(entry); | ||
| } | ||
| }, { | ||
| rootMargin, | ||
| threshold | ||
| }) | ||
| }; | ||
| observerRegistry.set(key, entry); | ||
| observerEntry = entry; | ||
| const elements = /* @__PURE__ */ new Map(); | ||
| const observer = new IntersectionObserver((entries) => { | ||
| for (const intersectingEntry of entries) { | ||
| if (!intersectingEntry.isIntersecting) continue; | ||
| const callbacks = elements.get(intersectingEntry.target); | ||
| if (!callbacks) continue; | ||
| callbacks.forEach((callback) => callback()); | ||
| elements.delete(intersectingEntry.target); | ||
| observer.unobserve(intersectingEntry.target); | ||
| cleanupVisibleObserverEntry(key, observer, elements); | ||
| } | ||
| }, { | ||
| rootMargin, | ||
| threshold | ||
| }); | ||
| observerEntry = [observer, elements]; | ||
| observerRegistry.set(key, observerEntry); | ||
| } | ||
| let callbacks = observerEntry.elements.get(element); | ||
| const [observer, elements] = observerEntry; | ||
| let callbacks = elements.get(element); | ||
| if (!callbacks) { | ||
| callbacks = /* @__PURE__ */ new Set(); | ||
| observerEntry.elements.set(element, callbacks); | ||
| observerEntry.observer.observe(element); | ||
| elements.set(element, callbacks); | ||
| observer.observe(element); | ||
| } | ||
| callbacks.add(callback); | ||
| return () => { | ||
| const currentCallbacks = observerEntry.elements.get(element); | ||
| const currentCallbacks = elements.get(element); | ||
| currentCallbacks?.delete(callback); | ||
| if (currentCallbacks?.size === 0) { | ||
| observerEntry.elements.delete(element); | ||
| observerEntry.observer.unobserve(element); | ||
| elements.delete(element); | ||
| observer.unobserve(element); | ||
| } | ||
| cleanupVisibleObserverEntry(observerEntry); | ||
| cleanupVisibleObserverEntry(key, observer, elements); | ||
| }; | ||
@@ -62,0 +60,0 @@ } |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"visible.js","names":[],"sources":["../../../src/hydration/visible.ts"],"sourcesContent":["import type { HydrationPrefetchStrategy } from './types'\n\nconst visibleType = 'visible'\n\nexport type VisibleHydrationOptions = {\n rootMargin?: string\n threshold?: number | Array<number>\n}\n\ntype VisibleObserverEntry = {\n key: string\n observer: IntersectionObserver\n elements: Map<Element, Set<() => void>>\n}\n\nconst observerRegistry = /* @__PURE__ */ new Map<string, VisibleObserverEntry>()\n\nfunction cleanupVisibleObserverEntry(observerEntry: VisibleObserverEntry) {\n if (observerEntry.elements.size > 0) return\n observerEntry.observer.disconnect()\n observerRegistry.delete(observerEntry.key)\n}\n\n/* @__NO_SIDE_EFFECTS__ */\nexport function visible(\n options: VisibleHydrationOptions = {},\n): HydrationPrefetchStrategy<typeof visibleType> {\n const rootMargin = options.rootMargin ?? '600px'\n const threshold = options.threshold ?? 0\n\n return {\n _t: visibleType,\n _s: ({ element, gate, prefetch }) => {\n const callback = prefetch ?? gate!.resolve\n\n if (!element) {\n callback()\n return\n }\n\n const key = `${rootMargin}|${\n Array.isArray(threshold) ? threshold.join(',') : String(threshold)\n }`\n let observerEntry = observerRegistry.get(key)\n\n if (!observerEntry) {\n const entry: VisibleObserverEntry = {\n key,\n elements: new Map<Element, Set<() => void>>(),\n observer: new IntersectionObserver(\n (entries) => {\n for (const intersectingEntry of entries) {\n if (!intersectingEntry.isIntersecting) continue\n\n const callbacks = entry.elements.get(intersectingEntry.target)\n if (!callbacks) continue\n\n callbacks.forEach((callback) => callback())\n entry.elements.delete(intersectingEntry.target)\n entry.observer.unobserve(intersectingEntry.target)\n cleanupVisibleObserverEntry(entry)\n }\n },\n { rootMargin, threshold },\n ),\n }\n observerRegistry.set(key, entry)\n observerEntry = entry\n }\n\n let callbacks = observerEntry.elements.get(element)\n if (!callbacks) {\n callbacks = new Set()\n observerEntry.elements.set(element, callbacks)\n observerEntry.observer.observe(element)\n }\n callbacks.add(callback)\n\n return () => {\n const currentCallbacks = observerEntry.elements.get(element)\n currentCallbacks?.delete(callback)\n if (currentCallbacks?.size === 0) {\n observerEntry.elements.delete(element)\n observerEntry.observer.unobserve(element)\n }\n cleanupVisibleObserverEntry(observerEntry)\n }\n },\n }\n}\n"],"mappings":";AAEA,IAAM,cAAc;AAapB,IAAM,mCAAmC,IAAI,IAAkC;AAE/E,SAAS,4BAA4B,eAAqC;CACxE,IAAI,cAAc,SAAS,OAAO,GAAG;CACrC,cAAc,SAAS,WAAW;CAClC,iBAAiB,OAAO,cAAc,GAAG;AAC3C;;AAGA,SAAgB,QACd,UAAmC,CAAC,GACW;CAC/C,MAAM,aAAa,QAAQ,cAAc;CACzC,MAAM,YAAY,QAAQ,aAAa;CAEvC,OAAO;EACL,IAAI;EACJ,KAAK,EAAE,SAAS,MAAM,eAAe;GACnC,MAAM,WAAW,YAAY,KAAM;GAEnC,IAAI,CAAC,SAAS;IACZ,SAAS;IACT;GACF;GAEA,MAAM,MAAM,GAAG,WAAW,GACxB,MAAM,QAAQ,SAAS,IAAI,UAAU,KAAK,GAAG,IAAI,OAAO,SAAS;GAEnE,IAAI,gBAAgB,iBAAiB,IAAI,GAAG;GAE5C,IAAI,CAAC,eAAe;IAClB,MAAM,QAA8B;KAClC;KACA,0BAAU,IAAI,IAA8B;KAC5C,UAAU,IAAI,sBACX,YAAY;MACX,KAAK,MAAM,qBAAqB,SAAS;OACvC,IAAI,CAAC,kBAAkB,gBAAgB;OAEvC,MAAM,YAAY,MAAM,SAAS,IAAI,kBAAkB,MAAM;OAC7D,IAAI,CAAC,WAAW;OAEhB,UAAU,SAAS,aAAa,SAAS,CAAC;OAC1C,MAAM,SAAS,OAAO,kBAAkB,MAAM;OAC9C,MAAM,SAAS,UAAU,kBAAkB,MAAM;OACjD,4BAA4B,KAAK;MACnC;KACF,GACA;MAAE;MAAY;KAAU,CAC1B;IACF;IACA,iBAAiB,IAAI,KAAK,KAAK;IAC/B,gBAAgB;GAClB;GAEA,IAAI,YAAY,cAAc,SAAS,IAAI,OAAO;GAClD,IAAI,CAAC,WAAW;IACd,4BAAY,IAAI,IAAI;IACpB,cAAc,SAAS,IAAI,SAAS,SAAS;IAC7C,cAAc,SAAS,QAAQ,OAAO;GACxC;GACA,UAAU,IAAI,QAAQ;GAEtB,aAAa;IACX,MAAM,mBAAmB,cAAc,SAAS,IAAI,OAAO;IAC3D,kBAAkB,OAAO,QAAQ;IACjC,IAAI,kBAAkB,SAAS,GAAG;KAChC,cAAc,SAAS,OAAO,OAAO;KACrC,cAAc,SAAS,UAAU,OAAO;IAC1C;IACA,4BAA4B,aAAa;GAC3C;EACF;CACF;AACF"} | ||
| {"version":3,"file":"visible.js","names":[],"sources":["../../../src/hydration/visible.ts"],"sourcesContent":["import type { HydrationPrefetchStrategy } from './types'\n\nconst visibleType = 'visible'\n\nexport type VisibleHydrationOptions = {\n rootMargin?: string\n threshold?: number | Array<number>\n}\n\ntype VisibleObserverEntry = [\n observer: IntersectionObserver,\n elements: Map<Element, Set<() => void>>,\n]\n\nconst observerRegistry = /* @__PURE__ */ new Map<string, VisibleObserverEntry>()\n\nfunction cleanupVisibleObserverEntry(\n key: string,\n observer: IntersectionObserver,\n elements: Map<Element, Set<() => void>>,\n) {\n if (elements.size > 0) {\n return\n }\n observer.disconnect()\n observerRegistry.delete(key)\n}\n\n/* @__NO_SIDE_EFFECTS__ */\nexport function visible(\n options: VisibleHydrationOptions = {},\n): HydrationPrefetchStrategy<typeof visibleType> {\n const rootMargin = options.rootMargin ?? '600px'\n const threshold = options.threshold ?? 0\n\n return {\n _t: visibleType,\n _s: ({ element, gate, prefetch }) => {\n const callback = prefetch ?? gate!.resolve\n\n if (!element) {\n callback()\n return\n }\n\n const key = `${rootMargin}|${\n Array.isArray(threshold) ? threshold.join(',') : String(threshold)\n }`\n let observerEntry = observerRegistry.get(key)\n\n if (!observerEntry) {\n const elements = new Map<Element, Set<() => void>>()\n const observer = new IntersectionObserver(\n (entries) => {\n for (const intersectingEntry of entries) {\n if (!intersectingEntry.isIntersecting) {\n continue\n }\n\n const callbacks = elements.get(intersectingEntry.target)\n if (!callbacks) {\n continue\n }\n\n callbacks.forEach((callback) => callback())\n elements.delete(intersectingEntry.target)\n observer.unobserve(intersectingEntry.target)\n cleanupVisibleObserverEntry(key, observer, elements)\n }\n },\n { rootMargin, threshold },\n )\n observerEntry = [observer, elements]\n observerRegistry.set(key, observerEntry)\n }\n\n const [observer, elements] = observerEntry\n let callbacks = elements.get(element)\n if (!callbacks) {\n callbacks = new Set()\n elements.set(element, callbacks)\n observer.observe(element)\n }\n callbacks.add(callback)\n\n return () => {\n const currentCallbacks = elements.get(element)\n currentCallbacks?.delete(callback)\n if (currentCallbacks?.size === 0) {\n elements.delete(element)\n observer.unobserve(element)\n }\n cleanupVisibleObserverEntry(key, observer, elements)\n }\n },\n }\n}\n"],"mappings":";AAEA,IAAM,cAAc;AAYpB,IAAM,mCAAmC,IAAI,IAAkC;AAE/E,SAAS,4BACP,KACA,UACA,UACA;CACA,IAAI,SAAS,OAAO,GAClB;CAEF,SAAS,WAAW;CACpB,iBAAiB,OAAO,GAAG;AAC7B;;AAGA,SAAgB,QACd,UAAmC,CAAC,GACW;CAC/C,MAAM,aAAa,QAAQ,cAAc;CACzC,MAAM,YAAY,QAAQ,aAAa;CAEvC,OAAO;EACL,IAAI;EACJ,KAAK,EAAE,SAAS,MAAM,eAAe;GACnC,MAAM,WAAW,YAAY,KAAM;GAEnC,IAAI,CAAC,SAAS;IACZ,SAAS;IACT;GACF;GAEA,MAAM,MAAM,GAAG,WAAW,GACxB,MAAM,QAAQ,SAAS,IAAI,UAAU,KAAK,GAAG,IAAI,OAAO,SAAS;GAEnE,IAAI,gBAAgB,iBAAiB,IAAI,GAAG;GAE5C,IAAI,CAAC,eAAe;IAClB,MAAM,2BAAW,IAAI,IAA8B;IACnD,MAAM,WAAW,IAAI,sBAClB,YAAY;KACX,KAAK,MAAM,qBAAqB,SAAS;MACvC,IAAI,CAAC,kBAAkB,gBACrB;MAGF,MAAM,YAAY,SAAS,IAAI,kBAAkB,MAAM;MACvD,IAAI,CAAC,WACH;MAGF,UAAU,SAAS,aAAa,SAAS,CAAC;MAC1C,SAAS,OAAO,kBAAkB,MAAM;MACxC,SAAS,UAAU,kBAAkB,MAAM;MAC3C,4BAA4B,KAAK,UAAU,QAAQ;KACrD;IACF,GACA;KAAE;KAAY;IAAU,CAC1B;IACA,gBAAgB,CAAC,UAAU,QAAQ;IACnC,iBAAiB,IAAI,KAAK,aAAa;GACzC;GAEA,MAAM,CAAC,UAAU,YAAY;GAC7B,IAAI,YAAY,SAAS,IAAI,OAAO;GACpC,IAAI,CAAC,WAAW;IACd,4BAAY,IAAI,IAAI;IACpB,SAAS,IAAI,SAAS,SAAS;IAC/B,SAAS,QAAQ,OAAO;GAC1B;GACA,UAAU,IAAI,QAAQ;GAEtB,aAAa;IACX,MAAM,mBAAmB,SAAS,IAAI,OAAO;IAC7C,kBAAkB,OAAO,QAAQ;IACjC,IAAI,kBAAkB,SAAS,GAAG;KAChC,SAAS,OAAO,OAAO;KACvB,SAAS,UAAU,OAAO;IAC5B;IACA,4BAA4B,KAAK,UAAU,QAAQ;GACrD;EACF;CACF;AACF"} |
+1
-1
| { | ||
| "name": "@tanstack/start-client-core", | ||
| "version": "1.170.18", | ||
| "version": "1.170.19", | ||
| "description": "Modern and scalable routing for React applications", | ||
@@ -5,0 +5,0 @@ "author": "Tanner Linsley", |
@@ -28,5 +28,5 @@ /** | ||
| /** Gets or creates a raw stream by ID (for use by deserialize plugin) */ | ||
| getOrCreateStream: (id: number) => ReadableStream<Uint8Array> | ||
| getStream: (id: number) => ReadableStream<Uint8Array> | ||
| /** Stream of JSON strings (NDJSON lines) */ | ||
| jsonChunks: ReadableStream<string> | ||
| chunks: ReadableStream<string> | ||
| } | ||
@@ -408,3 +408,3 @@ | ||
| return { getOrCreateStream, jsonChunks } | ||
| return { getStream: getOrCreateStream, chunks: jsonChunks } | ||
| } |
@@ -276,9 +276,6 @@ import { | ||
| const { getOrCreateStream, jsonChunks } = createFrameDecoder( | ||
| response.body, | ||
| ) | ||
| const { getStream, chunks } = createFrameDecoder(response.body) | ||
| // Create deserialize plugin that wires up the raw streams | ||
| const rawStreamPlugin = | ||
| createRawStreamDeserializePlugin(getOrCreateStream) | ||
| const rawStreamPlugin = createRawStreamDeserializePlugin(getStream) | ||
| const plugins = [rawStreamPlugin, ...(serovalPlugins || [])] | ||
@@ -288,3 +285,3 @@ | ||
| result = await processFramedResponse({ | ||
| jsonStream: jsonChunks, | ||
| jsonStream: chunks, | ||
| onMessage: (msg: any) => fromCrossJSON(msg, { refs, plugins }), | ||
@@ -291,0 +288,0 @@ onError(msg, error) { |
+12
-10
@@ -118,14 +118,15 @@ import { hydrateIdAttribute, hydrateWhenAttribute } from './constants' | ||
| return new Promise((resolve) => { | ||
| const state = { disposed: false } | ||
| const cleanupStrategyRef: { current: void | (() => void) } = { | ||
| current: undefined, | ||
| } | ||
| let disposed = false | ||
| // The strategy may finish synchronously before returning its cleanup. | ||
| let cleanupStrategy: void | (() => void) = undefined | ||
| let cleanupHydrate = () => {} | ||
| const finish = (reason: HydrationPrefetchWaitReason) => { | ||
| if (state.disposed) return | ||
| state.disposed = true | ||
| if (disposed) { | ||
| return | ||
| } | ||
| disposed = true | ||
| options.signal.removeEventListener('abort', onAbort) | ||
| cleanupHydrate() | ||
| runHydrationStrategyCleanup(cleanupStrategyRef.current)?.() | ||
| runHydrationStrategyCleanup(cleanupStrategy)?.() | ||
| resolve(reason) | ||
@@ -138,8 +139,9 @@ } | ||
| cleanupHydrate = options.onHydrate(() => finish('hydrate')) | ||
| const cleanupStrategy = strategy._s?.({ | ||
| cleanupStrategy = strategy._s?.({ | ||
| element: options.element, | ||
| prefetch: () => finish('prefetch'), | ||
| }) | ||
| cleanupStrategyRef.current = cleanupStrategy | ||
| if (state.disposed) { | ||
| // A synchronous finish must immediately run the cleanup just returned. | ||
| // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition | ||
| if (disposed) { | ||
| runHydrationStrategyCleanup(cleanupStrategy)?.() | ||
@@ -146,0 +148,0 @@ } |
+43
-36
@@ -10,14 +10,19 @@ import type { HydrationPrefetchStrategy } from './types' | ||
| type VisibleObserverEntry = { | ||
| key: string | ||
| observer: IntersectionObserver | ||
| elements: Map<Element, Set<() => void>> | ||
| } | ||
| type VisibleObserverEntry = [ | ||
| observer: IntersectionObserver, | ||
| elements: Map<Element, Set<() => void>>, | ||
| ] | ||
| const observerRegistry = /* @__PURE__ */ new Map<string, VisibleObserverEntry>() | ||
| function cleanupVisibleObserverEntry(observerEntry: VisibleObserverEntry) { | ||
| if (observerEntry.elements.size > 0) return | ||
| observerEntry.observer.disconnect() | ||
| observerRegistry.delete(observerEntry.key) | ||
| function cleanupVisibleObserverEntry( | ||
| key: string, | ||
| observer: IntersectionObserver, | ||
| elements: Map<Element, Set<() => void>>, | ||
| ) { | ||
| if (elements.size > 0) { | ||
| return | ||
| } | ||
| observer.disconnect() | ||
| observerRegistry.delete(key) | ||
| } | ||
@@ -48,31 +53,33 @@ | ||
| if (!observerEntry) { | ||
| const entry: VisibleObserverEntry = { | ||
| key, | ||
| elements: new Map<Element, Set<() => void>>(), | ||
| observer: new IntersectionObserver( | ||
| (entries) => { | ||
| for (const intersectingEntry of entries) { | ||
| if (!intersectingEntry.isIntersecting) continue | ||
| const elements = new Map<Element, Set<() => void>>() | ||
| const observer = new IntersectionObserver( | ||
| (entries) => { | ||
| for (const intersectingEntry of entries) { | ||
| if (!intersectingEntry.isIntersecting) { | ||
| continue | ||
| } | ||
| const callbacks = entry.elements.get(intersectingEntry.target) | ||
| if (!callbacks) continue | ||
| const callbacks = elements.get(intersectingEntry.target) | ||
| if (!callbacks) { | ||
| continue | ||
| } | ||
| callbacks.forEach((callback) => callback()) | ||
| entry.elements.delete(intersectingEntry.target) | ||
| entry.observer.unobserve(intersectingEntry.target) | ||
| cleanupVisibleObserverEntry(entry) | ||
| } | ||
| }, | ||
| { rootMargin, threshold }, | ||
| ), | ||
| } | ||
| observerRegistry.set(key, entry) | ||
| observerEntry = entry | ||
| callbacks.forEach((callback) => callback()) | ||
| elements.delete(intersectingEntry.target) | ||
| observer.unobserve(intersectingEntry.target) | ||
| cleanupVisibleObserverEntry(key, observer, elements) | ||
| } | ||
| }, | ||
| { rootMargin, threshold }, | ||
| ) | ||
| observerEntry = [observer, elements] | ||
| observerRegistry.set(key, observerEntry) | ||
| } | ||
| let callbacks = observerEntry.elements.get(element) | ||
| const [observer, elements] = observerEntry | ||
| let callbacks = elements.get(element) | ||
| if (!callbacks) { | ||
| callbacks = new Set() | ||
| observerEntry.elements.set(element, callbacks) | ||
| observerEntry.observer.observe(element) | ||
| elements.set(element, callbacks) | ||
| observer.observe(element) | ||
| } | ||
@@ -82,9 +89,9 @@ callbacks.add(callback) | ||
| return () => { | ||
| const currentCallbacks = observerEntry.elements.get(element) | ||
| const currentCallbacks = elements.get(element) | ||
| currentCallbacks?.delete(callback) | ||
| if (currentCallbacks?.size === 0) { | ||
| observerEntry.elements.delete(element) | ||
| observerEntry.observer.unobserve(element) | ||
| elements.delete(element) | ||
| observer.unobserve(element) | ||
| } | ||
| cleanupVisibleObserverEntry(observerEntry) | ||
| cleanupVisibleObserverEntry(key, observer, elements) | ||
| } | ||
@@ -91,0 +98,0 @@ }, |
8444
0.04%542425
-0.03%