Module / v2.0.0-beta.7

@typed/fx

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

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

Accumulation

  • RefSubject.scan

    variable · Re-export

    Stateful scan over a RefSubject / Computed, producing a Computed of the accumulated state.

    Fx subscriptions follow Fx.scan semantics (emit initial, then fold each source value). Effect sampling accumulates across source versions via a private state ref (do not mix heavy subscribe + sample on the same scan if you need a single shared accumulator).

  • RefSubject.scanEffect

    variable · Re-export

    Effectful stateful scan over a RefSubject / Computed, producing a Computed of the accumulated state.

Arithmetic queries

  • RefBigDecimal.abs

    variable · Re-export

    Get the absolute value of the current state of a RefBigDecimal.

  • RefBigDecimal.add

    variable · Re-export

    Add a BigDecimal to the current state of a RefBigDecimal.

  • RefBigDecimal.ceil

    variable · Re-export

    Calculate the ceiling of the current state of a RefBigDecimal.

  • RefBigDecimal.divide

    variable · Re-export

    Divide the current state of a RefBigDecimal by a BigDecimal.

  • RefBigDecimal.equals

    variable · Re-export

    Check if the current state of a RefBigDecimal equals a BigDecimal.

  • RefBigDecimal.floor

    variable · Re-export

    Calculate the floor of the current state of a RefBigDecimal.

  • RefBigDecimal.multiply

    variable · Re-export

    Multiply the current state of a RefBigDecimal by a BigDecimal.

  • RefBigDecimal.negate

    variable · Re-export

    Negate the current state of a RefBigDecimal.

  • RefBigDecimal.round

    variable · Re-export

    Round the current state of a RefBigDecimal.

  • RefBigDecimal.subtract

    variable · Re-export

    Subtract a BigDecimal from the current state of a RefBigDecimal.

  • RefBigDecimal.truncate

    variable · Re-export

    Truncate the current state of a RefBigDecimal.

  • RefBigInt.abs

    variable · Re-export

    Get the absolute value of the current state of a RefBigInt.

  • RefBigInt.add

    variable · Re-export

    Add a BigInt to the current state of a RefBigInt.

  • RefBigInt.divide

    variable · Re-export

    Divide the current state of a RefBigInt by a BigInt.

  • RefBigInt.equals

    variable · Re-export

    Check if the current state of a RefBigInt equals a BigInt.

  • RefBigInt.mod

    variable · Re-export

    Get the remainder of dividing the current state of a RefBigInt by a BigInt.

  • RefBigInt.multiply

    variable · Re-export

    Multiply the current state of a RefBigInt by a BigInt.

  • RefBigInt.negate

    variable · Re-export

    Negate the current state of a RefBigInt.

  • RefBigInt.subtract

    variable · Re-export

    Subtract a BigInt from the current state of a RefBigInt.

  • RefDuration.add

    variable · Re-export

    Add a Duration to the current state of a RefDuration.

  • RefDuration.divide

    variable · Re-export

    Divide the current state of a RefDuration by a number.

  • RefDuration.multiply

    variable · Re-export

    Multiply the current state of a RefDuration by a number.

  • RefDuration.subtract

    variable · Re-export

    Subtract a Duration from the current state of a RefDuration.

Bidirectional contracts

  • Push.Push

    interface · Re-export

    A bidirectional value that is both a Sink<A, E, R> and an Fx<B, E2, R2>.

    Calling onSuccess or onFailure sends exactly one input notification to the wrapped Sink. The returned Effect is the acknowledgment: a producer that runs and awaits it waits for the consumer callback to finish. Push adds no queue, buffering, replay, or demand protocol of its own.

    Running the Fx side preserves that Fx’s cardinality and ordering. Input values are not automatically forwarded to the output; any relationship between the two sides belongs to the supplied Sink and Fx (for example, a shared Subject).

Callback protocol

  • Fx.Emit

    type-alias · Re-export

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

Callback sources

  • Fx.callback

    variable · Re-export

    Creates an Fx from a callback-based source.

  • Fx.make

    variable · Re-export

    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.

Channel transformations

  • Versioned.map

    variable · Re-export

    Transform a Versioned’s output value as both an Fx and Effect.

  • Versioned.mapEffect

    variable · Re-export

    Transform a Versioned’s output value as both an Fx and Effect using an Effect.

  • Versioned.transform

    function · Re-export

    Transforms a Versioned value into another Versioned value.

Collecting values

  • Fx.collectAll

    variable · Re-export

    Collects all values emitted by an Fx into an array.

  • Fx.collectAllFork

    variable · Re-export

    Forks the collection of all values from an Fx.

  • Fx.collectUpTo

    variable · Re-export

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

  • Fx.collectUpToFork

    variable · Re-export

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

  • Fx.first

    function · Re-export

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

  • Sink.collect

    function · Re-export

    Collects all values into an array. Pass a Ref<ReadonlyArray<A>> (e.g. Ref.make([])); after running, read the result with Ref.get(ref).

  • Sink.head

    function · Re-export

    Keeps only the first value. Pass a Ref<Option.Option<A>> (e.g. Ref.make(Option.none())); after running, read the result with Ref.get(ref).

  • Sink.last

    function · Re-export

    Keeps only the last value. Pass a Ref<Option.Option<A>> (e.g. Ref.make(Option.none())); after running, read the result with Ref.get(ref).

  • Sink.reduce

    function · Re-export

    Reduces values into a single result using a pure function. Pass a Ref<B> (e.g. from Ref.make(initial)); after running, read the result with Ref.get(ref).

  • Sink.reduceEffect

    function · Re-export

    Reduces values into a single result using an effectful function. Pass a Ref<B>; after running, read the result with Ref.get(ref). If the reducer effect fails, the ref is left unchanged (Sink onSuccess is typed as never failing).

