variable / @typed/fx/Push

mapAccum

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.

Package version
2.0.0-beta.7
Category
Stateful outputs
Since
1.0.0

Import

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

Signatures

export declare const mapAccum: {
    <S, B, C>(initial: S, f: (s: S, b: B) => readonly [
        S,
        C
    ]): <A, E, R, E2, R2>(push: Push<A, E, R, B, E2, R2>) => Push<A, E, R, C, E2, R2>;
    <A, E, R, B, E2, R2, S, C>(push: Push<A, E, R, B, E2, R2>, initial: S, f: (s: S, b: B) => readonly [
        S,
        C
    ]): Push<A, E, R, C, E2, R2>;
};

Why

Stateful output projection can derive running totals or protocol state without moving that state into the bidirectional input boundary.

Ownership and lifetime

Each run owns its own synchronous accumulator. It is discarded when that run completes or is interrupted and acquires no Scope or external resource.

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 push = Push.make(Sink.make(() => Effect.void, (_: string) => Effect.void), Fx.fromIterable([1, 2, 3]))
const totals = Push.mapAccum(push, 0, (total, value) => [total + value, total + value] as const)

Other public imports

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

Source