defmodule Codex.Runtime.Exec do @moduledoc """ Session-oriented runtime kit for the common Codex exec CLI family. """ @behaviour Codex.RuntimeKit require Logger alias CliSubprocessCore.Event, as: CoreEvent alias CliSubprocessCore.Payload alias CliSubprocessCore.ProcessExit, as: CoreProcessExit alias CliSubprocessCore.ProviderProfiles.Codex, as: CoreCodex alias CliSubprocessCore.Session alias Codex.ApprovalPolicy alias Codex.Config.Overrides alias Codex.Events alias Codex.Exec.Options, as: ExecOptions alias Codex.Files.Attachment alias Codex.GovernedAuthority alias Codex.IO.Buffer alias Codex.Options alias Codex.ProcessExit alias Codex.Runtime.Env, as: RuntimeEnv alias Codex.Runtime.Exec.Profile, as: ExecProfile @default_session_event_tag :codex_sdk_exec_session @session_control_capabilities [ :session_history, :session_resume, :session_pause, :session_intervene ] @known_reason_codes %{ "cli_not_found" => :cli_not_found, "command_not_found" => :command_not_found, "invalid_input" => :invalid_input, "parse_error" => :parse_error, "permission_denied" => :permission_denied, "provider_error" => :provider_error, "transport_error" => :transport_error } @impl true def start_session(opts) when is_list(opts) do with {:ok, session_opts} <- build_session_options(opts) do Session.start_session(session_opts) end end @impl true def subscribe(session, pid, ref) when is_pid(session) and is_pid(pid) and is_reference(ref) do Session.subscribe(session, pid, ref) end @impl true def send_input(session, input, opts \\ []) when is_pid(session) do Session.send_input(session, input, opts) end @impl true def end_input(session) when is_pid(session), do: Session.end_input(session) @impl true def interrupt(session) when is_pid(session), do: Session.interrupt(session) @impl true def close(session) when is_pid(session), do: Session.close(session) @impl true def info(session) when is_pid(session), do: Session.info(session) @impl true def capabilities do (CoreCodex.capabilities() ++ @session_control_capabilities) |> Enum.uniq() end @spec list_provider_sessions(keyword()) :: {:ok, [map()]} | {:error, term()} def list_provider_sessions(opts \\ []) when is_list(opts) do with {:ok, sessions} <- Codex.list_sessions(opts) do {:ok, Enum.map(sessions, fn session -> %{ id: session.id, label: session.originator || session.metadata["title"] || session.metadata[:title] || session.id, cwd: session.cwd, updated_at: session.updated_at, source_kind: :thread_history, metadata: %{ path: session.path, started_at: session.started_at, originator: session.originator, cli_version: session.cli_version }, raw: session } end)} end end @doc false @spec session_event_tag() :: :codex_sdk_exec_session def session_event_tag, do: @default_session_event_tag @impl true def project_event(%CoreEvent{kind: :run_started}, state), do: {[], state} def project_event(%CoreEvent{raw: %{exit: _exit}}, state), do: {[], state} def project_event( %CoreEvent{ kind: :error, payload: %Payload.Error{code: "parse_error", metadata: metadata} }, state ) do line = Map.get(metadata, :line) || Map.get(metadata, "line") || "" Logger.warning("Failed to decode codex event: #{Buffer.format_binary_for_log(line)}") {[], state} end def project_event(%CoreEvent{raw: raw}, state) when is_map(raw) do case decode_public_event(raw, state) do {:ok, event} -> {[event], state} :drop -> {[], state} end end def project_event(_event, state), do: {[], state} @spec session_error(CoreEvent.t(), binary(), boolean()) :: {:error, term()} | nil def session_error( %CoreEvent{ kind: :error, raw: %{exit: exit}, payload: %Payload.Error{} = payload }, stderr, stderr_truncated? ) do if CoreProcessExit.match?(exit) do {:error, Codex.TransportError.new(exit_code(exit), message: payload.message || exit_message(exit), stderr: stderr, stderr_truncated?: stderr_truncated?, retryable?: retryable_exit?(exit), reason_code: normalize_reason_code(payload.code) )} end end def session_error(_event, _stderr, _stderr_truncated?), do: nil @spec stderr_chunk(CoreEvent.t()) :: binary() | nil def stderr_chunk(%CoreEvent{kind: :stderr, payload: %Payload.Stderr{content: content}}) when is_binary(content), do: content def stderr_chunk(_event), do: nil @spec build_session_options(keyword()) :: {:ok, keyword()} | {:error, term()} def build_session_options(opts) when is_list(opts) do exec_opts = Keyword.fetch!(opts, :exec_opts) input = Keyword.get(opts, :input) command_args = Keyword.get(opts, :command_args) subscriber = Keyword.get(opts, :subscriber) with %ExecOptions{} = exec_opts <- exec_opts, {:ok, command_spec} <- Options.codex_command_spec(exec_opts.codex_opts, exec_opts.execution_surface), {:ok, config_values} <- config_values(exec_opts), :ok <- validate_governed_runtime(exec_opts, config_values) do subcommand_args = command_args || command_args_for_run(exec_opts) session_opts = opts |> base_session_options(exec_opts, config_values, subcommand_args) |> Keyword.put(:subscriber, subscriber) |> Keyword.put(:command_spec, command_spec) |> Keyword.put(:stdin, normalize_prompt(input)) {:ok, session_opts} else {:error, _} = error -> error _other -> {:error, :invalid_exec_options} end end @doc false @spec render_for_test(keyword()) :: {:ok, %{ provider: :codex, args: [term()], stdin: String.t() | nil, cwd: term(), env: map(), clear_env?: boolean() | nil, execution_surface: term(), config_values: [term()], provider_native: map() }} | {:error, term()} def render_for_test(opts) when is_list(opts) do exec_opts = Keyword.fetch!(opts, :exec_opts) input = Keyword.get(opts, :input) command_args = Keyword.get(opts, :command_args) with %ExecOptions{} = exec_opts <- exec_opts, {:ok, config_values} <- config_values(exec_opts), :ok <- validate_governed_runtime(exec_opts, config_values) do subcommand_args = command_args || command_args_for_run(exec_opts) session_opts = base_session_options(opts, exec_opts, config_values, subcommand_args) {:ok, %{ provider: :codex, args: ExecProfile.render_args(session_opts), stdin: normalize_prompt(input), cwd: Keyword.get(session_opts, :cwd), env: build_env(exec_opts), clear_env?: exec_opts.clear_env?, execution_surface: exec_opts.execution_surface, config_values: config_values, provider_native: %{ sandbox: fetch_thread_opt(exec_opts.thread, :sandbox), working_directory: fetch_thread_opt(exec_opts.thread, :working_directory), additional_directories: normalize_string_list(fetch_thread_opt(exec_opts.thread, :additional_directories)), skip_git_repo_check: fetch_thread_opt(exec_opts.thread, :skip_git_repo_check) == true, output_schema: exec_opts.output_schema_path, full_auto: exec_opt(exec_opts, :full_auto), dangerously_bypass_approvals_and_sandbox: exec_opt(exec_opts, :dangerously_bypass_approvals_and_sandbox) } }} else {:error, _reason} = error -> error _other -> {:error, :invalid_exec_options} end end defp base_session_options(opts, %ExecOptions{} = exec_opts, config_values, subcommand_args) do [ provider: :codex, profile: Codex.Runtime.Exec.Profile, metadata: session_metadata(opts, exec_opts), cli_profile: exec_opt(exec_opts, :profile), oss: payload_oss?(exec_opts), local_provider: payload_local_provider(exec_opts), full_auto: exec_opt(exec_opts, :full_auto), dangerously_bypass_approvals_and_sandbox: exec_opt(exec_opts, :dangerously_bypass_approvals_and_sandbox), model: normalize_option_string(exec_opts.codex_opts.model), model_payload: exec_opts.codex_opts.model_payload, color: normalize_option_string(exec_opt(exec_opts, :color)), output_last_message: exec_opt(exec_opts, :output_last_message), sandbox: sandbox_mode(fetch_thread_opt(exec_opts.thread, :sandbox)), working_directory: fetch_thread_opt(exec_opts.thread, :working_directory), additional_directories: normalize_string_list(fetch_thread_opt(exec_opts.thread, :additional_directories)), skip_git_repo_check: fetch_thread_opt(exec_opts.thread, :skip_git_repo_check) == true, subcommand_args: subcommand_args, continuation_token: exec_opts.continuation_token, cancellation_token: exec_opts.cancellation_token, images: attachment_paths(exec_opts.attachments), output_schema: exec_opts.output_schema_path, config_values: config_values, env: build_env(exec_opts), clear_env?: exec_opts.clear_env?, session_event_tag: @default_session_event_tag, headless_timeout_ms: :infinity, max_stderr_buffer_size: transport_stderr_buffer_size(exec_opts) ] ++ Options.execution_surface_options(exec_opts.execution_surface) end defp session_metadata(opts, %ExecOptions{codex_opts: %Options{} = codex_opts}) do model = normalize_option_string(codex_opts.model) reasoning_effort = normalize_reasoning_value(codex_opts.reasoning_effort) opts |> Keyword.get(:metadata, %{}) |> normalize_session_metadata() |> Map.put_new(:lane, :codex_sdk) |> put_if_missing("model", model) |> put_reasoning_if_missing(reasoning_effort) |> put_reasoning_config_if_missing(reasoning_effort) end defp normalize_session_metadata(metadata) when is_map(metadata), do: metadata defp normalize_session_metadata(_metadata), do: %{} defp decode_public_event(raw, state) do event = Events.parse!(raw) {:ok, enrich_event(event, codex_options(state))} rescue error in ArgumentError -> Logger.warning("Unsupported codex event: #{Exception.message(error)}") :drop end defp codex_options(%{exec_opts: %ExecOptions{codex_opts: %Options{} = opts}}), do: opts defp codex_options(%{exec_opts: %Options{} = opts}), do: opts defp codex_options(%{codex_opts: %Options{} = opts}), do: opts defp codex_options(_state), do: nil defp enrich_event(%Events.ThreadStarted{} = event, %Options{} = opts) do %Events.ThreadStarted{ event | metadata: enrich_thread_started_metadata(event.metadata, opts) } end defp enrich_event(event, _opts), do: event defp enrich_thread_started_metadata(metadata, %Options{} = opts) do metadata = case metadata do value when is_map(value) -> value _ -> %{} end model = normalize_option_string(opts.model) reasoning_effort = normalize_reasoning_value(opts.reasoning_effort) metadata |> put_if_missing("model", model) |> put_reasoning_if_missing(reasoning_effort) |> put_reasoning_config_if_missing(reasoning_effort) end defp put_if_missing(map, _key, nil), do: map defp put_if_missing(map, key, value), do: if(Map.has_key?(map, key), do: map, else: Map.put(map, key, value)) defp put_reasoning_if_missing(map, nil), do: map defp put_reasoning_if_missing(map, value) do if Map.has_key?(map, "reasoning_effort") or Map.has_key?(map, "reasoningEffort") do map else Map.put(map, "reasoning_effort", value) end end defp put_reasoning_config_if_missing(map, nil), do: map defp put_reasoning_config_if_missing(map, value) do case Map.get(map, "config") do config when is_map(config) -> if Map.has_key?(config, "model_reasoning_effort") do map else Map.put(map, "config", Map.put(config, "model_reasoning_effort", value)) end _ -> Map.put(map, "config", %{"model_reasoning_effort" => value}) end end defp command_args_for_run(%ExecOptions{} = exec_opts) do resume_args(exec_opts.thread) end defp config_values(%ExecOptions{} = exec_opts) do with {:ok, approval_values} <- approval_config_values(exec_opts.thread), {:ok, override_values} <- override_config_values(exec_opts) do {:ok, reasoning_config_values(exec_opts) ++ network_access_config_values(exec_opts.thread) ++ payload_config_values(exec_opts) ++ approval_values ++ override_values} end end defp override_config_values(exec_opts) do with {:ok, global_overrides} <- global_config_overrides(exec_opts), {:ok, thread_overrides} <- thread_config_overrides(exec_opts), {:ok, turn_overrides} <- turn_config_overrides(exec_opts) do derived_overrides = derived_config_overrides(exec_opts) {:ok, (global_overrides ++ derived_overrides ++ thread_overrides ++ turn_overrides) |> Overrides.cli_args() |> config_values_from_cli_args()} end end defp reasoning_config_values(%ExecOptions{codex_opts: %Options{} = opts}) do case normalize_reasoning_value(opts.reasoning_effort) do nil -> [] reasoning -> [~s(model_reasoning_effort="#{reasoning}")] end end defp approval_config_values(%{thread_opts: %Codex.Thread.Options{ask_for_approval: nil}}), do: {:ok, []} defp approval_config_values(%{thread_opts: %Codex.Thread.Options{ask_for_approval: policy}}) do case ApprovalPolicy.to_external(policy) do {:ok, nil} -> {:ok, []} {:ok, value} when is_binary(value) -> {:ok, [~s(approval_policy="#{value}")]} {:ok, %{} = value} -> %{"approval_policy" => value} |> Overrides.normalize_config_overrides() |> case do {:ok, overrides} -> {:ok, overrides |> Overrides.cli_args() |> config_values_from_cli_args()} {:error, reason} -> {:error, reason} end {:error, reason} -> {:error, reason} end end defp approval_config_values(_thread), do: {:ok, []} defp normalize_reasoning_value(nil), do: nil defp normalize_reasoning_value(value) when is_atom(value) do Codex.Models.reasoning_effort_to_string(value) end defp normalize_reasoning_value(value) when is_binary(value), do: value defp network_access_config_values(%{ thread_opts: %Codex.Thread.Options{sandbox: {:external_sandbox, network_access}} }) do case network_access do :enabled -> ["sandbox_external.network_access=true"] :restricted -> ["sandbox_external.network_access=false"] _ -> [] end end defp network_access_config_values(%{ thread_opts: %Codex.Thread.Options{network_access_enabled: value} }) when value in [true, false] do ["sandbox_workspace_write.network_access=#{value}"] end defp network_access_config_values(_thread), do: [] defp global_config_overrides(%ExecOptions{codex_opts: %Options{} = opts}) do opts |> Map.get(:config_overrides, []) |> Overrides.normalize_config_overrides() end defp thread_config_overrides(%ExecOptions{ thread: %{thread_opts: %Codex.Thread.Options{} = opts} }) do opts |> Map.get(:config_overrides, []) |> Overrides.normalize_config_overrides() end defp thread_config_overrides(_), do: {:ok, []} defp turn_config_overrides(%ExecOptions{turn_opts: %{} = opts}) do opts |> fetch_turn_opt(:config_overrides) |> Overrides.normalize_config_overrides() end defp derived_config_overrides(%ExecOptions{codex_opts: %Options{} = opts} = exec_opts) do thread_opts = case exec_opts.thread do %{thread_opts: %Codex.Thread.Options{} = thread_opts} -> thread_opts _ -> nil end Overrides.derived_overrides(opts, thread_opts) end defp build_env(%ExecOptions{codex_opts: %Options{} = opts, env: env}) do RuntimeEnv.base_overrides(opts.api_key, opts.base_url) |> Map.merge(payload_env_overrides(opts), fn _key, _base, payload -> payload end) |> Map.merge(env, fn _key, _base, custom -> custom end) end defp validate_governed_runtime( %ExecOptions{codex_opts: %Options{governed_authority: authority} = opts} = exec_opts, config_values ) do env = build_env(exec_opts) with :ok <- GovernedAuthority.validate_clear_env(authority, exec_opts.clear_env?, :exec), :ok <- GovernedAuthority.validate_runtime_env(authority, env), :ok <- GovernedAuthority.reject_config_overrides(authority, config_values, :exec) do GovernedAuthority.validate_command_override(authority, opts.codex_path_override, :exec) end end defp attachment_paths(attachments) do attachments |> List.wrap() |> Enum.flat_map(fn %Attachment{path: path} when is_binary(path) and path != "" -> [path] _attachment -> [] end) end defp resume_args(%{thread_id: thread_id}) when is_binary(thread_id), do: ["resume", thread_id] defp resume_args(%{resume: :last}), do: ["resume", "--last"] defp resume_args(_), do: [] defp sandbox_mode(:strict), do: "read-only" defp sandbox_mode(:default), do: nil defp sandbox_mode(:permissive), do: "danger-full-access" defp sandbox_mode(:read_only), do: "read-only" defp sandbox_mode(:workspace_write), do: "workspace-write" defp sandbox_mode(:danger_full_access), do: "danger-full-access" defp sandbox_mode(:external_sandbox), do: "external-sandbox" defp sandbox_mode({:external_sandbox, _network_access}), do: "external-sandbox" defp sandbox_mode("read-only"), do: "read-only" defp sandbox_mode("workspace-write"), do: "workspace-write" defp sandbox_mode("danger-full-access"), do: "danger-full-access" defp sandbox_mode("external-sandbox"), do: "external-sandbox" defp sandbox_mode(nil), do: nil defp sandbox_mode(value) when is_binary(value), do: value defp sandbox_mode(_), do: nil defp exec_opt(%ExecOptions{} = exec_opts, key) when is_atom(key) do case fetch_turn_opt(exec_opts.turn_opts, key) do nil -> fetch_thread_opt(exec_opts.thread, key) value -> value end end defp payload_oss?(%ExecOptions{codex_opts: %Options{} = opts} = exec_opts) do case payload_provider_backend(opts.model_payload) do backend when backend in [:oss, "oss"] -> true _ -> exec_opt(exec_opts, :oss) == true end end defp payload_local_provider(%ExecOptions{codex_opts: %Options{} = opts} = exec_opts) do case payload_backend_metadata(opts.model_payload) do %{"oss_provider" => provider} when is_binary(provider) and provider != "" -> provider _ -> exec_opt(exec_opts, :local_provider) end end defp payload_config_values(%ExecOptions{codex_opts: %Options{} = opts}) do opts.model_payload |> payload_backend_metadata() |> Map.get("config_values", []) |> List.wrap() |> Enum.filter(&(is_binary(&1) and &1 != "")) end defp payload_env_overrides(%Options{model_payload: payload}) do payload |> case do payload when is_map(payload) -> Map.get(payload, :env_overrides, Map.get(payload, "env_overrides", %{})) _ -> %{} end |> case do env when is_map(env) -> Map.new(env, fn {key, value} -> {to_string(key), to_string(value)} end) _ -> %{} end end defp payload_backend_metadata(payload) when is_map(payload) do Map.get(payload, :backend_metadata, Map.get(payload, "backend_metadata", %{})) |> case do metadata when is_map(metadata) -> metadata _ -> %{} end end defp payload_backend_metadata(_payload), do: %{} defp payload_provider_backend(payload) when is_map(payload) do Map.get(payload, :provider_backend, Map.get(payload, "provider_backend")) end defp payload_provider_backend(_payload), do: nil defp fetch_turn_opt(%{} = opts, key) when is_atom(key) do case Map.fetch(opts, key) do {:ok, value} -> value :error -> case Map.fetch(opts, Atom.to_string(key)) do {:ok, value} -> value :error -> nil end end end defp fetch_thread_opt(%{thread_opts: %Codex.Thread.Options{} = opts}, key) when is_atom(key) do case Map.fetch(opts, key) do {:ok, value} -> value :error -> nil end end defp fetch_thread_opt(_thread, _key), do: nil defp normalize_option_string(value) when is_atom(value), do: Atom.to_string(value) defp normalize_option_string(value) when is_binary(value) and value != "", do: value defp normalize_option_string(_value), do: nil defp normalize_string_list(values) when is_list(values) do Enum.filter(values, &(is_binary(&1) and &1 != "")) end defp normalize_string_list(_values), do: [] defp normalize_prompt(value) when is_binary(value) and value != "", do: value defp normalize_prompt(_value), do: nil defp config_values_from_cli_args(args) when is_list(args) do args |> Enum.chunk_every(2) |> Enum.flat_map(fn ["--config", value] -> [value] _other -> [] end) end defp transport_stderr_buffer_size(%ExecOptions{max_stderr_buffer_bytes: nil}), do: 262_145 defp transport_stderr_buffer_size(%ExecOptions{max_stderr_buffer_bytes: max_bytes}) when is_integer(max_bytes) and max_bytes > 0 do max_bytes + 1 end defp transport_stderr_buffer_size(_exec_opts), do: 262_145 defp exit_code(exit) do case CoreProcessExit.status(exit) do :success -> 0 :exit -> core_exit_code(exit) :signal -> normalized_exit_code(exit) _other -> core_exit_code(exit) end end defp core_exit_code(exit) do case CoreProcessExit.code(exit) do code when is_integer(code) -> code _other -> -1 end end defp normalized_exit_code(exit) do case ProcessExit.exit_status(exit) do {:ok, status} -> status :unknown -> -1 end end defp exit_message(exit) do code = CoreProcessExit.code(exit) status = CoreProcessExit.status(exit) signal = CoreProcessExit.signal(exit) reason = CoreProcessExit.reason(exit) cond do status == :exit and is_integer(code) -> "codex executable exited with status #{code}" status == :signal -> "codex executable exited due to signal #{inspect(signal)}" true -> "codex executable exited: #{inspect(reason)}" end end defp normalize_reason_code(nil), do: nil defp normalize_reason_code(code) when is_atom(code), do: code defp normalize_reason_code(code) when is_binary(code) do normalized = code |> String.downcase() |> String.replace("-", "_") Map.get(@known_reason_codes, normalized, :provider_error) end defp retryable_exit?(exit), do: Codex.TransportError.retryable_status?(exit_code(exit)) end