Browse documentation

Fx

Flatten Fx with an explicit policy

Declare how inner work runs over time: overlap, wait, replace, or ignore new arrivals.

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.

NeedUsePrimary contract
Allow independent jobs to overlapflatMapConcurrentlyBound concurrent inner runs.
Process every job in orderconcatMapWait for each inner run to complete.
Replace obsolete workswitchMapInterrupt the prior inner run when new input arrives.
Ignore arrivals while busyexhaustMapFinish the active run and drop intervening inputs.
Finish active work, then use the newest inputexhaustLatestMapRetain 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));
Fx timelineflatMap runs every inner and lets their values interleaveflatMap
OperatorflatMap(load)
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));
Fx timelineflatMapConcurrently waits when every permit is occupiedflatMapConcurrently
OperatorflatMapConcurrently(load, 2)
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));
Fx timelineconcatMap finishes each inner before starting the nextconcatMap
OperatorconcatMap(save)
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));
Fx timelineswitchMap interrupts the old inner exactly when its replacement arrivesswitchMap
OperatorswitchMap(preview)
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));
Fx timelineexhaustMap ignores arrivals until the active inner completesexhaustMap
OperatorexhaustMap(submit)
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));
Fx timelineexhaustLatestMap keeps only the newest value waiting behind active workexhaustLatestMap
OperatorexhaustLatestMap(index)
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.

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.