Current section
Files
Jump to
Current section
Files
lib/comms.ex
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