Combining sources

  • Fx.append

    variable · Re-export

    Appends a value to the end of an Fx.

  • Fx.concat

    variable · Re-export

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

    variable · Re-export

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

  • Fx.delimit

    variable · Re-export

    Wraps an Fx with a start and end value.

  • Fx.merge

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.mergeLeft

    variable · Re-export

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

  • Fx.mergeOrdered

    function · Re-export

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

  • Fx.mergeRight

    variable · Re-export

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

  • Fx.prepend

    variable · Re-export

    Prepends a value to the beginning of an Fx.

  • Fx.struct

    function · Re-export

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

    function · Re-export

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

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.zip

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.zipRight

    variable · Re-export

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

  • Fx.zipWith

    variable · Re-export

    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 output work

  • Push.exhaustLatestMap

    variable · Re-export

    Runs one inner Fx at a time and retains only the latest value received while busy.

    A value while idle starts immediately. While its inner runs, newer outer values replace a single pending slot. f(value) is evaluated and an inner Fx is constructed before every replacement; a superseded pending Fx is never run. After completion, only the latest pending Fx starts. Accepted inners preserve their own order and all values. The input Sink is unchanged.

  • Push.exhaustLatestMapEffect

    variable · Re-export

    Runs one mapped Effect at a time and retains only the latest value received while busy.

    Each accepted Effect can emit one result. While it runs, one pending Effect is repeatedly overwritten; after completion only the latest pending Effect starts. The mapping callback still runs and constructs an Effect for every value before replacement; superseded Effects are not run. The input side is unchanged.

  • Push.exhaustMap

    variable · Re-export

    Runs at most one inner Fx and ignores outer values while it is active.

    The first value seen while idle starts an inner; every value arriving before that inner completes is dropped. f(value) is still evaluated and its inner Fx is constructed before the busy check; dropping means the returned Fx is not run. Accepted inners preserve their own order and all values. Input is unchanged.

  • Push.exhaustMapEffect

    variable · Re-export

    Runs at most one mapped Effect and ignores values while it is active.

    The first value while idle starts one Effect and can emit one result; all values received before it completes are dropped. The mapping callback is still invoked and constructs an Effect for every value before the busy check; dropped Effects are not run. The input side is unchanged.

  • Push.flatMap

    variable · Re-export

    Transforms each output value into an inner Fx and merges all inners concurrently.

    Every outer value starts one inner. All inner values are emitted; order within each inner is preserved, but values from different inners may interleave. The outer stream waits for all inners before normal completion. The input Sink is unchanged.

  • Push.flatMapEffect

    variable · Re-export

    Transforms every output value into an Effect and merges their results concurrently.

    One Effect starts per outer value and can emit one result. All successful results are emitted, but concurrent completion order may differ from input order. The input side is unchanged.

  • Push.switchMap

    variable · Re-export

    Transforms each output value into an inner Fx, observing only the latest one.

    A new outer value interrupts the previous inner fiber before starting the next. Output cardinality is the cardinality of the successive active inners; values from an interrupted inner stop. Outer order determines replacement order, while each active inner preserves its own order. The input Sink is unchanged.

  • Push.switchMapEffect

    variable · Re-export

    Transforms each output value into an Effect, keeping only the latest Effect.

    Each new outer value interrupts the previous Effect before starting its own. Every Effect can emit at most one value; interrupted Effects emit none. The input side is unchanged.

Concurrent work

  • Fx.concatMap

    variable · Re-export

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

  • Fx.concatMapEffect

    variable · Re-export

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

  • Fx.exhaustLatestMap

    variable · Re-export

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

  • Fx.exhaustLatestMapEffect

    variable · Re-export

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

  • Fx.exhaustMap

    variable · Re-export

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

  • Fx.exhaustMapEffect

    variable · Re-export

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

  • Fx.flatMap

    variable · Re-export

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

  • Fx.flatMapConcurrently

    variable · Re-export

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

  • Fx.flatMapConcurrentlyEffect

    variable · Re-export

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

  • Fx.flatMapEffect

    variable · Re-export

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

  • Fx.race

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.switchMap

    variable · Re-export

    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

    variable · Re-export

    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

  • Fx.if

    variable · Re-export

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

  • Fx.when

    variable · Re-export

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

Construction options

Constructors

  • RefArray.make

    function · Re-export

    Creates a new RefArray from an array, Effect, or Fx.

  • RefBigDecimal.make

    function · Re-export

    Creates a new RefBigDecimal from a BigDecimal, Effect, or Fx.

  • RefBigInt.make

    function · Re-export

    Creates a new RefBigInt from a BigInt, Effect, or Fx.

  • RefBoolean.make

    function · Re-export

    Creates a new RefBoolean from a boolean, Effect, or Fx.

  • RefCause.make

    function · Re-export

    Creates a new RefCause from a Cause, Effect, or Fx.

  • RefChunk.make

    function · Re-export

    Creates a new RefChunk from a Chunk, Effect, or Fx.

  • RefDateTime.make

    function · Re-export

    Creates a new RefDateTime from a DateTime, Effect, or Fx.

  • RefDuration.make

    function · Re-export

    Creates a new RefDuration from a Duration, Effect, or Fx.

  • RefGraph.directed

    function · Re-export

    Creates a new empty directed RefGraph.

  • RefGraph.make

    function · Re-export

    Creates a new RefGraph from a Graph, Effect, or Fx.

  • RefGraph.undirected

    function · Re-export

    Creates a new empty undirected RefGraph.

  • RefHashMap.make

    function · Re-export

    Creates a new RefHashMap from a HashMap, Effect, or Fx.

  • RefHashRing.empty

    function · Re-export

    Creates a new empty RefHashRing.

  • RefHashRing.make

    function · Re-export

    Creates a new RefHashRing from a HashRing, Effect, or Fx.

  • RefHashSet.make

    function · Re-export

    Creates a new RefHashSet from a HashSet, Effect, or Fx.

  • RefIterable.make

    function · Re-export

    Creates a new RefIterable from an Iterable, Effect, or Fx.

  • RefOption.make

    function · Re-export

    Creates a new RefOption from an Option, Effect, or Fx.

  • RefRecord.make

    function · Re-export

    Creates a new RefRecord from a Record, Effect, or Fx.

  • RefResult.make

    function · Re-export

    Creates a new RefResult from a Result, Effect, or Fx.

  • RefString.make

    function · Re-export

    Creates a new RefString from a string, Effect, or Fx.

  • RefStruct.make

    function · Re-export

    Creates a new RefStruct from a struct, Effect, or Fx.

  • RefSubject.make

    function · Re-export

    Creates a new RefSubject from a value, Effect, or Fx.

  • RefTrie.make

    function · Re-export

    Creates a new RefTrie from a Trie, Effect, or Fx.

  • RefTuple.make

    function · Re-export

    Creates a new RefTuple from a tuple, Effect, or Fx.

  • Versioned.make

    function · Re-export

    Creates a Versioned value from its components.

  • Versioned.of

    function · Re-export

    Creates a Versioned value from a constant.

