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
Member of
Push

Import

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

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