variable / @typed/fx

Push.flatMap

Transforms each output value into an inner Fx and merges all inners concurrently.

Every outer value starts one inner. All inner values are emitted; order within each inner is preserved, but values from different inners may interleave. The outer stream waits for all inners before normal completion. The input Sink 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.flatMap 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 flatMap: {
    <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

Use concurrent flattening when every produced task matters and independent work should overlap.

Ownership and lifetime

The output Scope owns a fiber set containing all active inners. Outer interruption interrupts that set and runs inner cleanup. Inner failures and services join the outer Fx channels; no ordering buffer is added.

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

Other public imports

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

Source