Consumer contracts

  • Sink.Sink

    interface · Re-export

    Consumes pushed successes and failures through effectful callbacks.

Date arithmetic

Date formatting

Derived queries

Effect interop

  • Fx.fromEffect

    variable · Re-export

    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

  • Fx.catch

    variable · Re-export

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

  • Fx.catchAll

    variable · Re-export

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

  • Fx.catchCause

    variable · Re-export

    Recovers from any failure cause by running a fallback Fx.

  • Fx.catchCauseIf

    variable · Re-export

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

  • Fx.catchIf

    variable · Re-export

    Recovers a typed failure only when a predicate accepts it.

  • Fx.catchTag

    variable · Re-export

    Recovers selected tagged typed failures by running a fallback Fx.

  • Fx.catchTags

    variable · Re-export

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

  • Fx.catch_

    variable · Re-export

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

  • Fx.causes

    variable · Re-export

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

  • Fx.exit

    variable · Re-export

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

  • Fx.flip

    variable · Re-export

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

  • Fx.mapBoth

    variable · Re-export

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

    variable · Re-export

    Transforms typed failures while preserving defects and interruption.

  • Fx.result

    variable · Re-export

    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

    variable · Re-export

    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 handling

  • Sink.exit

    variable · Re-export

    Materializes both sink channels as successful Effect Exit values.

  • Sink.flip

    variable · Re-export

    Exchanges a sink’s typed success and failure channels.

  • Sink.mapError

    function · Re-export

    Maps the error channel of a sink using the provided function. Failures are mapped via Cause.map; defects and interrupts are preserved.

  • Sink.skipInterrupt

    variable · Re-export

    Suppresses failure causes made entirely of interruption reasons.

Failure sources

  • Fx.die

    variable · Re-export

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

  • Fx.fail

    variable · Re-export

    Creates an Fx that immediately fails with the specified error.

  • Fx.failCause

    variable · Re-export

    Creates an Fx that immediately terminates with the specified Cause.

  • Fx.fromFailures

    variable · Re-export

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

  • Fx.interrupt

    variable · Re-export

    Creates an Fx that immediately interrupts.

Generator composition

  • Fx.fn

    namespace · Re-export

    Callable contracts implemented by fn.

  • Fx.gen

    variable · Re-export

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

  • Fx.genScoped

    variable · Re-export

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

  • Fx.unwrap

    variable · Re-export

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

  • Fx.unwrapScoped

    variable · Re-export

    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.

Hydration construction

  • RefSubject.hydrate

    function · Re-export

    Creates a named hydrated RefSubject using a string-encoded Schema codec.

  • RefSubject.hydrateAll

    function · Re-export

    Combines hydrated refs into one serialization and restoration boundary.

Hydration protocol

Hydration types

Keyed work

  • Fx.keyed

    variable · Re-export

    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.

Observation policy

  • RefSubject.CurrentComputedBehavior

    variable · Re-export

    Selects whether a computed Fx observes only its current value or also follows later pushes.

  • RefSubject.slice

    variable · Re-export

    Limits which pushed versions a RefSubject view observes without changing its state.

Observing failures

  • Fx.onError

    variable · Re-export

    Runs cleanup after the source reports a failure cause.

  • Fx.withSpan

    variable · Re-export

    Traces the whole subscription and each success or failure delivery.

Operator options

  • Fx.Bounds

    interface · Re-export

    Defines the bounds for slicing an Fx stream.

  • Fx.FromStreamOptions

    type-alias · Re-export

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

  • Fx.KeyedOptions

    interface · Re-export

    Configuration options for the keyed combinator.

  • Fx.ThrottleOptions

    type-alias · Re-export

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

  • Fx.ToStreamOptions

    type-alias · Re-export

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

  • Sink.Bounds

    interface · Re-export

    Zero-based skip count and maximum take count used by slice.

Optional channel transformations

  • Versioned.filterMap

    variable · Re-export

    Filter-maps a Versioned’s output as both an Fx and Effect; the Effect value becomes Option (Some when the predicate holds, None otherwise).

  • Versioned.filterMapEffect

    variable · Re-export

    Filter-maps a Versioned’s output as both an Fx and Effect using an Effect; the Effect value becomes Option.

