Browse documentation

Fx

Subject: publish events to many consumers

Connect independently owned producers and consumers through one scoped, typed publication boundary.

An invoice save can notify an activity panel and a banner without making either consumer part of the save workflow. A Subject is the scoped publication boundary: producers publish values or Causes; consumers subscribe as Fx. Start with Consuming Fx when active observation is unfamiliar.

Decide what a late observer should receive

An invoice-saved event is an occurrence, so a late banner normally does not replay it. A selected invoice is current state, so use RefSubject for that capability instead. Subject.make(0) is the default no-replay policy; capacity one or more replays that many publications in FIFO order.

Fx timelinea late Subject subscriber sees only future events without replay
OperatorSubject.make(0)
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.

Subscribe before publishing and give both sides an owner

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

const program = Effect.scoped(Effect.gen(function* () {
  const events = yield* Subject.make<string>(0)

  const activity = yield* Fx.collectAllFork(Fx.take(events, 2))
  const notifications = yield* Fx.collectAllFork(Fx.take(events, 2))
  yield* Effect.sleep(0)

  yield* events.onSuccess("saved")
  yield* events.onSuccess("published")

  return {
    activity: yield* Fiber.join(activity),
    notifications: yield* Fiber.join(notifications),
  }
}))

Both observers receive ["saved", "published"]. Effect.sleep(0) lets the forked consumers run before publishing; zero replay cannot recover a publication made before subscription. make requires Scope; closing its owner releases subscriptions and replay.

Publishing a failure does not close the Subject

onSuccess snapshots current subscribers and serializes concurrent or reentrant publications in FIFO order. A new subscriber sees later publications; replay can race an already-queued live delivery, so it is not an exactly-once protocol. onFailure publishes a Cause but does not close the Subject; an individual failing consumer may still end itself. Model recoverable availability as ordinary data when a consumer must remain live.

import { Cause, Data, Effect, Ref } from "effect"
import { Sink } from "@typed/fx"
import * as Subject from "@typed/fx/Subject"

class AuditUnavailable extends Data.TaggedError("AuditUnavailable")<{}> {}

const program = Effect.scoped(Effect.gen(function* () {
  const events = yield* Subject.make<string, AuditUnavailable>()

  const values = yield* Ref.make<ReadonlyArray<string>>([])
  const failures = yield* Ref.make(0)
  const sink = Sink.make<string, AuditUnavailable>(
    () => Ref.update(failures, (count) => count + 1),
    (value) => Ref.update(values, (all) => [...all, value]),
  )

  yield* Effect.forkScoped(events.run(sink))
  yield* Effect.sleep(0)

  yield* events.onFailure(Cause.fail(new AuditUnavailable()))
  yield* events.onSuccess("saved")

  return { failures: yield* Ref.get(failures), values: yield* Ref.get(values) }
}))

await Effect.runPromise(program)

Provide a named publication boundary

Subject.Service<Self, A, E>()(id) exposes both sides of one event channel: the class is an Fx for observation and a Sink for publication. Its make(replay) constructs the Subject in a Layer; it does not wrap an existing source.

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

class Saved extends Subject.Service<Saved, string>()("docs/Saved") {}

const SavedLive = Saved.make(0)
const receive = Fx.collectAllFork(Fx.take(Saved, 2))

const program = Effect.gen(function* () {
  const received = yield* receive

  yield* Effect.sleep(0)
  yield* Saved.onSuccess("invoice-42")
  yield* Saved.onSuccess("invoice-43")

  return yield* Fiber.join(received)
}).pipe(Effect.provide(SavedLive), Effect.scoped)

const result = await Effect.runPromise(program)
// ["invoice-42", "invoice-43"]

Provide the Layer around publishers and subscribers together so they resolve the same Subject. Saved.onFailure(cause) publishes a failure through the same channel. Saved.subscriberCount and Saved.interrupt delegate to that instance; Saved.service retrieves the underlying Subject when an integration needs it. Interrupting the shared Subject affects its subscribers, so that operation belongs with its owner.

A fresh Saved.make(0) gives a test its own event channel. A replay capacity of one changes late subscription behavior; it is not a substitute for establishing readiness in a zero-replay test. The Layer exposes invalid replay configuration as an acquisition error, separate from the E failures published later. Keep its Scope open for the whole publication journey.

Publishing events versus sharing a source

Subject.Service supplies a publication capability through a Layer when independently assembled features need one. To share an existing producer instead, see sharing one source execution. The distinction is who produces the events: callers publish into a Subject; a shared wrapper runs its source while subscribers need it.

Check both observers receive each publication and that a late observer receives no old events with capacity zero. Continue with dynamic producers when observation must choose or acquire the source.