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