Hyperlinkv0.8.0-beta.28

Stream

Stream.toAsyncIterableWithconsteffect/Stream.ts:11322
<XR>(context: Context.Context<XR>): <A, E, R extends XR>(
  self: Stream<A, E, R>
) => AsyncIterable<A>
<A, E, XR, R extends XR>(
  self: Stream<A, E, R>,
  context: Context.Context<XR>
): AsyncIterable<A>

Converts the stream to an AsyncIterable using the provided services.

When to use

Use when converting outside an Effect and you already have the Context needed to run the stream.

Example (Converting to an AsyncIterable with services)

import { Context, Stream } from "effect"

const stream = Stream.make(1, 2, 3)
const iterable = Stream.toAsyncIterableWith(stream, Context.empty())

const collect = async () => {
  const results: Array<number> = []
  for await (const value of iterable) {
    results.push(value)
  }
  console.log(results)
}

collect()
// [ 1, 2, 3 ]
destructors
Source effect/Stream.ts:1132245 lines
export const toAsyncIterableWith: {
  <XR>(context: Context.Context<XR>): <A, E, R extends XR>(self: Stream<A, E, R>) => AsyncIterable<A>
  <A, E, XR, R extends XR>(
    self: Stream<A, E, R>,
    context: Context.Context<XR>
  ): AsyncIterable<A>
} = dual(
  2,
  <A, E, XR, R extends XR>(
    self: Stream<A, E, R>,
    context: Context.Context<XR>
  ): AsyncIterable<A> => ({
    [Symbol.asyncIterator]() {
      const runPromise = Effect.runPromiseWith(context)
      const runPromiseExit = Effect.runPromiseExitWith(context)
      const scope = Scope.makeUnsafe()
      let pull: Pull.Pull<Arr.NonEmptyReadonlyArray<A>, E, void, R> | undefined
      let currentIter: Iterator<A> | undefined
      return {
        async next(): Promise<IteratorResult<A>> {
          if (currentIter) {
            const next = currentIter.next()
            if (!next.done) return next
            currentIter = undefined
          }
          pull ??= await runPromise(Channel.toPullScoped(self.channel, scope))
          const exit = await runPromiseExit(pull)
          if (Exit.isSuccess(exit)) {
            currentIter = exit.value[Symbol.iterator]()
            return currentIter.next()
          } else if (Pull.isDoneCause(exit.cause)) {
            return { done: true, value: undefined }
          }
          throw Cause.squash(exit.cause)
        },
        return(_) {
          return runPromise(Effect.as(
            Scope.close(scope, Exit.void),
            { done: true, value: undefined }
          ))
        }
      }
    }
  })
)
Referenced by 2 symbols