Browse documentation

Fx

Derive transitions and bounded batches

Carry only the local history a transform needs, then expose transitions, changes, and groups explicitly.

A shipment import page needs a running balance, numbered progress messages, meaningful status transitions, and small batches for persistence. One input alone cannot answer those questions. Each needs a different piece of history—and retaining the entire import would be unnecessary.

After Transforming Fx, choose the smallest state that answers the page’s question. Every accumulator below belongs to one subscription. A second run starts fresh; none of these operators creates shared writable application state.

Display the initial balance and each adjustment

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

const balances = Fx.fromIterable([12, -4, 7]).pipe(
  Fx.scan(100, (balance, adjustment) => balance + adjustment),
);

const values = await Effect.runPromise(Fx.collectAll(balances));
// [100, 112, 108, 115]
Fx timelinescan emits its seed and every accumulated valuescan
Operatorscan(100, add)
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.

scan first emits its seed 100. The adjustment 12 produces 112, -4 produces 108, and 7 produces 115. This seed matters: the page can show a balance before any input exists. A test that asserts only 115 misses the displayed history.

If computing the next balance needs an Effect, scanEffect exposes the same accumulated value after its reducer completes:

Fx timelinescanEffect emits each accumulated value when its reducer Effect resolvesscanEffect
OperatorscanEffect(100, oneTurnAdd)
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 output moves one logical turn after each input because the illustrated reducer takes one turn. The initial seed still appears first. Effectful reducers add their errors and required services; failed computation is not a new balance.

Produce a label while keeping the counter private

A progress label needs a position but should not expose that counter as its whole output:

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

const labels = Fx.fromIterable(["received", "packed", "shipped"]).pipe(
  Fx.loop(1, (position, event) => [`${position}. ${event}`, position + 1] as const),
);

const values = await Effect.runPromise(Fx.collectAll(labels));
// ["1. received", "2. packed", "3. shipped"]
Fx timelineloop separates its private state from its one output per eventloop
Operatorloop(position, label)
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.

loop returns [output, nextState]. For received, state 1 produces 1.received and stores 2; for packed, it produces 2.packed and stores 3. Unlike scan, it emits nothing before the first input. The accumulator lane is explanatory private state, not another subscribed producer.

Fx timelineloopEffect emits after each one-turn state transition resolvesloopEffect
OperatorloopEffect(position, oneTurnLabel)
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.

loopEffect separates the same two values after an Effect resolves. Do not assume every stateful Effect operator serializes concurrent deliveries: the producer can overlap callback Effects. Use a serialized input boundary when every transition must see the previous completed state.

Advance state even when a message is omitted

A progress display may deliberately report every other record while still counting all records. filterMapLoop returns [Option<output>, nextState]; None suppresses output but stores next state:

Fx timelinefilterMapLoop can update state without emitting a valuefilterMapLoop
OperatorfilterMapLoop(0, everyOther)
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 advances the position without producing a label, so c is labeled 2:c, not 1:c. Dropping b before an ordinary loop would be a different count.

Fx timelinefilterMapLoopEffect makes each zero-or-one decision after its Effect resolvesfilterMapLoopEffect
OperatorfilterMapLoopEffect(0, oneTurnEveryOther)
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 Effect variant makes this zero-or-one decision after asynchronous work. It also does not serialize concurrent input automatically. If multiple independent parts of the page must read and update one current count, move that responsibility to RefSubject rather than observing the same loop twice.

Highlight transitions rather than repeated reports

The import can report received several times without changing its status:

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

const transitions = Fx.fromIterable(["received", "received", "packed", "packed", "shipped"]).pipe(
  Fx.skipRepeats,
  Fx.pairwise,
);

const values = await Effect.runPromise(Fx.collectAll(transitions));
// [["received", "packed"], ["packed", "shipped"]]
Fx timelineskipRepeats removes only adjacent equivalentsskipRepeatsskipRepeatsWith
OperatorskipRepeats / skipRepeatsWith(Eq)
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.

