Module / v2.0.0-beta.13

@typed/fx/Fx

176 unique exports available through this public import path, including their members and aliases.

Explore every operation with marble diagrams →

import * as Fx from "@typed/fx/Fx";

Callback protocol

  • Emit

    type-alias

    Operations supplied to a callback producer for emitting values or ending its run.

Callback sources

  • callback

    variable

    Creates an Fx from a callback-based source.

  • make

    variable

    Creates an Fx from a function that provides values to a Sink.

    This is the lowest-level constructor for Fx, giving you full control over the stream’s behavior.

Collecting values

  • collectAll

    variable

    Collects all values emitted by an Fx into an array.

  • collectAllFork

    variable

    Forks the collection of all values from an Fx.

  • collectUpTo

    variable

    Collects the first n values emitted by an Fx into an array.

  • collectUpToFork

    variable

    Forks the collection of up to n values from an Fx.

  • first

    function

    Returns the first value emitted by the Fx wrapped in an Option. If the Fx is empty, returns None.

Combining sources

  • append

    variable

    Appends a value to the end of an Fx.

  • concat

    variable

    Concatenates two Fx streams: runs the first to completion, then runs the second. Emits all values from the first stream in order, then all values from the second stream in order.

    Completion: The concatenated stream completes when the second stream completes (the first must complete before the second starts).

    Failures: A source Cause is delivered to the Sink. Because Fx.run is infallible, delivery alone does not suppress the continuation: the right source is run after the left run returns.

  • continueWith

    variable

    Continues an Fx with a lazily created Fx after the first run returns.

  • delimit

    variable

    Wraps an Fx with a start and end value.

  • merge

    variable

    Merges two Fx streams into a single Fx that emits values from both streams concurrently. Order of emission is non-deterministic.

    Completion: The merged stream completes when both input streams have completed.

    Failures: Every failure Cause is delivered to the downstream Sink. Delivery does not make Fx.run fail, so merge does not itself cancel the sibling; a terminal observer may choose to.

  • mergeAll

    variable

    Merges multiple Fx streams into a single Fx that emits values from all input streams concurrently.

  • mergeLeft

    variable

    Merges two Fx streams and emits only values from the left stream. Both streams run concurrently; completion when both complete.

  • mergeOrdered

    function

    Runs multiple Fx streams concurrently while draining their values in argument order.

  • mergeRight

    variable

    Merges two Fx streams and emits only values from the right stream. Both streams run concurrently; completion when both complete.

  • prepend

    variable

    Prepends a value to the beginning of an Fx.

  • struct

    function

    Combines a record of Fx streams into a single Fx that emits a record of the latest values. Similar to tuple, but for objects.

  • tuple

    function

    Combines multiple Fx streams into a single Fx that emits a tuple of the latest values from each stream. The resulting Fx waits for all input streams to emit at least once before emitting the first tuple. Afterwards, it emits a new tuple whenever any input stream emits a new value.

  • withLatestFrom

    variable

    Emits [source, latest] whenever the source emits, using the latest value from that. Source values are dropped until that has emitted at least once.

    Unlike {@link zipLatest}, this does not emit when that updates.

    Completion: Completes when the source completes. Errors: The first failure from either stream fails the result.

  • withLatestFromWith

    variable

    Like {@link withLatestFrom}, but combines the pair with f.

  • zip

    variable

    Zips two Fx streams in strict lockstep: emits a pair [a, b] only when both streams have produced their next value. Emits the i-th pair when both have produced at least i values.

    Completion: The zipped stream completes when the first of the two streams completes (no further pairs are emitted). The other stream is interrupted.

    Errors: The first failure from either stream fails the zipped stream.

  • zipLatest

    variable

    Zips two Fx streams by latest values: waits for both to emit at least once, then emits [left, right] whenever either stream emits (using the latest value from the other). No strict pairing; output count is the sum of emissions from both after the first pair.

    Completion: Completes when both streams have completed. Errors: The first failure from either stream fails the result.

  • zipLatestWith

    variable

    Zips two Fx streams by latest values and combines each pair with a function. Waits for both to emit at least once, then emits f(left, right) whenever either stream emits.

    Completion: Completes when both streams have completed. Errors: The first failure from either stream fails the result.

  • zipLeft

    variable

    Zips two Fx streams in strict lockstep and emits only the left value. Completes when the first of the two streams completes.

  • zipRight

    variable

    Zips two Fx streams in strict lockstep and emits only the right value. Completes when the first of the two streams completes.

  • zipWith

    variable

    Zips two Fx streams in strict lockstep and combines each pair with a function. Emits f(a, b) when both streams have produced their next value.

    Completion: Completes when the first of the two streams completes. Errors: The first failure from either stream fails the result.

