Packages

Librería base para ejecución de comandos del sistema y gestión de tareas asíncronas con resultados estructurados.

Retired package: Release invalid - Versión publicada por error

Current section

Files

Jump to
argos lib parallel leader.ex
Raw

lib/parallel/leader.ex

defmodule Argos.Parallel.Leader do
@moduledoc """
Coordinador que supervisa el ciclo de vida lógico de los workers.
Arranca workers a través del `DynamicSupervisor`, recibe sus eventos y los
retransmite a los suscriptores, además de mantener un estado mínimo por worker.
"""
use GenServer
require Logger
alias Argos.Parallel.{LeaderState, Worker, WorkerState}
alias Argos.Parallel.Logger, as: ParallelLogger
@default_call_timeout 30_000
def start_link(opts \\ []) do
name = Keyword.get(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, opts, name: name)
end
@doc """
Inicia múltiples workers con las especificaciones dadas.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
## Ejemplos
iex> Argos.Parallel.Leader.start_workers([%{id: :w1, tasks: [fn -> :ok end]}])
[{:ok, #PID<0.123.0>}]
"""
@spec start_workers(GenServer.server(), [term()], keyword()) :: [term()]
def start_workers(server \\ __MODULE__, worker_specs, opts \\ [])
when is_list(worker_specs) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, {:start_workers, worker_specs}, timeout)
end
@doc """
Suscribe el proceso actual para recibir eventos de workers.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec subscribe(GenServer.server(), keyword()) :: :ok
def subscribe(server \\ __MODULE__, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, {:subscribe, self()}, timeout)
end
@doc """
Desuscribe el proceso actual de los eventos de workers.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec unsubscribe(GenServer.server(), keyword()) :: :ok
def unsubscribe(server \\ __MODULE__, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, {:unsubscribe, self()}, timeout)
end
@doc """
Obtiene el estado interno del líder.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_state(GenServer.server(), keyword()) :: map()
def get_state(server \\ __MODULE__, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, :get_state, timeout)
end
def stop(server \\ __MODULE__) do
GenServer.stop(server, :normal)
end
@doc """
Detiene un worker específico por su ID.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
## Retorna
`{:ok, :stopped}` si el worker se detuvo correctamente,
`{:error, :not_found}` si el worker no existe, o `{:error, reason}` si hubo otro error.
"""
@spec stop_worker(GenServer.server(), term(), keyword()) ::
{:ok, :stopped} | {:error, term()}
def stop_worker(server \\ __MODULE__, worker_id, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, {:stop_worker, worker_id}, timeout)
end
@impl true
def init(opts) do
dyn_sup = Keyword.get(opts, :dyn_sup_name, Argos.Parallel.DynamicSupervisor)
name = Keyword.get(opts, :name, __MODULE__)
state = %LeaderState{dyn_sup: dyn_sup, name: name}
{:ok, state}
end
@impl true
def handle_call({:subscribe, pid}, _from, %LeaderState{subscribers: subs} = state) do
Process.monitor(pid)
{:reply, :ok, %{state | subscribers: MapSet.put(subs, pid)}}
end
def handle_call({:unsubscribe, pid}, _from, %LeaderState{subscribers: subs} = state) do
{:reply, :ok, %{state | subscribers: MapSet.delete(subs, pid)}}
end
def handle_call(:get_state, _from, state) do
public_state =
state
|> Map.from_struct()
|> Map.update!(:workers, & &1)
{:reply, public_state, state}
end
def handle_call({:start_workers, worker_specs}, _from, %LeaderState{dyn_sup: dyn_sup} = state) do
:telemetry.execute(
[:argos, :parallel, :workers, :start],
%{count: length(worker_specs)},
%{leader: state.name}
)
{results, new_workers} =
Enum.map_reduce(worker_specs, %{}, fn spec, acc ->
{id, tasks, log_enabled} = normalize_spec(spec)
if log_enabled do
try do
ParallelLogger.register_worker(id, true)
rescue
e ->
Logger.warning("Failed to register worker logger: worker_id=#{inspect(id)}, error=#{inspect(e)}")
:ok
end
end
child_spec = %{
id: {:parallel_worker, id},
start:
{Worker, :start_link,
[
%{
leader: self(),
id: id,
tasks: tasks
}
]},
restart: :transient,
shutdown: 5_000,
type: :worker
}
case DynamicSupervisor.start_child(dyn_sup, child_spec) do
{:ok, pid} ->
Logger.debug("Started worker #{inspect(id)} -> #{inspect(pid)}")
worker_state = WorkerState.new(id) |> WorkerState.running(0, length(tasks), nil)
{{:ok, pid}, Map.put(acc, id, worker_state)}
{:ok, pid, _info} ->
Logger.debug("Started worker #{inspect(id)} (info) -> #{inspect(pid)}")
worker_state = WorkerState.new(id) |> WorkerState.running(0, length(tasks), nil)
{{:ok, pid}, Map.put(acc, id, worker_state)}
{:error, reason} ->
Logger.warning("Could not start worker #{inspect(id)}: #{inspect(reason)}")
maybe_unregister_logger(log_enabled, id)
worker_state = WorkerState.new(id) |> WorkerState.error(reason, nil)
{{:error, reason}, Map.put(acc, id, worker_state)}
end
end)
new_state = %{state | workers: Map.merge(state.workers, new_workers)}
successful = Enum.count(results, fn {status, _} -> status == :ok end)
failed = length(results) - successful
:telemetry.execute(
[:argos, :parallel, :workers, :started],
%{successful: successful, failed: failed, total: length(results)},
%{leader: state.name}
)
{:reply, results, new_state}
end
def handle_call(
{:stop_worker, worker_id},
_from,
%LeaderState{dyn_sup: dyn_sup, workers: workers, subscribers: subs, name: name} = state
) do
child_id = {:parallel_worker, worker_id}
result = stop_worker_internal(dyn_sup, child_id, worker_id, name, subs)
new_workers = update_workers_after_stop(workers, worker_id, result)
new_state = %{state | workers: new_workers}
{:reply, result, new_state}
end
defp stop_worker_internal(dyn_sup, child_id, worker_id, name, subs) do
children = DynamicSupervisor.which_children(dyn_sup)
find_and_stop_worker(children, child_id, worker_id, dyn_sup, name, subs)
end
defp find_and_stop_worker(children, child_id, worker_id, dyn_sup, name, subs) do
case Enum.find(children, fn {id, _pid, _type, _modules} -> id == child_id end) do
{^child_id, pid, _type, _modules} when is_pid(pid) ->
terminate_and_notify(dyn_sup, worker_id, pid, name, subs)
nil ->
Logger.warning("Worker #{inspect(worker_id)} not found")
{:error, :not_found}
end
end
defp terminate_and_notify(dyn_sup, worker_id, pid, name, subs) do
case DynamicSupervisor.terminate_child(dyn_sup, pid) do
:ok ->
Logger.debug("Stopped worker #{inspect(worker_id)} -> #{inspect(pid)}")
try do
ParallelLogger.unregister_worker(worker_id)
rescue
e ->
Logger.warning("Failed to unregister worker logger during stop: worker_id=#{inspect(worker_id)}, error=#{inspect(e)}")
:ok
end
notify_subscribers(subs, name, worker_id)
{:ok, :stopped}
{:error, :not_found} ->
Logger.warning("Worker #{inspect(worker_id)} not found in supervisor")
{:error, :not_found}
end
end
defp notify_subscribers(subs, name, worker_id) do
event = %{
worker_id: worker_id,
type: :stopped,
data: %{reason: :stopped_manually},
leader: name,
timestamp: DateTime.utc_now()
}
Enum.each(subs, fn sub_pid ->
send(sub_pid, {:parallel_event, name, event})
end)
end
defp update_workers_after_stop(workers, worker_id, result) do
case result do
{:ok, :stopped} ->
Map.update(workers, worker_id, nil, fn worker ->
current_task = worker.current_task || 0
WorkerState.error(worker, :stopped_manually, current_task)
end)
{:error, _reason} ->
workers
end
end
@impl true
def handle_cast({:worker_msg, worker_id, msg}, %LeaderState{subscribers: subs} = state) do
new_state = update_state_with_worker_msg(state, worker_id, msg)
event = build_event(state.name, worker_id, msg)
Enum.each(subs, fn pid ->
send(pid, {:parallel_event, state.name, event})
end)
{:noreply, new_state}
end
@impl true
def handle_info({:DOWN, _ref, :process, pid, _reason}, %LeaderState{subscribers: subs} = state) do
{:noreply, %{state | subscribers: MapSet.delete(subs, pid)}}
end
def handle_info(msg, state) do
Logger.debug("Leader received unexpected message: #{inspect(msg)}")
{:noreply, state}
end
defp normalize_spec({id, tasks}) when is_list(tasks), do: {id, tasks, false}
defp normalize_spec(%{id: id, tasks: tasks} = spec) when is_list(tasks) do
log_enabled = Map.get(spec, :log, false)
{id, tasks, log_enabled}
end
defp normalize_spec(tasks) when is_list(tasks), do: {make_ref(), tasks, false}
defp normalize_spec(other), do: {make_ref(), [other], false}
defp build_event(leader_name, worker_id, {:progress, task_index, total, percent}) do
%{
worker_id: worker_id,
type: :progress,
data: %{task_index: task_index, total: total, percent: percent},
leader: leader_name,
timestamp: DateTime.utc_now()
}
end
defp build_event(leader_name, worker_id, {:result, task_index, result}) do
%{
worker_id: worker_id,
type: :result,
data: %{task_index: task_index, result: result},
leader: leader_name,
timestamp: DateTime.utc_now()
}
end
defp build_event(leader_name, worker_id, {:error, task_index, reason}) do
%{
worker_id: worker_id,
type: :error,
data: %{task_index: task_index, reason: inspect(reason)},
leader: leader_name,
timestamp: DateTime.utc_now()
}
end
defp build_event(leader_name, worker_id, :started) do
%{
worker_id: worker_id,
type: :started,
data: %{},
leader: leader_name,
timestamp: DateTime.utc_now()
}
end
defp build_event(leader_name, worker_id, :finished) do
%{
worker_id: worker_id,
type: :finished,
data: %{},
leader: leader_name,
timestamp: DateTime.utc_now()
}
end
defp update_state_with_worker_msg(
%LeaderState{workers: workers} = state,
worker_id,
{:progress, task_index, total, percent}
) do
workers =
Map.update(
workers,
worker_id,
WorkerState.new(worker_id) |> WorkerState.running(percent, total, task_index),
fn worker ->
WorkerState.running(worker, percent, total, task_index)
end
)
%{state | workers: workers}
end
defp update_state_with_worker_msg(
%LeaderState{workers: workers} = state,
worker_id,
{:result, task_index, result}
) do
workers =
Map.update(
workers,
worker_id,
WorkerState.new(worker_id) |> WorkerState.add_result(task_index, result),
fn worker ->
WorkerState.add_result(worker, task_index, result)
end
)
%{state | workers: workers}
end
defp update_state_with_worker_msg(
%LeaderState{workers: workers} = state,
worker_id,
{:error, task_index, reason}
) do
workers =
Map.update(
workers,
worker_id,
WorkerState.new(worker_id) |> WorkerState.error(reason, task_index),
fn worker ->
WorkerState.error(worker, reason, task_index)
end
)
%{state | workers: workers}
end
defp update_state_with_worker_msg(state, worker_id, :started) do
workers = Map.put_new_lazy(state.workers, worker_id, fn -> WorkerState.new(worker_id) end)
%{state | workers: workers}
end
defp update_state_with_worker_msg(state, worker_id, :finished) do
workers =
Map.update(state.workers, worker_id, WorkerState.new(worker_id), fn worker ->
WorkerState.finished(worker)
end)
%{state | workers: workers}
end
defp maybe_unregister_logger(false, _id), do: :ok
defp maybe_unregister_logger(true, id) do
ParallelLogger.unregister_worker(id)
rescue
e ->
Logger.warning("Failed to unregister worker logger: worker_id=#{inspect(id)}, error=#{inspect(e)}")
:ok
end
end