Optional queries

  • RefArray.getIndex

    variable · Re-export

    Get a value contained a particular index of a RefArray.

  • RefArray.head

    variable · Re-export

    Gets the first element of a RefArray as a Filtered.

  • RefArray.last

    variable · Re-export

    Gets the last element of a RefArray as a Filtered.

  • RefChunk.findFirst

    variable · Re-export

    Find the first value satisfying a predicate.

  • RefChunk.findLast

    variable · Re-export

    Find the last value satisfying a predicate.

  • RefChunk.getIndex

    variable · Re-export

    Get a value at a particular index of a RefChunk.

  • RefChunk.head

    variable · Re-export

    Get the first element of a RefChunk as a Filtered.

  • RefChunk.last

    variable · Re-export

    Get the last element of a RefChunk as a Filtered.

  • RefGraph.findEdge

    variable · Re-export

    Find an edge matching a predicate.

  • RefGraph.findNode

    variable · Re-export

    Find a node matching a predicate.

  • RefGraph.getEdge

    variable · Re-export

    Get an edge’s data.

  • RefGraph.getNode

    variable · Re-export

    Get a node’s data.

  • RefHashMap.findFirst

    variable · Re-export

    Find the first entry satisfying a predicate.

  • RefHashMap.get

    variable · Re-export

    Get the value at a key as a Filtered.

  • RefHashRing.get

    variable · Re-export

    Get the node which should handle a given input string as a Filtered. Fails if the ring is empty.

  • RefIterable.findFirst

    variable · Re-export

    Find the first value satisfying a predicate.

  • RefIterable.findLast

    variable · Re-export

    Find the last value satisfying a predicate.

  • RefIterable.head

    variable · Re-export

    Get the first element of a RefIterable as a Filtered.

  • RefOption.getValue

    variable · Re-export

    Get the value from the Option as a Filtered (fails if None).

  • RefRecord.findFirst

    variable · Re-export

    Find the first entry satisfying a predicate.

  • RefRecord.get

    variable · Re-export

    Get the value at a key as a Filtered.

  • RefRecord.pop

    variable · Re-export

    Pop a value at a key as a Filtered.

  • RefResult.getFailure

    variable · Re-export

    Get the failure value from the Result as a Filtered (fails if Success).

  • RefResult.getSuccess

    variable · Re-export

    Get the success value from the Result as a Filtered (fails if Failure).

  • RefSubject.compact

    variable · Re-export

    Converts a Computed or Filtered of Option<A> into a Filtered<A>, filtering out None values.

  • RefSubject.filterMap

    variable · Re-export

    Filters and transforms a RefSubject, Computed, or Filtered using a pure function that returns an Option.

  • RefSubject.filterMapEffect

    variable · Re-export

    Filters and transforms a RefSubject, Computed, or Filtered using an Effectful function that returns an Option.

  • RefSubject.getOrElse

    variable · Re-export

    Returns a Computed that yields the value inside the Option, or the fallback when None. Works with Computed<Option<A>> (e.g. from fromOption / fromNullable) and with Filtered<A>.

  • RefSubject.makeFiltered

    function · Re-export

    Builds a Filtered view from the three channels of a Versioned source.

  • RefTrie.get

    variable · Re-export

    Get the value at a key as a Filtered.

Output failures

  • Push.mapError

    variable · Re-export

    Transforms the output (Fx) error channel of a Push using the provided function.

    Failures (Cause) are mapped via Cause.map, so only the typed failure (Fail) is transformed; defects and interrupts are preserved unchanged.

    Mirrors Effect.mapError on the Fx side. Cardinality, value order, input callbacks, and output service requirements are unchanged.

Providing services

  • Fx.Fx.Service

    interface · Re-export

    An Fx whose implementation is obtained from an Effect service.

  • Fx.Service

    function · Re-export

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

  • Fx.provide

    variable · Re-export

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

  • Fx.provideContext

    variable · Re-export

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

  • Fx.provideService

    variable · Re-export

    Provides one existing service value to the entire Fx run.

  • Fx.provideServiceEffect

    variable · Re-export

    Acquires one service with an Effect before running the Fx.

Publication contracts

  • Subject.Subject

    interface · Re-export

    A multicast boundary that is both an Fx of its publications and a Sink that accepts them. Successes and failures are pushed to every subscriber present when that publication begins.

Push construction

  • Push.make

    variable · Re-export

    Couples a Sink input with an independent Fx output.

    The result forwards each input callback directly to sink and delegates every output subscription to fx. It does not connect the two values, change output cardinality or ordering, buffer inputs, or start either side eagerly.

Push services

  • Push.Push.Class

    interface · Re-export

    Constructable static type produced by Push.Service.

  • Push.Push.Service

    interface · Re-export

    The static and Effect service surface returned by Push.Service.

    Service lookup supplies the same bidirectional value to run, onSuccess, and onFailure. The Self service appears in both required-service channels; the installed Push itself has those requirements captured by its Layer.

  • Push.Service

    function · Re-export

    Defines a named Effect service whose value is a Push.

    The returned class exposes onSuccess, onFailure, and run as Effects that first resolve the service from Context. make captures the Sink construction context and combines it with each output subscriber’s context; it does not start the Fx or send an input while building the Layer.

Read-only state

  • RefSubject.Computed

    interface · Re-export

    A Computed is a read-only view of a value that can change over time. It is an Fx that emits the current value and subsequent updates. It is also an Effect that samples the current value.

  • RefSubject.Filtered

    interface · Re-export

    A Filtered is a Computed that may not always have a value. It is essentially a Computed<Option<A>> with helper methods.

Resource lifetime

  • Fx.ensuring

    variable · Re-export

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

  • Fx.onExit

    variable · Re-export

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

  • Fx.onInterrupt

    variable · Re-export

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

Running effects

  • Fx.drain

    variable · Re-export

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

  • Fx.drainLayer

    variable · Re-export

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

  • Fx.fork

    variable · Re-export

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

  • Fx.observe

    variable · Re-export

    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

    variable · Re-export

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

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.runPromiseExit

    variable · Re-export

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

Runtime inspection

  • Fx.FxTypeId

    variable · Re-export

    Runtime symbol carried by every Fx implementation.

  • Fx.isFx

    function · Re-export

    Checks whether a value carries the FxTypeId protocol property.

