Bidirectional contracts
Pushinterface
A bidirectional value that is both a
Sink<A, E, R>and anFx<B, E2, R2>.Calling
onSuccessoronFailuresends exactly one input notification to the wrapped Sink. The returnedEffectis the acknowledgment: a producer that runs and awaits it waits for the consumer callback to finish.Pushadds 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
exhaustLatestMapvariable
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.exhaustLatestMapEffectvariable
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.
exhaustMapvariable
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.exhaustMapEffectvariable
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.
flatMapvariable
Transforms each output value into an inner
Fxand 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.
flatMapEffectvariable
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.
switchMapvariable
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.
switchMapEffectvariable
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
mapErrorvariable
Transforms the output (Fx) error channel of a
Pushusing 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.mapErroron the Fx side. Cardinality, value order, input callbacks, and output service requirements are unchanged.
Push construction
makevariable
Couples a
Sinkinput with an independentFxoutput.The result forwards each input callback directly to
sinkand delegates every output subscription tofx. It does not connect the two values, change output cardinality or ordering, buffer inputs, or start either side eagerly.
Push services
Push.Classinterface
Constructable static type produced by
Push.Service.Push.Serviceinterface
The static and Effect service surface returned by
Push.Service.Service lookup supplies the same bidirectional value to
run,onSuccess, andonFailure. TheSelfservice appears in both required-service channels; the installed Push itself has those requirements captured by its Layer.Servicefunction
Defines a named Effect service whose value is a
Push.The returned class exposes
onSuccess,onFailure, andrunas Effects that first resolve the service from Context.makecaptures 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
filterInputvariable
Keeps successful inputs that satisfy
fand discards the rest.Calling
onSuccessruns 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.filterInputEffectvariable
Effectfully decides whether each successful input reaches the Sink.
Calling
onSuccess(value)invokesf(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.trueforwards one value,falsenone, and failure sends its Cause to the Sink failure callback. Calls are not serialized; the producer controls order and concurrency. Output behavior is unchanged.filterMapInputvariable
Transforms an input and forwards it only when
freturnsSome.Calling
onSuccessevaluatesfimmediately.Some(a)constructs one Sink callback Effect;Noneimmediately 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.filterMapInputEffectvariable
Effectfully transforms an input and forwards only a resulting
Somevalue.Calling
onSuccess(value)invokesf(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,Nonenone, and failure sends its full Cause to the Sink failure callback. Concurrent calls are not serialized.
Selecting outputs
filtervariable
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.
filterEffectvariable
Effectfully decides which Fx output values are emitted.
Each upstream value runs one predicate.
trueemits that value,falseemits 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.filterMapvariable
Maps each Fx output and emits only resulting
Somevalues.The upstream Sink evaluates
fsynchronously for each value, before its returned downstream callback Effect runs. Mapping allocation and throws occur at that callback invocation.Some(c)emits exactly onec;Noneemits nothing. Relative order, errors, services, and input behavior are preserved.filterMapEffectvariable
Effectfully maps each Fx output and emits only resulting
Somevalues.One mapper Effect runs per upstream value.
Some(c)emits once,Noneemits nothing, and failures join the output error channel. Sequential sources retain order; the input Sink remains unchanged.
Stateful outputs
mapAccumvariable
Maps over the output (Fx) side of a
Pushwith an accumulator: for each emitted valueb, appliesf(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 oneCper upstream value, in order. Accumulator state is private to each output subscription; the input Sink is unchanged.mapAccumEffectvariable
Maps over the output (Fx) side of a
Pushwith an effectful accumulator: for each emitted valueb, runsf(state, b)to get[nextState, emitted]and emits the second element. The adapter does not serialize callbacks. CallingonSuccessinvokesfimmediately 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
mapInputvariable
Synchronously transforms each successful input before sending it to the Sink.
One input maps to exactly one downstream input. Calling
onSuccessevaluatesfimmediately, 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.mapInputEffectvariable
Effectfully transforms each successful input before sending it to the Sink.
Calling
onSuccess(value)invokesf(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 singleAreaches 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
mapvariable
Synchronously transforms every value emitted by the Fx output side.
It emits exactly one
Cfor every upstreamB, preserving order and the complete input Sink. The mapping does not buffer or introduce concurrency.mapBothvariable
Transforms both the output (Fx) success and error channels of a
Pushusing the provided options.Mirrors
Effect.mapBothon the Fx side:onSuccessmaps every emitted value one-to-one andonFailuremaps typed failures viaCause.map; defects and interrupts are preserved. Ordering, services, and the input side are unchanged.mapEffectvariable
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 joinR2. The input side is unchanged.
Type contracts
Push.Anytype-alias
Matches any
Pushwhen its six channel types are intentionally unknown.