variable / @typed/fx

Push.exhaustLatestMapEffect

Runs one mapped Effect at a time and retains only the latest value received while busy.

Each accepted Effect can emit one result. While it runs, one pending Effect is repeatedly overwritten; after completion only the latest pending Effect starts. The mapping callback still runs and constructs an Effect for every value before replacement; superseded 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.exhaustLatestMapEffect 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 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.

Source