defmodule Argos.Parallel.Example do @moduledoc """ Ejemplo completo de uso de `Argos.Parallel` para ejecución paralela de tareas. Este módulo sirve como referencia y guía para implementar procesamiento paralelo usando `Argos.Parallel`. Demuestra cómo: - Iniciar el sistema paralelo - Crear workers con múltiples tareas - Suscribirse a eventos del Leader (eventos individuales) - Suscribirse a eventos del Monitor (estado consolidado) - Procesar eventos en tiempo real - Obtener estadísticas y resultados finales - Detener el sistema correctamente ## Tipos de Suscripción `Argos.Parallel` ofrece dos tipos de suscripción: ### Leader (Eventos Individuales) - Recibe eventos individuales de cada worker en tiempo real - Útil para visualización detallada o procesamiento específico por evento - Mensaje: `{:parallel_event, leader, event}` - Eventos: `:result`, `:progress`, `:finished`, `:error` ### Monitor (Estado Consolidado) - Recibe actualizaciones del estado consolidado de todos los workers - Útil para visualización general o dashboards - Mensaje: `{:parallel_monitor_update, monitor_state}` - Contiene el estado completo de todos los workers ## Ejecución Para ejecutar este ejemplo: mix run -e "Argos.Parallel.Example.start()" ## Estructura del Ejemplo El ejemplo crea 4 workers con diferentes características: 1. **data_processor**: Procesamiento de datos con diferentes duraciones (con logging habilitado) 2. **calculator**: Cálculos matemáticos (suma, factorial, potencia) (con logging habilitado) 3. **io_simulator**: Simulación de operaciones de E/S con posibilidad de error (sin logging) 4. **fast_worker**: Operaciones rápidas múltiples (sin logging) Los workers con logging habilitado generan archivos de log individuales en `~/.argos/logs/parallel_workers/worker_.jsonl` que se eliminan automáticamente cuando se detiene el sistema. ## Flujo de Trabajo 1. Iniciar el sistema con `Argos.Parallel.start_system()` 2. Suscribirse a los canales deseados (`subscribe_leader()` y/o `subscribe_monitor()`) 3. Crear especificaciones de workers con `create_worker_spec/3` (puede incluir `log: true`) 4. Iniciar workers con `start_workers/1` 5. Escuchar eventos en un loop con `receive` 6. Procesar eventos según el tipo 7. Detener el sistema con `stop_system()` cuando termine (los logs se eliminan automáticamente) ## Ejemplo de Worker ### Worker sin logging worker_spec = Argos.Parallel.create_worker_spec(:my_worker, [ fn -> Process.sleep(1000) {:ok, :result_1} end, fn -> Process.sleep(2000) {:ok, :result_2} end ]) ### Worker con logging habilitado worker_spec_with_log = Argos.Parallel.create_worker_spec(:my_worker, [ fn -> Process.sleep(1000) {:ok, :result_1} end, fn -> Process.sleep(2000) {:ok, :result_2} end ], log: true) Los workers con `log: true` generan un archivo de log individual en formato JSON Lines en `~/.argos/logs/parallel_workers/worker_.jsonl`. Todos los logs se eliminan automáticamente cuando se llama a `Argos.Parallel.stop_system()`. ## Procesamiento de Eventos ### Eventos del Leader receive do {:parallel_event, _leader, event} -> case event do %{type: :result, worker_id: worker_id, data: data} -> IO.puts("Worker " <> to_string(worker_id) <> " completó tarea " <> Integer.to_string(data.task_index)) %{type: :progress, worker_id: worker_id, data: data} -> IO.puts("Worker " <> to_string(worker_id) <> " progreso: " <> Float.to_string(data.percent) <> "%") %{type: :finished, worker_id: worker_id} -> IO.puts("Worker " <> to_string(worker_id) <> " terminó") %{type: :error, worker_id: worker_id, data: data} -> IO.puts("Worker " <> to_string(worker_id) <> " error: " <> to_string(data.reason)) end end ### Eventos del Monitor receive do {:parallel_monitor_update, monitor_state} -> stats = Argos.Parallel.get_stats() IO.puts("Progreso general: " <> Float.to_string(stats.progress_percentage) <> "%") monitor_state.workers |> Enum.each(fn {id, worker} -> IO.puts(to_string(id) <> ": " <> to_string(worker.status) <> " - " <> Integer.to_string(worker.progress) <> "%") end) end ## Manejo de Errores Cuando un worker encuentra un error: - El worker se detiene inmediatamente - Se envía un evento `:error` con los detalles - Las tareas restantes del worker no se ejecutan - Otros workers continúan normalmente ## Estadísticas Obtener estadísticas generales: stats = Argos.Parallel.get_stats() ## Sistema de Logging `Argos.Parallel` incluye un sistema de logging automático que permite generar archivos de log individuales por worker. Para habilitar el logging en un worker, usa la opción `log: true` al crear la especificación: Argos.Parallel.create_worker_spec(:worker_id, tasks, log: true) Los logs se guardan en formato JSON Lines (una línea JSON por evento) en: `~/.argos/logs/parallel_workers/worker_.jsonl` Todos los logs se eliminan automáticamente cuando se detiene el sistema con `Argos.Parallel.stop_system()`. ## Ver También - `Argos.Parallel`: API principal del sistema paralelo - `Argos.Parallel.Leader`: Gestión de eventos individuales - `Argos.Parallel.Monitor`: Gestión de estado consolidado - `Argos.Parallel.WorkerState`: Estado de un worker individual - `Argos.Parallel.Logger`: Sistema de logging automático """ @doc """ Inicia el ejemplo completo de procesamiento paralelo. Este es el punto de entrada principal que demuestra todo el flujo de trabajo. """ alias Argos.Parallel.WorkerState @spec start() :: :ok def start do case Argos.Parallel.start_system() do {:ok, _pid} -> IO.puts("✅ Sistema paralelo iniciado") {:error, {:already_started, _pid}} -> IO.puts("✅ Sistema paralelo ya estaba ejecutándose") {:error, reason} -> IO.puts("❌ Error iniciando sistema: #{inspect(reason)}") end :ok = Argos.Parallel.subscribe_leader() :ok = Argos.Parallel.subscribe_monitor() IO.puts("✅ Suscrito a Leader y Monitor") worker_specs = create_sample_workers() IO.puts("🚀 Iniciando workers...") results = Argos.Parallel.start_workers(worker_specs) IO.puts("✅ Workers iniciados: #{length(results)}") listen() end @spec create_sample_workers() :: [%{id: any, tasks: list, log: boolean}] defp create_sample_workers do [ Argos.Parallel.create_worker_spec( :data_processor, [ fn -> Process.sleep(1000) {:ok, :processed_data_1, %{size: 100, timestamp: DateTime.utc_now()}} end, fn -> Process.sleep(2000) {:ok, :processed_data_2, %{size: 200, timestamp: DateTime.utc_now()}} end, fn -> Process.sleep(1500) {:ok, :processed_data_3, %{size: 150, timestamp: DateTime.utc_now()}} end ], log: true ), Argos.Parallel.create_worker_spec( :calculator, [ fn -> Process.sleep(800) result = Enum.sum(1..1000) {:sum, result} end, fn -> Process.sleep(1200) result = Enum.reduce(1..100, 1, &(&1 * &2)) {:factorial, result} end, fn -> Process.sleep(500) result = :math.pow(2, 16) |> round {:power, result} end ], log: true ), Argos.Parallel.create_worker_spec(:io_simulator, [ fn -> Process.sleep(3000) {:read, "file1.txt", "Content of file 1"} end, fn -> Process.sleep(2500) {:write, "file2.txt", 2048} end, fn -> Process.sleep(1800) if :rand.uniform(10) == 1 do raise "Simulated IO error: Device not ready" else {:read, "file3.txt", "Content of file 3"} end end ]), Argos.Parallel.create_worker_spec(:fast_worker, [ fn -> Process.sleep(100) :quick_task_1 end, fn -> Process.sleep(150) :quick_task_2 end, fn -> Process.sleep(200) :quick_task_3 end, fn -> Process.sleep(100) :quick_task_4 end, fn -> Process.sleep(50) :quick_task_5 end ]) ] end @doc """ Loop principal de escucha de eventos. Procesa eventos del Leader (individuales) y del Monitor (consolidados). Termina después de 30 segundos o cuando todos los workers terminan. """ @spec listen() :: :ok def listen do receive do {:parallel_event, _leader, event} -> handle_parallel_event(event) listen() {:parallel_monitor_update, monitor_state} -> handle_monitor_update(monitor_state) listen() after 30_000 -> IO.puts("\n🎉 Proceso terminado (timeout)") show_final_results() Argos.Parallel.stop_system() IO.puts("\n🛑 Sistema paralelo detenido") IO.puts("📝 Logs de workers eliminados automáticamente") end end defp handle_parallel_event(event) do System.cmd("clear", []) |> elem(0) |> IO.puts() IO.puts("") IO.puts("\n--- EVENTO INDIVIDUAL ---") IO.puts("Evento: #{format_event_type(event)}") case event do %{type: :result, worker_id: worker_id, data: data} -> IO.puts("✅ Worker #{inspect(worker_id)} - Tarea #{data.task_index} completada: #{inspect(data.result)}") %{type: :progress, worker_id: worker_id, data: data} -> IO.puts("📊 Worker #{inspect(worker_id)} - Progreso: #{data.percent}% (tarea #{data.task_index}/#{data.total})") %{type: :finished, worker_id: worker_id} -> IO.puts("🎉 Worker #{inspect(worker_id)} - COMPLETADO") %{type: :error, worker_id: worker_id, data: data} -> IO.puts("❌ Worker #{inspect(worker_id)} - ERROR en tarea #{data.task_index}: #{data.reason}") _ -> :ok end end defp handle_monitor_update(monitor_state) do System.cmd("clear", []) |> elem(0) |> IO.puts() IO.puts("\n--- ESTADO CONSOLIDADO ---") stats = Argos.Parallel.get_stats() IO.puts("📈 Progreso general: #{stats.progress_percentage}%") IO.puts("👥 Workers totales: #{stats.total_workers}") IO.puts("Estados workers: #{format_status_counts(stats.status_counts)}") monitor_state.workers |> Enum.each(&print_worker_status/1) end defp print_worker_status({worker_id, worker}) do status_icon = get_status_icon(worker.status) elapsed = WorkerState.elapsed_time(worker) IO.puts("#{status_icon} #{inspect(worker_id)}: #{worker.status} - #{worker.progress || 0}% (#{elapsed}ms)") end defp get_status_icon(:running), do: "🔄" defp get_status_icon(:finished), do: "✅" defp get_status_icon(:error), do: "❌" defp get_status_icon(:started), do: "🚀" defp get_status_icon(_), do: "⏸️" defp format_event_type(%{type: type}), do: "#{type}" defp format_event_type(_), do: "unknown" defp format_status_counts(counts) do Enum.map_join(counts, ", ", fn {status, count} -> "#{status}: #{count}" end) end defp show_final_results do IO.puts("\n📋 RESULTADOS FINALES:") final_state = Argos.Parallel.get_monitor_state() final_state.workers |> Enum.each(&print_worker_result/1) end defp print_worker_result({worker_id, worker}) do IO.puts("\n=== Worker: #{inspect(worker_id)} ===") IO.puts("Estado: #{worker.status}") IO.puts("Tareas completadas: #{length(worker.results)}/#{worker.total}") IO.puts("Tiempo total: #{WorkerState.elapsed_time(worker)}ms") print_worker_error(worker) print_worker_results(worker) end defp print_worker_error(%{error: nil}), do: :ok defp print_worker_error(%{error: error}) do IO.puts("❌ Error: #{inspect(error)}") end defp print_worker_results(%{results: []}), do: :ok defp print_worker_results(%{results: results}) do IO.puts("Resultados:") results |> Enum.reverse() |> Enum.each(fn result -> IO.puts(" 📝 Tarea #{result.task_index}: #{inspect(result.result)}") end) end end