ObanCodex.Agent (oban_codex v0.3.2)

Copy Markdown View Source

A facade over long-lived agent processes whose turns run as Oban jobs.

Experimental

The agent layer is the stateful floor above the stateless ObanCodex seam. It is opt-in (nothing runs unless ObanCodex.Agent.Supervisor is in your tree) and experimental: the API may change between minor releases while it carries this marker. See the Agent lifecycle guide.

One agent is one ObanCodex.Agent.Instance (:gen_statem) registered under a caller-chosen id. A prompt does not block on Codex: it enqueues an ObanCodex.Worker job and the state machine parks in :running until the worker reports back through job_finished/3. All interaction goes through this module; nothing here messages a process that is not running.

{:ok, _pid} = ObanCodex.Agent.start_agent("triage-7",
  args: ObanCodex.Args.defaults(model: "gpt-5"))

:processing = ObanCodex.Agent.submit_prompt("triage-7", "triage the new issues")
{:ok, :running} = ObanCodex.Agent.status("triage-7")

{:ok, {:awaiting_permission, %{id: id}}} =
  ObanCodex.Agent.await("triage-7", [:idle, :awaiting_permission, :waiting_for_user])
:processing = ObanCodex.Agent.approve_action("triage-7", id)

Requires ObanCodex.Agent.Supervisor in the host supervision tree.

Summary

Types

The caller-chosen agent identity, unique per running agent.

What status/1 returns: a bare state atom, except the two gated states, which atomically carry what they are gated on (the action map holds the :id that approve_action/2 / reject_action/3 take). :offline when the agent is not running.

Functions

Approve the pending action by id (read it off status/1 or await/3).

Block until the agent settles into one of states (a state atom or list of them), returning the full status/0 it landed on -- so a gated state arrives with its payload

The fire-and-forget form of submit_prompt/3: never blocks the caller.

Asynchronously force the agent into :paused lockdown, from any state. Drops any pending action or question.

Fork one named conversation arc into another and run prompt on the fork.

The agent's event log, oldest first. Works in every state. Result entries are {:result, structured_output_map} for --json-schema turns and {:result, text} otherwise.

The agent's bookkeeping in one map. :session_id remains the legacy default arc handle; :session_arcs, :active_arc_id, and :continuation expose the named-arc state and the current or most recent fresh/resume decision. Also includes :state, :turns, accumulated :cost_usd, and any pending gate.

The return path for workers: report a finished turn back to its agent.

The retry half of the return path: report a failed-but-retryable attempt.

All running agents as {agent_id, status} pairs, straight off the registry (no process is messaged). Order is unspecified. The fleet-inventory read for dashboards and supervisors.

Reject the pending action: the denial is recorded and the agent returns to :idle.

Release a :paused agent back to :idle.

Spawn a new agent under the dynamic supervisor.

Read the agent's lifecycle status from registry metadata, without messaging the process: one atomic read of the state and, in the gated states, the pending action or question it is gated on -- so there is no torn status-then-payload sequence.

Cleanly stop a running agent (:transient, so it is not restarted).

Send a prompt to a running agent. Synchronous: replies :processing once the turn's Oban job is enqueued.

Types

agent_id()

@type agent_id() :: term()

The caller-chosen agent identity, unique per running agent.

status()

@type status() ::
  :idle
  | :running
  | :paused
  | :offline
  | {:awaiting_permission, %{id: String.t(), description: String.t()}}
  | {:waiting_for_user, String.t() | nil}

What status/1 returns: a bare state atom, except the two gated states, which atomically carry what they are gated on (the action map holds the :id that approve_action/2 / reject_action/3 take). :offline when the agent is not running.

Functions

approve_action(agent_id, action_id, opts \\ [])

@spec approve_action(agent_id(), String.t(), keyword()) ::
  :processing | {:error, term()}

Approve the pending action by id (read it off status/1 or await/3).

