variable / @typed/fx

Fx.merge

Merges two Fx streams into a single Fx that emits values from both streams concurrently. Order of emission is non-deterministic.

Completion: The merged stream completes when both input streams have completed.

Failures: Every failure Cause is delivered to the downstream Sink. Delivery does not make Fx.run fail, so merge does not itself cancel the sibling; a terminal observer may choose to.

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

Follow one execution

Both subscriptions run together; every value is forwarded and completion waits for both.

Fx timelinemerge(left, right)merge
Operatormerge(left, right)
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.

Both runs are children of the consumer and completion waits for both. Downstream interruption cancels the remaining runs. A source Cause is only sent to sink.onFailure; whether that callback ends observation is Sink policy, not an intrinsic terminal rule of merge.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

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

Access Fx.merge 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 merge: {
    <A2, E2, R2>(that: Fx<A2, E2, R2>): <A, E, R>(self: Fx<A, E, R>) => Fx<A | A2, E | E2, R | R2>;
    <A, E, R, A2, E2, R2>(self: Fx<A, E, R>, that: Fx<A2, E2, R2>): Fx<A | A2, E | E2, R | R2>;
};

Why

merge preserves both producers instead of imposing a pairing relationship. Each input value is emitted once as soon as its source pushes it, so inter-source ordering is intentionally undefined.

Ownership and lifetime

Both runs are children of the consumer and completion waits for both. Downstream interruption cancels the remaining runs. A source Cause is only sent to sink.onFailure; whether that callback ends observation is Sink policy, not an intrinsic terminal rule of merge.

Examples

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

const program = Effect.gen(function* () {
  const deliveries = yield* Ref.make<Array<string>>([])
  const sink = Sink.make(
    () => Ref.update(deliveries, (xs) => [...xs, "failure"]),
    (value: string) => Ref.update(deliveries, (xs) => [...xs, value])
  )
  yield* Fx.merge(Fx.fail("left failed"), Fx.succeed("right continued")).run(sink)
  return yield* Ref.get(deliveries) // both "failure" and "right continued"
})

Other public imports

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

Source