CodexEx.AppServer.TurnStream (codex_ex v0.1.0)

Copy Markdown View Source

Typed collector for a single app-server turn.

The stream subscribes to the client, scopes notifications to one turn, and accumulates typed turn state until completion.

Summary

Types

state()

@type state() :: %{
  direct?: boolean(),
  client: CodexEx.AppServer.Client.t(),
  done?: boolean(),
  event_final_turn: CodexEx.AppServer.Turn.t() | nil,
  final_text: binary(),
  initial_turn: CodexEx.AppServer.Turn.t() | nil,
  item_order_rev: [term()],
  item_ref_by_id: %{optional(binary()) => term()},
  items_by_ref: %{optional(term()) => CodexEx.AppServer.ThreadItem.t()},
  next_item_ref: non_neg_integer(),
  pending_messages_rev: [CodexEx.AppServer.Message.t()],
  rpc_final_turn: CodexEx.AppServer.Turn.t() | nil,
  request_task: Task.t(),
  result: {:ok, CodexEx.AppServer.Turn.t()} | {:error, term()} | nil,
  text_deltas_rev: [binary()],
  thread_id: binary(),
  turn_id: binary() | nil,
  usage: CodexEx.AppServer.TokenUsage.t() | nil,
  waiters: [GenServer.from()],
  stop_when_done?: boolean(),
  subscribed?: boolean()
}

t()

@type t() :: %CodexEx.AppServer.TurnStream{
  final_text: binary(),
  final_turn: CodexEx.AppServer.Turn.t() | nil,
  initial_turn: CodexEx.AppServer.Turn.t() | nil,
  items: [CodexEx.AppServer.ThreadItem.t()],
  pid: pid(),
  text_deltas: [binary()],
  thread_id: binary(),
  turn_id: binary() | nil,
  usage: CodexEx.AppServer.TokenUsage.t() | nil
}

Functions

child_spec(init_arg)

Returns a specification to start this module under a supervisor.

See Supervisor.

ensure_success(turn_stream)

@spec ensure_success(t()) :: :ok | {:error, term()}

final_json(turn_stream)

@spec final_json(t()) :: {:ok, term()} | {:error, term()}

snapshot(pid)

@spec snapshot(
  %CodexEx.AppServer.TurnStream{
    final_text: term(),
    final_turn: term(),
    initial_turn: term(),
    items: term(),
    pid: term(),
    text_deltas: term(),
    thread_id: term(),
    turn_id: term(),
    usage: term()
  }
  | pid()
) :: %CodexEx.AppServer.TurnStream{
  final_text: term(),
  final_turn: term(),
  initial_turn: term(),
  items: term(),
  pid: term(),
  text_deltas: term(),
  thread_id: term(),
  turn_id: term(),
  usage: term()
}

start(client, thread_id, params, direct? \\ true)

@spec start(CodexEx.AppServer.Client.t(), binary(), map(), boolean()) ::
  {:ok,
   %CodexEx.AppServer.TurnStream{
     final_text: term(),
     final_turn: term(),
     initial_turn: term(),
     items: term(),
     pid: term(),
     text_deltas: term(),
     thread_id: term(),
     turn_id: term(),
     usage: term()
   }}
  | {:error, term()}

start_request(client, thread_id, request_fun, direct? \\ true)

@spec start_request(
  CodexEx.AppServer.Client.t(),
  binary(),
  (-> {:ok, CodexEx.AppServer.Turn.t()} | {:error, term()}),
  boolean()
) ::
  {:ok,
   %CodexEx.AppServer.TurnStream{
     final_text: term(),
     final_turn: term(),
     initial_turn: term(),
     items: term(),
     pid: term(),
     text_deltas: term(),
     thread_id: term(),
     turn_id: term(),
     usage: term()
   }}
  | {:error, term()}

wait(turn_stream, timeout \\ 1_800_000)

@spec wait(term(), term()) ::
  {:ok,
   %CodexEx.AppServer.TurnStream{
     final_text: term(),
     final_turn: term(),
     initial_turn: term(),
     items: term(),
     pid: term(),
     text_deltas: term(),
     thread_id: term(),
     turn_id: term(),
     usage: term()
   }}
  | {:error, term()}