defmodule Argos.Parallel.Worker do @moduledoc """ GenServer que ejecuta una lista de tareas en secuencia y notifica al LĂ­der. Cada tarea puede ser `fun/0` o `{mod, fun, args}`. Emite mensajes de `:started`, `{:progress, ...}`, `{:result, ...}` y `:error`. """ use GenServer require Logger @type task_spec :: (-> any()) | {module(), atom(), [any()]} @start_delay 0 @spec start_link(%{leader: pid, id: any, tasks: list}) :: GenServer.on_start() def start_link(%{leader: leader, id: id, tasks: tasks}) when is_list(tasks) do GenServer.start_link(__MODULE__, %{leader: leader, id: id, tasks: tasks}) end @spec start_link(term) :: {:error, :invalid_args} def start_link(_other) do {:error, :invalid_args} end @impl true def init(%{leader: leader, id: id, tasks: tasks}) do Process.send_after(self(), :start_work, @start_delay) {:ok, %{leader: leader, id: id, tasks: tasks, total: length(tasks)}} end @impl true def handle_info(:start_work, %{leader: leader, id: id} = state) do safe_cast(leader, {:worker_msg, id, :started}) state = Enum.with_index(state.tasks, 1) |> Enum.reduce_while(state, fn {task, idx}, acc -> if Map.get(acc, :errored) do {:halt, acc} else total = acc.total try do result = execute_task(task) percent = percent(idx, total) safe_cast(leader, {:worker_msg, id, {:progress, idx, total, percent}}) safe_cast(leader, {:worker_msg, id, {:result, idx, result}}) new_acc = Map.put(acc, :last_result, result) |> Map.put(:current_task_index, idx) {:cont, new_acc} rescue e -> reason = {e, __STACKTRACE__} safe_cast(leader, {:worker_msg, id, {:error, idx, reason}}) errored_acc = Map.put(acc, :last_error, reason) |> Map.put(:errored, true) |> Map.put(:current_task_index, idx) {:halt, errored_acc} catch kind, value -> reason = {kind, value} safe_cast(leader, {:worker_msg, id, {:error, idx, reason}}) errored_acc = Map.put(acc, :last_error, reason) |> Map.put(:errored, true) |> Map.put(:current_task_index, idx) {:halt, errored_acc} end end end) if Map.get(state, :errored) do :telemetry.execute( [:argos, :parallel, :worker, :finished_with_error], %{tasks_completed: Map.get(state, :current_task_index, 0)}, %{worker_id: id, error: Map.get(state, :last_error)} ) else safe_cast(leader, {:worker_msg, id, :finished}) :telemetry.execute( [:argos, :parallel, :worker, :finished], %{tasks_completed: state.total}, %{worker_id: id} ) end {:stop, :normal, state} end @doc """ Ejecuta una tarea individual. """ @spec execute_task((-> any) | {module, atom, [any]}) :: any def execute_task(fun) when is_function(fun, 0) do fun.() end def execute_task({mod, fun, args}) when is_atom(mod) and is_atom(fun) and is_list(args) do apply(mod, fun, args) end def execute_task(other) do raise ArgumentError, "Invalid task spec: #{inspect(other)}" end defp percent(idx, total) when total > 0 do idx * 100 / total end defp percent(_idx, _total), do: 0 defp safe_cast(leader, msg) do GenServer.cast(leader, msg) end end