effect-io-ai

Package: effect
Module: Stream

Stream.fromQueue

Creates a stream that pulls values from a Queue.Dequeue.

Details

The stream emits non-empty batches of queued values and ends when the queue fails with Cause.Done; other queue failures are propagated.

Example (Creating a stream from a queue of values)

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

const program = Effect.gen(function*() {
  const queue = yield* Queue.unbounded<number>()
  yield* Queue.offer(queue, 1)
  yield* Queue.offer(queue, 2)
  yield* Queue.offer(queue, 3)
  yield* Queue.shutdown(queue)

  const stream = Stream.fromQueue(queue)
  const values = yield* Stream.runCollect(stream)
  yield* Console.log(values)
})

Effect.runPromise(program)
// Output: [ 1, 2, 3 ]

Signature

declare const fromQueue: <A, E>(queue: Queue.Dequeue<A, E>) => Stream<A, Exclude<E, Cause.Done>>

Source

Since v2.0.0