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