Import
import { flatMap } from "@typed/fx/Push";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.