Import
import { exhaustMapEffect } from "@typed/fx/Push";Signatures
export declare const exhaustMapEffect: {
<B, C, E3, R3>(f: (b: B) => Effect.Effect<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 | Scope.Scope>;
<A, E, R, B, E2, R2, C, E3, R3>(push: Push<A, E, R, B, E2, R2>, f: (b: B) => Effect.Effect<C, E3, R3>): Push<A, E, R, C, E2 | E3, R2 | R3 | Scope.Scope>;
};Why
This guards non-overlapping Effect work without manual busy state.
Ownership and lifetime
The output Scope owns and joins the active Effect fiber. Subscription interruption stops it and runs finalizers. Its errors and services join the output channels; no value is buffered for later.
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 constructed: Array<number> = []
const finalized: Array<number> = []
const push = Push.make(
Sink.make(() => Effect.void, (_: string) => Effect.void),
Fx.fromIterable([1, 2])
)
const saving = Push.exhaustMapEffect(push, (id) => {
constructed.push(id)
return Effect.sleep("10 millis").pipe(
Effect.as(id),
Effect.ensuring(Effect.sync(() => finalized.push(id)))
)
})
const program = Fx.collectAll(saving).pipe(
Effect.map((values) => ({ values, constructed, finalized })),
Effect.scoped
)
// Effect.runPromise(program) => { values: [1], constructed: [1, 2], finalized: [1] }Other public imports
These import paths expose the same declaration. Each page retains its own public name and signature.