<A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}): Stream<A, E>Creates a stream from a lazily supplied Web ReadableStream.
Details
The stream reads from a ReadableStreamDefaultReader, maps read failures
with onError, and closes the reader when the stream finalizes. By default
the reader is canceled; set releaseLockOnEnd to release the lock instead.
Example (Creating a stream from a ReadableStream)
import { Console, Data, Effect, Stream } from "effect"
class StreamError extends Data.TaggedError("StreamError")<{ readonly cause: unknown }> {}
const readableStream = new ReadableStream({
start(controller) {
controller.enqueue(1)
controller.enqueue(2)
controller.enqueue(3)
controller.close()
}
})
const program = Effect.gen(function*() {
const stream = Stream.fromReadableStream({
evaluate: () => readableStream,
onError: (cause) => new StreamError({ cause })
})
const values = yield* Stream.runCollect(stream)
yield* Console.log(values)
})
Effect.runPromise(program)
// Output: [ 1, 2, 3 ]export const const fromReadableStream: <
A,
E
>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}) => Stream<A, E>
Creates a stream from a lazily supplied Web ReadableStream.
Details
The stream reads from a ReadableStreamDefaultReader, maps read failures
with onError, and closes the reader when the stream finalizes. By default
the reader is canceled; set releaseLockOnEnd to release the lock instead.
Example (Creating a stream from a ReadableStream)
import { Console, Data, Effect, Stream } from "effect"
class StreamError extends Data.TaggedError("StreamError")<{ readonly cause: unknown }> {}
const readableStream = new ReadableStream({
start(controller) {
controller.enqueue(1)
controller.enqueue(2)
controller.enqueue(3)
controller.close()
}
})
const program = Effect.gen(function*() {
const stream = Stream.fromReadableStream({
evaluate: () => readableStream,
onError: (cause) => new StreamError({ cause })
})
const values = yield* Stream.runCollect(stream)
yield* Console.log(values)
})
Effect.runPromise(program)
// Output: [ 1, 2, 3 ]
fromReadableStream = <function (type parameter) A in <A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean | undefined;
}): Stream<A, E>
A, function (type parameter) E in <A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean | undefined;
}): Stream<A, E>
E>(
options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}
options: {
readonly evaluate: LazyArg<ReadableStream<A>>evaluate: import LazyArgLazyArg<interface ReadableStream<R = any>The ReadableStream interface of the Streams API represents a readable stream of byte data.
ReadableStream<function (type parameter) A in <A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean | undefined;
}): Stream<A, E>
A>>
readonly onError: (error: unknown) => EonError: (error: unknownerror: unknown) => function (type parameter) E in <A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean | undefined;
}): Stream<A, E>
E
readonly releaseLockOnEnd?: boolean | undefinedreleaseLockOnEnd?: boolean | undefined
}
): interface Stream<out A, out E = never, out R = never>A Stream<A, E, R> describes a program that can emit many A values, fail
with E, and require R.
Details
Streams are pull-based with backpressure and emit chunks to amortize effect
evaluation. They support monadic composition and error handling similar to
Effect, adapted for multiple values.
Example (Creating and consuming streams)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
yield* Stream.make(1, 2, 3).pipe(
Stream.map((n) => n * 2),
Stream.runForEach((n) => Console.log(n))
)
})
Effect.runPromise(program)
// Output:
// 2
// 4
// 6
Stream<function (type parameter) A in <A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean | undefined;
}): Stream<A, E>
A, function (type parameter) E in <A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean | undefined;
}): Stream<A, E>
E> =>
const fromChannel: <
Arr extends Arr.NonEmptyReadonlyArray<any>,
E,
R
>(
channel: Channel.Channel<
Arr,
E,
void,
unknown,
unknown,
unknown,
R
>
) => Stream<
Arr extends Arr.NonEmptyReadonlyArray<infer A>
? A
: never,
E,
R
>
Creates a stream from a array-emitting Channel.
Example (Creating a stream from an array-emitting channel)
import { Channel, Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const channel = Channel.succeed([1, 2, 3] as const)
const stream = Stream.fromChannel(channel)
const result = yield* Stream.runCollect(stream)
yield* Console.log(result)
})
// Output: [ 1, 2, 3 ]
fromChannel(import ChannelChannel.fromTransform(import EffectEffect.fnUntraced(function*(_: Pull.Pull<
unknown,
unknown,
unknown,
never
>
(parameter) _: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
_, scope: Scope.Scope(parameter) scope: {
strategy: "sequential" | "parallel";
state: State.Open | State.Closed | State.Empty;
}
scope) {
const const reader: anyreader = options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}
options.evaluate: LazyArg<ReadableStream<A>>evaluate().getReader()
yield* import ScopeScope.addFinalizer(
scope: Scope.Scope(parameter) scope: {
strategy: "sequential" | "parallel";
state: State.Open | State.Closed | State.Empty;
}
scope,
options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}
options.releaseLockOnEnd?: boolean | undefinedreleaseLockOnEnd
? import EffectEffect.sync(() => const reader: anyreader.releaseLock())
: import EffectEffect.promise(() => const reader: anyreader.cancel().catch(import constVoidconstVoid))
)
return import EffectEffect.flatMap(
import EffectEffect.tryPromise({
try: () => anytry: () => const reader: anyreader.read(),
catch: (reason: any) => Ecatch: (reason: anyreason) => options: {
readonly evaluate: LazyArg<ReadableStream<A>>
readonly onError: (error: unknown) => E
readonly releaseLockOnEnd?: boolean | undefined
}
options.onError: (error: unknown) => EonError(reason: anyreason)
}),
({ done: anydone, value: anyvalue }) => done: anydone ? import CauseCause.done() : import EffectEffect.succeed(import ArrArr.of(value: anyvalue))
)
})))