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 monitor.ex
Raw

lib/parallel/monitor.ex

defmodule Argos.Parallel.Monitor do
@moduledoc """
Monitor de estado agregado de workers paralelos.
Se suscribe al `Leader`, transforma eventos en un mapa de `WorkerState` y
emite actualizaciones consolidadas a sus suscriptores. Ofrece estadísticas.
"""
use GenServer
alias Argos.Parallel.{Leader, MonitorState, WorkerState}
@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 """
Obtiene el estado completo del monitor.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_state(GenServer.server(), keyword()) :: MonitorState.t()
def get_state(server \\ __MODULE__, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, :get_state, timeout)
end
@doc """
Obtiene el estado de un worker específico.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
## Retorna
`WorkerState.t()` si el worker existe, `nil` si no existe.
"""
@spec get_worker_state(term(), GenServer.server(), keyword()) :: WorkerState.t() | nil
def get_worker_state(worker_id, server \\ __MODULE__, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, {:get_worker_state, worker_id}, timeout)
end
@doc """
Suscribe el proceso actual para recibir actualizaciones consolidadas del monitor.
## 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 del monitor.
## 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 estadísticas agregadas de todos los workers.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_stats(GenServer.server(), keyword()) :: map()
def get_stats(server \\ __MODULE__, opts \\ []) do
timeout = Keyword.get(opts, :timeout, @default_call_timeout)
GenServer.call(server, :get_stats, timeout)
end
@impl true
def init(opts) do
leader = Keyword.get(opts, :leader, Argos.Parallel.Leader)
state = %MonitorState{leader: leader}
Leader.subscribe(leader)
{:ok, state}
end
@impl true
def handle_call(:get_state, _from, state) do
{:reply, state, state}
end
@impl true
def handle_call({:get_worker_state, worker_id}, _from, state) do
worker_state = Map.get(state.workers, worker_id)
{:reply, worker_state, state}
end
@impl true
def handle_call({:subscribe, pid}, _from, %MonitorState{subscribers: subs} = state) do
Process.monitor(pid)
{:reply, :ok, %{state | subscribers: MapSet.put(subs, pid)}}
end
@impl true
def handle_call({:unsubscribe, pid}, _from, %MonitorState{subscribers: subs} = state) do
{:reply, :ok, %{state | subscribers: MapSet.delete(subs, pid)}}
end
@impl true
def handle_call(:get_stats, _from, state) do
stats = calculate_stats(state.workers)
{:reply, stats, state}
end
@impl true
def handle_info({:parallel_event, _leader, event}, state) do
state = update_state_from_event(state, event)
notify_subscribers(state)
{:noreply, state}
end
@impl true
def handle_info({:DOWN, _ref, :process, pid, _reason}, %MonitorState{subscribers: subs} = state) do
{:noreply, %{state | subscribers: MapSet.delete(subs, pid)}}
end
@impl true
def handle_info(_msg, state) do
{:noreply, state}
end
defp update_state_from_event(state, %{worker_id: worker_id, type: type, data: data}) do
current_workers = state.workers
updated_workers =
case type do
:started ->
Map.put(current_workers, worker_id, WorkerState.new(worker_id))
:progress ->
Map.update(
current_workers,
worker_id,
WorkerState.new(worker_id)
|> WorkerState.running(data.percent, data.total, data.task_index),
fn worker ->
WorkerState.running(worker, data.percent, data.total, data.task_index)
end
)
:result ->
Map.update(
current_workers,
worker_id,
WorkerState.new(worker_id),
fn worker ->
WorkerState.add_result(worker, data.task_index, data.result)
end
)
:finished ->
Map.update(
current_workers,
worker_id,
WorkerState.new(worker_id) |> WorkerState.finished(),
fn worker ->
WorkerState.finished(worker)
end
)
:error ->
Map.update(
current_workers,
worker_id,
WorkerState.new(worker_id) |> WorkerState.error(data.reason, data.task_index),
fn worker ->
WorkerState.error(worker, data.reason, data.task_index)
end
)
_ ->
current_workers
end
%{state | workers: updated_workers, last_update: DateTime.utc_now()}
end
defp notify_subscribers(%MonitorState{subscribers: subs} = state) do
Enum.each(subs, fn pid ->
send(pid, {:parallel_monitor_update, state})
end)
end
defp calculate_stats(workers) do
total_workers = map_size(workers)
status_counts =
Enum.reduce(workers, %{}, fn {_, worker}, acc ->
status = worker.status
Map.update(acc, status, 1, &(&1 + 1))
end)
total_tasks =
Enum.reduce(workers, 0, fn {_, worker}, acc ->
acc + (worker.total || 0)
end)
completed_tasks =
Enum.reduce(workers, 0, fn {_, worker}, acc ->
if worker.current_task do
acc + worker.current_task
else
acc
end
end)
%{
total_workers: total_workers,
status_counts: status_counts,
total_tasks: total_tasks,
completed_tasks: completed_tasks,
progress_percentage: if(total_tasks > 0, do: (completed_tasks / total_tasks * 100) |> Float.round(1), else: 0)
}
end
end