Stream
Describes effectful sources that emit values over time.
A Stream<A, E, R> can emit many A values, fail with E, and require services R while it is being consumed. Streams are useful for data that is pulled in steps, such as values from collections, queues, pubsubs, schedules, callbacks, async iterables, or platform streams. The APIs here cover the full stream lifecycle: create a stream, transform or combine it, control buffering and timing, handle failures, and finally consume it.
Accessors
Signature
declare function service<I, S>(service: Key<I, S>): Stream<S, never, I>;serviceOption
Optionally accesses a service from the context and emits the result as a single element.
When to use
Use when you need a stream that emits an optional service from the context without requiring that service to be present.
Signature
declare function serviceOption<I, S>(service: Key<I, S>): Stream<Option<S>>;Accumulation
accumulate
Accumulates elements into a growing array, emitting the cumulative array for each input chunk.
Signature
declare function accumulate<A, E, R>(self: Stream<A, E, R>): Stream<[A, ...Array<A>], E, R>;Collects all elements into an array and emits it as a single element.
Signature
declare function collect<A, E, R>(self: Stream<A, E, R>): Stream<Array<A>, E, R>;Accumulates state across the stream, emitting the initial state and each updated state.
Signature
declare const scan: {
<S, A>(initial: S, f: (s: S, a: A) => S): <E, R>(self: Stream<A, E, R>) => Stream<S, E, R>;
<A, E, R, S>(self: Stream<A, E, R>, initial: S, f: (s: S, a: A) => S): Stream<S, E, R>;
};scanEffect
Accumulates state effectfully and emits the initial state plus each accumulated state.
Signature
declare const scanEffect: {
<S, A, E2, R2>(
initial: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<S, E2 | E, R2 | R>;
<A, E, R, S, E2, R2>(
self: Stream<A, E, R>,
initial: S,
f: (s: S, a: A) => Effect<S, E2, R2>,
): Stream<S, E | E2, R | R2>;
};Aggregation
Aggregates elements using the provided sink and emits each sink result as a stream element.
Details
The stream runs the upstream and downstream in separate fibers, so the sink can keep consuming input while downstream is busy processing the previous output.
Signature
declare const aggregate: {
<B, A, A2, E2, R2>(
sink: Sink<B, A | A2, A2, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R>;
<A, E, R, B, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<B, A | A2, A2, E2, R2>,
): Stream<B, E | E2, R | R2>;
};aggregateWithin
Aggregates elements with a sink, emitting each result when the sink completes or the schedule triggers.
Details
The schedule can flush the current aggregation even if the sink has not finished.
Signature
declare const aggregateWithin: {
<B, A, A2, E2, R2, C, E3, R3>(
sink: Sink<B, A | A2, A2, E2, R2>,
schedule: Schedule<C, Option<B>, E3, R3>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E3 | E, R2 | R3 | R>;
<A, E, R, B, A2, E2, R2, C, E3, R3>(
self: Stream<A, E, R>,
sink: Sink<B, A | A2, A2, E2, R2>,
schedule: Schedule<C, Option<B>, E3, R3>,
): Stream<B, E | E2 | E3, R | R2 | R3>;
};Applies a sink transducer to the stream and emits each sink result.
Signature
declare const transduce: <A2, A, E2, R2>(sink: Sink<A2, A, A, E2, R2>) => <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R> & <A, E, R, A2, E2, R2>(self: Stream<A, E, R>, sink: Sink<A2, A, A, E2, R2>) => Stream<A2, E | E2, R | R2>Broadcasting
Creates a PubSub-backed stream that multicasts the source to all subscribers.
Details
The returned stream is scoped and uses the provided PubSub capacity and replay settings.
Signature
declare const broadcast: {
(
options:
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<Stream<A, E, never>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<Stream<A, E, never>, never, Scope | R>;
};broadcastN
Creates a fixed-size tuple of streams that each emit the same elements as the source stream.
Details
The source stream starts after all downstream streams have been subscribed. With the default suspend strategy, the source can only advance capacity chunks ahead of the slowest downstream stream. If a downstream stream is interrupted, it unsubscribes from the broadcast so it no longer contributes backpressure.
Signature
declare const broadcastN: {
<N extends number>(
options:
| {
readonly capacity: "unbounded";
readonly n: N;
readonly replay?: number;
}
| {
readonly capacity: number;
readonly n: N;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<TupleOf<N, Stream<A, E, never>>, never, Scope | R>;
<A, E, R, N extends number>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
readonly n: N;
readonly replay?: number;
}
| {
readonly capacity: number;
readonly n: N;
readonly replay?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Effect<TupleOf<N, Stream<A, E, never>>, never, Scope | R>;
};Buffering
Buffers up to capacity elements so a faster producer can progress independently of a slower consumer.
Details
Finite buffers use the configured queue strategy: "suspend" applies backpressure, while "dropping" and "sliding" may discard elements when the buffer is full. This combinator destroys chunking; use Stream.rechunk afterward if you need fixed chunk sizes.
Signature
declare const buffer: {
(
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Stream<A, E, R>;
};bufferArray
Allows a faster producer to progress independently of a slower consumer by buffering up to capacity chunks in a queue.
Details
Finite buffers use the configured queue strategy: "suspend" applies backpressure, while "dropping" and "sliding" may discard chunks when the buffer is full. This combinator preserves chunking and is best with power-of-2 capacities.
Signature
declare const bufferArray: {
(
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Stream<A, E, R>;
};Constants
DefaultChunkSize
The default chunk size used by Stream constructors and combinators.
Signature
declare const DefaultChunkSize: number;Constructors
Creates a stream from a callback that can emit values into a queue.
When to use
Use when you need callback-based code to emit stream values by offering to a Queue, or signal stream completion through the Queue module APIs.
By default it uses an "unbounded" buffer size. You can customize the buffer size and strategy by passing an object as the second argument with the bufferSize and strategy fields.
Signature
declare function callback<A, E = never, R = never>(
f: (queue: Queue<A, Done<void> | E>) => Effect<unknown, E, Scope | R>,
options?: {
readonly bufferSize?: number;
readonly strategy?: "sliding" | "dropping" | "suspend";
},
): Stream<A, E, Exclude<R, Scope>>;The stream that dies with the specified defect.
Signature
declare function die(defect: unknown): Stream<never>;Provides the entry point for do-notation style stream composition.
Signature
declare const Do: Stream<{}>;Creates an empty stream.
Signature
declare const empty: Stream<never>;Terminates with the specified error.
Signature
declare function fail<E>(error: E): Stream<never, E>;Creates a stream that fails with the specified Cause.
Signature
declare function failCause<E>(cause: Cause<E>): Stream<never, E>;failCauseSync
The stream that always fails with the specified lazily evaluated Cause.
Signature
declare function failCauseSync<E>(evaluate: LazyArg<Cause<E>>): Stream<never, E>;Terminates with the specified lazily evaluated error.
Signature
declare function failSync<E>(evaluate: LazyArg<E>): Stream<never, E>;Creates a stream from an array of values.
Signature
declare function fromArray<A>(array: readonly Array<A>): Stream<A>fromArrayEffect
Creates a stream from an effect that produces an array of values.
When to use
Use when the array must be acquired from an Effect before the stream emits, and acquisition services or failures should be part of the stream.
Signature
declare function fromArrayEffect<A, E, R>(effect: Effect<readonly Array<A>, E, R>): Stream<A, Exclude<E, Done<any>>, R>fromArrays
Creates a stream from an arbitrary number of arrays.
Signature
declare function fromArrays<Arr extends readonly Array<readonly Array<any>>>(...arrays: Arr): Stream<Arr[number][number]>fromAsyncIterable
Creates a stream from an AsyncIterable.
Signature
declare function fromAsyncIterable<A, E>(
iterable: AsyncIterable<A>,
onError: (error: unknown) => E,
): Stream<A, E>;fromChannel
Creates a stream from a array-emitting Channel.
Signature
declare 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>;fromEffect
Creates a stream from an effect.
Signature
declare function fromEffect<A, E, R>(effect: Effect<A, E, R>): Stream<A, E, R>;fromEffectDrain
Creates a stream that runs the effect and emits no elements.
Signature
declare function fromEffectDrain<A, E, R>(effect: Effect<A, E, R>): Stream<never, E, R>;fromEffectRepeat
Creates a stream from an effect producing a value of type A which repeats forever.
Signature
declare function fromEffectRepeat<A, E, R>(
effect: Effect<A, E, R>,
): Stream<A, Exclude<E, Done<any>>, R>;fromEffectSchedule
Creates a stream from an effect producing a value of type A, which is repeated using the specified schedule.
Signature
declare function fromEffectSchedule<A, E, R, X, AS, ES, RS>(
effect: Effect<A, E, R>,
schedule: Schedule<X, AS, ES, RS>,
): Stream<A, E | ES, R | RS>;fromEventListener
Creates a stream from an event listener.
Signature
declare function fromEventListener<A = unknown>(
target: EventListener<A>,
type: string,
options?:
| boolean
| {
readonly bufferSize?: number;
readonly capture?: boolean;
readonly once?: boolean;
readonly passive?: boolean;
},
): Stream<A>;fromIterable
Creates a new Stream from an iterable collection of values.
Details
- chunkSize: Maximum number of values emitted per chunk.
Signature
declare function fromIterable<A>(
iterable: Iterable<A>,
options?: {
readonly chunkSize?: number;
},
): Stream<A>;fromIterableEffect
Creates a stream from an effect producing an iterable of values.
When to use
Use when the iterable must be acquired from an Effect before the stream emits, and acquisition services or failures should be part of the stream.
Signature
declare function fromIterableEffect<A, E, R>(
iterable: Effect<Iterable<A, any, any>, E, R>,
): Stream<A, E, R>;fromIterableEffectRepeat
Creates a stream by repeatedly running an effect that yields an iterable of values.
Signature
declare function fromIterableEffectRepeat<A, E, R>(
iterable: Effect<Iterable<A, any, any>, E, R>,
): Stream<A, Exclude<E, Done<any>>, R>;fromIteratorSucceed
Creates a stream that consumes values from an iterator.
Details
The maxChunkSize parameter controls how many values are pulled per chunk.
Signature
declare function fromIteratorSucceed<A>(
iterator: IterableIterator<A>,
maxChunkSize?: number,
): Stream<A>;fromPubSub
Creates a stream from a subscription to a PubSub.
Signature
declare function fromPubSub<A>(pubsub: PubSub<A>): Stream<A>;fromPubSubTake
Creates a stream from a PubSub of Take values.
Details
Take values include end and failure signals.
Signature
declare function fromPubSubTake<A, E>(pubsub: PubSub<Take<A, E, void>>): Stream<A, E>;Creates a stream from a pull effect, such as one produced by Stream.toPull.
Details
A pull effect yields chunks on demand and completes when the upstream stream ends. See Stream.toPull for a matching producer.
Signature
declare function fromPull<A, E, R, EX, RX>(
pull: Effect<Pull<readonly [A, A], E, void, R>, EX, RX>,
): Stream<A, EX | Exclude<E, Done<any>>, R | RX>;Creates a stream that pulls values from a Queue.Dequeue.
Details
The stream emits non-empty batches of queued values and ends when the queue fails with Cause.Done; other queue failures are propagated.
Signature
declare function fromQueue<A, E>(queue: Dequeue<A, E>): Stream<A, Exclude<E, Done<void>>>;fromReadableStream
Creates a stream from a lazily supplied Web ReadableStream.
Details
The stream reads from a ReadableStreamDefaultReader, maps read failures with onError, and closes the reader when the stream finalizes. By default the reader is canceled; set releaseLockOnEnd to release the lock instead.
Signature
declare function fromReadableStream<A, E>(options: {
readonly evaluate: LazyArg<ReadableStream<A>>;
readonly onError: (error: unknown) => E;
readonly releaseLockOnEnd?: boolean;
}): Stream<A, E>;fromSchedule
Creates a stream that emits each output of a schedule that does not require input, for as long as the schedule continues.
Signature
declare function fromSchedule<O, E, R>(schedule: Schedule<O, unknown, E, R>): Stream<O, E, R>;fromSubscription
Creates a stream from a PubSub subscription.
When to use
Use when you already have a PubSub.Subscription and want to expose its messages as a Stream, with Stream.take or cancellation controlling how many values are consumed.
Signature
declare function fromSubscription<A>(pubsub: Subscription<A>): Stream<A>;Creates an infinite stream by repeatedly applying a function to a seed value.
Signature
declare function iterate<A>(value: A, next: (value: A) => A): Stream<A>;Creates a stream from a sequence of values.
Signature
declare function make<As extends readonly Array<any>>(...values: As): Stream<As[number]>The stream that never produces any value or fails with any error.
Signature
declare const never: Stream<never>;Creates a stream by repeatedly evaluating an effectful page function.
When to use
Use to consume paginated APIs where each step returns a batch of values together with an optional next state.
Details
This is similar to unfold, but each step can emit zero or more values and independently decide whether another state should be requested.
Signature
declare function paginate<S, A, E = never, R = never>(s: S, f: (s: S) => Effect<readonly [readonly Array<A>, Option<S>], E, R>): Stream<A, E, R>Constructs a stream from a range of integers, including both endpoints.
Details
If the provided min is greater than max, the stream will not emit any values.
Signature
declare function range(min: number, max: number, chunkSize: number): Stream<number>;Runs a stream that requires Scope in a managed scope, ensuring its finalizers are run when the stream completes.
Signature
declare function scoped<A, E, R>(self: Stream<A, E, R>): Stream<A, E, Exclude<R, Scope>>;Creates a single-valued pure stream.
Signature
declare function succeed<A>(value: A): Stream<A>;Creates a lazily constructed stream.
Details
The stream factory is evaluated each time the stream is run.
Signature
declare function suspend<A, E, R>(stream: LazyArg<Stream<A, E, R>>): Stream<A, E, R>;Creates a stream that synchronously evaluates a function and emits the result as a single value.
Details
The function is evaluated each time the stream is run.
Signature
declare function sync<A>(evaluate: LazyArg<A>): Stream<A>;Creates a stream that emits void immediately once, then emits another void after each specified interval.
Signature
declare function tick(interval: Input): Stream<void>;Creates a channel from a stream.
Signature
declare function toChannel<A, E, R>(
stream: Stream<A, E, R>,
): Channel<readonly [A, A], E, void, unknown, unknown, unknown, R>;transformPull
Derives a stream by transforming its pull effect.
Signature
declare function transformPull<A, E, R, B, E2, R2, EX, RX>(
self: Stream<A, E, R>,
f: (
pull: Pull<readonly [A, A], E, void>,
scope: Scope,
) => Effect<Pull<readonly [B, B], E2, void, R2>, EX, RX>,
): Stream<B, EX | Exclude<E2, Done<any>>, R | R2 | RX>;transformPullBracket
Transforms a stream by effectfully transforming its pull effect.
Details
A forked scope is also provided to the transformation function, which is closed once the resulting stream has finished processing.
Signature
declare function transformPullBracket<A, E, R, B, E2, R2, EX, RX>(
self: Stream<A, E, R>,
f: (
pull: Pull<readonly [A, A], E, void, R>,
scope: Scope,
forkedScope: Scope,
) => Effect<Pull<readonly [B, B], E2, void, R2>, EX, RX>,
): Stream<B, EX | Exclude<E2, Done<any>>, R | R2 | RX>;Creates a stream by repeatedly applying an effectful step function to a state.
Details
Each readonly [value, nextState] result emits value and continues with nextState; returning undefined ends the stream.
Signature
declare function unfold<S, A, E, R>(
s: S,
f: (s: S) => Effect<readonly [A, S] | undefined, E, R>,
): Stream<A, E, R>;Creates a stream produced from an Effect.
Signature
declare function unwrap<A, E2, R2, E, R>(
effect: Effect<Stream<A, E2, R2>, E, R>,
): Stream<A, E2 | E, R2 | Exclude<R, Scope>>;Decoding
decodeText
Decodes Uint8Array chunks into strings using TextDecoder with an optional encoding.
Signature
declare const decodeText: <
Arg extends
| Stream<Uint8Array, any, any>
| {
readonly encoding?: string;
}
| undefined = {
readonly encoding?: string;
},
>(
streamOrOptions?: Arg,
options?: {
readonly encoding?: string;
},
) => [Arg] extends [Stream<Uint8Array, infer _E, infer _R>]
? Stream<string, _E, _R>
: <E, R>(self: Stream<Uint8Array, E, R>) => Stream<string, E, R>;Deduplication
Emits only elements that differ from the previous one.
Signature
declare function changes<A, E, R>(self: Stream<A, E, R>): Stream<A, E, R>;changesWith
Returns a stream that only emits elements that are not equal to the previously emitted element, as determined by the specified predicate.
Signature
declare const changesWith: {
<A>(f: (x: A, y: A) => boolean): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, f: (x: A, y: A) => boolean): Stream<A, E, R>;
};changesWithEffect
Emits only elements that differ from the previous element, using an effectful equality check.
Details
The predicate runs for each element after the first; returning true treats it as equal and skips it.
Signature
declare const changesWithEffect: {
<A, E2, R2>(
f: (x: A, y: A) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
f: (x: A, y: A) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Delays & Timeouts
Ends the stream if it does not produce a value within the specified duration.
Signature
declare const timeout: {
(duration: Input): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, duration: Input): Stream<A, E, R>;
};timeoutOrElse
Switches to a fallback stream if this stream does not emit a value within the specified duration.
When to use
Use when a stream should continue with another stream if an upstream pull waits longer than the allowed duration.
Details
The timeout is checked for each pull. A zero duration uses orElse immediately, while an infinite duration leaves the original stream unchanged.
Gotchas
The fallback stream is not timed after the switch.
See
timeoutfor ending the stream instead of switching to a fallback stream
Signature
declare const timeoutOrElse: {
<B, E2, R2>(options: {
readonly duration: Duration.Input;
readonly orElse: () => Stream<B, E2, R2>;
}): <A, E, R>(self: Stream<A, E, R>) => Stream<B | A, E2 | E, R2 | R>;
<A, E, R, B, E2, R2>(
self: Stream<A, E, R>,
options: {
readonly duration: Duration.Input;
readonly orElse: () => Stream<B, E2, R2>;
},
): Stream<A | B, E | E2, R | R2>;
};Destructors
mkArrayBuffer
Concatenates the stream's Uint8Array chunks into a single ArrayBuffer.
Signature
declare function mkArrayBuffer<E, R>(
self: Stream<Uint8Array<ArrayBufferLike>, E, R>,
): Effect<ArrayBuffer, E, R>;Concatenates all emitted strings into a single string.
Signature
declare function mkString<E, R>(self: Stream<string, E, R>): Effect<string, E, R>;mkUint8Array
Concatenates the stream's Uint8Array chunks into a single Uint8Array.
Signature
declare function mkUint8Array<E, R>(
self: Stream<Uint8Array<ArrayBufferLike>, E, R>,
): Effect<Uint8Array<ArrayBufferLike>, E, R>;Runs a sink to peel off enough elements to produce a value and returns that value with the remaining stream in a scope.
Details
The returned stream is only valid within the scope.
Signature
declare const peel: {
<A2, A, E2, R2>(
sink: Sink<A2, A, A, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<[A2, Stream<A, E, never>], E2 | E, Scope | R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, A, E2, R2>,
): Effect<[A2, Stream<A, E, never>], E | E2, Scope | R | R2>;
};Runs a stream with a sink and returns the sink result.
Signature
declare const run: {
<A2, A, L, E2, R2>(
sink: Sink<A2, A, L, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<A2, E2 | E, R2 | R>;
<A, E, R, L, A2, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, L, E2, R2>,
): Effect<A2, E | E2, R | R2>;
};runCollect
Runs the stream and collects all elements into an array.
Signature
declare function runCollect<A, E, R>(self: Stream<A, E, R>): Effect<Array<A>, E, R>;Runs the stream and returns the number of elements emitted.
Signature
declare function runCount<A, E, R>(self: Stream<A, E, R>): Effect<number, E, R>;Runs the stream for its effects, discarding emitted elements.
Signature
declare function runDrain<A, E, R>(self: Stream<A, E, R>): Effect<void, E, R>;Runs the stream and folds elements using a pure reducer.
Signature
declare const runFold: {
<Z, A>(
initial: LazyArg<Z>,
f: (acc: Z, a: A) => Z,
): <E, R>(self: Stream<A, E, R>) => Effect<Z, E, R>;
<A, E, R, Z>(self: Stream<A, E, R>, initial: LazyArg<Z>, f: (acc: Z, a: A) => Z): Effect<Z, E, R>;
};runFoldEffect
Runs the stream and folds elements using an effectful reducer.
When to use
Use when reducing stream elements needs Effects, services, or failures in the reducer.
Signature
declare const runFoldEffect: {
<Z, A, EX, RX>(
initial: LazyArg<Z>,
f: (acc: Z, a: A) => Effect<Z, EX, RX>,
): <E, R>(self: Stream<A, E, R>) => Effect<Z, EX | E, RX | R>;
<A, E, R, Z, EX, RX>(
self: Stream<A, E, R>,
initial: LazyArg<Z>,
f: (acc: Z, a: A) => Effect<Z, EX, RX>,
): Effect<Z, E | EX, R | RX>;
};runForEach
Runs the provided effectful callback for each element of the stream.
Signature
declare const runForEach: {
<A, X, E2, R2>(
f: (a: A) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<X, E2, R2>,
): Effect<void, E | E2, R | R2>;
};runForEachArray
Consumes the stream in chunks, passing each non-empty array to the callback.
Signature
declare const runForEachArray: {
<A, X, E2, R2>(
f: (a: readonly [A, A]) => Effect<X, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: readonly [A, A]) => Effect<X, E2, R2>,
): Effect<void, E | E2, R | R2>;
};runForEachWhile
Runs the stream, applying the effectful predicate to each element and stopping when it returns false.
Signature
declare const runForEachWhile: {
<A, E2, R2>(
f: (a: A) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Effect<void, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Effect<boolean, E2, R2>,
): Effect<void, E | E2, R | R2>;
};Runs the stream and returns the first element as an Option.
Signature
declare function runHead<A, E, R>(self: Stream<A, E, R>): Effect<Option<A>, E, R>;runIntoPubSub
Runs the stream, publishing elements into the provided PubSub.
Details
shutdownOnEnd controls whether the PubSub is shut down when the stream ends. It only shuts down when set to true.
Signature
declare const runIntoPubSub: {
<A>(
pubsub: PubSub<A>,
options?: {
readonly shutdownOnEnd?: boolean;
},
): <E, R>(self: Stream<A, E, R>) => Effect<void, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
pubsub: PubSub<A>,
options?: {
readonly shutdownOnEnd?: boolean;
},
): Effect<void, never, R>;
};runIntoQueue
Runs the stream, offering each element to the provided queue and ending it with Cause.Done when the stream completes.
Signature
declare const runIntoQueue: {
<A, E>(queue: Queue<A, Done<void> | E>): <R>(self: Stream<A, E, R>) => Effect<void, never, R>;
<A, E, R>(self: Stream<A, E, R>, queue: Queue<A, Done<void> | E>): Effect<void, never, R>;
};Runs the stream and returns the last element as an Option.
When to use
Use to consume a finite stream when only the final emitted element matters.
Details
Option.some contains the last emitted element. Option.none means the stream completed without emitting.
Gotchas
The returned effect waits for the stream to complete before it can produce a value.
See
runHeadfor consuming only the first emitted elementrunCollectfor collecting every emitted elementrunDrainfor consuming the stream while discarding emitted elements
Signature
declare function runLast<A, E, R>(self: Stream<A, E, R>): Effect<Option<A>, E, R>;Runs the stream and returns the numeric sum of its elements.
Signature
declare function runSum<E, R>(self: Stream<number, E, R>): Effect<number, E, R>;toAsyncIterable
Converts a stream to an AsyncIterable for for await...of consumption.
Signature
declare function toAsyncIterable<A, E>(self: Stream<A, E>): AsyncIterable<A>;toAsyncIterableEffect
Creates an effect that yields an AsyncIterable using the current services.
When to use
Use when the AsyncIterable should be created inside Effect with the current context supplying the stream's services.
Signature
declare function toAsyncIterableEffect<A, E, R>(
self: Stream<A, E, R>,
): Effect<AsyncIterable<A, any, any>, never, R>;toAsyncIterableWith
Converts the stream to an AsyncIterable using the provided services.
When to use
Use when converting outside an Effect and you already have the Context needed to run the stream.
Signature
declare const toAsyncIterableWith: {
<XR>(context: Context<XR>): <A, E, R>(self: Stream<A, E, R>) => AsyncIterable<A>;
<A, E, XR, R>(self: Stream<A, E, R>, context: Context<XR>): AsyncIterable<A>;
};Converts a stream to a PubSub of emitted values for concurrent consumption.
Details
shutdownOnEnd indicates whether the PubSub should be shut down when the stream ends. By default this is true.
Signature
declare const toPubSub: {
(
options:
| {
readonly capacity: "unbounded";
readonly replay?: number;
readonly shutdownOnEnd?: boolean;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly shutdownOnEnd?: boolean;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<PubSub<A>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
readonly replay?: number;
readonly shutdownOnEnd?: boolean;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly shutdownOnEnd?: boolean;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Effect<PubSub<A>, never, Scope | R>;
};toPubSubTake
Converts a stream to a PubSub of Take values for concurrent consumption.
Details
Take values include the stream's end and failure signals.
Signature
declare const toPubSubTake: {
(
options:
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<PubSub<Take<A, E, void>>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
readonly replay?: number;
}
| {
readonly capacity: number;
readonly replay?: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Effect<PubSub<Take<A, E, void>>, never, Scope | R>;
};Returns a scoped pull for manually consuming the stream's output chunks.
Details
The pull fails with Cause.Done when the stream ends and with the stream error on failure.
Signature
declare function toPull<A, E, R>(
self: Stream<A, E, R>,
): Effect<Pull<readonly [A, A], E, void, never>, never, Scope | R>;Creates a scoped dequeue that is fed by the stream for concurrent consumption.
Details
Elements are offered to the queue as the stream runs. Stream completion is signaled with Cause.Done, stream failures fail the queue, and the queue is shut down when the surrounding scope closes.
Signature
declare const toQueue: {
(
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): <A, E, R>(self: Stream<A, E, R>) => Effect<Dequeue<A, Done<void> | E>, never, Scope | R>;
<A, E, R>(
self: Stream<A, E, R>,
options:
| {
readonly capacity: "unbounded";
}
| {
readonly capacity: number;
readonly strategy?: "dropping" | "sliding" | "suspend";
},
): Effect<Dequeue<A, Done<void> | E>, never, Scope | R>;
};toReadableStream
Converts a stream to a ReadableStream.
Details
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const toReadableStream: {
<A>(options?: {
readonly strategy?: QueuingStrategy<A>;
}): <E>(self: Stream<A, E>) => ReadableStream<A>;
<A, E>(
self: Stream<A, E>,
options?: {
readonly strategy?: QueuingStrategy<A>;
},
): ReadableStream<A>;
};toReadableStreamEffect
Creates an Effect that builds a ReadableStream from the stream.
When to use
Use when bridging to Web Streams from inside an Effect so the required services can be captured from the current context.
Details
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const toReadableStreamEffect: {
<A>(options?: {
readonly strategy?: QueuingStrategy<A>;
}): <E, R>(self: Stream<A, E, R>) => Effect<ReadableStream<A>, never, R>;
<A, E, R>(
self: Stream<A, E, R>,
options?: {
readonly strategy?: QueuingStrategy<A>;
},
): Effect<ReadableStream<A>, never, R>;
};toReadableStreamWith
Converts the stream to a ReadableStream using the provided services.
When to use
Use when bridging to Web Streams and you already have the Context required to run the stream outside an Effect.
Details
See https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream.
Signature
declare const toReadableStreamWith: <A, XR>(context: Context<XR>, options?: {
readonly strategy?: QueuingStrategy<A>;
}) => <E, R>(self: Stream<A, E, R>) => ReadableStream<A> & <A, E, XR, R>(self: Stream<A, E, R>, context: Context<XR>, options?: {
readonly strategy?: QueuingStrategy<A>;
}) => ReadableStream<A>Encoding
encodeText
Encodes a stream of strings into UTF-8 Uint8Array chunks.
Signature
declare function encodeText<E, R>(
self: Stream<string, E, R>,
): Stream<Uint8Array<ArrayBufferLike>, E, R>;Error Handling
catchCause
Switches over to the stream produced by the provided function in case this one fails. Allows recovery from all causes of failure, including interruption if the stream is uninterruptible.
Signature
declare const catchCause: {
<E, A2, E2, R2>(
f: (cause: Cause<E>) => Stream<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (cause: Cause<E>) => Stream<A2, E2, R2>,
): Stream<A | A2, E2, R | R2>;
};catchCauseFilter
Recovers from stream failures by filtering the Cause and switching to a recovery stream.
When to use
Use when you need to recover a stream only from causes selected by a Filter, while giving the recovery both the selected value and the original Cause.
Details
The filter is applied to the full Cause. A successful filter result is passed to f together with the original cause; a failed filter result re-fails with the residual cause.
See
catchCauseIffor predicate-based cause selectioncatchFilterfor filtering typed error values instead of full causescatchCausefor recovering from every cause without filtering
Signature
declare const catchCauseFilter: {
<E, EB, A2, E2, R2, X extends Cause<any>>(
filter: Filter<Cause<E>, EB, X>,
f: (failure: EB, cause: Cause<E>) => Stream<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | Error<X>, R2 | R>;
<A, E, R, EB, A2, E2, R2, X extends Cause<any>>(
self: Stream<A, E, R>,
filter: Filter<Cause<E>, EB, X>,
f: (failure: EB, cause: Cause<E>) => Stream<A2, E2, R2>,
): Stream<A | A2, E2 | Error<X>, R | R2>;
};catchCauseIf
Recovers from stream failures by filtering the Cause and switching to a recovery stream. Non-matching causes are re-emitted as failures.
Signature
declare const catchCauseIf: {
<E, A2, E2, R2>(
predicate: Predicate<Cause<E>>,
f: (cause: Cause<E>) => Stream<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E | E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
predicate: Predicate<Cause<E>>,
f: (cause: Cause<E>) => Stream<A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};catchFilter
Recovers from errors that match a Filter by switching to a recovery stream.
When to use
Use to recover from stream errors with a reusable Filter when matching can also narrow or transform the error before choosing the recovery stream.
Details
Successful filter results are passed to f. Failed filter results go to orElse when provided; otherwise the filter failure is re-failed.
See
catchIffor predicate or refinement based recoverycatchTagfor_tagbased recovery from one tagged errorcatchTagsfor_tagbased recovery from multiple tagged errorscatchCauseFilterfor filtering full causes
Signature
declare const catchFilter: {
<E, EB, A2, E2, R2, X, A3 = unassigned, E3 = never, R3 = never>(
filter: Filter<NoInfer<E>, EB, X>,
f: (failure: EB) => Stream<A2, E2, R2>,
orElse?: (failure: X) => Stream<A3, E3, R3>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
A2 | A | Exclude<A3, unassigned>,
E2 | E3 | A3 extends unassigned ? X : never,
R2 | R3 | R
>;
<A, E, R, EB, A2, E2, R2, X, A3 = unassigned, E3 = never, R3 = never>(
self: Stream<A, E, R>,
filter: Filter<NoInfer<E>, EB, X>,
f: (failure: EB) => Stream<A2, E2, R2>,
orElse?: (failure: X) => Stream<A3, E3, R3>,
): Stream<
A | A2 | Exclude<A3, unassigned>,
E2 | E3 | A3 extends unassigned ? X : never,
R | R2 | R3
>;
};Recovers from errors that match a predicate by switching to a recovery stream.
Details
When a failure matches the filter, the stream switches to the recovery stream. Non-matching failures propagate downstream, so the error type is preserved unless the filter narrows it.
Signature
declare const catchIf: {
<E, EB, A2, E2, R2, A3 = unassigned, E3 = never, R3 = never>(
refinement: Refinement<NoInfer<E>, EB>,
f: (e: EB) => Stream<A2, E2, R2>,
orElse?: (e: Exclude<E, EB>) => Stream<A3, E3, R3>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
A2 | A | Exclude<A3, unassigned>,
E2 | E3 | A3 extends unassigned ? Exclude<E, EB> : never,
R2 | R3 | R
>;
<E, A2, E2, R2, A3 = unassigned, E3 = never, R3 = never>(
predicate: Predicate<NoInfer<E>>,
f: (e: NoInfer<E>) => Stream<A2, E2, R2>,
orElse?: (e: NoInfer<E>) => Stream<A3, E3, R3>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
A2 | A | Exclude<A3, unassigned>,
E2 | E3 | A3 extends unassigned ? E : never,
R2 | R3 | R
>;
<A, E, R, EB, A2, E2, R2, A3 = unassigned, E3 = never, R3 = never>(
self: Stream<A, E, R>,
refinement: Refinement<E, EB>,
f: (e: EB) => Stream<A2, E2, R2>,
orElse?: (e: Exclude<E, EB>) => Stream<A3, E3, R3>,
): Stream<
A | A2 | Exclude<A3, unassigned>,
E2 | E3 | A3 extends unassigned ? Exclude<E, EB> : never,
R | R2 | R3
>;
<A, E, R, A2, E2, R2, A3 = unassigned, E3 = never, R3 = never>(
self: Stream<A, E, R>,
predicate: Predicate<E>,
f: (e: E) => Stream<A2, E2, R2>,
orElse?: (e: E) => Stream<A3, E3, R3>,
): Stream<
A | A2 | Exclude<A3, unassigned>,
E2 | E3 | A3 extends unassigned ? E : never,
R | R2 | R3
>;
};catchReason
Catches a specific reason within a tagged error.
When to use
Use to handle nested error causes without removing the parent error from the error channel.
Details
The handler receives the unwrapped reason.
Signature
declare const catchReason: {
<K extends string, E, RK extends string, A2, E2, R2, A3 = unassigned, E3 = never, R3 = never>(
errorTag: K,
reasonTag: RK,
f: (
reason: ExtractReason<ExtractTag<NoInfer<E>, K>, RK>,
error: NarrowReason<ExtractTag<NoInfer<E>, K>, RK>,
) => Stream<A2, E2, R2>,
orElse?: (
reason: ExcludeReason<ExtractTag<NoInfer<E>, K>, RK>,
error: OmitReason<ExtractTag<NoInfer<E>, K>, RK>,
) => Stream<A3, E3, R3>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
A2 | A | Exclude<A3, unassigned>,
E2 | E3 | ExcludeTag<E, K> | A3 extends unassigned ? ExtractTag<E, K> : never,
R2 | R3 | R
>;
<
A,
E,
R,
K extends string,
RK extends string,
A2,
E2,
R2,
A3 = unassigned,
E3 = never,
R3 = never,
>(
self: Stream<A, E, R>,
errorTag: K,
reasonTag: RK,
f: (
reason: ExtractReason<ExtractTag<E, K>, RK>,
error: NarrowReason<ExtractTag<E, K>, RK>,
) => Stream<A2, E2, R2>,
orElse?: (
reason: ExcludeReason<ExtractTag<E, K>, RK>,
error: OmitReason<ExtractTag<E, K>, RK>,
) => Stream<A3, E3, R3>,
): Stream<
A | A2 | Exclude<A3, unassigned>,
E2 | E3 | ExcludeTag<E, K> | A3 extends unassigned ? ExtractTag<E, K> : never,
R | R2 | R3
>;
};catchReasons
Catches multiple reasons within a tagged error using an object of handlers.
Signature
declare const catchReasons: {
<
K extends string,
E,
Cases extends {
[RK in string]: (
reason: ExtractReason<ExtractTag<NoInfer<E>, K>, RK>,
error: NarrowReason<ExtractTag<NoInfer<E>, K>, RK>,
) => Stream<any, any, any>;
},
A2 = unassigned,
E2 = never,
R2 = never,
>(
errorTag: K,
cases: Cases,
orElse?: (
reason: ExcludeReason<ExtractTag<NoInfer<E>, K>, Extract<keyof Cases, string>>,
error: OmitReason<ExtractTag<NoInfer<E>, K>, Extract<keyof Cases, string>>,
) => Stream<A2, E2, R2>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
| A
| Exclude<A2, unassigned>
| {
[RK in string | number | symbol]: Cases[RK] extends (
...args: Array<any>
) => Stream<A, any, any>
? A
: never;
}[keyof Cases],
E2 | ExcludeTag<E, K> | A2 extends unassigned
? ExtractTag<E, K>
:
| never
| {
[RK in string | number | symbol]: Cases[RK] extends (
...args: Array<any>
) => Stream<any, E, any>
? E
: never;
}[keyof Cases],
| R2
| R
| {
[RK in string | number | symbol]: Cases[RK] extends (
...args: Array<any>
) => Stream<any, any, R>
? R
: never;
}[keyof Cases]
>;
<
A,
E,
R,
K extends string,
Cases extends {
[RK in string]: (
reason: ExtractReason<ExtractTag<E, K>, RK>,
error: NarrowReason<ExtractTag<E, K>, RK>,
) => Stream<any, any, any>;
},
A2 = unassigned,
E2 = never,
R2 = never,
>(
self: Stream<A, E, R>,
errorTag: K,
cases: Cases,
orElse?: (
reason: ExcludeReason<ExtractTag<NoInfer<E>, K>, Extract<keyof Cases, string>>,
error: OmitReason<ExtractTag<NoInfer<E>, K>, Extract<keyof Cases, string>>,
) => Stream<A2, E2, R2>,
): Stream<
| A
| Exclude<A2, unassigned>
| {
[RK in string | number | symbol]: Cases[RK] extends (
...args: Array<any>
) => Stream<A, any, any>
? A
: never;
}[keyof Cases],
E2 | ExcludeTag<E, K> | A2 extends unassigned
? ExtractTag<E, K>
:
| never
| {
[RK in string | number | symbol]: Cases[RK] extends (
...args: Array<any>
) => Stream<any, E, any>
? E
: never;
}[keyof Cases],
| R
| R2
| {
[RK in string | number | symbol]: Cases[RK] extends (
...args: Array<any>
) => Stream<any, any, R>
? R
: never;
}[keyof Cases]
>;
};Recovers from failures whose _tag matches the provided value by switching to the stream returned by f.
When to use
Use when you need to handle a specific error case from a stream whose error type is a tagged union with a readonly _tag field.
Signature
declare const catchTag: {
<
K extends string | readonly [Tags<E>, Tags<E>],
E,
A1,
E1,
R1,
A2 = unassigned,
E2 = never,
R2 = never,
>(
k: K,
f: (
e: ExtractTag<NoInfer<E>, K extends readonly [string, string] ? K[number] : K>,
) => Stream<A1, E1, R1>,
orElse?: (
e: ExcludeTag<E, K extends readonly [string, string] ? K[number] : K>,
) => Stream<A2, E2, R2>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
A1 | A | Exclude<A2, unassigned>,
E1 | E2 | A2 extends unassigned
? ExcludeTag<E, K extends readonly [string, string] ? K[number] : K>
: never,
R1 | R2 | R
>;
<
A,
E,
R,
K extends string | readonly [Tags<E>, Tags<E>],
R1,
E1,
A1,
A2 = unassigned,
E2 = never,
R2 = never,
>(
self: Stream<A, E, R>,
k: K,
f: (
e: ExtractTag<E, K extends readonly [string, string] ? K[number] : K>,
) => Stream<A1, E1, R1>,
orElse?: (
e: ExcludeTag<E, K extends readonly [string, string] ? K[number] : K>,
) => Stream<A2, E2, R2>,
): Stream<
A | A1 | Exclude<A2, unassigned>,
E1 | E2 | A2 extends unassigned
? ExcludeTag<E, K extends readonly [string, string] ? K[number] : K>
: never,
R | R1 | R2
>;
};Switches to a recovery stream based on matching _tag handlers.
Signature
declare const catchTags: {
<
E,
Cases extends
| {}
| {
[K in string]: (
error: Extract<
E,
{
_tag: K;
}
>,
) => Stream<any, any, any>;
},
A2 = unassigned,
E2 = never,
R2 = never,
>(
cases: Cases,
orElse?: (
e: Exclude<
E,
{
_tag: keyof Cases;
}
>,
) => Stream<A2, E2, R2>,
): <A, R>(
self: Stream<A, E, R>,
) => Stream<
| A
| Exclude<A2, unassigned>
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Stream<A, any, any>
? A
: never;
}[keyof Cases],
E2 | A2 extends unassigned
? Exclude<
E,
{
_tag: keyof Cases;
}
>
:
| never
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Stream<any, E, any>
? E
: never;
}[keyof Cases],
| R2
| R
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Stream<any, any, R>
? R
: never;
}[keyof Cases]
>;
<
R,
E,
A,
Cases extends
| {}
| {
[K in string]: (
error: Extract<
E,
{
_tag: K;
}
>,
) => Stream<any, any, any>;
},
A2 = unassigned,
E2 = never,
R2 = never,
>(
self: Stream<A, E, R>,
cases: Cases,
orElse?: (
e: Exclude<
E,
{
_tag: keyof Cases;
}
>,
) => Stream<A2, E2, R2>,
): Stream<
| A
| Exclude<A2, unassigned>
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Stream<A, any, any>
? A
: never;
}[keyof Cases],
E2 | A2 extends unassigned
? Exclude<
E,
{
_tag: keyof Cases;
}
>
:
| never
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Stream<any, E, any>
? E
: never;
}[keyof Cases],
| R
| R2
| {
[K in string | number | symbol]: Cases[K] extends (
...args: Array<any>
) => Stream<any, any, R>
? R
: never;
}[keyof Cases]
>;
};Ignores failures and ends the stream on error.
When to use
Use when you want a failing stream to end gracefully rather than propagate the error.
Details
The log option controls whether the failure is logged before the stream terminates.
See
ignoreCausefor a variant that also ignores defects, not just typed failures
Signature
declare const ignore: <
Arg extends
| Stream<any, any, any>
| {
readonly log?: boolean | Severity;
}
| undefined,
>(
selfOrOptions: Arg,
options?: {
readonly log?: boolean | Severity;
},
) => [Arg] extends [Stream<infer A, infer _E, infer R>]
? Stream<A, never, R>
: <A, E, R>(self: Stream<A, E, R>) => Stream<A, never, R>;ignoreCause
Ignores the stream's failure cause, including defects, and ends the stream.
When to use
Use when you need to silently suppress a stream's entire failure cause, including both typed errors and defects, rather than propagate it downstream.
See
ignoreto ignore only typed failures without suppressing defects
Signature
declare const ignoreCause: <
Arg extends
| Stream<any, any, any>
| {
readonly log?: boolean | Severity;
}
| undefined,
>(
streamOrOptions: Arg,
options?: {
readonly log?: boolean | Severity;
},
) => [Arg] extends [Stream<infer A, infer _E, infer R>]
? Stream<A, never, R>
: <A, E, R>(self: Stream<A, E, R>) => Stream<A, never, R>;Transforms the errors emitted by this stream using f.
Signature
declare const mapError: {
<E, E2>(f: (error: E) => E2): <A, R>(self: Stream<A, E, R>) => Stream<A, E2, R>;
<A, E, R, E2>(self: Stream<A, E, R>, f: (error: E) => E2): Stream<A, E2, R>;
};Runs the provided effect when the stream fails, passing the failure cause.
Gotchas
Note: Unlike Effect.onError there is no guarantee that the provided effect will not be interrupted.
Signature
declare const onError: {
<E, X, R2>(
cleanup: (cause: Cause<E>) => Effect<X, never, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, X, R2>(
self: Stream<A, E, R>,
cleanup: (cause: Cause<E>) => Effect<X, never, R2>,
): Stream<A, E, R | R2>;
};Turns typed failures into defects, making the stream infallible.
Signature
declare function orDie<A, E, R>(self: Stream<A, E, R>): Stream<A, never, R>;orElseIfEmpty
Switches to a fallback stream if this stream is empty.
Signature
declare const orElseIfEmpty: {
<E, A2, E2, R2>(
orElse: LazyArg<Stream<A2, E2, R2>>,
): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E | E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
orElse: LazyArg<Stream<A2, E2, R2>>,
): Stream<A | A2, E | E2, R | R2>;
};orElseSucceed
Returns a stream that emits a fallback value when this stream fails.
Signature
declare const orElseSucceed: {
<E, A2>(f: (error: E) => A2): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, never, R>;
<A, E, R, A2>(self: Stream<A, E, R>, f: (error: E) => A2): Stream<A | A2, never, R>;
};Lifts failures and successes into a Result, yielding a stream that cannot fail.
Details
The stream ends after the first failure, emitting a Result.fail value.
Signature
declare function result<A, E, R>(self: Stream<A, E, R>): Stream<Result<A, E>, never, R>;Retries the stream according to the given schedule when it fails.
Details
This retries the entire stream, so will re-execute all of the stream's acquire operations.
The schedule is reset as soon as the first element passes through the stream again.
Signature
declare const retry: {
<E, X, E2, R2>(policy: Schedule<X, NoInfer<E>, E2, R2> | ($: <SO, SE, SR>(_: Schedule<SO, NoInfer<E>, SE, SR>) => Schedule<SO, E, SE, SR>) => Schedule<X, NoInfer<E>, E2, R2>): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>;
<A, E, R, X, E2, R2>(self: Stream<A, E, R>, policy: Schedule<X, NoInfer<E>, E2, R2> | ($: <SO, SE, SR>(_: Schedule<SO, NoInfer<E>, SE, SR>) => Schedule<SO, E, SE, SR>) => Schedule<X, NoInfer<E>, E2, R2>): Stream<A, E | E2, R | R2>;
}Runs an effect when the stream fails without changing its values or error, unless the tap effect itself fails.
Signature
declare const tapCause: {
<E, A2, E2, R2>(
f: (cause: Cause<E>) => Effect<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (cause: Cause<E>) => Effect<A2, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Peeks at errors effectfully without changing the stream unless the tap fails.
Signature
declare const tapError: {
<E, A2, E2, R2>(
f: (error: E) => Effect<A2, E2, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E | E2, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (error: E) => Effect<A2, E2, R2>,
): Stream<A, E | E2, R | R2>;
};withExecutionPlan
Applies an ExecutionPlan to a stream, retrying with step-provided resources until it succeeds or the plan is exhausted.
Details
By default, a failing step can fallback even after emitting elements; set preventFallbackOnPartialStream to fail instead of mixing partial output with a later fallback.
Attempts can be observed from outside the stream by passing options.onEvent, which receives an ExecutionPlan.Event before each attempt and after it settles; see Effect.withExecutionPlan for the handler semantics. When a downstream consumer stops pulling early (for example Stream.take outside the plan), the truncated attempt reports AttemptSuccess: the consumer stopped, not the source.
Signature
declare const withExecutionPlan: {
<Input, R2, Provides, PolicyE, RX = never>(
policy: ExecutionPlan<{
error: PolicyE;
input: Input;
provides: Provides;
requirements: R2;
}>,
options?: {
readonly onEvent?: (
event: ExecutionPlan.Event<Input | PolicyE>,
) => Effect.Effect<void, never, RX>;
readonly preventFallbackOnPartialStream?: boolean;
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, PolicyE | E, R2 | RX | Exclude<R, Provides>>;
<A, E, R, R2, Input, Provides, PolicyE, RX = never>(
self: Stream<A, E, R>,
policy: ExecutionPlan<{
error: PolicyE;
input: Input;
provides: Provides;
requirements: R2;
}>,
options?: {
readonly onEvent?: (
event: ExecutionPlan.Event<E | PolicyE>,
) => Effect.Effect<void, never, RX>;
readonly preventFallbackOnPartialStream?: boolean;
},
): Stream<A, E | PolicyE, R2 | RX | Exclude<R, Provides>>;
};Filtering
Drops the first n elements from this stream.
Signature
declare const drop: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Drops the last specified number of elements from this stream.
Details
Keeps the last n elements in memory to drop them on completion.
Signature
declare const dropRight: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Drops elements until the specified predicate evaluates to true, then drops that matching element.
Signature
declare const dropUntil: {
<A>(
predicate: (a: NoInfer<A>, index: number) => boolean,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>, index: number) => boolean,
): Stream<A, E, R>;
};dropUntilEffect
Drops all elements of the stream until the specified effectful predicate evaluates to true.
When to use
Use when dropping the leading prefix requires an Effect or service and the first matching element should also be dropped.
Signature
declare const dropUntilEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>, index: number) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>, index: number) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Drops elements from the stream while the specified predicate evaluates to true.
Signature
declare const dropWhile: {
<A>(
predicate: (a: NoInfer<A>, index: number) => boolean,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>, index: number) => boolean,
): Stream<A, E, R>;
};dropWhileEffect
Drops elements while the specified effectful predicate evaluates to true.
Signature
declare const dropWhileEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>, index: number) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: A, index: number) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};dropWhileFilter
Drops elements while the filter succeeds.
When to use
Use when you need to remove a leading stream prefix based on a synchronous Filter result while preserving the remaining original stream elements.
Details
Result.succeed drops the current element. The first Result.fail stops dropping, emits that original element, and the rest of the source stream is emitted without further filtering.
See
dropWhilefor boolean predicate prefix droppingtakeWhileFilterfor keeping the accepted prefix as filter success valuesdropWhileEffectfor effectful predicate prefix dropping
Signature
declare const dropWhileFilter: {
<A, B, X>(filter: Filter<NoInfer<A>, B, X>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R, B, X>(self: Stream<A, E, R>, filter: Filter<NoInfer<A>, B, X>): Stream<A, E, R>;
};Filters a stream to the elements that satisfy a predicate.
Signature
declare const filter: {
<A, B>(refinement: Refinement<NoInfer<A>, B>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A>(predicate: Predicate<NoInfer<A>>): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, refinement: Refinement<A, B>): Stream<B, E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<A, E, R>;
};filterEffect
Filters elements in a single pass effectfully.
Signature
declare const filterEffect: {
<A, EX, RX>(
predicate: (a: NoInfer<A>, i: number) => Effect<boolean, EX, RX>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, EX | E, RX | R>;
<A, E, R, EX, RX>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>, i: number) => Effect<boolean, EX, RX>,
): Stream<A, E | EX, R | RX>;
};Filters and maps stream elements in one pass using a Filter.
When to use
Use to keep only stream elements accepted by a Filter and emit each filter success value.
Details
Result.succeed values are emitted and Result.fail values are skipped.
See
filterfor keeping original elements with a boolean predicate or refinementfilterMapEffectfor an effectfulFilterpartitionfor consuming both filter success and failure values
Signature
declare const filterMap: {
<A, B, X>(filter: Filter<NoInfer<A>, B, X>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B, X>(self: Stream<A, E, R>, filter: Filter<A, B, X>): Stream<B, E, R>;
};filterMapEffect
Filters and maps elements in one pass effectfully using a FilterEffect.
When to use
Use to apply effectful logic that can reject stream elements or emit transformed values before they continue downstream.
Details
Result.succeed values are emitted, Result.fail values are skipped, and effect failures fail the stream.
See
filterMapfor the synchronousFiltervariantfilterEffectfor effectfully keeping original elementsmapEffectfor effectfully transforming every element
Signature
declare const filterMapEffect: {
<A, B, X, EX, RX>(
filter: FilterEffect<NoInfer<A>, B, X, EX, RX>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, EX | E, RX | R>;
<A, E, R, B, X, EX, RX>(
self: Stream<A, E, R>,
filter: FilterEffect<A, B, X, EX, RX>,
): Stream<B, E | EX, R | RX>;
};limitBytes
Emits byte chunks until the configured limit would be exceeded, then drops the crossing chunk and switches to a fallback stream.
Signature
declare const limitBytes: {
<E, R>(
bytes: SizeInput,
onLimitReached: LazyArg<Stream<Uint8Array<ArrayBufferLike>, E, R>>,
): (self: Stream<Uint8Array<ArrayBufferLike>, E, R>) => Stream<Uint8Array<ArrayBufferLike>, E, R>;
<E, R>(
self: Stream<Uint8Array<ArrayBufferLike>, E, R>,
bytes: SizeInput,
onLimitReached: LazyArg<Stream<Uint8Array<ArrayBufferLike>, E, R>>,
): Stream<Uint8Array<ArrayBufferLike>, E, R>;
};Splits a stream into scoped excluded and satisfying substreams using a Filter.
Details
The returned streams are backed by queues in the current scope and should be consumed while that scope remains open. The faster stream may advance up to bufferSize elements ahead of the slower one.
Signature
declare const partition: {
<A, Pass, Fail>(
filter: Filter<NoInfer<A>, Pass, Fail>,
options?: {
readonly bufferSize?: number;
},
): <E, R>(
self: Stream<A, E, R>,
) => Effect<
[excluded: Stream<Fail, E, never>, satisfying: Stream<Pass, E, never>],
never,
Scope | R
>;
<A, E, R, Pass, Fail>(
self: Stream<A, E, R>,
filter: Filter<NoInfer<A>, Pass, Fail>,
options?: {
readonly bufferSize?: number;
},
): Effect<
[excluded: Stream<Fail, E, never>, satisfying: Stream<Pass, E, never>],
never,
Scope | R
>;
};partitionEffect
Splits a stream with an effectful Filter, returning scoped streams for filter successes and failures.
When to use
Use when you need to classify each stream element with an effectful Filter and consume both passing and failing mapped values as streams.
Details
The returned streams are backed by queues in the current scope and should be consumed while that scope remains open. The first stream emits success values from the filter, and the second emits failure values.
See
partitionfor the pureFiltervariant, which returns the failing stream before the passing streampartitionQueuefor the lower-level queue resultfilterMapEffectfor effectful filtering that discards failed filter results
Signature
declare const partitionEffect: {
<A, Pass, Fail, EX, RX>(
filter: FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>,
options?: {
readonly capacity?: number | "unbounded";
readonly concurrency?: number | "unbounded";
},
): <E, R>(
self: Stream<A, E, R>,
) => Effect<
[passes: Stream<Pass, EX | E, never>, fails: Stream<Fail, EX | E, never>],
never,
Scope | RX | R
>;
<A, E, R, Pass, Fail, EX, RX>(
self: Stream<A, E, R>,
filter: FilterEffect<NoInfer<A>, Pass, Fail, EX, RX>,
options?: {
readonly capacity?: number | "unbounded";
readonly concurrency?: number | "unbounded";
},
): Effect<
[passes: Stream<Pass, E | EX, never>, fails: Stream<Fail, E | EX, never>],
never,
Scope | R | RX
>;
};partitionQueue
Partitions a stream using a Filter and exposes passing and failing values as scoped queues.
Details
The queues are backed by a fiber in the current scope and should be consumed while that scope remains open. Each queue fails with the stream error or Cause.Done when the source ends.
Signature
declare const partitionQueue: {
<A, Pass, Fail>(
filter: Filter<NoInfer<A>, Pass, Fail>,
options?: {
readonly capacity?: number | "unbounded";
},
): <E, R>(
self: Stream<A, E, R>,
) => Effect<
[passes: Dequeue<Pass, Done<void> | E>, fails: Dequeue<Fail, Done<void> | E>],
never,
Scope | R
>;
<A, E, R, Pass, Fail>(
self: Stream<A, E, R>,
filter: Filter<NoInfer<A>, Pass, Fail>,
options?: {
readonly capacity?: number | "unbounded";
},
): Effect<
[passes: Dequeue<Pass, Done<void> | E>, fails: Dequeue<Fail, Done<void> | E>],
never,
Scope | R
>;
};Takes the first n elements from this stream, returning Stream.empty when n < 1.
Signature
declare const take: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Keeps the last n elements from this stream.
Signature
declare const takeRight: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<A, E, R>;
};Takes elements until the predicate matches.
Details
When excludeLast is true, the matching element is dropped.
Signature
declare const takeUntil: {
<A>(
predicate: (a: NoInfer<A>, n: number) => boolean,
options?: {
readonly excludeLast?: boolean;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
predicate: (a: A, n: number) => boolean,
options?: {
readonly excludeLast?: boolean;
},
): Stream<A, E, R>;
};takeUntilEffect
Takes stream elements until an effectful predicate returns true.
When to use
Use when the stopping condition needs an Effect or service and predicate failure should fail the stream.
Signature
declare const takeUntilEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>, n: number) => Effect<boolean, E2, R2>,
options?: {
readonly excludeLast?: boolean;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: A, n: number) => Effect<boolean, E2, R2>,
options?: {
readonly excludeLast?: boolean;
},
): Stream<A, E | E2, R | R2>;
};Takes the longest initial prefix of elements that satisfy the predicate.
Signature
declare const takeWhile: {
<A, B>(
refinement: (a: NoInfer<A>, n: number) => a is B,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A>(
predicate: (a: NoInfer<A>, n: number) => boolean,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R, B>(
self: Stream<A, E, R>,
refinement: (a: NoInfer<A>, n: number) => a is B,
): Stream<B, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>, n: number) => boolean,
): Stream<A, E, R>;
};takeWhileEffect
Takes elements from the stream while the effectful predicate is true.
When to use
Use when the leading-prefix predicate needs an Effect or service and the stream should stop before the first false result.
Signature
declare const takeWhileEffect: {
<A, E2, R2>(
predicate: (a: NoInfer<A>, n: number) => Effect<boolean, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
predicate: (a: NoInfer<A>, n: number) => Effect<boolean, E2, R2>,
): Stream<A, E | E2, R | R2>;
};takeWhileFilter
Takes the longest initial prefix accepted by a Filter and emits the filter's success values.
When to use
Use to keep the leading stream elements that a Filter accepts, emit the filter's success values, and stop at the first filter failure.
Details
The stream stops at the first Result.fail returned by the filter.
See
takeWhilefor keeping original elements with a boolean predicate or refinementfilterMapfor filtering across the whole stream instead of only the leading prefixdropWhileFilterfor dropping the accepted prefix and keeping the remaining original elements
Signature
declare const takeWhileFilter: {
<A, B, X>(f: Filter<NoInfer<A>, B, X>): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B, X>(self: Stream<A, E, R>, f: Filter<NoInfer<A>, B, X>): Stream<B, E, R>;
};Returns the specified stream if the given condition is satisfied, otherwise returns an empty stream.
Signature
declare const when: {
<EX = never, RX = never>(
test: Effect<boolean, EX, RX>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, EX | E, RX | R>;
<A, E, R, EX = never, RX = never>(
self: Stream<A, E, R>,
test: Effect<boolean, EX, RX>,
): Stream<A, E | EX, R | RX>;
};Grouping
Exposes the underlying chunks as a stream of non-empty arrays.
Signature
declare function chunks<A, E, R>(self: Stream<A, E, R>): Stream<readonly [A, A], E, R>;groupAdjacentBy
Groups consecutive elements that have equal keys into non-empty arrays.
When to use
Use when you already have a stream ordered by the grouping key and want to emit each consecutive run as a non-empty array while keeping later non-adjacent runs separate.
Details
The key is computed with f; adjacent elements whose keys are equal by Equal.equals are emitted as one [key, group]. Later non-adjacent runs with the same key are emitted separately.
See
groupByKeyfor grouping all elements with the same key across the streamgroupByfor custom grouped stream construction
Signature
declare const groupAdjacentBy: {
<A, K>(
f: (a: NoInfer<A>) => K,
): <E, R>(self: Stream<A, E, R>) => Stream<readonly [K, [A, ...Array<A>]], E, R>;
<A, E, R, K>(
self: Stream<A, E, R>,
f: (a: NoInfer<A>) => K,
): Stream<readonly [K, [A, ...Array<A>]], E, R>;
};Groups elements into keyed substreams using an effectful classifier.
Signature
declare const groupBy: {
<A, K, V, E2, R2>(
f: (a: NoInfer<A>) => Effect<readonly [K, V], E2, R2>,
options?: {
readonly bufferSize?: number;
readonly idleTimeToLive?: Duration.Input;
},
): <E, R>(self: Stream<A, E, R>) => Stream<readonly [K, Stream<V, never, never>], E2 | E, R2 | R>;
<A, E, R, K, V, E2, R2>(
self: Stream<A, E, R>,
f: (a: NoInfer<A>) => Effect<readonly [K, V], E2, R2>,
options?: {
readonly bufferSize?: number;
readonly idleTimeToLive?: Duration.Input;
},
): Stream<readonly [K, Stream<V, never, never>], E | E2, R | R2>;
};groupByKey
Groups elements by a key and emits a stream per key.
Signature
declare const groupByKey: {
<A, K>(
f: (a: NoInfer<A>) => K,
options?: {
readonly bufferSize?: number;
readonly idleTimeToLive?: Duration.Input;
},
): <E, R>(self: Stream<A, E, R>) => Stream<readonly [K, Stream<A, never, never>], E, R>;
<A, E, R, K>(
self: Stream<A, E, R>,
f: (a: NoInfer<A>) => K,
options?: {
readonly bufferSize?: number;
readonly idleTimeToLive?: Duration.Input;
},
): Stream<readonly [K, Stream<A, never, never>], E, R>;
};Partitions the stream into non-empty arrays of the specified size.
Details
The final array may be smaller if there are not enough elements to fill it.
Signature
declare const grouped: {
(n: number): <A, E, R>(self: Stream<A, E, R>) => Stream<readonly [A, A], E, R>;
<A, E, R>(self: Stream<A, E, R>, n: number): Stream<readonly [A, A], E, R>;
};groupedWithin
Partitions the stream into arrays, emitting when the chunk size is reached or the duration passes.
Signature
declare const groupedWithin: {
(chunkSize: number, duration: Input): <A, E, R>(self: Stream<A, E, R>) => Stream<Array<A>, E, R>;
<A, E, R>(self: Stream<A, E, R>, chunkSize: number, duration: Input): Stream<Array<A>, E, R>;
};Groups the stream into arrays of the specified size, preserving element order.
Details
The size is clamped to at least 1.
Signature
declare const rechunk: {
(size: number): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, size: number): Stream<A, E, R>;
};Emits a sliding window of n elements.
Signature
declare const sliding: {
(chunkSize: number): <A, E, R>(self: Stream<A, E, R>) => Stream<readonly [A, A], E, R>;
<A, E, R>(self: Stream<A, E, R>, chunkSize: number): Stream<readonly [A, A], E, R>;
};slidingSize
Emits sliding windows of chunkSize elements, advancing by stepSize.
Signature
declare const slidingSize: {
(
chunkSize: number,
stepSize: number,
): <A, E, R>(self: Stream<A, E, R>) => Stream<readonly [A, A], E, R>;
<A, E, R>(
self: Stream<A, E, R>,
chunkSize: number,
stepSize: number,
): Stream<readonly [A, A], E, R>;
};Splits the stream into non-empty groups whenever the predicate matches.
Details
Matching elements act as delimiters and are not included in the output.
Signature
declare const split: {
<A, B>(
refinement: Refinement<NoInfer<A>, B>,
): <E, R>(self: Stream<A, E, R>) => Stream<readonly [Exclude<A, B>, Exclude<A, B>], E, R>;
<A>(
predicate: Predicate<NoInfer<A>>,
): <E, R>(self: Stream<A, E, R>) => Stream<readonly [A, A], E, R>;
<A, E, R, B>(
self: Stream<A, E, R>,
refinement: Refinement<A, B>,
): Stream<readonly [Exclude<A, B>, Exclude<A, B>], E, R>;
<A, E, R>(self: Stream<A, E, R>, predicate: Predicate<A>): Stream<readonly [A, A], E, R>;
};Guards
Interruption
Stops a stream after the current pull when an effect completes.
When to use
Use to stop before the next pull after an external signal completes.
Details
The effect is forked, its success value is discarded, and its failure fails the stream.
Gotchas
This does not interrupt or truncate an in-progress pull. A pull may emit multiple elements in a single chunk, in which case the entire chunk is emitted. Use interruptWhen when the stream should be interrupted immediately.
Signature
declare const haltWhen: {
<X, E2, R2>(
effect: Effect<X, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(self: Stream<A, E, R>, effect: Effect<X, E2, R2>): Stream<A, E | E2, R | R2>;
};interruptWhen
Interrupts the evaluation of this stream when the provided effect completes. The given effect will be forked as part of this stream, and its success will be discarded. This combinator will also interrupt any in-progress element being pulled from upstream.
Details
If the effect completes with a failure before the stream completes, the returned stream will emit that failure.
Signature
declare const interruptWhen: {
<X, E2, R2>(
effect: Effect<X, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(self: Stream<A, E, R>, effect: Effect<X, E2, R2>): Stream<A, E | E2, R | R2>;
};Mapping
Maps each element into a record keyed by the provided name.
Signature
declare const bindTo: {
<N extends string>(
name: N,
): <A, E, R>(self: Stream<A, E, R>) => Stream<{ [K in string]: A }, E, R>;
<A, E, R, N extends string>(self: Stream<A, E, R>, name: N): Stream<{ [K in string]: A }, E, R>;
};Maps each element to a stream and flattens the resulting streams.
Details
With the default sequential concurrency, inner streams are concatenated in input order. When concurrency is greater than 1 or "unbounded", multiple inner streams may run at the same time and their outputs are merged as they arrive.
Signature
declare const flatMap: {
<A, A2, E2, R2>(
f: (a: A) => Stream<A2, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Stream<A2, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): Stream<A2, E | E2, R | R2>;
};Flattens a stream of streams into a single stream.
Details
With the default sequential concurrency, inner streams are concatenated in strict order. When concurrency is greater than 1 or "unbounded", multiple inner streams may run at the same time and their outputs are merged as they arrive.
Signature
declare const flatten: <
Arg extends
| Stream<Stream<any, any, any>, any, any>
| {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
}
| undefined = {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
>(
selfOrOptions?: Arg,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
) => [Arg] extends [Stream<Stream<infer _A, infer _E, infer _R>, infer _E2, infer _R2>]
? Stream<_A, _E | _E2, _R | _R2>
: <A, E, R, E2, R2>(self: Stream<Stream<A, E, R>, E2, R2>) => Stream<A, E | E2, R | R2>;flattenEffect
Flattens a stream of Effect values into a stream of their results.
When to use
Use when stream elements already are effects and their successes should become stream elements while their failures enter the stream error channel.
Signature
declare const flattenEffect: <
Arg extends
| Stream<Effect.Effect<any, any, any>, any, any>
| {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
}
| undefined = {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
>(
selfOrOptions?: Arg,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
) => [Arg] extends [Stream<Effect.Effect<infer _A, infer _EX, infer _RX>, infer _E, infer _R>]
? Stream<_A, _EX | _E, _RX | _R>
: <A, EX, RX, E, R>(self: Stream<Effect.Effect<A, EX, RX>, E, R>) => Stream<A, EX | E, RX | R>;flattenIterable
Flattens the iterables emitted by this stream into the stream's structure.
Signature
declare function flattenIterable<A, E, R>(
self: Stream<Iterable<A, any, any>, E, R>,
): Stream<A, E, R>;Transforms the elements of this stream using the supplied function.
Signature
declare const map: {
<A, B>(f: (a: A, i: number) => B): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, f: (a: A, i: number) => B): Stream<B, E, R>;
};Maps elements statefully, emitting zero or more outputs per input.
Signature
declare const mapAccum: {
<S, A, B>(initial: LazyArg<S>, f: (s: S, a: A) => readonly [S, readonly Array<B>], options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, S, B>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => readonly [S, readonly Array<B>], options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): Stream<B, E, R>;
}mapAccumArray
Maps over non-empty chunk arrays statefully, emitting zero or more values per chunk.
Details
The mapping function runs once per chunk and the state is threaded across chunks.
Signature
declare const mapAccumArray: {
<S, A, B>(initial: LazyArg<S>, f: (s: S, a: readonly [A, A]) => readonly [S, readonly Array<B>], options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, S, B>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: readonly [A, A]) => readonly [S, readonly Array<B>], options?: {
readonly onHalt?: (state: S) => Array<B>;
}): Stream<B, E, R>;
}mapAccumArrayEffect
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.
Signature
declare const mapAccumArrayEffect: {
<S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: readonly [A, A]) => Effect<readonly [S, readonly Array<B>], E2, R2>, options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R>;
<A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: readonly [A, A]) => Effect<readonly [S, readonly Array<B>], E2, R2>, options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): Stream<B, E | E2, R | R2>;
}mapAccumEffect
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.
Signature
declare const mapAccumEffect: {
<S, A, B, E2, R2>(initial: LazyArg<S>, f: (s: S, a: A) => Effect<readonly [S, readonly Array<B>], E2, R2>, options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R>;
<A, E, R, S, B, E2, R2>(self: Stream<A, E, R>, initial: LazyArg<S>, f: (s: S, a: A) => Effect<readonly [S, readonly Array<B>], E2, R2>, options?: {
readonly onHalt?: (state: S) => ReadonlyArray<B>;
}): Stream<B, E | E2, R | R2>;
}Transforms each emitted chunk using the provided function, which receives the chunk and its index.
Signature
declare const mapArray: {
<A, B>(
f: (a: readonly [A, A], i: number) => readonly [B, B],
): <E, R>(self: Stream<A, E, R>) => Stream<B, E, R>;
<A, E, R, B>(
self: Stream<A, E, R>,
f: (a: readonly [A, A], i: number) => readonly [B, B],
): Stream<B, E, R>;
};mapArrayEffect
Maps over non-empty array chunks emitted by the stream effectfully.
When to use
Use when transformation needs to see and replace each non-empty emitted chunk effectfully instead of mapping individual stream elements.
Signature
declare const mapArrayEffect: {
<A, B, E2, R2>(
f: (a: readonly [A, A], i: number) => Effect<readonly [B, B], E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<B, E2 | E, R2 | R>;
<A, E, R, B, E2, R2>(
self: Stream<A, E, R>,
f: (a: readonly [A, A], i: number) => Effect<readonly [B, B], E2, R2>,
): Stream<B, E | E2, R | R2>;
};Maps both the failure and success channels of a stream.
Signature
declare const mapBoth: {
<E, E2, A, A2>(options: {
readonly onFailure: (e: E) => E2;
readonly onSuccess: (a: A) => A2;
}): <R>(self: Stream<A, E, R>) => Stream<A2, E2, R>;
<A, E, R, E2, A2>(
self: Stream<A, E, R>,
options: {
readonly onFailure: (e: E) => E2;
readonly onSuccess: (a: A) => A2;
},
): Stream<A2, E2, R>;
};Maps over elements of the stream with the specified effectful function.
When to use
Use when each stream element transformation needs an Effect, service dependency, failure channel, or configured concurrency.
Signature
declare const mapEffect: {
<A, A2, E2, R2>(
f: (a: A, i: number) => Effect<A2, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A, i: number) => Effect<A2, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): Stream<A2, E | E2, R | R2>;
};Merging
Combines elements from this stream and the specified stream by repeatedly applying a stateful function that can pull from either side.
Details
Where possible, prefer Stream.combineArray for a more efficient implementation.
Signature
declare const combine: {
<A2, E2, R2, S, E, A, A3, E3, R3>(
that: Stream<A2, E2, R2>,
s: LazyArg<S>,
f: (
s: S,
pullLeft: Pull<A, E, void>,
pullRight: Pull<A2, E2, void>,
) => Effect<readonly [A3, S], E3, R3>,
): <R>(self: Stream<A, E, R>) => Stream<A3, E3, R2 | R3 | R>;
<A, E, R, A2, E2, R2, S, A3, E3, R3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
s: LazyArg<S>,
f: (
s: S,
pullLeft: Pull<A, E, void>,
pullRight: Pull<A2, E2, void>,
) => Effect<readonly [A3, S], E3, R3>,
): Stream<A3, E3, R | R2 | R3>;
};interleave
Interleaves this stream with the specified stream by alternating pulls from each stream; when one ends, the remaining values from the other stream are emitted.
Signature
declare const interleave: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};interleaveWith
Interleaves two streams deterministically by following a boolean decider stream.
Details
The decider controls how many elements are pulled; if one side ends, pulls for that side are ignored.
Signature
declare const interleaveWith: {
<A2, E2, R2, E3, R3>(
that: Stream<A2, E2, R2>,
decider: Stream<boolean, E3, R3>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E3 | E, R2 | R3 | R>;
<A, E, R, A2, E2, R2, E3, R3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
decider: Stream<boolean, E3, R3>,
): Stream<A | A2, E | E2 | E3, R | R2 | R3>;
};Merges two streams, emitting elements from both as they arrive.
Details
By default, the merged stream ends when both streams end. Use haltStrategy to change the termination behavior.
Signature
declare const merge: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
options?: {
readonly haltStrategy?: HaltStrategy;
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
options?: {
readonly haltStrategy?: HaltStrategy;
},
): Stream<A | A2, E | E2, R | R2>;
};Merges a collection of streams, running up to the specified number concurrently.
When to use
Use to merge an iterable of already-created streams while bounding how many inner streams may run at the same time.
Details
The concurrency option is required and may be a number or "unbounded". bufferSize controls buffering between inner streams, and outputs are emitted as they arrive under concurrent merging.
See
Signature
declare const mergeAll: {
(options: {
readonly bufferSize?: number;
readonly concurrency: number | "unbounded";
}): <A, E, R>(streams: Iterable<Stream<A, E, R>>) => Stream<A, E, R>;
<A, E, R>(
streams: Iterable<Stream<A, E, R>>,
options: {
readonly bufferSize?: number;
readonly concurrency: number | "unbounded";
},
): Stream<A, E, R>;
};mergeEffect
Merges this stream with a background effect, keeping the stream's elements.
When to use
Use when an effect should run concurrently for the lifetime of a stream while only the stream's elements remain in the output.
Details
The effect runs concurrently, fails the stream if it fails, and is interrupted when the stream completes.
Signature
declare const mergeEffect: {
<A2, E2, R2>(
effect: Effect<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
effect: Effect<A2, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Merges two streams while emitting only the values from the left stream.
When to use
Use when the right stream is needed for its effects or failures, but downstream consumers should only receive values from the left stream.
Details
The right stream still runs for its effects, and any failures from the right stream are propagated. The merged stream completes when the left stream completes, interrupting the right stream.
Signature
declare const mergeLeft: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL, EL | ER, RL | RR>;
};mergeResult
Merges this stream and the specified stream together, tagging values from the left stream as Result.succeed and values from the right stream as Result.fail.
When to use
Use when values from both streams should be emitted and downstream code needs left values wrapped as successful Result values and right values wrapped as failed Result values.
Signature
declare const mergeResult: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<Result<A, A2>, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<Result<A, A2>, E | E2, R | R2>;
};mergeRight
Merges this stream and the specified stream together, emitting only the values from the right stream while the left stream runs for its effects.
When to use
Use when the left stream is needed for its effects or failures, but downstream consumers should only receive values from the right stream.
Details
The merged stream ends when the right stream completes, interrupting the left stream. Failures from the left stream still fail the merged stream.
Signature
declare const mergeRight: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AR, EL | ER, RL | RR>;
};Models
EventListener interface
Interface representing an event listener target.
Signature
interface EventListener<A> {
addEventListener(
event: string,
f: (event: A) => void,
options?:
| boolean
| {
readonly capture?: boolean;
readonly once?: boolean;
readonly passive?: boolean;
readonly signal?: AbortSignal;
},
): void;
removeEventListener(
event: string,
f: (event: A) => void,
options?:
| boolean
| {
readonly capture?: boolean;
},
): void;
}HaltStrategy type
Describes how merged streams decide when to halt.
Signature
type HaltStrategy = Channel.HaltStrategy;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.
Signature
interface Stream<out A, out E = never, out R = never> extends Variance<A, E, R>, Pipeable {
[ignoreSymbol]?: StreamUnifyIgnore;
[typeSymbol]?: unknown;
[unifySymbol]?: StreamUnify<Stream<A, E, R>>;
readonly channel: Channel<readonly [A, A], E, void, unknown, unknown, unknown, R>;
}StreamUnify interface
Type-level unification hook for Stream within the Effect type system.
Signature
interface StreamUnify<
A extends {
[typeSymbol]?: any;
},
> extends EffectUnify<A> {
Stream?: () => A[typeof typeSymbol] extends Stream<A0, E0, R0> | _ ? Stream<A0, E0, R0> : never;
}StreamUnifyIgnore interface
Type-level marker that excludes Stream from unification.
Signature
interface StreamUnifyIgnore {
Effect?: true;
}Type-level variance marker for Stream.
Details
The emitted value A, error E, and service requirement R type parameters are covariant.
Signature
interface Variance<out A, out E, out R> {
readonly "~effect/Stream": VarianceStruct<A, E, R>;
}VarianceStruct interface
Structural encoding used by Variance to record each Stream type parameter's variance.
Details
_A, _E, and _R are covariant markers.
Signature
interface VarianceStruct<out A, out E, out R> {
readonly _A: Covariant<A>;
readonly _E: Covariant<E>;
readonly _R: Covariant<R>;
}Other
Signature
declare const catch: {
<E, A2, E2, R2>(f: (error: E) => Stream<A2, E2, R2>): <A, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2, R2 | R>;
<A, E, R, A2, E2, R2>(self: Stream<A, E, R>, f: (error: E) => Stream<A2, E2, R2>): Stream<A | A2, E2, R | R2>;
}Signature
declare const let: {
<N extends string, A extends object, B>(
name: Exclude<N, keyof A>,
f: (a: NoInfer<A>) => B,
): <E, R>(
self: Stream<A, E, R>,
) => Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E, R>;
<A extends object, E, R, N extends string, B>(
self: Stream<A, E, R>,
name: Exclude<N, keyof A>,
f: (a: NoInfer<A>) => B,
): Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E, R>;
};Providing Services
Provides a layer or context to the stream, removing the corresponding service requirements. Use options.local to build the layer every time; by default, layers are shared between provide calls.
Signature
declare const provide: {
<AL, EL = never, RL = never>(
layer: Layer<AL, EL, RL> | Context<AL>,
options?: {
readonly local?: boolean;
},
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, EL | E, RL | Exclude<R, AL>>;
<A, E, R, AL, EL = never, RL = never>(
self: Stream<A, E, R>,
layer: Layer<AL, EL, RL> | Context<AL>,
options?: {
readonly local?: boolean;
},
): Stream<A, E | EL, RL | Exclude<R, AL>>;
};provideContext
Provides multiple services to the stream using a context.
Signature
declare const provideContext: {
<R2>(context: Context<R2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, Exclude<R, R2>>;
<A, E, R, R2>(self: Stream<A, E, R>, context: Context<R2>): Stream<A, E, Exclude<R, R2>>;
};provideService
Provides the stream with a single required service, eliminating that requirement from its environment.
Signature
declare const provideService: {
<I, S>(
key: Key<I, S>,
service: NoInfer<S>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, Exclude<R, I>>;
<A, E, R, I, S>(
self: Stream<A, E, R>,
key: Key<I, S>,
service: NoInfer<S>,
): Stream<A, E, Exclude<R, I>>;
};provideServiceEffect
Provides a service to the stream using an effect, removing the requirement and adding the effect's error and environment.
Signature
declare const provideServiceEffect: {
<I, S, ES, RS>(
key: Key<I, S>,
service: Effect<NoInfer<S>, ES, RS>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, ES | E, RS | Exclude<R, I>>;
<A, E, R, I, S, ES, RS>(
self: Stream<A, E, R>,
key: Key<I, S>,
service: Effect<NoInfer<S>, ES, RS>,
): Stream<A, E | ES, RS | Exclude<R, I>>;
};updateContext
Transforms the stream's required services by mapping the current context to a new one.
Signature
declare const updateContext: {
<R, R2>(
f: (context: Context<R2>) => Context<R>,
): <A, E>(self: Stream<A, E, R>) => Stream<A, E, R2>;
<A, E, R, R2>(self: Stream<A, E, R>, f: (context: Context<R2>) => Context<R>): Stream<A, E, R2>;
};updateService
Updates a single service in the stream environment by applying a function.
Signature
declare const updateService: {
<I, S>(
key: Key<I, S>,
f: (service: NoInfer<S>) => S,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, I | R>;
<A, E, R, I, S>(
self: Stream<A, E, R>,
key: Key<I, S>,
f: (service: NoInfer<S>) => S,
): Stream<A, E, R | I>;
};Racing
Runs both streams concurrently until one stream emits its first value, then mirrors that winning stream and interrupts the other.
Details
A failure or completion from one side before the other side emits does not win the race unless both sides fail or complete before emitting. After a winner is chosen, that stream's later failures are propagated.
Signature
declare const race: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR | AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL | AR, EL | ER, RL | RR>;
};Runs all streams concurrently until one stream emits its first value, then mirrors that winning stream and interrupts the rest.
Details
Failures or completion from losing streams before a winner is chosen are ignored unless every stream fails or completes before emitting. After a winner is chosen, that stream's later failures are propagated.
Signature
declare function raceAll<S extends readonly Array<Stream<any, any, any>>>(...streams: S): Stream<Success<S[number]>, Error<S[number]>, Services<S[number]>>Rate Limiting
Drops earlier elements within the debounce window and emits only the latest element after the pause.
Signature
declare const debounce: {
(duration: Input): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(self: Stream<A, E, R>, duration: Input): Stream<A, E, R>;
};Schedules the stream's elements according to the provided schedule.
Signature
declare const schedule: {
<X, E2, R2, A>(
schedule: Schedule<X, NoInfer<A>, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
schedule: Schedule<X, NoInfer<A>, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Rate-limits stream chunks with a synchronous cost function.
When to use
Use to throttle chunks when each chunk's cost can be computed synchronously.
Details
Uses a token bucket. The bucket can accumulate up to units + burst tokens, and each chunk consumes the cost returned by cost.
If using the "enforce" strategy, arrays that do not meet the bandwidth constraints are dropped. If using the "shape" strategy, arrays are delayed until they can be emitted without exceeding the bandwidth constraints.
Defaults to the "shape" strategy.
Signature
declare const throttle: {
<A>(options: {
readonly burst?: number;
readonly cost: (arr: Arr.NonEmptyReadonlyArray<A>) => number;
readonly duration: Duration.Input;
readonly strategy?: "enforce" | "shape";
readonly units: number;
}): <E, R>(self: Stream<A, E, R>) => Stream<A, E, R>;
<A, E, R>(
self: Stream<A, E, R>,
options: {
readonly burst?: number;
readonly cost: (arr: Arr.NonEmptyReadonlyArray<A>) => number;
readonly duration: Duration.Input;
readonly strategy?: "enforce" | "shape";
readonly units: number;
},
): Stream<A, E, R>;
};throttleEffect
Rate-limits stream chunks with an effectful cost function.
When to use
Use to throttle chunks when computing each chunk's cost requires an effect.
Details
Uses a token bucket. The bucket can accumulate up to units + burst tokens, and each chunk consumes the cost returned by the effectful cost function.
If using the "enforce" strategy, arrays that do not meet the bandwidth constraints are dropped. If using the "shape" strategy, arrays are delayed until they can be emitted without exceeding the bandwidth constraints.
Defaults to the "shape" strategy.
Signature
declare const throttleEffect: {
<A, E2, R2>(options: {
readonly burst?: number;
readonly cost: (arr: Arr.NonEmptyReadonlyArray<A>) => Effect.Effect<number, E2, R2>;
readonly duration: Duration.Input;
readonly strategy?: "enforce" | "shape";
readonly units: number;
}): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
options: {
readonly burst?: number;
readonly cost: (arr: Arr.NonEmptyReadonlyArray<A>) => Effect.Effect<number, E2, R2>;
readonly duration: Duration.Input;
readonly strategy?: "enforce" | "shape";
readonly units: number;
},
): Stream<A, E | E2, R | R2>;
};Resource Management
Executes the provided finalizer after this stream's finalizers run.
Signature
declare const ensuring: {
<R2>(
finalizer: Effect<unknown, never, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, R2>(self: Stream<A, E, R>, finalizer: Effect<unknown, never, R2>): Stream<A, E, R | R2>;
};Runs the provided finalizer when the stream exits, passing the exit value.
Signature
declare const onExit: {
<E, R2>(
finalizer: (exit: Exit<unknown, E>) => Effect<unknown, never, R2>,
): <A, R>(self: Stream<A, E, R>) => Stream<A, E, R2 | R>;
<A, E, R, R2>(
self: Stream<A, E, R>,
finalizer: (exit: Exit<unknown, E>) => Effect<unknown, never, R2>,
): Stream<A, E, R | R2>;
};Sequencing
Binds the result of a stream to a field in the do-notation record.
Signature
declare const bind: {
<N extends string, A, B, E2, R2>(
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Stream<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): <E, R>(
self: Stream<A, E, R>,
) => Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E2 | E, R2 | R>;
<A, E, R, N extends string, B, E2, R2>(
self: Stream<A, E, R>,
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Stream<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E | E2, R | R2>;
};bindEffect
Binds an Effect-produced value into the do-notation record for each stream element.
Signature
declare const bindEffect: {
<N extends string, A, B, E2, R2>(
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Effect<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): <E, R>(
self: Stream<A, E, R>,
) => Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E2 | E, R2 | R>;
<A, E, R, N extends string, B, E2, R2>(
self: Stream<A, E, R>,
tag: Exclude<N, keyof A>,
f: (_: NoInfer<A>) => Effect<B, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
readonly unordered?: boolean;
},
): Stream<{ [K in string | number | symbol]: K extends keyof A ? A[K] : B }, E | E2, R | R2>;
};combineArray
Combines two streams chunk-by-chunk with a stateful pull function.
When to use
Use to coordinate pulling chunks from two streams when each emitted chunk depends on both sides and local state.
Details
The combining function receives the current state and pull functions for the left and right streams. It returns the next non-empty chunk together with the next state.
Signature
declare const combineArray: {
<A2, E2, R2, S, E, A, A3, E3, R3>(
that: Stream<A2, E2, R2>,
s: LazyArg<S>,
f: (
s: S,
pullLeft: Pull<readonly [A, A], E, void>,
pullRight: Pull<readonly [A2, A2], E2, void>,
) => Effect<readonly [readonly [A3, A3], S], E3, R3>,
): <R>(self: Stream<A, E, R>) => Stream<A3, Exclude<E3, Done<any>>, R2 | R3 | R>;
<R, A2, E2, R2, S, E, A, A3, E3, R3>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
s: LazyArg<S>,
f: (
s: S,
pullLeft: Pull<readonly [A, A], E, void>,
pullRight: Pull<readonly [A2, A2], E2, void>,
) => Effect<readonly [readonly [A3, A3], S], E3, R3>,
): Stream<A3, Exclude<E3, Done<any>>, R | R2 | R3>;
};Concatenates two streams, emitting all elements from the first stream followed by all elements from the second stream.
Signature
declare const concat: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<A | A2, E | E2, R | R2>;
};Converts this stream to one that runs its effects but emits no elements.
Signature
declare function drain<A, E, R>(self: Stream<A, E, R>): Stream<never, E, R>;Runs the provided stream in the background while this stream runs, interrupting it when this stream completes and failing if the background stream fails or defects.
Signature
declare const drainFork: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(self: Stream<A, E, R>, that: Stream<A2, E2, R2>): Stream<A, E | E2, R | R2>;
};flattenArray
Flattens a stream of non-empty arrays into a stream of elements.
Signature
declare function flattenArray<A, E, R>(self: Stream<readonly [A, A], E, R>): Stream<A, E, R>;flattenTake
Unwraps Take values, emitting elements from non-empty arrays and ending or failing when the Exit signals completion.
Signature
declare function flattenTake<A, E, E2, R>(
self: Stream<Take<A, E, void>, E2, R>,
): Stream<A, E | E2, R>;Repeats this stream forever.
Signature
declare function forever<A, E, R>(self: Stream<A, E, R>): Stream<A, E, R>;intersperse
Inserts the provided element between emitted elements.
Signature
declare const intersperse: {
<A2>(element: A2): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A, E, R>;
<A, E, R, A2>(self: Stream<A, E, R>, element: A2): Stream<A | A2, E, R>;
};intersperseAffixes
Adds a start value, middle value, and end value around stream elements.
Details
The start and end values are always emitted, even when the stream is empty.
Signature
declare const intersperseAffixes: {
<A2, A3, A4>(options: {
readonly end: A4;
readonly middle: A3;
readonly start: A2;
}): <A, E, R>(self: Stream<A, E, R>) => Stream<A2 | A3 | A4 | A, E, R>;
<A, E, R, A2, A3, A4>(
self: Stream<A, E, R>,
options: {
readonly end: A4;
readonly middle: A3;
readonly start: A2;
},
): Stream<A | A2 | A3 | A4, E, R>;
};Runs the provided effect when the stream ends successfully.
Signature
declare const onEnd: {
<X, EX, RX>(
onEnd: Effect<X, EX, RX>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, EX | E, RX | R>;
<A, E, R, X, EX, RX>(self: Stream<A, E, R>, onEnd: Effect<X, EX, RX>): Stream<A, E | EX, R | RX>;
};Runs the provided effect with the first element emitted by the stream.
Signature
declare const onFirst: {
<A, X, EX, RX>(
onFirst: (element: NoInfer<A>) => Effect<X, EX, RX>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, EX | E, RX | R>;
<A, E, R, X, EX, RX>(
self: Stream<A, E, R>,
onFirst: (element: NoInfer<A>) => Effect<X, EX, RX>,
): Stream<A, E | EX, R | RX>;
};Runs the provided effect before this stream starts.
Signature
declare const onStart: {
<X, EX, RX>(
onStart: Effect<X, EX, RX>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, EX | E, RX | R>;
<A, E, R, X, EX, RX>(
self: Stream<A, E, R>,
onStart: Effect<X, EX, RX>,
): Stream<A, E | EX, R | RX>;
};pipeThrough
Pipes the stream through Sink.toChannel, emitting only the sink leftovers.
Details
If the sink completes mid-chunk, the remaining elements become the output stream.
Signature
declare const pipeThrough: {
<A2, A, L, E2, R2>(
sink: Sink<A2, A, L, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<L, E2 | E, R2 | R>;
<A, E, R, A2, L, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<A2, A, L, E2, R2>,
): Stream<L, E | E2, R | R2>;
};pipeThroughChannel
Pipes this stream through a channel that consumes and emits chunked elements.
Details
The channel receives NonEmptyReadonlyArray chunks and can transform both the output elements and error type.
Signature
declare const pipeThroughChannel: {
<R2, E, E2, A, A2>(
channel: Channel<readonly [A2, A2], E2, unknown, readonly [A, A], E, unknown, R2>,
): <R>(self: Stream<A, E, R>) => Stream<A2, E2, R2 | R>;
<R, R2, E, E2, A, A2>(
self: Stream<A, E, R>,
channel: Channel<readonly [A2, A2], E2, unknown, readonly [A, A], E, unknown, R2>,
): Stream<A2, E2, R | R2>;
};pipeThroughChannelOrFail
Pipes values through the provided channel while preserving this stream's failures alongside any channel failures.
Details
Upstream failures are not passed to the channel, so the resulting stream can fail with either the original stream error or the channel error.
Signature
declare const pipeThroughChannelOrFail: {
<R2, E, E2, A, A2>(
channel: Channel<readonly [A2, A2], E2, unknown, readonly [A, A], E, unknown, R2>,
): <R>(self: Stream<A, E, R>) => Stream<A2, E | E2, R2 | R>;
<R, R2, E, E2, A, A2>(
self: Stream<A, E, R>,
channel: Channel<readonly [A2, A2], E2, unknown, readonly [A, A], E, unknown, R2>,
): Stream<A2, E | E2, R | R2>;
};Prepends the values from the provided iterable before the stream's elements.
Signature
declare const prepend: {
<B>(values: Iterable<B>): <A, E, R>(self: Stream<A, E, R>) => Stream<B | A, E, R>;
<A, E, R, B>(self: Stream<A, E, R>, values: Iterable<B>): Stream<A | B, E, R>;
};Repeats the entire stream according to the provided schedule.
Signature
declare const repeat: {
<B, E2, R2>(schedule: Schedule<B, void, E2, R2> | ($: <SO, SE, SR>(_: Schedule<SO, void, SE, SR>) => Schedule<SO, void, SE, SR>) => Schedule<B, void, E2, R2>): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, B, E2, R2>(self: Stream<A, E, R>, schedule: Schedule<B, void, E2, R2> | ($: <SO, SE, SR>(_: Schedule<SO, void, SE, SR>) => Schedule<SO, void, SE, SR>) => Schedule<B, void, E2, R2>): Stream<A, E | E2, R | R2>;
}repeatElements
Repeats each element of the stream according to the provided schedule, including the original emission.
Signature
declare const repeatElements: {
<B, E2, R2>(
schedule: Schedule<B, unknown, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, B, E2, R2>(
self: Stream<A, E, R>,
schedule: Schedule<B, unknown, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Switches to the latest stream produced by the mapping function, interrupting the previous stream when a new element arrives.
Signature
declare const switchMap: {
<A, A2, E2, R2>(
f: (a: A) => Stream<A2, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): <E, R>(self: Stream<A, E, R>) => Stream<A2, E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
f: (a: A) => Stream<A2, E2, R2>,
options?: {
readonly bufferSize?: number;
readonly concurrency?: number | "unbounded";
},
): Stream<A2, E | E2, R | R2>;
};Runs the provided effect for each element while preserving the elements.
Signature
declare const tap: {
<A, X, E2, R2>(
f: (a: NoInfer<A>) => Effect<X, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
},
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, X, E2, R2>(
self: Stream<A, E, R>,
f: (a: NoInfer<A>) => Effect<X, E2, R2>,
options?: {
readonly concurrency?: number | "unbounded";
},
): Stream<A, E | E2, R | R2>;
};Returns a stream that effectfully "peeks" at elements and failures.
Signature
declare const tapBoth: {
<A, E, X, E2, R2, Y, E3, R3>(options: {
readonly concurrency?: number | "unbounded";
readonly onElement: (a: NoInfer<A>) => Effect.Effect<X, E2, R2>;
readonly onError: (a: NoInfer<E>) => Effect.Effect<Y, E3, R3>;
}): <R>(self: Stream<A, E, R>) => Stream<A, E | E2 | E3, R2 | R3 | R>;
<A, E, R, X, E2, R2, Y, E3, R3>(
self: Stream<A, E, R>,
options: {
readonly concurrency?: number | "unbounded";
readonly onElement: (a: NoInfer<A>) => Effect.Effect<X, E2, R2>;
readonly onError: (a: NoInfer<E>) => Effect.Effect<Y, E3, R3>;
},
): Stream<A, E | E2 | E3, R | R2 | R3>;
};Runs a sink for all stream elements while still emitting them downstream.
Signature
declare const tapSink: {
<A, E2, R2>(
sink: Sink<unknown, A, unknown, E2, R2>,
): <E, R>(self: Stream<A, E, R>) => Stream<A, E2 | E, R2 | R>;
<A, E, R, E2, R2>(
self: Stream<A, E, R>,
sink: Sink<unknown, A, unknown, E2, R2>,
): Stream<A, E | E2, R | R2>;
};Splitting
splitLines
Splits a stream of strings into lines, handling \n, \r, and \r\n delimiters across chunks.
Signature
declare function splitLines<E, R>(self: Stream<string, E, R>): Stream<string, E, R>;Tracing
Wraps the stream with a new span for tracing.
Signature
declare const withSpan: {
(
name: string,
options?: SpanOptions,
): <A, E, R>(self: Stream<A, E, R>) => Stream<A, E, Exclude<R, ParentSpan>>;
<A, E, R>(
self: Stream<A, E, R>,
name: string,
options?: SpanOptions,
): Stream<A, E, Exclude<R, ParentSpan>>;
};Type IDs
Runtime identifier stored on Stream values and used by isStream to recognize them.
Details
This marker is part of the runtime representation of Stream values. Prefer isStream when narrowing unknown values.
See
isStreamfor the public guard that checks this identifier
Signature
declare const TypeId: "~effect/Stream";String literal type used as the unique brand for Stream values.
Signature
type TypeId = "~effect/Stream";Utility Types
Extract the error type from a Stream type.
Signature
type Error<T extends Stream<any, any, any>> = [T] extends [Stream<infer _A, infer _E, infer _R>]
? _E
: never;Extract the services type from a Stream type.
Signature
type Services<T extends Stream<any, any, any>> = [T] extends [Stream<infer _A, infer _E, infer _R>]
? _R
: never;StreamTypeLambda interface
Type lambda for Stream used in higher-kinded type operations.
Signature
interface StreamTypeLambda extends TypeLambda {
readonly type: Stream<unknown, unknown, unknown>;
}Extract the success type from a Stream type.
Signature
type Success<T extends Stream<any, any, any>> = [T] extends [Stream<infer _A, infer _E, infer _R>]
? _A
: never;Zipping
Creates the cartesian product of two streams, running the right stream for each element in the left stream.
Details
See also Stream.zip for the more common point-wise variant.
Signature
declare const cross: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<[AL, AR], ER | EL, RR | RL>;
<AL, ER, RR, AR, EL, RL>(
left: Stream<AL, ER, RR>,
right: Stream<AR, EL, RL>,
): Stream<[AL, AR], ER | EL, RR | RL>;
};Creates a cartesian product of elements from two streams using a function.
Details
The right stream is rerun for every element in the left stream.
See also Stream.zipWith for the more common point-wise variant.
Signature
declare const crossWith: {
<AR, ER, RR, AL, A>(
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): Stream<A, EL | ER, RL | RR>;
};Zips this stream with another point-wise and emits tuples of elements from both streams. The new stream ends when either stream ends.
Signature
declare const zip: {
<A2, E2, R2>(
that: Stream<A2, E2, R2>,
): <A, E, R>(self: Stream<A, E, R>) => Stream<[A, A2], E2 | E, R2 | R>;
<A, E, R, A2, E2, R2>(
self: Stream<A, E, R>,
that: Stream<A2, E2, R2>,
): Stream<[A, A2], E | E2, R | R2>;
};zipFlatten
Zips this stream with another point-wise and emits tuples of elements from both streams, flattening the left tuple.
Details
The new stream will end when one of the sides ends.
Signature
declare const zipFlatten: {
<A2, E2, R2>(that: Stream<A2, E2, R2>): <A extends readonly Array<any>, E, R>(self: Stream<A, E, R>) => Stream<[...Array<A>, A2], E2 | E, R2 | R>;
<A extends readonly Array<any>, E, R, A2, E2, R2>(self: Stream<A, E, R>, that: Stream<A2, E2, R2>): Stream<[...Array<A>, A2], E | E2, R | R2>;
}Combines two streams by emitting each new element with the latest value from the other stream.
When to use
Use when two streams should start emitting combined pairs after both have produced at least one value.
Gotchas
Note: tracking the latest value is done on a per-array basis. That means that emitted elements that are not the last value in arrays will never be used for zipping.
Signature
declare const zipLatest: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<[AL, AR], ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<[AL, AR], EL | ER, RL | RR>;
};zipLatestAll
Zips multiple streams so that when a value is emitted by any stream, it is combined with the latest values from the other streams to produce a result.
When to use
Use when each stream should contribute its latest value after all streams have emitted at least once.
Gotchas
Note: tracking the latest value is done on a per-array basis. That means that emitted elements that are not the last value in arrays will never be used for zipping.
Signature
declare function zipLatestAll<T extends readonly Array<Stream<any, any, any>>>(...streams: T): Stream<[T[number]] extends [never] ? never : { [K in string | number | symbol]: T[K] extends Stream<A, _E, _R> ? A : never }, [T[number]] extends [never] ? never : T[number] extends Stream<_A, _E, _R> ? _E : never, [T[number]] extends [never] ? never : T[number] extends Stream<_A, _E, _R> ? _R : never>zipLatestWith
Combines the latest values from both streams whenever either emits, using the provided function.
When to use
Use when two streams should start emitting custom combined values after both have produced at least one value.
Gotchas
Note: tracking the latest value is done on a per-array basis. That means that emitted elements that are not the last value in arrays will never be used for zipping.
Signature
declare const zipLatestWith: {
<AR, ER, RR, AL, A>(
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): Stream<A, EL | ER, RL | RR>;
};Zips this stream with another point-wise and keeps only the values from the left stream.
Details
The resulting stream ends when either side ends.
Signature
declare const zipLeft: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AL, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AL, EL | ER, RL | RR>;
};Zips this stream with another point-wise, keeping only right values and ending when either stream ends.
Signature
declare const zipRight: {
<AR, ER, RR>(
right: Stream<AR, ER, RR>,
): <AL, EL, RL>(left: Stream<AL, EL, RL>) => Stream<AR, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
): Stream<AR, EL | ER, RL | RR>;
};Zips two streams point-wise with a combining function, ending when either stream ends.
Signature
declare const zipWith: {
<AR, ER, RR, AL, A>(
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(
left: Stream<AL, EL, RL>,
right: Stream<AR, ER, RR>,
f: (left: AL, right: AR) => A,
): Stream<A, EL | ER, RL | RR>;
};zipWithArray
Zips two streams by applying a function to non-empty arrays of elements.
Details
The function returns output plus leftover arrays that carry into the next pull.
Signature
declare const zipWithArray: {
<AR, ER, RR, AL, A>(right: Stream<AR, ER, RR>, f: (left: readonly [AL, AL], right: readonly [AR, AR]) => readonly [readonly [A, A], readonly Array<AL>, readonly Array<AR>]): <EL, RL>(left: Stream<AL, EL, RL>) => Stream<A, ER | EL, RR | RL>;
<AL, EL, RL, AR, ER, RR, A>(left: Stream<AL, EL, RL>, right: Stream<AR, ER, RR>, f: (left: readonly [AL, AL], right: readonly [AR, AR]) => readonly [readonly [A, A], readonly Array<AL>, readonly Array<AR>]): Stream<A, EL | ER, RL | RR>;
}zipWithIndex
Zips this stream together with the index of elements.
Signature
declare function zipWithIndex<A, E, R>(self: Stream<A, E, R>): Stream<[A, number], E, R>;zipWithNext
Zips each element with the next element, pairing the final element with Option.none().
Signature
declare function zipWithNext<A, E, R>(self: Stream<A, E, R>): Stream<[A, Option<A>], E, R>;zipWithPrevious
Zips each element with its previous element, starting with None.
Signature
declare function zipWithPrevious<A, E, R>(self: Stream<A, E, R>): Stream<[Option<A>, A], E, R>;zipWithPreviousAndNext
Zips each element with its previous and next values.
Signature
declare function zipWithPreviousAndNext<A, E, R>(
self: Stream<A, E, R>,
): Stream<[Option<A>, A, Option<A>], E, R>;
Accesses a service from the context and emits it as a single element.