beryl/pubsub
Distributed PubSub with Erlang pg
This module provides topic-based PubSub backed by Erlang's built-in pg
module. Subscribers are tracked by process group, so messages are
delivered to all connected nodes in the cluster.
The payload is generic: PubSub(payload) and Message(payload) carry
whatever Gleam type a given scope is started with. A broadcast sends that
value as a scope-tagged native BEAM term. There is no encoding step, even
across nodes, because Erlang's distribution protocol marshals arbitrary
terms. Use a gleam/json payload only when the data is also
destined for a JSON-speaking client (e.g. relayed on to a WebSocket
browser); payloads that never leave the cluster are cheaper and safer as
plain Gleam types.
Same-scope payload contract
Section titled “Same-scope payload contract”Payload compatibility is a caller obligation, not a runtime check. Every
handle and subscriber for one scope, on every connected node, must use the
same payload type and compatible native term representation. Repeated
start calls for one scope do not create isolated instances. Different
topics in that scope do not provide type isolation either.
Give incompatible payload types distinct, fixed configuration scopes:
let numbers: pubsub.PubSub(Int) = pubsub.start(pubsub.config_with_scope("my_app_numbers"))let text: pubsub.PubSub(String) = pubsub.start(pubsub.config_with_scope("my_app_text"))These names are a bounded set of deployment constants, not names derived from requests or tenants. A type annotation on a handle does not register or validate a schema. Violating this contract can deliver a value of the wrong type or crash a subscriber. Only trusted BEAM peers may participate.
Scope recovery and delivery
Section titled “Scope recovery and delivery”Each node runs a shared beryl PubSub supervisor. Each scope has its own
supervisor, membership registry, and pg process. The service outlives
the process that calls start. A pg restart preserves the registry,
which restores the topics of live local subscriber owners. Existing
handles remain valid, and repeated joins still create one membership.
Broadcasts are best-effort. During recovery, broadcasts and subscriber
queries can see empty or partial membership. Broadcasts are not buffered
or replayed, and a Nil return does not confirm delivery. Each node must
start its own scope; local recovery is not a cluster-wide readiness barrier.
Startup and membership operations can exit with an OTP error during
service failure or recovery. A failed or timed-out join or leave may have
recorded its intent; failure does not roll it back. Retrying is idempotent.
See start for startup conflicts and loss of the membership registry.
Quick start
Section titled “Quick start”let pubsub_handle = pubsub.start(pubsub.default_config())
// Sending: broadcast to all subscribers of a topicpubsub.broadcast(pubsub_handle, "room:lobby", "new_msg", "hello")
// Receiving: create a `subscriber`, join topics, and fold its// `selecting` into an actor's own `Selector`. `RemoteBroadcast` here is// the actor's own message constructor that wraps a `pubsub.Message`.let subscriber = pubsub.subscriber(pubsub_handle)pubsub.join(subscriber, "room:lobby")let selector = process.new_selector() |> process.select(subject) |> pubsub.selecting(subscriber, RemoteBroadcast)Message
Section titled “Message”pub type Message(a) { Message( topic: String, event: String, payload: a, from: PubSubFrom )}A PubSub message delivered to subscribers.
This type is transparent. Subscribers can inspect the topic, event, payload, and sender metadata in their process mailbox.
Frozen wire contract
Section titled “Frozen wire contract”beryl sends broadcasts raw between nodes through pg as a five-element
tuple. The PubSub scope atom comes first. The four message fields follow in
order. This scope-tagged runtime shape is the frozen wire contract for each
payload type. The contract also applies to PubSubFrom.
Payloads travel as native terms, not in a self-describing format such as
JSON. A change to the shape of your payload type is therefore a wire
change. Add your own version, for example an explicit v field, if the
shape must change during a rolling upgrade. Receive broadcasts with
selecting. It adds the subscriber's raw mailbox messages to a typed
Selector under the same-scope payload contract. It does not validate
field types or payloads. Do not match the raw process message yourself.
PubSub
Section titled “PubSub”pub type PubSub(a)A running PubSub instance.
This handle is opaque. Callers cannot forge pg scopes or depend on the
runtime representation. payload sets the Gleam type for every Message
sent through this instance. The scope identifies the runtime instance.
All handles for one scope, across all connected nodes, must use the same
payload type and compatible native term representation. This is a caller
obligation, not a runtime-enforced guarantee.
PubSubConfig
Section titled “PubSubConfig”pub type PubSubConfigPubSub configuration.
Build with default_config or config_with_scope so the underlying pg
scope representation can evolve without exposing record fields.
PubSubFrom
Section titled “PubSubFrom”pub type PubSubFrom { System FromPid(process.Pid) FromSocket( process.Pid, String )}Identifies the sender of a broadcast.
Part of the frozen wire contract described on Message.
Constructors
Section titled “Constructors”System
Section titled “System”SystemThe broadcast came from the system and has no sender PID.
FromPid
Section titled “FromPid”FromPid(process.Pid)The broadcast came from a specific process.
FromSocket
Section titled “FromSocket”FromSocket( process.Pid, String)The broadcast came from a process and must exclude a socket ID.
Subscriber
Section titled “Subscriber”pub type Subscriber(a)A typed subscription handle owned by a single process.
A Subscriber contains the pg scope and owning process. It also carries
the payload type at compile time. Join it to any number of topics with
join. One selecting call receives their frozen raw Message(payload)
records as typed values.
Create it in the receiving process, such as an actor's initializer.
A Subject delivers messages only to its owner.
Functions
Section titled “Functions”broadcast
Section titled “broadcast”pub fn broadcast( PubSub(a), String, String, a) -> NilBroadcast a message to all topic subscribers on all nodes.
broadcast_from
Section titled “broadcast_from”pub fn broadcast_from( PubSub(a), process.Pid, String, String, a) -> NilBroadcast a message to all subscribers except a specific PID.
broadcast_from_socket
Section titled “broadcast_from_socket”pub fn broadcast_from_socket( PubSub(a), process.Pid, String, String, String, a) -> NilBroadcast a message to all subscribers except a process.
Preserve a socket ID that receiving runtimes must exclude locally.
config_with_scope
Section titled “config_with_scope”pub fn config_with_scope(String) -> PubSubConfigCreate a PubSub configuration with a custom scope name.
The scope name is converted to an Erlang atom via binary_to_atom.
Atoms are never garbage-collected, so the scope name must be a
static, bounded deployment or configuration value. Never use raw
user-derived, per-request, per-tenant, database-derived, or otherwise
unbounded high-cardinality runtime input. A deployment-controlled value
is acceptable only when validated or selected from a fixed bounded set.
A malicious or high-cardinality source can exhaust the BEAM atom table
and crash the VM.
// Correct: static deployment constantpubsub.config_with_scope("my_app_pubsub")
// Correct: deployment-controlled and selected from a fixed bounded set// pubsub.config_with_scope(config.pubsub_scope())
// WRONG: never do this// pubsub.config_with_scope(user_request.tenant_id)// pubsub.config_with_scope(database_row.name)default_config
Section titled “default_config”pub fn default_config() -> PubSubConfigCreate a default PubSub configuration with scope beryl_pubsub.
pub fn join( Subscriber(a), String) -> NilJoin a topic so this subscriber receives broadcasts sent to it.
A subscriber can join many topics. All topics deliver to its owner's mailbox. Repeated or concurrent joins create one membership per owner process, scope, and topic, including joins through different handles. Each broadcast delivers once to that owner, and subscriber counts include it once.
The registry retains the membership through pg restarts and removes it
when its owner exits. During recovery this call can exit with an OTP error.
The call uses a five-second registry timeout; an error or timeout does not
roll back intent already recorded by the registry. Retrying is idempotent.
pub fn leave( Subscriber(a), String) -> NilLeave a topic previously joined with join.
One call removes the owner's membership for this scope and topic, even after repeated joins through different handles. Repeated leaves are harmless. Other owners, scopes, and topics are unaffected.
A leave removes recovery intent as well as live membership. Like join,
it can exit during recovery or after a five-second registry timeout.
A failed call may still take effect; retrying is idempotent.
local_broadcast
Section titled “local_broadcast”pub fn local_broadcast( PubSub(a), String, String, a) -> NilBroadcast a message only to subscribers on the current node.
selecting
Section titled “selecting”pub fn selecting( process.Selector(a), Subscriber(b), fn(Message(b)) -> a) -> process.Selector(a)Add a subscriber's PubSub message delivery to a Selector, alongside a
process's own subjects.
pg tracks bare pids, so broadcasts arrive as raw process messages.
selecting validates the subscriber's scope tag and four-field arity
before it recovers the compile-time payload type. It does not validate
topic, event, sender, or payload field types. Add it once. Every joined
topic uses the same mailbox.
Subscribers for different scopes may safely use different payload types in one process. All subscribers for the same scope must use the same payload type.
let subscriber = pubsub.subscriber(pubsub_handle)pubsub.join(subscriber, "room:lobby")let selector = process.new_selector() |> process.select(subject) |> pubsub.selecting(subscriber, RemoteBroadcast)pub fn start(PubSubConfig) -> PubSub(a)Start a PubSub instance.
This starts or attaches to a node-owned supervised scope. Repeated calls
share its membership registry; the first caller does not own its lifetime.
The registry restores live local subscriptions after a pg process restart.
Each participating node must start the same scope.
Startup errors cause an OTP exit instead of returning a handle. A scope
already owned by an external pg process is a startup conflict; beryl does
not adopt or stop it. Startup waits for supervision and local membership
reconciliation, but can fail if the service is recovering.
Loss of the registry itself, including supervisor restart-intensity
exhaustion, invalidates existing handles. Their membership, broadcast, and
subscriber-query calls exit rather than silently using an empty replacement.
Call start again, create new subscribers, and rejoin their topics.
This is not persistent storage across service or node failure.
payload is fixed by how the returned value is used or annotated at the
call site. For example: pubsub.start(config) : PubSub(MySyncPayload).
Starting the same scope again returns another handle to the same runtime
instance, so every use of that scope must choose the same payload type.
These calls do not create isolated instances or check payload compatibility.
Use distinct, fixed scopes for incompatible types, as in the module example.
subscriber
Section titled “subscriber”pub fn subscriber(PubSub(a)) -> Subscriber(a)Create a subscription handle owned by the current process.
Call this function from the process that will receive broadcasts, such as
an actor initializer or test process. Then join topics and add selecting
to that process's Selector.
A process may create subscribers for multiple scopes and payload types.
selecting uses each subscriber's scope to keep their raw mailbox messages
separate.
subscriber_count
Section titled “subscriber_count”pub fn subscriber_count( PubSub(a), String) -> IntReturn the number of topic subscribers on all nodes.
subscribers
Section titled “subscribers”pub fn subscribers( PubSub(a), String) -> List(process.Pid)Return all topic subscribers on all nodes.