skipRepeats compares with the last emitted value. It drops adjacent equivalents, not every value seen previously: received → packed → received still emits all three. For records, use skipRepeatsWith with the fields whose changes matter to the page. Ignoring revision data can hide real updates; comparing fresh object identity can expose meaningless repeats.

Fx timelinechangesWithEffect waits for each equivalence check before deciding the next outputchangesWithEffect
OperatorchangesWithEffect(sameStatus)
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.

changesWithEffect performs that equivalence through an Effect and serializes its comparisons. That specific guarantee is useful when comparison needs a service; it is not a guarantee shared by all Effectful state transforms.

Fx timelinepairwise waits for a prior value, then emits adjacent transitionspairwise
Operatorpairwise
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.

pairwise waits for two accepted values, then emits [previous, current]. Filtering repeated status before pairing yields received → packed and packed → shipped. Pairing raw reports first would create transitions containing duplicate statuses. This is why normalization and equivalence belong before the transition the page highlights.

Flush records without retaining the full import

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

const writes = Fx.fromIterable(["a", "b", "c", "d", "e"]).pipe(Fx.grouped(2));

const batches = await Effect.runPromise(Fx.collectAll(writes));
// [["a", "b"], ["c", "d"], ["e"]]
Fx timelinegrouped emits full batches and flushes the final partial batchgrouped
Operatorgrouped(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.

grouped(2) emits [a,b], then [c,d], then flushes [e] at normal completion. A batch bound must be a positive safe integer. Test an empty input, an exact multiple of the bound, and one extra record; the partial final batch is part of the contract, not an exceptional leftover.

For an open source, normal completion may be far away. Bound waiting time as well as count:

Fx timelinegroupedWithin flushes when its timer wins and again at source completiongroupedWithin
OperatorgroupedWithin(3, 2 turns)
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.

groupedWithin flushes a when the timer wins and c when the source ends. The retained aggregation buffer is one batch. That does not bound a downstream backlog of slow writes: use an explicit work policy and distinguish buffered records from queued persistence jobs. The timer requires a scoped owner.

Adapt repeated failure reports only at a cause boundary

Most imports use ordinary value state above and typed recovery. A lower-level consumer may instead need to transform delivered Causes with private state. The following diagrams show terminal-source examples; a Subject can deliver Causes repeatedly without permanently closing itself.

Fx timelineloopCause rewrites a terminal cause after passing earlier values throughloopCause
OperatorloopCause(0, prefix)
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.

loopCause passes successes through and transforms a Cause together with its next private state.

Fx timelineloopCauseEffect forwards its transformed terminal cause when its Effect resolvesloopCauseEffect
OperatorloopCauseEffect(0, oneTurnPrefix)
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.

loopCauseEffect waits for its Effectful transformation before forwarding the Cause.

Fx timelinefilterMapLoopCause can suppress a terminal causefilterMapLoopCause
OperatorfilterMapLoopCause(0, suppress)
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.

filterMapLoopCause can choose None, suppressing that Cause; in this terminal example the run then completes normally. That is an error policy, not a harmless formatting change.

Fx timelinefilterMapLoopCauseEffect completes only after its one-turn suppression decisionfilterMapLoopCauseEffect
OperatorfilterMapLoopCauseEffect(0, oneTurnSuppress)
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 Effect variant delays that decision and does not serialize concurrent Cause delivery. Use these operations only when a boundary truly needs stateful failure handling; ordinary progress state should not encode failures as artificial counter updates.

The page now has four deliberate histories: a seeded balance, a private position, one previous status, and one bounded batch. Check those independently when behavior diverges. A missing first transition may simply mean pairwise has only one value; a missing final batch may mean the source never completed. Time and rate adds explicit clock boundaries, and Subject explains publication state versus current readable state.