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.
Tick 0. Early subscriber: start.
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.
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.