Packages
finitomata
0.30.1
0.41.0
0.40.0
0.35.0
0.34.0
0.33.0
0.32.0
0.31.1
0.30.3
0.30.2
0.30.1
0.30.0
0.29.10
0.29.9
0.29.8
0.29.7
0.29.6
0.29.5
0.29.4
0.29.3
0.29.2
0.29.1
0.29.0
0.28.1
0.28.0
0.27.1
0.27.0
0.26.4
0.26.3
0.26.2
0.26.1
0.26.0
0.25.0
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.7
0.23.6
0.23.5
0.23.4
0.23.3
0.23.2
0.23.1
0.23.0
0.22.1
0.22.0
0.21.4
0.21.3
0.21.2
0.21.1
0.21.0
0.20.2
0.20.1
0.20.0
0.19.6
0.19.5
0.19.4
0.19.3
0.19.2
0.19.1
0.19.0
0.18.4
0.18.3
0.18.2
0.18.1
0.18.0
0.17.1
0.17.0
0.16.0
0.15.1
0.15.0
0.14.6
0.14.5
0.14.4
0.14.3
0.14.2
0.14.1
0.14.0
0.13.0
0.12.1
0.12.0
0.11.3
0.11.2
0.11.1
0.11.0
0.10.0
0.9.1
0.9.0
0.8.2
0.8.1
0.8.0
0.7.2
0.7.1
0.7.0
0.6.3
0.6.2
0.6.1
0.6.0
0.5.2
0.5.1
0.5.0
0.4.0
0.3.0
0.2.0
0.1.1
0.1.0
The FSM implementation generated from PlantUML textual representation.
Current section
Files
Jump to
Current section
Files
lib/finitomata/distributed/supervisor.ex
defmodule Finitomata.Distributed.Supervisor do
@moduledoc false
require Logger
use Supervisor
def start_link(id, nodes \\ Node.list()) do
__MODULE__
|> Supervisor.start_link(id, name: sup_name(id))
|> then(fn
{:ok, pid} when is_pid(pid) ->
Enum.each(nodes, fn node ->
with {:badrpc, error} <- :rpc.block_call(node, __MODULE__, :start_link, [id, []]) do
Logger.error(
"[♻️] Remote start: " <>
inspect(id: id, name: sup_name(id), node: node, error: error)
)
end
end)
{:ok, _task_pid} = Task.start(fn -> synch(id, nodes) end)
{:ok, pid}
{:error, {:already_started, pid}} ->
{:ok, pid}
other ->
other
end)
end
@impl true
def init(id) do
id = Finitomata.Supervisor.infinitomata_name(id)
children = [
%{id: {:pg, id}, start: {__MODULE__, :start_pg, []}},
{Finitomata, id},
%{id: {Agent, id}, start: {Agent, :start_link, [fn -> %{} end, [name: agent(id)]]}},
{Finitomata.Distributed.GroupMonitor, id}
]
Supervisor.init(children, strategy: :one_for_one)
end
defp sup_name(id), do: id |> Finitomata.Supervisor.infinitomata_name() |> Module.concat(Sup)
defp agent(nil), do: agent(__MODULE__)
defp agent(id), do: Module.concat(id, IdLookup)
def group(nil), do: group(__MODULE__)
def group(id), do: Module.concat(id, Group)
def ungroup(id), do: id |> Module.split() |> Enum.slice(0..-2//1) |> Module.concat()
def synch(id, nodes \\ Node.list(), fqn_id \\ nil) do
fqn_id = if is_nil(fqn_id), do: Finitomata.Supervisor.infinitomata_name(id), else: fqn_id
known_processes_alive = fqn_id |> group() |> :pg.get_members()
domestic = Map.new(Infinitomata.all(id), &fix_pid(&1, known_processes_alive))
known_fsms_alive = fn _ ->
merger =
fn
_k, %{node: node, pid: pid}, %{node: node, pid: pid} ->
%{node: node, pid: pid, ref: make_ref()}
_k, %{node: _node, pid: nil}, %{node: node, pid: pid} ->
%{node: node, pid: pid, ref: make_ref()}
_k, %{node: node, pid: pid}, %{node: _node, pid: _pid} ->
%{node: node, pid: pid, ref: make_ref()}
end
call_handler = fn result, id, node ->
with {:badrpc, error} <- result do
Logger.warning("[♻️] Synch Error: " <> inspect(id: id, node: node, error: error))
[]
end
end
Enum.reduce(nodes, domestic, fn node, acc ->
node
|> :rpc.block_call(Infinitomata, :all, [id])
|> call_handler.(id, node)
|> Map.new(&fix_pid(&1, known_processes_alive))
|> Map.merge(acc, merger)
end)
end
fqn_id |> agent() |> Agent.update(known_fsms_alive)
end
def all(id) do
if is_nil(GenServer.whereis(agent(id))) do
Logger.debug("[♻️] Not healthy ‹#{inspect(id)}›, probably it’s terminating")
%{}
else
with empty when empty == %{} <- Agent.get(agent(id), & &1) do
Logger.debug("[♻️] Empty pool for ‹#{inspect(id)}›")
%{}
end
end
end
def del(id, name), do: Agent.update(agent(id), &Map.delete(&1, name))
def get(id, name), do: Agent.get(agent(id), &Map.get(&1, name))
def select_by_pids(id, pids) when is_list(pids) do
Agent.get(
agent(id),
&Map.filter(&1, fn {_, %{pid: pid}} -> pid in pids end)
)
end
def delete_by_pids(id, pids) when is_list(pids) do
{match, mismatch} =
Agent.get(agent(id), &split_with(&1, fn {_, %{pid: pid}} -> pid in pids end))
Agent.update(agent(id), fn _ -> mismatch end)
match
end
def delete_by_pids(id, %MapSet{} = pids) do
{match, mismatch} =
Agent.get(agent(id), &split_with(&1, fn {_, %{pid: pid}} -> MapSet.member?(pids, pid) end))
Agent.update(agent(id), fn _ -> mismatch end)
match
end
def put(id, name, %{pid: pid} = data) when is_pid(pid) do
Agent.update(agent(id), &Map.put(&1, name, data))
end
defp skip_node(nil), do: nil
defp skip_node(pid) when is_pid(pid), do: pid |> :erlang.pid_to_list() |> skip_node()
defp skip_node([?. | tail]), do: tail
defp skip_node([_ | tail]), do: skip_node(tail)
defp fix_pid({name, %{pid: pid} = value}, pids),
do: {name, %{value | pid: fix_pid(pid, pids)}}
defp fix_pid(pid, pids) do
pid_tail = skip_node(pid)
Enum.find(pids, &(skip_node(&1) == pid_tail))
end
if Version.compare(System.version(), "1.15.0") == :lt do
defp split_with(%{} = map, fun) when is_function(fun, 1) do
{truthy, falsey} = Enum.split_with(map, fun)
{Map.new(truthy), Map.new(falsey)}
end
else
defdelegate split_with(map, fun), to: Map
end
@doc false
def start_pg do
with {:error, {:already_started, _pid}} <- :pg.start_link(), do: :ignore
end
end