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/group.ex
defmodule SuperWorker.Supervisor.Group do
@moduledoc """
Documentation for `SuperWorker.Supervisor.Group`.
"""
# Parameters for group.
@group_params [:id, :restart_strategy, :type, :max_restarts, :max_seconds, :auto_restart_time]
@enforce_keys [:id]
defstruct [
# group id, unique in supervior.
:id,
# default restart strategy for group is :one_for_all.
restart_strategy: :one_for_all,
# supervisor id (atom)
supervisor: nil,
# table data of supervisors.
table: nil
]
@type t :: %__MODULE__{
id: any,
restart_strategy: atom,
supervisor: atom,
table: atom
}
alias __MODULE__
alias SuperWorker.Supervisor.{Worker, Db, Validator, Constants}
require Logger
require SuperWorker.Log
## Public functions
@doc """
Check, validate and convert key-value pairs to struct.
"""
@spec check_options([keyword]) :: {:ok, %Group{}} | {:error, atom | {atom, any}}
def check_options(options) do
with {:ok, options} <- Validator.normalize_options(options, @group_params),
{:ok, options} <- validate_options(options),
{:ok, group} <- to_struct(options) do
{:ok, group}
else
{:error, reason} = error ->
Logger.error("SuperWorker, Group, incorrect options, #{inspect(reason)}")
error
end
end
@doc """
Get worker from the group.
"""
def get_worker(%Group{} = group, worker_id) do
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, supervisor #{inspect(group.supervisor)}, group #{inspect(group.id)}, get_worker: #{inspect(worker_id)}"
end)
worker_id =
case worker_id do
{_, id} -> id
_ -> worker_id
end
Db.get_worker_info(group.table, worker_id, {:group, group.id})
end
@doc """
Get all workers from the group.
"""
def get_all_workers(%Group{} = group) do
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, get_all_workers for supervisor #{inspect(group.supervisor)}"
end)
Db.get_worker_infos_by_parent(group.table, {:group, group.id})
end
@spec count_workers(%Group{}) :: non_neg_integer()
def count_workers(%Group{} = group) do
{:ok, workers} = get_all_workers(group)
length(workers)
end
@doc """
Check if worker exists in the group.
"""
def worker_exists?(group, worker_id) do
case get_worker(group, worker_id) do
{:ok, _} -> true
{:error, _} -> false
end
end
@doc """
A internal function. Add a worker to the group.
"""
def add_worker(group = %Group{}, %Worker{} = worker) do
case get_worker(group, worker.id) do
{:ok, _} ->
{:error, :worker_exists}
{:error, _} ->
worker = %Worker{worker | parent: group.id}
worker =
if !worker.id do
%Worker{worker | id: SuperWorker.Supervisor.Utils.random_id()}
else
worker
end
Db.put_worker_info(group.table, worker)
spawn_worker(group, worker)
end
end
@doc """
A internal function. Restart a worker in the group.
"""
def restart_worker(group = %Group{}, worker = %Worker{}) do
SuperWorker.Log.debug(fn -> "SuperWorker, Group, restart worker #{inspect(worker)}" end)
case kill_worker(group, worker, :restart) do
{:ok, _} ->
# Worker was alive and killed, now spawn a new one
spawn_worker(group, worker)
{:error, :not_alive} ->
# Worker process is already dead, clean up ETS and spawn new one
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, worker #{inspect(worker.id)} already dead, spawning new one"
end)
Db.delete_worker_info(group.table, worker.id, {:group, group.id})
spawn_worker(group, worker)
{:error, :not_found} ->
# Worker not in ETS, just spawn a new one
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, worker #{inspect(worker.id)} not found in ETS, spawning new one"
end)
spawn_worker(group, worker)
{:error, reason} ->
Logger.error(
"SuperWorker, Group, failed to kill worker #{inspect(worker.id)} before restart, reason: #{inspect(reason)}"
)
{:error, :kill_failed}
end
end
def restart_worker(group = %Group{}, worker_id) do
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, restart worker by id #{inspect(worker_id)}"
end)
case get_worker(group, worker_id) do
{:ok, worker} ->
restart_worker(group, worker)
{:error, _} = error ->
Logger.error(
"SuperWorker, Group, cannot get worker #{inspect(worker_id)}, #{inspect(error)}"
)
{:error, :worker_not_found}
end
end
def remove_worker(group = %Group{}, worker_id) do
if worker_exists?(group, worker_id) do
with {:ok, worker} <- get_worker(group, worker_id),
{:ok, _} <- kill_worker(group, worker, :removed) do
table = group.table
parent = {:group, group.id}
Db.delete_worker_by_id(table, worker_id, parent)
Db.delete_worker_info(table, worker_id, parent)
{:ok, :worker_removed}
else
{:error, reason} = error ->
Logger.error(
"SuperWorker, Group, failed to kill worker #{inspect(worker_id)} in group #{inspect(group.id)}, error: #{inspect(reason)}"
)
error
end
else
{:error, :worker_not_found}
end
end
def kill_worker(group = %Group{}, worker = %Worker{}, reason) do
with {:ok, {_, pid}} <- Db.get_worker_by_id(group.table, worker.id, {:group, group.id}) do
if Process.alive?(pid) do
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, group: #{inspect(group.id)}, kill_worker: #{inspect(worker)}, reason: #{inspect(reason)}"
end)
Process.exit(pid, reason)
{:ok, :killed}
else
{:error, :not_alive}
end
else
_ ->
{:error, :not_found}
end
end
def kill_worker(group = %Group{}, worker_id, reason) do
case get_worker(group, worker_id) do
{:ok, worker} ->
kill_worker(group, worker, reason)
{:error, _} ->
{:error, :worker_not_found}
end
end
def kill_all_workers(group = %Group{}, reason \\ :kill) do
{:ok, list_worker} = get_all_workers(group)
results =
Enum.map(list_worker, fn worker ->
case kill_worker(group, worker, reason) do
{:ok, _} ->
:ok
{:error, kill_reason} ->
Logger.warning(
"SuperWorker, Group, failed to kill worker #{inspect(worker.id)} in group #{inspect(group.id)}, reason: #{inspect(kill_reason)}"
)
{:error, worker.id, kill_reason}
end
end)
errors = Enum.filter(results, &match?({:error, _, _}, &1))
if Enum.empty?(errors) do
:ok
else
{:error, errors}
end
end
defp spawn_worker(group = %Group{}, %Worker{} = worker) do
SuperWorker.Log.debug(fn -> "SuperWorker, Group, spawn_worker: #{inspect(worker)}" end)
try do
do_spawn_worker(group, worker)
{:ok, group}
catch
:exit, reason ->
Logger.error(
"SuperWorker, Group, failed to spawn worker #{inspect(worker.id)}: #{inspect(reason)}"
)
{:error, :spawn_failed}
end
end
@spec broadcast(%Group{}, any()) :: :ok | {:error, list()}
def broadcast(group = %Group{}, message) do
case Group.get_all_workers(group) do
{:ok, workers} ->
results =
Enum.map(
workers,
fn %Worker{id: worker_id} ->
case Db.get_worker_by_id(group.table, worker_id, {:group, group.id}) do
{:ok, {_ref, pid}} ->
send(pid, message)
:ok
{:error, _reason} = error ->
Logger.error(
"SuperWorker, Group, cannot get worker pid for #{inspect(worker_id)} in group #{inspect(group.id)}, reason: #{inspect(error)}"
)
error
end
end
)
errors = Enum.filter(results, &match?({:error, _}, &1))
if Enum.empty?(errors) do
:ok
else
{:error, errors}
end
{:error, reason} ->
{:error, reason}
end
end
def send_message(group = %Group{}, worker_id, message) do
with {:ok, {_, pid}} <- Db.get_worker_by_id(group.table, worker_id, {:group, group.id}) do
send(pid, message)
else
error ->
Logger.error(
"SuperWorker, Group, send to worker #{inspect(worker_id)} failed, #{inspect(error)}"
)
{:error, :cannot_send}
end
end
## Private functions
defp do_spawn_worker(group, %Worker{} = worker) do
{pid, ref} =
case worker.fun do
{:gen_server, {m, f, a}} ->
{:ok, pid} = apply(m, f, a)
ref = Process.monitor(pid)
{pid, ref}
_ ->
spawn_monitor(fn ->
Process.put({:supervisor, :sup_id}, group.supervisor)
Process.put({:supervisor, :group_id}, group.id)
Process.put({:supervisor, :worker_id}, worker.id)
if worker.name do
if Process.whereis(worker.name) do
Logger.warning(
"SuperWorker, Group, worker name already registered: #{inspect(worker.name)}"
)
else
Process.register(self(), worker.name)
end
end
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, worker #{inspect(worker.id)} started, fun: #{inspect(worker.fun)}"
end)
case worker.fun do
{m, f, a} ->
apply(m, f, a)
{:fun, fun} ->
fun.()
end
end)
end
Db.put_worker(group.table, ref, worker.id, {:group, group.id}, pid)
SuperWorker.Log.debug(fn ->
"SuperWorker, Group, spawned worker #{inspect(worker.id)}, pid: #{inspect(pid)}, ref: #{inspect(ref)}"
end)
try do
Process.link(pid)
catch
:exit, reason ->
Logger.error(
"SuperWorker, Group, failed to link worker #{inspect(worker.id)}: #{inspect(reason)}"
)
Process.exit(pid, :kill)
exit(reason)
end
worker
|> Map.put(:pid, pid)
|> Map.put(:ref, ref)
end
defp validate_restart_strategy(options) do
if options.restart_strategy in Constants.Strategies.group_restart_strategies() do
{:ok, options}
else
{:error, "Invalid group restart strategy, #{inspect(options.restart_strategy)}"}
end
end
defp validate_options(options) do
validate_restart_strategy(options)
end
defp to_struct(options) when is_map(options) do
fields =
%Group{id: nil}
|> Map.from_struct()
|> Map.keys()
result =
%Group{} =
Enum.reduce(fields, %Group{id: nil}, fn field, acc ->
if Map.has_key?(options, field) do
%{acc | field => Map.get(options, field)}
else
acc
end
end)
{:ok, result}
end
end