Selecting inputs

  • Push.filterInput

    variable · Re-export

    Keeps successful inputs that satisfy f and discards the rest.

    Calling onSuccess runs the predicate immediately. A match constructs one Sink callback Effect; a non-match immediately returns an empty acknowledgment. Predicate allocation and throws therefore occur before the returned Effect is run. The producer controls call order and concurrency; output is preserved.

  • Push.filterInputEffect

    variable · Re-export

    Effectfully decides whether each successful input reaches the Sink.

    Calling onSuccess(value) invokes f(value) immediately to construct the predicate Effect; allocation and throws happen before an acknowledgment is returned. Running the acknowledgment later executes that Effect in the caller’s fiber. true forwards one value, false none, and failure sends its Cause to the Sink failure callback. Calls are not serialized; the producer controls order and concurrency. Output behavior is unchanged.

  • Push.filterMapInput

    variable · Re-export

    Transforms an input and forwards it only when f returns Some.

    Calling onSuccess evaluates f immediately. Some(a) constructs one Sink callback Effect; None immediately returns an empty acknowledgment. Mapping allocation and throws therefore happen before the returned Effect runs. Calls are not serialized, so ordering follows the producer. Output is unchanged.

  • Push.filterMapInputEffect

    variable · Re-export

    Effectfully transforms an input and forwards only a resulting Some value.

    Calling onSuccess(value) invokes f(value) immediately to construct an Effect; allocation and throws happen before an acknowledgment is returned. Running that acknowledgment later executes the constructed Effect in the caller’s fiber. Some(a) produces one Sink callback, None none, and failure sends its full Cause to the Sink failure callback. Concurrent calls are not serialized.

  • Sink.compact

    function · Re-export

    Forwards Some values and discards None values.

  • Sink.filter

    function · Re-export

    Filters values before they reach the sink using a predicate function.

  • Sink.filterEffect

    variable · Re-export

    Runs an effectful predicate and forwards inputs for which it succeeds with true.

  • Sink.filterMap

    function · Re-export

    Filters and transforms values before they reach the sink using a function that returns an Option.

  • Sink.filterMapEffect

    variable · Re-export

    Runs an Effect for each input and forwards its optional successful value.

Selecting outputs

  • Push.filter

    variable · Re-export

    Keeps Fx output values that satisfy f.

    The upstream Sink invokes the predicate synchronously for each output value, before the returned downstream callback Effect runs. Predicate allocation and throws therefore occur at upstream callback invocation. Matches preserve their relative order; non-matches produce no output. No buffer is added and every input callback is unchanged.

  • Push.filterEffect

    variable · Re-export

    Effectfully decides which Fx output values are emitted.

    Each upstream value runs one predicate. true emits that value, false emits none, and predicate failures join the output error channel. For a sequential source the predicate is acknowledged before the next delivery, preserving order; the input side is unchanged.

  • Push.filterMap

    variable · Re-export

    Maps each Fx output and emits only resulting Some values.

    The upstream Sink evaluates f synchronously for each value, before its returned downstream callback Effect runs. Mapping allocation and throws occur at that callback invocation. Some(c) emits exactly one c; None emits nothing. Relative order, errors, services, and input behavior are preserved.

  • Push.filterMapEffect

    variable · Re-export

    Effectfully maps each Fx output and emits only resulting Some values.

    One mapper Effect runs per upstream value. Some(c) emits once, None emits nothing, and failures join the output error channel. Sequential sources retain order; the input Sink remains unchanged.

Selecting values

  • Fx.changesWithEffect

    variable · Re-export

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

    variable · Re-export

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

  • Fx.dropAfter

    variable · Re-export

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

  • Fx.dropUntil

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

    Alias of skipWhile for Effect parity (dropWhile naming).

  • Fx.dropWhileEffect

    variable · Re-export

    Alias of skipWhileEffect for Effect parity (dropWhileEffect naming).

  • Fx.filter

    variable · Re-export

    Filters elements of an Fx using a predicate function.

  • Fx.filterEffect

    variable · Re-export

    Filters elements of an Fx using an effectful predicate function.

  • Fx.filterMap

    variable · Re-export

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

  • Fx.filterMapEffect

    variable · Re-export

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

  • Fx.skip

    variable · Re-export

    Skips the first n elements of an Fx.

  • Fx.skipEffect

    variable · Re-export

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

  • Fx.skipRepeats

    variable · Re-export

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

  • Fx.skipRepeatsWith

    variable · Re-export

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

  • Fx.skipWhile

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.sliceEffect

    variable · Re-export

    Slices an Fx with bounds produced by an Effect.

  • Fx.take

    variable · Re-export

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

  • Fx.takeEffect

    variable · Re-export

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

  • Fx.takeUntil

    variable · Re-export

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

  • Fx.takeUntilEffect

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

    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.

Services

Shared observation

Sharing sources

  • Subject.Share

    class · Re-export

    The concrete lazy Fx returned by share.

  • Subject.hold

    function · Re-export

    Shares an Fx and immediately replays its latest success or failure to each new subscriber.

  • Subject.multicast

    function · Re-export

    Multicasts an Fx without replaying values that arrived before a subscriber joined.

  • Subject.replay

    variable · Re-export

    Shares an Fx and replays up to the last capacity successes or failures to new subscribers.

  • Subject.share

    function · Re-export

    Shares one active execution of an Fx among subscribers through the supplied Subject.

Sink construction

  • Sink.make

    function · Re-export

    Creates a Sink from success and failure callbacks.

Sink services

  • Sink.Service

    function · Re-export

    Defines a class-shaped Effect Context service for a Sink.

  • Sink.Sink.Class

    interface · Re-export

    Constructor-shaped Context service returned by Sink.Service.

  • Sink.Sink.Service

    interface · Re-export

    Describes a Sink available through an Effect Context service.

Source adapters

