namespace / @typed/fx

Fx

Public namespace exposed as Fx from @typed/fx.

Package version
2.0.0-beta.13
Category
namespace

Import

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

This public exposure is a re-export. Its import path is supported; declaration documentation is shared with the other public exposures below.

Signatures

export * as Fx from "./Fx.js";

Members

  • Fx.Bounds

    Defines the bounds for slicing an Fx stream.

  • Fx.Emit

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

  • Fx.Error

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

  • Fx.FlatMapEffectLike

    Describes a dual flattening operator whose callback returns an Effect.

  • Fx.FlatMapLike

    Describes a dual flattening operator whose callback returns an Fx.

  • Fx.FromStreamOptions

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

  • Fx.Fx

    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.FxTypeId

    Runtime symbol carried by every Fx implementation.

  • Fx.KeyedOptions

    Configuration options for the keyed combinator.

  • Fx.Service

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

  • Fx.Services

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

  • Fx.Success

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

  • Fx.ThrottleOptions

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

  • Fx.ToStreamOptions

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

  • Fx.append

    Appends a value to the end of an Fx.

  • Fx.as

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

  • Fx.at

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

  • Fx.callback

    Creates an Fx from a callback-based source.

  • Fx.catch

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

  • Fx.catchAll

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

  • Fx.catchCause

    Recovers from any failure cause by running a fallback Fx.

  • Fx.catchCauseIf

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

  • Fx.catchIf

    Recovers a typed failure only when a predicate accepts it.

  • Fx.catchTag

    Recovers selected tagged typed failures by running a fallback Fx.

  • Fx.catchTags

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

  • Fx.catch_

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

  • Fx.causes

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

  • Fx.changesWithEffect

    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).

  • Fx.collectAll

    Collects all values emitted by an Fx into an array.

  • Fx.collectAllFork

    Forks the collection of all values from an Fx.

  • Fx.collectUpTo

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

  • Fx.collectUpToFork

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

  • Fx.compact

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

  • Fx.concat

    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.

  • Fx.concatMap

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

  • Fx.concatMapEffect

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

  • Fx.continueWith

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

  • Fx.debounce

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

  • Fx.delay

    Sleeps for duration before forwarding each successful source delivery.

  • Fx.delimit

    Wraps an Fx with a start and end value.

  • Fx.die

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

  • Fx.drain

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

  • Fx.drainLayer

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

  • Fx.dropAfter

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

  • Fx.dropUntil

    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.

  • Fx.dropUntilEffect

    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.

  • Fx.dropWhile

    Alias of skipWhile for Effect parity (dropWhile naming).

  • Fx.dropWhileEffect

    Alias of skipWhileEffect for Effect parity (dropWhileEffect naming).

  • Fx.during

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

  • Fx.empty

    An Fx that emits no values and completes immediately.

  • Fx.ensuring

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

  • Fx.exhaustLatestMap

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

  • Fx.exhaustLatestMapEffect

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

  • Fx.exhaustMap

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

  • Fx.exhaustMapEffect

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

  • Fx.exit

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

  • Fx.fail

    Creates an Fx that immediately fails with the specified error.

  • Fx.failCause

    Creates an Fx that immediately terminates with the specified Cause.

  • Fx.filter

    Filters elements of an Fx using a predicate function.

  • Fx.filterEffect

    Filters elements of an Fx using an effectful predicate function.

  • Fx.filterMap

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

  • Fx.filterMapEffect

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

  • Fx.filterMapLoop

    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.

  • Fx.filterMapLoopCause

    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.

  • Fx.filterMapLoopCauseEffect

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

  • Fx.filterMapLoopEffect

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

  • Fx.first

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

  • Fx.flatMap

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

  • Fx.flatMapConcurrently

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

  • Fx.flatMapConcurrentlyEffect

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

  • Fx.flatMapEffect

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

  • Fx.flip

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

  • Fx.fn

    Callable contracts implemented by fn.

  • Fx.fork

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

  • Fx.fromEffect

    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.

  • Fx.fromFailures

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

  • Fx.fromIterable

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

  • Fx.fromSchedule

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

  • Fx.fromStream

    Adapts an Effect Stream to a push-based Fx.

  • Fx.gen

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

  • Fx.genScoped

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

  • Fx.grouped

    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.

  • Fx.groupedWithin

    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.

  • Fx.if

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

  • Fx.interrupt

    Creates an Fx that immediately interrupts.

  • Fx.isFx

    Checks whether a value carries the FxTypeId protocol property.

  • Fx.keyed

    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.
  • Fx.loop

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

  • Fx.loopCause

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

  • Fx.loopCauseEffect

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

  • Fx.loopEffect

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

  • Fx.make

    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.

  • Fx.map

    Transforms the elements of an Fx using a provided function.

  • Fx.mapBoth

    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.

  • Fx.mapEffect

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

  • Fx.mapError

    Transforms typed failures while preserving defects and interruption.

  • Fx.merge

    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.

  • Fx.mergeAll

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

  • Fx.mergeLeft

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

  • Fx.mergeOrdered

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

  • Fx.mergeRight

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

  • Fx.never

    An Fx that waits forever without emitting a value.

  • Fx.null

    An Fx that emits null exactly once and then completes.

  • Fx.observe

    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.

  • Fx.observeLayer

    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.

  • Fx.onError

    Runs cleanup after the source reports a failure cause.

  • Fx.onExit

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

  • Fx.onInterrupt

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

  • Fx.pairwise

    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.

  • Fx.periodic

    Creates an Fx that emits a void value periodically.

  • Fx.prepend

    Prepends a value to the beginning of an Fx.

  • Fx.provide

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

  • Fx.provideContext

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

  • Fx.provideService

    Provides one existing service value to the entire Fx run.

  • Fx.provideServiceEffect

    Acquires one service with an Effect before running the Fx.

  • Fx.race

    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.

  • Fx.raceAll

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

  • Fx.repeat

    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.

  • Fx.result

    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).

  • Fx.retry

    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.

  • Fx.runFork

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

  • Fx.runPromise

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

  • Fx.runPromiseExit

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

  • Fx.sample

    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.

  • Fx.scan

    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), …

  • Fx.scanEffect

    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.

  • Fx.since

    Drops events until signal emits, then forwards the rest.

  • Fx.skip

    Skips the first n elements of an Fx.

  • Fx.skipEffect

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

  • Fx.skipRepeats

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

  • Fx.skipRepeatsWith

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

  • Fx.skipWhile

    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.

  • Fx.skipWhileEffect

    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.

  • Fx.slice

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

  • Fx.sliceEffect

    Slices an Fx with bounds produced by an Effect.

  • Fx.struct

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

  • Fx.succeed

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

  • Fx.succeedNull

    An Fx that emits null exactly once and then completes.

  • Fx.succeedUndefined

    An Fx that emits undefined exactly once and then completes.

  • Fx.succeedVoid

    An Fx that emits void exactly once and then completes.

  • Fx.suspend

    Defers creation of an Fx until each run begins.

  • Fx.switchMap

    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.

  • Fx.switchMapEffect

    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.

  • Fx.sync

    Lazily evaluates a synchronous function once for each Fx run.

  • Fx.take

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

  • Fx.takeEffect

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

  • Fx.takeUntil

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

  • Fx.takeUntilEffect

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

  • Fx.takeWhile

    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.

  • Fx.takeWhileEffect

    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.

  • Fx.tap

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

  • Fx.throttle

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

  • Fx.timeout

    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.

  • Fx.timeoutTo

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

  • Fx.toStream

    Adapts a push-based Fx to an Effect Stream.

  • Fx.tuple

    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.

  • Fx.undefined

    An Fx that emits undefined exactly once and then completes.

  • Fx.until

    Forwards events until signal emits, then interrupts events.

  • Fx.unwrap

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

  • Fx.unwrapScoped

    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.

  • Fx.void

    An Fx that emits void exactly once and then completes.

  • Fx.when

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

  • Fx.withLatestFrom

    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.

  • Fx.withLatestFromWith

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

  • Fx.withSpan

    Traces the whole subscription and each success or failure delivery.

  • Fx.zip

    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.

  • Fx.zipLatest

    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.

  • Fx.zipLatestWith

    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.

  • Fx.zipLeft

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

  • Fx.zipRight

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

  • Fx.zipWith

    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.

Source