variable / @typed/fx/Push

exhaustMap

Runs at most one inner Fx and ignores outer values while it is active.

The first value seen while idle starts an inner; every value arriving before that inner completes is dropped. f(value) is still evaluated and its inner Fx is constructed before the busy check; dropping means the returned Fx is not run. Accepted inners preserve their own order and all values. Input is unchanged.

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

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.

Source