11import { createServer as createHttpServer } from 'node:http' ;
2- import { createSecureServer , type Http2SecureServer } from 'node:http2' ;
2+ import {
3+ createSecureServer ,
4+ type Http2SecureServer ,
5+ constants as http2Constants ,
6+ } from 'node:http2' ;
37import { type AddressInfo , connect , createServer , type Server } from 'node:net' ;
48import type { TLSSocket } from 'node:tls' ;
59import { NODE_HTTP_ENV_VAR } from '@workflow/world' ;
@@ -25,6 +29,7 @@ import {
2529 EVENTS_AGENT_OPTIONS ,
2630 EVENTS_AGENT_OPTIONS_NO_H2 ,
2731 EVENTS_RECYCLE_AFTER_CONSECUTIVE_FAILURES ,
32+ EVENTS_REQUEST_TIMEOUT_MS ,
2833 getDispatcher ,
2934 getEventsDispatcher ,
3035 getNodeHttpAgents ,
@@ -33,6 +38,7 @@ import {
3338 getQueueRequestTimeoutMs ,
3439 getStreamCloseDispatcher ,
3540 getStreamDispatcher ,
41+ h2MultiplexInterceptor ,
3642 isRecyclableTransportError ,
3743 NODE_HTTP_BODY_TIMEOUT_MS ,
3844 NODE_HTTP_HEADERS_TIMEOUT_MS ,
@@ -145,6 +151,17 @@ describe('agent transport', () => {
145151 // fault, so it has to leave nothing of H2 behind: `allowH2: true` with the
146152 // interceptor skipped still keeps a wedged session (see
147153 // EVENTS_AGENT_OPTIONS_NO_H2).
154+ // undici's 300s defaults outlast the runtime's 240s replay budget, so a
155+ // silent stream would cost the whole delivery before it even failed.
156+ it ( 'arms events deadlines well inside the replay budget' , ( ) => {
157+ const REPLAY_BUDGET_MS = 240_000 ;
158+ for ( const options of [ EVENTS_AGENT_OPTIONS , EVENTS_AGENT_OPTIONS_NO_H2 ] ) {
159+ expect ( options . headersTimeout ) . toBe ( EVENTS_REQUEST_TIMEOUT_MS ) ;
160+ expect ( options . bodyTimeout ) . toBe ( EVENTS_REQUEST_TIMEOUT_MS ) ;
161+ }
162+ expect ( EVENTS_REQUEST_TIMEOUT_MS ) . toBeLessThan ( REPLAY_BUDGET_MS / 2 ) ;
163+ } ) ;
164+
148165 it ( 'gives the kill switch an events agent with no HTTP/2 at all' , ( ) => {
149166 expect ( EVENTS_AGENT_OPTIONS_NO_H2 . allowH2 ) . toBe ( false ) ;
150167 expect ( EVENTS_AGENT_OPTIONS_NO_H2 . pipelining ) . toBe ( 1 ) ;
@@ -476,6 +493,222 @@ describe('HTTP/2 multiplexing (events vs stream-write agents)', () => {
476493 } ) ;
477494} ) ;
478495
496+ // undici dispatches a streamed request body on HTTP/2 only once the connection
497+ // has no stream in flight, and stops its whole queue behind it until then. The
498+ // multiplexing interceptor re-buffers small bodies to avoid that, so the
499+ // production shape that still hit it was a body too large to re-buffer (a
500+ // multi-MiB run input or event batch) on a connection that never drained.
501+ describe ( 'events dispatcher on a busy HTTP/2 connection' , ( ) => {
502+ const LARGE_BODY_BYTES = 2 * 1024 * 1024 ;
503+ const LOOPBACK = { connect : { rejectUnauthorized : false } } ;
504+
505+ let server : Http2SecureServer ;
506+ let origin : string ;
507+ let seen : Array < { path : string ; httpVersion : string ; bytes : number } > ;
508+ let releaseHold : ( ( ) => void ) | undefined ;
509+ let holdArrived : Promise < void > ;
510+ let resetAttempts : number ;
511+
512+ beforeAll ( async ( ) => {
513+ // allowHTTP1 so one origin serves both the h2 stream and the h1 fallback.
514+ server = createSecureServer ( {
515+ key : TEST_KEY ,
516+ cert : TEST_CERT ,
517+ allowHTTP1 : true ,
518+ } ) ;
519+ server . on ( 'sessionError' , ( ) => undefined ) ;
520+ server . on ( 'clientError' , ( ) => undefined ) ;
521+ server . on ( 'request' , ( req , res ) => {
522+ req . stream ?. on ( 'error' , ( ) => undefined ) ;
523+ const path = String ( req . url ) ;
524+ let bytes = 0 ;
525+ req . on ( 'data' , ( chunk : Buffer ) => {
526+ bytes += chunk . length ;
527+ } ) ;
528+ req . on ( 'end' , ( ) => {
529+ seen . push ( { path, httpVersion : req . httpVersion , bytes } ) ;
530+ if ( path === '/hold' ) {
531+ holdArrivedResolve ( ) ;
532+ releaseHold = ( ) => {
533+ res . writeHead ( 200 ) ;
534+ res . end ( 'held' ) ;
535+ } ;
536+ return ;
537+ }
538+ // The edge's per-stream reset on a busy connection. `/reset-once`
539+ // answers the retry; `/reset-always` never does.
540+ if (
541+ ( path === '/reset-once' && ++ resetAttempts === 1 ) ||
542+ path === '/reset-always'
543+ ) {
544+ if ( path === '/reset-always' ) resetAttempts ++ ;
545+ req . stream . close ( http2Constants . NGHTTP2_ENHANCE_YOUR_CALM ) ;
546+ return ;
547+ }
548+ res . writeHead ( 200 ) ;
549+ res . end ( path ) ;
550+ } ) ;
551+ } ) ;
552+ await new Promise < void > ( ( resolve ) => {
553+ server . listen ( 0 , '127.0.0.1' , resolve ) ;
554+ } ) ;
555+ origin = `https://127.0.0.1:${ ( server . address ( ) as AddressInfo ) . port } ` ;
556+ } ) ;
557+
558+ afterAll ( async ( ) => {
559+ await new Promise < void > ( ( resolve ) => {
560+ server . close ( ( ) => resolve ( ) ) ;
561+ } ) ;
562+ } ) ;
563+
564+ let holdArrivedResolve : ( ) => void ;
565+ beforeEach ( ( ) => {
566+ seen = [ ] ;
567+ releaseHold = undefined ;
568+ resetAttempts = 0 ;
569+ holdArrived = new Promise < void > ( ( resolve ) => {
570+ holdArrivedResolve = resolve ;
571+ } ) ;
572+ } ) ;
573+
574+ /**
575+ * Opens a GET that the origin holds, so the H2 connection never drains.
576+ * Resolves once the origin has it, to a wrapper (not the response promise
577+ * itself, which an async return would await).
578+ */
579+ async function holdStreamOpen ( dispatcher : unknown ) {
580+ const held = fetch ( `${ origin } /hold` , {
581+ dispatcher,
582+ // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
583+ } as any ) . then ( ( r ) => r . text ( ) ) ;
584+ await holdArrived ;
585+ return { held } ;
586+ }
587+
588+ function postLarge ( dispatcher : unknown ) {
589+ return fetch ( `${ origin } /large` , {
590+ method : 'POST' ,
591+ body : new Uint8Array ( LARGE_BODY_BYTES ) ,
592+ dispatcher,
593+ // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
594+ } as any ) . then ( ( r ) => r . text ( ) ) ;
595+ }
596+
597+ /** Resolves to 'settled' or 'pending' after `ms`. */
598+ function stateAfter ( promise : Promise < unknown > , ms : number ) {
599+ return Promise . race ( [
600+ promise . then ( ( ) => 'settled' as const ) ,
601+ new Promise < 'pending' > ( ( resolve ) => {
602+ setTimeout ( ( ) => resolve ( 'pending' ) , ms ) ;
603+ } ) ,
604+ ] ) ;
605+ }
606+
607+ // Establishes the premise with the interceptor alone (no fallback): the
608+ // large POST cannot start while the held stream is in flight.
609+ it ( 'queues behind an in-flight H2 stream without the fallback (the failure this fixes)' , async ( ) => {
610+ const agent = new Agent ( { ...EVENTS_AGENT_OPTIONS , ...LOOPBACK } ) ;
611+ const dispatcher = agent . compose ( h2MultiplexInterceptor ) ;
612+ try {
613+ const { held } = await holdStreamOpen ( dispatcher ) ;
614+ const large = postLarge ( dispatcher ) ;
615+ expect ( await stateAfter ( large , 1_000 ) ) . toBe ( 'pending' ) ;
616+ expect ( seen . map ( ( s ) => s . path ) ) . toEqual ( [ '/hold' ] ) ;
617+
618+ releaseHold ?.( ) ;
619+ expect ( await held ) . toBe ( 'held' ) ;
620+ expect ( await large ) . toBe ( '/large' ) ;
621+ } finally {
622+ await agent . close ( ) ;
623+ }
624+ } ) ;
625+
626+ it ( 'sends it over HTTP/1.1 so it does not wait for the H2 connection to drain' , async ( ) => {
627+ const agent = createEventsDispatcher ( LOOPBACK ) ;
628+ try {
629+ const { held } = await holdStreamOpen ( agent ) ;
630+ const large = postLarge ( agent ) ;
631+ expect ( await stateAfter ( large , 5_000 ) ) . toBe ( 'settled' ) ;
632+ expect ( await large ) . toBe ( '/large' ) ;
633+ expect ( seen ) . toContainEqual ( {
634+ path : '/large' ,
635+ httpVersion : '1.1' ,
636+ bytes : LARGE_BODY_BYTES ,
637+ } ) ;
638+
639+ releaseHold ?.( ) ;
640+ expect ( await held ) . toBe ( 'held' ) ;
641+ expect ( seen ) . toContainEqual ( {
642+ path : '/hold' ,
643+ httpVersion : '2.0' ,
644+ bytes : 0 ,
645+ } ) ;
646+ } finally {
647+ await agent . close ( ) ;
648+ }
649+ } ) ;
650+
651+ it ( 'retries an event-log read whose stream the peer reset' , async ( ) => {
652+ const agent = createEventsDispatcher ( LOOPBACK ) ;
653+ try {
654+ const response = await fetch ( `${ origin } /reset-once` , {
655+ dispatcher : agent ,
656+ // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
657+ } as any ) ;
658+ expect ( await response . text ( ) ) . toBe ( '/reset-once' ) ;
659+ expect ( resetAttempts ) . toBe ( 2 ) ;
660+ } finally {
661+ await agent . close ( ) ;
662+ }
663+ } ) ;
664+
665+ // undici before 7.30.0 retired the wrong request when an H2 stream failed out
666+ // of order, leaving a phantom "running" slot on the connection for good
667+ // (nodejs/undici#5410, #5569). A streamed body then never dispatched on that
668+ // connection again, since it waits for zero running streams. Pins the undici
669+ // version this package depends on.
670+ it ( 'leaves no phantom running stream after a peer reset' , async ( ) => {
671+ const agent = new Agent ( { ...EVENTS_AGENT_OPTIONS , ...LOOPBACK } ) ;
672+ try {
673+ await (
674+ await fetch ( `${ origin } /warm` , {
675+ dispatcher : agent ,
676+ // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
677+ } as any )
678+ ) . text ( ) ;
679+ await expect (
680+ fetch ( `${ origin } /reset-always` , {
681+ dispatcher : agent ,
682+ // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
683+ } as any )
684+ ) . rejects . toThrow ( ) ;
685+ await new Promise ( ( resolve ) => setTimeout ( resolve , 100 ) ) ;
686+ expect ( agent . stats [ origin ] ?. running ) . toBe ( 0 ) ;
687+ } finally {
688+ await agent . close ( ) ;
689+ }
690+ } ) ;
691+
692+ // An event write may already have committed when its stream is reset, so it
693+ // must surface the error rather than be replayed.
694+ it ( 'does not retry an event write whose stream the peer reset' , async ( ) => {
695+ const agent = createEventsDispatcher ( LOOPBACK ) ;
696+ try {
697+ await expect (
698+ fetch ( `${ origin } /reset-always` , {
699+ method : 'POST' ,
700+ body : 'x' ,
701+ dispatcher : agent ,
702+ // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici dispatcher type doesn't match @types/node's RequestInit
703+ } as any )
704+ ) . rejects . toThrow ( ) ;
705+ expect ( resetAttempts ) . toBe ( 1 ) ;
706+ } finally {
707+ await agent . close ( ) ;
708+ }
709+ } ) ;
710+ } ) ;
711+
479712// The transport fault this whole mechanism exists for: an HTTP/2 session whose
480713// TCP connection stays established while no bytes cross it. undici keeps such a
481714// session in service — on a stream timeout it deliberately does not destroy the
0 commit comments