Skip to content

Use a WebSocket transport

beryl provides WebSocket transports for Mist and Ewe. Both transports use the same beryl runtime, configuration, authentication, origin checks, connection limits, and wire codecs.

Use mist_transport.upgrade to add WebSocket support:

import beryl
import beryl/transport/server
import beryl_mist as mist_transport
import gleam/bytes_tree
import gleam/http/request
import gleam/http/request.{type Request}
import gleam/http/response
import mist
fn handle_request(
http_request: Request(mist.Connection),
channels: beryl.Sockets,
) -> response.Response(mist.ResponseData) {
// Upgrade /socket/websocket requests to WebSocket
use <- mist_transport.upgrade(
http_request,
channels,
server.default_config("/socket/websocket"),
)
// Non-WebSocket requests fall through here
case request.path_segments(http_request) {
[] -> response.new(200) |> response.set_body(mist.Bytes(bytes_tree.new()))
_ -> response.new(404) |> response.set_body(mist.Bytes(bytes_tree.new()))
}
}

The upgrade function checks the request path. It upgrades a matching request and connects it to the beryl runtime.

Use ewe_transport.handler to combine WebSocket upgrades and regular HTTP routing in one Ewe listener:

import beryl
import beryl/transport/server
import beryl_ewe as ewe_transport
import ewe
import gleam/http/request.{type Request}
import gleam/http/response.{type Response}
pub fn start_ewe(sockets: beryl.Sockets) -> Nil {
let assert Ok(_) =
ewe_transport.handler(
sockets,
server.default_config("/socket/websocket"),
fn(
http_request: Request(ewe.Connection),
) -> Response(ewe.ResponseBody) {
case request.path_segments(http_request) {
[] ->
response.new(200)
|> response.set_body(ewe.TextData("Hello!"))
_ ->
response.new(404)
|> response.set_body(ewe.Empty)
}
},
)
|> ewe.new
|> ewe.listening(port: 8000)
|> ewe.start
Nil
}

handler sends matching WebSocket requests through the beryl admission pipeline and sends every other request to the fallback function. Use ewe_transport.upgrade instead when an existing Ewe handler already controls the fallback.

Both transports work with both beryl APIs. A handle from channel.child_spec is the same beryl.Sockets type as one from beryl.child_spec, so the setup is identical for both.

Use with_on_connect to authenticate a connection before the upgrade. It is similar to Phoenix UserSocket.connect/3. The hook runs once for each socket before any channel join. It can reject the connection.

The example below uses Mist request types. With Ewe, use Request(ewe.Connection) in the callback and pass the resulting configuration to ewe_transport.handler or ewe_transport.upgrade. The configuration builders and rejection behavior are the same.

let websocket_config =
server.default_config("/socket/websocket")
|> server.with_on_connect(fn(http_request: Request(mist.Connection)) {
// Check auth token, session, etc.
case validate_token(http_request) {
Ok(_user) -> Ok([]) // Allow; no connect metadata
Error(_) -> Error(server.ConnectRejected) // Reject with 403
}
})
use <- mist_transport.upgrade(http_request, channels, websocket_config)

Return Error(server.ConnectRejected) to send HTTP 403 before the WebSocket upgrade. See Reject a connection during authentication for the client error. See Authentication failures for diagnostic steps.

Block cross-site WebSocket hijacking (CSWSH)

Section titled “Block cross-site WebSocket hijacking (CSWSH)”

Browsers include cookies on WebSocket handshakes. If your socket authentication uses cookies, a malicious site can open a WebSocket to your application from a victim's browser unless you validate the Origin header. This is Cross-Site WebSocket Hijacking (CSWSH).

Use with_allowed_origins to allow only your application origins. Values match the full Origin header exactly: scheme, host, and port when present.

let websocket_config =
server.default_config("/socket/websocket")
|> server.with_allowed_origins(["https://app.example.com"])
|> server.with_on_connect(fn(http_request: Request(mist.Connection)) {
validate_cookie_session(http_request)
})

Requests with missing or non-matching origins are rejected with HTTP 403 before the WebSocket handshake. Without an explicit allow-list, the default SameOrigin policy rejects cross-site browser handshakes while allowing non-browser clients that omit Origin.

If you cannot use an origin allow-list, avoid cookie-based WebSocket authentication. Use a token passed explicitly to on_connect and reject invalid tokens before upgrading.

on_connect accepts or rejects the upgrade. The transport puts the request path, query parameters, and headers in a ConnectSeed, and your init function receives it as ConnectInfo.seed:

