Follow one execution
Maps each element of an Fx to a new Fx, running them concurrently with a limit. This scenario uses concurrent source deliveries: a pending source callback can wait for admission while another callback is already active.
flatMapConcurrentlyTick 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.
A non-positive, fractional, infinite, or unsafe-integer limit fails through the Sink with Cause.IllegalArgumentError. Source and inner failures and services remain typed. The required Scope owns waiting and active fibers; source completion waits for all of them, and interruption cancels both groups and runs active inner finalizers.
Import
import { flatMapConcurrently } from "@typed/fx/Fx";Signatures
export declare const flatMapConcurrently: {
<A, B, E2, R2>(f: (a: A) => Fx<B, E2, R2>, concurrency: number): <E, R>(self: Fx<A, E, R>) => Fx<B, E | E2 | Cause.IllegalArgumentError, R | R2 | Scope.Scope>;
<A, E, R, B, E2, R2>(self: Fx<A, E, R>, f: (a: A) => Fx<B, E2, R2>, concurrency: number): Fx<B, E | E2 | Cause.IllegalArgumentError, R | R2 | Scope.Scope>;
};Why
This is the admission-controlled form of {@link flatMap }. It bounds active work while retaining every source value, which is appropriate when a remote service or local resource has a known concurrency budget.
Concurrency, ordering, and buffering
Each source value creates one child fiber. A semaphore admits at most
concurrency inners into execution; excess fibers wait for a permit rather
than being dropped. Values from each inner retain local order, but admitted
inners may interleave and output is not reordered into source order.
Ownership and lifetime
A non-positive, fractional, infinite, or unsafe-integer limit fails through
the Sink with Cause.IllegalArgumentError. Source and inner failures and
services remain typed. The required Scope owns waiting and active fibers;
source completion waits for all of them, and interruption cancels both groups
and runs active inner finalizers.
Examples
import { Fx } from "@typed/fx"
import { Effect } from "effect"
const jobs = Fx.fromIterable([
{ id: "one", wait: "20 millis" as const },
{ id: "two", wait: "20 millis" as const },
{ id: "three", wait: "1 millis" as const }
])
const bounded = Fx.flatMapConcurrently(
jobs,
({ id, wait }) => Fx.at(id, wait),
2
)
Effect.runPromise(Effect.scoped(Fx.collectAll(bounded))).then(console.log)
// "three" waits for one of the first two permits despite its shorter delayOther public imports
These import paths expose the same declaration. Each page retains its own public name and signature.