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"]
mergemergeAllmergeLeftmergeRightTick 0. Local: local-1; Merge: local-1; MergeLeft: local-1.
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"]
concatcontinueWithTick 0. Cached: cached; Output: cached.
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:
appendprependdelimitTick 0. Prepend: start; Delimit: 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.
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));
tuplestructzipLatestzipLatestWithTick 0. Query: effect.
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:
withLatestFromwithLatestFromWithTick 0. Source: source-1.
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.
sampleTick 0. Values: value-1.
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:
zipzipWithzipLeftzipRightTick 0. Left: left-1.
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.
mergeOrderedTick 0. Second: 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.