FieldTypeDescription
pathStringThe request path the client connected to.
queryList(#(String, String))Parsed query parameters.
headersList(#(String, String))Request headers.
metadataList(#(String, String))Whatever on_connect returned.

on_connect does not return Ok(Nil). It returns Ok(metadata), a list of string pairs that reaches init as seed.metadata. Resolve the identity once during the handshake and return it as metadata. Then init can read it instead of decoding the request again:

let websocket_config =
server.default_config("/socket/websocket")
|> server.with_on_connect(fn(http_request: Request(mist.Connection)) {
// Validate once; reject the whole connection on failure.
case validate_token(http_request) {
Ok(user_id) -> Ok([#("user_id", user_id)])
Error(_) -> Error(server.ConnectRejected) // Reject with 403
}
})
// init reads what on_connect already resolved. It does not decode again.
beryl.child_spec(
config,
init: fn(info: socket.ConnectInfo(Message)) {
let user_id =
list.key_find(info.seed.metadata, "user_id")
|> result.unwrap("anonymous")
#(Model(user_id: user_id), [])
},
update: update,
)

Return Ok([]) when there is nothing to pass on. The list keeps the order on_connect produced and keeps duplicate keys, so list.key_find returns the first pair for a key. Values are strings only: encode anything richer, such as a list of roles, into one. Transports never log metadata values, but the seed reaches every join callback on that socket, so put an identity there rather than a secret.

With the channel layer, the same seed arrives in every handler's join callback as channel.JoinContext.seed, metadata included; there is no app-level init. See Authentication with beryl/channel.

Pass wire.phoenix_codec() to beryl.config to use the Phoenix JSON array format:

[join_ref, ref, topic, event, payload]

Applications can pass a custom codec to beryl.config(codec). The codec can use another text or binary message format. The transport sends each outbound frame as the type that the codec returns.

wire.phoenix_codec() uses beryl's Phoenix wire implementation and adds no dependencies. The public beryl/wire/codec.Codec API and wire format are stable, so applications can supply a codec for another message format.

FieldJSON typeDescription
join_refstring | nullReference from the join (for reply routing)
refstring | nullUnique message reference (for reply matching)
topicstringTopic name (e.g., "room:lobby")
eventstringEvent name (e.g., "phx_join", "new_message")
payloadanyJSON payload
EventDirectionPurpose
phx_joinClient -> ServerJoin a channel
phx_leaveClient -> ServerLeave a channel
heartbeatClient -> ServerKeepalive ping
phx_replyServer -> ClientReply to a client message
phx_errorServer -> ClientError notification
phx_closeServer -> ClientChannel closed

Client sends:

["1", "1", "room:lobby", "phx_join", {"user": "alice"}]

Server replies:

["1", "1", "room:lobby", "phx_reply", {"status": "ok", "response": {}}]
  1. Client connects via WebSocket to the configured path
  2. The on_connect callback runs, if configured. Rejection returns HTTP 403.
  3. The transport connection process builds the ConnectSeed, generates a unique socket ID, and starts a socket actor. The router admits and monitors that actor, which runs your init.
  4. The client sends phx_join messages to subscribe to topics. Each one arrives at update as a Join event.
  5. Messages are routed through the runtime to update as Message events
  6. On disconnect, update receives a Closed event for every joined topic

Clients should send periodic heartbeat messages to stay connected:

[null, "ref_123", "phoenix", "heartbeat", {}]

Configure heartbeat timing in the beryl config:

let config =
beryl.config(wire.phoenix_codec())
|> beryl.with_heartbeat(
timeout_ms: 60_000, // Server evicts after 60s silence (must be >= 2)
)

Protect against flood attacks with built-in rate limiting:

let config =
beryl.config(wire.phoenix_codec())
|> beryl.with_frame_rate(per_second: 150, burst: 300)
|> beryl.with_message_rate(per_second: 100, burst: 200)
|> beryl.with_join_rate(per_second: 5, burst: 10)
|> beryl.with_channel_rate(per_second: 50, burst: 100)
LimiterApplies toEnforced at
frame_ratePer connection, all complete framesTransport, before decode
message_ratePer socket, decoded non-join trafficRuntime
join_ratePer socket, joinsRuntime
channel_ratePer socket+topicRuntime
topic_ratesTopics matching a pattern; overrides channel_rateRuntime

Frame and message buckets are independent. Malformed frames and joins consume frame tokens; joins do not consume message tokens.

Cap both the connection-attempt rate and the number of concurrent connections a single client IP may hold. Both controls default to unlimited.

let config =
beryl.config(wire.phoenix_codec())
|> beryl.with_connection_rate_per_ip(per_second: 2, burst: 5)
|> beryl.with_max_connections_per_ip(max_connections: 5)

When a peer reaches either limit, the transport rejects the upgrade with 429 Too Many Requests before the handshake. A closed connection releases its concurrent capacity. The per-IP rate bucket remains after reconnects and app runtime restarts. Thus, a reconnect does not provide a new rate allowance.

Both controls use the socket peer IP, which is the TCP address that Mist or Ewe accepts. beryl does not trust forwarded headers such as X-Forwarded-For. Clients can forge these headers and bypass the limit.

This affects beryl when it runs behind a reverse proxy or load balancer (nginx, HAProxy, a cloud LB, etc.): every connection arrives from the proxy's IP, so a per-IP limit sees all clients as one address and throttles them together. In that setup:

  • Enforce per-IP limits at the proxy layer, where the real client IP is known, or
  • Terminate connections directly (no intermediary) if you want beryl's built-in per-IP limit to apply to individual clients.

beryl does not have a trusted-proxy option that reads a client IP from a forwarded header only when the immediate peer is trusted. Treat X-Forwarded-For as untrusted input.