He realizado unos cambios en el "example" de aegis, pero creo que realmente debreria implementarse dentro de TaskRunner para que sea transparaente. Si analizamos el "Example" de argos: 1º Start: - Primero se arranca el sistema de "ArgosParallel - Despues te suscribes al leader o al monitor, segun quieras estar al corriente. > Leader te llegan los eventos sueltos y en el momento, > monitor te llega el estado completo de la ejecucion con el estado actual de todos los procesos sueltos. > Puedes suscribirte a ambos a la vez sin problema. - A continuacion, se preparan la lista de workers. El worker se genera con Argos.Parallel.create_worker_spec() y se le pasan 2 parmetros: > El priumero, un atomo unico para diferenciar a cada worker > Lista de "task" de cada worker. · Cada elemento del task son funciones anonimas que van realizando las diferentes tareas que pertenencen a ese worker. · Al final de cada tarea se devuelve el resultado de ese task, para comunicar el estado de cada tarea. · Lo que devuelven los task puede ser lo que sea, pero debe ser un formato el cual sirva para ser procesado por la funcion que procesa las respuestas. - Una vez preparada, la lista de workers se manda a "Argos.Parallel.start_workers()". Aqui empezará la ejecucion en paralelo. - Para poder escuchar los eventos de cada worker, hay que crear una funcion como la que tiene el ejemplo de argos, llamada "listen()": def listen do receive do # Si nos hemos suscrito a "Argos.Parallel.subscribe_leader()" # habrá que hacer pattern matchin por esta tupla, que son los eventos del leader {:parallel_event, _leader, event} -> # La estructura del event es la siguiente: # %{ # data: %{}, # timestamp: ~U[2025-10-29 15:32:34.010701Z], # type: :finished, # leader: Argos.Parallel.Leader, # worker_id: :io_simulator # } # segun la respuesta, se puede realizar una tarea u otra. Aqui seria donde se realizaria el renderizado diferencial, refrescando la fila que pertenece al worker 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}") _ -> nil end #Se vuelve a lanzar el listener para recibir el siguiente mensaje de la cola listen() # Si nos hemos suscrito a "Argos.Parallel.subscribe_monitor()" # habrá que hacer pattern matchin por esta tupla, que son los eventos del monitor {:parallel_monitor_update, monitor_state} -> # La estructura del monitor_state es la siguiente: # monitor_state =>: %Argos.Parallel.MonitorState{ # leader: Argos.Parallel.Leader, # workers: %{ # mapa con el estado actual de cada worker. # data_processor: %Argos.Parallel.WorkerState{ # id: :data_processor, # status: :running, # progress: 66.66666666666667, # total: 3, # current_task: 2, # started_at: ~U[2025-10-29 15:35:32.858587Z], # finished_at: nil, # error: nil, # error_task: nil, # last_update: ~U[2025-10-29 15:35:35.848042Z], # results: [ # aqui esta el historico de los mensajes enviados por este worker # %{ # timestamp: ~U[2025-10-29 15:35:35.848042Z], # result: {:ok, :processed_data_2, # %{size: 200, timestamp: ~U[2025-10-29 15:35:35.847962Z]}}, # task_index: 2 # }, # %{ # timestamp: ~U[2025-10-29 15:35:33.847381Z], # result: {:ok, :processed_data_1, # %{size: 100, timestamp: ~U[2025-10-29 15:35:33.847178Z]}}, # task_index: 1 # } # ] # }, # calculator: %Argos.Parallel.WorkerState{ # id: :calculator, # status: :finished, # progress: 100, # total: 3, # current_task: 3, # started_at: ~U[2025-10-29 15:35:32.858598Z], # finished_at: ~U[2025-10-29 15:35:35.366297Z], # error: nil, # error_task: nil, # last_update: ~U[2025-10-29 15:35:35.366297Z], # results: [ # %{ # timestamp: ~U[2025-10-29 15:35:35.366291Z], # result: {:power, 65536}, # task_index: 3 # }, # %{ # timestamp: ~U[2025-10-29 15:35:34.865216Z], # result: {:factorial, # 93326215443944152681699238856266700490715968264381621468592963895217599993229915608941463976156518286253697920827223758251185210916864000000000000000000000000}, # task_index: 2 # }, # %{ # timestamp: ~U[2025-10-29 15:35:33.664337Z], # result: {:sum, 500500}, # task_index: 1 # } # ] # }, # io_simulator: %Argos.Parallel.WorkerState{ # id: :io_simulator, # status: :running, # progress: 33.333333333333336, # total: 3, # current_task: 1, # started_at: ~U[2025-10-29 15:35:32.858601Z], # finished_at: nil, # error: nil, # error_task: nil, # last_update: ~U[2025-10-29 15:35:35.859057Z], # results: [ # %{ # timestamp: ~U[2025-10-29 15:35:35.859057Z], # result: {:read, "file1.txt", "Content of file 1"}, # task_index: 1 # } # ] # }, # fast_worker: %Argos.Parallel.WorkerState{ # id: :fast_worker, # status: :finished, # progress: 100, # total: 5, # current_task: 5, # started_at: ~U[2025-10-29 15:35:32.858607Z], # finished_at: ~U[2025-10-29 15:35:33.463210Z], # error: nil, # error_task: nil, # last_update: ~U[2025-10-29 15:35:33.463210Z], # results: [ # %{ # timestamp: ~U[2025-10-29 15:35:33.463199Z], # result: :quick_task_5, # task_index: 5 # }, # %{ # timestamp: ~U[2025-10-29 15:35:33.412127Z], # result: :quick_task_4, # task_index: 4 # }, # %{ # timestamp: ~U[2025-10-29 15:35:33.311107Z], # result: :quick_task_3, # task_index: 3 # }, # %{ # timestamp: ~U[2025-10-29 15:35:33.110156Z], # result: :quick_task_2, # task_index: 2 # }, # %{ # timestamp: ~U[2025-10-29 15:35:32.959244Z], # result: :quick_task_1, # task_index: 1 # } # ] # } # }, # subscribers: MapSet.new([#PID<0.209.0>]), # last_update: ~U[2025-10-29 15:35:35.859057Z] # } # stats => %{ # esto es el estado global de la ejecucion completa # progress_percentage: 85.7, # total_workers: 4, # status_counts: %{running: 1, finished: 3}, # completed_tasks: 12, # total_tasks: 14 # } # para ver el estado de cada worker monitor_state.workers |> Enum.each(fn {worker_id, worker} -> status_icon = case worker.status do :running -> "🔄" :finished -> "✅" :error -> "❌" :started -> "🚀" _ -> "⏸️" end # para ver el tiempo que lleva ejecutandose un worker elapsed = Argos.Parallel.WorkerState.elapsed_time(worker) end) listen() after # aqui se se configura el timeout de la ejecucion 30_000 -> IO.puts("\n🎉 Proceso terminado (timeout)") # con esto puedes pedir el estado del monitor de ese momento final_state = Argos.Parallel.get_monitor_state() final_state.workers |> Enum.each(fn {worker_id, worker} -> # para procesar los workers como se quiera end) # Detener el sistema Argos.Parallel.stop_system() end end