<
Arg extends
| Stream<Stream<any, any, any>, any, any>
| {
readonly concurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefined
}
| undefined = {
readonly concurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefined
}
>(
selfOrOptions?: Arg,
options?:
| {
readonly concurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefined
}
| undefined
): [Arg] extends [
Stream<Stream<infer _A, infer _E, infer _R>, infer _E2, infer _R2>
]
? Stream<_A, _E | _E2, _R | _R2>
: <A, E, R, E2, R2>(
self: Stream<Stream<A, E, R>, E2, R2>
) => Stream<A, E | E2, R | R2>Flattens a stream of streams into a single stream.
Details
With the default sequential concurrency, inner streams are concatenated in
strict order. When concurrency is greater than 1 or "unbounded",
multiple inner streams may run at the same time and their outputs are merged
as they arrive.
Example (Flattening nested streams)
import { Console, Effect, Stream } from "effect"
const streamOfStreams = Stream.make(
Stream.make(1, 2),
Stream.make(3, 4),
Stream.make(5, 6)
)
const program = Effect.gen(function*() {
const values = yield* Stream.runCollect(Stream.flatten(streamOfStreams))
yield* Console.log(values)
})
Effect.runPromise(program)
// Output: [ 1, 2, 3, 4, 5, 6 ]export const const flatten: <
Arg extends
| Stream<Stream<any, any, any>, any, any>
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | undefined
}
| undefined = {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | undefined
}
>(
selfOrOptions?: Arg,
options?:
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | undefined
}
| undefined
) => [Arg] extends [
Stream<
Stream<infer _A, infer _E, infer _R>,
infer _E2,
infer _R2
>
]
? Stream<_A, _E | _E2, _R | _R2>
: <A, E, R, E2, R2>(
self: Stream<Stream<A, E, R>, E2, R2>
) => Stream<A, E | E2, R | R2>
Flattens a stream of streams into a single stream.
Details
With the default sequential concurrency, inner streams are concatenated in
strict order. When concurrency is greater than 1 or "unbounded",
multiple inner streams may run at the same time and their outputs are merged
as they arrive.
Example (Flattening nested streams)
import { Console, Effect, Stream } from "effect"
const streamOfStreams = Stream.make(
Stream.make(1, 2),
Stream.make(3, 4),
Stream.make(5, 6)
)
const program = Effect.gen(function*() {
const values = yield* Stream.runCollect(Stream.flatten(streamOfStreams))
yield* Console.log(values)
})
Effect.runPromise(program)
// Output: [ 1, 2, 3, 4, 5, 6 ]
flatten: <
function (type parameter) Arg in <Arg extends Stream<Stream<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined = {
...;
}>(selfOrOptions?: Arg, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): [Arg] extends [Stream<Stream<infer _A, infer _E, infer _R>, infer _E2, infer _R2>] ? Stream<_A, _E | _E2, _R | _R2> : <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>) => Stream<A, E | E2, R | R2>
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<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<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefinedbufferSize?: number | undefined
} | undefined = {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefinedbufferSize?: number | undefined
}
>(
selfOrOptions: Arg | undefinedselfOrOptions?: function (type parameter) Arg in <Arg extends Stream<Stream<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined = {
...;
}>(selfOrOptions?: Arg, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): [Arg] extends [Stream<Stream<infer _A, infer _E, infer _R>, infer _E2, infer _R2>] ? Stream<_A, _E | _E2, _R | _R2> : <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>) => Stream<A, E | E2, R | R2>
Arg,
options: | {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | undefined
}
| undefined
options?: {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefinedbufferSize?: number | undefined
} | undefined
) => [function (type parameter) Arg in <Arg extends Stream<Stream<any, any, any>, any, any> | {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined = {
...;
}>(selfOrOptions?: Arg, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): [Arg] extends [Stream<Stream<infer _A, infer _E, infer _R>, infer _E2, infer _R2>] ? Stream<_A, _E | _E2, _R | _R2> : <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>) => Stream<A, E | E2, R | R2>
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<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<infer function (type parameter) _A_A, infer function (type parameter) _E_E, infer function (type parameter) _R_R>, infer function (type parameter) _E2_E2, infer function (type parameter) _R2_R2>] ? 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) _E_E | function (type parameter) _E2_E2, function (type parameter) _R_R | function (type parameter) _R2_R2>
: <function (type parameter) A in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>A, function (type parameter) E in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>E, function (type parameter) R in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>R, function (type parameter) E2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>E2, function (type parameter) R2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>R2>(self: Stream<Stream<A, E, R>, E2, R2>(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<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, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>A, function (type parameter) E in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>E, function (type parameter) R in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>R>, function (type parameter) E2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>E2, function (type parameter) R2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>R2>) => 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, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>A, function (type parameter) E in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>E | function (type parameter) E2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>E2, function (type parameter) R in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>R | function (type parameter) R2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>): Stream<A, E | E2, R | R2>R2> = 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, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
E, function (type parameter) R in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
R, function (type parameter) E2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
E2, function (type parameter) R2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
R2>(
self: Stream<Stream<A, E, R>, E2, R2>(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<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, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
E, function (type parameter) R in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
R>, function (type parameter) E2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
E2, function (type parameter) R2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
R2>,
options: | {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | undefined
}
| undefined
options?: {
readonly concurrency?: number | "unbounded" | undefinedconcurrency?: number | "unbounded" | undefined
readonly bufferSize?: number | undefinedbufferSize?: number | 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, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
E | function (type parameter) E2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
E2, function (type parameter) R in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
R | function (type parameter) R2 in <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>, options?: {
readonly concurrency?: number | "unbounded" | undefined;
readonly bufferSize?: number | undefined;
} | undefined): Stream<A, E | E2, R | R2>
R2> => const flatMap: {
<A, A2, E2, R2>(
f: (a: A) => Stream<A2, E2, R2>,
options?:
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | 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) => Stream<A2, E2, R2>,
options?:
| {
readonly concurrency?:
| number
| "unbounded"
| undefined
readonly bufferSize?: number | undefined
}
| undefined
): Stream<A2, E | E2, R | R2>
}
flatMap(self: Stream<Stream<A, E, R>, E2, R2>(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 bufferSize?: number | undefined
}
| undefined
options)
)