<S, A, B, E2, R2>(
initial: LazyArg<S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>,
options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined
}
): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
<A, E, R, S, B, E2, R2>(
self: Stream<A, E, R>,
initial: LazyArg<S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>,
options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined
}
): Stream<B, E | E2, R | R2>Maps each element statefully and effectfully, emitting zero or more output values per input.
When to use
Use when stateful element mapping needs Effects or can fail while emitting zero or more values per input element.
Details
The mapping effect receives the current state and element, then returns the next state plus the values to emit. The state is threaded through the stream.
Example (Effectfully mapping stream values with state)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const result = yield* Stream.make(1, 1, 1).pipe(
Stream.mapAccumEffect(() => 0, (total, n) =>
Effect.succeed([total + n, [total + n]])
),
Stream.runCollect
)
yield* Console.log(result)
})
Effect.runPromise(program)
// Output: [ 1, 2, 3 ]export const const mapAccumEffect: {
<S, A, B, E2, R2>(
initial: LazyArg<S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [
state: S,
values: ReadonlyArray<B>
],
E2,
R2
>,
options?: {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
): <E, R>(
self: Stream<A, E, R>
) => Stream<B, E | E2, R | R2>
<A, E, R, S, B, E2, R2>(
self: Stream<A, E, R>,
initial: LazyArg<S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [
state: S,
values: ReadonlyArray<B>
],
E2,
R2
>,
options?: {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
): Stream<B, E | E2, R | R2>
}
Maps each element statefully and effectfully, emitting zero or more output
values per input.
When to use
Use when stateful element mapping needs Effects or can fail while emitting
zero or more values per input element.
Details
The mapping effect receives the current state and element, then returns the
next state plus the values to emit. The state is threaded through the
stream.
Example (Effectfully mapping stream values with state)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const result = yield* Stream.make(1, 1, 1).pipe(
Stream.mapAccumEffect(() => 0, (total, n) =>
Effect.succeed([total + n, [total + n]])
),
Stream.runCollect
)
yield* Console.log(result)
})
Effect.runPromise(program)
// Output: [ 1, 2, 3 ]
mapAccumEffect: {
<function (type parameter) S in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
S, function (type parameter) A in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
A, function (type parameter) B in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
B, function (type parameter) E2 in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
E2, function (type parameter) R2 in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
R2>(
initial: LazyArg<S>initial: import LazyArgLazyArg<function (type parameter) S in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>
f: (s: Ss: function (type parameter) S in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
S, a: Aa: function (type parameter) A in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
A) => import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<readonly [Sstate: function (type parameter) S in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
S, readonly B[]values: interface ReadonlyArray<T>ReadonlyArray<function (type parameter) B in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
B>], function (type parameter) E2 in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
E2, function (type parameter) R2 in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
R2>,
options: | {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
| undefined
options?: {
readonly onHalt?: | ((state: S) => ReadonlyArray<B>)
| undefined
onHalt?: ((state: Sstate: function (type parameter) S in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
S) => interface ReadonlyArray<T>ReadonlyArray<function (type parameter) B in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
B>) | undefined
}
): <function (type parameter) E in <E, R>(self: Stream<A, E, R>): Stream<B, E | E2, R | R2>E, function (type parameter) R in <E, R>(self: Stream<A, E, R>): Stream<B, E | E2, R | R2>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 <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
A, function (type parameter) E in <E, R>(self: Stream<A, E, R>): Stream<B, E | E2, R | R2>E, function (type parameter) R in <E, R>(self: Stream<A, E, R>): Stream<B, E | E2, R | R2>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) B in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
B, function (type parameter) E in <E, R>(self: Stream<A, E, R>): Stream<B, E | E2, R | R2>E | function (type parameter) E2 in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
E2, function (type parameter) R in <E, R>(self: Stream<A, E, R>): Stream<B, E | E2, R | R2>R | function (type parameter) R2 in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E | E2, R | R2>
R2>
<function (type parameter) A in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E, function (type parameter) R in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R, function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S, function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B, function (type parameter) E2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E2, function (type parameter) R2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R2>(
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, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E, function (type parameter) R in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R>,
initial: LazyArg<S>initial: import LazyArgLazyArg<function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>
f: (s: Ss: function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S, a: Aa: function (type parameter) A in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
A) => import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<readonly [Sstate: function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S, readonly B[]values: interface ReadonlyArray<T>ReadonlyArray<function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B>], function (type parameter) E2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E2, function (type parameter) R2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R2>,
options: | {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
| undefined
options?: {
readonly onHalt?: | ((state: S) => ReadonlyArray<B>)
| undefined
onHalt?: ((state: Sstate: function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S) => interface ReadonlyArray<T>ReadonlyArray<function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B>) | 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) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B, function (type parameter) E in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E | function (type parameter) E2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E2, function (type parameter) R in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R | function (type parameter) R2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, 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, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E, function (type parameter) R in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R, function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S, function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B, function (type parameter) E2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E2, function (type parameter) R2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R2>(
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, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
A, function (type parameter) E in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E, function (type parameter) R in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R>,
initial: LazyArg<S>initial: import LazyArgLazyArg<function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S>,
f: (
s: S,
a: A
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>
f: (s: Ss: function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S, a: Aa: function (type parameter) A in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
A) => import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<readonly [Sstate: function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S, readonly B[]values: interface ReadonlyArray<T>ReadonlyArray<function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B>], function (type parameter) E2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E2, function (type parameter) R2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R2>,
options: | {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
| undefined
options?: {
readonly onHalt?: | ((state: S) => ReadonlyArray<B>)
| undefined
onHalt?: ((state: Sstate: function (type parameter) S in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
S) => interface ReadonlyArray<T>ReadonlyArray<function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B>) | 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) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B, function (type parameter) E in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E | function (type parameter) E2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
E2, function (type parameter) R in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R | function (type parameter) R2 in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
R2> =>
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.pipe(
import ChannelChannel.flattenArray,
import ChannelChannel.mapAccum(
initial: LazyArg<S>initial,
(state: anystate, a: NoInfer<A>a) =>
import EffectEffect.map(
f: (
s: S,
a: A
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>
f(state: anystate, a: NoInfer<A>a),
([state: anystate, values: anyvalues]) => [
state: anystate,
import ArrArr.isReadonlyArrayNonEmpty(values: anyvalues) ? import ArrArr.of(values: any(parameter) values: {
0: B;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<B>>): Array<B>; (...items: Array<B | ConcatArray<B>>): Array<B> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<B>;
indexOf: (searchElement: B, fromIndex?: number) => number;
lastIndexOf: (searchElement: B, fromIndex?: number) => number;
every: { (predicate: (value: B, index: number, array: ReadonlyArray<B>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: B, index: number, array: ReadonlyArray<B>) => void, thisArg?: any) => void;
map: (callbackfn: (value: B, index: number, array: ReadonlyArray<B>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: B, index: number, array: ReadonlyArray<B>) => value is S, thisArg?: any): Array<S>; (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any): Array<B> };
reduce: { (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B): B; (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B, initialValue: B): B; (callbac…;
reduceRight: { (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B): B; (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B, initialValue: B): B; (callbac…;
find: { (predicate: (value: B, index: number, obj: ReadonlyArray<B>) => value is S, thisArg?: any): S | undefined; (predicate: (value: B, index: number, obj: ReadonlyArray<B>) => unknown, thisArg?: any): B | undefined };
findIndex: (predicate: (value: B, index: number, obj: ReadonlyArray<B>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, B]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<B>;
includes: (searchElement: B, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: B, index: number, array: Array<B>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => B | undefined;
findLast: { (predicate: (value: B, index: number, array: ReadonlyArray<B>) => value is S, thisArg?: any): S | undefined; (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any): B | undefined };
findLastIndex: (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any) => number;
toReversed: () => Array<B>;
toSorted: (compareFn?: ((a: B, b: B) => number) | undefined) => Array<B>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<B>): Array<B>; (start: number, deleteCount?: number): Array<B> };
with: (index: number, value: B) => Array<B>;
}
values) : import ArrArr.empty<import ArrArr.type Arr.NonEmptyReadonlyArray = /*unresolved*/ anyNonEmptyReadonlyArray<function (type parameter) B in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect.Effect<readonly [state: S, values: ReadonlyArray<B>], E2, R2>, options?: {
readonly onHalt?: ((state: S) => ReadonlyArray<B>) | undefined;
}): Stream<B, E | E2, R | R2>
B>>()
]
),
options: | {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
| undefined
options?.onHalt?: | ((state: S) => ReadonlyArray<B>)
| undefined
onHalt ?
{
function onHalt(state: any): anyonHalt(state: anystate) {
const const arr: readonly B[]arr = options: {
readonly onHalt?:
| ((state: S) => ReadonlyArray<B>)
| undefined
}
options.onHalt?: | ((state: S) => ReadonlyArray<B>)
| undefined
onHalt!(state: anystate)
return import ArrArr.isReadonlyArrayNonEmpty(const arr: readonly B[]arr) ? import ArrArr.of(const arr: readonly B[]const arr: {
0: B;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<B>>): Array<B>; (...items: Array<B | ConcatArray<B>>): Array<B> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<B>;
indexOf: (searchElement: B, fromIndex?: number) => number;
lastIndexOf: (searchElement: B, fromIndex?: number) => number;
every: { (predicate: (value: B, index: number, array: ReadonlyArray<B>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: B, index: number, array: ReadonlyArray<B>) => void, thisArg?: any) => void;
map: (callbackfn: (value: B, index: number, array: ReadonlyArray<B>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: B, index: number, array: ReadonlyArray<B>) => value is S, thisArg?: any): Array<S>; (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any): Array<B> };
reduce: { (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B): B; (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B, initialValue: B): B; (callbac…;
reduceRight: { (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B): B; (callbackfn: (previousValue: B, currentValue: B, currentIndex: number, array: ReadonlyArray<B>) => B, initialValue: B): B; (callbac…;
find: { (predicate: (value: B, index: number, obj: ReadonlyArray<B>) => value is S, thisArg?: any): S | undefined; (predicate: (value: B, index: number, obj: ReadonlyArray<B>) => unknown, thisArg?: any): B | undefined };
findIndex: (predicate: (value: B, index: number, obj: ReadonlyArray<B>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, B]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<B>;
includes: (searchElement: B, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: B, index: number, array: Array<B>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => B | undefined;
findLast: { (predicate: (value: B, index: number, array: ReadonlyArray<B>) => value is S, thisArg?: any): S | undefined; (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any): B | undefined };
findLastIndex: (predicate: (value: B, index: number, array: ReadonlyArray<B>) => unknown, thisArg?: any) => number;
toReversed: () => Array<B>;
toSorted: (compareFn?: ((a: B, b: B) => number) | undefined) => Array<B>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<B>): Array<B>; (start: number, deleteCount?: number): Array<B> };
with: (index: number, value: B) => Array<B>;
}
arr) : const emptyArr: anyemptyArr
}
} :
var undefinedundefined
),
const fromChannel: <
Arr extends Arr.NonEmptyReadonlyArray<any>,
E,
R
>(
channel: Channel.Channel<
Arr,
E,
void,
unknown,
unknown,
unknown,
R
>
) => Stream<
Arr extends Arr.NonEmptyReadonlyArray<infer A>
? A
: never,
E,
R
>
Creates a stream from a array-emitting Channel.
Example (Creating a stream from an array-emitting channel)
import { Channel, Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const channel = Channel.succeed([1, 2, 3] as const)
const stream = Stream.fromChannel(channel)
const result = yield* Stream.runCollect(stream)
yield* Console.log(result)
})
// Output: [ 1, 2, 3 ]
fromChannel
))