variable / @typed/fx/Push

mapInput

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.

Package version
2.0.0-beta.7
Category
Transforming inputs
Since
1.0.0

Import

import { mapInput } from "@typed/fx/Push";

Signatures

export declare const mapInput: {
    <C, A>(f: (c: C) => A): <E, R, B, E2, R2>(push: Push<A, E, R, B, E2, R2>) => Push<C, E, R, B, E2, R2>;
    <A, E, R, B, E2, R2, C>(push: Push<A, E, R, B, E2, R2>, f: (c: C) => A): Push<C, E, R, B, E2, R2>;
};

Why

Adapt an external command shape to an existing consumer without rebuilding or changing the independent output stream.

Ownership and lifetime

f runs synchronously when onSuccess is called, before its returned Effect exists, and acquires no resource. Downstream handling retains the original Sink’s E, R, caller fiber, interruption, and cleanup behavior.

Examples

import { Effect } from "effect"
import * as Fx from "@typed/fx/Fx"
import * as Push from "@typed/fx/Push"
import * as Sink from "@typed/fx/Sink"

const events: Array<string> = []
const base = Push.make(
  Sink.make(() => Effect.void, (n: number) => Effect.sync(() => events.push(`sink:${n}`))),
  Fx.empty
)
const mapped = Push.mapInput(base, (text: string) => {
  events.push(`map:${text}`)
  return Number(text)
})
const acknowledgement = mapped.onSuccess("42")
// events is already ["map:42"]; the Sink Effect has not run yet.
const program = Effect.as(acknowledgement, events)

Other public imports

These import paths expose the same declaration. Each page retains its own public name and signature.

Source