<
Arg extends
| Stream<Effect.Effect<any, any, any>, any, any>
| {
readonly concurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefined
}
| undefined = {
readonly concurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefined
}
>(
selfOrOptions?: Arg,
options?:
| {
readonly concurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefined
}
| undefined
): [Arg] extends [
Stream<
Effect.Effect<infer _A, infer _EX, infer _RX>,
infer _E,
infer _R
>
]
? Stream<_A, _EX | _E, _RX | _R>
: <A, EX, RX, E, R>(
self: Stream<Effect.Effect<A, EX, RX>, E, R>
) => Stream<A, EX | E, RX | R>Flattens a stream of Effect values into a stream of their results.
When to use
Use when stream elements already are effects and their successes should become stream elements while their failures enter the stream error channel.
Example (Flattening a stream of Effect values into a stream of their results)
import { Console, Effect, Stream } from "effect"
const stream = Stream.make(Effect.succeed(1), Effect.succeed(2), Effect.succeed(3))
const program = Effect.gen(function*() {
const result = yield* Stream.runCollect(stream.pipe(Stream.flattenEffect()))
yield* Console.log(result)
})
Effect.runPromise(program)
// Output: [1, 2, 3]export const const flattenEffect: <
Arg extends
| Stream<
Effect.Effect<any, any, any>,
any,
any
>
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined = {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
>(
selfOrOptions?: Arg,
options?:
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined
) => [Arg] extends [
Stream<
Effect.Effect<infer _A, infer _EX, infer _RX>,
infer _E,
infer _R
>
]
? Stream<_A, _EX | _E, _RX | _R>
: <A, EX, RX, E, R>(
self: Stream<Effect.Effect<A, EX, RX>, E, R>
) => Stream<A, EX | E, RX | R>
Flattens a stream of Effect values into a stream of their results.
When to use
Use when stream elements already are effects and their successes should become
stream elements while their failures enter the stream error channel.
Example (Flattening a stream of Effect values into a stream of their results)
import { Console, Effect, Stream } from "effect"
const stream = Stream.make(Effect.succeed(1), Effect.succeed(2), Effect.succeed(3))
const program = Effect.gen(function*() {
const result = yield* Stream.runCollect(stream.pipe(Stream.flattenEffect()))
yield* Console.log(result)
})
Effect.runPromise(program)
// Output: [1, 2, 3]
flattenEffect: <
function (type parameter) Arg in <Arg extends Stream<Effect.Effect<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined = {
...;
}>(selfOrOptions?: Arg, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): [Arg] extends [Stream<Effect.Effect<infer _A, infer _EX, infer _RX>, infer _E, infer _R>] ? Stream<_A, _EX | _E, _RX | _R> : <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>) => Stream<A, EX | E, RX | R>
Arg extends 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<import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefinedunordered?: boolean | undefined
} | undefined = {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefinedunordered?: boolean | undefined
}
>(
selfOrOptions: Arg | undefinedselfOrOptions?: function (type parameter) Arg in <Arg extends Stream<Effect.Effect<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined = {
...;
}>(selfOrOptions?: Arg, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): [Arg] extends [Stream<Effect.Effect<infer _A, infer _EX, infer _RX>, infer _E, infer _R>] ? Stream<_A, _EX | _E, _RX | _R> : <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>) => Stream<A, EX | E, RX | R>
Arg,
options: | {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined
options?: {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefinedunordered?: boolean | undefined
} | undefined
) => [function (type parameter) Arg in <Arg extends Stream<Effect.Effect<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined = {
...;
}>(selfOrOptions?: Arg, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): [Arg] extends [Stream<Effect.Effect<infer _A, infer _EX, infer _RX>, infer _E, infer _R>] ? Stream<_A, _EX | _E, _RX | _R> : <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>) => Stream<A, EX | E, RX | R>
Arg] extends [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<import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<infer function (type parameter) _A_A, infer function (type parameter) _EX_EX, infer function (type parameter) _RX_RX>, infer function (type parameter) _E_E, infer function (type parameter) _R_R>] ?
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_A, function (type parameter) _EX_EX | function (type parameter) _E_E, function (type parameter) _RX_RX | function (type parameter) _R_R>
: <function (type parameter) A in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>A, function (type parameter) EX in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>EX, function (type parameter) RX in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>RX, function (type parameter) E in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>E, function (type parameter) R in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>R>(self: Stream<Effect.Effect<A, EX, RX>, 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<import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<function (type parameter) A in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>A, function (type parameter) EX in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>EX, function (type parameter) RX in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>RX>, function (type parameter) E in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>E, function (type parameter) R in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>R>) => 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, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>A, function (type parameter) EX in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>EX | function (type parameter) E in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>E, function (type parameter) RX in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>RX | function (type parameter) R in <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>): Stream<A, EX | E, RX | R>R> = import dualdual(
(args: anyargs) => const isStream: (
u: unknown
) => u is Stream<unknown, unknown, unknown>
Checks whether a value is a Stream.
Example (Checking whether a value is a Stream)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const stream = Stream.make(1, 2, 3)
const notStream = { data: [1, 2, 3] }
yield* Console.log(Stream.isStream(stream))
// true
yield* Console.log(Stream.isStream(notStream))
// false
})
Effect.runPromise(program)
isStream(args: anyargs[0]),
<function (type parameter) A in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
A, function (type parameter) E in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
E, function (type parameter) R in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
R, function (type parameter) EX in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
EX, function (type parameter) RX in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
RX>(
self: Stream<Effect.Effect<A, EX, RX>, 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<import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<function (type parameter) A in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
A, function (type parameter) EX in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
EX, function (type parameter) RX in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
RX>, function (type parameter) E in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
E, function (type parameter) R in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
R>,
options: | {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined
options?: {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly unordered?: boolean | undefinedunordered?: boolean | undefined
} | 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, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
A, function (type parameter) EX in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
EX | function (type parameter) E in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
E, function (type parameter) RX in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
RX | function (type parameter) R in <A, E, R, EX, RX>(self: Stream<Effect.Effect<A, EX, RX>, E, R>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly unordered?: boolean | undefined;
} | undefined): Stream<A, EX | E, RX | R>
R> => const mapEffect: {
<A, A2, E2, R2>(
f: (
a: A,
i: number
) => Effect.Effect<A2, E2, R2>,
options?:
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined
): <E, R>(
self: Stream<A, E, R>
) => Stream<A2, E2 | E, R2 | R>
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (
a: A,
i: number
) => Effect.Effect<A2, E2, R2>,
options?:
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined
): Stream<A2, E | E2, R | R2>
}
mapEffect(self: Stream<Effect.Effect<A, EX, RX>, 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, import identityidentity, options: | {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly unordered?: boolean | undefined
}
| undefined
options)
)