Skip to content

beryl

beryl: type-safe real-time communication

beryl is a standalone Gleam library for building real-time applications on the BEAM. It provides app-side WebSocket dispatch, distributed presence tracking, PubSub messaging, and topic groups.

  • Sockets: App-side dispatch with topic-based WebSocket messaging routed by your update function (beryl, beryl/socket)
  • PubSub: Distributed publish/subscribe via Erlang pg (beryl/pubsub)
  • Presence: Distributed presence tracking backed by a causal-context CRDT (add-wins observed-remove set) (beryl/presence)
  • Groups: Named collections of topics for multi-topic broadcasting (beryl/group)
import beryl
import beryl/socket.{
AcceptJoin, Binary, Broadcast, Closed, Info, Join, Message, Next,
}
import beryl/pubsub
import beryl/wire
import gleam/json
import gleam/option
import gleam/otp/static_supervisor
pub fn main() -> Nil {
// Optional: start PubSub for distributed messaging
let pubsub_handle = pubsub.start(pubsub.default_config())
// Build the supervised system. The app supplies `init` (the per-socket
// model) and `update` (which routes every event by topic).
let config =
beryl.config(wire.phoenix_codec())
|> beryl.with_pubsub(pubsub_handle)
let assert Ok(#(sockets, child_specification)) =
beryl.child_spec(
config,
init: fn(_info) { #(Nil, []) },
update: fn(model, event) {
case event {
Join("room:" <> _, _payload, ref) ->
Next(model, [AcceptJoin(ref, option.None)])
Message(topic, "new_msg", payload, _ref) ->
Next(model, [Broadcast(topic, "new_msg", payload)])
Join(..) | Message(..) | Binary(..) | Closed(..) | Info(..) ->
Next(model, [])
}
},
)
let assert Ok(_root) =
static_supervisor.new(static_supervisor.OneForOne)
|> static_supervisor.add(child_specification)
|> static_supervisor.start()
// Broadcast to all subscribers of a topic
beryl.broadcast(sockets, "room:lobby", "announce", json.object([]))
}
pub type Config

Configuration for an app-side socket runtime.

This type is opaque. Construct it with config and adjust it with the with_* functions. beryl can then add configuration options without a breaking change.

pub type ConfigError {
HeartbeatTimeoutTooLow(minimum: Int)
RateLimitTooHigh(maximum: Int)
InvalidTopicPattern(
pattern: String,
reason: topic.TopicError
)
}

Why an eagerly validated Config was rejected before any process started.

child_spec validates the configuration before it allocates names or starts the runtime. It returns an invalid configuration instead of crashing a supervised child during initialization.

HeartbeatTimeoutTooLow(minimum: Int)

heartbeat_timeout_ms was below the minimum. The server derives its staleness check interval as heartbeat_timeout_ms / 2 (integer division). A timeout of 1 would round down to a check interval of 0 and disable heartbeat eviction. The wrapped Int is the smallest accepted timeout.

RateLimitTooHigh(maximum: Int)

A configured per-second rate was above the largest value that can still consume at least one token. The wrapped Int is the maximum accepted rate.

InvalidTopicPattern(
pattern: String,
reason: topic.TopicError
)

A per-topic-pattern rate limit used a pattern string that is not a valid topic pattern. pattern is the offending pattern and reason is the beryl/topic error nested rather than flattened to a string, so it stays matchable.

Match each topic.TopicError variant explicitly. Adding a variant affects API compatibility and requires updating exhaustive matches.

pub type LoggingConfig

Logging configuration for beryl diagnostics.

This type is opaque. Construct it with logging_config and adjust it with the with_* functions. beryl can then add logging options without a breaking change.

pub type LogLevel {
DebugLevel
InfoLevel
WarnLevel
ErrorLevel
}

Logging verbosity for beryl's internal loggers.

The variants carry a Level suffix so ErrorLevel does not shadow the prelude's Result Error constructor when imported unqualified.

pub type Sockets

A runtime system handle.

child_spec returns this opaque handle with the supervised subtree. Pass it to broadcast, group, and transport functions. beryl hides its internals so they can change without breaking application code.

The handle is non-generic. An app-side dispatch system is generic over the application's model/message, but those types are sealed inside monomorphic closures at construction time. They never appear in this handle or in any transport signature.

pub type StopError {
NotRunning
StopTimeout
}

Errors when stopping a beryl system with stop.

NotRunning

The handle referred to a system that was not running: it was never started (for example, a child_spec handle whose supervisor was never added to a running tree) or it has already been stopped. stop is safe to call in these cases; it reports NotRunning rather than crashing.

StopTimeout

The runtime did not acknowledge the stop request within the shutdown window. The system may still be terminating.

pub fn broadcast(
Sockets,
String,
String,
json.Json
) -> Result(Nil, overload.AdmissionError)

Broadcast a message to all subscribers of a topic.

This function sends the message to all sockets subscribed to the topic. When the system was started with PubSub, it also sends the broadcast to subscribers on other nodes.

beryl.broadcast(
sockets,
"room:lobby",
"new_message",
json.object([#("text", json.string("Hello!"))]),
)
pub fn broadcast_from(
Sockets,
String,
String,
String,
json.Json
) -> Result(Nil, overload.AdmissionError)

Broadcast a message to all subscribers except one socket.

Useful for broadcasting a message to everyone except the sender. When PubSub is configured, the excluded socket ID is preserved across nodes so clustered deployments do not echo the event back to that socket on another node.

beryl.broadcast_from(
sockets,
socket_id,
"room:lobby",
"user_typing",
json.object([#("user", json.string("alice"))]),
)
pub fn broadcast_presence_diff(
Sockets,
String,
presence.Diff
) -> Result(Nil, overload.AdmissionError)

Broadcast a Phoenix-compatible presence_diff event for a topic.

This encodes the topic's joins and leaves as:

{
"joins": { "user:1": { "metas": [{ "status": "online" }] } },
"leaves": { "user:2": { "metas": [{ "status": "offline" }] } }
}

Honors presence.diff_scope: application-mutation and explicitly constructed Cluster diffs use the same distributed semantics as broadcast. LocalNode diffs from replication, failure detection, and recovery reach only local socket subscribers, including other local runtimes in the same PubSub scope. They must not change healthy clients on another node.

Pass the original diff from presence.with_on_diff. Encoding or rebuilding the diff before an unconditional broadcast loses this routing metadata.

pub fn child_spec(
Config,
init: fn(socket.ConnectInfo(a)) -> #(b, List(socket.Effect)),
update: fn(b, socket.Input(a)) -> socket.Next(b)
) -> Result(#(Sockets, supervision.ChildSpecification(static_supervisor.Supervisor)), ConfigError)

Build the app-side dispatch supervision child specification.

Add the returned specification to the application's supervision tree. This function validates the configuration before the application's supervisor starts. It returns an error instead of crashing a supervised child during initialization.

The returned Sockets handle is name-backed and usable immediately, even before the supervision tree that owns the returned child specification is started. Before startup, during a runtime restart window, and after shutdown, sends return admission errors and connection admission fails.

let assert Ok(#(sockets, child_specification)) =
beryl.child_spec(beryl.config(wire.phoenix_codec()), init:, update:)
let assert Ok(_root) =
static_supervisor.new(static_supervisor.OneForOne)
|> static_supervisor.add(child_specification)
|> static_supervisor.start()
// `sockets` is usable once the tree above is running.
pub fn config(codec.Codec) -> Config

Build a configuration with the default settings.

Pass wire.phoenix_codec() for Phoenix framing or a custom Codec for another framing.

pub fn logging_config(
level: LogLevel,
include_payloads: Bool
) -> LoggingConfig

Build a logging configuration.

Payloads are excluded by default to avoid accidental sensitive-data exposure. Use with_payload_preview_bytes to adjust the bounded preview size when payload previews are enabled.

pub fn stop(Sockets) -> Result(Nil, StopError)

Stop a beryl system.

This function drains and stops the supervised runtime. It delivers Closed to every joined topic before it closes each transport connection. Presence is application-owned and is not stopped by this function. The runtime is a Transient child, so it is not restarted after a graceful stop.

You can call stop more than once or use a handle whose system never started. In these cases, it returns Error(NotRunning) and does not crash. It returns Error(StopTimeout) if the app runtime does not acknowledge the stop within the shutdown window. After a successful stop the handle should no longer be used.

pub fn validate_config(Config) -> Result(Nil, ConfigError)

Validate a Config without starting any process.

This checks that heartbeat_timeout_ms is at least 2, that every configured rate limit is representable, and that every per-topic rate-limit pattern is valid.

pub fn with_channel_rate(
Config,
per_second: Int,
burst: Int
) -> Config

Configure per-channel message rate limiting.

The limiter applies only after a socket has joined a topic. Active per-socket channel buckets are capped by default; use with_channel_rate_max_keys_per_socket to adjust the cap.

pub fn with_channel_rate_max_keys_per_socket(
Config,
max_keys: Int
) -> Config

Configure the maximum active per-channel rate-limit buckets per socket.

Values <= 0 disable the cap. The default is 1000.

pub fn with_connection_rate_per_ip(
Config,
per_second: Int,
burst: Int
) -> Config

Configure a per-IP connection-attempt rate limit.

Each WebSocket upgrade attempt that passes the concurrent connection ceilings consumes one token before authentication and handshake setup. A non-positive per_second disables the limit. A burst of 0 uses per_second as the burst capacity.

Unlike per-connection frame and message buckets, this allowance is keyed by the real socket peer IP and lives in beryl's supervised connection limiter. Disconnecting or restarting the app runtime therefore does not refresh it. Idle IP buckets are removed once their allowance has fully refilled.

This uses the same peer IP and trusted-proxy caveats as with_max_connections_per_ip.

pub fn with_effect_limits(
Config,
overload.Limits
) -> Config

Bound each callback's effect count and inspectable result bytes.

pub fn with_frame_rate(
Config,
per_second: Int,
burst: Int
) -> Config

Configure per-connection frame-rate limiting at the transport edge.

Every complete inbound text or binary frame consumes this independent bucket before decoding. Configure it alongside with_message_rate to combine edge shedding with a runtime cap on decoded non-join traffic. An over-rate heartbeat is shed before it can refresh the socket's heartbeat deadline, so a sustained flood is eventually closed by heartbeat eviction.

pub fn with_heartbeat(
Config,
timeout_ms: Int
) -> Config

Configure the server-side heartbeat staleness window.

The runtime evicts a socket that sends no heartbeat within timeout_ms. It checks at half this window, so values below 2 are rejected by validate_config with HeartbeatTimeoutTooLow. The default is 60000 ms.

pub fn with_join_rate(
Config,
per_second: Int,
burst: Int
) -> Config

Configure per-socket join rate limiting.

pub fn with_logging(
Config,
LoggingConfig
) -> Config

Configure beryl's internal logging.

pub fn with_max_connections(
Config,
max_connections: Int
) -> Config

Configure the maximum number of concurrent connections for this beryl system on one BEAM node, regardless of source IP.

A value of 0, the default, means unlimited. When a limit is set, a transport admits a connection only while this system is below the limit. It rejects other connections before it allocates long-lived per-socket state. The transport frees the slot when the connection closes, its process dies, or its handshake or setup fails. The limiter actor performs the check and increment atomically. Concurrent opens cannot materially exceed the ceiling.

Each independently constructed Sockets system owns a separate limiter. Two systems on the same node therefore each have their own configured capacity; this is not one shared ceiling for every beryl system on the node.

This per-system ceiling works with with_max_connections_per_ip. When both are set, a connection must be under both limits. The per-IP limit throttles any single abusive peer, while this total ceiling bounds the system's resource use so that many distinct source addresses (for example, a botnet or IPv6 address rotation) cannot exhaust its process, socket, and runtime budget. A per-IP limit alone cannot stop this case.

Each beryl system enforces this ceiling independently on its BEAM node. With one system per node, a load-balanced cluster's effective ceiling is roughly max_connections × node_count (subject to how the balancer distributes connections). Multiple systems on one node contribute separate allowances. Size their combined limits against that node's capacity. Use the load balancer's global connection/rate controls when you need a cluster-wide cap.

pub fn with_max_connections_per_ip(
Config,
max_connections: Int
) -> Config

Configure the maximum number of concurrent connections allowed per client IP address.

A value of 0, the default, means unlimited. When a limit is set, a transport admits a new connection only while the peer is below the limit. It rejects other connections. The transport frees the slot when the connection closes.

The limit is enforced on the real socket peer IP as reported by the transport (for the Mist transport, the address of the TCP connection). beryl does not trust or parse forwarded headers such as X-Forwarded-For. A client can set these headers and spoof its address to bypass this limit.

If beryl runs behind a trusted reverse proxy or load balancer, every connection shares the proxy's address, so a per-IP limit throttles all clients as a single IP. In that topology you must resolve the real client IP yourself at the proxy layer, for example by enforcing limits there. See the WebSocket transport guide for deployment guidance.

pub fn with_max_event_length(
Config,
max_length: Int
) -> Config

Configure the maximum allowed byte length for client-supplied event name strings.

Event names longer than max_length bytes are dropped before reaching the app's update function. The default is 64.

pub fn with_max_inbound_frame_bytes(
Config,
max_bytes: Int
) -> Config

Configure the maximum allowed inbound WebSocket frame size in bytes.

beryl enforces the limit post-assembly. The transport (Mist or Ewe) buffers and assembles a complete frame first. beryl then measures it and closes the connection if it exceeds max_bytes. This bounds per-message processing cost (decode, routing, rate-limit accounting), but it does not by itself bound transport memory. A hostile client can declare a huge payload and stream it slowly, or send many fragmented continuation frames, and the transport's receive buffer grows before this check runs. This setting alone does not stop one connection from exhausting node memory.

For a true transport memory bound you must place an edge proxy or load balancer in front of beryl and configure a WebSocket frame-size limit there (and a matching request/body size limit). beryl's connection, frame-rate, and message-rate limits all run after frame assembly and do not mitigate this vector. See the README's "Security" section.

Values <= 0 disable the cap. The default is 1 MiB.

pub fn with_max_joined_topics_per_socket(
Config,
max_topics: Int
) -> Config

Configure the maximum number of topics a socket may join at once.

Values <= 0 disable the cap. The default is 1000.

pub fn with_max_topic_length(
Config,
max_length: Int
) -> Config

Configure the maximum allowed byte length for client-supplied topic strings.

Topics longer than max_length bytes are rejected with a phx_reply error before reaching your update function, bounding the size of keys tracked per socket. The default is 256.

pub fn with_message_rate(
Config,
per_second: Int,
burst: Int
) -> Config

Configure per-socket decoded message-rate limiting in the runtime.

Joins use with_join_rate; decoded leaves, heartbeats, events, decoded binary, and raw Binary inputs consume this bucket. It is independent of with_frame_rate. An over-rate heartbeat does not refresh the socket's heartbeat deadline, so a sustained flood is eventually closed by heartbeat eviction. Leave enough rate and burst headroom for legitimate heartbeats.

pub fn with_payload_preview_bytes(
LoggingConfig,
bytes: Int
) -> LoggingConfig

Configure the maximum payload and frame preview length for logs.

The bytes argument is measured in grapheme clusters.

pub fn with_presence_handle(
Config,
presence: presence.Presence
) -> Config

Attach the presence actor used by socket presence effects and channel.on_presence callbacks.

pub fn with_pubsub(
Config,
pubsub.PubSub(json.Json)
) -> Config

Add PubSub to a configuration for distributed broadcasts.

pub fn with_router_queue_limits(
Config,
overload.Limits
) -> Config

Bound outstanding local router requests.

pub fn with_socket_queue_limits(
Config,
overload.Limits
) -> Config

Bound outstanding socket work, including presence-suspended continuations.

pub fn with_telemetry(Config) -> Config

Enable beryl's :telemetry events.

pub fn with_topic_rate(
Config,
pattern: String,
per_second: Int,
burst: Int
) -> Config

Configure a per-topic-pattern message rate limit for an app-dispatch runtime built with child_spec.

Patterns use the same syntax as topic routing ("room:*", "document:*:ops", "*"). The runtime checks limits in the order they were added. The first matching pattern wins. Topics that match no pattern use the global with_channel_rate limit. The limiter applies only after a socket has joined the topic. A non-positive per_second explicitly disables limiting for matching topics, including any global channel limit, and allocates no bucket.

pub fn with_worker_queue_limits(
Config,
overload.Limits
) -> Config

Bound outstanding input for each topic-worker incarnation.