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