Skip to content

beryl/channel

The channel composition surface: a channel is a topic pattern paired with a typed join callback and callbacks over private state.

import beryl/channel
import gleam/json
pub type Note {
Announce(String)
}
pub fn room() -> channel.Handler {
channel.handler("room:*", fn(context) {
channel.notify(context.self, Announce("later, on this topic"))
channel.accept(0)
|> channel.on_message(fn(count, message) {
channel.next(count + 1, [
channel.broadcast(message.event, json.int(count + 1)),
])
})
|> channel.on_info(fn(count, note) {
let Announce(text) = note
channel.next(count, [
channel.push("announce", json.string(text)),
])
})
|> channel.on_terminate(fn(_count, _reason) {
[channel.broadcast("left", json.string(context.topic))]
})
|> channel.with_actions([
channel.push("welcome", json.string(context.topic)),
])
})
}

A channel picks two types of its own: state, its private model, and info, the type of server-side messages it accepts. Neither escapes: handler seals both inside its registration closure, so the resulting Handler is not generic and handlers with unrelated state and info types compose in one list. The runtime does not erase values to Dynamic or use unchecked coercion: a closure carries each typed info value. Only the worker for the join that created the closure can open it. After the join ends, the worker no longer exists and the runtime drops the value.

Each accepted topic runs in its own worker process under a supervisor that its socket actor owns. The socket actor owns the protocol state, refs, subscriptions, presence data, and frame writes. The worker runs join during startup. The socket actor waits for a maximum of five seconds. The worker also runs on_message and on_info, and sends its actions to the socket actor.

The runtime applies action lists from left to right. Actions always target the channel's own topic. The runtime converts them to beryl's core Effect values and applies them in list order. Action order is therefore wire order for one topic. beryl does not define an order between different topics on one socket. Their workers run concurrently, as Phoenix channel processes do. An asynchronous presence effect can pause this socket while other sockets continue. The runtime resumes the remaining actions after that effect completes.

A join's actions (see with_actions) are emitted with the join acknowledgment, immediately after it: the socket is already subscribed, so a push cannot precede its own join reply. This ordering does not make an asynchronous presence mutation a cross-socket reservation; use application-owned synchronous state for atomic capacity checks.

The worker processes queued messages before a close. Thus, the runtime still delivers a push or reply that the worker computed before a leave. The runtime then applies on_terminate actions in the turn that closes the topic, after the channel instance is gone. Its closing-phase action type permits only operations that remain meaningful then.

pub type Action(a)

One operation on the channel's own topic.

The phase parameter prevents on_terminate from returning active-only operations. Put actions in a list in wire order.

pub type Active

Marker for actions valid while a channel is active.

pub type ChildSpecError {
InvalidPattern(
pattern: String,
reason: topic.TopicError
)
DuplicatePattern(pattern: String)
InvalidConfig(reason: beryl.ConfigError)
}

Why building a channel-system child specification failed.

The function validates handler patterns before the core configuration. It checks each pattern's syntax in registration order. It then checks for exact duplicates in the same order. Overlapping patterns are allowed when they are not identical because routing uses the first match.

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

A handler used an invalid topic pattern.

DuplicatePattern(pattern: String)

Two handlers registered the same pattern string.

InvalidConfig(reason: beryl.ConfigError)

The core beryl.Config failed eager validation.

pub type Closing

Marker for actions valid while a channel is closing.

pub type Handler

A registered channel: a topic pattern plus its sealed join callback.

Handler is not generic. The closure contains the channel's sealed state and info types. A single List(Handler) can therefore hold channels with unrelated types.

pub type JoinContext(a) {
JoinContext(
socket_id: String,
seed: socket.ConnectSeed,
self: Sender(a),
topic: String,
parameters: List(String),
payload: dynamic.Dynamic
)
}

Information about one join attempt.

parameters contains wildcard captures in pattern order and is empty for exact patterns. self is this channel's generation-scoped Sender, for scheduling a later turn.

pub type JoinResult(a, b)

A join callback's answer: join this channel, or refuse.

Start an accepted result with accept, then pipe it through the on_* functions for the inputs that the channel handles. Unhandled inputs keep the channel joined and produce no actions. handler seals the private state and info types.

pub type Message {
Message(
event: String,
payload: dynamic.Dynamic,
reply: option.Option(socket.ReplyRef)
)
}

A client message delivered to a joined channel's on_message callback.

reply is present only when the client asked for a reply; pass it to reply_ok, reply_error, or discard_reply.

pub type Next(a)

What a channel callback decided to do next.

Build one with next or close. Use stay when the state does not change and there are no actions.

