variable / @typed/fx/Fx

toStream

Adapts a push-based Fx to an Effect Stream.

Package version
2.0.0-beta.7
Category
Stream interop
Since
1.0.0

Follow one execution

Adapt push deliveries through a scoped queue for Stream consumption; this trace assumes capacity is available and the reader keeps up.

Fx timelinetoStream(fx)toStream
OperatortoStream(fx)
Read this diagram

Follow each lane from left to right. Events stacked vertically share a tick; the green cursor marks the current time across every lane.

  • A value
    The text inside the pill is the emitted value.
  • Work starts
    The raised chevron starts an inner run (^ in the source).
  • The run returns
    The vertical bar ends this lane’s run.
  • A cause is delivered
    The exclamation mark belongs to this lane.
  • Work is interrupted
    The cross marks cancellation of this run.
  • Current time
    The line and diamond move together across all lanes.
  • Happening now
    A highlighted event is at the current tick.
  • Still ahead
    Muted, dashed values have not happened yet.
  • Time continues
    The lane’s arrow is not a return marker. An empty stretch can be quiet work that is still running.

Illustrated ticks start at 0. At 1×, one illustrated tick takes one second; captions specify real durations when timing matters. A cause or interruption belongs to its lane, and other work may continue. Scroll horizontally to inspect the rest of a long timeline.

Conversion is lazy. Running the Stream allocates the callback queue and runs the Fx; the Stream’s scope owns both. Interruption closes the callback subscription through Effect Stream’s lifecycle. Buffering follows options.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

import { toStream } from "@typed/fx/Fx";

Signatures

export declare const toStream: <A, E, R>(fx: Fx.Fx<A, E, R>, options?: ToStreamOptions) => Stream.Stream<A, E, R>;

Why

Existing Effect Stream consumers can use an Fx without losing the producer’s typed values, failures, or services. Values and causes cross a queue in emission order; the queue ends when the Fx run completes.

Ownership and lifetime

Conversion is lazy. Running the Stream allocates the callback queue and runs the Fx; the Stream’s scope owns both. Interruption closes the callback subscription through Effect Stream’s lifecycle. Buffering follows options.

Examples

import { Stream } from "effect"
import { fromIterable, toStream } from "@typed/fx/Fx"

const values = toStream(fromIterable([1, 2, 3]))
const total = Stream.runFold(values, () => 0, (sum, value) => sum + value)

Other public imports

These import paths expose the same declaration. Each page retains its own public name and signature.

Source