PubSub
beryl's PubSub API uses Erlang pg process groups. Publishers send messages
to a topic, and each process subscribed to that topic receives them.
Starting PubSub
Section titled “Starting PubSub”import beryl/pubsub
// Default scope ("beryl_pubsub")let pubsub_handle = pubsub.start(pubsub.default_config())
// Custom scope (isolates process groups)let pubsub_handle = pubsub.start(pubsub.config_with_scope("my_app_pubsub"))The scope maps to a pg scope atom and identifies the PubSub instance.
Different scopes are isolated and can use different payload types in one
process mailbox. All handles in one scope must use the same payload type.
beryl starts a node-owned supervisor with a separate subtree for each scope.
The subtree owns a membership registry and the pg process. Repeated start
calls share that registry; the service outlives its first caller. Start the
scope on each participating node rather than sending a handle between nodes.
Startup errors cause an OTP exit. beryl rejects a scope name already owned by
an external pg process; it does not adopt or stop that process.
Same-scope payload contract
Section titled “Same-scope payload contract”Payload compatibility is a caller obligation, not a runtime-enforced
guarantee. All handles and subscribers for one scope on all connected nodes
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.
Give incompatible payload types distinct, bounded 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 fixed deployment constants. A type annotation does not register a schema or check other handles. An incompatible same-scope handle can deliver a value of the wrong type or crash a subscriber. Only trusted BEAM peers may participate.
Use a separate fixed scope for presence replication: presence.with_pubsub
carries presence sync payloads, while beryl.with_pubsub carries JSON
broadcasts. Presence's versioned sync requests reject unknown versions; they
do not validate arbitrary BEAM terms or snapshot state.
This contract retains the existing five-element wire tuple and same-scope sharing. It introduces no wire migration. Applications that mix incompatible payload types must move them to distinct fixed scopes on every participating node; coordinate that configuration change because different scopes cannot communicate.
Subscribing
Section titled “Subscribing”Create one typed subscriber in the process that owns the mailbox. Join the required topics. Add PubSub delivery to the process selector:
let subscriber = pubsub.subscriber(pubsub_handle)pubsub.join(subscriber, "room:lobby")
let selector = process.new_selector() |> process.select(app_subject) |> pubsub.selecting(subscriber, RemoteBroadcast)
// Later:pubsub.leave(subscriber, "room:lobby")PubSub records arrive as raw BEAM messages. selecting checks only the scope
tag and tuple arity, not the topic, event, sender, or payload field types.
Malformed outer tuples do not match the selector and remain in the mailbox.
One process can select subscribers with different payload types if their
scopes differ.
Repeated or concurrent joins to the same topic create one membership per
owner process and scope, even when you use multiple subscriber handles.
The owner receives each broadcast once and counts as one subscriber.
One leave removes that membership; further leaves are harmless. Leaving
does not affect other owners, scopes, or topics.
Scope recovery
Section titled “Scope recovery”After a pg process crash, the scope supervisor restarts it. The surviving
membership registry restores the topics of live local subscriber owners.
You can keep existing handles. New joins and repeated joins use the same
registry, and one leave removes the owner's membership and recovery intent.
The registry monitors subscriber owners and drops their topics when they exit.
An ordinary pg crash does not restart another scope.
During recovery, a broadcast or subscriber query can see empty or partial
membership. A broadcast returns Nil without confirming delivery. beryl does
not buffer or replay broadcasts, and local recovery does not wait for
cluster-wide convergence.
Startup, joins, and leaves can exit with an OTP error during recovery or service failure. Membership calls use a five-second registry timeout. A failed or timed-out join or leave may have recorded its intent; failure does not roll it back. Retrying the operation is idempotent. Handle these exits at your application's OTP supervision boundary.
The registry is in-memory state. If the registry itself is lost, including
when repeated failures exhaust the supervisor's restart budget, old handles
become invalid. Their membership, broadcast, and subscriber-query calls exit
instead of silently using a replacement with no subscriptions. Call start
again, create new subscribers, and rejoin. Selector-only receivers do not get a
separate recovery notification. Service and node failures do not preserve
subscriptions.
Message format
Section titled “Message format”Subscribers receive typed Message(payload) records:
pub type Message(payload) { Message( topic: String, event: String, payload: payload, from: PubSubFrom, )}
pub type PubSubFrom { System // Broadcast with no sender FromPid(Pid) // Broadcast from a specific process FromSocket(Pid, String) // Broadcast from a process, excluding a socket ID}On the wire, the PubSub scope atom replaces the public record tag and is followed by these four fields in order. This five-element tuple is a frozen cross-node wire contract. Nodes using the old unscoped message shape do not interoperate, so upgrade the cluster together when adopting this version. Version changes to your own payload type explicitly when rolling upgrades must accept old and new nodes concurrently.
FromSocket contains the sender PID and a socket ID to exclude. Receiving
runtimes do not send the message to that socket. Thus,
beryl.broadcast_from excludes the sender across cluster nodes.
Send broadcasts
Section titled “Send broadcasts”import gleam/json
// Broadcast to all subscribers (all nodes)pubsub.broadcast( pubsub_handle, "room:lobby", "new_message", json.string("hello"),)
// Broadcast to all except the sender processpubsub.broadcast_from( pubsub_handle, process.self(), "room:lobby", "new_message", json.string("hello"),)
// Broadcast to all except a specific socket ID (clustered "broadcast except this socket")pubsub.broadcast_from_socket( pubsub_handle, process.self(), // sending runtime process socket_id, // socket ID to exclude on receiving runtimes "room:lobby", "new_message", json.string("hello"),)
// Broadcast to local node onlypubsub.local_broadcast( pubsub_handle, "room:lobby", "new_message", json.string("hello"),)Use broadcast_from_socket to send to all cluster subscribers except one
socket. The socket can be on another node. beryl.broadcast_from calls this
function.
List subscribers
Section titled “List subscribers”// All subscribers across all nodeslet pids = pubsub.subscribers(pubsub_handle, "room:lobby")
// Count subscriberslet count = pubsub.subscriber_count(pubsub_handle, "room:lobby")Use PubSub across Erlang nodes
Section titled “Use PubSub across Erlang nodes”Erlang pg works across connected nodes. After your application establishes
Erlang distribution between the nodes, pg merges their process groups and
sends messages to subscribers across the cluster. beryl PubSub needs no
additional configuration, but it does not connect the nodes for you.
Automated tests currently exercise distributed behavior with multiple actors on one BEAM node. Issue #365 tracks integration coverage across separate distributed Erlang nodes for PubSub delivery and presence convergence.
Use PubSub with beryl
Section titled “Use PubSub with beryl”The channel system uses PubSub internally for distributed broadcasts when configured:
import berylimport beryl/wire
let pubsub_handle = pubsub.start(pubsub.default_config())let config = beryl.config(wire.phoenix_codec()) |> beryl.with_pubsub(pubsub_handle)let assert Ok(#(channels, runtime_specification)) = beryl.child_spec(config, init: init, update: update)// Add `runtime_specification` to your application supervisor before using// `channels`.
// beryl.broadcast() sends to all nodes automaticallyberyl.broadcast(channels, "room:lobby", "event", payload)Next steps
Section titled “Next steps”- Supervision guide: supervised startup and multi-node deployment
- Architecture overview: PubSub's place in beryl
- Troubleshooting: diagnose cluster broadcast failures and different presence state across nodes
