Skip to content

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.

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:

  1. Queues carry an error type E. A queue isn’t just a buffer anymore — it can fail, and every taker observes that failure through take’s error channel. This is what makes streams-over-queues work without wrapper types.
  2. Completion is explicit. You signal “no more messages” or “everything after this is broken” by completing the queue itself.
import { Effect, Queue } from "effect"
const q1 = yield* Queue.bounded<string>(100) // capacity 100, suspend when full
const q2 = yield* Queue.unbounded<number>() // infinite
const q3 = yield* Queue.dropping<number>(50) // drop NEW messages when full
const 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"
})
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.

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.)

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.

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 not
const metrics = yield* Queue.dropping<MetricSample>(1024)
// order events: losing any is unacceptable — slow the publisher instead
const 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.

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:

  • forkScoped ties each worker to the surrounding scope. If the caller is interrupted mid-production, workers die with it — no orphaned pool.
  • The E | Cause.Done convention: ending a queue fails subsequent takes with a Done marker. Type it into your queue’s error channel and match on it explicitly (Effect.result + Result.isFailure + Cause.isDone), as above — or let Stream.fromQueue do exactly this loop for you and exclude Done from 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.end before joining guarantees workers see termination after draining everything already offered — the suspend strategy makes the final end block until space exists if producers outran consumers.

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 one
const drained = yield* PubSub.takeAll(sub) // non-empty batch

Each 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))

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, not PubSub. The implementation can swap to Redis streams or Kafka later without touching call sites.
  • subscribe as a Stream means 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 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.)

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.