@@ -198,11 +198,36 @@ describe("custom seroval plugins", () => {
198198 } ) ;
199199} ) ;
200200
201+ function streamOf ( pieces : ( string | Uint8Array ) [ ] ) {
202+ const encoder = new TextEncoder ( ) ;
203+ return new ReadableStream < Uint8Array > ( {
204+ start ( controller ) {
205+ for ( const piece of pieces ) {
206+ controller . enqueue ( typeof piece === "string" ? encoder . encode ( piece ) : piece ) ;
207+ }
208+ controller . close ( ) ;
209+ } ,
210+ } ) ;
211+ }
212+
213+ function frame ( data : string ) {
214+ const size = new TextEncoder ( ) . encode ( data ) . length ;
215+ return `;0x${ size . toString ( 16 ) . padStart ( 8 , "0" ) } ;${ data } ` ;
216+ }
217+
218+ /** Splits bytes into pieces of `size`, the way a network read can. */
219+ function split ( text : string , size : number ) {
220+ const bytes = new TextEncoder ( ) . encode ( text ) ;
221+ const pieces : Uint8Array [ ] = [ ] ;
222+ for ( let i = 0 ; i < bytes . length ; i += size ) {
223+ pieces . push ( bytes . subarray ( i , i + size ) ) ;
224+ }
225+ return pieces ;
226+ }
227+
201228/** Frames one serialized node the way `serializeToJSONStream` does. */
202- function frame ( node : unknown ) {
203- const data = JSON . stringify ( node ) ;
204- const size = new TextEncoder ( ) . encode ( data ) . length . toString ( 16 ) . padStart ( 8 , "0" ) ;
205- return `;0x${ size } ;${ data } ` ;
229+ function frameNode ( node : unknown ) {
230+ return frame ( JSON . stringify ( node ) ) ;
206231}
207232
208233/** An argument list whose only item is a promise that a later frame was meant to settle. */
@@ -238,13 +263,30 @@ async function collectUnhandledRejections(run: () => Promise<void>) {
238263 return reasons ;
239264}
240265
266+ /** Reads one whole frame off a serialized stream; its header and data can arrive as separate pieces. */
267+ async function readFirstFrame ( stream : ReadableStream < Uint8Array > ) {
268+ const reader = stream . getReader ( ) ;
269+ let bytes = new Uint8Array ( 0 ) ;
270+ let end = Infinity ;
271+ while ( bytes . length < end ) {
272+ const { done, value } = await reader . read ( ) ;
273+ if ( done ) break ;
274+ const joined = new Uint8Array ( bytes . length + value . length ) ;
275+ joined . set ( bytes ) ;
276+ joined . set ( value , bytes . length ) ;
277+ bytes = joined ;
278+ if ( end === Infinity && bytes . length >= 12 ) {
279+ end = 12 + Number . parseInt ( new TextDecoder ( ) . decode ( bytes . subarray ( 3 , 11 ) ) , 16 ) ;
280+ }
281+ }
282+ await reader . cancel ( ) ;
283+ return new TextDecoder ( ) . decode ( bytes . subarray ( 0 , end ) ) ;
284+ }
285+
241286/** The first frame of a JSON stream that never completes on its own. */
242287async function firstFrame ( value : unknown ) {
243288 const { serializeToJSONStream } = await loadSerialization ( true ) ;
244- const reader = serializeToJSONStream ( value ) . getReader ( ) ;
245- const { value : chunk } = await reader . read ( ) ;
246- await reader . cancel ( ) ;
247- return new TextDecoder ( ) . decode ( chunk ) ;
289+ return await readFirstFrame ( serializeToJSONStream ( value ) ) ;
248290}
249291
250292describe ( "values waiting on a later frame" , ( ) => {
@@ -254,12 +296,13 @@ describe("values waiting on a later frame", () => {
254296
255297 afterEach ( ( ) => {
256298 vi . unstubAllEnvs ( ) ;
299+ vi . restoreAllMocks ( ) ;
257300 } ) ;
258301
259302 it ( "rejects a promise that is still pending when the body ends" , async ( ) => {
260303 const { deserializeJSONStream } = await loadSerialization ( true ) ;
261304
262- const [ arg ] = ( await deserializeJSONStream ( new Response ( frame ( PENDING_PROMISE_ARGS ) ) ) ) as [
305+ const [ arg ] = ( await deserializeJSONStream ( new Response ( frameNode ( PENDING_PROMISE_ARGS ) ) ) ) as [
263306 Promise < unknown > ,
264307 ] ;
265308
@@ -283,25 +326,30 @@ describe("values waiting on a later frame", () => {
283326
284327 it ( "rejects pending values with the failure when a later frame is malformed" , async ( ) => {
285328 const { deserializeJSONStream } = await loadSerialization ( true ) ;
329+ const consoleError = vi . spyOn ( console , "error" ) . mockImplementation ( ( ) => { } ) ;
286330 let outcome : Awaited < ReturnType < typeof settleWithin > > | undefined ;
287331
288332 const unhandled = await collectUnhandledRejections ( async ( ) => {
289333 const [ arg ] = ( await deserializeJSONStream (
290- new Response ( frame ( PENDING_PROMISE_ARGS ) + ";0xZZZZZZZZ;junk" ) ,
334+ new Response ( frameNode ( PENDING_PROMISE_ARGS ) + ";0xZZZZZZZZ;junk" ) ,
291335 ) ) as [ Promise < unknown > ] ;
292336 outcome = await settleWithin ( arg ) ;
293337 } ) ;
294338
295339 expect ( unhandled ) . toEqual ( [ ] ) ;
296340 expect ( outcome ?. status ) . toBe ( "rejected" ) ;
297341 expect ( ( outcome as { reason : Error } ) . reason . message ) . toBe ( "Malformed server function stream." ) ;
342+ expect ( consoleError ) . toHaveBeenCalledWith (
343+ expect . stringContaining ( "server function stream" ) ,
344+ ( outcome as { reason : Error } ) . reason ,
345+ ) ;
298346 } ) ;
299347
300348 it ( "does not report a pending promise nobody awaits once the body ends" , async ( ) => {
301349 const { deserializeJSONStream } = await loadSerialization ( true ) ;
302350
303351 const unhandled = await collectUnhandledRejections ( async ( ) => {
304- await deserializeJSONStream ( new Response ( frame ( PENDING_PROMISE_ARGS ) ) ) ;
352+ await deserializeJSONStream ( new Response ( frameNode ( PENDING_PROMISE_ARGS ) ) ) ;
305353 } ) ;
306354
307355 expect ( unhandled ) . toEqual ( [ ] ) ;
@@ -313,7 +361,7 @@ describe("values waiting on a later frame", () => {
313361 let arg : Promise < unknown > | undefined ;
314362
315363 const unhandled = await collectUnhandledRejections ( async ( ) => {
316- [ arg ] = ( await deserializeJSONStream ( new Response ( frame ( rejectedArgs ) ) ) ) as [
364+ [ arg ] = ( await deserializeJSONStream ( new Response ( frameNode ( rejectedArgs ) ) ) ) as [
317365 Promise < unknown > ,
318366 ] ;
319367 } ) ;
@@ -362,16 +410,14 @@ describe("values waiting on a later JS frame", () => {
362410
363411 afterEach ( ( ) => {
364412 vi . unstubAllEnvs ( ) ;
413+ vi . restoreAllMocks ( ) ;
365414 delete ( globalThis as any ) . self ;
366415 delete ( globalThis as any ) . $R ;
367416 } ) ;
368417
369418 async function firstJSFrame ( id : string , value : unknown ) {
370419 const { serializeToJSStream } = await loadSerialization ( true ) ;
371- const reader = serializeToJSStream ( id , value ) . getReader ( ) ;
372- const { value : chunk } = await reader . read ( ) ;
373- await reader . cancel ( ) ;
374- return new TextDecoder ( ) . decode ( chunk ) ;
420+ return await readFirstFrame ( serializeToJSStream ( id , value ) ) ;
375421 }
376422
377423 it ( "rejects a promise that is still pending when the body ends" , async ( ) => {
@@ -406,6 +452,7 @@ describe("values waiting on a later JS frame", () => {
406452 it ( "rejects pending values with the failure when a later frame is malformed" , async ( ) => {
407453 const body = await firstJSFrame ( "server-fn:1" , [ new Promise ( ( ) => { } ) ] ) ;
408454 const { deserializeJSStream } = await loadSerialization ( true ) ;
455+ const consoleError = vi . spyOn ( console , "error" ) . mockImplementation ( ( ) => { } ) ;
409456 let outcome : Awaited < ReturnType < typeof settleWithin > > | undefined ;
410457
411458 const unhandled = await collectUnhandledRejections ( async ( ) => {
@@ -419,6 +466,134 @@ describe("values waiting on a later JS frame", () => {
419466 expect ( unhandled ) . toEqual ( [ ] ) ;
420467 expect ( outcome ?. status ) . toBe ( "rejected" ) ;
421468 expect ( ( outcome as { reason : Error } ) . reason . message ) . toBe ( "Malformed server function stream." ) ;
469+ expect ( consoleError ) . toHaveBeenCalledWith (
470+ expect . stringContaining ( "server function stream" ) ,
471+ ( outcome as { reason : Error } ) . reason ,
472+ ) ;
422473 expect ( ( globalThis as any ) . $R [ "server-fn:1" ] ) . toBeUndefined ( ) ;
423474 } ) ;
475+
476+ it ( "releases the scope and cancels the body when the first frame fails" , async ( ) => {
477+ const body = await firstJSFrame ( "server-fn:3" , [ new Promise ( ( ) => { } ) ] ) ;
478+ const source = body . slice ( 12 ) + ';throw new Error("first frame failed")' ;
479+ const { deserializeJSStream } = await loadSerialization ( true ) ;
480+ let cancelled = false ;
481+ const stream = new ReadableStream < Uint8Array > ( {
482+ start ( controller ) {
483+ controller . enqueue ( new TextEncoder ( ) . encode ( frame ( source ) ) ) ;
484+ } ,
485+ cancel ( ) {
486+ cancelled = true ;
487+ } ,
488+ } ) ;
489+
490+ await expect ( deserializeJSStream ( "server-fn:3" , new Response ( stream ) ) ) . rejects . toThrow (
491+ "first frame failed" ,
492+ ) ;
493+ expect ( cancelled ) . toBe ( true ) ;
494+ expect ( ( globalThis as any ) . $R [ "server-fn:3" ] ) . toBeUndefined ( ) ;
495+ } ) ;
496+ } ) ;
497+
498+ describe ( "SerovalChunkReader" , ( ) => {
499+ beforeEach ( ( ) => {
500+ vi . resetModules ( ) ;
501+ } ) ;
502+
503+ afterEach ( ( ) => {
504+ vi . unstubAllEnvs ( ) ;
505+ vi . restoreAllMocks ( ) ;
506+ } ) ;
507+
508+ it ( "writes a fixed-size header before each chunk" , async ( ) => {
509+ const { serializeToJSONString } = await loadSerialization ( true ) ;
510+
511+ const payload = await serializeToJSONString ( 1 ) ;
512+
513+ expect ( payload ) . toMatch ( / ^ ; 0 x [ 0 - 9 a - f ] { 8 } ; / ) ;
514+ expect ( payload ) . toBe ( frame ( payload . slice ( 12 ) ) ) ;
515+ } ) ;
516+
517+ it ( "reads chunks split at any byte, including inside a character" , async ( ) => {
518+ const { SerovalChunkReader } = await loadSerialization ( true ) ;
519+ const body = frame ( "héllo wörld ✓" ) + frame ( "" ) + frame ( "second" ) ;
520+
521+ const reader = new SerovalChunkReader ( streamOf ( split ( body , 1 ) ) ) ;
522+ const chunks : string [ ] = [ ] ;
523+ await reader . drain ( chunk => chunks . push ( chunk ) ) ;
524+
525+ expect ( chunks ) . toEqual ( [ "héllo wörld ✓" , "" , "second" ] ) ;
526+ } ) ;
527+
528+ it ( "reads a large chunk in small pieces in linear time" , async ( ) => {
529+ const { SerovalChunkReader } = await loadSerialization ( true ) ;
530+ const data = "x" . repeat ( 16 * 1024 * 1024 ) ;
531+
532+ const start = performance . now ( ) ;
533+ const result = await new SerovalChunkReader ( streamOf ( split ( frame ( data ) , 16 * 1024 ) ) ) . next ( ) ;
534+
535+ expect ( result . value ) . toHaveLength ( data . length ) ;
536+ // Copying the whole buffer on every piece took about 2 seconds here.
537+ expect ( performance . now ( ) - start ) . toBeLessThan ( 500 ) ;
538+ } ) ;
539+
540+ it . each ( [
541+ [ "a missing delimiter" , "X0x00000003Yabc" ] ,
542+ [ "a non-hex size" , ";0x0000zz03;abc" ] ,
543+ [ "a missing 0x prefix" , ";0000000003;abc" ] ,
544+ [ "a truncated header" , ";0x0000" ] ,
545+ [ "truncated data" , ";0x000000ff;abc" ] ,
546+ ] ) ( "rejects %s" , async ( _ , body ) => {
547+ const { SerovalChunkReader } = await loadSerialization ( true ) ;
548+
549+ await expect ( new SerovalChunkReader ( streamOf ( [ body ] ) ) . next ( ) ) . rejects . toThrow (
550+ "Malformed server function stream." ,
551+ ) ;
552+ } ) ;
553+
554+ it ( "rejects a chunk over the size limit before buffering it" , async ( ) => {
555+ const { SerovalChunkReader } = await loadSerialization ( true ) ;
556+ const reader = new SerovalChunkReader ( streamOf ( [ ";0xffffffff;" ] ) , { maxChunkSize : 1024 } ) ;
557+
558+ await expect ( reader . next ( ) ) . rejects . toThrow ( / l a r g e r t h a n t h e l i m i t / ) ;
559+ } ) ;
560+
561+ it ( "reports a bad later chunk instead of leaving an unhandled rejection" , async ( ) => {
562+ const { serializeToJSONString, deserializeJSONStream } = await loadSerialization ( true ) ;
563+ const consoleError = vi . spyOn ( console , "error" ) . mockImplementation ( ( ) => { } ) ;
564+ const unhandled : unknown [ ] = [ ] ;
565+ const onUnhandled = ( error : unknown ) => unhandled . push ( error ) ;
566+ process . on ( "unhandledRejection" , onUnhandled ) ;
567+
568+ try {
569+ const first = await serializeToJSONString ( [ 1 , 2 ] ) ;
570+ const value = await deserializeJSONStream ( new Response ( first + frame ( "{nope" ) ) ) ;
571+ await new Promise ( resolve => setTimeout ( resolve , 20 ) ) ;
572+
573+ expect ( value ) . toEqual ( [ 1 , 2 ] ) ;
574+ expect ( unhandled ) . toEqual ( [ ] ) ;
575+ expect ( consoleError ) . toHaveBeenCalledWith (
576+ expect . stringContaining ( "server function stream" ) ,
577+ expect . any ( SyntaxError ) ,
578+ ) ;
579+ } finally {
580+ process . off ( "unhandledRejection" , onUnhandled ) ;
581+ }
582+ } ) ;
583+
584+ it ( "cancels the body when the first chunk cannot be parsed" , async ( ) => {
585+ const { deserializeJSONStream } = await loadSerialization ( true ) ;
586+ let cancelled = false ;
587+ const body = new ReadableStream < Uint8Array > ( {
588+ start ( controller ) {
589+ controller . enqueue ( new TextEncoder ( ) . encode ( frame ( "{nope" ) ) ) ;
590+ } ,
591+ cancel ( ) {
592+ cancelled = true ;
593+ } ,
594+ } ) ;
595+
596+ await expect ( deserializeJSONStream ( new Response ( body ) ) ) . rejects . toThrow ( SyntaxError ) ;
597+ expect ( cancelled ) . toBe ( true ) ;
598+ } ) ;
424599} ) ;
0 commit comments