defmodule ViaUtils.Comms do use GenServer require Logger @spec start_unsupervised_operator(atom(), integer()) :: tuple() def start_unsupervised_operator(name, refresh_groups_interval_ms \\ 1000) do start_link(name: name, refresh_groups_loop_interval_ms: refresh_groups_interval_ms) end def start_link(config) do name = Keyword.fetch!(config, :name) Logger.debug("Start ViaUtils.Comms: #{inspect(name)}") ViaUtils.Process.start_link_singular(GenServer, __MODULE__, config, via_tuple(name)) end @impl GenServer def init(config) do state = %{ groups: %{}, # purely for dianostics name: Keyword.fetch!(config, :name) } ViaUtils.Process.start_loop( self(), Keyword.fetch!(config, :refresh_groups_loop_interval_ms), :refresh_groups ) {:ok, state} end @impl GenServer def handle_cast({:join_group, group, process_id}, state) do # We will be added to our own record of the group during the # :refresh_groups cycle Logger.warn("#{inspect(state.name)} is joining group: #{inspect(group)}") # :pg.create(group) if !is_in_group?(group, process_id) do :pg.join(group, process_id) end {:noreply, state} end @impl GenServer def handle_cast({:leave_group, group, process_id}, state) do # We will be remove from our own record of the group during the # :refresh_groups cycle if is_in_group?(group, process_id) do :pg.leave(group, process_id) end {:noreply, state} end @impl GenServer def handle_cast({:send_msg_to_group, message, group, sender, global_or_local}, state) do # Logger.debug("send_msg. group: #{inspect(group)}") group_members = get_group_members(state.groups, group, global_or_local) # Logger.debug("op pid: #{inspect(self())}") # Logger.debug("Group members: #{inspect(group_members)}") send_msg_to_group_members(message, group_members, sender) {:noreply, state} end @impl GenServer def handle_info(:refresh_groups, state) do groups = Enum.reduce(:pg.which_groups(), %{}, fn group, acc -> all_group_members = :pg.get_members(group) local_group_members = :pg.get_local_members(group) Map.put(acc, group, %{global: all_group_members, local: local_group_members}) end) # Logger.debug("#{inspect(state.name)} groups after refresh: #{inspect(groups)}") {:noreply, %{state | groups: groups}} end def join_group(operator_name, group, process_id) do GenServer.cast(via_tuple(operator_name), {:join_group, group, process_id}) end def join_group(operator_name, group) do GenServer.cast(via_tuple(operator_name), {:join_group, group, self()}) end def leave_group(operator_name, group, process_id) do GenServer.cast(via_tuple(operator_name), {:leave_group, group, process_id}) end def leave_group(operator_name, group) do GenServer.cast(via_tuple(operator_name), {:leave_group, group, self()}) end @spec send_local_msg_to_group(atom(), any(), any(), any()) :: atom() def send_local_msg_to_group(operator_name, message, group, sender) do GenServer.cast(via_tuple(operator_name), {:send_msg_to_group, message, group, sender, :local}) end @spec send_local_msg_to_group(atom(), tuple(), any()) :: atom() def send_local_msg_to_group(operator_name, message, sender) do # Logger.debug("send to group: #{elem(message, 0)}: #{inspect(message)}") GenServer.cast( via_tuple(operator_name), {:send_msg_to_group, message, elem(message, 0), sender, :local} ) end @spec send_global_msg_to_group(atom(), any(), any(), any()) :: atom() def send_global_msg_to_group(operator_name, message, group, sender) do # Logger.debug("send global: #{inspect(message)}") GenServer.cast( via_tuple(operator_name), {:send_msg_to_group, message, group, sender, :global} ) end @spec send_global_msg_to_group(atom(), tuple(), any()) :: atom() def send_global_msg_to_group(operator_name, message, sender) do # Logger.debug("send global: #{inspect(message)}") GenServer.cast( via_tuple(operator_name), {:send_msg_to_group, message, elem(message, 0), sender, :global} ) end defp send_msg_to_group_members(message, group_members, sender) do Enum.each(group_members, fn dest -> if dest != sender do # Logger.debug("Send #{inspect(message)} to #{inspect(dest)}") GenServer.cast(dest, message) end end) end def is_in_group?(group, pid) do :pg.get_members(group) |> Enum.member?(pid) end def get_group_members(groups, group, global_or_local) do Map.get(groups, group, %{}) |> Map.get(global_or_local, []) end def via_tuple(name) do ViaUtils.Registry.via_tuple(__MODULE__, name) end @spec start_operator(atom) :: {:error, any} | {:ok, pid} | {:ok, pid, any} defdelegate start_operator(name), to: ViaUtils.Comms.Supervisor end