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 instruccions.txt
Raw

lib/instruccions.txt

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