Follow one execution
Maps each source value to an inner Fx and merges every inner concurrently.
flatMapTick 0. Input: a; A: start.
Read this diagram
Follow each lane from left to right. Events stacked vertically share a tick; the green cursor marks the current time across every lane.
- A value
The text inside the pill is the emitted value. - Work starts
The raised chevron starts an inner run (^in the source). - The run returns
The vertical bar ends this lane’s run. - A cause is delivered
The exclamation mark belongs to this lane. - Work is interrupted
The cross marks cancellation of this run. - Current time
The line and diamond move together across all lanes. - Happening now
A highlighted event is at the current tick. - Still ahead
Muted, dashed values have not happened yet. - Time continues
The lane’s arrow is not a return marker. An empty stretch can be quiet work that is still running.
Illustrated ticks start at 0. At 1×, one illustrated tick takes one second; captions specify real durations when timing matters. A cause or interruption belongs to its lane, and other work may continue. Scroll horizontally to inspect the rest of a long timeline.
Source and inner failures are forwarded; their required services are unioned. A FiberSet in the required Scope owns all inner fibers. Source completion waits for the set to empty. Interrupting observation closes the Scope and interrupts every active inner, running each inner Scope’s finalizers.
Import
import { flatMap } from "@typed/fx/Fx";Signatures
export declare const flatMap: FlatMapLike;Why
flatMap is the unconstrained merge policy for independent pushed work. It
preserves push-based production instead of collecting an inner Fx before the
next source value can be handled.
Concurrency, ordering, and cardinality
Every source value creates exactly one inner Fx with no concurrency limit. Each inner retains its own emission order, but values from different inners interleave according to arrival time. There is no output buffer or global ordering guarantee.
Ownership and lifetime
Source and inner failures are forwarded; their required services are unioned.
A FiberSet in the required Scope owns all inner fibers. Source completion
waits for the set to empty. Interrupting observation closes the Scope and
interrupts every active inner, running each inner Scope’s finalizers.
Examples
import { Fx } from "@typed/fx"
import { Effect } from "effect"
const rows = Fx.flatMap(Fx.fromIterable([
{ id: "slow", wait: "20 millis" as const },
{ id: "fast", wait: "1 millis" as const }
]), ({ id, wait }) => Fx.at(id, wait))
Effect.runPromise(Effect.scoped(Fx.collectAll(rows))).then(console.log)
// ["fast", "slow"]: inners overlap and results arrive by production timeOther public imports
These import paths expose the same declaration. Each page retains its own public name and signature.