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