Skip to content

Channels

beryl/channel is the recommended default for applications that serve more than one topic namespace on a socket. It is also the default for a Phoenix style design. Register a list of channel handlers. The layer routes each join, message, server message, and close to the correct channel.

If you are not sure which layer you want, read Choose an API first. If you want one topic family and complete control over routing, use Raw Dispatch instead. It is the core API under the channel layer.

A channel is a topic pattern plus a typed join callback. The join callback receives one context value. It rejects the join or accepts it with private state and callbacks.

src/my_app/room_channel.gleam
import beryl/channel
import gleam/json
/// This channel's private state. Each joined topic has one value.
type State {
State(room_id: String, username: String, sent: Int)
}
/// This channel's server-side message type.
type Note {
Tick(Int)
}
pub fn room() -> channel.Handler {
channel.handler("room:*", fn(context: channel.JoinContext(Note)) {
let state =
State(
room_id: context.topic,
username: context.socket_id,
sent: 0,
)
channel.accept(state)
|> channel.on_message(fn(state: State, message: channel.Message) {
channel.next(
State(..state, sent: state.sent + 1),
[
channel.broadcast_from(message.event, json.int(state.sent + 1)),
],
)
})
|> channel.on_info(fn(state: State, note: Note) {
let Tick(at) = note
channel.next(state, [channel.push("tick", json.int(at))])
})
|> channel.on_terminate(fn(state: State, _reason) {
[channel.broadcast("left", json.string(state.username))]
})
|> channel.with_reply(
json.object([#("room", json.string(context.topic))]),
)
|> channel.with_actions([
channel.broadcast("joined", json.string(state.username)),
])
})
}

State and Note stay private. channel.handler returns a channel.Handler, not a generic value. Unrelated channels can therefore use one handler list.

channel.child_spec takes the same beryl.Config as beryl.child_spec. The config contains the codec, rate limits, presence handle, PubSub, and logging. channel.child_spec also takes the handler table. It returns a beryl.Sockets handle and a child specification for the application supervision tree.

src/my_app.gleam
import beryl
import beryl/transport/server
import beryl/wire
import beryl/channel
import beryl_mist as mist_transport
import gleam/bytes_tree
import gleam/erlang/process
import gleam/http/request
import gleam/http/response
import gleam/otp/static_supervisor
import mist
import my_app/room_channel
pub fn handlers() -> List(channel.Handler) {
[room_channel.room()]
}
pub fn main() -> Nil {
let config =
beryl.config(wire.phoenix_codec())
|> beryl.with_frame_rate(per_second: 35, burst: 70)
|> beryl.with_message_rate(per_second: 30, burst: 60)
let assert Ok(#(sockets, channel_specification)) =
channel.child_spec(config, handlers: handlers())
let assert Ok(_root) =
static_supervisor.new(static_supervisor.OneForOne)
|> static_supervisor.add(channel_specification)
|> static_supervisor.start()
// `sockets` is an ordinary core handle: hand it to a transport, to
// `beryl.broadcast`, and to `beryl.stop`.
let assert Ok(_) =
mist_transport.handler(
sockets,
server.default_config("/socket/websocket"),
handle_http,
)
|> mist.new
|> mist.port(8000)
|> mist.start
process.sleep_forever()
}
fn handle_http(
_req: request.Request(mist.Connection),
) -> response.Response(mist.ResponseData) {
response.new(404)
|> response.set_body(mist.Bytes(bytes_tree.new()))
}

handle_http is the HTTP fallback. mist_transport.handler routes WebSocket upgrades on the configured path to beryl. It sends all other requests to handle_http. The Quick Start shows how to serve pages from this function.

Both APIs use the same runtime, wire codec, presence, abuse controls, and transports. The channel layer supplies only the init and update pair.

child_spec reports eager validation failures as channel.ChildSpecError:

VariantMeaning
InvalidPattern(pattern, reason)A handler pattern is invalid
DuplicatePattern(pattern)The same pattern was registered twice
InvalidConfig(beryl.ConfigError)The core's eager config validation failed

Patterns use beryl topic syntax, such as "room:lobby", "room:*", "document:*:ops", and "*". The layer checks patterns in registration order. The first match owns the topic.

A bigger handlers() returns one entry per channel module:

// A fragment of `handlers()`. See "Starting a channel system" above.
[
// Special-case one topic by putting it ahead of the wildcard.
lobby_channel.lobby(), // "room:lobby"
room_channel.room(), // "room:*"
document_channel.document(), // "document:*"
]

Patterns can overlap. Put "room:lobby" before "room:*" to give the lobby a separate channel. Put specific patterns first. If "room:*" comes first, the lobby handler cannot receive a join. The layer cannot detect this routing error. It rejects duplicate pattern strings because the second handler cannot receive a join. It also rejects unmatched topics with {"reason": "unmatched topic"}.

InvalidPattern contains the core topic.TopicError instead of a string, so you can match on the reason. Name each variant explicitly so the compiler identifies the matches that need an update when the API adds a variant:

Error(channel.InvalidPattern(pattern, topic.EmptyTopic)) ->
panic as { "channel pattern " <> pattern <> " is empty" }
Error(channel.InvalidPattern(pattern, topic.InvalidFormat(detail))) ->
panic as { "invalid channel pattern " <> pattern <> ": " <> detail }

The same rule applies to beryl.InvalidTopicPattern(pattern, reason), which nests the identical topic.TopicError.

child_spec first checks pattern syntax in registration order. It then checks for duplicate strings in the same order. Both checks happen before it builds the supervised child processes.

Phoenix keeps per-channel state in socket.assigns, a map of atoms to untyped terms. A beryl channel keeps a value of its own type, chosen by the channel and known to the compiler:

type State {
State(room_id: String, username: String, sent: Int)
}
// Inside the `join` callback:
channel.accept(State(room_id: context.topic, username: name, sent: 0))

channel.handler keeps the state in its callbacks. Each channel.next result creates the next callbacks with the new state. The layer does not convert the state to Dynamic or use unchecked conversion. One List(Handler) can therefore contain channels with unrelated state types.

The layer keeps one instance per joined topic. A socket joined to room:general and room:random has two independent State values, and the layer prunes an instance when its topic closes. You do not write cleanup code for the state itself.

Channel state is private to one accepted join, not shared by the room. Two clients on room:general get two State values, even when they join the same topic. Use channel state for per-join data such as the current user, transient UI state, or a typed sender for that join.

When several joins must share one mutable value or one writer, move that state into an application-owned process or store and capture its handle in the handler closure. Add a room-wide or domain owner only when you have a real shared invariant to protect.

The examples show two common choices. live_poll/store.gleam keeps one in-memory actor that serializes get, vote, and close calls with join and leave updates. Callers that use get, vote, or close wait for that actor before they continue. collab_document/document_store.gleam also keeps one in-memory actor, but merge_state only enqueues a merge and returns. OR-Map merges still converge, but another process can read before that merge runs, so there is no cross-process read-after-write guarantee.

Do not shard these example stores unless measurement shows real contention. Start with one owner for the invariant you need.

Presence is not that owner. Use presence for ephemeral online state, not for locks, authorization, durable membership, or atomic capacity checks. See Presence.

The join callback receives a channel.JoinContext(info):

FieldTypeDescription
socket_idStringUnique id of the socket that is joining
seedsocket.ConnectSeedRequest data the transport assembled before the upgrade: path, query, headers, and any on_connect metadata
selfchannel.Sender(info)This channel instance's own typed sender
topicStringConcrete topic being joined
parametersList(String)Wildcard captures in pattern order; exact patterns receive []
payloadDynamicRaw client join payload

The layer builds one JoinContext for each join. It is not the core socket.ConnectInfo. The layer owns the socket model and message type. This lets channels keep private state and private message types. Each join receives the required connection data and a sender for that join.

channel.notify(sender, message) delivers message to that channel's on_info callback with its type intact:

pub type Note {
Tick(Int)
}
// Call from a timer, application actor, or HTTP handler:
let assert Ok(Nil) = channel.notify(sender, Tick(1))

This mechanism does not use casts. notify keeps the value in a typed function and sends it to the worker process of that join. Only that join can read the value during delivery. No mailbox stores the typed value between turns.

Because notify admits work to the channel's worker queue, it reserves worker capacity before publication. Gleam does not guarantee how many bytes a sealed function retains, so limit application payload sizes as well as queue item counts.

A sender applies only to the join that produced it. notify returns an admission Result. Queue saturation and a closed or unavailable worker return an error. Ok(Nil) confirms admission, not callback completion. Accepted messages retain their queue order; the runtime does not combine them.

A long-lived process can keep a sender but cannot use it to reach a different join. Handle errors when the target is gone or full; do not build an unlimited retry queue. See overload handling.

Use notify to schedule a later turn, including from another process. Put work that must occur during the join in actions after a join.

An action is one operation on this channel's own topic. Constructor functions return one action. Put actions in a list in the order clients should observe them:

ActionEffect
push(event, payload)Server-initiated message to this socket on this topic
broadcast(event, payload)To every subscriber of this topic, including this socket
broadcast_from(event, payload)To every subscriber except this socket
reply_ok(reply, payload)Success reply when Message.reply is Some; no effect for None
reply_error(reply, payload)Error reply with the same optional-ref behavior
discard_reply(reply)Release an intentionally unanswered reply handle without a wire reply
presence_track(key, meta)Track this socket under key and emit the presence_diff join
presence_untrack(key)Untrack and emit the presence_diff leave
push_presence(event, encode)Presence snapshot for this topic, to this socket
broadcast_presence(event, encode)Presence snapshot for this topic, to every subscriber

No action names a topic. Each action applies to the channel that returned it. To send across topics, use the external APIs described in When to use raw dispatch or another process.

Presence actions need a presence handle on the config (beryl.with_presence_handle); without one they are dropped with a warning, exactly as the equivalent core effects are.

For the standard Phoenix presence flow, add channel.with_presence to an accepted join instead of assembling the track and snapshot actions yourself:

channel.accept(state)
|> channel.with_presence(key: state.username, meta: meta(state))

This is shorthand for presence_track followed by push_presence with the Phoenix event name and encoder. It adds no new lifecycle behavior: diff delivery and automatic cleanup also work with the explicit actions. See the equivalent code in Add presence to a channel.

To react on the server, register channel.on_presence. It receives a presence.Snapshot first, then topic-scoped presence.Changed events, and returns the same Next(state) as other callbacks. Observation is independent of tracking. See React to presence changes for the initial roster, update, and failure contracts.

The runtime applies channel actions in list order. Each action maps to one core socket.Effect. An asynchronous presence effect can pause this socket while other sockets continue. The remaining actions resume after the effect completes.

The encode callbacks of push_presence and broadcast_presence run when the action is applied, so a snapshot already reflects any presence_track or presence_untrack earlier in the same list:

[
channel.presence_track(state.username, meta(state)),
channel.broadcast_presence("presence_list", presence_helpers.encode_users),
]

channel.with_actions attaches ordered actions to an accepted join. They are emitted with the acknowledgment and applied strictly after it:

channel.accept(state)
|> channel.with_reply(reply)
|> channel.with_actions([
channel.presence_track(state.username, meta(state)),
channel.broadcast("new_msg", joined_message(state)),
channel.broadcast_presence("presence_list", encode_users),
])

This ordering has two consequences:

  • The acknowledgment always reaches the wire first. The socket is already subscribed to the topic when the join's own actions run, so a push cannot overtake its own join reply.
  • Effect order is per socket, not a cross-socket transaction. If an action starts asynchronous presence work, the runtime may process other sockets while this one waits. Use application-owned synchronous state for an atomic capacity reservation.

with_actions appends, so it composes with itself, and it returns channel.reject results unchanged: a refused join has no topic to act on.

Use join actions instead of sending notify to the same channel from join. notify schedules a later input. Join actions stay directly after the join acknowledgment.

CallbackInputSignature
on_messageA client message on this topicfn(state, channel.Message) -> Next(state)
on_infoA notify addressed to this joinfn(state, info) -> Next(state)
on_presenceThis topic's initial roster and later changesfn(state, presence.Event) -> Next(state)
on_terminateThis channel ending, for any reasonfn(state, socket.StopReason) -> List(Action(Closing))

channel.accept(state) stays joined until you add callbacks with the on_* builders. Unhandled messages and server-side notifications have no effect. For raw binary frames, use Raw Dispatch.

A channel.Message has event, the raw payload as Dynamic, and reply: Option(socket.ReplyRef), which is present only when the client asked for a reply. Store context.topic in channel state when a callback needs it:

|> channel.on_message(fn(state: State, message: channel.Message) {
case message.event {
"new_msg" ->
channel.next(state, [
channel.broadcast("new_msg", body(message.payload)),
channel.reply_ok(message.reply, json.object([])),
])
"typing" ->
channel.next(state, [
channel.broadcast_from("typing", json.object([])),
])
_ -> channel.stay(state)
}
})

Every callback answers with a channel.Next(state):

ResultBehavior
next(state, actions)Stay joined, applying the active-phase actions in order
stay(state)Stay joined with no actions
close(actions)Apply active-phase actions, then leave this channel

close applies its actions first and then closes the topic, so a farewell broadcast still reaches the topic's subscribers.

on_terminate runs once for each accepted join. It runs after phx_leave, close([]), socket disconnect, heartbeat timeout, and beryl.stop. A rejected join does not create an instance, so it does not run on_terminate.

The runtime converts the returned actions to effects during the turn that closes the topic, right after the instance has been removed. Closing-phase lists can contain broadcast, broadcast_from, presence_untrack, and broadcast_presence. Active-only pushes, replies, presence tracking, and presence pushes do not type-check in on_terminate.

Put a leave announcement and updated roster here instead of sending them from code outside the channel:

|> channel.on_terminate(fn(state: State, _reason) {
[
channel.broadcast("new_msg", departure(state)),
channel.presence_untrack(state.username),
channel.broadcast_presence("presence_list", encode_users),
]
})

Put presence_untrack before broadcast_presence. The runtime encodes the snapshot after it applies the untrack, so the roster is current. The runtime also removes presence entries when the topic closes. That removal occurs after Closed, not before your actions.

sequenceDiagram
  participant Client
  participant Router as beryl/channel router
  participant Ch as your channel
  Client->>Router: phx_join "room:lobby"
  Router->>Router: first matching pattern wins
  Router->>Ch: join(JoinContext)
  Ch-->>Router: accept(state) |> on_message(..) |> with_reply(reply)
  Router-->>Client: phx_reply ok, then the join's actions
  Client->>Router: event on "room:lobby"
  Router->>Ch: on_message(state, Message)
  Ch-->>Router: next(state', actions)
  Router-->>Client: actions, in order
  Client->>Router: phx_leave / disconnect
  Router->>Router: remove the instance
  Router->>Ch: on_terminate(state, reason)
  Ch-->>Router: closing-phase actions

An app panic does not stop the runtime. The core limits the effect based on where the panic occurred:

Panic inEffect
joinThat join is rejected; the socket survives
on_messageThe runtime closes that topic with phx_error and runs on_terminate. Other topics on the socket continue
on_infoThe runtime closes that topic with phx_error and runs on_terminate. Other topics on the socket continue
on_terminateThe runtime logs the panic and completes the close without its actions. It still runs termination actions for sibling channels

Each topic runs in its own worker process. A panic in a callback keeps the state from before that callback, so on_terminate still sees it, and once on_terminate finishes the worker stops, so a sender for that join delivers nothing. If the worker stops unexpectedly instead, the runtime closes the topic with phx_error without running on_terminate (it held the channel state), and the client must rejoin; the runtime also kills a worker that does not finish its queued work and on_terminate within five seconds, then completes the close without the termination actions. During a graceful beryl.stop, this bound is one second per worker. Several blocked workers on one socket can still exceed the total stop budget.

Crash isolation stops at the socket for other faults: a fault in the socket actor loses only that socket and its workers, while a router crash loses every socket on that beryl.Sockets handle. See Socket Processes & Restarts.

channel.child_spec mirrors beryl.child_spec: it returns a beryl.Sockets handle and a child specification for applications that own their supervision tree, reporting the same validation errors as channel.ChildSpecError (see the table in Starting a channel system).

The channel layer uses the core child processes, restart policy, and beryl.stop behavior. See the Supervision guide. The application must start and supervise PubSub, presence, and group actors, and pass their handles to beryl.Config.

After a runtime crash the handle keeps working for new connections, but every live channel instance is gone: clients reconnect and rejoin, and each join runs afresh.

The compiler guides the package migration:

  1. Remove the beryl_channels dependency from gleam.toml.
  2. Replace imports of beryl_channels/channel with beryl/channel.
  3. Remove the root beryl_channels import and call channel.child_spec(config, handlers:).
  4. Match startup errors as channel.InvalidPattern, channel.DuplicatePattern, and channel.InvalidConfig.

Handler construction, private typed state and server messages, action ordering, and callback behavior are unchanged. Run gleam check to find every remaining old import or qualified error name.

When to use raw dispatch or another process

Section titled “When to use raw dispatch or another process”

The layer handles one topic at a time, so each action applies only to the channel that returned it: a room:general channel cannot broadcast on lobby. Cross-topic publishing goes through the external Sockets APIs (beryl.broadcast, beryl.broadcast_from, or a beryl/group actor) using the handle channel.child_spec returns. Because handlers are built before that handle exists, a channel cannot capture it directly; instead, keep the handle in a small actor that exposes a publish(topic, event, payload) function, like Phoenix's Endpoint.broadcast/3, and bind it after child_spec returns and before the transport starts accepting connections. examples/showcase does exactly this for its lobby room list. Channels do not share state with each other either — anything two channels both need (a document store, presence handle, groups actor) is a dependency you capture in the handler closures when you build the table.

The layer owns the socket-level model and message type, so pick raw dispatch or the channel layer per socket endpoint; mixing hand-written update logic into a channel system is not supported. Raw binary frames require raw dispatch, though binary frames the configured codec decodes into normal events still reach on_message. beryl/bridge targets the core sender, not a channel's: bridge.start(to:, with:) wants a beryl/socket.Sender, which a channel never sees, so adapt an existing actor's messages into on_info by forwarding them yourself with channel.notify(context.self, ..) — the typed sender is safe to hold and is dropped after the join ends.

Each channel runs in its own worker, so a slow callback delays only its own topic while other topics on the socket continue; the socket actor still waits up to five seconds for join, so keep it short and return long-running results through channel.notify. beryl does not define an order between topics: actions for one topic keep their order, but replies and pushes for different topics on one socket can interleave.