variable / @typed/fx

Push.mapInputEffect

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.

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

Import

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

Access Push.mapInputEffect through the imported Push export. Its declaration below describes the member.

This public exposure is a re-export. Its import path is supported; declaration documentation is shared with the other public exposures below.

Signatures

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

Why

Use this at an input boundary that must decode, validate, or load Effect services before the existing consumer can accept a value.

Ownership and lifetime

Effect construction happens eagerly at onSuccess; only execution is deferred. The fiber running the returned acknowledgment owns execution and interruption of the constructed Effect and downstream callback. E is handled by the Sink failure channel, while R3 joins the input service requirements.

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 numbers: Array<number> = []
const base = Push.make(
  Sink.make(() => Effect.void, (n: number) => Effect.sync(() => numbers.push(n))),
  Fx.empty
)
const parsed = Push.mapInputEffect(base, (text: string) =>
  text === "" ? Effect.fail("empty" as const) : Effect.succeed(Number(text))
)

Other public imports

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

Source