variable / @typed/fx/Fx

race

Runs two streams concurrently until one emits, then mirrors the winner and interrupts the other.

A failure or completion from one side before the other emits does not win unless every side ends without emitting. After a winner is chosen, that stream’s later failures are propagated.

Package version
2.0.0-beta.13
Category
Concurrent work
Since
1.0.0

Follow one execution

Runs two streams concurrently until one emits, then mirrors the winner and interrupts the other. A failure or completion from one side before the other emits does not win unless every side ends without emitting. After a winner is chosen, that stream’s later failures are propagated.

Fx timelinerace cancels slow once fast emits firstrace
Operatorrace(slow, fast)
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.

Before a winner, a non-interruption failure is remembered but does not win; it is reported only if both inputs end without a value. After selection, winner failures are forwarded. Both environments remain required. The observing fiber owns both child fibers and interruption cancels the race and finalizers.

Source implementation · Learn the surrounding model

Compare all Fx timelines →

Import

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

Signatures

export declare const race: {
    <AR, ER, RR>(that: Fx<AR, ER, RR>): <AL, EL, RL>(self: Fx<AL, EL, RL>) => Fx<AL | AR, EL | ER, RL | RR>;
    <AL, EL, RL, AR, ER, RR>(self: Fx<AL, EL, RL>, that: Fx<AR, ER, RR>): Fx<AL | AR, EL | ER, RL | RR>;
};

Why

race selects a live producer by its first useful value rather than by setup, completion, or a fast failure. This is suitable for redundant sources where a producer that ends silently should not prevent another from becoming useful.

Selection, ordering, and cardinality

Both inputs start concurrently. The first emitted value atomically selects its producer; that value and all later values from the winner are forwarded in order. The loser emits nothing after selection and is interrupted. There is no buffering or replay.

Ownership and lifetime

Before a winner, a non-interruption failure is remembered but does not win; it is reported only if both inputs end without a value. After selection, winner failures are forwarded. Both environments remain required. The observing fiber owns both child fibers and interruption cancels the race and finalizers.

Examples

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

const response = Fx.race(
  Fx.ensuring(Fx.at("slow", "50 millis"), Effect.log("slow closed")),
  Fx.ensuring(Fx.at("fast", "5 millis"), Effect.log("fast closed"))
)
Effect.runPromise(Fx.collectAll(response)).then(console.log)
// "slow closed" proves loser cleanup; the winning finalizer also runs
// resolves ["fast"]

Other public imports

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

Source