Browse documentation

Fx

Select values and bound cardinality

Keep, omit, gate, and stop pushed values without confusing selection with cancellation policy.

An import screen receives raw status lines. It should omit malformed records, show the next two useful messages after a banner, and include the final “complete” record before stopping. These are three decisions: admission, a counted window, and a terminal boundary. Treating them as one filter makes it easy to stop at the wrong moment.

Transforming Fx introduced zero-or-one output. Here we connect that choice to how long the producer remains subscribed. A source can be active while every value is rejected; “nothing visible” does not mean “nothing running.”

Parse useful records before counting them

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

const messages = Fx.fromIterable(["", "notice: connected", "ready", "done"]).pipe(
  Fx.filterMap((line) => {
    const separator = line.indexOf(": ");
    return separator < 0 ? Option.none() : Option.some(line.slice(separator + 2));
  }),
);
// Emits: ["connected"]

filterMap emits Some and omits None; the example extracts only the structured notice. If the original value should remain unchanged, use filter. An Effectful admission rule exposes its failures and service requirements:

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

const visible = Fx.fromIterable([
  { text: "private", allowed: false },
  { text: "queued", allowed: true },
]).pipe(Fx.filterEffect((message) => Effect.succeed(message.allowed)));
// Emits only the allowed message.

An Effect returning false omits one value. An Effect failure reports a Cause instead; it is not a negative predicate result. On a concurrent producer, Effectful checks may finish out of input order, so choose an explicit serialized boundary if record order is part of the contract.

Select the useful window

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

const page = Fx.fromIterable(["banner", "connected", "indexing", "complete", "ignored"]).pipe(
  Fx.slice({ skip: 1, take: 2 }),
);
// Emits: ["connected", "indexing"]
Fx timelineskip removes only its fixed prefixskipskipEffect
Operatorskip(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.

skip removes a fixed prefix while keeping the source live afterward.

Fx timelinetake completes after its fixed prefixtaketakeEffect
Operatortake(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.

take closes after its accepted prefix; later source values are no longer useful work.

Fx timelineslice keeps one bounded index window and then completesslicesliceEffect
Operatorslice({ skip: 1, take: 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.

slice combines both counters. In this run, banner is skipped, connected and indexing are emitted, then upstream stops. The Effect variants obtain their bounds before subscribing to the source.

Operator order changes the count. For blank, connected, indexing, filtering blanks before take(2) returns both useful messages. Taking two raw records before filtering returns only connected. Choose whether the bound means “inspect two inputs” or “show two useful outputs.” This is an event window, not server-side pagination unless the producer supplies that dataset and ordering contract.

Include or exclude the terminal record

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

const beforeComplete = Fx.fromIterable(["connected", "indexing", "complete", "ignored"]).pipe(
  Fx.takeUntil((line) => line === "complete"),
);
// Emits: ["connected", "indexing"]

const throughComplete = Fx.fromIterable(["connected", "indexing", "complete", "ignored"]).pipe(
  Fx.dropAfter((line) => line === "complete"),
);
// Emits: ["connected", "indexing", "complete"]
Fx timelineskipWhile drops a matching prefix, and dropWhile is its aliasskipWhileskipWhileEffectdropWhiledropWhileEffect
OperatorskipWhile(isBanner)
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.

skipWhile/dropWhile omit the true prefix. After the gate opens, later matching values are no longer part of that prefix.

Fx timelinedropUntil includes the boundary that opens its gatedropUntildropUntilEffect
OperatordropUntil(isConnected)
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.

dropUntil includes the value that first satisfies its predicate and all later values.

Fx timelinetakeWhile stops before its first false valuetakeWhiletakeWhileEffect
OperatortakeWhile(isInProgress)
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.

takeWhile excludes the first false value and stops.

Fx timelinetakeUntil excludes its matching sentineltakeUntiltakeUntilEffect
OperatortakeUntil(isComplete)
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.

takeUntil excludes its true sentinel. Here the completion marker is control-only.

Fx timelinedropAfter includes its matching sentineldropAfter
OperatordropAfter(isComplete)
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.

dropAfter includes that sentinel before closing. Use this for the import screen’s final visible status. Read the last occupied output slot, not merely the method name: “until” and “after” make opposite promises about that boundary value.

Effectful variants use the same boundary after resolving their checks, and add the checks’ errors and requirements. skipWhileEffect and dropUntilEffect still evaluate after their gate opens; choose a pure predicate when later service calls would be unintended work.

Let another producer open or close the window

A separate user action can own the window independently of record content:

Fx timelinesince opens when its named start signal emitssince
Operatorsince(events, 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.

since(events, start) already runs the event source, discarding values until start emits. It does not buffer draft for later. A failed start signal leaves the event source alive.

Fx timelineuntil stops when its named stop signal emitsuntil
Operatoruntil(events, stop)
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.

until(events, stop) closes when the stop lane emits. Its control value never reaches output, and its failure propagates because the signal owns stopping work.

Fx timelineduring forwards only while its named window is activeduring
Operatorduring(events, drag)
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.

during(events, starts) uses the first start value as an inner stop Fx. Read the inner up token as the end of that selected window. Signal failures propagate. For repeated drag windows carrying start coordinates, use the explicit switchMap/until composition in Time and rate.

A boolean chooses branch values rather than directly counting or gating source events:

Fx timelinewhen selects a value for each spaced conditionwhen
Operatorwhen(condition, { onTrue, onFalse })
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.

when selects the corresponding constant branch. Closely spaced conditions can replace a branch before it emits, so this is not guaranteed one output per pushed boolean.

The import is finished when its terminal rule is met, even if the underlying listener could keep producing. Its callback cleanup still needs a real subscription owner. Test an absent start signal, an absent stop signal, and a signal failure as well as the happy path. A silent gate can remain live forever; use a timeout only when the product defines a time limit. For one optional answer instead of a bounded Fx, continue with Fx.first.