Skip to content

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.

244 exports Added in v2.0.0 Source

Accessors

service

Added in v4.0.0 Source

Accesses a service from the context and emits it as a single element.

Signature

declare function service<I, S>(service: Key<I, S>): Stream<S, never, I>;

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

Added in v2.0.0 Source

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>;

collect

Added in v4.0.0 Source

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>;

scan

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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

aggregate

Added in v2.0.0 Source

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>;
};

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>;
};

transduce

Added in v2.0.0 Source

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

broadcast

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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>;
};

share

Added in v3.8.0 Source

Returns a new Stream that multicasts the original stream, subscribing when the first consumer starts.

Details

The upstream continues running while there is at least one consumer and is finalized after the last one exits. If idleTimeToLive is set, the upstream is kept alive for that duration so a later subscriber can continue from the next element instead of restarting.

Signature

declare const share: {
  (
    options:
      | {
          readonly capacity: "unbounded";
          readonly idleTimeToLive?: Duration.Input;
          readonly replay?: number;
        }
      | {
          readonly capacity: number;
          readonly idleTimeToLive?: Duration.Input;
          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 idleTimeToLive?: Duration.Input;
          readonly replay?: number;
        }
      | {
          readonly capacity: number;
          readonly idleTimeToLive?: Duration.Input;
          readonly replay?: number;
          readonly strategy?: "sliding" | "dropping" | "suspend";
        },
  ): Effect<Stream<A, E, never>, never, Scope | R>;
};

Buffering

buffer

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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

The default chunk size used by Stream constructors and combinators.

Signature

declare const DefaultChunkSize: number;

Constructors

callback

Added in v4.0.0 Source

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>>;

die

Added in v2.0.0 Source

The stream that dies with the specified defect.

Signature

declare function die(defect: unknown): Stream<never>;

Do

Added in v2.0.0 Source

Provides the entry point for do-notation style stream composition.

Signature

declare const Do: Stream<{}>;

empty

Added in v2.0.0 Source

Creates an empty stream.

Signature

declare const empty: Stream<never>;

fail

Added in v2.0.0 Source

Terminates with the specified error.

Signature

declare function fail<E>(error: E): Stream<never, E>;

failCause

Added in v2.0.0 Source

Creates a stream that fails with the specified Cause.

Signature

declare function failCause<E>(cause: Cause<E>): Stream<never, E>;

The stream that always fails with the specified lazily evaluated Cause.

Signature

declare function failCauseSync<E>(evaluate: LazyArg<Cause<E>>): Stream<never, E>;

failSync

Added in v2.0.0 Source

Terminates with the specified lazily evaluated error.

Signature

declare function failSync<E>(evaluate: LazyArg<E>): Stream<never, E>;

fromArray

Added in v4.0.0 Source

Creates a stream from an array of values.

Signature

declare function fromArray<A>(array: readonly Array<A>): Stream<A>

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

Added in v4.0.0 Source

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]>

Creates a stream from an AsyncIterable.

Signature

declare function fromAsyncIterable<A, E>(
  iterable: AsyncIterable<A>,
  onError: (error: unknown) => E,
): Stream<A, E>;

fromChannel

Added in v2.0.0 Source

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

Added in v2.0.0 Source

Creates a stream from an effect.

Signature

declare function fromEffect<A, E, R>(effect: Effect<A, E, R>): Stream<A, E, R>;

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>;

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>;

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>;

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

Added in v2.0.0 Source

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>;

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>;

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>;

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

Added in v2.0.0 Source

Creates a stream from a subscription to a PubSub.

Signature

declare function fromPubSub<A>(pubsub: PubSub<A>): Stream<A>;

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>;

fromPull

Added in v2.0.0 Source

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>;

fromQueue

Added in v2.0.0 Source

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>>>;

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

Added in v2.0.0 Source

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>;

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>;

iterate

Added in v2.0.0 Source

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>;

make

Added in v2.0.0 Source

Creates a stream from a sequence of values.

Signature

declare function make<As extends readonly Array<any>>(...values: As): Stream<As[number]>

never

Added in v2.0.0 Source

The stream that never produces any value or fails with any error.

Signature

declare const never: Stream<never>;

paginate

Added in v2.0.0 Source

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>

range

Added in v2.0.0 Source

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>;

scoped

Added in v2.0.0 Source

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>>;

succeed

Added in v2.0.0 Source

Creates a single-valued pure stream.

Signature

declare function succeed<A>(value: A): Stream<A>;

suspend

Added in v2.0.0 Source

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>;

sync

Added in v2.0.0 Source

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>;

tick

Added in v2.0.0 Source

