<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.interface Effect<out A, out E = never, out R = never>The Effect interface defines a value that lazily describes a workflow or
job. The workflow requires some context R, and may fail with an error of
type E, or succeed with a value of type A.
When to use
Use when you need to represent a lazy, composable workflow that can require
services, fail with a typed error, or succeed with a typed value.
Details
Effect values model resourceful interaction with the outside world,
including synchronous, asynchronous, concurrent, and parallel interaction.
They use a fiber-based concurrency model, with built-in support for
scheduling, fine-grained interruption, structured concurrency, and high
scalability.
To run an Effect value, you need a Runtime, which is a type that is
capable of executing Effect values.
Effect<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.const map: {
<A, B>(f: (a: A) => B): <E, R>(
self: Effect<A, E, R>
) => Effect<B, E, R>
<A, E, R, B>(
self: Effect<A, E, R>,
f: (a: A) => B
): Effect<B, E, R>
}
map(
import EffectEffect.const context: <R = never>() => Effect<
Context.Context<R>,
never,
R
>
Returns the complete context.
When to use
Use to read the complete Context available to the current effect.
Details
This function allows you to access all services that are currently available
in the effect's environment. This can be useful for debugging, introspection,
or when you need to pass the entire context to another function.
Example (Reading the full context)
import { Console, Context, Effect, Option } from "effect"
const Logger = Context.Service<{
log: (msg: string) => void
}>("Logger")
const Database = Context.Service<{
query: (sql: string) => string
}>("Database")
const program = Effect.gen(function*() {
const allServices = yield* Effect.context()
// Check if specific services are available
const loggerOption = Context.getOption(allServices, Logger)
const databaseOption = Context.getOption(allServices, Database)
yield* Console.log(`Logger available: ${Option.isSome(loggerOption)}`)
yield* Console.log(`Database available: ${Option.isSome(databaseOption)}`)
})
const context = Context.make(Logger, { log: console.log })
.pipe(Context.add(Database, { query: () => "result" }))
const provided = Effect.provideContext(program, context)
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)
)