Browse documentation

Fx

Composing Fx

Coordinate independent producers, then recognize when a value starts work of its own.

Imagine a search screen. Local edits and server notifications feed its activity log. The current query and category define its search input. A submit button should use that input without submitting again whenever the category changes. Each relationship calls for different composition.

Transforming Fx changed one value at a time. Here the question is which independent producers belong together and which of them may trigger output. Starting a request for that output is a later decision about competing work.

NeedUsePrimary contract
Receive every event from several sourcesmergeAllForward arrivals without pairing them.
Run one phase after anotherconcatStart the next source after the previous completes.
Combine the latest valuesstructWait for every input, then update the combined value.
Let one source trigger a snapshot of anotherwithLatestFromContext updates alone do not trigger output.
Pair values by their positionzipMatch corresponding emissions rather than latest values.

Merge events when each occurrence matters

Local and server activity are peers. Both should appear when they arrive:

import { Effect } from "effect";
import { Fx } from "@typed/fx";

const local = Fx.at("saved locally", "5 millis");
const server = Fx.at("saved on the server", "1 millis");

const activity = Fx.merge(local, server);

const messages = await Effect.runPromise(Fx.collectAll(activity));

// ["saved on the server", "saved locally"]
Fx timelinemerge preserves the timing of independent producersmergemergeAllmergeLeftmergeRight
Operatormerge / mergeAll / mergeLeft / mergeRight
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.

Read each event down into its output slot: neither lane waits for the other. mergeAll generalizes to several lanes. mergeLeft and mergeRight run both but expose only their named side. Completion waits for both, so a silent live peer can keep the result open after the other finishes. Interrupting the observation stops both.

Show cached output before starting the live phase

A cache-first feed has a different promise: finish the snapshot, then subscribe to live updates.

import { Effect } from "effect";
import { Fx } from "@typed/fx";

const cached = Fx.fromIterable(["snapshot: 1", "snapshot: 2"]);
const live = Fx.fromIterable(["update: 3"]);

const feed = Fx.concat(cached, live);

const values = await Effect.runPromise(Fx.collectAll(feed));

// ["snapshot: 1", "snapshot: 2", "update: 3"]
Fx timelineconcat and continueWith start the next lane after the first endsconcatcontinueWith
Operatorconcat / continueWith
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 live lane’s start chevron follows the cache lane’s return bar. An infinite cache lane would prevent live subscription entirely. continueWith chooses the next producer at that boundary; concat already has it. For constant status markers, framing operators express the same sequence:

Fx timelineappend, prepend, and delimit frame one producerappendprependdelimit
Operatorappend / prepend / delimit
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.

prepend emits start before the source, append emits end after successful source completion, and delimit does both. An appended value is a normal event, not a finalizer that is guaranteed on failure or interruption. Keep cleanup in the source’s scoped resource contract.

Combine current query and category

The next request needs a latest value from every input. It must change when either input changes:

import { Effect } from "effect";
import { Fx } from "@typed/fx";

const queries = Fx.fromIterable(["effect", "effect v4"]);
const filters = Fx.fromIterable(["guides", "api"]);

const searchInput = Fx.zipLatest(queries, filters);

const states = await Effect.runPromise(Fx.collectAll(searchInput));
Fx timelinezipLatest emits after both inputs have a current valuetuplestructzipLatestzipLatestWith
Operatortuple / struct / zipLatest / zipLatestWith
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.

zipLatest first waits for both lanes. E+G means query effect with category guides; V4+API is the newer query with category api. tuple keeps positional values, struct names them, and zipLatestWith projects them. Each emits again when any lane changes after all lanes have supplied an initial value.

If nothing appears, inspect the initial-value requirement before changing the operator. A field that has not emitted can block the combination while another field changes repeatedly. Supply actual current state from its owner; a fabricated seed can accidentally trigger a request with invalid data. The finite fixture illustrates the API, while the diagram makes independently timed changes explicit.

Make the click the trigger and form data the context

When an action should occur only on a click, latest-value recombination is too eager:

Fx timelinewithLatestFrom emits only when its source changes after state is readywithLatestFromwithLatestFromWith
OperatorwithLatestFrom / withLatestFromWith
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.

withLatestFrom emits only for the left/source event after supporting state exists. source-1 arrives too early and is dropped; source-2 combines with ready. Updating state to revised does not replay a click. Keep the action unavailable until its supporting state is ready when losing an early click would violate the interaction contract.

Sampling reverses the trigger relationship: it retains source values and emits on sampler ticks.

Fx timelinesample reads the latest source value on each sampler ticksample
Operatorsample(values, sampler)
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 tick reads the latest value; changes between ticks only replace that retained value. Source completion ends the result and cancels the sampler. This is useful for periodic snapshots of a changing measurement, not for commands where every occurrence must be handled.

When order is the requirement

zip pairs inputs by position and queues unmatched values. mergeOrdered subscribes to all lanes but buffers later lanes behind earlier ones. Those policies can retain substantial work when one lane is slow or never finishes. They answer an ordering question rather than a latest-value question; consult the operator atlas when position or lane order is part of the contract.

Hand a ready input to the admission policy

Once query and category form a useful input, a request is a separate job-admission decision. Continue with higher-order policies for overlap, ordering, replacement, and busy-input behavior. A write that must finish is a different promise from an obsolete search read.