Skip to content

beryl/socket

Types for building app-side dispatch systems with beryl.child_spec.

With app-side dispatch the application owns routing: beryl delivers every wire event for a socket to one update function, and the function returns the next model plus a list of Effects for beryl to apply. There are no channel modules, no registry, and no type erasure. Each socket has a single model and a single message type.

Effects are applied strictly in list order, and every frame for a socket is written by that socket's own runtime actor — so list order is wire order. An AcceptJoin followed by a Push in the same list is guaranteed to arrive as the join acknowledgment first and the push second.

Most effects are applied in one actor turn. PresenceTrack and PresenceUntrack are the exception: they are applied by the presence actor. beryl holds the rest of the list and every later input for that socket until the mutation has been applied. It then continues exactly where it left off. The visible order is unchanged (a PushPresence after a PresenceTrack still sees the track), and the socket's own inputs still arrive in the order the client sent them. No other socket, broadcast, or heartbeat waits on that mutation. Those continue, so a broadcast from elsewhere may arrive between two effects from this socket.

pub type ConnectInfo(a) {
ConnectInfo(
socket_id: String,
seed: ConnectSeed,
self: Sender(a)
)
}

Everything the app's init receives when a socket connects.

pub type ConnectSeed {
ConnectSeed(
path: String,
query: List(#(String, String)),
headers: List(#(String, String)),
metadata: List(#(String, String))
)
}

Connection metadata that the transport builds before the WebSocket upgrade.

The app's init function receives it through ConnectInfo.

pub type Effect {
AcceptJoin(
ref: JoinRef,
reply: option.Option(json.Json)
)
RejectJoin(
ref: JoinRef,
reason: json.Json
)
ReplyOk(
ref: ReplyRef,
payload: json.Json
)
ReplyError(
ref: ReplyRef,
payload: json.Json
)
DiscardReply(ref: ReplyRef)
Push(
topic: String,
event: String,
payload: json.Json
)
Broadcast(
topic: String,
event: String,
payload: json.Json
)
BroadcastFrom(
topic: String,
event: String,
payload: json.Json
)
PresenceTrack(
topic: String,
key: String,
meta: json.Json
)
PresenceUntrack(
topic: String,
key: String
)
PushPresence(
topic: String,
event: String,
encode: fn(List(presence.PresenceEntry)) -> json.Json
)
BroadcastPresence(
topic: String,
event: String,
encode: fn(List(presence.PresenceEntry)) -> json.Json
)
KickTopic(topic: String)
}

One update may return several effects, applied strictly in list order (see the module docs for the ordering guarantee).

AcceptJoin(
ref: JoinRef,
reply: option.Option(json.Json)
)

Accept a pending join. This subscribes the socket to the topic and sends the join acknowledgment with an optional reply payload. The effect is valid only while the Join input's ref is pending.

RejectJoin(
ref: JoinRef,
reason: json.Json
)

Reject a pending join with an error payload.

ReplyOk(
ref: ReplyRef,
payload: json.Json
)

Reply successfully to a client message ref.

ReplyError(
ref: ReplyRef,
payload: json.Json
)

Reply with an error to a client message ref.

DiscardReply(ref: ReplyRef)

Discard a client message ref without sending a reply.

Use this when the application intentionally will not answer a referenced message. It releases the retained socket capacity and permits later reuse of the same wire ref. A completed or stale ref has no effect.

Push(
topic: String,
event: String,
payload: json.Json
)

Push a server-initiated message to this socket on a joined topic. The runtime drops pushes to topics that this socket has not joined and logs a warning. Put a Push after its topic's AcceptJoin.

Broadcast(
topic: String,
event: String,
payload: json.Json
)

Broadcast to every subscriber of a topic (including this socket, when joined). Distributed via PubSub when configured.

BroadcastFrom(
topic: String,
event: String,
payload: json.Json
)

Broadcast to every subscriber of a topic except this socket.

PresenceTrack(
topic: String,
key: String,
meta: json.Json
)

Track this socket's presence under a key in a topic and broadcast the corresponding presence_diff join. This effect requires a presence handle on the config (beryl.with_presence_handle). Without a handle, the runtime drops the effect and logs a warning.

Tracking an existing key replaces the previous entry atomically. The key is never absent during the replacement. One presence_diff contains both the leave and the join. Later effects wait for the mutation, as described in the module documentation. Other sockets do not wait.

PresenceUntrack(
topic: String,
key: String
)

Untrack a presence previously tracked with PresenceTrack and broadcast the corresponding presence_diff leave. When the topic closes, the runtime removes the remaining tracked keys in one batch and produces one aggregate leave diff. Later effects wait for the mutation, as described in the module documentation. Other sockets do not wait. This effect requires a presence handle (beryl.with_presence_handle). Without a handle, the runtime drops the effect and logs a warning.

PushPresence(
topic: String,
event: String,
encode: fn(List(presence.PresenceEntry)) -> json.Json
)

Push a presence snapshot for a topic to this socket. A payload built inside update sees presence from before this effects list. In contrast, encode runs when the effect is applied. It runs after earlier PresenceTrack and PresenceUntrack effects in the same list. The entries therefore include those changes. This effect requires a presence handle (beryl.with_presence_handle). Without a handle, the runtime drops the effect and logs a warning. Like Push, the runtime drops it if the topic is not joined.

BroadcastPresence(
topic: String,
event: String,
encode: fn(List(presence.PresenceEntry)) -> json.Json
)

Broadcast a presence snapshot for a topic to all its subscribers, with the same apply-time encode semantics as PushPresence. Order it after the PresenceTrack/PresenceUntrack it should reflect. This effect requires a presence handle (beryl.with_presence_handle). Without a handle, the runtime drops the effect and logs a warning.

KickTopic(topic: String)

Close this socket's subscription to a topic. The topic receives Closed(topic, Shutdown). If the codec has a close encoder, the client also receives its terminal frame.

pub type Input(a) {
Join(
topic: String,
payload: dynamic.Dynamic,
ref: JoinRef
)
Message(
topic: String,
event: String,
payload: dynamic.Dynamic,
ref: option.Option(ReplyRef)
)
Binary(
topic: String,
data: BitArray
)
Closed(
topic: String,
reason: StopReason
)
Info(a)
}

Everything the runtime delivers to the app's update function.

Join(
topic: String,
payload: dynamic.Dynamic,
ref: JoinRef
)

A client asked to join a topic. Return an AcceptJoin or RejectJoin effect. The runtime rejects a Join that is unanswered at the end of the update turn.

Message(
topic: String,
event: String,
payload: dynamic.Dynamic,
ref: option.Option(ReplyRef)
)

A client message on a joined topic. ref is present for messages that expect a reply.

Binary(
topic: String,
data: BitArray
)

A binary frame on a joined topic (codecs without a binary decoder deliver the raw frame once per joined topic).

Closed(
topic: String,
reason: StopReason
)

A joined topic ended because of a client leave, kick, crash, or socket close. The runtime sends this input on every exit path. Use it to remove per-topic state from the model. Frames pushed to the closing topic are dropped; broadcasts still reach the topic's remaining subscribers.

Info(a)

A typed server-side message, sent via the socket's Sender (see ConnectInfo.self and notify).

pub type JoinRef

A pending join correlation handle.

Pass it back in AcceptJoin or RejectJoin. A join ref is valid only for its pending join. It carries a unique runtime token. A delayed completion for an older same-topic join cannot answer a replacement or retry.

pub type Next(a) {
Next(
model: a,
effects: List(Effect)
)
Stop(reason: StopReason)
}

The result of one update call.

It contains the next model and effects, or an instruction to stop the socket.

Next(
model: a,
effects: List(Effect)
)

Continue with the given model, applying the effects in order.

Stop(reason: StopReason)

Tear down the socket: every joined topic receives a Closed input, configured terminal frames are sent, and the transport connection is closed.

pub type ReplyRef

A client message reply correlation handle.

Pass it back in ReplyOk, ReplyError, or DiscardReply. You can store reply refs in the model and answer them in a later update turn, for example after an asynchronous lookup. They do not expire. They are single-use and remain valid only while the topic instance that received the message stays open.

pub type Sender(a)

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

Get this handle from ConnectInfo.self in init. Any process can call notify with it. The socket's update function receives the message as an Info event. This typed send does not erase the message type.

pub type StopReason {
Normal
Shutdown
HeartbeatTimeout
Errored(String)
AdmissionRejected(overload.AdmissionError)
}

Why a socket or topic is stopping.

The runtime delivers this reason in Closed inputs, and Stop accepts it. The variants are exhaustive; match each reason explicitly. Adding a variant affects API compatibility and requires updating exhaustive matches.

Normal

Normal shutdown (client left or disconnected cleanly).

Shutdown

Server-initiated shutdown (system stop, KickTopic).

HeartbeatTimeout

The client failed to send a heartbeat within the configured timeout.

Errored(String)

An error stopped the socket or topic. The name Errored prevents an unqualified import from shadowing the prelude's Result Error constructor.

AdmissionRejected(overload.AdmissionError)

A local queue or callback result exceeded its admission budget.

pub fn empty_seed() -> ConnectSeed

Return an empty connect seed for tests and transports with no request data.

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

Send a typed server-side message to a socket.

The socket's update function receives Info(message). The runtime reports an admission error if the socket is closed or its queue is full. Ok means admitted, not handled by the application's callback.

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

Read socket queue accounting without waiting for its callback.

pub fn reply_ok(
option.Option(ReplyRef),
json.Json
) -> List(Effect)

Return a ReplyOk effect when the client supplied a ref.

Message inputs carry Option(ReplyRef) (refless messages expect no reply) while the ReplyOk effect demands a ReplyRef, so every handler that replies conditionally needs this check. This function returns no effects when the client did not supply a ref.