State composition

  • RefSubject.struct

    function · Re-export

    Combines multiple RefSubject, Computed, or Filtered instances into a single struct.

  • RefSubject.tuple

    function · Re-export

    Combines multiple RefSubject, Computed, or Filtered instances into a single tuple.

  • Versioned.struct

    function · Re-export

    Combines multiple Versioned values into a single struct.

  • Versioned.tuple

    function · Re-export

    Combines multiple Versioned values into a single tuple.

State models

  • RefArray.RefArray

    interface · Re-export

    A RefArray is a RefSubject that is specialized over an array of values.

  • RefBigDecimal.RefBigDecimal

    interface · Re-export

    A RefBigDecimal is a RefSubject specialized over a BigDecimal value.

  • RefBigInt.RefBigInt

    interface · Re-export

    A RefBigInt is a RefSubject specialized over a BigInt value.

  • RefBoolean.RefBoolean

    interface · Re-export

    A RefBoolean is a RefSubject specialized over a boolean value.

  • RefCause.RefCause

    interface · Re-export

    A RefCause is a RefSubject specialized over a Cause value.

  • RefChunk.RefChunk

    interface · Re-export

    A RefChunk is a RefSubject specialized over a Chunk of values.

  • RefDateTime.RefDateTime

    interface · Re-export

    A RefDateTime is a RefSubject specialized over a DateTime value.

  • RefDuration.RefDuration

    interface · Re-export

    A RefDuration is a RefSubject specialized over a Duration value.

  • RefGraph.RefGraph

    interface · Re-export

    A RefGraph is a RefSubject specialized over a Graph.

  • RefHashMap.RefHashMap

    interface · Re-export

    A RefHashMap is a RefSubject specialized over a HashMap.

  • RefHashRing.RefHashRing

    interface · Re-export

    A RefHashRing is a RefSubject specialized over a HashRing.

  • RefHashSet.RefHashSet

    interface · Re-export

    A RefHashSet is a RefSubject specialized over a HashSet.

  • RefIterable.RefIterable

    interface · Re-export

    A RefIterable is a RefSubject specialized over an Iterable of values.

  • RefOption.RefOption

    interface · Re-export

    A RefOption is a RefSubject specialized over an Option value.

  • RefRecord.RefRecord

    interface · Re-export

    A RefRecord is a RefSubject specialized over a Record.

  • RefResult.RefResult

    interface · Re-export

    A RefResult is a RefSubject specialized over a Result value.

  • RefString.RefString

    interface · Re-export

    A RefString is a RefSubject specialized over a string value.

  • RefStruct.RefStruct

    interface · Re-export

    A RefStruct is a RefSubject specialized over a struct value.

  • RefTrie.RefTrie

    interface · Re-export

    A RefTrie is a RefSubject specialized over a Trie.

  • RefTuple.RefTuple

    interface · Re-export

    A RefTuple is a RefSubject specialized over a tuple value.

State predicates

State protocol

  • Versioned.Versioned

    interface · Re-export

    A Versioned value is a value that changes over time, and each change is associated with a version number. It combines the capabilities of an Fx (to observe changes) and an Effect (to get the current value).

State updates

Stateful delivery

  • Sink.Sink.WithState

    interface · Re-export

    An early-exit sink with a mutable Effect Ref for consumer-local state.

  • Sink.Sink.WithStateSemaphore

    interface · Re-export

    An early-exit sink with serialized effectful access to consumer-local state.

  • Sink.filterMapLoop

    variable · Re-export

    Threads pure state through successes and optionally forwards a derived value.

  • Sink.filterMapLoopEffect

    variable · Re-export

    Threads state through an effectful success transition and optionally forwards its value.

  • Sink.loop

    variable · Re-export

    Threads pure state through successful inputs and forwards one derived value per input.

  • Sink.loopEffect

    variable · Re-export

    Threads state through an effectful success transformation and forwards one value on success.

  • Sink.withState

    function · Re-export

    Runs a callback with an early-exit sink and a private Effect Ref initialized to state.

  • Sink.withStateSemaphore

    function · Re-export

    Runs a callback with private state whose effectful reads and writes are serialized.

Stateful failure handling

  • Sink.filterMapLoopCause

    variable · Re-export

    Threads pure state through failures and optionally forwards a transformed cause.

  • Sink.filterMapLoopCauseEffect

    function · Re-export

    Threads state through an effectful failure transition and optionally forwards a cause.

  • Sink.loopCause

    variable · Re-export

    Threads pure state through failure causes and forwards one transformed cause per failure.

  • Sink.loopCauseEffect

    variable · Re-export

    Threads state through an effectful failure transformation.

Stateful outputs

  • Push.mapAccum

    variable · Re-export

    Maps over the output (Fx) side of a Push with an accumulator: for each emitted value b, applies f(state, b) to get [nextState, emitted] and emits the second element. The first element is the initial state; subsequent states are updated by each step. It emits exactly one C per upstream value, in order. Accumulator state is private to each output subscription; the input Sink is unchanged.

  • Push.mapAccumEffect

    variable · Re-export

    Maps over the output (Fx) side of a Push with an effectful accumulator: for each emitted value b, runs f(state, b) to get [nextState, emitted] and emits the second element. The adapter does not serialize callbacks. Calling onSuccess invokes f immediately with the current seed and constructs its Effect; overlapping calls can therefore observe the same seed. Successful completion commits the returned seed and emits in completion order, so later completion may overwrite newer state. Reducer failure is sent to the output Sink, emits nothing, restores that call’s previous seed, and completes normally so a continuing producer can send later values. The input Sink is unchanged.

Stateful transforms

  • Fx.filterMapLoop

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

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

  • Fx.filterMapLoopEffect

    variable · Re-export

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

  • Fx.grouped

    variable · Re-export

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

    variable · Re-export

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

  • Fx.loopCause

    variable · Re-export

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

  • Fx.loopCauseEffect

    variable · Re-export

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

  • Fx.loopEffect

    variable · Re-export

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

  • Fx.pairwise

    variable · Re-export

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

    variable · Re-export

    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

    variable · Re-export

    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.

