Module / v2.0.0-beta.13

@typed/fx/Push

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

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

Bidirectional contracts

  • Push

    interface

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

Concurrent output work

  • exhaustLatestMap

    variable

    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.

  • exhaustLatestMapEffect

    variable

    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.

  • exhaustMap

    variable

    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.

  • exhaustMapEffect

    variable

    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.

  • flatMap

    variable

    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.

  • flatMapEffect

    variable

    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.

  • switchMap

    variable

    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.

  • switchMapEffect

    variable

    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.

Output failures

  • mapError

    variable

    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.

Push construction

  • make

    variable

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

    interface

    Constructable static type produced by Push.Service.

  • Push.Service

    interface

    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.

  • Service

    function

    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.

Selecting inputs

  • filterInput

    variable

    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.

  • filterInputEffect

    variable

    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.

  • filterMapInput

    variable

    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.

  • filterMapInputEffect

    variable

    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.

Selecting outputs

  • filter

    variable

    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.

  • filterEffect

    variable

    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.

  • filterMap

    variable

    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.

  • filterMapEffect

    variable

    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.

Stateful outputs

  • mapAccum

    variable

    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.

  • mapAccumEffect

    variable

    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.

Transforming inputs

  • mapInput

    variable

    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.

  • mapInputEffect

    variable

    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.

Transforming outputs

  • map

    variable

    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.

  • mapBoth

    variable

    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.

  • mapEffect

    variable

    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.

Type contracts

  • Push.Any

    type-alias

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