An Fx can require a service and acquire a resource during observation. Providing the service satisfies its dependency; the subscription’s scope owns what it acquires. This guide shows where to provide producer and observer services, and how stopping the observer releases the resource.
Begin with dynamic producers and Consuming Fx.
Expose a source with Fx.Service
Fx.Service<Self, A, E>()(id) names an observable capability. Self is the dependency
consumers require; A and E are the source’s value and failure types. The class itself is
an Fx, so consumers compose it before an implementation is provided.
import { Effect } from "effect"
import { Fx } from "@typed/fx"
class Quotes extends Fx.Service<Quotes, number>()("docs/Quotes") {}
const QuotesLive = Quotes.make(Fx.fromIterable([100, 102]))
const prices = Fx.collectAll(Fx.map(Quotes, (cents) => cents / 100))
const result = await Effect.runPromise(prices.pipe(Effect.provide(QuotesLive)))
// [1, 1.02]
const QuotesTest = Quotes.make(Fx.succeed(250))
const testResult = await Effect.runPromise(prices.pipe(Effect.provide(QuotesTest)))
// [2.5]
Quotes.make(source) returns a Layer. It also accepts an Effect that constructs the source;
that construction runs during Layer acquisition. Dependencies needed by the implementation
are captured there, while a downstream observer still supplies its own dependencies.
Use Quotes.service when you need the underlying Context key or actual source instance.
Providing one source value does not share its execution: two observations can still run it twice. Choose an explicit sharing operator when they must share one connection. The named facade adds neither replay nor state, and declaring the class starts no work.
Give the monitor an explicit acquisition and shutdown path
A quote source needs MarketFeed; its consumer needs PriceAudit. Keep those requirements
separate so the application can provide each implementation:
import { Context, Effect, Fiber, Layer } from "effect";
import { Fx } from "@typed/fx";
interface Quote {
readonly symbol: string;
readonly cents: number;
}
class MarketFeed extends Context.Service<
MarketFeed,
{
readonly open: Effect.Effect<{
readonly quotes: Fx.Fx<Quote>;
readonly close: Effect.Effect<void>;
}>;
}
>()("app/MarketFeed") {}
class PriceAudit extends Context.Service<
PriceAudit,
{ readonly write: (quote: Quote) => Effect.Effect<void> }
>()("app/PriceAudit") {}
const quotes: Fx.Fx<Quote, never, MarketFeed> = Fx.genScoped(function* () {
const feed = yield* MarketFeed;
const socket = yield* Effect.acquireRelease(feed.open, (socket) => socket.close);
return socket.quotes;
});
// A periodic source stands in for a live connection.
const MarketFeedLive = Layer.succeed(MarketFeed, {
open: Effect.succeed({
quotes: Fx.periodic("1 second").pipe(Fx.map(() => ({ symbol: "TYPED", cents: 12_345 }))),
close: Effect.log("market socket closed"),
}),
});
const PriceAuditLive = Layer.succeed(PriceAudit, {
write: (quote: Quote) => Effect.log(`${quote.symbol}: ${quote.cents}`),
});
const observeQuotes = Fx.observe(
quotes.pipe(Fx.provide(MarketFeedLive)),
Effect.fn(function* (quote: Quote) {
const audit = yield* PriceAudit;
yield* audit.write(quote);
}),
).pipe(Effect.provide(PriceAuditLive));
// The host owns this root Fiber and interrupts it during shutdown.
const monitorFiber = Effect.runFork(observeQuotes);
const stopMarketMonitor = () => Effect.runPromise(Fiber.interrupt(monitorFiber));
Running observeQuotes provides MarketFeed, opens its handle, and observes quotes. The downstream
callback separately requires PriceAudit; its Layer supplies the destination. The service channels
remain visible until those providers are installed. stopMarketMonitor interrupts the root Fiber,
which closes the source scope and runs socket.close.
A service instance is not necessarily its resource. One MarketFeed service can open multiple
connections; providing it does not automatically share the Fx. Conversely, an already-open resource
may have an application owner that outlives this particular monitor.
genScoped keeps acquisition alive through observation and waits for cleanup before completing.
Dynamic producers
explains that scope placement.
Choose whether the provider builds or reuses the service
provideTick 0. Service Layer: 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.
provide builds the Layer for this subscription,
then releases that Layer’s Scope when the run ends. The Layer’s own errors and dependencies remain
part of the resulting type contract. Supplying a Layer is acquisition, not merely a cast removing R.
provideContext and provideService reuse existing instances; the caller keeps ownership of them.
They do not acquire or finalize those instances. provideServiceEffect runs a construction Effect
before the source starts; if it requires Scope, that requirement remains for the caller.
For application-wide setup, compose the running work with Fx.drainLayer and supply its
dependencies through Layer.provide. A single Layer graph gives the setup one owning Scope and
allows shared dependency builds. Repeated Fx.provide boundaries each build and own their own
Layers; use them when independent subscription lifetimes are intended. The
linked routing application shows the Layer form.
Trace a second observer before choosing sharing
A chart and a status badge observing ordinary quotes each open a connection. Removing the chart
releases only its connection; the badge keeps running. If both should use one connection, construct
one Subject.multicast(quotes) wrapper and expose it to both. Two independently constructed wrappers
still represent two sharing populations. The first subscriber starts the shared source, the last
leaving interrupts it, and a later subscriber starts a fresh execution.
Sharing decides the source population; Scope decides its owner. Do not fork an observer into a scope that immediately returns and assume the connection remains live. Keep the scope open for the actual feature lifetime, or use the existing application scope.
Check the shutdown promise
Count acquisitions and releases: observing ordinary quotes twice should open two handles.
Interrupt one observer and expect only its handle to close; interrupt the other and expect the
second to close. Include a source that stays silent, so cleanup cannot accidentally depend on
receiving a value.
For related boundaries, use keyed collections to retain work across collection updates, Sink services to expose an output capability, and the operator atlas to look up lifecycle hooks, tracing, and Layer-owned background runners.