variable / @typed/fx

Subject.replay

Shares an Fx and replays up to the last capacity successes or failures to new subscribers.

Package version
2.0.0-beta.7
Category
Sharing sources
Since
1.0.0
Member of
Subject

Import

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

Access Subject.replay through the imported Subject 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 replay: {
    (capacity: number): <A, E, R>(fx: Fx.Fx<A, E, R>) => Fx.Fx<A, E | Cause.IllegalArgumentError, R | Scope.Scope>;
    <A, E, R>(fx: Fx.Fx<A, E, R>, capacity: number): Fx.Fx<A, E | Cause.IllegalArgumentError, R | Scope.Scope>;
};

Why

replay makes retention an explicit caller-selected policy rather than an implicit property of every shared stream. Capacity 0 is equivalent to multicast retention; capacity 1 has hold semantics; larger values replay the retained window from oldest to newest.

Errors and interruption

A non-integer, negative, or larger-than-32-bit capacity produces an Fx that fails with Cause.IllegalArgumentError without starting the source. Valid buffers retain Exit values, so typed failures keep their original causes and order. The last subscriber interrupts the source and clears the buffer; source services remain in the return type.

Ownership and lifetime

Scope owns each subscription and the final subscriber owns shutdown of the active source session. Buffer memory belongs to that session and is cleared when it is interrupted.

Join-during-publication behavior

Each Exit enters the replay buffer before it enters the serialized publication queue. A subscriber joining while an earlier publication is still draining can replay a later queued exit and then receive the same exit again when normal delivery reaches it. Replay guarantees retained order, not exactly-once delivery across a concurrent subscribe/publish race.

Examples

import { Fx } from "@typed/fx"
import * as Subject from "@typed/fx/Subject"

const recent = Fx.fromIterable([1, 2, 3]).pipe(Subject.replay(2))

Other public imports

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

Source