<A, E, R>(self: Stream<A, E, R>): Effect.Effect<
AsyncIterable<A>,
never,
R
>Creates an effect that yields an AsyncIterable using the current services.
When to use
Use when the AsyncIterable should be created inside Effect with the current
context supplying the stream's services.
Example (Creating an AsyncIterable effect)
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
const program = Effect.gen(function*() {
const iterable = yield* Stream.toAsyncIterableEffect(stream)
const values = yield* Effect.promise(async () => {
const collected: Array<number> = []
for await (const value of iterable) {
collected.push(value)
}
return collected
})
yield* Effect.sync(() => console.log(values))
})
Effect.runPromise(program)
// [ 1, 2, 3 ]export const const toAsyncIterableEffect: <A, E, R>(
self: Stream<A, E, R>
) => Effect.Effect<AsyncIterable<A>, never, R>
Creates an effect that yields an AsyncIterable using the current services.
When to use
Use when the AsyncIterable should be created inside Effect with the current
context supplying the stream's services.
Example (Creating an AsyncIterable effect)
import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
const program = Effect.gen(function*() {
const iterable = yield* Stream.toAsyncIterableEffect(stream)
const values = yield* Effect.promise(async () => {
const collected: Array<number> = []
for await (const value of iterable) {
collected.push(value)
}
return collected
})
yield* Effect.sync(() => console.log(values))
})
Effect.runPromise(program)
// [ 1, 2, 3 ]
toAsyncIterableEffect = <function (type parameter) A in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>R>(self: Stream<A, E, R>(parameter) self: {
channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>;
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; <…;
}
self: 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, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>A, function (type parameter) E in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>E, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>R>): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<interface AsyncIterable<T, TReturn = any, TNext = any>AsyncIterable<function (type parameter) A in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>A>, never, function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>R> =>
import EffectEffect.map(
import EffectEffect.context<function (type parameter) R in <A, E, R>(self: Stream<A, E, R>): Effect.Effect<AsyncIterable<A>, never, R>R>(),
(context: Context.Context<R>(parameter) context: {
mapUnsafe: ReadonlyMap<string, any>;
mutable: boolean;
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;
}
context) => 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>
}
toAsyncIterableWith(self: Stream<A, E, R>(parameter) self: {
channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>;
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; <…;
}
self, context: Context.Context<R>(parameter) context: {
mapUnsafe: ReadonlyMap<string, any>;
mutable: boolean;
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;
}
context)
)