variable / @typed/fx/Push

flatMapEffect

Transforms every output value into an Effect and merges their results concurrently.

One Effect starts per outer value and can emit one result. All successful results are emitted, but concurrent completion order may differ from input order. The input side is unchanged.

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

Import

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

Signatures

export declare const flatMapEffect: {
    <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 runs independent Effect work concurrently without manually lifting each Effect into Fx.

Ownership and lifetime

Active Effects are fibers owned by the output Scope. Completion waits for all of them; interruption stops them and runs their finalizers. E3 and R3 join the output failure and service 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 finalized: Array<number> = []
const push = Push.make(
  Sink.make(() => Effect.void, (_: string) => Effect.void),
  Fx.fromIterable([1, 2])
)
const loaded = Push.flatMapEffect(push, (id) =>
  Effect.sleep(id === 1 ? "20 millis" : "1 millis").pipe(
    Effect.as(id),
    Effect.ensuring(Effect.sync(() => finalized.push(id)))
  )
)
const program = Fx.collectAll(loaded).pipe(
  Effect.map((values) => ({ values, finalized })),
  Effect.scoped
)
// Effect.runPromise(program) => { values: [2, 1], finalized: [2, 1] }

Other public imports

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

Source