A document editor starts several kinds of work: load previews, upload attachments, save revisions, and submit a final command. A new input arriving while old work is active must have an intentional meaning. Running everything, canceling old work, and dropping repeated commands are different user promises, even when all three call the same server.
Composing Fx combines independent producers. Here an outer value creates an
inner Fx. The inner may emit progress and a final result, fail, or remain live. A flattening operator
owns the relationship between those runs. Effect supplies the ability to fork and interrupt fibers;
the higher-order combinator composes those operations into a declarative policy. For example,
switchMap says that new input replaces the old inner execution, including its interruption and
cleanup. You choose the relationship instead of maintaining a current-fiber variable yourself.
For arrivals at 0, 5, and 10 milliseconds and 20-millisecond jobs, immediate-finalization assumptions
give these outcomes: concatMap finishes all three at 60; switchMap locally finishes only c at
30; exhaustMap finishes only a at 20; exhaustLatestMap finishes a then c at 40. The choice
changes what the user ultimately saved, not just throughput.
| Need | Use | Primary contract |
|---|---|---|
| Allow independent jobs to overlap | flatMapConcurrently | Bound concurrent inner runs. |
| Process every job in order | concatMap | Wait for each inner run to complete. |
| Replace obsolete work | switchMap | Interrupt the prior inner run when new input arrives. |
| Ignore arrivals while busy | exhaustMap | Finish the active run and drop intervening inputs. |
| Finish active work, then use the newest input | exhaustLatestMap | Retain only the latest pending input. |
Let independent lookups overlap
When one input’s work does not invalidate another’s, flatMap starts every inner immediately:
import { Fx } from "@typed/fx";
const requests = Fx.mergeAll(
Fx.at("a", "0 millis"),
Fx.at("b", "5 millis"),
);
const load = (id: string) =>
Fx.mergeAll(
Fx.at(`${id}:cached`, "10 millis"),
Fx.at(`${id}:fresh`, "40 millis"),
);
const results = requests.pipe(Fx.flatMap(load));
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.
The b lane starts while a is still active. Both cached and fresh values survive, and the output
interleaves them. There is no ordering between different inners. This policy can create an unbounded
number of active jobs if input keeps arriving faster than jobs finish.
For uploads, bound active jobs while retaining every selected file:
import { Fx } from "@typed/fx";
const attachments = Fx.fromIterable(["a", "b", "c"]);
const upload = (file: string) =>
Fx.mergeAll(
Fx.succeed(`${file}:opened`),
Fx.at(`${file}:stored`, "1 second"),
);
const uploads = attachments.pipe(Fx.flatMapConcurrently(upload, 2));
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.
Two permits admit a and b; c waits until a permit becomes available. The tightly spaced input
lanes depict concurrent source deliveries that can wait for admission; a sequential producer cannot
issue its next input while its current delivery is still waiting. The bound limits active
work, not all retained inputs. A large waiting population still consumes memory. The concurrency
argument must be a positive safe integer; an invalid value fails with Cause.IllegalArgumentError.
Preserve every revision in order
A protocol that must apply revision a before b needs sequential admission:
import { Fx } from "@typed/fx";
const revisions = Fx.fromIterable(["a", "b", "c"]);
const save = (revision: string) =>
Fx.mergeAll(
Fx.succeed(`${revision}:accepted`),
Fx.at(`${revision}:stored`, "20 millis"),
);
const auditTrail = revisions.pipe(Fx.concatMap(save));
concatMapTick 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.
concatMap waits for each inner to complete before starting the next. Both its progress and stored
result arrive before the next revision starts. An infinite first inner prevents every later revision
from starting; a rapidly growing input backlog increases latency even with only one active job.
Use this because every command matters, not as a generic solution to concurrency.
Replace an obsolete preview
A revised document makes the previous preview irrelevant. switchMap stops future work from that
old input while preserving anything it already emitted:
import { Fx } from "@typed/fx";
const drafts = Fx.mergeAll(
Fx.at("initial", "0 millis"),
Fx.at("revised", "5 millis"),
);
const preview = (draft: string) =>
Fx.mergeAll(
Fx.at(`${draft}:started`, "1 millis"),
Fx.at(`${draft}:ready`, "20 millis"),
);
const previews = drafts.pipe(Fx.switchMap(preview));
switchMapTick 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.
b causes the x in a’s lane. a1 remains because it arrived before the switch; a2 never arrives.
Switching waits for interruption and finalizers before starting the replacement, so a slow finalizer
can delay the next raised start chevron beyond this idealized timeline.
This is a good contract for previews and current search results. It is not a rollback mechanism: interrupting local work does not undo a command already accepted by a server. Connect cancellation to the foreign API in the source adapter, and handle write idempotency or revision checks in the server protocol.
Ignore repeated submit clicks while the command is active
If an in-flight submit already represents the user’s intent, ignore busy arrivals:
import { Fx } from "@typed/fx";
const submits = Fx.mergeAll(
Fx.at("first", "0 millis"),
Fx.at("ignored", "5 millis"),
Fx.at("later", "30 millis"),
);
const submit = (command: string) => Fx.at(`saved:${command}`, "20 millis");
const accepted = submits.pipe(Fx.exhaustMap(submit));
exhaustMapTick 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.
There is no b inner because it never starts. exhaustMap does not keep it for later. When a
finishes, a later c may start. This prevents repeated local submissions while busy, but it does not
provide a global exactly-once guarantee across retries, devices, or server responses.
Finish the current save, then save only the newest snapshot
Some autosave protocols require nonoverlapping writes but do not need every intermediate snapshot:
import { Fx } from "@typed/fx";
const versions = Fx.mergeAll(
Fx.at("v1", "0 millis"),
Fx.at("v2", "5 millis"),
Fx.at("v3", "10 millis"),
);
const index = (version: string) => Fx.at(`indexed:${version}`, "20 millis");
const indexed = versions.pipe(Fx.exhaustLatestMap(index));
exhaustLatestMapTick 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.
b first becomes pending, then c replaces it. exhaustLatestMap finishes a and then starts c.
It never cancels the active write. This is appropriate only when intermediate snapshots are
replaceable; dropping an intermediate command that changes meaning is a different feature.
One-result jobs and related policies
When the inner job is an Effect, Fx.concatMapEffect(save) is shorthand for
Fx.concatMap((revision) => Fx.fromEffect(save(revision))). The other *Effect variants retain
their named admission policy too; each admitted callback produces at most one successful result.
Branch selection with if and choosing the first emitting source with race are separate decisions.
Use the operator atlas for those operators and convenience variants.
All policies combine outer and inner error/service channels and require a Scope owning admitted and waiting work. Put request recovery inside the mapper when later input should survive that failure; errors and recovery works through that placement. Services and lifetime gives these runs the owner that ends them when the editor closes.