Hyperlinkv0.8.0-beta.28

Queue

Queue.collectconsteffect/Queue.ts:1246
<A, E>(self: Dequeue<A, E | Done>): Effect<Array<A>, Pull.ExcludeDone<E>>

Takes all messages from the queue, until the queue has errored or is done.

Example (Collecting values until completion)

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

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

  // Add several messages
  yield* Queue.offerAll(queue, [1, 2, 3, 4, 5])
  // Some time later, end the queue
  yield* Effect.forkChild(Queue.end(queue))

  // Collect all available messages
  const messages = yield* Queue.collect(queue)
  console.log(messages) // [1, 2, 3, 4, 5]
})
taking
Source effect/Queue.ts:124619 lines
export const collect = <A, E>(self: Dequeue<A, E | Done>): Effect<Array<A>, Pull.ExcludeDone<E>> =>
  internalEffect.suspend(() => {
    const out = Arr.empty<A>()
    return internalEffect.as(
      Pull.catchDone(
        internalEffect.whileLoop({
          while: constTrue,
          body: constant(takeAll(self)),
          step(items: Arr.NonEmptyArray<A>) {
            for (let i = 0; i < items.length; i++) {
              out.push(items[i])
            }
          }
        }),
        () => internalEffect.void
      ),
      out
    )
  }) as any