beryl/presence
Distributed presence tracking with a CRDT
This module wraps the pure lattice_presence/presence_state CRDT in an
OTP actor that:
- Handles track/update/untrack calls
- Publishes an actor-owned ETS read model
- Requests snapshots at startup and periodically through PubSub
- Receives remote snapshots and merges them internally
- Hides unavailable remote replicas without forgetting their causal state
- Invokes
on_diffwhen local changes or merges produce non-empty diffs
Each actor owns its local entries. Remote snapshots are accepted only from
a live owner answering an outstanding request, never through a relay.
Unavailable state is retained for 60 seconds before safe compaction; see
with_pubsub for the recovery and incarnation rules.
Presence is independent of the beryl runtime and runs under your application's supervision tree.
Read consistency
Section titled “Read consistency”The presence actor is the only writer to a protected ETS read model. It atomically replaces each topic's snapshot after completing a mutation, merge, or prune. Synchronous mutations publish before replying, and runtime mutations acknowledge only after publishing, so a later read observes the completed mutation without waiting on the actor mailbox. A read-model deletion failure stops publication and enters the mutation or sync processing error path; it is never treated as a successful write.
A read concurrent with a queued or in-progress mutation can observe the
previous or new complete snapshot. Reads of separate topics do not form
one atomic view. If the actor stops, it also destroys the table, and reads
return Error(Nil) instead of stale data.
Example
Section titled “Example”let pubsub_handle = pubsub.start(pubsub.default_config())let config = presence.default_config("node1") |> presence.with_pubsub(pubsub_handle) |> presence.with_broadcast_interval(1500)let #(presence_handle, presence_specification) = presence.child_spec(config)let assert Ok(_root) = static_supervisor.new(static_supervisor.OneForOne) |> static_supervisor.add(presence_specification) |> static_supervisor.start()let ref = presence.track( presence_handle, "room:lobby", "user:1", "socket-1", meta, )let assert Ok(entries) = presence.list(presence_handle, "room:lobby")Config
Section titled “Config”pub type ConfigConfiguration for starting presence.
Build configurations with default_config and the with_* functions.
beryl can then add options without exposing record fields as public API.
pub type DiffAn opaque diff representing presence joins and leaves grouped by topic.
beryl passes this value to Config.on_diff.
beryl.broadcast_presence_diff preserves its delivery scope automatically.
DiffScope
Section titled “DiffScope”pub type DiffScope { Cluster LocalNode}The audience for a presence diff.
Constructors
Section titled “Constructors”Cluster
Section titled “Cluster”ClusterApplication mutations and explicitly constructed diffs may be broadcast to the cluster.
LocalNode
Section titled “LocalNode”LocalNodeReplication, failure detection, and recovery describe this node's view. Deliver these diffs only to socket subscribers on the observing node.
pub type Event { Snapshot(entries: List(PresenceEntry)) Changed( joins: List(PresenceEntry), leaves: List(PresenceEntry) )}One topic's presence stream, delivered by channel.on_presence.
A subscription starts with one snapshot, including an empty list for an empty topic. Later changes include this connection's own changes and changes merged from other nodes. Apply leaves before joins; a metadata update contains both the old entry's leave and the new entry's join.
Message
Section titled “Message”pub type MessageMessages that the presence actor handles.
Presence
Section titled “Presence”pub type PresenceA running Presence instance.
This handle is opaque. Callers cannot forge actor subjects or depend on the runtime representation. It contains the actor's stable registered subject and the name of its actor-owned ETS read model.
Node affinity
Section titled “Node affinity”The stable registered subject and ETS read model are resolved on the
caller's node. Keep a Presence handle on the node where its child
specification runs. From another BEAM node, synchronous mutations
(track, update, untrack, and untrack_all) cannot reach the owning
actor and panic as unavailable, while list, get_by_key, and count
return Error(Nil) because the read model is unavailable. Use PubSub
replication (with_pubsub) to share presence state across nodes instead of
moving the handle itself.
PresenceEntry
Section titled “PresenceEntry”pub type PresenceEntry { PresenceEntry( session_id: String, key: String, meta: json.Json )}A presence entry returned from queries and diff accessors.
This type is transparent. Callers can inspect query results and construct
entries for diff.
PresenceUpdateError
Section titled “PresenceUpdateError”pub type PresenceUpdateError { UnknownRef(ref: String) RequestFailed(overload.CallError)}Errors from an update to a tracked presence.
Constructors
Section titled “Constructors”UnknownRef
Section titled “UnknownRef”UnknownRef(ref: String)The ref is unknown, already removed, or was not returned by track.
Functions
Section titled “Functions”child_spec
Section titled “child_spec”pub fn child_spec(Config) -> #(Presence, supervision.ChildSpecification(process.Subject(Message)))Build the supervised presence actor.
Add the returned child specification to your application's supervisor. The returned handle is name-backed and works again after a supervised restart. A restart resets the in-memory presence entries and tracking refs.
pub fn count( Presence, String) -> Result(Int, Nil)Count presences in a topic.
This is equivalent to list(presence, topic) |> list.length, but O(1). It
reads the materialized count directly from the read model via
ets:lookup_element/4 instead of building (and copying) the entry list
just to measure it.
This function reads the actor-owned read model directly. The read model is an ETS snapshot created after each mutation, merge, or prune. This function does not wait on the actor mailbox. See the module's Read consistency section for ordering guarantees.
Returns Error(Nil) when the presence read model is unavailable. This can
occur when the presence actor is not running or when this handle is used
from a process on another BEAM node than the one it was started on (see
the node affinity note on Presence).
default_config
Section titled “default_config”pub fn default_config(String) -> ConfigDefault configuration (no PubSub).
The repair interval defaults to 1500 ms. Adding with_pubsub enables an
initial snapshot request and periodic repair, even without local changes.
Without PubSub, the interval is unused. A non-positive interval disables
periodic requests, but the initial exchange and replies remain enabled.
pub fn diff( joins: List(#(String, List(PresenceEntry))), leaves: List(#(String, List(PresenceEntry)))) -> DiffBuild a presence diff from topic-grouped joins and leaves.
Most applications receive diffs from Config.on_diff. Use this function
to construct an application diff with Cluster scope for
beryl.broadcast_presence_diff. Do not rebuild a replica-view diff with
this function: that would discard its LocalNode scope.
diff_joins
Section titled “diff_joins”pub fn diff_joins( Diff, String) -> List(PresenceEntry)Return presence joins for a topic in this diff.
diff_leaves
Section titled “diff_leaves”pub fn diff_leaves( Diff, String) -> List(PresenceEntry)Return presence leaves for a topic in this diff.
diff_scope
Section titled “diff_scope”pub fn diff_scope(Diff) -> DiffScopeReturn where this diff may be delivered.
Local application mutations produce Cluster diffs. Remote snapshots and
replica availability changes produce LocalNode diffs, even when one
update contains both causal changes and liveness changes. They repair the
observing node's view and must not be rebroadcast to other nodes.
Prefer beryl.broadcast_presence_diff, which handles this distinction.
Custom publishers must preserve it; the Phoenix JSON payload has no scope.
diff_topics
Section titled “diff_topics”pub fn diff_topics(Diff) -> List(String)List topics touched by this diff.
get_by_key
Section titled “get_by_key”pub fn get_by_key( Presence, String, String) -> Result(List(#(String, json.Json)), Nil)Get presences for a specific key within a topic.
This function reads the actor-owned read model directly. The read model is an ETS snapshot created after each mutation, merge, or prune. This function does not wait on the actor mailbox. See the module's Read consistency section for ordering guarantees.
Returns Error(Nil) when the presence read model is unavailable. This can
occur when the presence actor is not running or when this handle is used
from a process on another BEAM node than the one it was started on (see
the node affinity note on Presence).
pub fn list( Presence, String) -> Result(List(PresenceEntry), Nil)List all presences for a topic.
This function reads the actor-owned read model directly. The read model is an ETS snapshot created after each mutation, merge, or prune. This function does not wait on the actor mailbox. See the module's Read consistency section for ordering guarantees.
Returns Error(Nil) when the presence read model is unavailable. This can
occur when the presence actor is not running or when this handle is used
from a process on another BEAM node than the one it was started on (see
the node affinity note on Presence).
queue_snapshot
Section titled “queue_snapshot”pub fn queue_snapshot(Presence) -> Result(overload.Occupancy, overload.AdmissionError)Read mutation and reserved-cleanup accounting without waiting for the actor.
pub fn track( Presence, String, String, String, json.Json) -> Result(String, overload.CallError)Track a presence in a topic.
session_id identifies the session, such as a socket, that owns this
presence. untrack_all matches this value when the session disconnects.
Returns a server-generated tracking ref. It is an opaque, unique handle for
this presence. Pass it to untrack to remove this entry. The ref is not
the session ID. The presence actor creates it, and it is meaningful only
to that actor. The ref is also merged into object metas as phx_ref for
Phoenix client compatibility.
Returns a typed call error on admission failure, owner exit, or timeout (5 seconds by default). A timeout cancels pending work, but a running mutation may still complete.
untrack
Section titled “untrack”pub fn untrack( Presence, String) -> Result(Nil, overload.CallError)Untrack a specific presence using the ref returned by track.
Removing an unknown or already-removed ref is a harmless no-op.
Returns a typed call error on admission failure, owner exit, or timeout. A timeout cancels pending work, but a running mutation may still complete.
untrack_all
Section titled “untrack_all”pub fn untrack_all( Presence, String) -> Result(Nil, overload.CallError)Untrack all presences locally tracked for a session, such as when a socket disconnects.
Replicated entries owned by another presence actor are not removed, even when they use the same session ID.
Returns a typed call error on admission failure, owner exit, or timeout. A timeout cancels pending work, but a running mutation may still complete.
update
Section titled “update”pub fn update( Presence, String, json.Json) -> Result(String, PresenceUpdateError)Replace the meta of a presence created by track.
One diff contains the old ref's leave and the new ref's join. Subscribers do not observe an intermediate state without the presence key. Other tracked refs for the same key do not change.
Returns the replacement ref, which must be used for subsequent update
or untrack calls. Returns Error(UnknownRef(ref)) when ref is
unknown, already removed, or belongs to the internal runtime.
RequestFailed wraps admission, owner-exit, and timeout errors. A timeout
cancels pending work, but a running mutation may still complete.
with_broadcast_interval
Section titled “with_broadcast_interval”pub fn with_broadcast_interval( Config, Int) -> ConfigSet how often presence requests full snapshots from its PubSub peers.
Requests run at startup and every interval_ms thereafter, including when
no application state changes. Lost requests or replies are retried on later
ticks. The default is 1500 ms; this is a repair cadence, not a convergence
deadline. Delivery, membership propagation, and actor work can delay repair.
A non-positive value disables periodic requests, not the initial request or replies to peers. Without periodic requests, quiet recovery is not guaranteed.
with_call_timeout
Section titled “with_call_timeout”pub fn with_call_timeout( Config, Int) -> ConfigSet the timeout for synchronous presence mutations, in milliseconds.
This timeout applies to track, update, untrack, and untrack_all. These
functions panic if the actor does not reply before the timeout. The default
is 5000 ms.
with_on_diff
Section titled “with_on_diff”pub fn with_on_diff( Config, fn(Diff) -> Nil) -> ConfigSet the callback for diffs from local changes, remote merges, or replica availability changes.
Pass the original diff to beryl.broadcast_presence_diff to preserve its
delivery scope. Application mutations publish cluster-wide at their source.
Replication and availability callbacks repair local clients only. A custom
publisher must inspect diff_scope rather than broadcast encoded JSON
unconditionally. A local worker may handle the callback; do not move a
LocalNode diff to another node for publication.
The callback runs synchronously on the presence actor, for both local
mutations (track/update/untrack/untrack_all, and the asynchronous
mutations the runtime issues for presence effects) and remote merges,
before the affected topics' read-model snapshots are (re)published and
before the triggering call replies or the mutation is acknowledged.
This ordering is the same for local and remote diffs.
If the callback reads presence state through the same
Presence handle (list, get_by_key, count) for a topic this diff
changes, it observes the previous snapshot. It does not observe the
snapshot that the diff will produce. Read the entries and counts you need
directly from the Diff argument (via
diff_joins/diff_leaves) instead of re-reading through presence
inside the callback.
Keep the callback fast and non-blocking. It runs on the actor process.
A slow or blocking callback delays that topic's read-model publish,
the reply to (or acknowledgement of) the mutating operation, and every
other message behind it in the actor's mailbox. Concurrent
list/get_by_key/count calls from other processes do not use the
mailbox and are not delayed. A socket with an active presence effect waits
for the callback. Callers of synchronous mutations also wait for their
replies. Enqueue a small message to a bounded application-owned worker and
return. Do not make network calls or synchronously mutate the same presence
actor from this callback.
beryl catches and logs callback exceptions, exits, and throws. A callback failure does not veto an otherwise successful local mutation or remote merge: beryl still publishes the snapshot and replies or acknowledges. beryl does not retry the callback, and it cannot roll back callback effects that completed before the failure. Treat delivery as a notification, not exactly-once application processing.
Presence queue snapshots and occupancy telemetry retain an admitted local mutation while its callback runs. They do not impose a callback deadline, apply to remote sync, or bound an application worker's mailbox.
with_pubsub
Section titled “with_pubsub”pub fn with_pubsub( Config, pubsub.PubSub(SyncPayload)) -> ConfigEnable PubSub replication for presence.
Use a fixed scope reserved for presence sync payloads on every node.
Do not reuse the scope used by beryl.with_pubsub, which carries JSON
broadcasts. Different topics do not isolate incompatible payload types.
The PubSub same-scope caller obligation also applies to presence; its
envelope version is not runtime validation of arbitrary BEAM terms.
Remote visibility follows monitored actor ownership and the local pg
membership view. Actor exit, node disconnection, or membership loss hides
that replica and emits leaves. Its causal state remains available for repair.
A fresh snapshot from the same actor restores its current entries; a
replacement actor starts a new incarnation. A partition can therefore hide
sessions that remain connected to their local node.
Each snapshot contains only its sender's authoritative state. A receiver's request order, not a random suffix or message arrival order, determines whether a new incarnation can replace its known owner. Concurrent live actors sharing a replica base in one scope are unsupported.
A confirmed replacement retires its predecessor. Otherwise, unavailable state remains for 60 seconds, checked every second while the actor runs. Compaction invalidates outstanding requests. A returning actor must answer a new request with its current full local snapshot, so delayed replies and lagging peers cannot reintroduce compacted history. The retention check also runs when periodic snapshot requests are disabled. Actor work can delay it.
with_queue_limits
Section titled “with_queue_limits”pub fn with_queue_limits( Config, overload.Limits) -> ConfigBound queued and executing local mutations. Replication is not admitted here.
with_telemetry
Section titled “with_telemetry”pub fn with_telemetry(Config) -> ConfigEmit local mutation queue occupancy and overload events.