The continuation turn resumes the Codex session with the approved action as its prompt, and carries the agent's :approved_args (e.g. a sandbox elevation) merged over the default args -- so approval actually unlocks the tools the action needs, on that turn only.

Options

  • :args -- a string-keyed map of Codex args merged over :approved_args for this continuation only. The caller can size the elevation to the action that was approved. Nothing is remembered; a later approval uses the standing args again unless it supplies another override. Non-string keys are refused with {:error, {:invalid_args, keys}} and the action stays pending.

await(agent_id, states, timeout \\ 60000)

@spec await(agent_id(), atom() | [atom()], timeout :: pos_integer()) ::
  {:ok, status()} | {:error, :timeout}

Block until the agent settles into one of states (a state atom or list of them), returning the full status/0 it landed on -- so a gated state arrives with its payload:

case ObanCodex.Agent.await("live", [:idle, :awaiting_permission], 300_000) do
  {:ok, {:awaiting_permission, %{id: id}}} -> ObanCodex.Agent.approve_action("live", id)
  {:ok, :idle} -> :done
end

Polls the registry (25ms), so it never messages the agent. :offline is awaitable (e.g. after stop_agent/1). Returns {:error, :timeout} when the deadline passes first.

cast_prompt(agent_id, prompt, opts \\ [])

@spec cast_prompt(agent_id(), String.t(), keyword()) ::
  :ok | {:error, :agent_not_running}

The fire-and-forget form of submit_prompt/3: never blocks the caller.

Returns :ok as soon as the cast is sent -- there is no :processing acknowledgment and no error reply on a failed enqueue (that lands in history/1 as {:enqueue_failed, reason}). State handling matches submit_prompt/3 -- queued in :running / :awaiting_permission, the answer in :waiting_for_user -- except :paused, where the prompt is dropped (recorded as {:dropped_prompt, text}) since lockdown has no caller to refuse. Use this from callers that must not block on a busy agent (a LiveView event handler, a scheduler); use submit_prompt/3 when you want backpressure and the enqueue acknowledgment. Takes the same options.

emergency_pause(agent_id)

@spec emergency_pause(agent_id()) :: :ok | {:error, :agent_not_running}

Asynchronously force the agent into :paused lockdown, from any state. Drops any pending action or question.

fork_arc(agent_id, source_arc_id, target_arc_id, prompt, opts \\ [])

@spec fork_arc(agent_id(), String.t(), String.t(), String.t(), keyword()) ::
  :processing | {:error, term()}

Fork one named conversation arc into another and run prompt on the fork.

The current codex_wrapper command surface does not expose codex exec fork, so this returns {:error, :fork_unsupported} without enqueueing a turn. The stable API is present now so support can land without changing the host contract.

history(agent_id)

@spec history(agent_id()) :: {:ok, list()} | {:error, :agent_not_running}

The agent's event log, oldest first. Works in every state. Result entries are {:result, structured_output_map} for --json-schema turns and {:result, text} otherwise.

info(agent_id)

@spec info(agent_id()) :: {:ok, map()} | {:error, :agent_not_running}

The agent's bookkeeping in one map. :session_id remains the legacy default arc handle; :session_arcs, :active_arc_id, and :continuation expose the named-arc state and the current or most recent fresh/resume decision. Also includes :state, :turns, accumulated :cost_usd, and any pending gate.

job_finished(agent_id, payload)

This function is deprecated. pass the job metadata as the third argument.
@spec job_finished(agent_id(), term()) :: {:error, :turn_identity_required}

job_finished(agent_id, payload, captured_meta)

@spec job_finished(
  agent_id(),
  {:ok, CodexWrapper.Result.t()} | {:error, term(), term()},
  map()
) :: :ok | {:error, :agent_not_running | :turn_identity_required}

The return path for workers: report a finished turn back to its agent.

