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/channelimport 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)), ]) })}Type safety
Section titled “Type safety”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.
Processes and ordering
Section titled “Processes and ordering”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.
Action
Section titled “Action”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.
Active
Section titled “Active”pub type ActiveMarker for actions valid while a channel is active.
ChildSpecError
Section titled “ChildSpecError”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.
Constructors
Section titled “Constructors”InvalidPattern
Section titled “InvalidPattern”InvalidPattern( pattern: String, reason: topic.TopicError)A handler used an invalid topic pattern.
DuplicatePattern
Section titled “DuplicatePattern”DuplicatePattern(pattern: String)Two handlers registered the same pattern string.
InvalidConfig
Section titled “InvalidConfig”InvalidConfig(reason: beryl.ConfigError)The core beryl.Config failed eager validation.
Closing
Section titled “Closing”pub type ClosingMarker for actions valid while a channel is closing.
Handler
Section titled “Handler”pub type HandlerA 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.
JoinContext
Section titled “JoinContext”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.
JoinResult
Section titled “JoinResult”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.
Message
Section titled “Message”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.
Sender
Section titled “Sender”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.
Functions
Section titled “Functions”accept
Section titled “accept”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.
broadcast
Section titled “broadcast”pub fn broadcast( String, json.Json) -> Action(a)Broadcast to every subscriber of this channel's topic, including this socket.
broadcast_from
Section titled “broadcast_from”pub fn broadcast_from( String, json.Json) -> Action(a)Broadcast to every subscriber of this channel's topic except this socket.
broadcast_presence
Section titled “broadcast_presence”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.
child_spec
Section titled “child_spec”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.
Example
Section titled “Example”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.
discard_reply
Section titled “discard_reply”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.
handler
Section titled “handler”pub fn handler( String, fn(JoinContext(a)) -> JoinResult(b, a)) -> HandlerRegister 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.
notify
Section titled “notify”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.
on_info
Section titled “on_info”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.
on_message
Section titled “on_message”pub fn on_message( JoinResult(a, b), fn(a, Message) -> Next(a)) -> JoinResult(a, b)Handle client messages on this channel's topic.
on_presence
Section titled “on_presence”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/presenceimport 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)})on_terminate
Section titled “on_terminate”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.
presence_track
Section titled “presence_track”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).
presence_untrack
Section titled “presence_untrack”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.
push_presence
Section titled “push_presence”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.
queue_snapshot
Section titled “queue_snapshot”pub fn queue_snapshot(Sender(a)) -> Result(overload.Occupancy, overload.AdmissionError)Read this worker incarnation's queue accounting without waiting for it.
reject
Section titled “reject”pub fn reject(json.Json) -> JoinResult(a, b)Refuse the join, returning reason to the client.
reply_error
Section titled “reply_error”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.
reply_ok
Section titled “reply_ok”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.
with_actions
Section titled “with_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.
with_presence
Section titled “with_presence”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"))]),)with_reply
Section titled “with_reply”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.
