Queues, PubSub & Coordination
Effect v4's Queue with typed errors and completion signaling, backpressure strategies compared, a worker pool built from forkScoped and take loops, PubSub fan-out with replay, the event-bus service pattern, and Latch barriers.
Producers and consumers rarely finish at the same rate. Queues absorb the
difference; pubsubs broadcast it; latches and deferreds synchronize around it.
Effect ships all four as first-class values with real cancellation semantics —
and v4 reworked Queue substantially: it now carries an error type, splits
into producer/consumer views, and has explicit completion signaling.
Queue in v4
Section titled “Queue in v4”A queue is one object exposing two directional views:
interface Enqueue<in A, in E = never> // producers: offer, shutdown, fail...interface Dequeue<out A, out E = never> // consumers: take, poll, peek...interface Queue<in out A, in out E = never> extends Enqueue<A, E>, Dequeue<A, E> {}Two things to notice against your v3 instincts:
- Queues carry an error type
E. A queue isn’t just a buffer anymore — it can fail, and every taker observes that failure throughtake’s error channel. This is what makes streams-over-queues work without wrapper types. - Completion is explicit. You signal “no more messages” or “everything after this is broken” by completing the queue itself.
Constructors
Section titled “Constructors”import { Effect, Queue } from "effect"
const q1 = yield* Queue.bounded<string>(100) // capacity 100, suspend when fullconst q2 = yield* Queue.unbounded<number>() // infiniteconst q3 = yield* Queue.dropping<number>(50) // drop NEW messages when fullconst q4 = yield* Queue.sliding<number>(50) // drop OLDEST messages when full
// general form:const q5 = yield* Queue.make<Job, JobError>({ capacity: 100, strategy: "suspend" // "dropping" | "sliding"})Core operations
Section titled “Core operations”Producer (Enqueue) |
Consumer (Dequeue) |
|---|---|
offer(q, msg): Effect<boolean> — false if dropped |
take(q): Effect<A, E> — suspends until available |
offerAll(q, iter): Effect<Array<A>> — returns dropped items |
poll(q): Effect<Option<A>> — non-blocking take |
offerUnsafe(q, msg): boolean — sync best-effort |
peek(q): Effect<A, E> — look without consuming |
shutdown(q): Effect<boolean> — release all waiters |
takeAll(q) / takeN(n) / takeBetween(min,max) |
size, capacity |
clear(q): Effect<Array<A>> |
poll returning an Option replaces v3’s takeUpTo ceremony for the common
“grab work if there is any” case. takeAll returns a non-empty array (it
suspends until at least one message exists) — handy for batch-processing
consumers.
Completion: Done, fail, interrupt
Section titled “Completion: Done, fail, interrupt”The lifecycle operations are how queues terminate consumers cleanly:
import { Cause, Effect, Queue } from "effect"
const program = Effect.gen(function* () { const queue = yield* Queue.bounded<number, string>(10)
yield* Queue.offer(queue, 1)
// graceful: no more data coming. Takers drain remaining, then observe Done yield* Queue.end(queue)
const result = yield* Effect.flip(Queue.take(queue)) // result: Cause.Done() — take fails with the Done marker once empty
// failure: everyone downstream fails with this error yield* Queue.fail(queue, "upstream exploded")
// interruption: equivalent to shutdown semantics for waiters yield* Queue.shutdown(queue)})| Signal | What consumers see |
|---|---|
Queue.end(q) |
buffered messages first, then failure with Cause.Done |
Queue.fail(q, e) / failCause(q, c) |
remaining messages, then failure with e/c |
Queue.interrupt(q) / shutdown(q) |
pending takes/offers interrupted immediately; further ops shut down |
There is deliberately no awaitShutdown on queues — to wait for a queue’s
demise, either take until it completes, or register a finalizer. (PubSub,
below, does keep an explicit awaitShutdown.)
Batching consumers
Section titled “Batching consumers”High-throughput sinks (bulk DB inserts, batched HTTP APIs) want groups of messages per round trip. The consumer-side combinators express this without hand-rolled buffering:
import { Cause, Effect, Queue } from "effect"
const bulkIngest = Effect.gen(function* () { const queue = yield* Queue.bounded<Row, Cause.Done>(10_000)
const flushLoop = Effect.forever(Effect.gen(function* () { // wait for at least 1 message, then drain up to 500 in one shot const batch = yield* Queue.takeBetween(queue, 1, 500) yield* insertRows(batch) }))
// alternative shapes: // Queue.takeN(queue, 100) — exactly 100 (waits for a full batch) // Queue.takeAll(queue) — suspend for >=1, return everything buffered // Queue.peek(queue) — inspect head without consuming})takeBetween(queue, min, max) suspends until min messages are available,
then returns up to max — (1, 500) is the classic micro-batch. takeAll
returns everything currently buffered as a non-empty array; takeN enforces
exact batch sizes when the sink requires them; and completion propagates: if
the queue ends before the minimum arrives, the batch operation fails with the
terminal error so your flush loop exits through its normal error path.
Backpressure strategies
Section titled “Backpressure strategies”Choosing a strategy is choosing who absorbs overload. The three bounded strategies compared:
| Strategy | When full, offer… |
Message loss | Producer latency | Typical use |
|---|---|---|---|---|
"suspend" (default, bounded) |
suspends until space frees | none | grows under load | correctness-first pipelines; DB writes, job processing |
"dropping" |
returns immediately; new message discarded | newest lost | constant | metrics, sampling, telemetry where freshness beats completeness |
"sliding" |
evicts oldest, inserts new | oldest lost | constant | “latest state” feeds — price ticks, cursor positions |
Unbounded queues have no strategy because they never refuse: memory grows until the process dies. They’re appropriate when producers and consumers are logically decoupled but practically paced together, or when batch sizes are provably small. As a rule of thumb: bounded with suspend is the default; reach for the others only when you can articulate which messages deserve to die.
// telemetry: losing the newest sample is fine, blocking the hot path is notconst metrics = yield* Queue.dropping<MetricSample>(1024)
// order events: losing any is unacceptable — slow the publisher insteadconst orders = yield* Queue.bounded<OrderEvent>(4096)Note the interaction with interruption: a producer suspended in offer on a
full suspend-strategy queue is interruptible like any fiber. If its fiber dies
(or a scope closes), the pending offer unwinds — backpressure never traps a
fiber.
Worker pool: the canonical composition
Section titled “Worker pool: the canonical composition”Fork N consumers over a shared queue, feed it from anywhere, shut it down gracefully when done:
import { Cause, Effect, Fiber, Queue, Result, Schema } from "effect"
class JobError extends Schema.TaggedError<JobError>()( "JobError", { jobId: Schema.String, reason: Schema.String }) {}
interface Job { readonly id: string; readonly payload: string }
const processJob = Effect.fn("processJob")(function* (job: Job) { yield* Effect.log(`processing ${job.id}`) if (job.payload === "") { return yield* new JobError({ jobId: job.id, reason: "empty payload" }) } yield* Effect.sleep("50 millis")})
const runPool = Effect.fn("runPool")(function* (jobs: Iterable<Job>) { const queue = yield* Queue.bounded<Job, JobError | Cause.Done>(128) const deadLetters = yield* Queue.unbounded<JobError>()
// workers: loop take -> handle -> repeat; exit cleanly when queue ends const workers = yield* Effect.forEach( [0, 1, 2, 3], () => Effect.forkScoped(Effect.forever( Effect.gen(function* () { // E includes Done because we typed the queue that way — // end() surfaces as a failure value we can pattern-match const outcome = yield* Effect.result(Queue.take(queue)) if (Result.isFailure(outcome)) { const error = outcome.failure if (Cause.isDone(error)) return // graceful stop yield* Queue.offer(deadLetters, error) return } yield* processJob(outcome.success).pipe( Effect.catchAll((err) => Queue.offer(deadLetters, err)) ) }) )), { concurrency: "unbounded", discard: true } )
// produce for (const job of jobs) { yield* Queue.offer(queue, job) }
// graceful stop: end() lets workers drain buffered jobs first, // then every take observes Done and the worker loops exit yield* Queue.end(queue) yield* Fiber.joinAll(workers) yield* Queue.shutdown(deadLetters)})
await Effect.runPromise( runPool([{ id: "a", payload: "x" }, { id: "b", payload: "" }]).pipe(Effect.scoped))Points worth savoring:
forkScopedties each worker to the surrounding scope. If the caller is interrupted mid-production, workers die with it — no orphaned pool.- The
E | Cause.Doneconvention: ending a queue fails subsequent takes with aDonemarker. Type it into your queue’s error channel and match on it explicitly (Effect.result+Result.isFailure+Cause.isDone), as above — or letStream.fromQueuedo exactly this loop for you and excludeDonefrom the stream’s error type. - Errors flow through the queue’s
E, and per-job failures go to a dead-letter queue instead of killing workers — a decision made visible in types rather than buried in a try/catch. Queue.endbefore joining guarantees workers see termination after draining everything already offered — the suspend strategy makes the finalendblock until space exists if producers outran consumers.
PubSub: fan-out with replay
Section titled “PubSub: fan-out with replay”Where a queue delivers each message to exactly one consumer, a PubSub
delivers each published message to every active subscription:
import { Effect, PubSub } from "effect"
const bus = yield* PubSub.bounded<OrderEvent>({ capacity: 1024, replay: 64 })// plain: PubSub.bounded<OrderEvent>(1024)// unbounded with replay: PubSub.unbounded<OrderEvent>({ replay: 16 })
yield* PubSub.publish(bus, event) // dual: PubSub.publish(event)(bus)yield* PubSub.publishAll(bus, events)yield* PubSub.shutdown(bus)Subscription requires a scope — subscriptions register a finalizer so they can’t leak:
const sub = yield* PubSub.subscribe(bus) // Effect<Subscription<Event>, never, Scope>const next = yield* PubSub.take(sub) // suspend until this subscriber gets oneconst drained = yield* PubSub.takeAll(sub) // non-empty batchEach subscription sees messages published after it subscribed. The replay buffer changes that for late arrivals: the last N published messages are re-delivered to each new subscriber before live traffic resumes — the “catch up on what I missed” feature chat UIs and dashboards always need.
Strategies mirror queues: bounded (publishers suspend when the slowest
subscriber falls behind by more than capacity), dropping(capacity),
sliding(capacity). With multiple subscribers, bounded backpressure means
publish waits for all of them — fan-out amplifies the slowest consumer.
Streams integrate directly — Stream.fromPubSub(bus) subscribes on run
(within the stream’s scope) and emits until shutdown:
import { Effect, Stream } from "effect"
const consumer = Stream.fromPubSub(bus).pipe( Stream.filter((e) => e.kind === "paid"), Stream.runForEach((e) => Effect.log(`paid: ${e.orderId}`)))
await Effect.runPromise(Effect.scoped(consumer))The event-bus service pattern
Section titled “The event-bus service pattern”Raw pubsubs shouldn’t leak into business code. Wrap the bus in a service whose type advertises exactly two capabilities — publish, and subscribe-as-stream — and let a layer own the lifecycle:
import { Context, Effect, Layer, PubSub, Schema, Stream } from "effect"
class OrderPlaced extends Schema.TaggedError<OrderPlaced>()( "OrderPlaced", { orderId: Schema.String, totalCents: Schema.Number }) {}export type OrderEvent = OrderPlaced
export class OrderEvents extends Context.Service< OrderEvents, { readonly publishOrderEvent: (event: OrderEvent) => Effect.Effect<void> readonly subscribe: Stream.Stream<OrderEvent> }>()("OrderEvents") {}
export const OrderEventsLive = Layer.effect( OrderEvents, Effect.gen(function* () { const pubsub = yield* PubSub.bounded<OrderEvent>({ capacity: 4096, replay: 128 })
// the layer's own scope owns the pubsub: when the layer shuts down, // the bus does too — subscribers' finalizers unwind first (LIFO) yield* Effect.addFinalizer(() => PubSub.shutdown(pubsub))
return { publishOrderEvent: (event) => PubSub.publish(pubsub, event).pipe( Effect.flatMap((published) => published ? Effect.void : Effect.logWarning(`order event dropped: ${event.orderId}`) ) ), subscribe: Stream.fromPubSub(pubsub) } }))
// consumer side — a scoped stream, e.g., wired into the app's runtime:const auditTrail = Effect.gen(function* () { const events = yield* OrderEvents yield* events.subscribe.pipe( Stream.runForEach((event) => Effect.log(`audit ${event._tag}: ${event.orderId}`) ) )}).pipe(Effect.provide(OrderEventsLive))
// publisher side — anywhere in the app:const checkout = Effect.gen(function* () { const events = yield* OrderEvents yield* events.publishOrderEvent(new OrderPlaced({ orderId: "o-1", totalCents: 4200 }))})Why this shape earns its keep:
- Consumers depend on
OrderEvents, notPubSub. The implementation can swap to Redis streams or Kafka later without touching call sites. subscribeas aStreammeans consumers get filtering, batching, retries, and backpressure semantics for free, and the subscription’s scope requirement is absorbed into stream-running scopes.- Lifecycle is layered: bus created in the layer, closed by the layer, subscriptions closed by their own scopes — three lifetimes, zero manual coordination.
Latch: phase barriers
Section titled “Latch: phase barriers”Latch is a reusable gate: fibers suspend at await until the gate opens.
import { Effect, Latch } from "effect"
const raceStart = Effect.gen(function* () { const gate = yield* Latch.make() // starts closed
yield* Effect.all([ Effect.forkChild(gate.whenOpen(Effect.log("runner 1 going"))), Effect.forkChild(gate.whenOpen(Effect.log("runner 2 going"))) ])
yield* Effect.sleep("10 millis") // everyone is staged yield* gate.open // release current AND future waiters
// ...later phases: yield* gate.close // re-arm: future awaits suspend again})Three distinct operations matter:
| Operation | Effect |
|---|---|
latch.open |
opens the gate; wakes current waiters; future awaits pass immediately |
latch.close |
closes it again — latches in v4 are reusable, unlike v3’s one-shot |
latch.release |
wakes current waiters only; the gate stays closed for future ones |
latch.await |
suspend until open-or-released |
latch.whenOpen(effect) |
run effect only once the latch permits |
Use latches for orchestration (“don’t start phase 2 until setup signals”),
connection warmup gates, and tests that need deterministic interleavings.
For one-shot handoffs of values, use Deferred
(chapter 13); for counting down N units of work,
a semaphore plus latch combination covers what Go’s sync.WaitGroup does.
A wait-group-shaped helper built from verified primitives — workers signal a shared counter, the coordinator parks on the latch until the count hits zero:
import { Effect, Latch, Ref } from "effect"
const runCohort = Effect.fn("runCohort")(function* (jobs: Array<Effect.Effect<void>>) { const pending = yield* Ref.make(jobs.length) const done = yield* Latch.make()
const tracked = jobs.map((job) => job.pipe( Effect.onExit(() => Ref.modify(pending, (n) => [n - 1 === 0, n - 1]).pipe( Effect.flatMap((last) => (last ? Latch.open(done) : Effect.void)) ) ) ) )
yield* Effect.forEach(tracked, Effect.forkChild, { discard: true }) yield* done.await // parks until the last worker opens it})The interesting bit is Ref.modify returning [shouldOpen, newCount]
atomically — exactly one worker observes the transition to zero and opens the
latch, so there’s no double-open race and no polling loop. (For
transactional coordination across many counters and conditions at once,
STM in chapter 16 subsumes this pattern.)
Choosing: Queue vs PubSub vs Stream
Section titled “Choosing: Queue vs PubSub vs Stream”| Need | Reach for |
|---|---|
| Work distribution — each item processed once | Queue (+ worker forks, or Stream.fromQueue) |
| Broadcast — every listener sees every event | PubSub behind a service |
| Late subscribers must catch up (last N) | PubSub with replay |
| Composable pipeline over the data (map/filter/retry/window) | Stream sourced from a queue or pubsub |
| One-shot rendezvous between two parties | Deferred (chapter 13) |
| Phase gating / start signals | Latch |
| Limit concurrent access to a resource | Semaphore (chapter 13) |
| Multi-variable atomic updates across those structures | STM — chapter 16 |
Rule of thumb: queues move responsibility, pubsubs move information, streams express processing, and the small primitives (Deferred, Latch, Semaphore) coordinate timing. Most systems need one queue-shaped thing and one pubsub-shaped thing; resist building both out of the other.