Import
import { exhaustLatestMapEffect } from "@typed/fx/Push";Signatures
export declare const exhaustLatestMapEffect: {
<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 provides non-overlapping Effect work with eventual latest-state handling, avoiding an unbounded queue of obsolete requests.
Ownership and lifetime
The output Scope owns the active Effect fiber and pending slot. It waits for the final Effect, interrupts it with the subscription, and runs its finalizers. Effect failures and services join the output channels.
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, 3])
)
const latestSaved = Push.exhaustLatestMapEffect(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(latestSaved).pipe(
Effect.map((values) => ({ values, constructed, finalized })),
Effect.scoped
)
// Effect.runPromise(program) => { values: [1, 3], constructed: [1, 2, 3], finalized: [1, 3] }Other public imports
These import paths expose the same declaration. Each page retains its own public name and signature.