Import
import { Push } from "@typed/fx";Access Push.mapAccumEffect 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 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.