Hyperlinkv0.8.0-beta.28

Stream

Stream.raceAllconsteffect/Stream.ts:4209
<S extends ReadonlyArray<Stream<any, any, any>>>(...streams: S): Stream<
  Success<S[number]>,
  Error<S[number]>,
  Services<S[number]>
>

Runs all streams concurrently until one stream emits its first value, then mirrors that winning stream and interrupts the rest.

Details

Failures or completion from losing streams before a winner is chosen are ignored unless every stream fails or completes before emitting. After a winner is chosen, that stream's later failures are propagated.

Example (Racing multiple streams)

import { Console, Effect, Schedule, Stream } from "effect"

const program = Effect.gen(function*() {
  const result = yield* Stream.raceAll(
    Stream.fromSchedule(Schedule.spaced("1 second")),
    Stream.make(0, 1, 2)
  ).pipe(Stream.runCollect)
  yield* Console.log(result)
})

Effect.runPromise(program)
// Output: [ 0, 1, 2 ]
racing
Source effect/Stream.ts:420928 lines
export const raceAll = <S extends ReadonlyArray<Stream<any, any, any>>>(
  ...streams: S
): Stream<Success<S[number]>, Error<S[number]>, Services<S[number]>> =>
  fromChannel(Channel.fromTransform((_, scope) =>
    Effect.sync(() => {
      let winner:
        | Pull.Pull<Arr.NonEmptyReadonlyArray<Success<S[number]>>, Error<S[number]>, void, Services<S[number]>>
        | undefined
      const race = Effect.raceAll(streams.map((stream) => {
        const childScope = Scope.forkUnsafe(scope)
        return Channel.toPullScoped(stream.channel, childScope).pipe(
          Effect.flatMap((pull) => Effect.zip(Effect.succeed(pull), pull)),
          Effect.onExit((exit) => {
            if (exit._tag === "Success") {
              if (winner) {
                return Scope.close(childScope, exit)
              }
              winner = exit.value[0]
              return Effect.void
            }
            return Scope.close(childScope, exit)
          }),
          Effect.map(([, chunk]) => chunk)
        )
      }))
      return Effect.suspend(() => winner ?? race)
    })
  ))
Referenced by 1 symbols