Creates a stream that emits void immediately once, then emits another void after each specified interval.

Signature

declare function tick(interval: Input): Stream<void>;

toChannel

Added in v2.0.0 Source

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>;

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>;

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>;

unfold

Added in v2.0.0 Source

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>;

unwrap

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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

changes

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
};

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

timeout

Added in v2.0.0 Source

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>;
};

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

  • timeout for 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

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>;

mkString

Added in v2.0.0 Source

Concatenates all emitted strings into a single string.

Signature

declare function mkString<E, R>(self: Stream<string, E, R>): Effect<string, E, R>;

mkUint8Array

Added in v4.0.0 Source

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>;

peel

Added in v2.0.0 Source

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>;
};

run

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;

runCount

Added in v2.0.0 Source

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>;

runDrain

Added in v2.0.0 Source

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>;

runFold

Added in v2.0.0 Source

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>;
};

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

Added in v2.0.0 Source

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>;
};

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>;
};

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>;
};

runHead

Added in v2.0.0 Source

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>;

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

Added in v2.0.0 Source

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>;
};

runLast

Added in v2.0.0 Source

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

  • runHead for consuming only the first emitted element
  • runCollect for collecting every emitted element
  • runDrain for consuming the stream while discarding emitted elements

Signature

declare function runLast<A, E, R>(self: Stream<A, E, R>): Effect<Option<A>, E, R>;

runSum

Added in v2.0.0 Source

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>;

Converts a stream to an AsyncIterable for for await...of consumption.

Signature

declare function toAsyncIterable<A, E>(self: Stream<A, E>): AsyncIterable<A>;

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>;

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>;
};

toPubSub

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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>;
};

toPull

Added in v2.0.0 Source

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>;

toQueue

Added in v2.0.0 Source

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>;
};

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>;
};

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>;
};

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

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
};

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

  • catchCauseIf for predicate-based cause selection
  • catchFilter for filtering typed error values instead of full causes
  • catchCause for 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

Added in v4.0.0 Source

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

Added in v4.0.0 Source

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

  • catchIf for predicate or refinement based recovery
  • catchTag for _tag based recovery from one tagged error
  • catchTags for _tag based recovery from multiple tagged errors
  • catchCauseFilter for 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
  >;
};

catchIf

Added in v4.0.0 Source

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

Added in v4.0.0 Source

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

Added in v4.0.0 Source

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]
  >;
};

catchTag

Added in v2.0.0 Source

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
  >;
};

catchTags

Added in v2.0.0 Source

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]
  >;
};

ignore

Added in v4.0.0 Source

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

  • ignoreCause for 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

Added in v4.0.0 Source

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

  • ignore to 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>;

mapError

Added in v2.0.0 Source

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>;
};

onError

Added in v2.0.0 Source

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>;
};

orDie

Added in v2.0.0 Source

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>;

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>;
};

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>;
};

result

Added in v4.0.0 Source

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>;

retry

Added in v2.0.0 Source

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>;
}

tapCause

Added in v2.0.0 Source

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>;
};

tapError

Added in v2.0.0 Source

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>;
};

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

drop

Added in v2.0.0 Source

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>;
};

dropRight

Added in v2.0.0 Source

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>;
};

dropUntil

Added in v2.0.0 Source

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>;
};

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>;
};

dropWhile

Added in v2.0.0 Source

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>;
};

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>;
};

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

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>;
};

filter

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
};

filterMap

Added in v2.0.0 Source

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

  • filter for keeping original elements with a boolean predicate or refinement
  • filterMapEffect for an effectful Filter
  • partition for 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>;
};

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

  • filterMap for the synchronous Filter variant
  • filterEffect for effectfully keeping original elements
  • mapEffect for 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

Added in v4.0.0 Source

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>;
};

partition

Added in v2.0.0 Source

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
  >;
};

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

  • partition for the pure Filter variant, which returns the failing stream before the passing stream
  • partitionQueue for the lower-level queue result
  • filterMapEffect for 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
  >;
};

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
  >;
};

take

Added in v2.0.0 Source

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>;
};

takeRight

Added in v2.0.0 Source

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>;
};

takeUntil

Added in v2.0.0 Source

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>;
};

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>;
};

takeWhile

Added in v2.0.0 Source

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>;
};

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>;
};

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

  • takeWhile for keeping original elements with a boolean predicate or refinement
  • filterMap for filtering across the whole stream instead of only the leading prefix
  • dropWhileFilter for 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>;
};

when

Added in v2.0.0 Source

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

chunks

Added in v2.0.0 Source

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>;

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

  • groupByKey for grouping all elements with the same key across the stream
  • groupBy for 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>;
};