ObanCodex.Agent.Job calls this from handle_result/2 / handle_error/3 with the complete metadata captured by the job. Payload shapes: {:ok, %CodexWrapper.Result{}} on success, or {:error, oban_return, payload} for a failed turn. Fire-and-forget: if the agent is gone the outcome is dropped.

job_retrying(agent_id, retry)

This function is deprecated. pass the job metadata as the third argument.
@spec job_retrying(agent_id(), map()) :: {:error, :turn_identity_required}

job_retrying(agent_id, retry, captured_meta)

@spec job_retrying(agent_id(), map(), map()) ::
  :ok | {:error, :agent_not_running | :turn_identity_required}

The retry half of the return path: report a failed-but-retryable attempt.

ObanCodex.Agent.Job's error callback calls this when Oban will re-run the job (an {:error, _} with attempts remaining, or a {:snooze, _}). The machine stays in :running -- the logical turn is still in flight -- records {:retrying, retry} in history, and re-arms the :job_timeout watchdog to cover the retry's backoff plus execution. retry is %{attempt:, max_attempts:, verdict:}. Fire-and-forget, like job_finished/3.

list()

@spec list() :: [{agent_id(), status()}]

All running agents as {agent_id, status} pairs, straight off the registry (no process is messaged). Order is unspecified. The fleet-inventory read for dashboards and supervisors.

reject_action(agent_id, action_id, reason \\ "denied")

@spec reject_action(agent_id(), String.t(), String.t()) ::
  :rejected | {:error, term()}

Reject the pending action: the denial is recorded and the agent returns to :idle.

resume_agent(agent_id)

@spec resume_agent(agent_id()) :: :resumed | {:error, term()}

Release a :paused agent back to :idle.

start_agent(agent_id, config \\ [])

@spec start_agent(agent_id(), keyword() | map()) :: {:ok, pid()} | {:error, term()}

Spawn a new agent under the dynamic supervisor.

Starting an id that is already running returns the existing pid, so the call is idempotent. See ObanCodex.Agent.Instance for the config keys.

status(agent_id)

@spec status(agent_id()) :: {:ok, status()}

Read the agent's lifecycle status from registry metadata, without messaging the process: one atomic read of the state and, in the gated states, the pending action or question it is gated on -- so there is no torn status-then-payload sequence.

{:ok, :running} = status("live")
{:ok, {:awaiting_permission, %{id: id, description: d}}} = status("live")
{:ok, {:waiting_for_user, question}} = status("live")

An unknown or stopped agent reads as {:ok, :offline}. Registry cleanup after a process death is asynchronous, so :offline after stop_agent/1 is eventually-consistent -- await(id, :offline) if you need to block on it.

stop_agent(agent_id)

@spec stop_agent(agent_id()) :: :ok | {:error, :agent_not_running}

Cleanly stop a running agent (:transient, so it is not restarted).

submit_prompt(agent_id, prompt, opts \\ [])

@spec submit_prompt(agent_id(), String.t(), keyword()) ::
  :processing | {:error, term()}

Send a prompt to a running agent. Synchronous: replies :processing once the turn's Oban job is enqueued.

In :running and :awaiting_permission the call blocks (the machine postpones it) until the turn finishes / the gate clears; in :waiting_for_user the prompt is the answer to the pending question and resumes the Codex session.

Options (shared with cast_prompt/3):

  • :session -- :resume (default) continues the agent's Codex session; :fresh starts a new one for this turn (the resume handle is cleared at delivery time, so it composes with queued prompts); :fresh_fallback records that the host deliberately started fresh after a failed resume.
  • :arc_id -- an opaque non-empty string naming the provider conversation; defaults to "default". Arcs isolate their session handles while the agent remains single-turn-at-a-time.
  • :origin -- :operator (default) or :tick. A :tick prompt is a scheduled delivery: in :waiting_for_user it queues behind the pending question instead of being consumed as the answer. ObanCodex.Agent.Tick sets this; operators rarely should.