Browse documentation

Fx

Select values and bound cardinality

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

Selection decides which pushed values reach a consumer and when its subscription stops. The examples below demonstrate three independent choices: which values pass, how many accepted values to keep, and whether to include the value that ends the run.

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.”

Choose which values pass

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.

filterEffect makes the same admission decision with an Effect. Returning false omits one value; failure reports a Cause instead. On a concurrent producer, Effectful checks may finish out of input order. See concurrency policies when order matters.

Operator order changes the count. For blank, connected, indexing, filtering blanks before Fx.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.”

Select the useful window

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

const window = 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.

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 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 the matching value before closing. Use it when the terminal record is part of the result, rather than a control signal to discard.

The related prefix operators differ at the boundary:

OperatorValues forwarded
skipWhile / dropWhileEverything from the first false predicate result onward.
dropUntilThe first matching value and everything afterward.
takeWhileValues before the first false predicate result.

Effectful variants resolve their checks before applying the same boundary 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. The operator atlas compares these variants in detail.

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.

An unopened since gate still runs its event source; an absent until signal cannot stop an infinite source. Use a timeout when there is an actual time limit.

For a window with both a start and a stop, see during in the operator atlas. For one optional answer instead of a bounded Fx, continue with Fx.first.