variable / @typed/fx/Push

switchMap

Transforms each output value into an inner Fx, observing only the latest one.

A new outer value interrupts the previous inner fiber before starting the next. Output cardinality is the cardinality of the successive active inners; values from an interrupted inner stop. Outer order determines replacement order, while each active inner preserves its own order. The input Sink is unchanged.

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

Import

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

Signatures

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

Latest-wins composition models search, navigation, and other work where a newer request makes the previous result irrelevant.

Ownership and lifetime

Running the result requires Scope.Scope. Each inner runs in that Scope and is interrupted on replacement or outer interruption; the last inner is joined before normal completion. Outer and inner errors/services are combined.

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 finalized: Array<number> = []
const outer = Fx.make<number>((sink) =>
  sink.onSuccess(1).pipe(
    Effect.andThen(Effect.sleep("5 millis")),
    Effect.andThen(sink.onSuccess(2))
  )
)
const push = Push.make(Sink.make(() => Effect.void, (_: string) => Effect.void), outer)
const switched = Push.switchMap(push, (id) =>
  Fx.make<number>((sink) =>
    Effect.gen(function* () {
      yield* Effect.sleep("2 millis")
      yield* sink.onSuccess(id * 10 + 1)
      yield* Effect.sleep(id === 1 ? "20 millis" : "2 millis")
      yield* sink.onSuccess(id * 10 + 2)
    }).pipe(Effect.ensuring(Effect.sync(() => finalized.push(id))))
  )
)
const program = Fx.collectAll(switched).pipe(
  Effect.map((values) => ({ values, finalized })),
  Effect.scoped
)
// Effect.runPromise(program) => { values: [11, 21, 22], finalized: [1, 2] }

Other public imports

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

Source