<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>Converts the stream to an AsyncIterable using the provided services.
When to use
Use when converting outside an Effect and you already have the Context
needed to run the stream.
Example (Converting to an AsyncIterable with services)
import { Context, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
const iterable = Stream.toAsyncIterableWith(stream, Context.empty())
const collect = async () => {
const results: Array<number> = []
for await (const value of iterable) {
results.push(value)
}
console.log(results)
}
collect()
// [ 1, 2, 3 ]export const 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>
}
Converts the stream to an AsyncIterable using the provided services.
When to use
Use when converting outside an Effect and you already have the Context
needed to run the stream.
Example (Converting to an AsyncIterable with services)
import { Context, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
const iterable = Stream.toAsyncIterableWith(stream, Context.empty())
const collect = async () => {
const results: Array<number> = []
for await (const value of iterable) {
results.push(value)
}
console.log(results)
}
collect()
// [ 1, 2, 3 ]
toAsyncIterableWith: {
<function (type parameter) XR in <XR>(context: Context.Context<XR>): <A, E, R extends XR>(self: Stream<A, E, R>) => AsyncIterable<A>XR>(context: Context.Context<XR>(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: import ContextContext.type Context.Context = /*unresolved*/ anyContext<function (type parameter) XR in <XR>(context: Context.Context<XR>): <A, E, R extends XR>(self: Stream<A, E, R>) => AsyncIterable<A>XR>): <function (type parameter) A in <A, E, R extends XR>(self: Stream<A, E, R>): AsyncIterable<A>A, function (type parameter) E in <A, E, R extends XR>(self: Stream<A, E, R>): AsyncIterable<A>E, function (type parameter) R in <A, E, R extends XR>(self: Stream<A, E, R>): AsyncIterable<A>R extends function (type parameter) XR in <XR>(context: Context.Context<XR>): <A, E, R extends XR>(self: Stream<A, E, R>) => AsyncIterable<A>XR>(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 extends XR>(self: Stream<A, E, R>): AsyncIterable<A>A, function (type parameter) E in <A, E, R extends XR>(self: Stream<A, E, R>): AsyncIterable<A>E, function (type parameter) R in <A, E, R extends XR>(self: Stream<A, E, R>): AsyncIterable<A>R>) => interface AsyncIterable<T, TReturn = any, TNext = any>AsyncIterable<function (type parameter) A in <A, E, R extends XR>(self: Stream<A, E, R>): AsyncIterable<A>A>
<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A, function (type parameter) E in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>E, function (type parameter) XR in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>XR, function (type parameter) R in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>R extends function (type parameter) XR in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>XR>(
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, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A, function (type parameter) E in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>E, function (type parameter) R in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>R>,
context: Context.Context<XR>(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: import ContextContext.type Context.Context = /*unresolved*/ anyContext<function (type parameter) XR in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>XR>
): interface AsyncIterable<T, TReturn = any, TNext = any>AsyncIterable<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A>
} = import dualdual(
2,
<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A, function (type parameter) E in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>E, function (type parameter) XR in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>XR, function (type parameter) R in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>R extends function (type parameter) XR in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>XR>(
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, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A, function (type parameter) E in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>E, function (type parameter) R in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>R>,
context: Context.Context<XR>(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: import ContextContext.type Context.Context = /*unresolved*/ anyContext<function (type parameter) XR in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>XR>
): interface AsyncIterable<T, TReturn = any, TNext = any>AsyncIterable<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A> => ({
[var Symbol: SymbolConstructorSymbol.SymbolConstructor.asyncIterator: typeof Symbol.asyncIteratorA method that returns the default async iterator for an object. Called by the semantics of
the for-await-of statement.
asyncIterator]() {
const const runPromise: (
effect: Effect<A, E, XR>,
options?: RunOptions | undefined
) => Promise<A>
runPromise = import EffectEffect.runPromiseWith(context: Context.Context<XR>(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 const runPromiseExit: (
effect: Effect<A, E, XR>,
options?: RunOptions | undefined
) => Promise<Exit.Exit<A, E>>
runPromiseExit = import EffectEffect.runPromiseExitWith(context: Context.Context<XR>(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 const scope: Scope.Closeableconst scope: {
strategy: "sequential" | "parallel";
state: State.Open | State.Closed | State.Empty;
}
scope = import ScopeScope.makeUnsafe()
let let pull:
| Pull.Pull<
Arr.NonEmptyReadonlyArray<A>,
E,
void,
R
>
| undefined
pull: import PullPull.type Pull.Pull = /*unresolved*/ anyPull<import ArrArr.type Arr.NonEmptyReadonlyArray = /*unresolved*/ anyNonEmptyReadonlyArray<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A>, function (type parameter) E in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>E, void, function (type parameter) R in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>R> | undefined
let let currentIter:
| Iterator<A, any, any>
| undefined
currentIter: interface Iterator<T, TReturn = any, TNext = any>Iterator<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A> | undefined
return {
async AsyncIterator<A, any, any>.next(...[value]: [] | [any]): Promise<IteratorResult<A, any>>next(): interface Promise<T>Represents the completion of an asynchronous operation
Promise<type IteratorResult<T, TReturn = any> =
| IteratorYieldResult<T>
| IteratorReturnResult<TReturn>
IteratorResult<function (type parameter) A in <A, E, XR, R extends XR>(self: Stream<A, E, R>, context: Context.Context<XR>): AsyncIterable<A>A>> {
if (let currentIter:
| Iterator<A, any, any>
| undefined
currentIter) {
const const next: IteratorResult<A, any>next = let currentIter: Iterator<A, any, any>currentIter.Iterator<A, any, any>.next(...[value]: [] | [any]): IteratorResult<A, any>next()
if (!const next: IteratorResult<A, any>next.done?: boolean | undefineddone) return const next: IteratorYieldResult<A>next
let currentIter:
| Iterator<A, any, any>
| undefined
currentIter = var undefinedundefined
}
let pull:
| Pull.Pull<
Arr.NonEmptyReadonlyArray<A>,
E,
void,
R
>
| undefined
pull ??= await const runPromise: (
effect: Effect<A, E, XR>,
options?: RunOptions | undefined
) => Promise<A>
runPromise(import ChannelChannel.toPullScoped(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.Stream<out A, out E = never, out R = never>.channel: Channel.Channel<Arr.NonEmptyReadonlyArray<A>, E, void, unknown, unknown, unknown, R>(property) Stream<out A, out E = never, out R = never>.channel: {
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; <…;
}
channel, const scope: Scope.Closeableconst scope: {
strategy: "sequential" | "parallel";
state: State.Open | State.Closed | State.Empty;
}
scope))
const const exit: Exit.Exit<
readonly [A, ...A[]],
Cause.Done<void> | E
>
exit = await const runPromiseExit: (
effect: Effect<A, E, XR>,
options?: RunOptions | undefined
) => Promise<Exit.Exit<A, E>>
runPromiseExit(let pull:
| Pull.Pull<
Arr.NonEmptyReadonlyArray<A>,
E,
void,
R
>
| undefined
let pull: {
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;
}
pull)
if (import ExitExit.isSuccess(const exit: Exit.Exit<
readonly [A, ...A[]],
Cause.Done<void> | E
>
exit)) {
let currentIter:
| Iterator<A, any, any>
| undefined
currentIter = const exit: Exit.Success<
readonly [A, ...A[]],
Cause.Done<void> | E
>
const exit: {
_tag: "Success";
value: A;
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;
}
exit.value[var Symbol: SymbolConstructorSymbol.SymbolConstructor.iterator: typeof Symbol.iteratorA method that returns the default iterator for an object. Called by the semantics of the
for-of statement.
iterator]()
return let currentIter:
| Iterator<A, any, any>
| undefined
currentIter.Iterator<A, any, any>.next(...[value]: [] | [any]): IteratorResult<A, any>next()
} else if (import PullPull.isDoneCause(const exit: Exit.Failure<
readonly [A, ...A[]],
Cause.Done<void> | E
>
const exit: {
_tag: "Failure";
cause: Cause.Cause<E>;
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;
}
exit.cause)) {
return { IteratorReturnResult<any>.done: truedone: true, IteratorReturnResult<any>.value: anyvalue: var undefinedundefined }
}
throw import CauseCause.squash(const exit: Exit.Failure<
readonly [A, ...A[]],
Cause.Done<void> | E
>
const exit: {
_tag: "Failure";
cause: Cause.Cause<E>;
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;
}
exit.cause)
},
AsyncIterator<A, any, any>.return?(value?: any): Promise<IteratorResult<A, any>>return(_: any_) {
return const runPromise: (
effect: Effect<A, E, XR>,
options?: RunOptions | undefined
) => Promise<A>
runPromise(import EffectEffect.as(
import ScopeScope.close(const scope: Scope.Closeableconst scope: {
strategy: "sequential" | "parallel";
state: State.Open | State.Closed | State.Empty;
}
scope, import ExitExit.void),
{ done: booleandone: true, value: undefinedvalue: var undefinedundefined }
))
}
}
}
})
)