pub type Sender(a)

A typed handle for sending server-side messages to one joined channel.

Get this handle from JoinContext in the join callback. You can share it with any process. The channel's on_info callback receives each message with its type intact.

A sender is scoped to the join that produced it. Sending reserves worker capacity and returns an admission Result. A closed or stale sender returns an error; it cannot send to a later join of the same topic.

A sealed function carries each message to the worker. The worker opens the function and uses a selective receive in the same turn. Function environments are not included in accounted bytes. Bound typed message payloads in the application as well as configuring item limits.

pub fn accept(a) -> JoinResult(a, b)

Accept the join with an empty acknowledgment.

Pipe the result through on_message, on_info, and on_terminate to register the callbacks that the channel needs.

pub fn broadcast(
String,
json.Json
) -> Action(a)

Broadcast to every subscriber of this channel's topic, including this socket.

pub fn broadcast_from(
String,
json.Json
) -> Action(a)

Broadcast to every subscriber of this channel's topic except this socket.

pub fn broadcast_presence(
String,
fn(List(presence.PresenceEntry)) -> json.Json
) -> Action(a)

Broadcast a presence snapshot for this channel's topic to every subscriber, with the same apply-time encode semantics as push_presence.

pub fn child_spec(
beryl.Config,
handlers: List(Handler)
) -> Result(#(beryl.Sockets, supervision.ChildSpecification(static_supervisor.Supervisor)), ChildSpecError)

Build a channel system's supervision child specification for embedding in an application's supervision tree.

Like beryl.child_spec, this function reports only errors that it can detect before the tree starts. It validates the handler table first and then validates beryl.Config. You can use the returned beryl.Sockets after the owning tree starts.

let assert Ok(#(sockets, child_specification)) =
channel.child_spec(
beryl.config(wire.phoenix_codec()),
handlers: [room.channel()],
)
let assert Ok(_root) =
static_supervisor.new(static_supervisor.OneForOne)
|> static_supervisor.add(child_specification)
|> static_supervisor.start()
pub fn close(List(Action(Active))) -> Next(a)

Leave this channel after applying actions in order.

The socket stays connected. Its other channels do not change. This channel's on_terminate callback still runs.

pub fn discard_reply(option.Option(socket.ReplyRef)) -> Action(Active)

Discard a client message reply handle without sending a wire reply.

Use this for messages the application intentionally will not answer. option.None produces no effect.

pub fn handler(
String,
fn(JoinContext(a)) -> JoinResult(b, a)
) -> Handler

Register a channel for every topic matching pattern.

pattern uses beryl's topic pattern syntax ("room:lobby", "room:*", "document:*:ops", "*") and is validated when the handler table is used by channel.child_spec.

join receives one JoinContext containing connection data, the concrete topic, wildcard captures, and the payload.

pub fn next(
a,
List(Action(Active))
) -> Next(a)

Stay joined with the given state, applying actions in order.

pub fn notify(
Sender(a),
a
) -> Result(Nil, overload.AdmissionError)

Send a typed server-side message to the channel that owns sender.

Each accepted call queues one message. The runtime does not combine sends. The worker processes accepted messages in queue order.

Ok(Nil) confirms admission, not callback completion. Queue saturation, oversized input, and a closed or unavailable worker return an error. Rejection alone does not close the channel. See Sender for delivery cost and the limits of sealed-message byte accounting.

pub fn on_info(
JoinResult(a, b),
fn(a, b) -> Next(a)
) -> JoinResult(a, b)

Handle typed server-side messages sent through this channel's Sender.

pub fn on_message(
JoinResult(a, b),
fn(a, Message) -> Next(a)
) -> JoinResult(a, b)

Handle client messages on this channel's topic.

pub fn on_presence(
JoinResult(a, b),
fn(a, presence.Event) -> Next(a)
) -> JoinResult(a, b)

Observe this topic's initial presence roster and later changes.

The callback runs in this channel's worker with its private state. It receives presence.Snapshot first, then presence.Changed events, including this connection's changes. Ordinary message and info callbacks wait for the initial callback. The stream reflects the local presence replica, not a globally consistent cluster snapshot.

Observation does not track this connection. Add with_presence to track it as well. Its initial track may be in the snapshot or a later change, depending on actor ordering.

A missing presence handle rejects the join. Source failure, subscription timeout, callback panic, or more than 64 pending change batches closes only this topic with an error. Rejoin to obtain a fresh snapshot. One large snapshot or metadata value is not bounded by that batch limit. Subscription startup waits up to five seconds by default.

Apply leaves before joins when maintaining a roster. A callback that changes presence can trigger itself again; avoid unconditional updates. Repeated calls replace the callback. A rejected result stays rejected.

import beryl/presence
import gleam/list
channel.accept(0)
|> channel.on_presence(fn(online_sessions, event) {
let next = case event {
presence.Snapshot(entries) -> list.length(entries)
presence.Changed(joins, leaves) ->
online_sessions + list.length(joins) - list.length(leaves)
}
channel.stay(next)
})
pub fn on_terminate(
JoinResult(a, b),
fn(a, socket.StopReason) -> List(Action(Closing))
) -> JoinResult(a, b)

Run cleanup when the channel ends for any reason: client leave, a close result, a socket teardown, or a disconnect.

The runtime applies the returned closing-phase actions in the turn that closes this topic, after it removes the channel instance. This phase allows broadcasts, presence untracking, and presence broadcasts. It does not allow pushes, replies, or presence tracking.

An on_terminate panic does not stop the socket. The runtime logs the panic and completes the close without this callback's actions. The worker stops, so a Sender for this join cannot deliver more messages.

pub fn presence_track(
String,
json.Json
) -> Action(Active)

Track this socket's presence under key on this channel's topic and broadcast the matching presence_diff join.

Requires a presence handle on the Config (beryl.with_presence_handle).

pub fn presence_untrack(String) -> Action(a)

Untrack a presence previously tracked with presence_track and broadcast the matching presence_diff leave.

pub fn push(
String,
json.Json
) -> Action(Active)

Push a server-initiated message to this socket on this channel's topic.

pub fn push_presence(
String,
fn(List(presence.PresenceEntry)) -> json.Json
) -> Action(Active)

Push a presence snapshot for this channel's topic to this socket.

encode runs when the action is applied, so it already sees any earlier presence_track or presence_untrack in the same list.

pub fn queue_snapshot(Sender(a)) -> Result(overload.Occupancy, overload.AdmissionError)

Read this worker incarnation's queue accounting without waiting for it.

pub fn reject(json.Json) -> JoinResult(a, b)

Refuse the join, returning reason to the client.

pub fn reply_error(
option.Option(socket.ReplyRef),
json.Json
) -> Action(Active)

Reply with an error when a client message supplied a reply handle.

option.None produces no effect.

pub fn reply_ok(
option.Option(socket.ReplyRef),
json.Json
) -> Action(Active)

Reply successfully when a client message supplied a reply handle.

option.None produces no effect.

pub fn stay(a) -> Next(a)

Stay joined with the given state and no actions.

pub fn with_actions(
JoinResult(a, b),
List(Action(Active))
) -> JoinResult(a, b)

Add ordered actions to an accepted join.

The runtime emits the actions with the acknowledgment and applies them after it. The socket is therefore already subscribed to the topic. A push cannot overtake its own join reply. If an action becomes an asynchronous presence effect, the runtime may process other sockets while this socket waits. A check followed by presence_track is not an atomic cross-socket capacity reservation.

Use this function instead of notifying the channel from join: notify schedules a later input, while actions preserve their declared position immediately after the join acknowledgment.

Existing actions stay before the actions added here. A refused join has no topic, so this function returns reject results unchanged.

pub fn with_presence(
JoinResult(a, b),
key: String,
meta: json.Json
) -> JoinResult(a, b)

Track this connection and send a Phoenix-compatible presence snapshot after an accepted join.

A shorthand for existing presence actions, not a new presence lifecycle. Diff delivery and automatic cleanup also apply when using those actions directly.

Requires a running presence actor attached with beryl.with_presence_handle. Without a handle, the runtime logs warnings and drops the presence actions; it does not reject the join.

Appends presence_track, then push_presence with event presence_state and beryl/presence/wire.encode_state. The snapshot includes the new entry after tracking succeeds. Existing join actions stay before these actions. A rejected join remains unchanged.

Tracking broadcasts a presence_diff, which can arrive before the initial snapshot. Phoenix Presence clients buffer diffs until presence_state. The runtime removes this connection's entries when the topic closes. Connections with the same key remain separate entries under that key.

To replace metadata later, return presence_track with the same key from a callback. This builder does not observe state changes or register server-side callbacks for presence changes; use on_presence for those callbacks. Use the actions directly for a custom snapshot event name or encoder. Neither approach reserves room capacity.

channel.accept(state)
|> channel.with_presence(
key: "user:alice",
meta: json.object([#("status", json.string("online"))]),
)
pub fn with_reply(
JoinResult(a, b),
json.Json
) -> JoinResult(a, b)

Add a payload to an accepted join's acknowledgment.

A rejected join remains rejected.