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.BoundsDefines the bounds for slicing an Fx stream.
Fx.EmitOperations supplied to a callback producer for emitting values or ending its run.
Fx.ErrorAlias of
Fx.Errorfor extracting anFxerror channel.Fx.FlatMapEffectLikeDescribes a dual flattening operator whose callback returns an Effect.
Fx.FlatMapLikeDescribes a dual flattening operator whose callback returns an Fx.
Fx.FromStreamOptionsEffect Stream mapping options used while delivering elements to an
Fxsink.Fx.FxFxis 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
Effectwhich produces a single value,Fxcan 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.- Emits values of type
Fx.FxTypeIdRuntime symbol carried by every
Fximplementation.Fx.KeyedOptionsConfiguration options for the
keyedcombinator.Fx.ServiceDefines an Effect service whose value is also a directly runnable
Fx.Fx.ServicesAlias of
Fx.Servicesfor extracting anFxservice channel.Fx.SuccessAlias of
Fx.Successfor extracting anFxvalue channel.Fx.ThrottleOptionsOptions for {@link throttle}. A duration-only call is leading-edge (
{ leading: true, trailing: false }).Fx.ToStreamOptionsBuffering and callback options accepted while adapting an
Fxto an EffectStream.Fx.appendAppends a value to the end of an Fx.
Fx.asReplaces all emitted values from the Fx with the provided value
b.Fx.atCreates an Fx that emits a single value after a specified delay.
Fx.callbackCreates an Fx from a callback-based source.
Fx.catchRecovers from the first typed failure of an Fx by running a fallback Fx.
Fx.catchAllUses the Effect-style
catchAllname for {@link catch}.Fx.catchCauseRecovers from any failure cause by running a fallback Fx.
Fx.catchCauseIfRecovers a failure cause only when a predicate accepts the complete cause.
Fx.catchIfRecovers a typed failure only when a predicate accepts it.
Fx.catchTagRecovers selected tagged typed failures by running a fallback Fx.
Fx.catchTagsRecovers 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.causesEmits the source’s terminal failure cause and discards every successful value.
Fx.changesWithEffectDrops consecutive elements that are considered equal by an effectful predicate. When the effect returns
true, the element is skipped; whenfalse, it is emitted.This is the effectful variant of
skipRepeatsWith: instead of a pureEquivalence<A>, you supply(prev, next) => Effect<boolean>wheretruemeans “equal” (skip) andfalsemeans “changed” (emit).Fx.collectAllCollects all values emitted by an
Fxinto an array.Fx.collectAllForkForks the collection of all values from an
Fx.Fx.collectUpToCollects the first
nvalues emitted by anFxinto an array.Fx.collectUpToForkForks the collection of up to
nvalues from anFx.Fx.compactCompacts an Fx of Options, discarding
Nonevalues and unwrappingSomevalues.Fx.concatConcatenates 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.runis infallible, delivery alone does not suppress the continuation: the right source is run after the left run returns.Fx.concatMapMaps each element to an inner Fx and concatenates the results sequentially.
Fx.concatMapEffectMaps each element to an Effect and concatenates the results sequentially.
Fx.continueWithContinues an Fx with a lazily created Fx after the first run returns.
Fx.debounceEmits a value only after no newer source value arrives for
duration.Fx.delaySleeps for
durationbefore forwarding each successful source delivery.Fx.delimitWraps an Fx with a start and end value.
Fx.dieCreates an Fx that immediately terminates with a defect (unexpected error).
Fx.drainRuns an
Fxstream to completion, discarding all values. Useful when the side effects of the stream are all that matter.Fx.drainLayerRuns an
Fxstream as a Layer. The stream is forked in the background when the layer is acquired.Fx.dropAfterDrops elements from an Fx after a predicate returns true. The element that satisfies the predicate is included in the output.
Fx.dropUntilDrops 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.dropUntilEffectDrops 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.dropWhileAlias of
skipWhilefor Effect parity (dropWhilenaming).Fx.dropWhileEffectAlias of
skipWhileEffectfor Effect parity (dropWhileEffectnaming).Fx.duringForwards
eventsonly between a start signal and that signal’s first stop event.Fx.emptyAn Fx that emits no values and completes immediately.
Fx.ensuringRuns a finalizer with no typed error after the Fx run ends for any reason.
Fx.exhaustLatestMapMaps each element to an inner Fx, running one now and retaining only the latest waiting value.
Fx.exhaustLatestMapEffectMaps each element to an Effect, running one now and retaining only the latest waiting value.
Fx.exhaustMapMaps each element of an Fx to a new Fx, ignoring new elements until the current inner Fx completes.
Fx.exhaustMapEffectMaps each element of an Fx to an Effect, ignoring new elements until the current effect completes.
Fx.exitMaterializes every success and the terminal failure as infallible
Exitvalues.Fx.failCreates an Fx that immediately fails with the specified error.
Fx.failCauseCreates an Fx that immediately terminates with the specified Cause.
Fx.filterFilters elements of an Fx using a predicate function.
Fx.filterEffectFilters elements of an Fx using an effectful predicate function.
Fx.filterMapMaps and filters elements of an Fx in a single operation.
Fx.filterMapEffectMaps and filters elements of an Fx using an effectful function.
Fx.filterMapLoopLoops 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.filterMapLoopCauseLoops 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.filterMapLoopCauseEffectEffectfully loops over the failure causes of an Fx with an accumulator.
Fx.filterMapLoopEffectEffectfully loops over an Fx with an accumulator, producing an optional new value.
Fx.firstReturns the first value emitted by the
Fxwrapped in anOption. If theFxis empty, returnsNone.Fx.flatMapMaps each source value to an inner Fx and merges every inner concurrently.
Fx.flatMapConcurrentlyMaps each element of an Fx to a new Fx, running them concurrently with a limit.
Fx.flatMapConcurrentlyEffectMaps each element of an Fx to an Effect, running them concurrently with a limit.
Fx.flatMapEffectMaps each element of an Fx to an Effect, and merges the results.
Fx.flipEmits typed failures as values and fails with the first successful value.
Fx.fnCallable contracts implemented by
fn.Fx.forkForks the execution of an
Fxinto a background fiber. The stream will run until it completes or the fiber is interrupted.Fx.fromEffectCreates 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.fromFailuresCreates an Fx from a collection of failures (errors).
Fx.fromIterableCreates an Fx from an Iterable. Emits each value from the iterable in order and then completes.
Fx.fromScheduleCreates an Fx that emits values according to a Schedule. The Fx emits
voideach time the schedule fires.Fx.fromStreamAdapts an Effect
Streamto a push-basedFx.Fx.genBuilds an Fx by yielding Effects and returning the Fx to run afterward.
Fx.genScopedBuilds an Fx with a subscription-owned Scope shared by setup and streaming.
Fx.groupedPartitions 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 tonvalues, so callers own the memory policy for valid sizes. Invalid sizes fail withCause.IllegalArgumentError.Matches Effect
Stream.grouped.Fx.groupedWithinPartitions the stream into arrays, emitting when
nis reached ordurationelapses after the first element of the current group. The size must be a positive safe integer. A group can retain up tonvalues, so callers own the memory policy for valid sizes. Invalid sizes fail withCause.IllegalArgumentError.Matches Effect
Stream.groupedWithin.Fx.ifConditionally runs one of two Fx streams based on the boolean value emitted by the condition stream.
Fx.interruptCreates an Fx that immediately interrupts.
Fx.isFxChecks whether a value carries the
FxTypeIdprotocol property.Fx.keyedEfficiently 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
onValueto be called. - Existing keys have their
RefSubjectupdated with the new value. - Removed keys close the supplied child Scope and clean resources registered
through it; the
onValuerun fiber remains owned by the outer parent Scope.
- New keys cause
Fx.loopLoops over an Fx with an accumulator, producing a new value for each element and updating the accumulator.
Fx.loopCauseLoops over the failure causes of an Fx with an accumulator.
Fx.loopCauseEffectEffectfully loops over the failure causes of an Fx with an accumulator.
Fx.loopEffectEffectfully loops over an Fx with an accumulator, producing a new value for each element.
Fx.makeCreates 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.mapTransforms the elements of an Fx using a provided function.
Fx.mapBothTransforms both the success and error channels of an Fx using the provided options.
Mirrors
Effect.mapBoth:onSuccessmaps emitted values,onFailuremaps the typed failure (viaCause.map); defects and interrupts are preserved.Fx.mapEffectTransforms the elements of an Fx using a provided Effectful function.
Fx.mapErrorTransforms typed failures while preserving defects and interruption.
Fx.mergeMerges 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.runfail, somergedoes not itself cancel the sibling; a terminal observer may choose to.Fx.mergeAllMerges multiple Fx streams into a single Fx that emits values from all input streams concurrently.
Fx.mergeLeftMerges two Fx streams and emits only values from the left stream. Both streams run concurrently; completion when both complete.
Fx.mergeOrderedRuns multiple Fx streams concurrently while draining their values in argument order.
Fx.mergeRightMerges two Fx streams and emits only values from the right stream. Both streams run concurrently; completion when both complete.
Fx.neverAn Fx that waits forever without emitting a value.
Fx.nullAn Fx that emits
nullexactly once and then completes.Fx.observeObserves the values of an
Fxstream using a callback function. The callback can returnvoidor anEffectwhich will be executed for each value.Fx.observeLayerObserves the values of an
Fxstream using a callback function and returns aLayer. The callback can returnvoidor anEffectwhich will be executed for each value.Fx.onErrorRuns cleanup after the source reports a failure cause.
Fx.onExitObserves the Fx’s final success or failure with an Effect finalizer.
Fx.onInterruptRuns a finalizer when the Fx reports or externally receives interruption.
Fx.pairwiseEmits consecutive pairs
[previous, current]. The first value is not emitted until a second value arrives.Equivalent to RxJS
pairwiseand EffectStream.sliding(2)for pairs.Fx.periodicCreates an Fx that emits a
voidvalue periodically.Fx.prependPrepends a value to the beginning of an Fx.
Fx.provideBuilds a Layer for each subscription and provides it to the entire Fx run.
Fx.provideContextProvides an already-built Effect Context to the entire Fx run.
Fx.provideServiceProvides one existing service value to the entire Fx run.
Fx.provideServiceEffectAcquires one service with an Effect before running the Fx.
Fx.raceRuns 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.raceAllRaces many streams: the first to emit wins and the rest are interrupted.
Fx.repeatRepeats the entire stream according to
scheduleafter each successful completion. Failures are not repeated.Schedule.recurs(n)runs the streamn + 1times (the original plusnrepeats), matching EffectStream.repeat.Fx.resultMaterializes success and failure of an Fx as
Resultvalues.- 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 isCause<E>, so defects and interrupts are explicitly represented in theResultand the resulting Fx has error typenever.
The resulting Fx never fails at the stream level; all outcomes are emitted as
Result<A, Cause<E>>. Consumers can useResult.matchorResult.isSuccess/Result.isFailureto handle success vs failure (including defect/interrupt).- Success: each emitted value is wrapped as
Fx.retryRetries the entire stream when its Cause contains a typed
Failaccepted byschedule.The schedule is reset as soon as the first element of an attempt is emitted, matching Effect
Stream.retry.Fx.runForkRuns an
Fxin a new fiber, using the standardEffect.runFork. This is useful for integrating with the top-level Effect runtime.Fx.runPromiseRuns an
Fxstream to completion and returns a Promise. Rejects if the stream fails.Fx.runPromiseExitRuns an
Fxstream to completion and returns a Promise of the Exit.Fx.sampleEmits the latest source value whenever
sampleremits. 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.scanScans the stream with a pure function, emitting the accumulated state after each element. Emits the initial value first, then for each input
aemitsf(state, a)and updates state.Semantics align with Effect Stream’s
scan: output isinitial,f(initial, a1),f(..., a2), …Fx.scanEffectScans the stream with an effectful function, emitting the accumulated state after each element. Emits the initial value first, then for each input
arunsf(state, a)and emits the resulting state.Fx.sinceDrops
eventsuntilsignalemits, then forwards the rest.Fx.skipSkips the first
nelements of an Fx.Fx.skipEffectSkips the first
nelements wherenis produced by an Effect.Fx.skipRepeatsDrops elements that are equal to the previous element using standard equality.
Fx.skipRepeatsWithDrops elements that are equal to the previous element using a custom equivalence function.
Fx.skipWhileSkips 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.skipWhileEffectSkips 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.sliceSlices an Fx by skipping a number of elements and then taking a number of elements.
Fx.sliceEffectSlices an Fx with bounds produced by an Effect.
Fx.structCombines a record of Fx streams into a single Fx that emits a record of the latest values. Similar to
tuple, but for objects.Fx.succeedCreates an Fx that emits a single value and then completes.
Fx.succeedNullAn Fx that emits
nullexactly once and then completes.Fx.succeedUndefinedAn Fx that emits
undefinedexactly once and then completes.Fx.succeedVoidAn Fx that emits
voidexactly once and then completes.Fx.suspendDefers creation of an
Fxuntil each run begins.Fx.switchMapMaps 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.switchMapEffectMaps 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.syncLazily evaluates a synchronous function once for each Fx run.
Fx.takeTakes the first
nelements from an Fx and then completes.Fx.takeEffectTakes the first
nelements wherenis produced by an Effect.Fx.takeUntilTakes elements from an Fx until a predicate returns true. The element that satisfies the predicate is not included in the output.
Fx.takeUntilEffectTakes elements from an Fx until an effectful predicate returns true. The element that satisfies the predicate is not included in the output.
Fx.takeWhileTakes 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.takeWhileEffectTakes 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.tapPerforms an effect for each element of the Fx, without changing the elements.
Fx.throttleLimits 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.timeoutCompletes the stream if it does not produce a value (or complete) within
durationof the previous event. Matches EffectStream.timeout.The timeout is reset after each emission. An infinite duration is a no-op; a zero duration completes immediately.
Fx.timeoutToSwitches to
fallbackif the source does not produce a value withindurationof the previous event. Matches EffectStream.timeoutOrElseand RxJStimeoutTo.Fx.toStreamAdapts a push-based
Fxto an EffectStream.Fx.tupleCombines 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.undefinedAn Fx that emits
undefinedexactly once and then completes.Fx.untilForwards
eventsuntilsignalemits, then interruptsevents.Fx.unwrapUnwraps an Effect that produces an Fx into a single Fx.
Fx.unwrapScopedUnwraps 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.voidAn Fx that emits
voidexactly once and then completes.Fx.whenConditionally emits one of two values based on the boolean value emitted by the condition stream.
Fx.withLatestFromEmits
[source, latest]whenever the source emits, using the latest value fromthat. Source values are dropped untilthathas emitted at least once.Unlike {@link zipLatest}, this does not emit when
thatupdates.Completion: Completes when the source completes. Errors: The first failure from either stream fails the result.
Fx.withLatestFromWithLike {@link withLatestFrom}, but combines the pair with
f.Fx.withSpanTraces the whole subscription and each success or failure delivery.
Fx.zipZips 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.zipLatestZips 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.zipLatestWithZips 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.zipLeftZips two Fx streams in strict lockstep and emits only the left value. Completes when the first of the two streams completes.
Fx.zipRightZips two Fx streams in strict lockstep and emits only the right value. Completes when the first of the two streams completes.
Fx.zipWithZips 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.