Packages
SuperWorker is a powerful Elixir library for working with supervisors and background jobs. It provides a much simpler approach than traditional supervisors. This library is currently under development and is unstable, so it is not recommended for production use.
Current section
Files
Jump to
Current section
Files
lib/supervisor/supervisor.ex
defmodule SuperWorker.Supervisor do
@moduledoc """
Documentation for `SuperWorker.Supervisor`.
This module is a new model for the Supervisor module.
Supervisor supports the following features:
- Group processes
- Chain processes
- Freedom processes
## Group processes
Group processes are a set of processes that are started together. If one of the processes dies, all the processes in the group will be stopped.
## Chain processes
Chain processes are a set of processes that are started one after another. The output of the previous process is passed to the next process.
## Freedom processes
Freedom processes are independent processes that are started separately.
All type of processes can be started in parallel & can be stopped individually or in a group.
## Examples
# Start a supervisor with 2 partitions & 2 groups:
alias SuperWorker.Supervisor, as: Sup
opts = [id: :sup1, number_of_partitions: 2, link: false]
Sup.start(opts)
Sup.add_group(:sup1, [id: :group1, restart_strategy: :one_for_all])
Sup.add_group_worker(:sup1, :group1, {Dev, :task, [15]}, [id: :g1_1])
Sup.add_group(:sup1, [id: :group2, restart_strategy: :one_for_all])
Sup.add_group_worker(:sup1, :group2, fn ->
receice do
msg ->
:ok
end
end, [id: :g2_2])
"""
require Logger
alias SuperWorker.Supervisor.{Group, Chain, Worker, Message}
alias :ets, as: Ets
import SuperWorker.Supervisor.Utils
defstruct [
:id, # supervisor id
:owner, # owner of the supervisor
:master,
:number_of_partitions, # number of partitions, default is number of online schedulers
link: true, # link the supervisor to the caller
report_to: [], # list of pid or callback function, for reporting worker crashed or worker finished.
linked_pids: [], # list of linked external pids
children: [] # list of children (group, chain, worker) to start when supervisor starts.
]
@sup_params [:id, :number_of_partitions, :link, :report_to, :children]
@me __MODULE__
alias __MODULE__
# Default timeout (miliseconds) for API calls.
@default_time 5_000
# List message from api.
@api_messages [:start_worker, :get_group, :remove_group_worker, :restart_group_worker,
:get_chain, :send_to_group, :send_to_group_random, :add_data_to_chain, :send_to_worker,
:remove_group_worker, :add_group, :add_chain, :stop]
# List internal message.
@internal_messages []
## Public APIs
@doc """
Start supervisor for run standalone please set option :link to false.
result format: {:ok, pid} or {:error, reason}
"""
@spec start([id: atom(), link: boolean() | pid(), number_of_partitions: integer(),
report_to: list()]) :: {:ok, pid} | {:error, any()}
def start(opts, timeout \\ 5_000) when is_list(opts) do
with {:ok, opts} <- check_opts(opts),
{:ok, sup} <- map_to_struct(opts),
false <- is_running?(sup.id) do
start_supervisor(sup, timeout)
else
true ->
Logger.error("Supervisor is already running.")
{:error, :already_running}
{:error, _} = error ->
Logger.error("Error when starting supervisor: #{inspect error}")
error
end
end
@doc """
Stop supervisor.
Type of shutdown:
- :normal supervisor will send a message to worker for graceful shutdown. Not support for spawn process by function.
- :kill supervisor will kill worker.
"""
@spec stop(atom(), shutdown_type :: atom(), timeout :: integer()) :: {:ok, atom()} | {:error, any()}
def stop(sup_id, shutdown_type \\ :kill, timeout \\ @default_time) do
case get_pid(sup_id) do
{:error, _} = err ->
Logger.error("Supervisor is not running.")
err
{:ok, pid} ->
Logger.debug("Stopping supervisor: #{inspect pid}, shutdown type: #{inspect shutdown_type}")
call_api(pid, :stop, shutdown_type, timeout)
end
end
@doc """
Check if supervisor is running.
return true if supervisor is running, otherwise return false.
"""
@spec is_running?(atom()) :: boolean()
def is_running?(sup_id) do
case get_pid(sup_id) do
{:ok, _} -> true
{:error, _} -> false
end
end
@doc """
Add a standalone worker process to the supervisor.
function for start worker can be a function or a {module, function, arguments}.
Standalone worker is run independently from other workers follow :one_to_one strategy.
If worker crashes, it will check the restart strategy of worker then act accordingly.
"""
@spec add_standalone_worker(atom(), {module(), atom(), list()} | fun(), list(), integer()) :: {:ok, atom()} | {:error, any()}
def add_standalone_worker(sup_id, mfa_or_fun, opts \\ [], timeout \\ @default_time)
def add_standalone_worker(sup_id, {m, f, a} = mfa, opts, timeout)
when is_list(opts) and is_atom(m) and is_atom(f) and is_list(a) do
do_add_worker(sup_id, :standalone, [{:fun, mfa} | opts], timeout)
end
def add_standalone_worker(sup_id, fun, opts, timeout) when is_list(opts) and is_function(fun, 0) do
do_add_worker(sup_id, :standalone, [{:fun, {:fun, fun}} | opts], timeout)
end
@doc """
Add a worker to a group in the supervisor.
Function's options follow `Worker` module.
"""
@spec add_group_worker(atom(), atom(), {module(), atom(), list()} | fun(), list(), integer()) :: {:ok, atom()} | {:error, any()}
def add_group_worker(sup_id, group_id, mfa_or_fun, opts, timeout \\ @default_time)
def add_group_worker(sup_id, group_id, {m, f, a} = mfa, opts, timeout)
when is_list(opts) and is_atom(m) and is_atom(f) and is_list(a) do
do_add_worker(sup_id, {:group_id, group_id}, [{:fun, mfa} | opts], timeout)
end
def add_group_worker(sup_id, group_id, fun, opts, timeout) when is_list(opts) and is_function(fun, 0) do
do_add_worker(sup_id, {:group_id, group_id}, [{:fun, {:fun, fun}} | opts], timeout)
end
@doc """
Add a worker to the chain in supervisor.
"""
@spec add_chain_worker(atom(), atom(), {module(), atom(), list()} | fun(), list(), integer()) :: {:ok, atom()} | {:error, any()}
def add_chain_worker(sup_id, chain_id, mfa_or_fun, opts, timeout \\ @default_time)
def add_chain_worker(sup_id, chain_id, {m, f, a} = mfa, opts, timeout)
when is_list(opts) and is_atom(m) and is_atom(f) and is_list(a) do
do_add_worker(sup_id, {:chain_id, chain_id}, [ {:fun, mfa} | opts], timeout)
end
def add_chain_worker(sup_id, chain_id, fun, opts, timeout) when is_list(opts) and is_function(fun, 0) do
do_add_worker(sup_id, {:chain_id, chain_id}, [{:fun, {:fun, fun}} | opts], timeout)
end
@doc """
Add a group to the supervisor.
Group's options follow docs in `Group` module.
"""
@spec add_group(atom(), list(), integer()) :: {:ok, atom()} | {:error, any()}
def add_group(sup_id, opts, timeout \\ @default_time) do
with {:ok, group} <- Group.check_options(opts),
true <- is_running?(sup_id) do
case Ets.lookup(get_table_name(sup_id), {:group, group.id}) do
[] ->
with {:ok, pid} <- get_host_partition(sup_id, group.id) do
call_api(pid, :add_group, group, timeout)
end
_ ->
{:error, :already_exists}
end
else
false ->
Logger.error("Supervisor is not running.")
{:error, :not_running}
{:error, _} = error ->
Logger.error("Error when adding group: #{inspect error}")
error
end
end
@doc """
get group structure from supervisor.
"""
@spec get_group(atom(), atom(), integer()) :: {:ok, map()} | {:error, any()}
def get_group(sup_id, group_id, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, group_id) do
call_api(pid, :get_group, group_id, timeout)
end
end
@doc """
Add a chain to the supervisor.
Chain's options follow docs in `Chain` module.
"""
def add_chain(sup_id, opts, timeout \\ 5_000) do
with {:ok, chain} <- Chain.check_options(opts),
true <- is_running?(sup_id) do
case Registry.lookup(sup_id, {:chain, chain.id}) do
[] ->
with {:ok, pid} <- get_host_partition(sup_id, chain.id) do
call_api(pid, :add_chain, chain, timeout)
end
_ ->
{:error, :already_exists}
end
end
end
@doc """
Send data to the entry worker in the chain.
If chain doesn't has any worker, it will be dropped.
"""
def send_to_chain(sup_id, chain_id, data, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, chain_id) do
call_api(pid, :add_data_to_chain, {chain_id, data}, timeout)
end
end
@doc """
Send data directly to the worker (standalone, group, chain) in the supervisor.
"""
def send_to_worker(sup_id, worker_id, data, timeout \\ @default_time) do
with true <- is_running?(sup_id),
[{pid, _}] <- Registry.lookup(sup_id, {:worker, worker_id}) do
call_api(pid, :send_to_worker, {worker_id, data}, timeout)
end
end
@doc """
Send data to all workers in a group.
"""
def broadcast_to_group(sup_id, group_id, data, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, group_id) do
call_api(pid, :broadcast_to_group, {group_id, data}, timeout)
end
end
@doc """
Send data to all workers in current group of worker.
Using for communite between workers in the same group.
"""
def broadcast_to_my_group(data) do
group_id = get_my_group()
sup_id = get_my_supervisor()
cond do
group_id == nil ->
Logger.error("Group not found.")
{:error, :not_found}
sup_id == nil ->
Logger.error("Supervisor not found.")
{:error, :not_found}
true ->
broadcast_to_group(sup_id, group_id, data)
end
end
@doc """
Send data to a worker in the group.
"""
def send_to_group(sup_id, group_id, worker_id, data, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, group_id) do
call_api(pid, :send_to_group, {group_id, worker_id, data}, timeout)
end
end
@doc """
Send data to a random worker in the group.
"""
def send_to_group_random(sup_id, group_id, data, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, group_id) do
call_api(pid, :send_to_group_random, {group_id, data}, timeout)
end
end
@doc """
Send data to other worker in the same group.
"""
def send_to_my_group(worker_id, data) do
group_id = get_my_group()
sup_id = get_my_supervisor()
cond do
group_id == nil ->
Logger.error("Group not found.")
{:error, :not_found}
sup_id == nil ->
Logger.error("Supervisor not found.")
{:error, :not_found}
true ->
send_to_group(sup_id, group_id, worker_id, data)
end
end
def send_to_my_group_random(data) do
group_id = get_my_group()
sup_id = get_my_supervisor()
cond do
group_id == nil ->
Logger.error("Group not found.")
{:error, :not_found}
sup_id == nil ->
Logger.error("Supervisor not found.")
{:error, :not_found}
true ->
send_to_group_random(sup_id, group_id, data)
end
end
def get_my_group() do
Process.get({:supervisor, :group_id})
end
def get_my_supervisor() do
Process.get({:supervisor, :sup_id})
end
def get_chain(sup_id, chain_id, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, chain_id) do
call_api(pid, :get_chain, chain_id, timeout)
end
end
def remove_group_worker(sup_id, group_id, worker_id, timeout \\ @default_time) do
with {:ok, pid} <- verify_and_get_pid(sup_id, group_id) do
call_api(pid, :remove_group_worker, {worker_id, group_id}, timeout)
end
end
## Internal public functions
def init(opts = %Supervisor{}, ref) do
state = %{
groups: MapSet.new(), # storage for group processes
chains: MapSet.new(), # storage for chain processes
standalone: MapSet.new(), # storage for standalone processes
id: opts.id, # supervisor id
owner: opts.owner,
data_table: get_table_name(opts.id),
master: opts.id,
prefix: "[#{inspect opts.id}, master]"
}
# Register the supervisor process.
Process.register(self(), get_master_id(opts.id))
# Link to remote pid if link is a pids
case opts.link do
pid when is_pid(pid) ->
Process.link(pid)
list_pid when is_list(list_pid) ->
Enum.each(list_pid, fn pid ->
Process.link(pid)
end)
bool when bool in [true, false] ->
:ok
end
# Create Registry for map worker/group/chain id to pid.
{:ok, _} = Registry.start_link(keys: :duplicate, name: opts.id, partitions: opts.number_of_partitions)
Registry.register(state.master, :master, opts.number_of_partitions)
# Store group, chain, workers in Ets
create_table(state.data_table)
# Turn main process to system process.
Process.flag(:trap_exit, true)
Ets.insert(state.data_table, {:number_of_partitions, opts.number_of_partitions})
Ets.insert(state.data_table, {:owner, opts.owner})
Ets.insert(state.data_table, {:master, opts.id})
list_partitions = init_additional_partitions(opts)
state =
state
|> Map.put(:partitions , list_partitions)
|> Map.put(:role, :master)
Logger.debug("Supervisor #{inspect state.id} initialized: #{inspect state}")
# TO-DO: Add group, chain, worker from opts.
if opts.children != nil do
Enum.each(opts.children, fn child ->
case child do
{:group, group} ->
add_group(state.id, group)
{:chain, chain} ->
add_chain(state.id, chain)
{:standalone, worker} ->
if Keyword.get(worker, :options) == nil do
add_standalone_worker(state.id, worker.task)
else
add_standalone_worker(state.id, worker.task, worker.opts)
end
end
end)
end
api_response(ref, {:ok, self()})
# Start the main loop
main_loop(state, opts)
end
def child_spec(opts) do
%{
id: Keyword.get(opts, :id, @me), # default id is module name
start: {@me, :start, [opts]}
}
end
## Private functions
defp init_partition(opts = %Supervisor{}) do
# TO-DO: Implement partition supervisor.
state = %{
groups: %{}, # storage for group processes
chains: %{}, # storage for chain processes
standalone: %{}, # storage for standalone processes
ref_to_id: %{}, # storage for refs to id & type
id: opts.id, # supervisor id
owner: opts.owner,
role: :partition,
master: opts.master,
number_of_partitions: opts.number_of_partitions,
data_table: get_table_name(opts.master),
prefix: "[#{inspect opts.master}, #{inspect opts.id}]"
}
# Start the main loop
pid = spawn_link(@me, :main_loop, [state, opts])
Process.register(pid, state.id)
Logger.debug("#{inspect state.prefix} initialized, pid: #{inspect pid}")
Registry.register(state.master, {:partition , state.id}, [])
Ets.insert(state.data_table, {{:partition, state.id}, pid})
{:ok, opts.id, pid}
end
defp init_additional_partitions(opts) do
partitions = opts.number_of_partitions
Enum.map(0..partitions - 1, fn i ->
Logger.debug("[#{inspect opts.id}] add partition: #{inspect i}")
opts
|> Map.put(:master, opts.id)
|> Map.put(:id, String.to_atom("#{Atom.to_string(opts.id)}_#{i}"))
|> Map.put(:role, :partition)
|> init_partition()
end)
end
# Main loop for the supervisor & partition.
def main_loop(state, sup_opts) do
receive do
{msg_type, _, _} = msg when msg_type in @api_messages ->
Logger.debug("#{inspect state.prefix} received a api message: #{inspect(msg)}")
process_api_message(state, sup_opts, msg)
{:DOWN, _ref, :process, pid, reason} = msg ->
Logger.debug("#{inspect state.prefix} Worker died: #{inspect(pid)}, reason: #{inspect(reason)}")
process_worker_down(state, sup_opts, msg)
{:'EXIT', from, reason} ->
process_exit_message(state, sup_opts, from, reason)
{:stop_partition, type} ->
Logger.info("#{inspect state.prefix} Stopping supervisor partition, for #{inspect self()}")
# Stop the supervisor.
shutdown(state, type)
unknown ->
Logger.warning("#{inspect state.prefix} main_loop, unknown message: #{inspect(unknown)}")
main_loop(state, sup_opts)
end
Logger.debug("#{inspect state.prefix} #{inspect self()} main loop exited.")
end
defp shutdown(state, :kill) do
Logger.debug("Shutting down supervisor: #{inspect state.id}")
# TO-DO: Implement graceful shutdown for worker processes.
Enum.each(state.groups, fn {_, group} ->
Group.kill_all_workers(group)
end)
Enum.each(state.chains, fn {_, chain} ->
Chain.kill_all_workers(chain)
end)
Enum.each(state.standalone, fn {_, worker} ->
Process.exit(worker.pid, :kill)
end)
{:ok, :brutal_kill}
end
# process exit message for outside processes.
defp process_exit_message(state, sup_opts, from, reason) do
Logger.debug("#{inspect state.prefix} Exit message from: #{inspect from}, reason: #{inspect reason}")
if from in sup_opts.linked_pids do
Logger.warning("#{inspect state.prefix} exited follow external process (crashed): #{inspect from}")
raise "#{inspect state.master} crashed follow external process: #{inspect from}"
else
Logger.debug("#{inspect state.prefix} skipped exit for internal process: #{inspect from}")
main_loop(state, sup_opts)
end
end
# Add new worker to group/chain/standalone.
defp process_api_message(state, sup_opts, {:start_worker, ref, opts}) do
runable =
case opts.type do
:group ->
if has_group?(state, opts.group_id) or has_group_worker?(state, opts.group_id, opts.id) do
true
else
:group_not_found_or_worker_already_exists
end
:chain ->
if has_chain?(state, opts.chain_id) or has_chain_worker?(state, opts.chain_id, opts.id) do
true
else
:chain_not_found_or_worker_already_exists
end
:standalone ->
if has_worker?(state, opts.id) do
:worker_already_exists
else
true
end
end
state =
if runable == true do # start child process.
Logger.debug("#{inspect state.prefix} Everything is fine, starting child process with options: #{inspect(opts)}")
api_response(ref, {:ok, opts.id})
sup_start_child(state, opts)
else # not found group or chain, return error to the caller.
api_response(ref, {:error, runable})
state
end
main_loop(state, sup_opts)
end
# Get chain in supervisor and return to the caller.
defp process_api_message(state, sup_opts, {:get_chain, ref, chain_id}) do
result =
case get_group_or_chain(state, chain_id, :chain) do
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Not found chain with id #{inspect chain_id}")
error
{:ok, _} = res ->
res
end
api_response(ref, result)
main_loop(state, sup_opts)
end
# broadcast a data to all worker in group.
defp process_api_message(state, sup_opts, {:broadcast_to_group, ref, {group_id, data}}) do
result =
case get_group_or_chain(state, group_id, :group) do
{:ok, group} ->
Group.broadcast(group, data)
# TO-DO: Improve response
:ok
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Not found group with id #{inspect group_id}")
error
end
api_response(ref, result)
main_loop(state, sup_opts)
end
# send data directly to worker from api.
defp process_api_message(state, sup_opts, {:send_to_group, ref, {group_id, worker_id, data}}) do
with {:ok, group} <- get_group_or_chain(state, group_id, :group),
{:ok, worker} <- Group.get_worker(group, worker_id) do
send(worker.pid, data)
api_response(ref, :ok)
else
failed ->
Logger.error("#{inspect state.id}, send to worker #{inspect worker_id} in group #{inspect group_id}, error: #{inspect failed}")
api_response(ref, failed)
end
main_loop(state, sup_opts)
end
# send data to random worker from api.
defp process_api_message(state, sup_opts, {:send_to_group_random, ref, {group_id, data}}) do
case get_group_or_chain(state, group_id, :group) do
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Group not found: #{inspect group_id}, error: #{inspect error}")
api_response(ref, error)
{:ok, group} ->
worker_id = Enum.random(group.workers)
case Group.get_worker(group, worker_id) do
{:ok, worker} ->
send(worker.pid, data)
api_response(ref, :ok)
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Not found worker #{inspect worker_id} in group #{inspect group_id}")
api_response(ref, error)
end
end
main_loop(state, sup_opts)
end
# add data to chain from api.
defp process_api_message(state, sup_opts, {:add_data_to_chain, {from, _} = ref, {chain_id, data}}) do
case get_group_or_chain(state, chain_id, :chain) do
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Chain not found: #{inspect chain_id}")
api_response(ref, error)
{:ok, chain} ->
msg = Message.new(from, nil, data)
result = Chain.new_data(chain, msg)
Logger.debug("#{inspect state.prefix} Add data to chain: #{inspect chain_id}, result: #{inspect result}")
api_response(ref, result)
end
main_loop(state, sup_opts)
end
defp process_api_message(state, sup_opts, {:send_to_worker, ref, {worker_id, data}}) do
case get_group_or_chain(state, worker_id, :not_implement) do
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Worker not found: #{inspect worker_id}")
api_response(ref, error)
{:ok, worker} ->
send(worker.pid, data)
api_response(ref, :ok)
end
main_loop(state, sup_opts)
end
# remove worker from group.
defp process_api_message(state, sup_opts, {:remove_group_worker, ref, {worker_id, group_id}}) do
result =
case get_group_or_chain(state, group_id, :group) do
{:ok, group} ->
Group.remove_worker(group, worker_id)
{:error, _} = error ->
Logger.error("#{inspect state.prefix} Group not found: #{inspect group_id}")
error
end
api_response(ref, result)
main_loop(state, sup_opts)
end
defp process_api_message(state, sup_opts, {:restart_group_worker, worker_id , group_id}) do
Logger.debug("#{inspect state.prefix} Starting worker process, worker id: #{inspect(worker_id)}, group id: #{inspect(group_id)}")
[{_, group}] = Ets.lookup(state.data_table, {:group, group_id})
# Restart the worker process.
Group.restart_worker(group, worker_id)
main_loop(state, sup_opts)
end
# add group from api.
defp process_api_message(state, sup_opts, {:add_group, ref, group}) do
case Ets.lookup(state.data_table, {:gorup, group.id}) do
[_] ->
Logger.error("#{inspect state.prefix} Group already exists: #{inspect(group.id)}")
api_response(ref, {:error, :already_exists})
main_loop(state, sup_opts)
[] ->
Logger.debug("#{inspect state.prefix} Adding group: #{inspect(group.id)}")
# Send the response to the caller.
api_response(ref, {:ok, group.id})
state
|> add_new_group(group)
|> main_loop(sup_opts)
end
end
# get group info from api.
defp process_api_message(state, sup_opts, {:get_group, ref, group_id}) do
result =
case Ets.lookup(state.data_table, {:group, group_id}) do
[] ->
Logger.error("#{inspect state.prefix} Group not found: #{inspect(group_id)}")
{:error, :not_found}
[{_, group}] ->
{:ok, group}
end
api_response(ref, result)
main_loop(state, sup_opts)
end
# add chain from api.
defp process_api_message(state, sup_opts, {:add_chain, ref, chain}) do
case Ets.lookup(state.data_table, {:chain, chain.id}) do
[_] ->
Logger.error("#{inspect state.prefix} Chain already exists: #{inspect(chain.id)}")
api_response(ref, {:error, :already_exists})
main_loop(state, sup_opts)
[] ->
Logger.debug("#{inspect state.prefix} Adding chain: #{inspect(chain.id)}")
# Send the response to the caller.
api_response(ref, {:ok, chain.id})
state
|> add_new_chain(chain)
|> main_loop(sup_opts)
end
end
# Stop supervisor from api.
defp process_api_message(state, sup_opts, {:stop, ref, type}) do
Logger.info("#{inspect state.prefix} Stopping supervisor, request from #{inspect ref}")
# Send shutdown signal to all partitions.
Enum.each(0..sup_opts.number_of_partitions - 1, fn i ->
partition_id = get_partition_id(state.id, i)
case Ets.lookup(state.data_table, {:partition, partition_id}) do
[{_, pid}] ->
Logger.debug("#{inspect state.prefix} Sending shutdown signal to partition: #{inspect partition_id}")
send(pid, {:stop_partition, type})
_ ->
Logger.error("#{inspect state.prefix} Supervisor not found: #{inspect partition_id}")
end
end)
# stop worker on master.
shutdown(state, type)
# TO-DO: Clean KV store for supervisor.
api_response(ref, {:ok, :stopped})
exit(:normal)
end
defp process_api_message(state, sup_opts, unknown_msg) do
Logger.warning("#{inspect state.prefix} Unknown api message: #{inspect(unknown_msg)}")
main_loop(state, sup_opts)
end
defp process_worker_down(state, sup_opts, {:DOWN, _ref, :process, pid, :restart}) do
Logger.debug("#{inspect state.prefix} Ignore died process (process by other msg): #{inspect(pid)}")
main_loop(state, sup_opts)
end
defp process_worker_down(state, sup_opts, {:DOWN, ref, :process, pid, reason}) do
Logger.debug("Child process died: #{inspect(pid)}, ref: #{inspect ref}, reason: #{inspect(reason)}")
state =
case Ets.lookup(state.data_table, {:worker, :ref, ref}) do
[] ->
Logger.debug("Child is not found in table: #{inspect ref}, maybe already stopped.")
state
[{_, id, pid, type} = ref_data] ->
Logger.debug("Child found: #{inspect pid}, restarting. meta: #{inspect ref_data}")
Ets.delete(state.data_table, {:worker, :ref, ref})
case type do
:standalone ->
[{_, child}] = Ets.lookup(state.data_table, {:worker, id})
restart_standalone(state, child, {pid, reason})
{:group, group_id} ->
[{_, group}] = Ets.lookup(state.data_table, {:group, group_id})
restart_group(state, group, id, {pid, reason})
{:chain, chain_id} ->
[{_, chain}] = Ets.lookup(state.data_table, {:chain, chain_id})
restart_chain(state, chain, id, {pid, reason})
end
end
main_loop(state, sup_opts)
end
defp restart_standalone(state, %Worker{} = child, {pid, reason}) do
case child.restart_strategy do
:permanent ->
Logger.debug("#{inspect state.id}, :permanent, restarting #{inspect pid}")
# Restart the child process
sup_start_child(state, child)
:transient when reason != :normal ->
Logger.debug("#{inspect state.id}, :transient, reason down: #{inspect reason} restarting #{inspect pid}")
# Restart the child process
sup_start_child(state, child)
:transient ->
Logger.debug("#{inspect state.id}, :transient, ignore restarting #{inspect pid}")
state
:temporary ->
Logger.debug("#{inspect state.id}, :temporary, ignore restarting #{inspect pid}")
state
end
end
defp restart_group(state, %{restart_strategy: :one_for_one} = group, child_id, {pid, reason}) do
case reason do
:normal ->
Logger.debug("#{inspect state.id}, Child process(#{inspect(pid)}) is shutdown with reason :normal, ignore restarting.")
state
reason ->
Logger.debug("#{inspect state.id}, Child process(#{inspect(pid)}) is down with reason: #{inspect reason}, restarting.")
Group.restart_worker(group, child_id)
state
end
end
defp restart_group(state, %{restart_strategy: :one_for_all} = group, _child_id, {pid, reason}) do
case reason do
:normal ->
Logger.debug("#{inspect state.id}, Child process(#{inspect(pid)}) is :normal shutdown, ignore restarting.")
state
reason ->
Logger.debug("#{inspect state.id}, Child process(#{inspect(pid)}) is down with reason #{inspect reason}, restarting...")
old_ref_keys = Enum.reduce(Group.get_all_workers(group), [], fn worker, acc ->
[worker.ref | acc]
end)
Logger.debug("#{inspect state.id}, Old ref: #{inspect old_ref_keys}")
# Clean up old process.
# TO-DO: make sure pid, ref in worker struct is cleaned & correct after restart.
Group.kill_all_workers(group, :restart)
Enum.each(Group.get_all_workers(group), fn %Worker{id: worker_id} ->
{:ok, pid} = get_host_partition(state.master, worker_id)
send(pid, {:restart_group_worker, worker_id, group.id})
end)
state
end
end
# TO-DO: Follow restart strategy for child group & chain.
defp restart_chain(state, %{restart_strategy: :one_for_one} = chain, child_id, {pid, reason}) do
case reason do
:normal ->
Logger.debug("Worker(#{inspect child_id}) process(#{inspect(pid)}) is normal, ignore restarting.")
state
_ ->
Logger.debug("Worker(#{inspect child_id}) process(#{inspect(pid)}) is down, restarting.")
{:ok, chain} = Chain.restart_worker(chain, child_id)
{:ok, worker} = Chain.get_worker(chain, child_id)
chains = Map.put(state.chains, chain.id, chain)
ref_to_id = Map.put(state.ref_to_id, worker.ref, {child_id, {:chain, chain.id}})
state
|> Map.put(:chains, chains)
|> Map.put(:ref_to_id, ref_to_id)
end
end
defp restart_chain(state, %{restart_strategy: :one_for_all} = chain, _child_id, {pid, reason}) do
case reason do
:normal ->
Logger.debug("Child process(#{inspect(pid)}) is normal, ignore restarting.")
state
_ ->
Logger.debug("Child process(#{inspect(pid)}) is down, restarting...")
old_ref_keys = Enum.reduce(chain.workers, [], fn {_, worker}, acc -> [worker.ref | acc] end)
Logger.debug("Chai, Old ref: #{inspect old_ref_keys}")
{:ok, chain} = Chain.restart_all_workers(chain)
new_refs =
state.ref_to_id
|> Map.drop(old_ref_keys)
new_refs = Enum.reduce(chain.workers, new_refs, fn {child_id, worker}, acc ->
Map.put(acc, worker.ref, {child_id, {:chain, chain.id}})
end)
Logger.debug("New ref: #{inspect new_refs}")
state
|> Map.put(:chains, Map.put(state.chains, chain.id, chain))
|> Map.put(:ref_to_id, new_refs)
end
end
defp sup_start_child(state, %Worker{id: id, type: :standalone} = opts) do
# Start a child process
Logger.debug("Starting standalone worker process(#{inspect(id)})")
Ets.insert(state.data_table, {{:worker, id}, opts})
{pid, ref} =
spawn_monitor(fn ->
# Register the worker process.
Registry.register(state.master, {:worker, id}, [])
# Store for user can directly access to the worker.
Process.put({:supervisor, :sup_id}, state.id)
Process.put({:supervisor, :worker_id}, id)
case opts.fun do
{:fun, fun} ->
fun.()
{m, f, a} ->
apply(m, f, a)
end
end)
# Link to child for case supervisor is down.
Process.link(pid)
Ets.insert(state.data_table, {{:worker, :ref, ref}, id, pid, :standalone})
state
end
defp sup_start_child(state, %Worker{id: id, parent: group_id, type: :group} = opts) do
# Start a child process
Logger.debug("Starting child process(#{inspect(id)}) for group #{inspect(group_id)}")
[{_, group}] = Ets.lookup(state.data_table, {:group, group_id})
{:ok, _} = Group.add_worker(group, opts)
state
end
defp sup_start_child(state, %Worker{id: id, parent: chain_id, type: :chain} = opts) do
Logger.debug("Starting child process(#{inspect(id)}) for chain #{inspect(chain_id)}")
[{_, chain}] = Ets.lookup(state.data_table, {:chain, chain_id})
with {:ok, chain} <- Chain.add_worker(chain, opts) do
# Move to Chain module.
Ets.insert(state.data_table, {{:chain, chain_id}, chain})
Logger.debug("part: #{inspect state.id}, added worker to chain: #{inspect chain}")
end
state
end
defp get_group_or_chain(state, chain_id, type) do
case Ets.lookup(state.data_table, {type, chain_id}) do
[] ->
Logger.info("#{inspect state.id}, Chain not found: #{inspect chain_id}")
{:error, :not_found}
[{_, data}] ->
{:ok, data}
end
end
defp add_new_group(state, group) do
group = %Group{group | supervisor: state.master, partition: state.id, data_table: state.data_table}
Registry.register(state.master, {:group, group.id}, [])
Ets.insert(state.data_table, {{:group, group.id}, group})
state
end
defp add_new_chain(state, chain) do
chain = %Chain{chain | supervisor: state.master, partition: state.id, data_table: state.data_table}
Registry.register(state.master, {:chain, chain.id}, [])
Ets.insert(state.data_table, {{:chain, chain.id}, chain})
state
end
@spec do_add_worker(atom(), atom() | tuple(), list(), integer()) :: {:ok, any()} | {:error, any()}
defp do_add_worker(sup_id, :standalone, opts, timeout) do
Logger.debug("Starting child process with options: #{inspect(opts)}")
with {:ok, opts} <- Worker.check_standalone_options(opts),
{:ok, pid} <- verify_and_get_pid(sup_id, opts.id) do
opts =
opts
|> Map.put(:type, :standalone)
call_api(pid, :start_worker, opts, timeout)
else
other ->
Logger.error("cannot add standalone worker, something happened: #{inspect other}, options: #{inspect opts}")
other
end
end
defp do_add_worker(sup_id, {:group_id, group_id} = group, opts, timeout) do
Logger.debug("Starting child process with options: #{inspect(opts)}")
with {:ok, opts} <- Worker.check_group_options([group | opts]),
{:ok, pid} <-verify_and_get_pid(sup_id, opts.id) do
opts =
opts
|> Map.put(:group_id, group_id)
|> Map.put(:type, :group)
call_api(pid, :start_worker, opts, timeout)
else
error ->
Logger.error("cannot add group worker: #{inspect error}, options: #{inspect opts}")
error
end
end
defp do_add_worker(sup_id, {:chain_id, chain_id} = chain, opts, timeout) do
Logger.debug("Starting child process with options: #{inspect(opts)}")
with {:ok, opts} <- Worker.check_chain_options([chain | opts]),
{:ok, pid} <- verify_and_get_pid(sup_id, opts.id) do
opts =
opts
|> Map.put(:chain_id, chain_id)
|> Map.put(:type, :chain)
call_api(pid, :start_worker, opts, timeout)
else
error ->
Logger.error("cannot add chain worker: #{inspect error}, options: #{inspect opts}")
error
end
end
defp get_pid(id) when is_atom(id) do
master = get_master_id(id)
case Process.whereis(master) do
nil ->
{:error, :not_running}
pid ->
{:ok, pid}
end
end
defp get_master_id(id) do
String.to_atom("#{Atom.to_string(id)}_master")
end
defp check_opts(opts) do
with {:ok, opts} <- normalize_opts(opts, @sup_params),
{:ok, opts} <- validate_opts(opts),
{:ok, sup} <- map_to_struct(opts) do
{:ok, sup}
end
end
# Validate the type & value of options.
defp validate_opts(opts) do
with {:ok, opts} <- check_type(opts, :id, &is_atom/1),
opts <- default_sup_opts(opts),
opts <- generic_default_sup_opts(opts),
{:ok, opts} <- check_type(opts, :number_of_partitions, &is_integer/1),
{:ok, opts} <- check_type(opts, :number_of_partitions, &(&1 > 0)),
{:ok, opts} <- check_type(opts, :owner, &is_pid/1),
{:ok, opts} <- check_type(opts, :link, &is_boolean/1) do
{:ok, opts}
else
{:error, reason} = error ->
Logger.error("Error in validating options: #{inspect reason}")
error
end
end
# Set the default options if not provided.
# TO-DO: Merge with generic_default_sup_opts/1.
defp default_sup_opts(opts) do
if Map.has_key?(opts, :number_of_partitions) do
opts
else
Map.put(opts, :number_of_partitions, :erlang.system_info(:schedulers_online))
end
end
# Start the supervisor main processes.
defp start_supervisor(opts = %Supervisor{}, timeout) do
Logger.debug("Starting supervisor with options: #{inspect opts}")
ref = response_ref()
# Start main process of the supervisor
case opts.link do
true ->
Logger.debug("Starting supervisor with link.")
opts = Map.put(opts, :linked_pids, [self()])
spawn_link(@me, :init, [opts, ref])
false ->
Logger.debug("Starting supervisor without link.")
spawn(@me, :init, [opts, ref])
pid when is_pid(pid) ->
Logger.debug("Starting supervisor and link with remote pid.")
opts = Map.put(opts, :linked_pids, [pid])
spawn(@me, :init, [opts, ref])
list_pid when is_list(list_pid) ->
Logger.debug("Starting supervisor and link with remote pids.")
opts = Map.put(opts, :linked_pids, list_pid)
spawn(@me, :init, [opts, ref])
end
api_receiver(ref, timeout)
end
@spec get_partition_id(atom(), integer()) :: atom()
defp get_partition_id(sup_id, partition_id) do
if partition_id < 0 do
sup_id
else
String.to_atom("#{Atom.to_string(sup_id)}_#{inspect partition_id}")
end
end
@spec get_target_partition(atom(), any(), integer()) :: atom()
defp get_target_partition(prefix, data, num_partitions) when is_integer(num_partitions) do
partition_id = get_hash_order(data, num_partitions)
get_partition_id(prefix, partition_id)
end
@spec get_partition_pid(atom(), atom()) :: {:error, atom()} | {:ok, pid()}
defp get_partition_pid(sup_id, partition_id) do
get_table_name(sup_id)
|> Ets.lookup({:partition, partition_id})
|> case do
[{_, pid}] ->
{:ok, pid}
_ ->
{:error, :not_found}
end
end
@spec get_host_partition(atom(), any()) :: {:error, atom()} | {:ok, pid()}
def get_host_partition(sup_id, data) do
with [{_, num}] <- Ets.lookup(get_table_name(sup_id), :number_of_partitions),
partition_id <- get_target_partition(sup_id, data, num),
{:ok, pid} <- get_partition_pid(sup_id, partition_id) do
{:ok, pid}
else
[] ->
Logger.error("not found partition, data: #{inspect data}, sup_id: #{inspect sup_id}")
{:error, :not_found}
{:error, _} = error ->
Logger.error("Get partition pid failed: #{inspect error}, data: #{inspect data}, sup_id: #{inspect sup_id}")
error
end
end
@spec create_table(atom()) :: atom()
defp create_table(table_name) when is_atom(table_name) do
^table_name = Ets.new(table_name, [
:set,
:public,
:named_table,
{:keypos, 1},
{:read_concurrency, true},
{:decentralized_counters, true}
])
end
@spec has_group?(map(), any()) :: boolean()
defp has_group?(%{} = state, group_id) do
case Ets.lookup(state.data_table, {:group, group_id}) do
[] ->
false
_ ->
true
end
end
@spec has_chain?(map(), any()) :: boolean()
defp has_chain?(%{} = state, chain_id) do
case Ets.lookup(state.data_table, {:chain, chain_id}) do
[] ->
false
_ ->
true
end
end
@spec has_group_worker?(map(), any(), any()) :: boolean()
defp has_group_worker?(%{} = state, group_id, worker_id) do
case Ets.lookup(state.data_table, {:worker, {:group, group_id}, worker_id}) do
[] ->
false
_ ->
true
end
end
@spec has_chain_worker?(map(), any(), any()) :: boolean()
defp has_chain_worker?(%{} = state, chain_id, worker_id) do
case Ets.lookup(state.data_table, {:worker, {:chain, chain_id}, worker_id}) do
[] ->
false
_ ->
true
end
end
@spec has_worker?(map(), any()) :: boolean()
defp has_worker?(%{} = state, worker_id) do
case Ets.lookup(state.data_table, {:worker, worker_id}) do
[] ->
false
_ ->
true
end
end
# Check the supervisor is running or not.
# if running get pid of partition.
@spec verify_and_get_pid(atom(), any()) :: {:error, atom()} | {:ok, pid()}
defp verify_and_get_pid(sup_id, id) do
with true <- is_running?(sup_id),
{:ok, pid} <- get_host_partition(sup_id, id) do
{:ok, pid}
else
false ->
Logger.error("Supervisor not running.")
{:error, :not_running}
{:error, reason} = error ->
Logger.error("Get target partition failed: #{inspect reason}")
error
end
end
defp map_to_struct(opts) when is_map(opts) do
{:ok, struct(@me, opts)}
end
end