Hyperlinkv0.8.0-beta.28

Stream

Stream.partitionQueueconsteffect/Stream.ts:4466
<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
>

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))
filtering
Source effect/Stream.ts:446689 lines
export 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
  >
} = dual(
  (args) => isStream(args[0]),
  Effect.fnUntraced(
    function*<A, E, R, Pass, Fail>(
      self: Stream<A, E, R>,
      filter: Filter.Filter<NoInfer<A>, Pass, Fail>,
      options?: {
        readonly capacity?: number | "unbounded" | undefined
      }
    ): Effect.fn.Return<
      [
        passes: Queue.Dequeue<Pass, E | Cause.Done>,
        fails: Queue.Dequeue<Fail, E | Cause.Done>
      ],
      never,
      R | Scope.Scope
    > {
      const scope = yield* Effect.scope
      const pull = yield* Channel.toPullScoped(self.channel, scope)
      const capacity = options?.capacity === "unbounded" ? undefined : options?.capacity ?? DefaultChunkSize
      const passes = yield* Queue.make<Pass, E | Cause.Done>({ capacity })
      const fails = yield* Queue.make<Fail, E | Cause.Done>({ capacity })

      yield* Effect.gen(function*() {
        while (true) {
          const chunk = yield* pull
          const excluded: Array<Fail> = []
          const satisfying: Array<Pass> = []
          for (let i = 0; i < chunk.length; i++) {
            const result = filter(chunk[i] as NoInfer<A>)
            if (Result.isFailure(result)) {
              excluded.push(result.failure)
            } else {
              satisfying.push(result.success)
            }
          }
          let passFiber: Fiber.Fiber<any> | undefined = undefined
          if (satisfying.length > 0) {
            const leftover = Queue.offerAllUnsafe(passes, satisfying)
            if (leftover.length > 0) {
              passFiber = yield* Effect.forkChild(Queue.offerAll(passes, leftover))
            }
          }
          if (excluded.length > 0) {
            const leftover = Queue.offerAllUnsafe(fails, excluded)
            if (leftover.length > 0) {
              yield* Queue.offerAll(fails, leftover)
            }
          }
          if (passFiber) yield* Fiber.join(passFiber)
        }
      }).pipe(
        Effect.onError((cause) => {
          Queue.failCauseUnsafe(passes, cause)
          Queue.failCauseUnsafe(fails, cause)
          return Effect.void
        }),
        Effect.forkIn(scope)
      )

      return [passes, fails]
    }
  )
)
Referenced by 2 symbols