variable / @typed/fx

Push.exhaustMapEffect

Runs at most one mapped Effect and ignores values while it is active.

The first value while idle starts one Effect and can emit one result; all values received before it completes are dropped. The mapping callback is still invoked and constructs an Effect for every value before the busy check; dropped Effects are not run. The input side is unchanged.

Package version
2.0.0-beta.13
Category
Concurrent output work
Since
1.0.0
Member of
Push

Import

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

Access Push.exhaustMapEffect 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 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.

Source