beryl/bridge
Bridge - Forward an external OTP actor's message stream to a socket channel.
A common pattern is a long-lived domain actor (e.g. a per-document session)
that emits updates which need to be pushed to each joined 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 beryl.send_info, then tear the process down on
terminate.
bridge packages that plumbing into a single helper. Start a bridge inside
a channel join, store the handle in the socket assigns, subscribe the
returned Subject to your domain actor, and stop the bridge in terminate.
Calling stop from terminate is required for cleanup: channels are
dispatched by a shared coordinator rather than one process per channel, so
the process the forwarder monitors is the coordinator itself — the monitor
is a backstop for coordinator death, not a per-channel lifecycle. A bridge
whose stop is never called keeps running until the coordinator exits.
Example
Section titled “Example”import beryl.{type RegisteredChannel}import beryl/bridge.{type Bridge}import beryl/channel.{type Channel}import beryl/socketimport gleam/option.{None}import my_app/doc
// Messages emitted by your domain actor.pub type DocEvent { Updated(version: Int)}
pub type Assigns { Assigns(bridge: Bridge(DocEvent))}
// `registered_channel` is the handle returned by `beryl.register`.fn new_channel( registered_channel: RegisteredChannel(Assigns, DocEvent), doc_actor,) -> Channel(Assigns, DocEvent) { channel.new(fn(topic, _payload, socket) { // Forward each DocEvent to this socket's `handle_info` callback. let b = bridge.start( channel: registered_channel, socket_id: socket.id(socket), topic: topic, with: fn(event) { event }, ) // Subscribe the domain actor to the bridge's subject. doc.subscribe(doc_actor, bridge.subject(b)) channel.JoinOk(reply: None, socket: socket.set_assigns(socket, Assigns(b))) }) |> channel.with_terminate(fn(_reason, socket) { bridge.stop(socket.get_assigns(socket).bridge) })}Bridge
Section titled “Bridge”A handle to a running bridge forwarder process.
message is the type emitted by the external actor and received on the
bridge's Subject. Obtain that subject with subject to wire it up to a
domain actor, and call stop to tear the forwarder down.
pub type Bridge(a)Functions
Section titled “Functions”The forwarder process id.
Exposed for diagnostics and supervision; you normally only need subject
and stop.
pub fn pid(Bridge(a)) -> process.PidStart a bridge that forwards values from an external Subject to a socket's
channel as handle_info messages.
The returned Bridge owns a freshly spawned forwarder process. Pass
subject(bridge) to the external/domain actor so it delivers its stream to
the forwarder; each received value is mapped with transform and delivered
via beryl.send_info(channel, socket_id, topic, transform(value)).
Use transform to translate the domain message into whatever your channel's
handle_info expects. If no translation is needed, pass the identity
function fn(value) { value }.
Always call stop from your channel's terminate — that is the only
per-channel cleanup. The forwarder also monitors the calling process, but
because channel callbacks run inside the shared coordinator, that monitor
fires only if the coordinator itself dies; it does not detect an
individual channel ending.
pub fn start( channel: beryl.RegisteredChannel(a, b), socket_id: String, topic: String, with: fn(c) -> b) -> Bridge(c)Stop the bridge's forwarder process.
Call this from your channel's terminate callback. It is safe to call more
than once and after the forwarder has already exited.
pub fn stop(Bridge(a)) -> Nilsubject
Section titled “subject”The Subject the external actor should send its stream to.
Hand this to your domain actor (e.g. as its subscriber) so each emitted value is forwarded to the bridged socket/topic.
pub fn subject(Bridge(a)) -> process.Subject(a)