Packages

Graph-first runtime for building agent systems on the BEAM in Elixir

Current section

Files

Jump to
ex_ai lib mix tasks ex_ai.run.events.ex
Raw

lib/mix/tasks/ex_ai.run.events.ex

defmodule Mix.Tasks.ExAi.Run.Events do
@shortdoc "List persisted events for an ExAI run"
@moduledoc """
Loads persisted run events.
mix ex_ai.run.events <run_id>
mix ex_ai.run.events <run_id> --json
mix ex_ai.run.events <run_id> --jsonl
"""
use Mix.Task
alias ExAI.CLI.CommandTelemetry
alias ExAI.CLI.JSON
alias ExAI.CLI.Output
alias ExAI.CLI.Runtime
alias ExAI.Error
alias ExAI.Persistence.History
alias ExAI.Persistence.Store
@switches [dir: :string, json: :boolean, jsonl: :boolean]
@impl true
def run(argv) do
{opts, args, invalid} = OptionParser.parse(argv, strict: @switches)
opts = Map.new(opts)
persistence_opts = Runtime.persistence_opts(opts)
started_ms = CommandTelemetry.start("ex_ai.run.events")
with :ok <- validate_invalid_options(invalid),
:ok <- validate_mode(opts),
{:ok, run_id} <- parse_run_id(args),
:ok <- Runtime.boot(opts),
{:ok, run} <- Store.load_run(run_id, persistence_opts),
{:ok, events} <- History.list_events(run_id, persistence_opts) do
CommandTelemetry.stop("ex_ai.run.events", started_ms, %{run_id: run_id})
emit_success(run, run_id, events, opts)
else
{:error, %Error{} = error} ->
CommandTelemetry.error("ex_ai.run.events", started_ms, error)
emit_error(error, opts, %{})
{:error, reason} ->
error = Error.new(:validation_error, reason)
CommandTelemetry.error("ex_ai.run.events", started_ms, error)
emit_error(error, opts, %{})
end
end
@spec parse_run_id([String.t()]) :: {:ok, String.t()} | {:error, String.t()}
defp parse_run_id([run_id | _]) when is_binary(run_id), do: {:ok, run_id}
defp parse_run_id(_), do: {:error, "usage: mix ex_ai.run.events <run_id>"}
@spec validate_invalid_options([{String.t(), String.t() | nil}]) :: :ok | {:error, String.t()}
defp validate_invalid_options([]), do: :ok
defp validate_invalid_options(invalid) do
{:error, "invalid options: #{inspect(invalid)}"}
end
@spec validate_mode(map()) :: :ok | {:error, String.t()}
defp validate_mode(%{json: true, jsonl: true}),
do: {:error, "choose either --json or --jsonl, not both"}
defp validate_mode(_), do: :ok
@spec emit_success(map(), String.t(), [map()], map()) :: :ok
defp emit_success(run, run_id, events, %{jsonl: true}) do
lines =
events
|> Enum.with_index(1)
|> Enum.map(fn {event, index} ->
JSON.success(%{
run_id: run_id,
status: :event,
events: [event],
metadata: %{
index: index,
event_type: event_type(event),
run_status: Map.get(run, :status)
}
})
end)
Output.emit_jsonl(lines)
end
defp emit_success(run, run_id, events, %{json: true}) do
Output.emit_json(
JSON.success(%{
run_id: run_id,
status: Map.get(run, :status),
events: events,
metadata: Map.get(run, :metadata, %{})
})
)
end
defp emit_success(_run, run_id, events, _opts) do
Output.emit_info("run #{run_id} events=#{length(events)}")
:ok
end
@spec emit_error(ExAI.Error.t(), map(), map()) :: no_return()
defp emit_error(%Error{} = error, opts, payload_overrides) do
Output.halt_error(error, opts, payload_overrides)
end
@spec event_type(map()) :: atom() | String.t()
defp event_type(event) when is_map(event) do
metadata = Map.get(event, :metadata, %{})
Map.get(metadata, :event_type) || Map.get(metadata, "event_type") || :unknown
end
end