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_.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