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

lib/parallel/logger.ex

defmodule Argos.Parallel.Logger do
@moduledoc """
Logger personalizado para eventos de Argos.Parallel.
Gestiona logs individuales por worker. Cada worker puede tener su propio
archivo de log si se especifica `log: true` en su especificación.
Los logs se guardan en `~/.argos/logs/parallel_workers/` con el formato:
`worker_<worker_id>.jsonl`
Todos los logs se eliminan automáticamente cuando se detiene el sistema
paralelo con `Argos.Parallel.stop_system()`.
## Uso
Para habilitar logging en un worker:
Argos.Parallel.create_worker_spec(:my_worker, tasks, log: true)
## Formato de Eventos
Los eventos se escriben en formato JSON Lines (una línea JSON por evento):
```json
{"type":"result","worker_id":"worker1","timestamp":"2025-01-15T10:30:00Z",...}
{"type":"progress","worker_id":"worker1","timestamp":"2025-01-15T10:30:01Z",...}
```
## Acceso a Logs Durante la Ejecución
Puedes acceder a los logs mientras la ejecución está en proceso de varias formas:
### Desde Código Elixir
# Obtener la ruta del log de un worker
log_path = Argos.Parallel.Logger.get_log_path(:data_processor)
# Leer todo el log de un worker
events = Argos.Parallel.Logger.read_log(:data_processor)
# Leer las últimas 20 líneas
recent_events = Argos.Parallel.Logger.read_log_tail(:data_processor, 20)
# Listar todos los workers con logs activos
active_logs = Argos.Parallel.Logger.list_active_logs()
### Desde la Terminal
# Ver log en tiempo real (tail -f)
tail -f ~/.argos/logs/parallel_workers/worker_data_processor.jsonl
# Ver todo el log
cat ~/.argos/logs/parallel_workers/worker_data_processor.jsonl
# Ver últimas 20 líneas
tail -n 20 ~/.argos/logs/parallel_workers/worker_data_processor.jsonl
# Ver logs formateados con jq
tail -f ~/.argos/logs/parallel_workers/worker_data_processor.jsonl | jq .
"""
use GenServer
require Logger
alias Argos.Parallel.{Leader, Monitor}
defp log_base_dir do
System.user_home!()
|> Path.join(".argos")
|> Path.join("logs")
|> Path.join("parallel_workers")
end
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@doc """
Registra un worker para logging.
"""
@spec register_worker(any, boolean) :: :ok
def register_worker(worker_id, enabled) do
GenServer.call(__MODULE__, {:register_worker, worker_id, enabled})
end
@doc """
Desregistra un worker y cierra su archivo de log.
"""
@spec unregister_worker(any) :: :ok
def unregister_worker(worker_id) do
GenServer.call(__MODULE__, {:unregister_worker, worker_id})
end
@doc """
Limpia todos los logs de workers.
"""
@spec cleanup_logs() :: :ok
def cleanup_logs do
GenServer.call(__MODULE__, :cleanup_logs)
end
@doc """
Obtiene la ruta del archivo de log de un worker.
Retorna `nil` si el worker no tiene logging habilitado o no existe.
"""
@spec get_log_path(any) :: String.t() | nil
def get_log_path(worker_id) do
base_dir = log_base_dir()
log_file = worker_log_file(base_dir, worker_id)
if File.exists?(log_file) do
log_file
else
nil
end
end
@doc """
Lee el contenido completo del log de un worker.
Retorna una lista de mapas (eventos parseados) o `nil` si el log no existe.
"""
@spec read_log(any) :: [map()] | nil
def read_log(worker_id) do
case get_log_path(worker_id) do
nil ->
nil
log_file ->
try do
log_file
|> File.read!()
|> String.split("\n", trim: true)
|> Enum.map(&Jason.decode!/1)
rescue
_ -> []
end
end
end
@doc """
Lee las últimas N líneas del log de un worker.
Retorna una lista de mapas (eventos parseados) o `nil` si el log no existe.
"""
@spec read_log_tail(any, non_neg_integer()) :: [map()] | nil
def read_log_tail(worker_id, lines \\ 20) do
case get_log_path(worker_id) do
nil ->
nil
log_file ->
try do
log_file
|> File.read!()
|> String.split("\n", trim: true)
|> Enum.take(-lines)
|> Enum.map(&Jason.decode!/1)
rescue
_ -> []
end
end
end
@doc """
Lista todos los workers que tienen logs activos.
Retorna una lista de tuplas `{worker_id, log_file_path}`.
"""
@spec list_active_logs() :: [{any, String.t()}]
def list_active_logs do
base_dir = log_base_dir()
if File.exists?(base_dir) do
base_dir
|> File.ls!()
|> Enum.filter(&String.ends_with?(&1, ".jsonl"))
|> Enum.map(fn filename ->
worker_id =
filename
|> String.replace("worker_", "")
|> String.replace(".jsonl", "")
log_file = Path.join(base_dir, filename)
{worker_id, log_file}
end)
else
[]
end
end
@impl true
def init(opts) do
base_dir = log_base_dir()
File.mkdir_p!(base_dir)
leader = Keyword.get(opts, :leader, Leader)
monitor = Keyword.get(opts, :monitor, Monitor)
try do
:ok = Leader.subscribe(leader)
:ok = Monitor.subscribe(monitor)
rescue
e ->
Logger.warning("Argos.Parallel.Logger: Error al suscribirse: #{inspect(e)}")
end
state = %{
log_base_dir: base_dir,
leader: leader,
monitor: monitor,
worker_logs: %{}
}
{:ok, state}
end
@impl true
def handle_call({:register_worker, worker_id, true}, _from, state) do
log_file = worker_log_file(state.log_base_dir, worker_id)
File.write!(log_file, "")
new_worker_logs = Map.put(state.worker_logs, worker_id, log_file)
new_state = %{state | worker_logs: new_worker_logs}
{:reply, :ok, new_state}
end
def handle_call({:register_worker, _worker_id, false}, _from, state) do
{:reply, :ok, state}
end
@impl true
def handle_call({:unregister_worker, worker_id}, _from, state) do
case Map.get(state.worker_logs, worker_id) do
nil ->
{:reply, :ok, state}
log_file ->
try do
File.rm(log_file)
rescue
e ->
Logger.warning("Argos.Parallel.Logger: Error eliminando log de worker #{inspect(worker_id)}: #{inspect(e)}")
end
new_worker_logs = Map.delete(state.worker_logs, worker_id)
new_state = %{state | worker_logs: new_worker_logs}
{:reply, :ok, new_state}
end
end
@impl true
def handle_call(:cleanup_logs, _from, state) do
Enum.each(state.worker_logs, fn {worker_id, log_file} ->
try do
File.rm(log_file)
rescue
e ->
Logger.warning("Argos.Parallel.Logger: Error eliminando log de worker #{inspect(worker_id)}: #{inspect(e)}")
end
end)
try do
File.rmdir(state.log_base_dir)
rescue
_ -> :ok
end
new_state = %{state | worker_logs: %{}}
{:reply, :ok, new_state}
end
@impl true
def handle_info({:parallel_event, _leader, %{worker_id: worker_id} = event}, state) do
case Map.get(state.worker_logs, worker_id) do
nil ->
{:noreply, state}
log_file ->
write_event(log_file, event)
{:noreply, state}
end
end
@impl true
def handle_info({:parallel_monitor_update, _monitor_state}, state) do
{:noreply, state}
end
@impl true
def handle_info(_msg, state) do
{:noreply, state}
end
defp worker_log_file(base_dir, worker_id) do
safe_id = worker_id |> to_string() |> String.replace(~r/[^a-zA-Z0-9_-]/, "_")
Path.join(base_dir, "worker_#{safe_id}.jsonl")
end
defp write_event(log_file, event) do
json = Jason.encode!(event)
File.write!(log_file, json <> "\n", [:append])
rescue
e ->
Logger.error("Argos.Parallel.Logger: Error escribiendo evento: #{inspect(e)}")
end
end