Skip to content

beryl/bridge

Forward an external OTP actor's message stream to a socket through its Sender.

A common pattern is a long-lived domain actor (e.g. a per-document session) that emits updates which need to be pushed to a connected socket. Wiring this up by hand requires per-socket boilerplate: spawn a forwarder process holding a Subject, subscribe it to the domain actor, translate each message and call socket.notify, then tear the process down when the socket closes.

bridge packages that plumbing into a single helper. Start a bridge in your app's init (or when a topic is joined), store the handle in the socket model, subscribe the returned Subject to your domain actor, and stop the bridge when the socket or topic closes.

Calling stop when the socket/topic ends is required for cleanup: the forwarder monitors the process that started it only as a backstop for that owner's death, not as a per-topic lifecycle. A bridge whose stop is never called keeps running until its owner exits.

import beryl/bridge.{type Bridge}
import beryl/socket.{type ConnectInfo}
// Messages emitted by your domain actor.
pub type DocEvent {
Updated(version: Int)
}
// Your app's server-side message type, delivered to `update` as `Info`.
pub type Message {
DocUpdated(version: Int)
}
fn init(info: ConnectInfo(Message)) -> #(Model, List(socket.Effect)) {
// Forward each DocEvent to this socket as an `Info(Message)` event.
let assert Ok(bridge_handle) =
bridge.start(to: info.self, with: fn(event: DocEvent) {
let Updated(version) = event
DocUpdated(version)
})
// Subscribe the domain actor to the bridge's subject.
doc.subscribe(doc_actor, bridge.subject(bridge_handle))
#(Model(bridge: bridge_handle), [])
}
// Stop the bridge when the socket closes (e.g. from a `Closed` event).
bridge.stop(model.bridge)
pub type Bridge(a)

A handle to a running bridge forwarder.

message is the type that the external actor sends to the bridge's Subject. Get the subject with subject. Call stop to stop the forwarder.

pub type StartError {
ForwarderUnavailable
}

Why a bridge failed to start.

ForwarderUnavailable

The forwarder did not report its subjects within handshake_timeout_ms. It failed to spawn or reach readiness, and any timed-out child is cleaned up before start returns.

pub fn pid(Bridge(a)) -> process.Pid

Return the forwarder process ID.

Use this value for diagnostics and supervision. Most callers need only subject and stop.

pub fn start(
to: socket.Sender(a),
with: fn(b) -> a
) -> Result(Bridge(b), StartError)

Start a bridge from an external Subject to a socket's update function.

The returned Bridge owns a new forwarder process. Pass subject(bridge) to the external or domain actor. The forwarder maps each received value with transform and sends it as an Info event through socket.notify(sender, transform(value)).

Use transform to translate the domain message into your app's server-side message type (the Info payload). If no translation is needed, pass the identity function fn(value) { value }.

Always call stop when the owning socket or topic ends. The forwarder also monitors the calling process. This monitor stops the forwarder if the owner dies without calling stop, but it does not track topic lifecycles.

pub fn stop(Bridge(a)) -> Nil

Stop the bridge's forwarder.

Call this when the owning socket or topic ends. It is safe to call more than once and after the forwarder has already exited.

pub fn subject(Bridge(a)) -> process.Subject(a)

Return the Subject that receives the external actor's stream.

Give this subject to the domain actor, for example as its subscriber. Each value is then sent to the bridged socket.