Hyperlinkv0.8.0-beta.28

Stream

Stream.interleaveWithconsteffect/Stream.ts:9620
<A2, E2, R2, E3, R3>(
  that: Stream<A2, E2, R2>,
  decider: Stream<boolean, E3, R3>
): <A, E, R>(
  self: Stream<A, E, R>
) => Stream<A2 | A, E2 | E3 | E, R2 | R3 | R>
<A, E, R, A2, E2, R2, E3, R3>(
  self: Stream<A, E, R>,
  that: Stream<A2, E2, R2>,
  decider: Stream<boolean, E3, R3>
): Stream<A | A2, E | E2 | E3, R | R2 | R3>

Interleaves two streams deterministically by following a boolean decider stream.

Details

The decider controls how many elements are pulled; if one side ends, pulls for that side are ignored.

Example (Interleaving two streams deterministically by following a boolean decider stream)

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

const program = Effect.gen(function*() {
  const left = Stream.make(1, 3, 5)
  const right = Stream.make(2, 4, 6)
  const decider = Stream.make(true, false, false, true, true)

  const values = yield* Stream.runCollect(
    Stream.interleaveWith(left, right, decider)
  )

  yield* Console.log(values)
})

Effect.runPromise(program)
// [ 1, 2, 4, 3, 5 ]
merging
Source effect/Stream.ts:962054 lines
export const interleaveWith: {
  <A2, E2, R2, E3, R3>(
    that: Stream<A2, E2, R2>,
    decider: Stream<boolean, E3, R3>
  ): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E3 | E, R2 | R3 | R>
  <A, E, R, A2, E2, R2, E3, R3>(
    self: Stream<A, E, R>,
    that: Stream<A2, E2, R2>,
    decider: Stream<boolean, E3, R3>
  ): Stream<A | A2, E | E2 | E3, R | R2 | R3>
} = dual(3, <A, E, R, A2, E2, R2, E3, R3>(
  self: Stream<A, E, R>,
  that: Stream<A2, E2, R2>,
  decider: Stream<boolean, E3, R3>
): Stream<A | A2, E | E2 | E3, R | R2 | R3> =>
  fromChannel(Channel.fromTransform(Effect.fnUntraced(function*(upstream, scope) {
    const pullDecider = yield* Channel.toTransform(Channel.flattenArray(decider.channel))(upstream, scope)
    const retry = Symbol()
    type retry = typeof retry
    let leftDone = false
    let rightDone = false
    const pullLeft = (yield* Channel.toTransform(Channel.flattenArray(self.channel))(
      upstream,
      scope
    )).pipe(
      Pull.catchDone(() => {
        leftDone = true
        return Effect.succeed<retry>(retry)
      })
    )
    const pullRight = (yield* Channel.toTransform(Channel.flattenArray(that.channel))(
      upstream,
      scope
    )).pipe(
      Pull.catchDone(() => {
        rightDone = true
        return Effect.succeed<retry>(retry)
      })
    )

    return Effect.gen(function*() {
      while (true) {
        if (leftDone && rightDone) {
          return yield* Cause.done()
        }
        const side = yield* pullDecider
        if (side && leftDone) continue
        if (!side && rightDone) continue
        const elem = yield* (side ? pullLeft : pullRight)
        if (elem === retry) continue
        return Arr.of(elem)
      }
    })
  }))))
Referenced by 1 symbols