Stopping delivery

  • Sink.Sink.WithEarlyExit

    interface · Re-export

    A sink that can ask its producer callback to stop early.

  • Sink.dropAfter

    variable · Re-export

    Runs a producer callback until the first matching value has been forwarded.

  • Sink.slice

    variable · Re-export

    Runs a producer callback through a bounded view of a sink.

  • Sink.withEarlyExit

    function · Re-export

    Runs a producer-style callback with a sink that can complete the surrounding Effect early.

Stream interop

  • Fx.fromStream

    variable · Re-export

    Adapts an Effect Stream to a push-based Fx.

  • Fx.toStream

    variable · Re-export

    Adapts a push-based Fx to an Effect Stream.

Subject construction

  • Subject.make

    function · Re-export

    Acquires a Subject whose subscribers and replay state are released with the current Scope.

  • Subject.unsafeMake

    function · Re-export

    Immediately allocates a Subject with the requested replay capacity and manual ownership.

Subject services

Time and rate

  • Fx.at

    variable · Re-export

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

  • Fx.debounce

    variable · Re-export

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

  • Fx.delay

    variable · Re-export

    Sleeps for duration before forwarding each successful source delivery.

  • Fx.during

    function · Re-export

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

  • Fx.fromSchedule

    variable · Re-export

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

  • Fx.groupedWithin

    variable · Re-export

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

    variable · Re-export

    Creates an Fx that emits a void value periodically.

  • Fx.repeat

    variable · Re-export

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

    variable · Re-export

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

    variable · Re-export

    Drops events until signal emits, then forwards the rest.

  • Fx.throttle

    variable · Re-export

    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

    variable · Re-export

    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

    variable · Re-export

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

    variable · Re-export

    Forwards events until signal emits, then interrupts events.

Transactions

  • RefSubject.GetSetDelete

    interface · Re-export

    Interface for basic RefSubject operations: get, set, delete.

  • RefSubject.modify

    variable · Re-export

    Modifies a RefSubject using a pure function that returns both a result and a new value.

  • RefSubject.modifyEffect

    variable · Re-export

    Modifies a RefSubject using an Effectful function that returns both a result and a new value.

  • RefSubject.runUpdates

    variable · Re-export

    Runs an effect that can modify a RefSubject transactionally, with optional interrupt handling.

Transforming inputs

  • Push.mapInput

    variable · Re-export

    Synchronously transforms each successful input before sending it to the Sink.

    One input maps to exactly one downstream input. Calling onSuccess evaluates f immediately, before the returned downstream Effect is run; allocation and thrown exceptions therefore occur at callback invocation (and become defects only when that invocation itself occurs inside Effect evaluation). The Sink callback Effect remains the acknowledgment. Calls are not serialized; execution order and concurrency follow the producer. The Fx output is unchanged.

  • Push.mapInputEffect

    variable · Re-export

    Effectfully transforms each successful input before sending it to the Sink.

    Calling onSuccess(value) invokes f(value) immediately to construct an Effect. Allocation and thrown exceptions therefore occur before an acknowledgment Effect is returned. Running that returned Effect later executes the constructed Effect in the caller’s fiber. On success its single A reaches the Sink; typed failure, defect, or interruption sends its full Cause to the Sink failure callback. Calls are not serialized, so the producer controls order and concurrency. Output values and order are unchanged.

  • Sink.map

    function · Re-export

    Transforms values before they reach the sink using a pure function.

  • Sink.mapEffect

    variable · Re-export

    Runs an Effect for each input and forwards its successful value to the sink.

  • Sink.mapInput

    variable · Re-export

    Alias for map, named for the input-side direction of the transformation.

  • Sink.mapInputEffect

    variable · Re-export

    Alias for mapEffect, named for the input-side direction of the effectful transformation.

  • Sink.tapEffect

    variable · Re-export

    Runs an effectful observation before forwarding each successful input unchanged.

Transforming outputs

  • Push.map

    variable · Re-export

    Synchronously transforms every value emitted by the Fx output side.

    It emits exactly one C for every upstream B, preserving order and the complete input Sink. The mapping does not buffer or introduce concurrency.

  • Push.mapBoth

    variable · Re-export

    Transforms both the output (Fx) success and error channels of a Push using the provided options.

    Mirrors Effect.mapBoth on the Fx side: onSuccess maps every emitted value one-to-one and onFailure maps typed failures via Cause.map; defects and interrupts are preserved. Ordering, services, and the input side are unchanged.

  • Push.mapEffect

    variable · Re-export

    Effectfully transforms each Fx output value.

    The mapper runs once per upstream value and emits one result on success, in upstream order for a sequential source. Its typed failures join E2; its required services join R2. The input side is unchanged.

Transforming values

  • Fx.as

    variable · Re-export

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

  • Fx.map

    variable · Re-export

    Transforms the elements of an Fx using a provided function.

  • Fx.mapEffect

    variable · Re-export

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

  • Fx.tap

    variable · Re-export

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

