(
options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}
)
| 31 | * @since 4.0.0 |
| 32 | */ |
| 33 | export 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 | }))) |
nothing calls this directly
no test coverage detected