Packages

A distributed cache library for Elixir with pluggable backends, topologies, and near-cache support.

Current section

Files

Jump to
vela_cache lib vela distributed ring_manager.ex
Raw

lib/vela/distributed/ring_manager.ex

defmodule Vela.Distributed.RingManager do
@moduledoc """
Watches for cluster membership changes and updates the partitioned
topology's hash ring in persistent_term.
Subscribes to `Vela.Distributed.Cluster` events (:node_up / :node_down).
When the cluster changes, rebuilds the ring and atomically swaps it
into the topology state stored in persistent_term.
"""
use GenServer
alias Vela.Distributed.Ring
def start_link(opts) do
cache_name = Keyword.fetch!(opts, :name)
GenServer.start_link(__MODULE__, cache_name, name: server_name(cache_name))
end
def server_name(cache_name), do: :"vela_ring_manager_#{cache_name}"
@impl true
def init(cache_name) do
# Subscribe to cluster events via pg
Vela.Distributed.Cluster.subscribe()
# Sync the ring immediately with current cluster membership,
# in case nodes connected before this process started
sync_ring(cache_name)
{:ok, %{cache_name: cache_name}}
end
@impl true
def handle_info({:vela_cluster_event, event, changed_node}, state)
when event in [:node_up, :node_down] do
update_ring(state.cache_name, event, changed_node)
{:noreply, state}
end
def handle_info(_msg, state), do: {:noreply, state}
defp update_ring(cache_name, event, changed_node) do
ts = :persistent_term.get({:vela_topology_state, cache_name})
new_ring =
case event do
:node_up -> Ring.add_node(ts.ring, changed_node)
:node_down -> Ring.remove_node(ts.ring, changed_node)
end
:persistent_term.put(
{:vela_topology_state, cache_name},
%{ts | ring: new_ring}
)
end
defp sync_ring(cache_name) do
ts = :persistent_term.get({:vela_topology_state, cache_name})
current_nodes = [node() | Node.list()]
new_ring = Ring.new(current_nodes)
:persistent_term.put(
{:vela_topology_state, cache_name},
%{ts | ring: new_ring}
)
end
end