Type contracts

  • Fx.Error

    type-alias · Re-export

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

  • Fx.FlatMapEffectLike

    type-alias · Re-export

    Describes a dual flattening operator whose callback returns an Effect.

  • Fx.FlatMapLike

    type-alias · Re-export

    Describes a dual flattening operator whose callback returns an Fx.

  • Fx.Fx

    interface · Re-export

    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.Fx.Error

    type-alias · Re-export

    Extracts the typed error from an Fx.

  • Fx.Fx.Services

    type-alias · Re-export

    Extracts the services required to run an Fx.

  • Fx.Fx.Success

    type-alias · Re-export

    Extracts the emitted value type from an Fx.

  • Fx.Services

    type-alias · Re-export

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

  • Fx.Success

    type-alias · Re-export

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

  • Fx.fn.Gen

    type-alias · Re-export

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

  • Fx.fn.NonGen

    type-alias · Re-export

    Contract for functions whose body returns an Fx directly.

  • Push.Push.Any

    type-alias · Re-export

    Matches any Push when its six channel types are intentionally unknown.

  • Sink.Error

    type-alias · Re-export

    Alias of Sink.Error for direct imports.

  • Sink.Services

    type-alias · Re-export

    Alias of Sink.Services for direct imports.

  • Sink.Sink.Any

    type-alias · Re-export

    Matches any Sink when its channels do not need to be preserved.

  • Sink.Sink.Error

    type-alias · Re-export

    Extracts the typed error carried by causes consumed by a Sink.

  • Sink.Sink.Services

    type-alias · Re-export

    Extracts the Effect services required by a Sink’s callbacks.

  • Sink.Sink.Success

    type-alias · Re-export

    Extracts the successful input type consumed by a Sink.

  • Sink.Success

    type-alias · Re-export

    Alias of Sink.Success for direct imports.

Type guards

Type identity

Type utilities

Unit conversions

  • RefDuration.days

    variable · Re-export

    Get the days value of the current state of a RefDuration.

  • RefDuration.hours

    variable · Re-export

    Get the hours value of the current state of a RefDuration.

  • RefDuration.millis

    variable · Re-export

    Get the milliseconds value of the current state of a RefDuration.

  • RefDuration.minutes

    variable · Re-export

    Get the minutes value of the current state of a RefDuration.

  • RefDuration.seconds

    variable · Re-export

    Get the seconds value of the current state of a RefDuration.

Value sources

  • Fx.empty

    variable · Re-export

    An Fx that emits no values and completes immediately.

  • Fx.fromIterable

    variable · Re-export

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

  • Fx.never

    variable · Re-export

    An Fx that waits forever without emitting a value.

  • Fx.null

    variable · Re-export

    An Fx that emits null exactly once and then completes.

  • Fx.succeed

    variable · Re-export

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

  • Fx.succeedNull

    variable · Re-export

    An Fx that emits null exactly once and then completes.

  • Fx.succeedUndefined

    variable · Re-export

    An Fx that emits undefined exactly once and then completes.

  • Fx.succeedVoid

    variable · Re-export

    An Fx that emits void exactly once and then completes.

  • Fx.suspend

    variable · Re-export

    Defers creation of an Fx until each run begins.

  • Fx.sync

    variable · Re-export

    Lazily evaluates a synchronous function once for each Fx run.

  • Fx.undefined

    variable · Re-export

    An Fx that emits undefined exactly once and then completes.

  • Fx.void

    variable · Re-export

    An Fx that emits void exactly once and then completes.

Writable state

  • RefSubject.RefSubject

    interface · Re-export

    A RefSubject is a mutable reference that can be observed as an Fx. It combines the capabilities of a Ref (get/set/update) with a Subject (subscribe).

Writable views

  • RefSubject.transform

    variable · Re-export

    Transforms a RefSubject invariantly using bidirectional mapping functions.

models

  • Fx.Fx.Any

    type-alias · Re-export

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

  • Fx.Fx.Class

    interface · Re-export

    The constructible service class returned by Fx.Service.

  • Fx.Fx.Variance

    interface · Re-export

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

namespace

  • Fx

    namespace · Re-export

    Public namespace exposed as Fx from @typed/fx.

  • Push

    namespace · Re-export

    Public namespace exposed as Push from @typed/fx.

  • RefArray

    namespace · Re-export

    Public namespace exposed as RefArray from @typed/fx.

  • RefBigDecimal

    namespace · Re-export

    Public namespace exposed as RefBigDecimal from @typed/fx.

  • RefBigInt

    namespace · Re-export

    Public namespace exposed as RefBigInt from @typed/fx.

  • RefBoolean

    namespace · Re-export

    Public namespace exposed as RefBoolean from @typed/fx.

  • RefCause

    namespace · Re-export

    Public namespace exposed as RefCause from @typed/fx.

  • RefChunk

    namespace · Re-export

    Public namespace exposed as RefChunk from @typed/fx.

  • RefDateTime

    namespace · Re-export

    Public namespace exposed as RefDateTime from @typed/fx.

  • RefDuration

    namespace · Re-export

    Public namespace exposed as RefDuration from @typed/fx.

  • RefGraph

    namespace · Re-export

    Public namespace exposed as RefGraph from @typed/fx.

  • RefHashMap

    namespace · Re-export

    Public namespace exposed as RefHashMap from @typed/fx.

  • RefHashRing

    namespace · Re-export

    Public namespace exposed as RefHashRing from @typed/fx.

  • RefHashSet

    namespace · Re-export

    Public namespace exposed as RefHashSet from @typed/fx.

  • RefIterable

    namespace · Re-export

    Public namespace exposed as RefIterable from @typed/fx.

  • RefOption

    namespace · Re-export

    Public namespace exposed as RefOption from @typed/fx.

  • RefRecord

    namespace · Re-export

    Public namespace exposed as RefRecord from @typed/fx.

  • RefResult

    namespace · Re-export

    Public namespace exposed as RefResult from @typed/fx.

  • RefString

    namespace · Re-export

    Public namespace exposed as RefString from @typed/fx.

  • RefStruct

    namespace · Re-export

    Public namespace exposed as RefStruct from @typed/fx.

  • RefSubject

    namespace · Re-export

    Public namespace exposed as RefSubject from @typed/fx.

  • RefTrie

    namespace · Re-export

    Public namespace exposed as RefTrie from @typed/fx.

  • RefTuple

    namespace · Re-export

    Public namespace exposed as RefTuple from @typed/fx.

  • Sink

    namespace · Re-export

    Public namespace exposed as Sink from @typed/fx.

  • Subject

    namespace · Re-export

    Public namespace exposed as Subject from @typed/fx.

  • Versioned

    namespace · Re-export

    Public namespace exposed as Versioned from @typed/fx.