View Source EchoPubSub

CI hex.pm version hex.pm license

A Phoenix.PubSub adapter that distributes messages between nodes using the erlang :pg module, like the default adapter, however with the additional guarentees of "at least once" delivery.

This means that nodes can disconnect temporarily from the cluster - even for a blip as short as ~1ms - and then "catch up" when they rejoin, thanks to a buffer of messages and read cursors.

See the Docs for more information.

How it works

Phoenix.PubSub.PG2 is fire-and-forget: a broadcast reaches only the nodes connected at that instant. A blip as short as ~1ms silently drops messages for any node briefly unreachable.

EchoPubSub makes delivery at-least-once:

  • Buffer + cursors - each broadcaster keeps a ring buffer of recent messages, plus a per-node read cursor that advances only on an acked delivery.

  • Replay on reconnect - a reconnecting node is replayed exactly the messages it missed, in order.

  • Told if it fell behind - if it stayed gone long enough that those messages were overwritten in the bounded buffer, it gets {:cursor_expired, node_name} (see Usage) telling it to reload from a source of truth.

Core guarantee: either you receive every message in order, or you are told you fell behind - never a silent gap.

See how it works for diagrams, the cursor internals, and failure scenarios.

When to use it

EchoPubSub gives you at-least-once cross-node delivery without standing up a dedicated message broker. If you already run a BEAM cluster, you reuse it - no extra service to deploy, secure, monitor, or scale.

A good fit for small-to-mid projects that need reliable cross-node messaging but don't want the operational burden of Kafka / RabbitMQ / NATS:

  • Replicated in-memory caches - a missed invalidation means a node serves stale data forever. EchoPubSub replays it on reconnect, or sends {:cursor_expired, node} to trigger a reload - never a silent stale node.
  • Event logs / projections / derived state kept in sync across nodes.
  • Presence / state fan-out where a dropped update corrupts a peer's view.
  • Anything currently on plain Phoenix.PubSub that quietly breaks during network blips.

When to reach for a real broker instead: durable persistence across a full cluster restart, replay from disk / long retention, cross-language consumers, huge backlogs, or delivery to non-BEAM systems. EchoPubSub's buffer is in-memory and bounded - it closes the network-blip gap, it is not a durable log.

Caveats

At-least-once means possibly-more-than-once. If a message is delivered and handled but its ack is lost, the cursor doesn't advance and the message is re-sent

  • so a handler can see the same message twice. Make handlers tolerate duplicates: send absolute state rather than deltas, or dedupe by a per-message id. See handling duplicate deliveries for worked examples.

Usage

Note: I used LLM for typing - but ideas and decisions were mine

def deps do
  [
    {:echo_pubsub, "~> 0.1.0"}
  ]
end

Not a drop-in replacement for Phoenix.PubSub. At-least-once delivery costs more than fire-and-forget (buffering, acked cross-node calls). Keep the default PubSub for ordinary broadcasts and run EchoPubSub alongside it, using it only for cross-node data that must not be lost (replicated caches, event logs, derived state).

Add it to your supervision tree. Running EchoPubSub alongside your existing default PubSub, both children default to the same child id (that of the Phoenix.PubSub supervisor), so give each a distinct id::

# application.ex
children = [
  Supervisor.child_spec({Phoenix.PubSub, name: MyApp.PubSub}, id: MyApp.PubSub),
  Supervisor.child_spec(
    {Phoenix.PubSub, name: MyApp.EchoPubSub, adapter: EchoPubSub},
    id: MyApp.EchoPubSub
  )
]

Config Options

OptionDescriptionDefault
:-----------------------:------------------------------------------------------------------------:-------------
:nameThe required name to register the PubSub processes, ie: MyApp.PubSub
:pool_sizeThe number of workers and producers on each node1
:buffer_sizeThe numbers of messages to hold in memory for each producer in the pool10_000
:batch_intervalMilliseconds to batch writes before a flush; 0 flushes immediately200

With :pool_size > 1 there are independent producers and order/at-least-once is per producer (a broadcast routes by sender pid) - see Pools and ordering.

Subscribing processes should handle the message {:cursor_expired, node_name} which indicates that your client has been disconnected long enough that your position in the broadcaster's buffer has been overwritten. At this point it is the subscribing process's job to return to a valid state i.e. reloading state from source like database or another node.

Benchmarks

Verified end-to-end throughput on a real Fly.io cluster (fra, performance-4x), every node receives every message, no loss, at batch_interval=100, pool_size=1, publishers=4, 10 B payload:

nodesFly.io (fra)
3~79k msg/s
4~75k msg/s

Delivery is network-bound (Fly's private WireGuard mesh), and batching is the dominant lever - batch_interval=0 collapses to ~1k msg/s (each message becomes its own acked round-trip). Small payloads (10–30 B) barely move the numbers, but a 200 B whole-object payload cuts 4-node throughput ~42% - so send small deltas, not whole objects.

  • How to run (locally and on a Fly.io cloud cluster) - see the benchmark branch: bench/README.md.
  • Detailed results: Fly.

Credits

EchoPubSub is a fork of phoenix_pubsub_buffered by Eric Newbury, who designed and built the original at-least-once buffered PubSub adapter. All credit for the core design goes to him - this fork builds on that foundation with batched inter-node delivery, automatic replay on failure, flush-path expiry detection, telemetry, capacity warnings, and additional configuration options.