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.

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(["cached: Ada", "cached: Lin"]);
const live = Fx.fromIterable(["live: Grace"]);

const people = Fx.concat(cached, live);
const values = await Effect.runPromise(Fx.collectAll(people));
// ["cached: Ada", "cached: Lin", "live: Grace"]
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.

Preserve position only when position has meaning

zip pairs first with first, then second with second. It is appropriate for corresponding protocol records, not for pairing an event feed with a rarely changing setting:

Fx timelinezip variants pair each next value until a completed lane runs out of valueszipzipWithzipLeftzipRight
Operatorzip / zipWith / zipLeft / zipRight
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.

Every output waits for its matching partner. A completed lane can still supply its queued values; pairing ends when a completion marker is reached after those values. Remaining unmatched inputs are discarded. zipWith, zipLeft, and zipRight change the output representation, not this clock. Unmatched inputs wait in queues; a fast lane can therefore retain substantial work.

mergeOrdered solves a different ordering problem: subscribe to all lanes now but expose their values in lane order.

Fx timelinemergeOrdered buffers a faster later lane behind an earlier lanemergeOrdered
OperatormergeOrdered(first, second)
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.

Here second is already available but waits for the first lane to finish. A first lane that never finishes prevents later buffered values from appearing. Choose ordered buffering because the feature requires it, not merely to make a test’s output easier to compare.

Start the current request and replace obsolete ones

Once query and category form a useful input, another producer can do the request. A search result becomes obsolete when the input changes, so select a switching policy:

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

class SearchFailed extends Data.TaggedError("SearchFailed")<{
  readonly query: string;
}> {}

interface SearchClient {
  readonly search: (query: string) => Effect.Effect<ReadonlyArray<string>, SearchFailed>;
}

const SearchClient = Context.Service<SearchClient>("docs/SearchClient");

const queries = Fx.mergeAll(Fx.at("effect", "0 millis"), Fx.at("effect v4", "5 millis"));
const results = queries.pipe(
  Fx.switchMapEffect((query) => Effect.flatMap(SearchClient, ({ search }) => search(query))),
);

const program = results.pipe(
  Fx.provideService(SearchClient, {
    search: (query) => Effect.as(Effect.sleep("20 millis"), [`result: ${query}`]),
  }),
  Fx.collectAll,
  Effect.scoped,
);

const values = await Effect.runPromise(program);
// [["result: effect v4"]]

The first request sleeps for 20 milliseconds; the revised query arrives after 5 and interrupts it. Only the revised result is delivered. Errors and the SearchClient requirement remain visible until recovered or provided. Effect.scoped gives admitted work its owner.

This last step differs from combining independent facts: each input starts work of its own. Continue with higher-order policies for flatMapConcurrently, concatMap, switchMap, exhaustMap, and exhaustLatestMap. In particular, a write that must finish is a different promise from an obsolete search read.