Hyperlinkv0.8.0-beta.28

Stream

Stream.toReadableStreamWithconsteffect/Stream.ts:11145
<A, XR>(
  context: Context.Context<XR>,
  options?: { readonly strategy?: QueuingStrategy<A> | undefined }
): <E, R extends XR>(self: Stream<A, E, R>) => ReadableStream<A>
<A, E, XR, R extends XR>(
  self: Stream<A, E, R>,
  context: Context.Context<XR>,
  options?: { readonly strategy?: QueuingStrategy<A> | undefined }
): ReadableStream<A>

Converts the stream to a ReadableStream using the provided services.

When to use

Use when bridging to Web Streams and you already have the Context required to run the stream outside an Effect.

Details

See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.

Example (Converting to a ReadableStream with services)

import { Context, Stream } from "effect"

const stream = Stream.make(1, 2, 3, 4, 5)
const readableStream = Stream.toReadableStreamWith(stream, Context.empty())
destructors
Source effect/Stream.ts:1114556 lines
export const toReadableStreamWith = dual<
  <A, XR>(
    context: Context.Context<XR>,
    options?: { readonly strategy?: QueuingStrategy<A> | undefined }
  ) => <E, R extends XR>(self: Stream<A, E, R>) => ReadableStream<A>,
  <A, E, XR, R extends XR>(
    self: Stream<A, E, R>,
    context: Context.Context<XR>,
    options?: { readonly strategy?: QueuingStrategy<A> | undefined }
  ) => ReadableStream<A>
>(
  (args) => isStream(args[0]),
  <A, E, XR, R extends XR>(
    self: Stream<A, E, R>,
    context: Context.Context<XR>,
    options?: { readonly strategy?: QueuingStrategy<A> | undefined }
  ): ReadableStream<A> => {
    let currentResolve: (() => void) | undefined = undefined
    let fiber: Fiber.Fiber<void, E> | undefined = undefined
    const latch = Latch.makeUnsafe(false)

    return new ReadableStream<A>({
      start(controller) {
        fiber = Effect.runFork(Effect.provideContext(
          runForEachArray(self, (chunk) =>
            latch.whenOpen(Effect.sync(() => {
              latch.closeUnsafe()
              for (let i = 0; i < chunk.length; i++) {
                controller.enqueue(chunk[i])
              }
              currentResolve!()
              currentResolve = undefined
            }))),
          context
        ))
        fiber.addObserver((exit) => {
          if (exit._tag === "Failure") {
            controller.error(Cause.squash(exit.cause))
          } else {
            controller.close()
          }
        })
      },
      pull() {
        return new Promise<void>((resolve) => {
          currentResolve = resolve
          latch.openUnsafe()
        })
      },
      cancel() {
        if (!fiber) return
        return Effect.runPromise(Effect.asVoid(Fiber.interrupt(fiber)))
      }
    }, options?.strategy)
  }
)
Referenced by 2 symbols