Concurrent work

  • concatMap

    variable

    Maps each element to an inner Fx and concatenates the results sequentially.

  • concatMapEffect

    variable

    Maps each element to an Effect and concatenates the results sequentially.

  • exhaustLatestMap

    variable

    Maps each element to an inner Fx, running one now and retaining only the latest waiting value.

  • exhaustLatestMapEffect

    variable

    Maps each element to an Effect, running one now and retaining only the latest waiting value.

  • exhaustMap

    variable

    Maps each element of an Fx to a new Fx, ignoring new elements until the current inner Fx completes.

  • exhaustMapEffect

    variable

    Maps each element of an Fx to an Effect, ignoring new elements until the current effect completes.

  • flatMap

    variable

    Maps each source value to an inner Fx and merges every inner concurrently.

  • flatMapConcurrently

    variable

    Maps each element of an Fx to a new Fx, running them concurrently with a limit.

  • flatMapConcurrentlyEffect

    variable

    Maps each element of an Fx to an Effect, running them concurrently with a limit.

  • flatMapEffect

    variable

    Maps each element of an Fx to an Effect, and merges the results.

  • race

    variable

    Runs two streams concurrently until one emits, then mirrors the winner and interrupts the other.

    A failure or completion from one side before the other emits does not win unless every side ends without emitting. After a winner is chosen, that stream’s later failures are propagated.

  • raceAll

    variable

    Races many streams: the first to emit wins and the rest are interrupted.

  • switchMap

    variable

    Maps each element of an Fx to a new Fx, and switches to the latest inner Fx.

    When a new element is emitted, the previous inner Fx is cancelled.

  • switchMapEffect

    variable

    Maps each element of an Fx to an Effect, and switches to the latest effect.

    When a new element is emitted, the previous effect is cancelled.

Conditional sources

  • if

    variable · Re-export

    Conditionally runs one of two Fx streams based on the boolean value emitted by the condition stream.

  • when

    variable

    Conditionally emits one of two values based on the boolean value emitted by the condition stream.

Effect interop

  • fromEffect

    variable

    Creates an Fx from an Effect.

    If the Effect succeeds, the Fx emits the value and completes. If the Effect fails, the Fx fails with the same error.

