Skip to content

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.

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.

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.

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.

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.

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.

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 process
pubsub.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 only
pubsub.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.

// All subscribers across all nodes
let pids = pubsub.subscribers(pubsub_handle, "room:lobby")
// Count subscribers
let count = pubsub.subscriber_count(pubsub_handle, "room:lobby")

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.

The channel system uses PubSub internally for distributed broadcasts when configured:

import beryl
import 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 automatically
beryl.broadcast(channels, "room:lobby", "event", payload)