Hyperlinkv0.8.0-beta.28

Stream

Stream.repeatElementsconsteffect/Stream.ts:2908
<B, E2, R2>(schedule: Schedule.Schedule<B, unknown, E2, R2>): <A, E, R>(
  self: Stream<A, E, R>
) => Stream<A, E | E2, R2 | R>
<A, E, R, B, E2, R2>(
  self: Stream<A, E, R>,
  schedule: Schedule.Schedule<B, unknown, E2, R2>
): Stream<A, E | E2, R | R2>

Repeats each element of the stream according to the provided schedule, including the original emission.

Example (Repeating stream elements)

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

const program = Effect.gen(function*() {
  const values = yield* Stream.make("A", "B", "C").pipe(
    Stream.repeatElements(Schedule.recurs(1)),
    Stream.runCollect
  )
  yield* Console.log(values)
})

Effect.runPromise(program)
// Output: [ "A", "A", "B", "B", "C", "C" ]
sequencing
Source effect/Stream.ts:290844 lines
export const repeatElements: {
  <B, E2, R2>(
    schedule: Schedule.Schedule<B, unknown, E2, R2>
  ): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>
  <A, E, R, B, E2, R2>(
    self: Stream<A, E, R>,
    schedule: Schedule.Schedule<B, unknown, E2, R2>
  ): Stream<A, E | E2, R | R2>
} = dual(
  2,
  <A, E, R, B, E2, R2>(
    self: Stream<A, E, R>,
    schedule: Schedule.Schedule<B, unknown, E2, R2>
  ): Stream<A, E | E2, R | R2> =>
    fromChannel(Channel.fromTransform((upstream, scope) =>
      Effect.map(
        Channel.toTransform(Channel.flattenArray(self.channel))(upstream, scope),
        (pullElement) => {
          let pullRepeat: Pull.Pull<Arr.NonEmptyReadonlyArray<A>, E | E2, void, R | R2> | undefined = undefined

          const pull: Pull.Pull<
            Arr.NonEmptyReadonlyArray<A>,
            E,
            void,
            R | R2
          > = Effect.gen(function*() {
            const element = yield* pullElement
            const chunk = Arr.of(element)
            const step = yield* Schedule.toStepWithSleep(schedule)
            pullRepeat = step(element).pipe(
              Effect.as(chunk),
              Pull.catchDone((_) => {
                pullRepeat = undefined
                return pull
              })
            )
            return chunk
          })

          return Effect.suspend(() => pullRepeat ?? pull)
        }
      )
    ))
)