Current section
Files
Jump to
Current section
Files
lib/parallel/worker.ex
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