@@ -17,6 +17,9 @@ export interface StreamExecutionEvent {
1717 raw : string ;
1818}
1919
20+ const MAX_SSE_BUFFER_CHARS = 1_048_576 ;
21+ const SSE_CONNECT_TIMEOUT_MS = 15_000 ;
22+
2023export async function streamExecutionEvents ( options : StreamExecutionEventsOptions ) : Promise < void > {
2124 const startupDeadline = Date . now ( ) + ( options . startupRetryMs ?? 20_000 ) ;
2225 const phase = options . phase ?? 1 ;
@@ -29,39 +32,81 @@ export async function streamExecutionEvents(options: StreamExecutionEventsOption
2932 return ;
3033 }
3134
32- const response = await fetch ( url , {
33- method : 'GET' ,
34- headers : {
35- Accept : 'text/event-stream' ,
36- Authorization : `Bearer ${ options . token } ` ,
37- } ,
38- signal : options . signal ,
39- } ) ;
40-
41- if ( ! response . ok ) {
42- if (
43- ( response . status === 404 || response . status === 425 || response . status === 409 )
44- && Date . now ( ) < startupDeadline
45- ) {
46- await sleep ( 600 ) ;
47- continue ;
35+ const request = createStreamController ( options . signal ) ;
36+ let response : Response ;
37+ try {
38+ response = await fetch ( url , {
39+ method : 'GET' ,
40+ headers : {
41+ Accept : 'text/event-stream' ,
42+ Authorization : `Bearer ${ options . token } ` ,
43+ } ,
44+ signal : request . signal ,
45+ } ) ;
46+ request . clearConnectDeadline ( ) ;
47+ } catch ( error ) {
48+ request . dispose ( ) ;
49+ throw error ;
50+ }
51+
52+ try {
53+ if ( ! response . ok ) {
54+ if (
55+ ( response . status === 404 || response . status === 425 || response . status === 409 )
56+ && Date . now ( ) < startupDeadline
57+ ) {
58+ await sleep ( 600 , options . signal ) ;
59+ continue ;
60+ }
61+
62+ const body = await response . text ( ) ;
63+ throw new Error (
64+ `Event stream failed (${ response . status } ): ${ body || response . statusText || 'unknown error' } ` ,
65+ ) ;
4866 }
4967
50- const body = await response . text ( ) ;
51- throw new Error (
52- `Event stream failed (${ response . status } ): ${ body || response . statusText || 'unknown error' } ` ,
53- ) ;
54- }
68+ if ( ! response . body ) {
69+ throw new Error ( 'Event stream response did not include a readable body.' ) ;
70+ }
5571
56- if ( ! response . body ) {
57- throw new Error ( 'Event stream response did not include a readable body.' ) ;
72+ await consumeEventStream ( response . body , options . onEvent , request . signal ) ;
73+ return ;
74+ } finally {
75+ request . dispose ( ) ;
5876 }
59-
60- await consumeEventStream ( response . body , options . onEvent , options . signal ) ;
61- return ;
6277 }
6378}
6479
80+ function createStreamController ( parent ?: AbortSignal ) : {
81+ signal : AbortSignal ;
82+ clearConnectDeadline : ( ) => void ;
83+ dispose : ( ) => void ;
84+ } {
85+ const controller = new AbortController ( ) ;
86+ const onAbort = ( ) => controller . abort ( parent ?. reason ?? new Error ( 'Event stream cancelled.' ) ) ;
87+ if ( parent ?. aborted ) onAbort ( ) ;
88+ else parent ?. addEventListener ( 'abort' , onAbort , { once : true } ) ;
89+ const deadline = setTimeout (
90+ ( ) => controller . abort ( new Error ( 'Event stream connection timed out.' ) ) ,
91+ SSE_CONNECT_TIMEOUT_MS ,
92+ ) ;
93+ deadline . unref ?.( ) ;
94+ let deadlineCleared = false ;
95+ const clearConnectDeadline = ( ) => {
96+ if ( deadlineCleared ) return ;
97+ deadlineCleared = true ;
98+ clearTimeout ( deadline ) ;
99+ } ;
100+ return {
101+ signal : controller . signal ,
102+ clearConnectDeadline,
103+ dispose : ( ) => {
104+ clearConnectDeadline ( ) ;
105+ parent ?. removeEventListener ( 'abort' , onAbort ) ;
106+ } ,
107+ } ;
108+ }
109+
65110function buildEventsUrl (
66111 baseUrl : string ,
67112 trajectoryId : string ,
@@ -88,22 +133,33 @@ async function consumeEventStream(
88133 const decoder = new TextDecoder ( ) ;
89134 let buffer = '' ;
90135
91- while ( true ) {
92- if ( signal ?. aborted ) {
93- return ;
94- }
136+ const onAbort = ( ) => void reader . cancel ( signal ?. reason ) . catch ( ( ) => undefined ) ;
137+ signal ?. addEventListener ( 'abort' , onAbort , { once : true } ) ;
138+ try {
139+ while ( true ) {
140+ if ( signal ?. aborted ) {
141+ return ;
142+ }
95143
96- const { done, value } = await reader . read ( ) ;
97- if ( done ) {
98- break ;
144+ const { done, value } = await reader . read ( ) ;
145+ if ( done ) {
146+ break ;
147+ }
148+
149+ buffer += decoder . decode ( value , { stream : true } ) ;
150+ if ( buffer . length > MAX_SSE_BUFFER_CHARS ) {
151+ await reader . cancel ( 'SSE event exceeded the bounded buffer size.' ) . catch ( ( ) => undefined ) ;
152+ throw new Error ( 'Event stream sent an oversized or unterminated event.' ) ;
153+ }
154+ buffer = flushBuffer ( buffer , onEvent ) ;
99155 }
100156
101- buffer += decoder . decode ( value , { stream : true } ) ;
102- buffer = flushBuffer ( buffer , onEvent ) ;
157+ buffer += decoder . decode ( ) ;
158+ flushBuffer ( buffer , onEvent ) ;
159+ } finally {
160+ signal ?. removeEventListener ( 'abort' , onAbort ) ;
161+ reader . releaseLock ( ) ;
103162 }
104-
105- buffer += decoder . decode ( ) ;
106- flushBuffer ( buffer , onEvent ) ;
107163}
108164
109165function flushBuffer (
@@ -185,6 +241,19 @@ function parseEventChunk(chunk: string): StreamExecutionEvent | null {
185241 } ;
186242}
187243
188- async function sleep ( ms : number ) : Promise < void > {
189- await new Promise ( ( resolve ) => setTimeout ( resolve , ms ) ) ;
244+ async function sleep ( ms : number , signal ?: AbortSignal ) : Promise < void > {
245+ await new Promise < void > ( ( resolve , reject ) => {
246+ if ( signal ?. aborted ) { reject ( signal . reason ?? new Error ( 'Event stream cancelled.' ) ) ; return ; }
247+ const cleanup = ( ) => signal ?. removeEventListener ( 'abort' , onAbort ) ;
248+ const onAbort = ( ) => {
249+ clearTimeout ( timer ) ;
250+ cleanup ( ) ;
251+ reject ( signal ?. reason ?? new Error ( 'Event stream cancelled.' ) ) ;
252+ } ;
253+ const timer = setTimeout ( ( ) => {
254+ cleanup ( ) ;
255+ resolve ( ) ;
256+ } , ms ) ;
257+ signal ?. addEventListener ( 'abort' , onAbort , { once : true } ) ;
258+ } ) ;
190259}
0 commit comments