effect-io-ai

Package: effect
Module: Stream

Stream.partitionQueue

Partitions a stream using a Filter and exposes passing and failing values as scoped queues.

Details

The queues are backed by a fiber in the current scope and should be consumed while that scope remains open. Each queue fails with the stream error or Cause.Done when the source ends.

Example (Partitioning a stream into queues)

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

const program = Effect.gen(function*() {
  const [passes, fails] = yield* Stream.make(1, 2, 3, 4).pipe(
    Stream.partitionQueue((n) => n % 2 === 0 ? Result.succeed(n) : Result.fail(n))
  )

  const passValues = yield* Stream.fromQueue(passes).pipe(Stream.runCollect)
  const failValues = yield* Stream.fromQueue(fails).pipe(Stream.runCollect)

  yield* Console.log(passValues)
  // Output: [ 2, 4 ]
  yield* Console.log(failValues)
  // Output: [ 1, 3 ]
})

Effect.runPromise(Effect.scoped(program))

Signature

declare const partitionQueue: { <A, Pass, Fail>(filter: Filter.Filter<NoInfer<A>, Pass, Fail>, options?: { readonly capacity?: number | "unbounded" | undefined; }): <E, R>(self: Stream<A, E, R>) => Effect.Effect<[passes: Queue.Dequeue<Pass, E | Cause.Done>, fails: Queue.Dequeue<Fail, E | Cause.Done>], never, R | Scope.Scope>; <A, E, R, Pass, Fail>(self: Stream<A, E, R>, filter: Filter.Filter<NoInfer<A>, Pass, Fail>, options?: { readonly capacity?: number | "unbounded" | undefined; }): Effect.Effect<[passes: Queue.Dequeue<Pass, E | Cause.Done>, fails: Queue.Dequeue<Fail, E | Cause.Done>], never, R | Scope.Scope>; }

Source

Since v4.0.0