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