Errors and recovery

  • catch

    variable · Re-export

    Recovers from the first typed failure of an Fx by running a fallback Fx.

  • catchAll

    variable

    Uses the Effect-style catchAll name for {@link catch}.

  • catchCause

    variable

    Recovers from any failure cause by running a fallback Fx.

  • catchCauseIf

    variable

    Recovers a failure cause only when a predicate accepts the complete cause.

  • catchIf

    variable

    Recovers a typed failure only when a predicate accepts it.

  • catchTag

    variable

    Recovers selected tagged typed failures by running a fallback Fx.

  • catchTags

    variable

    Recovers several tagged typed-error variants with one handler table.

  • catch_

    variable

    Recovers from the first typed failure of an Fx by running a fallback Fx.

  • causes

    variable

    Emits the source’s terminal failure cause and discards every successful value.

  • exit

    variable

    Materializes every success and the terminal failure as infallible Exit values.

  • flip

    variable

    Emits typed failures as values and fails with the first successful value.

  • mapBoth

    variable

    Transforms both the success and error channels of an Fx using the provided options.

    Mirrors Effect.mapBoth: onSuccess maps emitted values, onFailure maps the typed failure (via Cause.map); defects and interrupts are preserved.

  • mapError

    variable

    Transforms typed failures while preserving defects and interruption.

  • result

    variable

    Materializes success and failure of an Fx as Result values.

    • Success: each emitted value is wrapped as Result.succeed(value).
    • Failure: any failure (including typed error, defect, and interrupt) is materialized as Result.fail(cause). The output error type is Cause<E>, so defects and interrupts are explicitly represented in the Result and the resulting Fx has error type never.

    The resulting Fx never fails at the stream level; all outcomes are emitted as Result<A, Cause<E>>. Consumers can use Result.match or Result.isSuccess / Result.isFailure to handle success vs failure (including defect/interrupt).

  • retry

    variable

    Retries the entire stream when its Cause contains a typed Fail accepted by schedule.

    The schedule is reset as soon as the first element of an attempt is emitted, matching Effect Stream.retry.

Failure sources

  • die

    variable

    Creates an Fx that immediately terminates with a defect (unexpected error).

  • fail

    variable

    Creates an Fx that immediately fails with the specified error.

  • failCause

    variable

    Creates an Fx that immediately terminates with the specified Cause.

  • fromFailures

    variable

    Creates an Fx from a collection of failures (errors).

  • interrupt

    variable

    Creates an Fx that immediately interrupts.

Generator composition

  • fn

    namespace

    Callable contracts implemented by fn.

  • gen

    variable

    Builds an Fx by yielding Effects and returning the Fx to run afterward.

  • genScoped

    variable

    Builds an Fx with a subscription-owned Scope shared by setup and streaming.

  • unwrap

    variable

    Unwraps an Effect that produces an Fx into a single Fx.

  • unwrapScoped

    variable

    Unwraps an Effect that produces an Fx into a single Fx, managing the scope of the effect.

    The scope of the effect is closed when the Fx completes or is interrupted.

Keyed work

  • keyed

    variable

    Efficiently transforms a list of values into a list of Fx streams, using keys to track identity.

    This is crucial for performance when rendering lists or managing collections of stateful entities. When the input list changes:

    • New keys cause onValue to be called.
    • Existing keys have their RefSubject updated with the new value.
    • Removed keys close the supplied child Scope and clean resources registered through it; the onValue run fiber remains owned by the outer parent Scope.

Observing failures

  • onError

    variable

    Runs cleanup after the source reports a failure cause.

  • withSpan

    variable

    Traces the whole subscription and each success or failure delivery.

Operator options

  • Bounds

    interface

    Defines the bounds for slicing an Fx stream.

  • FromStreamOptions

    type-alias

    Effect Stream mapping options used while delivering elements to an Fx sink.

  • KeyedOptions

    interface

    Configuration options for the keyed combinator.

  • ThrottleOptions

    type-alias

    Options for {@link throttle}. A duration-only call is leading-edge ({ leading: true, trailing: false }).

  • ToStreamOptions

    type-alias

    Buffering and callback options accepted while adapting an Fx to an Effect Stream.

Providing services

  • Fx.Service

    interface

    An Fx whose implementation is obtained from an Effect service.

  • Service

    function

    Defines an Effect service whose value is also a directly runnable Fx.

  • provide

    variable

    Builds a Layer for each subscription and provides it to the entire Fx run.

  • provideContext

    variable

    Provides an already-built Effect Context to the entire Fx run.

  • provideService

    variable

    Provides one existing service value to the entire Fx run.

  • provideServiceEffect

    variable

    Acquires one service with an Effect before running the Fx.

