Import
import { exhaustMap } from "@typed/fx/Push";Signatures
export declare const exhaustMap: {
<B, C, E3, R3>(f: (b: B) => Fx.Fx<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) => Fx.Fx<C, E3, R3>): Push<A, E, R, C, E2 | E3, R2 | R3 | Scope.Scope>;
};Why
Exhaust semantics prevent duplicate work, such as repeated submit clicks, while allowing another request after the active one completes.
Ownership and lifetime
The output Scope owns the active inner fiber, joins it before completion, and interrupts it with the subscription. Inner failures and services join the outer channels; there is no pending-value buffer.
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 started: Array<number> = []
const finalized: Array<number> = []
const push = Push.make(
Sink.make(() => Effect.void, (_: string) => Effect.void),
Fx.fromIterable([1, 2])
)
const exhausted = Push.exhaustMap(push, (id) => {
constructed.push(id)
return Fx.make<number>((sink) =>
Effect.gen(function* () {
yield* Effect.sync(() => started.push(id))
yield* Effect.sleep("5 millis")
yield* sink.onSuccess(id * 10 + 1)
yield* Effect.sleep("5 millis")
yield* sink.onSuccess(id * 10 + 2)
}).pipe(Effect.ensuring(Effect.sync(() => finalized.push(id))))
)
})
const program = Fx.collectAll(exhausted).pipe(
Effect.map((values) => ({ values, constructed, started, finalized })),
Effect.scoped
)
// => { values: [11, 12], constructed: [1, 2], started: [1], finalized: [1] }Other public imports
These import paths expose the same declaration. Each page retains its own public name and signature.