Skip to content

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_diff when 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.

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.

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")
pub type Config

Configuration 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 Diff

An 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.

pub type DiffScope {
Cluster
LocalNode
}

The audience for a presence diff.

Cluster

Application mutations and explicitly constructed diffs may be broadcast to the cluster.

LocalNode

Replication, 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.

pub type Message

Messages that the presence actor handles.

pub type Presence

A 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.

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.

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.

pub type PresenceUpdateError {
UnknownRef(ref: String)
RequestFailed(overload.CallError)
}

Errors from an update to a tracked presence.

UnknownRef(ref: String)

The ref is unknown, already removed, or was not returned by track.

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).

pub fn default_config(String) -> Config

Default 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)))
) -> Diff

Build 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.

pub fn diff_joins(
Diff,
String
) -> List(PresenceEntry)

Return presence joins for a topic in this diff.

pub fn diff_leaves(
Diff,
String
) -> List(PresenceEntry)

Return presence leaves for a topic in this diff.

pub fn diff_scope(Diff) -> DiffScope

Return 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.

pub fn diff_topics(Diff) -> List(String)

List topics touched by this diff.

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).

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.

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.

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.

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.

pub fn with_broadcast_interval(
Config,
Int
) -> Config

Set 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.

pub fn with_call_timeout(
Config,
Int
) -> Config

Set 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.

pub fn with_on_diff(
Config,
fn(Diff) -> Nil
) -> Config

Set 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.

pub fn with_pubsub(
Config,
pubsub.PubSub(SyncPayload)
) -> Config

Enable 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.

pub fn with_queue_limits(
Config,
overload.Limits
) -> Config

Bound queued and executing local mutations. Replication is not admitted here.

pub fn with_telemetry(Config) -> Config

Emit local mutation queue occupancy and overload events.