Resource lifetime

  • ensuring

    variable

    Runs a finalizer with no typed error after the Fx run ends for any reason.

  • onExit

    variable

    Observes the Fx’s final success or failure with an Effect finalizer.

  • onInterrupt

    variable

    Runs a finalizer when the Fx reports or externally receives interruption.

Running effects

  • drain

    variable

    Runs an Fx stream to completion, discarding all values. Useful when the side effects of the stream are all that matter.

  • drainLayer

    variable

    Runs an Fx stream as a Layer. The stream is forked in the background when the layer is acquired.

  • fork

    variable

    Forks the execution of an Fx into a background fiber. The stream will run until it completes or the fiber is interrupted.

  • observe

    variable

    Observes the values of an Fx stream using a callback function. The callback can return void or an Effect which will be executed for each value.

  • observeLayer

    variable

    Observes the values of an Fx stream using a callback function and returns a Layer. The callback can return void or an Effect which will be executed for each value.

  • runFork

    variable

    Runs an Fx in a new fiber, using the standard Effect.runFork. This is useful for integrating with the top-level Effect runtime.

  • runPromise

    variable

    Runs an Fx stream to completion and returns a Promise. Rejects if the stream fails.

  • runPromiseExit

    variable

    Runs an Fx stream to completion and returns a Promise of the Exit.

Runtime inspection

  • FxTypeId

    variable

    Runtime symbol carried by every Fx implementation.

  • isFx

    function

    Checks whether a value carries the FxTypeId protocol property.

Selecting values

  • changesWithEffect

    variable

    Drops consecutive elements that are considered equal by an effectful predicate. When the effect returns true, the element is skipped; when false, it is emitted.

    This is the effectful variant of skipRepeatsWith: instead of a pure Equivalence<A>, you supply (prev, next) => Effect<boolean> where true means “equal” (skip) and false means “changed” (emit).

  • compact

    variable

    Compacts an Fx of Options, discarding None values and unwrapping Some values.

  • dropAfter

    variable

    Drops elements from an Fx after a predicate returns true. The element that satisfies the predicate is included in the output.

  • dropUntil

    variable

    Drops elements from an Fx until a predicate returns true. Emits from the first element for which the predicate returns true (including that element) and all following elements.

  • dropUntilEffect

    variable

    Drops elements from an Fx until an effectful predicate returns true. Emits from the first element for which the predicate effect succeeds with true (including that element) and all following elements.

  • dropWhile

    variable

    Alias of skipWhile for Effect parity (dropWhile naming).

  • dropWhileEffect

    variable

    Alias of skipWhileEffect for Effect parity (dropWhileEffect naming).

  • filter

    variable

    Filters elements of an Fx using a predicate function.

  • filterEffect

    variable

    Filters elements of an Fx using an effectful predicate function.

  • filterMap

    variable

    Maps and filters elements of an Fx in a single operation.

  • filterMapEffect

    variable

    Maps and filters elements of an Fx using an effectful function.

  • skip

    variable

    Skips the first n elements of an Fx.

  • skipEffect

    variable

    Skips the first n elements where n is produced by an Effect.

  • skipRepeats

    variable

    Drops elements that are equal to the previous element using standard equality.

  • skipRepeatsWith

    variable

    Drops elements that are equal to the previous element using a custom equivalence function.

  • skipWhile

    variable

    Skips elements from an Fx while a predicate returns true. Emits from the first element for which the predicate returns false (including that element) and all following elements.

  • skipWhileEffect

    variable

    Skips elements from an Fx while an effectful predicate returns true. Emits from the first element for which the predicate effect succeeds with false (including that element) and all following elements.

  • slice

    variable

    Slices an Fx by skipping a number of elements and then taking a number of elements.

  • sliceEffect

    variable

    Slices an Fx with bounds produced by an Effect.

  • take

    variable

    Takes the first n elements from an Fx and then completes.

  • takeEffect

    variable

    Takes the first n elements where n is produced by an Effect.

  • takeUntil

    variable

    Takes elements from an Fx until a predicate returns true. The element that satisfies the predicate is not included in the output.

  • takeUntilEffect

    variable

    Takes elements from an Fx until an effectful predicate returns true. The element that satisfies the predicate is not included in the output.

  • takeWhile

    variable

    Takes elements from an Fx while a predicate returns true. Stops at the first element for which the predicate returns false; that element is not included.

  • takeWhileEffect

    variable

    Takes elements from an Fx while an effectful predicate returns true. Stops at the first element for which the predicate effect succeeds with false; that element is not included.

