variable / @typed/fx/Push

exhaustLatestMap

Runs one inner Fx at a time and retains only the latest value received while busy.

A value while idle starts immediately. While its inner runs, newer outer values replace a single pending slot. f(value) is evaluated and an inner Fx is constructed before every replacement; a superseded pending Fx is never run. After completion, only the latest pending Fx starts. Accepted inners preserve their own order and all values. The input Sink is unchanged.

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

Import

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

Signatures

export declare const exhaustLatestMap: {
    <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

This is useful when work must not overlap but the newest requested state must eventually be processed.

Ownership and lifetime

The output Scope owns one active inner fiber and an in-memory latest slot. It joins the final inner before completion and interrupts it with the subscription. Inner 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 started: Array<number> = []
const finalized: Array<number> = []
const push = Push.make(
  Sink.make(() => Effect.void, (_: string) => Effect.void),
  Fx.fromIterable([1, 2, 3])
)
const latest = Push.exhaustLatestMap(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(latest).pipe(
  Effect.map((values) => ({ values, constructed, started, finalized })),
  Effect.scoped
)
// => values [11, 12, 31, 32]; constructed [1, 2, 3]; started/finalized [1, 3].

Other public imports

These import paths expose the same declaration. Each page retains its own public name and signature.

Source