groupBy

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
};

grouped

Added in v2.0.0 Source

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>;
};

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>;
};

rechunk

Added in v2.0.0 Source

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>;
};

sliding

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
};

split

Added in v2.0.0 Source

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

isStream

Added in v4.0.0 Source

Checks whether a value is a Stream.

Signature

declare function isStream(u: unknown): u is Stream<unknown, unknown, unknown>;

Interruption

haltWhen

Added in v2.0.0 Source

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>;
};

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

bindTo

Added in v2.0.0 Source

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>;
};

flatMap

Added in v2.0.0 Source

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>;
};

flatten

Added in v2.0.0 Source

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>;

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>;

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>;

map

Added in v2.0.0 Source

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>;
};

mapAccum

Added in v2.0.0 Source

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>;
}

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>;
}

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>;
}

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>;
}

mapArray

Added in v4.0.0 Source

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>;
};

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>;
};

mapBoth

Added in v2.0.0 Source

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>;
};

mapEffect

Added in v2.0.0 Source

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

combine

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
};

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>;
};

merge

Added in v2.0.0 Source

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>;
};

mergeAll

Added in v2.0.0 Source

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

  • merge for merging exactly two streams and choosing a halt strategy
  • flatten for flattening a stream that already emits streams

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

Added in v4.0.0 Source

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>;
};

mergeLeft

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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

Added in v2.0.0 Source

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

Added in v3.4.0 Source

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

Added in v4.0.0 Source

Describes how merged streams decide when to halt.

Signature

type HaltStrategy = Channel.HaltStrategy;

Stream interface

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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

Added in v2.0.0 Source

Type-level marker that excludes Stream from unification.

Signature

interface StreamUnifyIgnore {
  Effect?: true;
}

Variance interface

Added in v2.0.0 Source

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

Added in v3.4.0 Source

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

provide

Added in v4.0.0 Source

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>>;
};

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>>;
};

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>>;
};

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>>;
};

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>;
};

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

race

Added in v3.7.0 Source

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>;
};

raceAll

Added in v3.5.0 Source

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

debounce

Added in v2.0.0 Source

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>;
};

schedule

Added in v2.0.0 Source

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>;
};

throttle

Added in v2.0.0 Source

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>;
};

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

ensuring

Added in v2.0.0 Source

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>;
};

onExit

Added in v4.0.0 Source

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

bind

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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>;
};

concat

Added in v2.0.0 Source

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>;
};

drain

Added in v2.0.0 Source

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>;

drainFork

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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

Added in v2.0.0 Source

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>;

forever

Added in v2.0.0 Source

Repeats this stream forever.

Signature

declare function forever<A, E, R>(self: Stream<A, E, R>): Stream<A, E, R>;

intersperse

Added in v2.0.0 Source

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>;
};

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>;
};

onEnd

Added in v3.6.0 Source

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>;
};

onFirst

Added in v4.0.0 Source

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>;
};

onStart

Added in v3.6.0 Source

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

Added in v2.0.0 Source

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>;
};

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>;
};

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>;
};

prepend

Added in v2.0.0 Source

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>;
};

repeat

Added in v2.0.0 Source

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>;
}

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>;
};

switchMap

Added in v4.0.0 Source

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>;
};

tap

Added in v2.0.0 Source

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>;
};

tapBoth

Added in v2.0.0 Source

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>;
};

tapSink

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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

withSpan

Added in v2.0.0 Source

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

TypeId

Added in v4.0.0 Source

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

  • isStream for the public guard that checks this identifier

Signature

declare const TypeId: "~effect/Stream";

TypeId type

Added in v4.0.0 Source

String literal type used as the unique brand for Stream values.

Signature

type TypeId = "~effect/Stream";

Utility Types

Error type

Added in v3.4.0 Source

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;

Services type

Added in v4.0.0 Source

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

Added in v2.0.0 Source

Type lambda for Stream used in higher-kinded type operations.

Signature

interface StreamTypeLambda extends TypeLambda {
  readonly type: Stream<unknown, unknown, unknown>;
}

Success type

Added in v3.4.0 Source

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

cross

Added in v2.0.0 Source

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>;
};

crossWith

Added in v2.0.0 Source

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>;
};

zip

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;
}

zipLatest

Added in v2.0.0 Source

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

Added in v3.3.0 Source

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>

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>;
};

zipLeft

Added in v2.0.0 Source

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>;
};

zipRight

Added in v2.0.0 Source

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>;
};

zipWith

Added in v2.0.0 Source

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

Added in v4.0.0 Source

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

Added in v2.0.0 Source

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

Added in v2.0.0 Source

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>;

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>;

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>;