Stateful transforms

  • filterMapLoop

    variable

    Loops over an Fx with an accumulator, producing an optional new value for each element. If the function returns None, the element is filtered out.

  • filterMapLoopCause

    variable

    Loops over the failure causes of an Fx with an accumulator, potentially transforming or filtering them. This allows for complex error handling logic that maintains state across failures.

  • filterMapLoopCauseEffect

    variable

    Effectfully loops over the failure causes of an Fx with an accumulator.

  • filterMapLoopEffect

    variable

    Effectfully loops over an Fx with an accumulator, producing an optional new value.

  • grouped

    variable

    Partitions the stream into non-empty arrays of size n. The final array may be smaller if there are leftover elements. The size must be a positive safe integer. A group can retain up to n values, so callers own the memory policy for valid sizes. Invalid sizes fail with Cause.IllegalArgumentError.

    Matches Effect Stream.grouped.

  • loop

    variable

    Loops over an Fx with an accumulator, producing a new value for each element and updating the accumulator.

  • loopCause

    variable

    Loops over the failure causes of an Fx with an accumulator.

  • loopCauseEffect

    variable

    Effectfully loops over the failure causes of an Fx with an accumulator.

  • loopEffect

    variable

    Effectfully loops over an Fx with an accumulator, producing a new value for each element.

  • pairwise

    variable

    Emits consecutive pairs [previous, current]. The first value is not emitted until a second value arrives.

    Equivalent to RxJS pairwise and Effect Stream.sliding(2) for pairs.

  • scan

    variable

    Scans the stream with a pure function, emitting the accumulated state after each element. Emits the initial value first, then for each input a emits f(state, a) and updates state.

    Semantics align with Effect Stream’s scan: output is initial, f(initial, a1), f(..., a2), …

  • scanEffect

    variable

    Scans the stream with an effectful function, emitting the accumulated state after each element. Emits the initial value first, then for each input a runs f(state, a) and emits the resulting state.

Stream interop

  • fromStream

    variable

    Adapts an Effect Stream to a push-based Fx.

  • toStream

    variable

    Adapts a push-based Fx to an Effect Stream.

Time and rate

  • at

    variable

    Creates an Fx that emits a single value after a specified delay.

  • debounce

    variable

    Emits a value only after no newer source value arrives for duration.

  • delay

    variable

    Sleeps for duration before forwarding each successful source delivery.

  • during

    function

    Forwards events only between a start signal and that signal’s first stop event.

  • fromSchedule

    variable

    Creates an Fx that emits values according to a Schedule. The Fx emits void each time the schedule fires.

  • groupedWithin

    variable

    Partitions the stream into arrays, emitting when n is reached or duration elapses after the first element of the current group. The size must be a positive safe integer. A group can retain up to n values, so callers own the memory policy for valid sizes. Invalid sizes fail with Cause.IllegalArgumentError.

    Matches Effect Stream.groupedWithin.

  • periodic

    variable

    Creates an Fx that emits a void value periodically.

  • repeat

    variable

    Repeats the entire stream according to schedule after each successful completion. Failures are not repeated.

    Schedule.recurs(n) runs the stream n + 1 times (the original plus n repeats), matching Effect Stream.repeat.

  • sample

    variable

    Emits the latest source value whenever sampler emits. Source values that arrive between sampler ticks are not forwarded until the next tick.

    Completion: Completes when the source completes. Errors: The first failure from either stream fails the result.

  • since

    variable

    Drops events until signal emits, then forwards the rest.

  • throttle

    variable

    Limits emissions to configured leading and trailing edges of fixed windows.

    Pass { duration, leading, trailing } for trailing or both-edge behavior. A duration alone defaults to { leading: true, trailing: false }.

  • timeout

    variable

    Completes the stream if it does not produce a value (or complete) within duration of the previous event. Matches Effect Stream.timeout.

    The timeout is reset after each emission. An infinite duration is a no-op; a zero duration completes immediately.

  • timeoutTo

    variable

    Switches to fallback if the source does not produce a value within duration of the previous event. Matches Effect Stream.timeoutOrElse and RxJS timeoutTo.

  • until

    variable

    Forwards events until signal emits, then interrupts events.

