interface / @typed/fx

Subject.Subject

A multicast boundary that is both an Fx of its publications and a Sink that accepts them. Successes and failures are pushed to every subscriber present when that publication begins.

Package version
2.0.0-beta.7
Category
Publication contracts
Since
1.0.0
Member of
Subject

Import

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

Access Subject.Subject 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 interface Subject<A, E = never, R = never> extends Fx.Fx<A, E, R | Scope.Scope>, Sink.Sink<A, E, R> {
    readonly subscriberCount: Effect.Effect<number, never, R>;
    readonly interrupt: Effect.Effect<void, never, R>;
}
export declare namespace Subject {
    interface Service<Self, Id extends string, A, E> extends Subject<A, E, Self> {
        readonly id: Id;
        readonly service: Context.Service<Self, Subject<A, E>>;
        readonly make: (replay?: number) => Layer.Layer<Self, Cause.IllegalArgumentError, Scope.Scope>;
    }
    interface Class<Self, Id extends string, A, E> extends Service<Self, Id, A, E> {
        new (): Service<Self, Id, A, E>;
    }
}

Why

Subject connects imperative or independently owned producers to Fx without changing the producer-driven direction of the work. Because the same value is also a Sink, another Fx can publish into it with run while any number of consumers subscribe through the ordinary Fx surface.

Publication order

Publications are serialized in FIFO order, including publications made concurrently or from inside a subscriber callback. The subscriber set is snapshotted once per publication: a subscriber removed during a delivery receives no later publication, and a subscriber added during a delivery begins with the next publication. A failure publication uses the error channel but does not permanently terminate the subject. A reentrant onSuccess or onFailure call from the fiber currently draining subscribers enqueues its publication and returns before that queued publication is delivered; the outer drain still delivers it in FIFO order.

Ownership and lifetime

Each call to run registers its sink in the caller’s Scope; closing that scope removes only that subscription. interrupt closes all current subscriber scopes and clears retained replay state. A subject created by make is interrupted automatically when its owning scope closes; unsafeMake leaves that responsibility with the caller.

Property: interrupt

Interrupts every current subscription and clears any retained replay values.

Property: interrupt: Why

Gives the subject owner one deterministic shutdown operation for all active consumers and retained publications.

Property: interrupt: Ownership and lifetime

Closing subscriber scopes runs their finalizers. The effect cannot fail, and the subject may be subscribed to and published through again after the interruption.

Property: subscriberCount

Samples the number of sinks currently registered with this subject.

Property: subscriberCount: Why

Exposes demand for diagnostics and demand-sensitive coordination without coupling producers to the subject implementation.

Property: subscriberCount: Ownership and lifetime

The returned Effect does not subscribe or retain anything. It reads the count when executed and requires the same services R as the subject implementation.

Examples

import { Effect, Fiber } from "effect"
import { Fx } from "@typed/fx"
import * as Subject from "@typed/fx/Subject"

const program = Effect.gen(function* () {
  const events = yield* Subject.make<string>(2)
  const collected = yield* Effect.forkScoped(Fx.collectAll(Fx.take(events, 2)))

  yield* events.onSuccess("connected")
  yield* events.onSuccess("ready")

  return yield* Fiber.join(collected)
}).pipe(Effect.scoped)

Members

Other public imports

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

Source