A stateful transform retains history for one subscription: an accumulated value, a counter, the previous value, or a bounded batch. A second subscription starts fresh. These operators do not create shared writable application state.
After Transforming Fx, choose the smallest history needed for each output. The independent examples below show what is retained and when it is emitted.
Emit the accumulated value, including its seed
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]
scanTick 0. Output: 100.
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. The seed provides an initial
value before any input exists; collecting only the final 115 would discard that history.
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"]
loopTick 0. Events: received; Accumulator: 1; Labels: 1.received.
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.
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:
filterMapLoopTick 0. Input: a; Output: 0:a.
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.
scanEffect, loopEffect, and filterMapLoopEffect compute their transitions with Effects.
Their outputs arrive after the transition resolves, and their errors and service requirements
become part of the Fx. They do not serialize concurrent input automatically: use a serialized
work policy when every transition must see the previous
completed state. scanEffect still emits its seed first.
For several consumers that must read and update one current value, use RefSubject rather than subscribing to the same loop twice.
Highlight transitions rather than repeated reports
Repeated reports do not necessarily represent a change:
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"]]
skipRepeatsskipRepeatsWithTick 0. Input: received; Output: received.
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 consumer. Ignoring revision data can hide
real updates; comparing fresh object identity can expose meaningless repeats.
changesWithEffectTick 0. Input: received; Output: received.
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.
pairwiseTick 0. Input: received.
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. Put deduplication before pairing when only
actual changes should produce a transition.
Retain one bounded batch
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"]]
groupedTick 0. Input: a.
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:
groupedWithinTick 0. Input: a.
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.
For time-based boundaries, continue with Time and rate. The operator atlas covers the full Effect and Cause variants of these stateful transforms.