Transforming values

  • as

    variable

    Replaces all emitted values from the Fx with the provided value b.

  • map

    variable

    Transforms the elements of an Fx using a provided function.

  • mapEffect

    variable

    Transforms the elements of an Fx using a provided Effectful function.

  • tap

    variable

    Performs an effect for each element of the Fx, without changing the elements.

Type contracts

  • Error

    type-alias

    Alias of Fx.Error for extracting an Fx error channel.

  • FlatMapEffectLike

    type-alias

    Describes a dual flattening operator whose callback returns an Effect.

  • FlatMapLike

    type-alias

    Describes a dual flattening operator whose callback returns an Fx.

  • Fx

    interface

    Fx is a reactive stream of values that supports concurrency, error handling, and context management, fully integrated with the Effect ecosystem.

    Conceptually, an Fx<A, E, R> is a push-based stream that:

    • Emits values of type A
    • Can fail with an error of type E
    • Requires a context/environment of type R

    Unlike a standard Effect which produces a single value, Fx can produce 0, 1, or many values over time. It is similar to RxJS Observables or AsyncIterables, but built on top of Effect’s fiber-based concurrency model.

  • Fx.Error

    type-alias

    Extracts the typed error from an Fx.

  • Fx.Services

    type-alias

    Extracts the services required to run an Fx.

  • Fx.Success

    type-alias

    Extracts the emitted value type from an Fx.

  • Services

    type-alias

    Alias of Fx.Services for extracting an Fx service channel.

  • Success

    type-alias

    Alias of Fx.Success for extracting an Fx value channel.

  • fn.Gen

    type-alias

    Contract for functions whose body yields Effects and returns an Fx.

  • fn.NonGen

    type-alias

    Contract for functions whose body returns an Fx directly.

Value sources

  • empty

    variable

    An Fx that emits no values and completes immediately.

  • fromIterable

    variable

    Creates an Fx from an Iterable. Emits each value from the iterable in order and then completes.

  • never

    variable

    An Fx that waits forever without emitting a value.

  • null

    variable · Re-export

    An Fx that emits null exactly once and then completes.

  • succeed

    variable

    Creates an Fx that emits a single value and then completes.

  • succeedNull

    variable

    An Fx that emits null exactly once and then completes.

  • succeedUndefined

    variable

    An Fx that emits undefined exactly once and then completes.

  • succeedVoid

    variable

    An Fx that emits void exactly once and then completes.

  • suspend

    variable

    Defers creation of an Fx until each run begins.

  • sync

    variable

    Lazily evaluates a synchronous function once for each Fx run.

  • undefined

    variable · Re-export

    An Fx that emits undefined exactly once and then completes.

  • void

    variable · Re-export

    An Fx that emits void exactly once and then completes.

models

  • Fx.Any

    type-alias

    Matches any Fx regardless of its value, error, or service channels.

  • Fx.Class

    interface

    The constructible service class returned by Fx.Service.

  • Fx.Variance

    interface

    Describes how an Fx varies in its value, error, and service channels.