variable / @typed/fx

Fx.mergeAll

Merges multiple Fx streams into a single Fx that emits values from all input streams concurrently.

Package version
2.0.0-beta.13
Category
Combining sources
Since
1.0.0
Member of
Fx

Follow one execution

All supplied producers run concurrently; this three-source example retains every delivery.

Fx timelinemergeAll(local, server, cache)mergeAll
OperatormergeAll(local, server, cache)
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.

Typed failures from any input are forwarded and every input environment is required. Interrupt-only causes from sibling cancellation are suppressed at the Sink boundary. The observing fiber owns all concurrent runs: completion waits for all inputs, and interruption cancels the remaining runs and their resource lifetimes.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

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

Access Fx.mergeAll through the imported Fx export. Its declaration below describes the member.

This public exposure is a re-export. Its import path is supported; declaration documentation is shared with the other public exposures below.

Signatures

export declare const mergeAll: <FX extends ReadonlyArray<Fx<any, any, any>>>(...fx: FX) => Fx<Fx.Success<FX[number]>, Fx.Error<FX[number]>, Fx.Services<FX[number]>>;

Why

mergeAll combines a known set of independent push producers without adding a source-of-sources or choosing a winner.

Concurrency, ordering, and cardinality

All inputs start concurrently. Every non-interruption value from every input is forwarded once. Each input preserves its own order, while values across inputs interleave by arrival; there is no buffering to restore argument order. An empty argument list completes without emitting.

Ownership and lifetime

Typed failures from any input are forwarded and every input environment is required. Interrupt-only causes from sibling cancellation are suppressed at the Sink boundary. The observing fiber owns all concurrent runs: completion waits for all inputs, and interruption cancels the remaining runs and their resource lifetimes.

Examples

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

const events = Fx.mergeAll(
  Fx.at("slow", "20 millis"),
  Fx.at("fast", "1 millis")
)
Effect.runPromise(Fx.collectAll(events)).then(console.log)
// ["fast", "slow"]

Other public imports

These import paths expose the same declaration. Each page retains its own public name and signature.

Source