Callback protocol
Emittype-alias
Operations supplied to a callback producer for emitting values or ending its run.
Callback sources
Collecting values
collectAllvariable
Collects all values emitted by an
Fxinto an array.collectAllForkvariable
Forks the collection of all values from an
Fx.collectUpTovariable
Collects the first
nvalues emitted by anFxinto an array.collectUpToForkvariable
Forks the collection of up to
nvalues from anFx.firstfunction
Returns the first value emitted by the
Fxwrapped in anOption. If theFxis empty, returnsNone.
Combining sources
appendvariable
Appends a value to the end of an Fx.
concatvariable
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.runis infallible, delivery alone does not suppress the continuation: the right source is run after the left run returns.continueWithvariable
Continues an Fx with a lazily created Fx after the first run returns.
delimitvariable
Wraps an Fx with a start and end value.
mergevariable
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.runfail, somergedoes not itself cancel the sibling; a terminal observer may choose to.mergeAllvariable
Merges multiple Fx streams into a single Fx that emits values from all input streams concurrently.
mergeLeftvariable
Merges two Fx streams and emits only values from the left stream. Both streams run concurrently; completion when both complete.
mergeOrderedfunction
Runs multiple Fx streams concurrently while draining their values in argument order.
mergeRightvariable
Merges two Fx streams and emits only values from the right stream. Both streams run concurrently; completion when both complete.
prependvariable
Prepends a value to the beginning of an Fx.
structfunction
Combines a record of Fx streams into a single Fx that emits a record of the latest values. Similar to
tuple, but for objects.tuplefunction
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.
withLatestFromvariable
Emits
[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.
withLatestFromWithvariable
Like {@link withLatestFrom}, but combines the pair with
f.zipvariable
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.
zipLatestvariable
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.
zipLatestWithvariable
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.
zipLeftvariable
Zips two Fx streams in strict lockstep and emits only the left value. Completes when the first of the two streams completes.
zipRightvariable
Zips two Fx streams in strict lockstep and emits only the right value. Completes when the first of the two streams completes.
zipWithvariable
Zips two Fx streams in strict lockstep and combines each pair with a function. Emits
f(a, b)when both streams have produced their next value.Completion: Completes when the first of the two streams completes. Errors: The first failure from either stream fails the result.
Concurrent work
concatMapvariable
Maps each element to an inner Fx and concatenates the results sequentially.
concatMapEffectvariable
Maps each element to an Effect and concatenates the results sequentially.
exhaustLatestMapvariable
Maps each element to an inner Fx, running one now and retaining only the latest waiting value.
exhaustLatestMapEffectvariable
Maps each element to an Effect, running one now and retaining only the latest waiting value.
exhaustMapvariable
Maps each element of an Fx to a new Fx, ignoring new elements until the current inner Fx completes.
exhaustMapEffectvariable
Maps each element of an Fx to an Effect, ignoring new elements until the current effect completes.
flatMapvariable
Maps each source value to an inner Fx and merges every inner concurrently.
flatMapConcurrentlyvariable
Maps each element of an Fx to a new Fx, running them concurrently with a limit.
flatMapConcurrentlyEffectvariable
Maps each element of an Fx to an Effect, running them concurrently with a limit.
flatMapEffectvariable
Maps each element of an Fx to an Effect, and merges the results.
racevariable
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.
raceAllvariable
Races many streams: the first to emit wins and the rest are interrupted.
switchMapvariable
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.
switchMapEffectvariable
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
Effect interop
fromEffectvariable
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
catchvariable · Re-export
Recovers from the first typed failure of an Fx by running a fallback Fx.
catchAllvariable
Uses the Effect-style
catchAllname for {@link catch}.catchCausevariable
Recovers from any failure cause by running a fallback Fx.
catchCauseIfvariable
Recovers a failure cause only when a predicate accepts the complete cause.
catchIfvariable
Recovers a typed failure only when a predicate accepts it.
catchTagvariable
Recovers selected tagged typed failures by running a fallback Fx.
catchTagsvariable
Recovers several tagged typed-error variants with one handler table.
catch_variable
Recovers from the first typed failure of an Fx by running a fallback Fx.
causesvariable
Emits the source’s terminal failure cause and discards every successful value.
exitvariable
Materializes every success and the terminal failure as infallible
Exitvalues.flipvariable
Emits typed failures as values and fails with the first successful value.
mapBothvariable
Transforms 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.mapErrorvariable
Transforms typed failures while preserving defects and interruption.
resultvariable
Materializes 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
retryvariable
Retries 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.
Failure sources
dievariable
Creates an Fx that immediately terminates with a defect (unexpected error).
failvariable
Creates an Fx that immediately fails with the specified error.
failCausevariable
Creates an Fx that immediately terminates with the specified Cause.
fromFailuresvariable
Creates an Fx from a collection of failures (errors).
interruptvariable
Creates an Fx that immediately interrupts.
Generator composition
fnnamespace
Callable contracts implemented by
fn.genvariable
Builds an Fx by yielding Effects and returning the Fx to run afterward.
genScopedvariable
Builds an Fx with a subscription-owned Scope shared by setup and streaming.
unwrapvariable
Unwraps an Effect that produces an Fx into a single Fx.
unwrapScopedvariable
Unwraps an Effect that produces an Fx into a single Fx, managing the scope of the effect.
The scope of the effect is closed when the Fx completes or is interrupted.
Keyed work
keyedvariable
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
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
Observing failures
Operator options
Boundsinterface
Defines the bounds for slicing an Fx stream.
FromStreamOptionstype-alias
Effect Stream mapping options used while delivering elements to an
Fxsink.KeyedOptionsinterface
Configuration options for the
keyedcombinator.ThrottleOptionstype-alias
Options for {@link throttle}. A duration-only call is leading-edge (
{ leading: true, trailing: false }).ToStreamOptionstype-alias
Buffering and callback options accepted while adapting an
Fxto an EffectStream.
Providing services
Fx.Serviceinterface
An
Fxwhose implementation is obtained from an Effect service.Servicefunction
Defines an Effect service whose value is also a directly runnable
Fx.providevariable
Builds a Layer for each subscription and provides it to the entire Fx run.
provideContextvariable
Provides an already-built Effect Context to the entire Fx run.
provideServicevariable
Provides one existing service value to the entire Fx run.
provideServiceEffectvariable
Acquires one service with an Effect before running the Fx.
Resource lifetime
ensuringvariable
Runs a finalizer with no typed error after the Fx run ends for any reason.
onExitvariable
Observes the Fx’s final success or failure with an Effect finalizer.
onInterruptvariable
Runs a finalizer when the Fx reports or externally receives interruption.
Running effects
drainvariable
Runs an
Fxstream to completion, discarding all values. Useful when the side effects of the stream are all that matter.drainLayervariable
Runs an
Fxstream as a Layer. The stream is forked in the background when the layer is acquired.forkvariable
Forks the execution of an
Fxinto a background fiber. The stream will run until it completes or the fiber is interrupted.observevariable
Observes the values of an
Fxstream using a callback function. The callback can returnvoidor anEffectwhich will be executed for each value.observeLayervariable
Observes the values of an
Fxstream using a callback function and returns aLayer. The callback can returnvoidor anEffectwhich will be executed for each value.runForkvariable
Runs an
Fxin a new fiber, using the standardEffect.runFork. This is useful for integrating with the top-level Effect runtime.runPromisevariable
Runs an
Fxstream to completion and returns a Promise. Rejects if the stream fails.runPromiseExitvariable
Runs an
Fxstream to completion and returns a Promise of the Exit.
Runtime inspection
Selecting values
changesWithEffectvariable
Drops 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).compactvariable
Compacts an Fx of Options, discarding
Nonevalues and unwrappingSomevalues.dropAftervariable
Drops elements from an Fx after a predicate returns true. The element that satisfies the predicate is included in the output.
dropUntilvariable
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.
dropUntilEffectvariable
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.
dropWhilevariable
Alias of
skipWhilefor Effect parity (dropWhilenaming).dropWhileEffectvariable
Alias of
skipWhileEffectfor Effect parity (dropWhileEffectnaming).filtervariable
Filters elements of an Fx using a predicate function.
filterEffectvariable
Filters elements of an Fx using an effectful predicate function.
filterMapvariable
Maps and filters elements of an Fx in a single operation.
filterMapEffectvariable
Maps and filters elements of an Fx using an effectful function.
skipvariable
Skips the first
nelements of an Fx.skipEffectvariable
Skips the first
nelements wherenis produced by an Effect.skipRepeatsvariable
Drops elements that are equal to the previous element using standard equality.
skipRepeatsWithvariable
Drops elements that are equal to the previous element using a custom equivalence function.
skipWhilevariable
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.
skipWhileEffectvariable
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.
slicevariable
Slices an Fx by skipping a number of elements and then taking a number of elements.
sliceEffectvariable
Slices an Fx with bounds produced by an Effect.
takevariable
Takes the first
nelements from an Fx and then completes.takeEffectvariable
Takes the first
nelements wherenis produced by an Effect.takeUntilvariable
Takes elements from an Fx until a predicate returns true. The element that satisfies the predicate is not included in the output.
takeUntilEffectvariable
Takes elements from an Fx until an effectful predicate returns true. The element that satisfies the predicate is not included in the output.
takeWhilevariable
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.
takeWhileEffectvariable
Takes elements from an Fx while an effectful predicate returns true. Stops at the first element for which the predicate effect succeeds with false; that element is not included.
Stateful transforms
filterMapLoopvariable
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.filterMapLoopCausevariable
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.
filterMapLoopCauseEffectvariable
Effectfully loops over the failure causes of an Fx with an accumulator.
filterMapLoopEffectvariable
Effectfully loops over an Fx with an accumulator, producing an optional new value.
groupedvariable
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 tonvalues, so callers own the memory policy for valid sizes. Invalid sizes fail withCause.IllegalArgumentError.Matches Effect
Stream.grouped.loopvariable
Loops over an Fx with an accumulator, producing a new value for each element and updating the accumulator.
loopCausevariable
Loops over the failure causes of an Fx with an accumulator.
loopCauseEffectvariable
Effectfully loops over the failure causes of an Fx with an accumulator.
loopEffectvariable
Effectfully loops over an Fx with an accumulator, producing a new value for each element.
pairwisevariable
Emits consecutive pairs
[previous, current]. The first value is not emitted until a second value arrives.Equivalent to RxJS
pairwiseand EffectStream.sliding(2)for pairs.scanvariable
Scans 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), …scanEffectvariable
Scans 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.
Stream interop
fromStreamvariable
Adapts an Effect
Streamto a push-basedFx.toStreamvariable
Adapts a push-based
Fxto an EffectStream.
Time and rate
atvariable
Creates an Fx that emits a single value after a specified delay.
debouncevariable
Emits a value only after no newer source value arrives for
duration.delayvariable
Sleeps for
durationbefore forwarding each successful source delivery.duringfunction
Forwards
eventsonly between a start signal and that signal’s first stop event.fromSchedulevariable
Creates an Fx that emits values according to a Schedule. The Fx emits
voideach time the schedule fires.groupedWithinvariable
Partitions 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.periodicvariable
Creates an Fx that emits a
voidvalue periodically.repeatvariable
Repeats 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.samplevariable
Emits 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.
sincevariable
Drops
eventsuntilsignalemits, then forwards the rest.throttlevariable
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 }.timeoutvariable
Completes 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.
timeoutTovariable
Switches to
fallbackif the source does not produce a value withindurationof the previous event. Matches EffectStream.timeoutOrElseand RxJStimeoutTo.untilvariable
Forwards
eventsuntilsignalemits, then interruptsevents.
Transforming values
asvariable
Replaces all emitted values from the Fx with the provided value
b.mapvariable
Transforms the elements of an Fx using a provided function.
mapEffectvariable
Transforms the elements of an Fx using a provided Effectful function.
tapvariable
Performs an effect for each element of the Fx, without changing the elements.
Type contracts
Errortype-alias
Alias of
Fx.Errorfor extracting anFxerror channel.FlatMapEffectLiketype-alias
Describes a dual flattening operator whose callback returns an Effect.
FlatMapLiketype-alias
Describes a dual flattening operator whose callback returns an Fx.
Fxinterface
Fxis 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.Errortype-alias
Extracts the typed error from an
Fx.Fx.Servicestype-alias
Extracts the services required to run an
Fx.Fx.Successtype-alias
Extracts the emitted value type from an
Fx.Servicestype-alias
Alias of
Fx.Servicesfor extracting anFxservice channel.Successtype-alias
Alias of
Fx.Successfor extracting anFxvalue channel.fn.Gentype-alias
Contract for functions whose body yields Effects and returns an
Fx.fn.NonGentype-alias
Contract for functions whose body returns an
Fxdirectly.
Value sources
emptyvariable
An Fx that emits no values and completes immediately.
fromIterablevariable
Creates an Fx from an Iterable. Emits each value from the iterable in order and then completes.
nevervariable
An Fx that waits forever without emitting a value.
nullvariable · Re-export
An Fx that emits
nullexactly once and then completes.succeedvariable
Creates an Fx that emits a single value and then completes.
succeedNullvariable
An Fx that emits
nullexactly once and then completes.succeedUndefinedvariable
An Fx that emits
undefinedexactly once and then completes.succeedVoidvariable
An Fx that emits
voidexactly once and then completes.suspendvariable
Defers creation of an
Fxuntil each run begins.syncvariable
Lazily evaluates a synchronous function once for each Fx run.
undefinedvariable · Re-export
An Fx that emits
undefinedexactly once and then completes.voidvariable · Re-export
An Fx that emits
voidexactly once and then completes.
models
Fx.Anytype-alias
Matches any
Fxregardless of its value, error, or service channels.Fx.Classinterface
The constructible service class returned by
Fx.Service.Fx.Varianceinterface
Describes how an
Fxvaries in its value, error, and service channels.