variable / @typed/fx/Push

mapAccumEffect

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.

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

Import

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

Signatures

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

Why

Use an effectful accumulator when each transition needs services or can fail, while keeping state local to observation rather than global application state.

Ownership and lifetime

Each output run owns one mutable seed, but concurrent producer callbacks may race over it; callers needing serialized state transitions must serialize the upstream deliveries. Each callback fiber runs and interrupts its own reducer Effect. E3 and R3 join the output channels. State is discarded when the subscription ends.

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 concurrent = Fx.make<number>((sink) =>
  Effect.all([sink.onSuccess(1), sink.onSuccess(2)], {
    concurrency: "unbounded",
    discard: true
  })
)
const push = Push.make(Sink.make(() => Effect.void, (_: string) => Effect.void), concurrent)
const totals = Push.mapAccumEffect(push, 0, (seed, value) =>
  Effect.sleep(value === 1 ? "20 millis" : "1 millis").pipe(
    Effect.as([seed + value, seed + value] as const)
  )
)
const program = Fx.collectAll(totals).pipe(Effect.scoped)
// Both reducers see seed 0; Effect.runPromise(program) resolves to [2, 1].

const values: Array<number> = []
let failures = 0
const sequential = Push.make(
  Sink.make(() => Effect.void, (_: string) => Effect.void),
  Fx.fromIterable([1, 2, 3])
)
const continued = Push.mapAccumEffect(sequential, 0, (seed, value) =>
  value === 2
    ? Effect.fail("rejected" as const)
    : Effect.succeed([seed + value, seed + value] as const)
)
const recovery = continued.run(Sink.make(
  () => Effect.sync(() => { failures += 1 }),
  (value) => Effect.sync(() => values.push(value))
))
// Effect.runPromise(recovery) leaves values [1, 4] and failures 1.

Other public imports

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

Source