Follow one execution
Two permits admit a and b; c waits until b releases a permit. Results follow completion order. This scenario uses concurrent source deliveries: a pending source callback can wait for admission while another callback is already active.
flatMapConcurrentlyEffectTick 0. Input: a; Load-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.
Invalid limits fail with Cause.IllegalArgumentError; source and callback failures remain typed and callback services are added to requirements. The required Scope owns waiting and active Effects. Completion drains them; interruption cancels them and runs their finalizers.
Import
import { Fx } from "@typed/fx";Access Fx.flatMapConcurrentlyEffect through the imported Fx 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 flatMapConcurrentlyEffect: {
<A, B, E2, R2>(f: (a: A) => Effect.Effect<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) => Effect.Effect<B, E2, R2>, concurrency: number): Fx<B, E | E2 | Cause.IllegalArgumentError, R | R2 | Scope.Scope>;
};Why
This Effect-producing bounded merge avoids manual Fx.fromEffect conversion
while making the callback’s concurrency budget part of the operation.
Concurrency, ordering, and cardinality
Every source value eventually starts one callback Effect, with no more than
concurrency active at once. Each success emits exactly one value. Waiting
work is retained, and results arrive by completion rather than source order.
Ownership and lifetime
Invalid limits fail with Cause.IllegalArgumentError; source and callback
failures remain typed and callback services are added to requirements. The
required Scope owns waiting and active Effects. Completion drains them;
interruption cancels them and runs their finalizers.
Examples
import { Fx } from "@typed/fx"
import { Effect } from "effect"
const loaded = Fx.flatMapConcurrentlyEffect(
Fx.fromIterable([
{ id: "one", wait: "20 millis" as const },
{ id: "two", wait: "20 millis" as const },
{ id: "three", wait: "1 millis" as const }
]),
({ id, wait }) => Effect.as(Effect.sleep(wait), id),
2
)
Effect.runPromise(Effect.scoped(Fx.collectAll(loaded))).then(console.log)
// "three" remains queued until a permit is releasedOther public imports
These import paths expose the same declaration. Each page retains its own public name and signature.