<S, A, B, E2, R2>(
initial: LazyArg<S>,
f: (
s: S,
a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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 non-empty input chunk statefully and effectfully, emitting zero or more output values per chunk.
When to use
Use when stateful mapping should process each emitted non-empty chunk with an Effect instead of each element separately.
Details
The mapping effect receives the current state and chunk, then returns the next state plus the values to emit. The state is threaded across chunks.
Example (Effectfully mapping stream chunks with state)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const totals = yield* Stream.make(1, 2, 3, 4).pipe(
Stream.rechunk(2),
Stream.mapAccumArrayEffect(() => 0, (total, chunk) =>
Effect.gen(function*() {
const next = chunk.reduce((sum, value) => sum + value, total)
return [next, [next]] as const
})
),
Stream.runCollect
)
yield* Console.log(totals)
})
Effect.runPromise(program)
// Output: [ 3, 10 ]export const const mapAccumArrayEffect: {
<S, A, B, E2, R2>(
initial: LazyArg<S>,
f: (
s: S,
a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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 non-empty input chunk statefully and effectfully, emitting zero or
more output values per chunk.
When to use
Use when stateful mapping should process each emitted non-empty chunk with an
Effect instead of each element separately.
Details
The mapping effect receives the current state and chunk, then returns the
next state plus the values to emit. The state is threaded across chunks.
Example (Effectfully mapping stream chunks with state)
import { Console, Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const totals = yield* Stream.make(1, 2, 3, 4).pipe(
Stream.rechunk(2),
Stream.mapAccumArrayEffect(() => 0, (total, chunk) =>
Effect.gen(function*() {
const next = chunk.reduce((sum, value) => sum + value, total)
return [next, [next]] as const
})
),
Stream.runCollect
)
yield* Console.log(totals)
})
Effect.runPromise(program)
// Output: [ 3, 10 ]
mapAccumArrayEffect: {
<function (type parameter) S in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<A>(parameter) a: {
0: A;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<A>>): Array<A>; (...items: Array<A | ConcatArray<A>>): Array<A> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<A>;
indexOf: (searchElement: A, fromIndex?: number) => number;
lastIndexOf: (searchElement: A, fromIndex?: number) => number;
every: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => void, thisArg?: any) => void;
map: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): Array<S>; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): Array<A> };
reduce: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
reduceRight: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
find: { (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findIndex: (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, A]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<A>;
includes: (searchElement: A, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: A, index: number, array: Array<A>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => A | undefined;
findLast: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findLastIndex: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
toReversed: () => Array<A>;
toSorted: (compareFn?: ((a: A, b: A) => number) | undefined) => Array<A>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<A>): Array<A>; (start: number, deleteCount?: number): Array<A> };
with: (index: number, value: A) => Array<A>;
}
a: import ArrArr.type Arr.NonEmptyReadonlyArray = /*unresolved*/ anyNonEmptyReadonlyArray<function (type parameter) A in <S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<A>(parameter) a: {
0: A;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<A>>): Array<A>; (...items: Array<A | ConcatArray<A>>): Array<A> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<A>;
indexOf: (searchElement: A, fromIndex?: number) => number;
lastIndexOf: (searchElement: A, fromIndex?: number) => number;
every: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => void, thisArg?: any) => void;
map: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): Array<S>; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): Array<A> };
reduce: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
reduceRight: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
find: { (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findIndex: (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, A]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<A>;
includes: (searchElement: A, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: A, index: number, array: Array<A>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => A | undefined;
findLast: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findLastIndex: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
toReversed: () => Array<A>;
toSorted: (compareFn?: ((a: A, b: A) => number) | undefined) => Array<A>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<A>): Array<A>; (start: number, deleteCount?: number): Array<A> };
with: (index: number, value: A) => Array<A>;
}
a: import ArrArr.type Arr.NonEmptyReadonlyArray = /*unresolved*/ anyNonEmptyReadonlyArray<function (type parameter) A in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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), <function (type parameter) A in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<A>(parameter) a: {
0: A;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<A>>): Array<A>; (...items: Array<A | ConcatArray<A>>): Array<A> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<A>;
indexOf: (searchElement: A, fromIndex?: number) => number;
lastIndexOf: (searchElement: A, fromIndex?: number) => number;
every: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => void, thisArg?: any) => void;
map: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): Array<S>; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): Array<A> };
reduce: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
reduceRight: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
find: { (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findIndex: (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, A]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<A>;
includes: (searchElement: A, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: A, index: number, array: Array<A>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => A | undefined;
findLast: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findLastIndex: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
toReversed: () => Array<A>;
toSorted: (compareFn?: ((a: A, b: A) => number) | undefined) => Array<A>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<A>): Array<A>; (start: number, deleteCount?: number): Array<A> };
with: (index: number, value: A) => Array<A>;
}
a: import ArrArr.type Arr.NonEmptyReadonlyArray = /*unresolved*/ anyNonEmptyReadonlyArray<function (type parameter) A in <A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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: Arr.NonEmptyReadonlyArray<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.mapAccum(
initial: LazyArg<S>initial,
(state: anystate, a: any(parameter) a: {
0: A;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<A>>): Array<A>; (...items: Array<A | ConcatArray<A>>): Array<A> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<A>;
indexOf: (searchElement: A, fromIndex?: number) => number;
lastIndexOf: (searchElement: A, fromIndex?: number) => number;
every: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => void, thisArg?: any) => void;
map: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): Array<S>; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): Array<A> };
reduce: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
reduceRight: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
find: { (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findIndex: (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, A]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<A>;
includes: (searchElement: A, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: A, index: number, array: Array<A>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => A | undefined;
findLast: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findLastIndex: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
toReversed: () => Array<A>;
toSorted: (compareFn?: ((a: A, b: A) => number) | undefined) => Array<A>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<A>): Array<A>; (start: number, deleteCount?: number): Array<A> };
with: (index: number, value: A) => Array<A>;
}
a) =>
import EffectEffect.map(
f: (
s: S,
a: Arr.NonEmptyReadonlyArray<A>
) => Effect.Effect<
readonly [state: S, values: ReadonlyArray<B>],
E2,
R2
>
f(state: anystate, a: any(parameter) a: {
0: A;
length: number;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<A>>): Array<A>; (...items: Array<A | ConcatArray<A>>): Array<A> };
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<A>;
indexOf: (searchElement: A, fromIndex?: number) => number;
lastIndexOf: (searchElement: A, fromIndex?: number) => number;
every: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): boolean };
some: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => void, thisArg?: any) => void;
map: (callbackfn: (value: A, index: number, array: ReadonlyArray<A>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): Array<S>; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): Array<A> };
reduce: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
reduceRight: { (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A): A; (callbackfn: (previousValue: A, currentValue: A, currentIndex: number, array: ReadonlyArray<A>) => A, initialValue: A): A; (callbac…;
find: { (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findIndex: (predicate: (value: A, index: number, obj: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, A]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<A>;
includes: (searchElement: A, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: A, index: number, array: Array<A>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => A | undefined;
findLast: { (predicate: (value: A, index: number, array: ReadonlyArray<A>) => value is S, thisArg?: any): S | undefined; (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any): A | undefined };
findLastIndex: (predicate: (value: A, index: number, array: ReadonlyArray<A>) => unknown, thisArg?: any) => number;
toReversed: () => Array<A>;
toSorted: (compareFn?: ((a: A, b: A) => number) | undefined) => Array<A>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<A>): Array<A>; (start: number, deleteCount?: number): Array<A> };
with: (index: number, value: A) => Array<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) : const emptyArr: anyemptyArr
]
),
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
))