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

lib/parallel.ex

defmodule Argos.Parallel do
@moduledoc """
API pública del sistema de ejecución en paralelo.
Ofrece arranque/detención del sistema, suscripción a eventos individuales y
a estado agregado, helpers de creación de workers y consulta de estado.
## Arquitectura
El sistema paralelo está compuesto por:
- `Leader`: Coordina el ciclo de vida de los workers
- `Monitor`: Agrega el estado de todos los workers
- `Supervisor`: Supervisa los procesos del sistema
- `Worker`: Ejecuta tareas secuencialmente
## Características
- Ejecución paralela de workers con tareas secuenciales
- Suscripción a eventos individuales o agregados
- Consulta de estado y estadísticas
- Timeouts configurables en todas las operaciones
- Validación de especificaciones de workers
## Ejemplos
# Iniciar sistema
Argos.Parallel.start_system()
# Crear y ejecutar workers
specs = [Argos.Parallel.create_worker_spec(:w1, [fn -> :ok end])]
Argos.Parallel.start_workers(specs)
# Suscribirse a eventos
Argos.Parallel.subscribe()
"""
alias Argos.Parallel.{Leader, Logger, Monitor, Supervisor}
alias Argos.Validation
@spec subscribe(keyword()) :: :ok
def subscribe(opts \\ []) do
ensure_system_started()
Leader.subscribe(Leader, opts)
end
@spec unsubscribe(keyword()) :: :ok
def unsubscribe(opts \\ []) do
case Process.whereis(Leader) do
nil -> :ok
_ -> Leader.unsubscribe(Leader, opts)
end
end
@spec get_state(keyword()) :: map()
def get_state(opts \\ []) do
ensure_system_started()
Leader.get_state(Leader, opts)
end
@spec stop() :: :ok
def stop do
case Process.whereis(Leader) do
nil -> :ok
_ -> Leader.stop()
end
end
@doc """
Inicia el sistema paralelo completo.
"""
@spec start_system(keyword) :: {:ok, pid} | {:error, term}
def start_system(opts \\ []) do
case Process.whereis(Supervisor) do
nil ->
Supervisor.start_link(opts)
pid ->
if system_running?() do
{:ok, pid}
else
stop_system()
Supervisor.start_link(opts)
end
end
end
@doc """
Verifica si el sistema está ejecutándose.
"""
@spec system_running?() :: boolean
def system_running? do
Process.whereis(Supervisor) != nil &&
Process.whereis(Leader) != nil &&
Process.whereis(Monitor) != nil
end
@doc """
Inicia múltiples workers con las especificaciones dadas.
"""
@type task :: (-> any()) | {module(), atom(), [any()]}
@type worker_spec :: %{id: any(), tasks: [task()], log: boolean()}
@type worker_result :: {:ok, pid()} | {:error, term()}
@spec start_workers([worker_spec()], keyword()) :: [worker_result()]
def start_workers(worker_specs, opts \\ []) when is_list(worker_specs) do
ensure_system_started()
Leader.start_workers(Leader, worker_specs, opts)
end
@doc """
Función de conveniencia para crear especificaciones de worker.
## Opciones
- `:log` (boolean): Si es `true`, el worker tendrá su propio archivo de log.
Por defecto es `false`.
## Ejemplos
Argos.Parallel.create_worker_spec(:worker1, [fn -> :ok end])
Argos.Parallel.create_worker_spec(:worker2, [fn -> :ok end], log: true)
## Validaciones
- `id` no puede ser nil
- `tasks` debe ser una lista no vacía
"""
@spec create_worker_spec(any(), [task()], keyword()) :: worker_spec()
def create_worker_spec(id, tasks, opts \\ []) when is_list(tasks) do
log_enabled = Keyword.get(opts, :log, false)
spec = %{
id: id,
tasks: tasks,
log: log_enabled
}
case Validation.validate_worker_spec(spec) do
:ok ->
spec
{:error, reason} ->
raise ArgumentError, "Invalid worker spec: #{reason}"
end
end
@doc """
Suscribe el proceso actual al líder para recibir eventos individuales.
"""
@spec subscribe_leader(keyword()) :: :ok
def subscribe_leader(opts \\ []) do
ensure_system_started()
Leader.subscribe(Leader, opts)
end
@doc """
Suscribe el proceso actual al monitor para recibir estados consolidados.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec subscribe_monitor(keyword()) :: :ok
def subscribe_monitor(opts \\ []) do
ensure_system_started()
Monitor.subscribe(Monitor, opts)
end
@doc """
Desuscribe el proceso actual del líder.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec unsubscribe_leader(keyword()) :: :ok
def unsubscribe_leader(opts \\ []) do
Leader.unsubscribe(Leader, opts)
end
@doc """
Desuscribe el proceso actual del monitor.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec unsubscribe_monitor(keyword()) :: :ok
def unsubscribe_monitor(opts \\ []) do
Monitor.unsubscribe(Monitor, opts)
end
@doc """
Obtiene el estado interno del líder.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_leader_state(keyword()) :: map
def get_leader_state(opts \\ []) do
Leader.get_state(Leader, opts)
end
@doc """
Obtiene el estado completo del monitor.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_monitor_state(keyword()) :: map
def get_monitor_state(opts \\ []) do
Monitor.get_state(Monitor, opts)
end
@doc """
Obtiene el estado de un worker específico.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_worker_state(any, keyword()) :: map | nil
def get_worker_state(worker_id, opts \\ []) do
Monitor.get_worker_state(worker_id, Monitor, opts)
end
@doc """
Obtiene estadísticas agregadas de todos los workers.
## Opciones
* `:timeout` - Timeout para la llamada GenServer (por defecto: 30 segundos)
"""
@spec get_stats(keyword()) :: map
def get_stats(opts \\ []) do
Monitor.get_stats(Monitor, opts)
end
@doc """
Detiene el sistema paralelo y limpia todos los logs de workers.
"""
@spec stop_system() :: :ok
def stop_system do
try do
if Process.whereis(Logger) != nil do
Logger.cleanup_logs()
end
rescue
_ -> :ok
end
case Process.whereis(Supervisor) do
nil ->
:ok
pid ->
try do
Supervisor.stop(pid)
rescue
_ -> :ok
catch
_, _ -> :ok
else
_ -> :ok
end
end
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(any, keyword()) :: {:ok, :stopped} | {:error, term}
def stop_worker(worker_id, opts \\ []) do
ensure_system_started()
Leader.stop_worker(Leader, worker_id, opts)
end
@system_startup_delay 100
defp ensure_system_started do
unless system_running?() do
start_system()
Process.sleep(@system_startup_delay)
end
end
end