@@ -88,6 +88,16 @@ const testLayer = ServerConfig.layerTest(process.cwd(), process.cwd()).pipe(
8888 Layer . provideMerge ( NodeServices . layer ) ,
8989) ;
9090
91+ function joinIterableFiber < A , E > ( fiber : Fiber . Fiber < Iterable < A > , E > ) {
92+ return Effect . gen ( function * ( ) {
93+ const result = yield * Fiber . join ( fiber ) . pipe ( Effect . timeoutOption ( "2 seconds" ) ) ;
94+ if ( result . _tag === "None" ) {
95+ return yield * Effect . die ( new Error ( "Timed out waiting for Droid test events." ) ) ;
96+ }
97+ return Array . from ( result . value ) ;
98+ } ) ;
99+ }
100+
91101it . effect ( "maps Droid SDK stream messages into canonical runtime events" , ( ) =>
92102 Effect . scoped (
93103 Effect . gen ( function * ( ) {
@@ -159,7 +169,7 @@ it.effect("maps Droid SDK stream messages into canonical runtime events", () =>
159169 } ) ;
160170 yield * adapter . sendTurn ( { threadId, input : "hello" } ) ;
161171
162- const events = Array . from ( yield * Fiber . join ( eventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ) ;
172+ const events = yield * joinIterableFiber ( eventsFiber ) ;
163173 assert . equal ( createOptions ?. modelId , "claude-sonnet" ) ;
164174 assert . equal ( createOptions ?. autonomyLevel , AutonomyLevel . High ) ;
165175 assert . equal ( createOptions ?. interactionMode , DroidInteractionMode . Auto ) ;
@@ -255,9 +265,7 @@ it.effect("keeps Droid token usage cumulative across turns", () =>
255265 runtimeMode : "full-access" ,
256266 } ) ;
257267 yield * adapter . sendTurn ( { threadId : usageThreadId , input : "first" } ) ;
258- const firstEvents = Array . from (
259- yield * Fiber . join ( firstEventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ,
260- ) ;
268+ const firstEvents = yield * joinIterableFiber ( firstEventsFiber ) ;
261269
262270 const secondEventsFiber = yield * adapter . streamEvents . pipe (
263271 Stream . filter ( ( event ) => event . threadId === usageThreadId ) ,
@@ -266,9 +274,7 @@ it.effect("keeps Droid token usage cumulative across turns", () =>
266274 Effect . forkChild ,
267275 ) ;
268276 yield * adapter . sendTurn ( { threadId : usageThreadId , input : "second" } ) ;
269- const secondEvents = Array . from (
270- yield * Fiber . join ( secondEventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ,
271- ) ;
277+ const secondEvents = yield * joinIterableFiber ( secondEventsFiber ) ;
272278 const events = [ ...firstEvents , ...secondEvents ] ;
273279 const usageEvents = events . filter ( ( event ) => event . type === "thread.token-usage.updated" ) ;
274280 const completedTurns = events . filter ( ( event ) => event . type === "turn.completed" ) ;
@@ -459,7 +465,7 @@ it.effect("uses final Droid create_message content when deltas are absent", () =
459465 } ) ;
460466 yield * adapter . sendTurn ( { threadId, input : "hello" } ) ;
461467
462- const events = Array . from ( yield * Fiber . join ( eventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ) ;
468+ const events = yield * joinIterableFiber ( eventsFiber ) ;
463469 const deltas = events . filter ( ( event ) => event . type === "content.delta" ) ;
464470 assert . deepEqual (
465471 deltas . map ( ( event ) => ( event . type === "content.delta" ? event . payload : undefined ) ) ,
@@ -524,7 +530,7 @@ it.effect("does not duplicate Droid final create_message text after streaming de
524530 } ) ;
525531 yield * adapter . sendTurn ( { threadId, input : "hello" } ) ;
526532
527- const events = Array . from ( yield * Fiber . join ( eventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ) ;
533+ const events = yield * joinIterableFiber ( eventsFiber ) ;
528534 const deltas = events . filter ( ( event ) => event . type === "content.delta" ) ;
529535 assert . deepEqual (
530536 deltas . map ( ( event ) => ( event . type === "content.delta" ? event . payload . delta : undefined ) ) ,
@@ -645,7 +651,7 @@ it.effect("does not duplicate Droid final thinking content after streaming delta
645651 } ) ;
646652 yield * adapter . sendTurn ( { threadId, input : "hello" } ) ;
647653
648- const events = Array . from ( yield * Fiber . join ( eventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ) ;
654+ const events = yield * joinIterableFiber ( eventsFiber ) ;
649655 const deltas = events . filter ( ( event ) => event . type === "content.delta" ) ;
650656 assert . deepEqual (
651657 deltas . map ( ( event ) => ( event . type === "content.delta" ? event . payload : undefined ) ) ,
@@ -894,17 +900,13 @@ it.effect("settles pending Droid permission and user-input waits when stopped",
894900 runtimeMode : "approval-required" ,
895901 } ) ;
896902 yield * adapter . sendTurn ( { threadId, input : "run lint" } ) ;
897- const openedEvents = Array . from (
898- yield * Fiber . join ( openedEventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ,
899- ) ;
903+ const openedEvents = yield * joinIterableFiber ( openedEventsFiber ) ;
900904 assert . deepEqual ( openedEvents . map ( ( event ) => event . type ) . toSorted ( ) , [
901905 "request.opened" ,
902906 "user-input.requested" ,
903907 ] ) ;
904908 yield * adapter . stopSession ( threadId ) ;
905- const resolvedEvents = Array . from (
906- yield * Fiber . join ( resolvedEventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ,
907- ) ;
909+ const resolvedEvents = yield * joinIterableFiber ( resolvedEventsFiber ) ;
908910 assert . deepEqual ( resolvedEvents . map ( ( event ) => event . type ) . toSorted ( ) , [
909911 "request.resolved" ,
910912 "user-input.resolved" ,
@@ -1010,7 +1012,7 @@ it.effect("marks Droid stream errors as failed turns", () =>
10101012 runtimeMode : "full-access" ,
10111013 } ) ;
10121014 yield * adapter . sendTurn ( { threadId, input : "hello" } ) ;
1013- const events = Array . from ( yield * Fiber . join ( eventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ) ;
1015+ const events = yield * joinIterableFiber ( eventsFiber ) ;
10141016 const runtimeError = events . find ( ( event ) => event . type === "runtime.error" ) ;
10151017 const turnCompleted = events . find ( ( event ) => event . type === "turn.completed" ) ;
10161018
@@ -1070,7 +1072,7 @@ it.effect("marks aborted Droid turns as interrupted without runtime error", () =
10701072 yield * Effect . promise ( ( ) => abortReady ) . pipe ( Effect . timeout ( "2 seconds" ) ) ;
10711073 yield * adapter . interruptTurn ( threadId ) ;
10721074
1073- const events = Array . from ( yield * Fiber . join ( eventsFiber ) . pipe ( Effect . timeout ( "2 seconds" ) ) ) ;
1075+ const events = yield * joinIterableFiber ( eventsFiber ) ;
10741076 assert . equal (
10751077 events . some ( ( event ) => event . type === "runtime.error" ) ,
10761078 false ,
0 commit comments