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
super_worker lib supervisor supervisor.ex
Raw

lib/supervisor/supervisor.ex

defmodule SuperWorker.Supervisor do
@moduledoc """
Documentation for `SuperWorker.Supervisor`.
This module is new model supervisor.
That fix some issues in the old model.
That is an all-in-one supervisor for Elixir application.
New 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.
Each group has a seperated restart strategy that determines how to restart the group when a process dies.
## Chain processes
Chain processes are a set of processes that support for chain prcessing.
Each process in a chain has order to process data.
The output of the previous process is passed to the next process.
## Freedom processes
Freedom processes are independent processes that are started separately.
Each process has its own restart strategy.
All type of processes can be started in parallel & can be stopped individually or in a group.
## Examples
```elixir
# Start a supervisor with 2 partitions & 2 groups:
alias SuperWorker.Supervisor, as: Sup
# Config for supervisor
opts = [id: :sup1, number_of_partitions: 2, link: false]
# Start supervisor
Sup.start(opts)
# Add group in runtime, you also can add group in config.
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_one])
Sup.add_group_worker(:sup1, :group2, fn ->
receice do
msg ->
:ok
end
end, [id: :g2_2])
```
"""
alias SuperWorker.Supervisor.{Group, Chain, Worker, Message}
import SuperWorker.Supervisor.Utils
require Logger
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]
alias __MODULE__
# Default timeout (miliseconds) for API calls.
@default_time 3_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]
## 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 \\ @default_time) when is_list(opts) do
with {:ok, sup} <- check_opts(opts),
false <- is_running?(sup.id) do
start_supervisor(sup, timeout)
else
true ->
Logger.error("SuperWorker, Supervisor, supervisor has id in #{inspect opts} is already running.")
{:error, :already_running}
{:error, _} = error ->
Logger.error("SuperWorker, Supervisor, Eeror 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("SuperWorker, Supervisor, supervisor is not running.")
err
{:ok, pid} ->
Logger.debug("SuperWorker, Supervisor, 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) and group_id != nil 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),
:error <- Registry.meta(sup_id, {:group, group.id}),
{:ok, parititon_id, pid} <- get_host_partition(sup_id, group.id)
do
group = %Group{group | supervisor: sup_id, partition: parititon_id}
call_api(pid, :add_group, group, timeout)
else
wrong ->
Logger.error("SuperWorker, Supervisor, error when adding group: #{inspect wrong}")
{:error, :supervisor_not_found_or_group_exists}
end
end
@doc """
get group structure from supervisor.
"""
@spec get_group(atom(), atom()) :: {:ok, Group.t} | {:error, any()}
def get_group(sup_id, group_id) do
case Registry.meta(sup_id, {:group, group_id}) do
{:ok, group} -> {:ok, group}
:error -> {:error, :group_not_found}
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),
:error <- Registry.meta(sup_id, {:chain, chain.id}),
{:ok, parition_id, pid} <- get_host_partition(sup_id, chain.id)
do
chain = %Chain{chain | supervisor: sup_id, partition: parition_id}
call_api(pid, :add_chain, chain, timeout)
else
wrong ->
Logger.error("SuperWorker, Supervisor, error when adding chain: #{inspect wrong}")
{:error, :supervisor_not_found_or_chain_exists}
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("SuperWorker, Supervisor, group not found.")
{:error, :group_not_found}
sup_id == nil ->
Logger.error("SuperWorker, Supervisor, supervisor not found.")
{:error, :supervisor_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("SuperWorker, Supervisor, group not found.")
{:error, :group_not_found}
sup_id == nil ->
Logger.error("SuperWorker, Supervisor, supervisor not found.")
{:error, :supervisor_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("SuperWorker, Supervisor, group not found.")
{:error, :grou_not_found}
sup_id == nil ->
Logger.error("SuperWorker, Supervisor, supervisor not found.")
{:error, :supervisor_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) do
case Registry.meta(sup_id, {:chain, chain_id}) do
{:ok, chain} -> {:ok, chain}
:error -> {:error, :chain_not_found}
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,
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.
pid =
case Registry.start_link(keys: :duplicate, name: opts.id, partitions: opts.number_of_partitions) do
{:ok, pid} -> pid
{:error, _} ->
Logger.error("SuperWorker, Supervisor, failed to start registry with name #{inspect opts.id}")
raise "SuperWorker, Supervisor, failed to start registry with name #{inspect opts.id}"
end
Registry.register(state.master, :master, opts)
Registry.put_meta(state.master, :supervisor, opts)
# Turn main process to system process.
Process.flag(:trap_exit, true)
list_partitions = init_additional_partitions(opts)
started_partitions =
Enum.reduce(1..length(list_partitions), [], fn _, acc ->
receive do
{:partition_started, id} ->
Logger.debug("SuperWorker, Supervisor, supervisor #{inspect state.id} received started partition msg from #{inspect id}")
[id | acc]
after 3_000 ->
Logger.debug("SuperWorker, Supervisor, supervisor #{inspect state.id} timeout when starting partition. Current list: #{inspect acc}")
acc
end
end)
if length(started_partitions) != length(list_partitions) do
Logger.error("SuperWorker, Supervisor, supervisor #{inspect state.id} failed to start partitions.")
api_response(ref, {:error, :failed_to_start_partitions})
else
state =
state
|> Map.put(:partitions , list_partitions)
|> Map.put(:role, :master)
|> Map.put(:registry_pid, pid)
Logger.debug("SuperWorker, Supervisor, 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
end
def child_spec(opts) do
%{
id: Keyword.get(opts, :id, Supervisor), # default id is module name
start: {Supervisor, :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,
prefix: "{#{inspect opts.master}, #{inspect opts.id}}"
}
# Start the main loop
pid = spawn_link(Supervisor, :start_partition, [state, opts])
Logger.debug("SuperWorker, Supervisor, #{state.prefix} initialized, pid: #{inspect pid}")
{:ok, opts.id, pid}
end
def start_partition(state, opts) do
{:ok, _} = Registry.register(state.master, {:partition , state.id}, state.id)
{:ok, _} = Registry.register(state.master, :all_partitions, state.id)
send(get_master_id(opts.master), {:partition_started, state.id})
# Turn partition process to system process.
Process.flag(:trap_exit, true)
main_loop(state, opts)
end
defp init_additional_partitions(opts) do
partitions = opts.number_of_partitions
Enum.map(0..partitions - 1, fn i ->
Logger.debug("SuperWorker, Supervisor, [#{inspect opts.id}] add partition: #{inspect i}")
opts =
opts
|> Map.put(:master, opts.id)
|> Map.put(:id, String.to_atom("#{Atom.to_string(opts.id)}_#{i}"))
|> Map.put(:role, :partition)
{:ok, partition, pid} = init_partition(opts)
{partition, pid}
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("SuperWorker, Supervisor, #{ state.prefix} received a api message: #{inspect(msg)}")
process_api_message(state, sup_opts, msg)
{:DOWN, _ref, :process, pid, reason} = msg ->
Logger.debug("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ state.prefix} Stopping supervisor partition, for #{inspect self()}")
# Stop the supervisor.
shutdown(state, type)
send(get_master_id(state.master), {:partition_stopped, state.id})
unknown ->
Logger.warning("SuperWorker, Supervisor, #{ state.prefix} main_loop, unknown message: #{inspect(unknown)}")
main_loop(state, sup_opts)
end
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} #{inspect self()} main loop exited.")
end
defp shutdown(state, :kill) do
Logger.debug("SuperWorker, Supervisor, 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("SuperWorker, Supervisor, #{ state.prefix} Exit message from: #{inspect from}, reason: #{inspect reason}")
if from in sup_opts.linked_pids do
Logger.warning("SuperWorker, Supervisor, #{ state.prefix} exited follow external process (crashed): #{inspect from}")
raise "#{inspect state.master} crashed follow external process: #{inspect from}"
else
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} skipped exit for internal process: #{inspect from}, reason: #{inspect reason}")
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 ->
cond do
!has_group?(state, opts.parent) ->
:group_not_found
has_group_worker?(state, opts.parent, opts.id) ->
:worker_already_exists
true ->
true
end
:chain ->
cond do
!has_chain?(state, opts.parent) ->
:chain_not_found
has_chain_worker?(state, opts.parent, opts.id) ->
:worker_already_exists
true ->
true
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("SuperWorker, Supervisor, #{ state.prefix} Everything is fine, starting child process with options: #{inspect(opts)}")
state = sup_start_child(state, opts)
api_response(ref, {:ok, opts.id})
state
else # not found group or chain, return error to the caller.
Logger.error("SuperWorker, Supervisor, #{ state.prefix} Error when starting worker: #{inspect runable}")
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("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, supervisor #{ 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("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ state.prefix} Chain not found: #{inspect chain_id}")
api_response(ref, error)
{:ok, chain} ->
msg = Message.new(from, nil, data)
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} Add data to chain: #{inspect chain_id}, msg: #{inspect msg}")
result = Chain.new_data(chain, msg)
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} Added 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("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ state.prefix} Starting worker process, worker id: #{inspect(worker_id)}, group id: #{inspect(group_id)}")
{:ok, group} = Registry.meta(state.master, {: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 Registry.meta(state.master, {:group, group.id}) do
:error ->
Logger.debug("SuperWorker,Supervisor, #{ state.prefix} Adding group: #{inspect(group.id)}")
state = add_new_group(state, group)
# Send the response to the caller.
api_response(ref, {:ok, group.id})
state
{:ok, _} ->
Logger.error("SuperWorker, Supervisor, #{ state.prefix} Group already exists: #{inspect(group.id)}")
api_response(ref, {:error, :already_exists})
state
end
|> main_loop(sup_opts)
end
# get group info from api.
defp process_api_message(state, sup_opts, {:get_group, ref, group_id}) do
result =
case Registry.meta(state.master, {:group, group_id}) do
:error ->
Logger.error("SuperWorker, Supervisor, #{ state.prefix} Group not found: #{inspect(group_id)}")
{:error, :group_not_found}
{:ok, _group} = result ->
result
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 Registry.meta(state.master, {:chain, chain.id}) do
:error ->
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} Adding chain: #{inspect(chain.id)}")
state = add_new_chain(state, chain)
# Send the response to the caller.
api_response(ref, {:ok, chain.id})
state
{:ok, _} ->
Logger.error("SuperWorker, Supervisor, #{ state.prefix} Chain already exists: #{inspect(chain.id)}")
api_response(ref, {:error, :already_exists})
state
end
|> main_loop(sup_opts)
end
# Stop supervisor from api.
defp process_api_message(state, _sup_opts, {:stop, ref, type}) do
Logger.info("SuperWorker, Supervisor, #{ state.prefix} Stopping supervisor, request from #{inspect ref}")
list_partitions = Registry.lookup(state.master, :all_partitions)
# Send shutdown signal to all partitions.
Enum.each(list_partitions, fn {pid, id} ->
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} Sending shutdown signal to partition: #{inspect id} (#{inspect pid})")
send(pid, {:stop_partition, type})
end)
stopped_partitions = Enum.reduce(1..length(list_partitions), [], fn _, acc ->
receive do
{:partition_stopped, id} ->
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} received stopped partition msg from #{inspect id}")
[id | acc]
after 3_000 ->
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} timeout when stopping partition. Current list: #{inspect acc}")
acc
end
end)
result =
if length(list_partitions) != length(stopped_partitions) do
Logger.error("SuperWorker, Supervisor, #{ state.prefix} failed to stop partitions.")
{:error, :failed_to_stop_partitions}
else
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} stopped all partitions.")
{:ok, :stopped}
end
# stop worker on master.
shutdown(state, type)
# TO-DO: Clean KV store for supervisor.
api_response(ref, result)
exit(:normal)
end
defp process_api_message(state, sup_opts, unknown_msg) do
Logger.warning("SuperWorker, Supervisor, #{ 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("SuperWorker, Supervisor, #{ state.prefix} Ignore died process (restarting): #{inspect(pid)}")
main_loop(state, sup_opts)
end
defp process_worker_down(state, sup_opts, {:DOWN, _ref, :process, pid, :removed}) do
Logger.debug("SuperWorker, Supervisor, #{ state.prefix} Ignore died process (removed by user): #{inspect(pid)}")
main_loop(state, sup_opts)
end
defp process_worker_down(state, sup_opts, {:DOWN, ref, :process, pid, reason}) do
Logger.debug("SuperWorker, Supervisor, child process died: #{inspect(pid)}, ref: #{inspect ref}, reason: #{inspect(reason)}")
state =
case Registry.lookup(state.master, {:worker, :ref, ref}) do
[] ->
Logger.debug("SuperWorker, Supervisor, child is not found in table: #{inspect ref}, maybe already stopped.")
state
[{pid, {type, worker_id}} = ref_data] ->
Logger.debug("SuperWorker, Supervisor, child found: #{inspect pid}, restarting. meta: #{inspect ref_data}")
case type do
:standalone ->
{:ok, child} = Registry.meta(state.master, {:worker, :standalone, worker_id})
restart_standalone(state, child, {pid, reason})
{:group, group_id} = group_type ->
with {:ok, child} <- Registry.meta(state.master, {:worker, group_type, worker_id}),
{:ok, group} <- Registry.meta(state.master, {:group, group_id}) do
restart_group(state, group, child, {pid, reason})
else
wrong ->
Logger.error("SuperWorker, Supervisor, #{ state.prefix} Error when restarting group worker: #{inspect wrong}")
state
end
{:chain, _chain_id} = chain_type ->
{:ok, chain} = Registry.meta(state.master, {:worker, chain_type, worker_id})
restart_chain(state, chain, worker_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("SuperWorker, Supervisor, #{inspect state.id}, :permanent, restarting #{inspect pid}")
# Restart the child process
sup_start_child(state, child)
:transient when reason != :normal ->
Logger.debug("SuperWorker, Supervisor, #{inspect state.id}, :transient, reason down: #{inspect reason} restarting #{inspect pid}")
# Restart the child process
sup_start_child(state, child)
:transient ->
Logger.debug("SuperWorker, Supervisor, #{inspect state.id}, :transient, ignore restarting #{inspect pid}")
state
:temporary ->
Logger.debug("SuperWorker, Supervisor, #{inspect state.id}, :temporary, ignore restarting #{inspect pid}")
state
end
end
defp restart_group(state, %{restart_strategy: :one_for_one} = group, worker = %Worker{}, {pid, reason}) do
case reason do
:normal ->
Logger.debug("SuperWorker, Supervisor, #{inspect state.id}, Child process(#{inspect(pid)}) is shutdown with reason :normal, ignore restarting.")
state
reason ->
Logger.debug("SuperWorker, Supervisor, #{inspect state.id}, Child process(#{inspect(pid)}) is down with reason: #{inspect reason}, restarting.")
Group.restart_worker(group, worker)
state
end
end
defp restart_group(state, %{restart_strategy: :one_for_all} = group, _child_id, {pid, reason}) do
case reason do
:normal ->
Logger.debug("SuperWorker, Supervisor, #{inspect state.id}, Child process(#{inspect(pid)}) is :normal shutdown, ignore restarting.")
state
reason ->
Logger.debug("SuperWorker, Supervisor, #{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("SuperWorker, Supervisor, #{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("SuperWorker, Supervisor, worker(#{inspect child_id}) process(#{inspect(pid)}) is normal, ignore restarting.")
state
_ ->
Logger.debug("SuperWorker, Supervisor, 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("SuperWorker, Supervisor, child process(#{inspect(pid)}) is normal, ignore restarting.")
state
_ ->
Logger.debug("SuperWorker, Supervisor, child process(#{inspect(pid)}) is down, restarting...")
old_ref_keys = Enum.reduce(chain.workers, [], fn {_, worker}, acc -> [worker.ref | acc] end)
Logger.debug("SuperWorker, Supervisor, chain, 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("SuperWorker, Supervisor, 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} = worker) do
# Start a child process
Logger.debug("SuperWorker, Supervisor, starting standalone worker process(#{inspect(id)})")
Registry.put_meta(state.master, {:standalone, {:worker, id}}, worker)
{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 worker.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)
Registry.register(state.master, {:worker, :ref, ref}, {:standalone, worker.id})
state
end
defp sup_start_child(state, %Worker{id: id, parent: group_id, type: :group} = opts) when group_id != nil do
# Start a child process
Logger.debug("SuperWorker, Supervisor, starting child process(#{inspect(id)}) for group #{inspect(group_id)}")
{:ok, group} = Registry.meta(state.master, {: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("SuperWorker, Supervisor, starting child process(#{inspect(id)}) for chain #{inspect(chain_id)}")
{:ok, chain} = Registry.meta(state.master, {:chain, chain_id})
with {:ok, chain} <- Chain.add_worker(chain, opts) do
Registry.put_meta(state.master, {:chain, chain_id}, chain)
Logger.debug("SuperWorker, Supervisor, part: #{inspect state.id}, added worker to chain: #{inspect chain}")
end
state
end
defp get_group_or_chain(state, id, type) do
case Registry.meta(state.master, {type, id}) do
:error ->
Logger.warning("SuperWorker, Supervisor, #{inspect state.id}, Chain/Group not found: #{inspect id}")
{:error, {:not_found, type}}
{:ok, _} = result ->
result
end
end
defp add_new_group(state, group) do
group = %Group{group | supervisor: state.master, partition: state.id}
Registry.put_meta(state.master, {:group, group.id}, group)
state
end
defp add_new_chain(state, chain) do
chain = %Chain{chain | supervisor: state.master, partition: state.id}
Registry.put_meta(state.master, {: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("SuperWorker, Supervisor, starting standalone 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("SuperWorker, Supervisor, 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) when group_id != nil do
Logger.debug("SuperWorker, Supervisor, starting worker group(#{inspect group_id}) 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(:parent, group_id)
|> Map.put(:type, :group)
Logger.debug("SuperWorker, Supervisor, start call :start_worker api with opts: #{inspect opts}")
call_api(pid, :start_worker, opts, timeout)
else
error ->
Logger.error("SuperWorker, Supervisor, 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("SuperWorker, Supervisor, 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(:parent, chain_id)
|> Map.put(:type, :chain)
call_api(pid, :start_worker, opts, timeout)
else
error ->
Logger.error("SuperWorker, Supervisor, 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("SuperWorker, Supervisor, 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("SuperWorker, Supervisor, starting supervisor with options: #{inspect opts}")
ref = response_ref()
# Start main process of the supervisor
case opts.link do
true ->
Logger.debug("SuperWorker, Supervisor, starting supervisor with link.")
opts = Map.put(opts, :linked_pids, [self()])
spawn_link(Supervisor, :init, [opts, ref])
false ->
Logger.debug("SuperWorker, Supervisor, starting supervisor without link.")
spawn(Supervisor, :init, [opts, ref])
pid when is_pid(pid) ->
Logger.debug("SuperWorker, Supervisor, starting supervisor and link with remote pid.")
opts = Map.put(opts, :linked_pids, [pid])
spawn(Supervisor, :init, [opts, ref])
list_pid when is_list(list_pid) ->
Logger.debug("SuperWorker, Supervisor, starting supervisor and link with remote pids.")
opts = Map.put(opts, :linked_pids, list_pid)
spawn(Supervisor, :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
case Registry.lookup(sup_id, {:partition, partition_id}) do
[{pid, ^partition_id}] ->
{:ok, pid}
[] ->
Logger.error("SuperWorker, Supervisor, supervisor #{inspect sup_id} partition not found: #{inspect partition_id}")
{:error, {:partition_not_found, partition_id}}
end
end
@spec get_host_partition(atom, any) :: {:error, atom} | {:ok, atom, pid}
def get_host_partition(sup_id, data) do
with {:ok, sup} <- Registry.meta(sup_id, :supervisor),
partition_id <- get_target_partition(sup_id, data, sup.number_of_partitions),
{:ok, pid} <- get_partition_pid(sup_id, partition_id) do
{:ok, partition_id, pid}
else
:error ->
Logger.error("SuperWorker, Supervisor, not found partition, data: #{inspect data}, sup_id: #{inspect sup_id}")
{:error, {:supervisor_not_found, sup_id}}
{:error, _} = error ->
Logger.error("SuperWorker, Supervisor, get partition pid failed: #{inspect error}, data: #{inspect data}, sup_id: #{inspect sup_id}")
error
end
end
@spec has_group?(map(), any()) :: boolean()
defp has_group?(%{} = state, group_id) do
case Registry.meta(state.master, {:group, group_id}) do
:error ->
false
{:ok, _} ->
true
end
end
@spec has_chain?(map(), any()) :: boolean()
defp has_chain?(%{} = state, chain_id) do
case Registry.meta(state.master, {:chain, chain_id}) do
:error ->
false
{:ok, _} ->
true
end
end
@spec has_group_worker?(map(), any(), any()) :: boolean()
defp has_group_worker?(%{} = state, group_id, worker_id) do
case Registry.lookup(state.master, {: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 Registry.lookup(state.master, {: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 Registry.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, _partition_id, pid} <- get_host_partition(sup_id, id) do
{:ok, pid}
else
false ->
Logger.error("SuperWorker, Supervisor, supervisor #{inspect sup_id} not running.")
{:error, :not_running}
{:error, reason} = error ->
Logger.error("SuperWorker, Supervisor, get target partition failed: #{inspect reason}")
error
end
end
defp map_to_struct(opts) when is_map(opts) do
{:ok, struct(Supervisor, opts)}
end
end