MCPcopy Create free account
hub / github.com/Effect-TS/effect / fromReadableStream

Function fromReadableStream

packages/platform/bun/src/BunStream.ts:33–67  ·  view source on GitHub ↗
(
  options: {
    readonly evaluate: LazyArg<ReadableStream<A>>
    readonly onError: (error: unknown) => E
    readonly releaseLockOnEnd?: boolean | undefined
  }
)

Source from the content-addressed store, hash-verified

31 * @since 4.0.0
32 */
33export const fromReadableStream = <A, E>(
34 options: {
35 readonly evaluate: LazyArg<ReadableStream<A>>
36 readonly onError: (error: unknown) => E
37 readonly releaseLockOnEnd?: boolean | undefined
38 }
39): Stream.Stream<A, E> =>
40 Stream.fromChannel(Channel.fromTransform(Effect.fnUntraced(function*(_, scope) {
41 const reader = options.evaluate().getReader()
42 yield* Scope.addFinalizer(
43 scope,
44 options.releaseLockOnEnd ? Effect.sync(() => reader.releaseLock()) : Effect.promise(() => reader.cancel())
45 )
46 function readMany(): Pull.Pull<Arr.NonEmptyReadonlyArray<A>, E> {
47 const result = reader.readMany()
48 if ("then" in result) {
49 return Effect.callback<Arr.NonEmptyReadonlyArray<A>, E | Cause.Done>((resume) => {
50 result.then((_) => resume(handleResult(_)), (e) => resume(Effect.fail(options.onError(e))))
51 })
52 }
53 return handleResult(result)
54 }
55 function handleResult(
56 result: Bun.ReadableStreamDefaultReadManyResult<A>
57 ): Pull.Pull<Arr.NonEmptyReadonlyArray<A>, E> {
58 if (result.done) {
59 return Cause.done()
60 } else if (!Arr.isReadonlyArrayNonEmpty(result.value)) {
61 return readMany()
62 }
63 return Effect.succeed(result.value)
64 }
65 // @effect-diagnostics-next-line returnEffectInGen:off
66 return Effect.suspend(readMany)
67 })))

Callers

nothing calls this directly

Calls 4

addFinalizerMethod · 0.80
evaluateMethod · 0.65
syncMethod · 0.45
cancelMethod · 0.45

Tested by

no test coverage detected