Hyperlinkv0.8.0-beta.28

Stream

Stream.runIntoQueueconsteffect/Stream.ts:11716
<A, E>(queue: Queue.Queue<A, E | Cause.Done>): <R>(
  self: Stream<A, E, R>
) => Effect.Effect<void, never, R>
<A, E, R>(
  self: Stream<A, E, R>,
  queue: Queue.Queue<A, E | Cause.Done>
): Effect.Effect<void, never, R>

Runs the stream, offering each element to the provided queue and ending it with Cause.Done when the stream completes.

Example (Running a stream into a queue)

import { Cause, Effect, Queue, Stream } from "effect"

const program = Effect.gen(function*() {
  const queue = yield* Queue.bounded<number, Cause.Done>(4)

  yield* Effect.forkChild(
    Stream.runIntoQueue(Stream.fromIterable([1, 2, 3]), queue)
  )

  const values = [
    yield* Queue.take(queue),
    yield* Queue.take(queue),
    yield* Queue.take(queue)
  ]
  const done = yield* Effect.flip(Queue.take(queue))

  return { values, done }
})
destructors
export const runIntoQueue: {
  <A, E>(queue: Queue.Queue<A, E | Cause.Done>): <R>(self: Stream<A, E, R>) => Effect.Effect<void, never, R>
  <A, E, R>(self: Stream<A, E, R>, queue: Queue.Queue<A, E | Cause.Done>): Effect.Effect<void, never, R>
} = dual(2, <A, E, R>(
  self: Stream<A, E, R>,
  queue: Queue.Queue<A, E | Cause.Done>
): Effect.Effect<void, never, R> => Channel.runIntoQueueArray(self.channel, queue))