function / @typed/fx

Fx.mergeOrdered

Runs multiple Fx streams concurrently while draining their values in argument order.

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

Follow one execution

Runs multiple Fx streams concurrently while draining their values in argument order.

Fx timelinemergeOrdered buffers a faster later lane behind an earlier lanemergeOrdered
OperatormergeOrdered(first, second)
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.

Non-interruption failures are forwarded. Interrupt-only causes mark that input ended so they cannot deadlock later buffers. All input services remain typed. The observing fiber owns all runs and buffers; completion waits for all inputs, while interruption discards buffers and interrupts remaining resource scopes.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

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

Access Fx.mergeOrdered 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 function mergeOrdered<FX extends ReadonlyArray<Fx<any, any, any>>>(...fx: FX): Fx<Fx.Success<FX[number]>, Fx.Error<FX[number]>, Fx.Services<FX[number]>>;

Why

mergeOrdered overlaps production latency while preserving the same output grouping as sequential concatenation. It is useful only when that ordering is worth retaining later results in memory.

Concurrency, ordering, and buffering

Every input starts immediately. Input zero forwards as it produces. Values from each later input are buffered until every earlier input ends, then drained in that input’s emission order before the next buffer is released. Buffering is unbounded: a fast or infinite later input can retain arbitrary values while an earlier input remains open.

Ownership and lifetime

Non-interruption failures are forwarded. Interrupt-only causes mark that input ended so they cannot deadlock later buffers. All input services remain typed. The observing fiber owns all runs and buffers; completion waits for all inputs, while interruption discards buffers and interrupts remaining resource scopes.

Examples

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

const ordered = Fx.mergeOrdered(
  Fx.at("slow first", "20 millis"),
  Fx.at("fast second", "1 millis")
)
Effect.runPromise(Fx.collectAll(ordered)).then(console.log)
// ["slow first", "fast second"]: the second value waits in its buffer

Other public imports

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

Source