Packages
langchain
0.9.2
0.9.2
0.9.1
0.9.0
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.0
0.6.3
0.6.2
0.6.1
0.6.0
0.5.2
0.5.1
0.5.0
0.4.1
0.4.0
0.4.0-rc.3
0.4.0-rc.2
0.4.0-rc.1
0.4.0-rc.0
0.3.3
0.3.2
0.3.1
0.3.0
0.3.0-rc.2
0.3.0-rc.1
0.3.0-rc.0
0.2.0
0.1.10
0.1.9
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
Elixir implementation of a LangChain style framework that lets Elixir projects integrate with and leverage LLMs.
Current section
Files
Jump to
Current section
Files
lib/chat_models/chat_open_ai_responses.ex
defmodule LangChain.ChatModels.ChatOpenAIResponses do
@moduledoc """
Represents the OpenAI Responses API
Parses and validates inputs for making requests to the OpenAI Responses API.
Converts responses into more specialized `LangChain` data structures.
## ContentPart Types
OpenAI's Responses API supports several types of content parts that can be combined in a single message:
### Text Content
Basic text content is the default and most common type:
Message.new_user!("Hello, how are you?")
### Image Content
OpenAI supports both base64-encoded images and image URLs:
# Using a base64 encoded image
Message.new_user!([
ContentPart.text!("What's in this image?"),
ContentPart.image!("base64_encoded_image_data", media: :jpg)
])
# Using an image URL
Message.new_user!([
ContentPart.text!("Describe this image:"),
ContentPart.image_url!("https://example.com/image.jpg")
])
# Using a file ID (after uploading to OpenAI)
Message.new_user!([
ContentPart.text!("Describe this image:"),
ContentPart.image!("file-1234", type: :file_id)
])
For images, you can specify the detail level which affects token usage:
- `detail: "low"` - Lower resolution, fewer tokens
- `detail: "high"` - Higher resolution, more tokens
- `detail: "auto"` - Let the model decide
### File Content
OpenAI supports both base64-encoded files and file IDs:
# Using a base64 encoded file
Message.new_user!([
ContentPart.text!("Process this file:"),
ContentPart.file!("base64_encoded_file_data",
type: :base64,
filename: "document.pdf"
)
])
# Using a file ID (after uploading to OpenAI)
Message.new_user!([
ContentPart.text!("Process this file:"),
ContentPart.file!("file-1234", type: :file_id)
])
## Callbacks
See the set of available callbacks: `LangChain.Chains.ChainCallbacks`
### Rate Limit API Response Headers
OpenAI returns rate limit information in the response headers. Those can be
accessed using the LLM callback `on_llm_ratelimit_info` like this:
handlers = %{
on_llm_ratelimit_info: fn _model, headers ->
IO.inspect(headers)
end
}
{:ok, chat} = ChatOpenAI.new(%{callbacks: [handlers]})
When a request is received, something similar to the following will be output
to the console.
%{
"x-ratelimit-limit-requests" => ["5000"],
"x-ratelimit-limit-tokens" => ["160000"],
"x-ratelimit-remaining-requests" => ["4999"],
"x-ratelimit-remaining-tokens" => ["159973"],
"x-ratelimit-reset-requests" => ["12ms"],
"x-ratelimit-reset-tokens" => ["10ms"],
"x-request-id" => ["req_1234"]
}
### Token Usage
OpenAI returns token usage information as part of the response body. The
`LangChain.TokenUsage` is added to the `metadata` of the `LangChain.Message`
and `LangChain.MessageDelta` structs that are processed under the `:usage`
key.
The OpenAI documentation instructs to provide the `stream_options` with the
`include_usage: true` for the information to be provided.
The `TokenUsage` data is accumulated for `MessageDelta` structs and the final usage information will be on the `LangChain.Message`.
NOTE: Of special note is that the `TokenUsage` information is returned once
for all "choices" in the response. The `LangChain.TokenUsage` data is added to
each message, but if your usage requests multiple choices, you will see the
same usage information for each choice but it is duplicated and only one
response is meaningful.
## Native Tools (Web Search)
Open AI's Responses API also supports built-in tools. Among those, we support Web Search currently.
### Example
To optionally permit the model to use web search:
native_web_tool = NativeTool.new!(%{name: "web_search_preview", configuration: %{}})
%{llm: ChatOpenAIResponses.new!(%{model: "gpt-4o"})}
|> LLMChain.new!()
|> LLMChain.add_message(Message.new_user!("Can you tell me something that happened today in Texas?"))
|> LLMChain.add_tools(web_tool)
|> LLMChain.run()
You may provide additional configuration per the OpenAI documentation:
web_config = %{
search_context_size: "medium",
user_location: %{
type: "approximate",
city: "Humble",
country: "US",
region: "Texas",
timezone: "America/Chicago"
}
}
native_web_tool = NativeTool.new!(%{name: "web_search_preview", configuration: web_config)
You may reference a prior web_search_call in subsequent runs as:
Message.new_assistant!([
ContentPart.new!(%{
type: :unsupported,
options: %{
id: "ws_123456789", # ID as provided from Open AI
status: "completed",
type: "web_search_call"
}
}
),
ContentPart.text!("The Astros won today 5-4...")
])
Note: Not all Open AI models support `web_search_preview`. OpenAI will return an error if you request web_search_preview for when using a model that doesn't support it.
## Tool Choice
OpenAI's ChatGPT API supports forcing a tool to be used.
- https://platform.openai.com/docs/api-reference/chat/create#chat-create-tool_choice
This is supported through the `tool_choice` options. It takes a plain Elixir
map to provide the configuration.
By default, the LLM will choose a tool call if a tool is available and it
determines it is needed. That's the "auto" mode.
### Example
For the LLM's response to make a tool call of the "get_weather" function.
ChatOpenAI.new(%{
model: "...",
tool_choice: %{"type" => "function", "function" => %{"name" => "get_weather"}}
})
...or to force a native tool (such as web search):
ChatOpenAI.new(%{
model: "...",
tool_choice: "web_search_preview"
})
## Verbosity
The `verbosity` option controls the length of the model's response. Accepted
values are `"low"`, `"medium"`, and `"high"`. When omitted, the API uses its
default behavior.
This is sent as part of the `text` parameter in the Responses API and can be
combined with JSON response formats.
Only supported for gpt-5 or newer models
### Example
ChatOpenAIResponses.new!(%{model: "gpt-5", verbosity: "low"})
## WebSocket Transport
Instead of HTTP, requests can be sent over a persistent WebSocket connection
for lower latency. Use `connect_websocket!/1` to open a connection and
`disconnect_websocket!/1` to close it:
model =
ChatOpenAIResponses.new!(%{model: "gpt-4o"})
|> ChatOpenAIResponses.connect_websocket!()
{:ok, chain} =
%{llm: model}
|> LLMChain.new!()
|> LLMChain.add_message(Message.new_user!("Hello"))
|> LLMChain.run()
ChatOpenAIResponses.disconnect_websocket!(model)
The WebSocket connection is reused across multiple LLM calls within the same
chain run (e.g. multi-turn tool calling with `:while_needs_response`).
### Lifecycle Management
**The application is responsible for managing the WebSocket lifecycle.**
`connect_websocket!/1` starts a `LangChain.WebSocket` GenServer via
`start_link/1`, linking it to the calling process. The PID is stored in
the model struct's `:websocket` field. There is no supervisor, automatic
reconnection, or health monitoring built in.
This means:
- If the calling process exits, the WebSocket is terminated (process link).
- The WebSocket PID cannot be serialized. If the model struct is persisted
to a database and restored later, the `:websocket` field will be stale.
- The server may close idle connections at any time. There is no automatic
reconnection.
- There is no retry logic for WebSocket failures (unlike the HTTP transport).
**The WebSocket transport is best suited for short-lived, synchronous
sessions** where you control the full lifecycle. It is not currently safe
for long-lived agent processes, human-in-the-loop workflows with
interruptions, or any scenario where the model struct is serialized and
restored across process boundaries.
For long-running or interruptible workloads, use the default HTTP transport.
### Known Limitation: temperature and top_p
The `:temperature` and `:top_p` parameters are currently excluded from
WebSocket payloads due to an
[OpenAI bug](https://community.openai.com/t/1375536) that silently closes
the connection when these are sent as decimals. A `Logger.warning` is
emitted when these values are dropped. This workaround will be removed
once OpenAI resolves the issue.
## Connection Retry Behavior
The `retry_count` option controls how many times a request is retried when
a pooled HTTP connection turns out to be stale (server closed it between
requests). This is a transport-level issue where retrying with a fresh
connection is the correct response.
**Only closed-connection errors are retried.** Timeouts, rate limits (429),
overloaded (529), authentication errors, and invalid requests all return
immediately -- they are not problems that a simple retry will fix.
| `retry_count` | Total HTTP requests |
|---|---|
| `0` | 1 (no retries) |
| `1` | 2 (1 initial + 1 retry) |
| `2` (default) | 3 (1 initial + 2 retries) |
Req's built-in HTTP retry is disabled to prevent the two retry layers from
compounding. See [GitHub issue #503](https://github.com/brainlid/langchain/issues/503).
When running LLM calls from a background job queue (e.g., Oban) that has its
own retry logic, set `retry_count: 0` so there are no hidden retries:
ChatOpenAIResponses.new!(%{model: "...", retry_count: 0})
"""
use Ecto.Schema
require Logger
import Ecto.Changeset
alias __MODULE__
alias LangChain.Config
alias LangChain.ChatModels.ChatModel
alias LangChain.PromptTemplate
alias LangChain.Message
alias LangChain.Message.Citation
alias LangChain.Message.ContentPart
alias LangChain.Message.ToolCall
alias LangChain.Message.ToolResult
alias LangChain.TokenUsage
alias LangChain.Function
alias LangChain.NativeTool
alias LangChain.FunctionParam
alias LangChain.LangChainError
alias LangChain.Utils
alias LangChain.MessageDelta
alias LangChain.Callbacks
alias LangChain.ChatModels.ReasoningOptions
@behaviour ChatModel
@current_config_version 1
@receive_timeout 60_000
@primary_key false
# https://platform.openai.com/docs/api-reference/responses/create
embedded_schema do
field :receive_timeout, :integer, default: @receive_timeout
field :api_key, :string, redact: true
field :endpoint, :string, default: "https://api.openai.com/v1/responses"
field :model, :string, default: "gpt-3.5-turbo"
field :include, {:array, :string}, default: []
# omit instructions becasue langchain assumes statelessness
field :max_output_tokens, :integer, default: nil
# omit metadata because chat_open_ai also omits it
# omit parallel_tool_calls because chat_open_ai also omits it
field :previous_response_id, :string, default: nil
# Reasoning options for gpt-5 and o-series models
embeds_one(:reasoning, ReasoningOptions)
# omit service_tier because chat_open_ai also omits it
field :store, :boolean, default: false
field :stream, :boolean, default: false
field :temperature, :float, default: nil
field :json_response, :boolean, default: false
field :json_schema, :map, default: nil
field :json_schema_name, :string, default: nil
# This can be a string or object. We will need to allow ["none", "auto", "required", "file_search", "web_search_preview", and "computer_use_preview"] and take any other string and turn it to %{name: value, type: "function"}
field :tool_choice, :any, default: nil, virtual: true
field :top_p, :float, default: 1.0
field :truncation, :string
field :verbosity, :string, default: nil
field :user, :string
field :callbacks, {:array, :map}, default: []
field :verbose_api, :boolean, default: false
# Number of retries on closed-connection errors (stale pool). The initial
# request always runs; this controls additional attempts only.
field :retry_count, :integer, default: 2
# Req options to merge into the request.
# Refer to `https://hexdocs.pm/req/Req.html#new/1-options` for
# `Req.new` supported set of options.
field :req_config, :map, default: %{}
# Optional WebSocket transport. When set to a PID of a
# `LangChain.WebSocket` process, requests will be sent over the
# persistent WebSocket connection instead of HTTP.
field :websocket, :any, virtual: true, default: nil
end
@type t :: %ChatOpenAIResponses{}
# Omits callbacks. Otherwise identical to above.
@create_fields [
:receive_timeout,
:api_key,
:endpoint,
:model,
:include,
:max_output_tokens,
:previous_response_id,
:store,
:stream,
:temperature,
:json_response,
:json_schema,
:json_schema_name,
:tool_choice,
:top_p,
:truncation,
:verbosity,
:user,
:verbose_api,
:retry_count,
:req_config,
:websocket
]
@required_fields [:endpoint, :model]
@spec get_api_key(t()) :: String.t()
defp get_api_key(%ChatOpenAIResponses{api_key: api_key}) do
# if no API key is set default to `""` which will raise a OpenAI API error
api_key || Config.resolve(:openai_key, "")
end
@spec get_org_id() :: String.t() | nil
defp get_org_id() do
Config.resolve(:openai_org_id)
end
@spec get_proj_id() :: String.t() | nil
defp get_proj_id() do
Config.resolve(:openai_proj_id)
end
@doc """
Setup a ChatOpenAI client configuration.
"""
@spec new(attrs :: map()) :: {:ok, t} | {:error, Ecto.Changeset.t()}
def new(%{} = attrs \\ %{}) do
%ChatOpenAIResponses{}
|> cast(attrs, @create_fields)
|> cast_embed(:reasoning)
|> common_validation()
|> apply_action(:insert)
end
@doc """
Setup a ChatOpenAI client configuration and return it or raise an error if invalid.
"""
@spec new!(attrs :: map()) :: t() | no_return()
def new!(attrs \\ %{}) do
case new(attrs) do
{:ok, chain} ->
chain
{:error, changeset} ->
raise LangChainError, changeset
end
end
@doc """
Open a `LangChain.WebSocket` connection using the model's endpoint and API key.
Returns `{:ok, model}` with the `:websocket` field set to the WebSocket PID,
or `{:error, reason}` on failure.
Requires the optional `mint_web_socket` dependency.
## Example
{:ok, model} = ChatOpenAIResponses.connect_websocket(model)
"""
@spec connect_websocket(t()) :: {:ok, t()} | {:error, String.t()}
if Code.ensure_loaded?(Mint.WebSocket) do
def connect_websocket(%ChatOpenAIResponses{} = model) do
api_key = model.api_key || Config.resolve(:openai_key, nil)
if api_key do
ws_url =
model.endpoint
|> String.replace_leading("https://", "wss://")
|> String.replace_leading("http://", "ws://")
case LangChain.WebSocket.start_link(
url: ws_url,
headers: [{"authorization", "Bearer #{api_key}"}],
receive_timeout: model.receive_timeout
) do
{:ok, pid} ->
{:ok, %{model | websocket: pid}}
{:error, reason} ->
{:error, "Failed to connect WebSocket: #{inspect(reason)}"}
end
else
{:error, "API key is required to open a WebSocket connection"}
end
end
else
def connect_websocket(%ChatOpenAIResponses{}) do
{:error, "WebSocket support requires the :mint_web_socket dependency"}
end
end
@doc """
Like `connect_websocket/1` but raises on failure.
## Example
model =
ChatOpenAIResponses.new!(%{model: "gpt-4o"})
|> ChatOpenAIResponses.connect_websocket!()
# ... use model in chains ...
ChatOpenAIResponses.disconnect_websocket!(model)
"""
@spec connect_websocket!(t()) :: t() | no_return()
if Code.ensure_loaded?(Mint.WebSocket) do
def connect_websocket!(%ChatOpenAIResponses{} = model) do
case connect_websocket(model) do
{:ok, model} -> model
{:error, reason} -> raise LangChainError, reason
end
end
else
def connect_websocket!(%ChatOpenAIResponses{}) do
raise LangChainError, "WebSocket support requires the :mint_web_socket dependency"
end
end
@doc """
Close the WebSocket connection associated with this model.
Returns the model with `:websocket` set to `nil`.
Safe to call even if the WebSocket is already closed.
"""
@spec disconnect_websocket!(t()) :: t()
if Code.ensure_loaded?(Mint.WebSocket) do
def disconnect_websocket!(%ChatOpenAIResponses{websocket: pid} = model) when is_pid(pid) do
if Process.alive?(pid), do: LangChain.WebSocket.close(pid)
%{model | websocket: nil}
end
else
def disconnect_websocket!(%ChatOpenAIResponses{websocket: pid} = model) when is_pid(pid) do
%{model | websocket: nil}
end
end
def disconnect_websocket!(%ChatOpenAIResponses{} = model), do: model
defp common_validation(changeset) do
changeset
|> validate_required(@required_fields)
|> validate_number(:temperature, greater_than_or_equal_to: 0, less_than_or_equal_to: 2)
|> validate_number(:top_p, greater_than_or_equal_to: 0, less_than_or_equal_to: 1)
|> validate_number(:receive_timeout, greater_than_or_equal_to: 0)
|> validate_inclusion(:verbosity, ~w(low medium high))
end
@doc """
Return the params formatted for an API request.
"""
@spec for_api(t | Message.t() | Function.t(), message :: [map()], ChatModel.tools()) :: %{
atom() => any()
}
def for_api(%ChatOpenAIResponses{} = openai, messages, tools) do
%{
model: openai.model,
stream: openai.stream,
store: if(openai.previous_response_id, do: true, else: openai.store),
input:
messages
|> Enum.reduce([], fn m, acc ->
case for_api(openai, m) do
%{} = data ->
[data | acc]
data when is_list(data) ->
Enum.reverse(data) ++ acc
end
end)
|> Enum.reverse()
}
|> Utils.conditionally_add_to_map(:include, openai.include)
|> Utils.conditionally_add_to_map(:max_output_tokens, openai.max_output_tokens)
|> Utils.conditionally_add_to_map(:previous_response_id, openai.previous_response_id)
|> Utils.conditionally_add_to_map(:reasoning, ReasoningOptions.to_api_map(openai.reasoning))
|> Utils.conditionally_add_to_map(:text, set_text_format(openai))
|> Utils.conditionally_add_to_map(:tool_choice, get_tool_choice(openai))
|> Utils.conditionally_add_to_map(:truncation, openai.truncation)
|> Utils.conditionally_add_to_map(:tools, get_tools_for_api(openai, tools))
|> Utils.conditionally_add_to_map(:user, openai.user)
|> Utils.conditionally_add_to_map(:temperature, openai.temperature)
|> maybe_add_top_p(openai)
end
# Build the payload for WebSocket mode. Wraps the standard API payload
# in a response.create envelope and removes transport-specific fields.
#
# NOTE: :temperature and :top_p are dropped because the WebSocket endpoint
# silently closes the connection (code 1000) when these are sent as decimals.
# This is an OpenAI bug — see:
# https://community.openai.com/t/responses-websocket-v1-responses-closes-with-code-1000-and-no-events-when-temperature-is-a-decimal-e-g-1-2/1375536
if Code.ensure_loaded?(Mint.WebSocket) do
defp for_api_websocket(%ChatOpenAIResponses{} = openai, messages, tools) do
payload = for_api(openai, messages, tools)
dropped = [:stream, :background, :temperature, :top_p]
if payload[:temperature] || payload[:top_p] do
Logger.warning(
"WebSocket transport: dropping :temperature and :top_p from payload " <>
"due to an OpenAI bug that silently closes the connection when these " <>
"are sent as decimals. See: https://community.openai.com/t/1375536"
)
end
payload
|> Map.drop(dropped)
|> Map.put(:type, "response.create")
|> Jason.encode!()
end
end
# gpt-5.2 and newer do not support the top_p parameter.
# Earlier models (gpt-4.x, gpt-5.0, gpt-5.1) accept top_p.
defp maybe_add_top_p(map, %ChatOpenAIResponses{model: model, top_p: top_p}) do
if supports_top_p?(model) do
Utils.conditionally_add_to_map(map, :top_p, top_p)
else
map
end
end
@doc false
@spec supports_top_p?(String.t()) :: boolean()
def supports_top_p?(model) when is_binary(model) do
# Match models known to support top_p. This set is fixed and won't grow.
cond do
String.starts_with?(model, "gpt-4") -> true
String.starts_with?(model, "gpt-5.0") -> true
String.starts_with?(model, "gpt-5.1") -> true
true -> false
end
end
defp get_tools_for_api(%ChatOpenAIResponses{} = _model, nil), do: []
defp get_tools_for_api(%ChatOpenAIResponses{} = model, tools) do
Enum.map(tools, fn
%Function{} = function ->
for_api(model, function)
%NativeTool{} = tool ->
for_api(model, tool)
end)
end
# JSON schema + optional verbosity
defp set_text_format(%ChatOpenAIResponses{
json_response: true,
json_schema: json_schema,
json_schema_name: json_schema_name,
verbosity: verbosity
})
when not is_nil(json_schema) and not is_nil(json_schema_name) do
%{
"format" => %{
"type" => "json_schema",
"name" => json_schema_name,
"schema" => json_schema,
"strict" => true
}
}
|> maybe_add_verbosity(verbosity)
end
# JSON object + optional verbosity
defp set_text_format(%ChatOpenAIResponses{json_response: true, verbosity: verbosity}) do
%{"format" => %{"type" => "json_object"}}
|> maybe_add_verbosity(verbosity)
end
# Plain text with verbosity
defp set_text_format(%ChatOpenAIResponses{json_response: false, verbosity: verbosity})
when is_binary(verbosity) do
%{"verbosity" => verbosity}
end
# Plain text, no verbosity (default)
defp set_text_format(%ChatOpenAIResponses{json_response: false}) do
# NOTE: The default handling when unspecified is `%{"type" => "text"}`
# This returns a `nil` which has the same effect.
nil
end
defp maybe_add_verbosity(map, nil), do: map
defp maybe_add_verbosity(map, verbosity) when is_binary(verbosity),
do: Map.put(map, "verbosity", verbosity)
defp get_tool_choice(%ChatOpenAIResponses{tool_choice: choice})
when choice in ["none", "auto", "required"],
do: choice
defp get_tool_choice(%ChatOpenAIResponses{tool_choice: choice})
when choice in ["file_search", "web_search_preview", "computer_use_preview"],
do: %{"type" => choice}
defp get_tool_choice(%ChatOpenAIResponses{tool_choice: choice})
when is_binary(choice) and byte_size(choice) > 0,
do: %{"type" => "function", "name" => choice}
defp get_tool_choice(%ChatOpenAIResponses{}), do: nil
@spec for_api(
struct(),
Message.t()
| PromptTemplate.t()
| ToolCall.t()
| ToolResult.t()
| ContentPart.t()
| Function.t()
| NativeTool.t()
) ::
%{String.t() => any()} | [%{String.t() => any()}]
# Function support
def for_api(%ChatOpenAIResponses{} = _model, %Function{} = fun) do
%{
"name" => fun.name,
"parameters" => get_parameters(fun),
"type" => "function"
}
|> Utils.conditionally_add_to_map("description", fun.description)
|> Utils.conditionally_add_to_map("strict", fun.strict)
end
def for_api(
%ChatOpenAIResponses{} = _model,
%NativeTool{name: name, configuration: config}
) do
Map.put_new(config, :type, name)
end
def for_api(%ChatOpenAIResponses{} = model, %Message{role: :system, content: content})
when is_list(content) do
%{
"role" => "system",
"type" => "message",
"content" => content_parts_for_api(model, content)
}
end
def for_api(%ChatOpenAIResponses{} = model, %Message{role: :user, content: content})
when is_list(content) do
%{
"role" => "user",
"type" => "message",
"content" => content_parts_for_api(model, content)
}
end
def for_api(
%ChatOpenAIResponses{} = model,
%Message{role: :tool, tool_results: tool_results}
)
when is_list(tool_results) do
Enum.map(tool_results, &for_api(model, &1))
end
# Native tool calls (such as web_search_call) need to get plucked
# out of the content parts and become their own input items.
def for_api(
%ChatOpenAIResponses{} = model,
%Message{role: :assistant, content: content} = msg
)
when is_list(content) do
native_tool_calls_for_api(model, content) ++
[
%{
"role" => "user",
"type" => "message",
"content" => content_parts_for_api(model, content)
}
] ++
Enum.map(msg.tool_calls || [], &for_api(model, &1))
end
def for_api(
%ChatOpenAIResponses{} = model,
%Message{role: :assistant, tool_calls: tool_calls}
)
when is_list(tool_calls) do
Enum.map(tool_calls, &for_api(model, &1))
end
def for_api(%ChatOpenAIResponses{} = _model, %ToolResult{type: :function} = result) do
# a ToolResult becomes a stand-alone %Message{role: :tool} response.
[%ContentPart{type: :text, content: output, options: []}] = result.content
%{
"call_id" => result.tool_call_id,
"output" => output,
"type" => "function_call_output"
}
end
# ToolCall support
def for_api(%ChatOpenAIResponses{} = _model, %ToolCall{type: :function} = fun) do
%{
"arguments" => Jason.encode!(fun.arguments),
"call_id" => fun.call_id,
"name" => fun.name,
"type" => "function_call"
}
|> Utils.conditionally_add_to_map("status", "completed")
end
def for_api(%ChatOpenAIResponses{} = _model, %PromptTemplate{} = _template) do
raise LangChainError, "PromptTemplates must be converted to messages."
end
def native_tool_calls_for_api(%ChatOpenAIResponses{} = model, content_parts)
when is_list(content_parts) do
Enum.map(content_parts, &native_tool_call_for_api(model, &1))
|> Enum.reject(&is_nil/1)
end
@spec native_tool_call_for_api(any(), any()) ::
nil | %{id: any(), status: any(), type: <<_::120>>}
def native_tool_call_for_api(%ChatOpenAIResponses{} = _model, %ContentPart{
type: :unsupported,
options: %{type: "web_search_call"} = opts
}) do
%{id: opts.id, type: "web_search_call", status: opts.status}
end
def native_tool_call_for_api(_, _), do: nil
@doc """
Convert a list of ContentParts to the expected map of data for the OpenAI API.
"""
def content_parts_for_api(%ChatOpenAIResponses{} = model, content_parts)
when is_list(content_parts) do
Enum.map(content_parts, &content_part_for_api(model, &1))
|> Enum.reject(&is_nil/1)
end
@doc """
Convert a ContentPart to the expected map of data for the OpenAI API.
"""
def content_part_for_api(%ChatOpenAIResponses{} = _model, %ContentPart{type: :text} = part) do
%{"type" => "input_text", "text" => part.content}
end
def content_part_for_api(
%ChatOpenAIResponses{} = _model,
%ContentPart{type: :file_url} = part
) do
%{
"type" => "input_file",
"file_url" => part.content
}
end
def content_part_for_api(
%ChatOpenAIResponses{} = _model,
%ContentPart{type: :file, options: opts} = part
) do
case Keyword.get(opts, :type, :base64) do
:file_id ->
%{
"type" => "input_file",
"file_id" => part.content
}
:base64 ->
%{
"type" => "input_file",
"filename" => Keyword.get(opts, :filename, "file.pdf"),
"file_data" => "data:application/pdf;base64," <> part.content
}
end
end
def content_part_for_api(%ChatOpenAIResponses{} = _model, %ContentPart{type: image} = part)
when image in [:image, :image_url] do
output =
if Keyword.get(part.options, :type) == :file_id do
%{"type" => "input_image", "file_id" => part.content}
else
media_prefix =
case Keyword.get(part.options, :media, nil) do
nil ->
""
type when is_binary(type) ->
"data:#{type};base64,"
type when type in [:jpeg, :jpg] ->
"data:image/jpg;base64,"
:png ->
"data:image/png;base64,"
:gif ->
"data:image/gif;base64,"
:webp ->
"data:image/webp;base64,"
other ->
raise LangChainError,
"Received unsupported media type for ContentPart: #{inspect(other)}"
end
%{
"type" => "input_image",
"image_url" => media_prefix <> part.content
}
end
detail_option = Keyword.get(part.options, :detail, nil)
Utils.conditionally_add_to_map(output, "detail", detail_option)
end
# Thinking content parts are output-only and should be omitted when sending to the API
def content_part_for_api(%ChatOpenAIResponses{} = _model, %ContentPart{type: :thinking}),
do: nil
# Ignore unknown, unsupported content parts
def content_part_for_api(%ChatOpenAIResponses{} = _model, %ContentPart{type: :unsupported}),
do: nil
@doc false
def get_parameters(%Function{parameters: [], parameters_schema: nil} = _fun) do
%{
"type" => "object",
"properties" => %{}
}
end
def get_parameters(%Function{parameters: [], parameters_schema: schema} = _fun)
when is_map(schema) do
schema
end
def get_parameters(%Function{parameters: params} = _fun) do
FunctionParam.to_parameters_schema(params)
end
@impl ChatModel
def call(openai, prompt, tools \\ [])
def call(%ChatOpenAIResponses{} = openai, prompt, tools) when is_binary(prompt) do
messages = [
Message.new_system!(),
Message.new_user!(prompt)
]
call(openai, messages, tools)
end
def call(%ChatOpenAIResponses{} = openai, messages, tools) when is_list(messages) do
metadata = %{
model: openai.model,
provider: provider(),
message_count: length(messages),
tools_count: length(tools)
}
ChatModel.llm_telemetry_span(openai, metadata, fn ->
try do
# Track the prompt being sent
LangChain.Telemetry.llm_prompt(
%{system_time: System.system_time()},
%{model: openai.model, messages: messages}
)
# make base api request and perform high-level success/failure checks
case do_api_request(openai, messages, tools) do
{:error, reason} ->
{:error, reason}
parsed_data ->
# Track the response being received
LangChain.Telemetry.llm_response(
%{system_time: System.system_time()},
%{model: openai.model, response: parsed_data}
)
{:ok, parsed_data}
end
rescue
err in LangChainError ->
{:error, err}
end
end)
end
@spec do_api_request(t(), [Message.t()], ChatModel.tools(), integer() | nil) ::
list() | struct() | {:error, LangChainError.t()}
def do_api_request(openai, messages, tools, retry_count \\ nil)
def do_api_request(_openai, _messages, _tools, 0) do
raise LangChainError, "Retries exceeded. Connection failed."
end
if Code.ensure_loaded?(Mint.WebSocket) do
def do_api_request(
%ChatOpenAIResponses{websocket: ws_pid} = openai,
messages,
tools,
_retry_count
)
when is_pid(ws_pid) do
payload = for_api_websocket(openai, messages, tools)
done_fn = fn event ->
event["type"] in ["response.completed", "response.failed"]
end
case openai.stream do
false ->
case LangChain.WebSocket.send_and_collect(ws_pid, payload, done_fn,
timeout: openai.receive_timeout
) do
{:ok, events} ->
# Find the completed/failed event and process its response
events
|> Enum.find(&match?(%{"type" => "response.completed"}, &1))
|> case do
%{"response" => response} ->
case do_process_response(openai, response) do
{:error, %LangChainError{} = reason} ->
{:error, reason}
result ->
Callbacks.fire(openai.callbacks, :on_llm_new_message, [result])
result
end
nil ->
# Check for failed response
events
|> Enum.find(&match?(%{"type" => "response.failed"}, &1))
|> case do
%{"type" => "response.failed"} = failed_event ->
do_process_response(openai, failed_event)
nil ->
{:error,
LangChainError.exception(
type: "unexpected_response",
message: "No completed or failed event received via WebSocket"
)}
end
end
{:error, reason} ->
{:error,
LangChainError.exception(
type: "websocket_error",
message: "WebSocket request failed: #{inspect(reason)}"
)}
end
true ->
callback_fn = fn event ->
case do_process_response(openai, event) do
:skip ->
:skip
result ->
Utils.fire_streamed_callback(openai, List.wrap(result))
result
end
end
case LangChain.WebSocket.send_and_stream(ws_pid, payload, callback_fn, done_fn,
timeout: openai.receive_timeout
) do
{:ok, results} ->
results = results |> Enum.reject(&(&1 == :skip)) |> List.flatten()
# Check if any result is an error and return the first one found
case Enum.find(results, &match?({:error, %LangChainError{}}, &1)) do
{:error, _} = error -> error
nil -> results
end
{:error, reason} ->
{:error,
LangChainError.exception(
type: "websocket_error",
message: "WebSocket streaming request failed: #{inspect(reason)}"
)}
end
end
end
end
def do_api_request(
%ChatOpenAIResponses{stream: false} = openai,
messages,
tools,
retry_count
) do
retry_count = retry_count || openai.retry_count + 1
raw_data = for_api(openai, messages, tools)
if openai.verbose_api do
IO.inspect(raw_data, label: "RAW DATA BEING SUBMITTED")
end
req =
Req.new(
url: openai.endpoint,
json: raw_data,
# required for OpenAI API
auth: {:bearer, get_api_key(openai)},
# required for Azure OpenAI version
headers: [
{"api-key", get_api_key(openai)}
],
receive_timeout: openai.receive_timeout,
# Disable Req-level retry to prevent compounding with LangChain's own
# :closed retry. See https://github.com/brainlid/langchain/issues/503
retry: false
)
req
|> maybe_add_org_id_header()
|> maybe_add_proj_id_header()
|> Req.merge(openai.req_config |> Keyword.new())
|> Req.post()
# parse the body and return it as parsed structs
|> case do
{:ok, %Req.Response{body: data} = response} ->
if openai.verbose_api do
IO.inspect(response, label: "RAW REQ RESPONSE")
end
Callbacks.fire(openai.callbacks, :on_llm_ratelimit_info, [
get_ratelimit_info(response.headers)
])
case do_process_response(openai, data) do
{:error, %LangChainError{} = reason} ->
{:error, reason}
result ->
Callbacks.fire(openai.callbacks, :on_llm_new_message, [result])
# Track non-streaming response completion
LangChain.Telemetry.emit_event(
[:langchain, :llm, :response, :non_streaming],
%{system_time: System.system_time()},
%{
model: openai.model,
response_size: byte_size(inspect(result))
}
)
result
end
{:error, %Req.TransportError{reason: :timeout} = err} ->
{:error,
LangChainError.exception(type: "timeout", message: "Request timed out", original: err)}
{:error, %Req.TransportError{reason: :closed}} ->
# Force a retry by making a recursive call decrementing the counter
Logger.debug(fn -> "Mint connection closed: retry count = #{inspect(retry_count)}" end)
do_api_request(openai, messages, tools, retry_count - 1)
other ->
Logger.warning(fn -> "Unexpected and unhandled API response! #{inspect(other)}" end)
other
end
end
def do_api_request(
%ChatOpenAIResponses{stream: true} = openai,
messages,
tools,
retry_count
) do
retry_count = retry_count || openai.retry_count + 1
Req.new(
url: openai.endpoint,
json: for_api(openai, messages, tools),
# required for OpenAI API
auth: {:bearer, get_api_key(openai)},
# required for Azure OpenAI version
headers: [
{"api-key", get_api_key(openai)}
],
receive_timeout: openai.receive_timeout
)
|> maybe_add_org_id_header()
|> maybe_add_proj_id_header()
|> Req.merge(openai.req_config |> Keyword.new())
|> Req.post(
into: Utils.handle_stream_fn(openai, &decode_stream/1, &do_process_response(openai, &1))
)
|> case do
{:ok, %Req.Response{body: {:error, %LangChainError{} = error}}} ->
{:error, error}
{:ok, %Req.Response{body: data} = response} ->
Callbacks.fire(openai.callbacks, :on_llm_ratelimit_info, [
get_ratelimit_info(response.headers)
])
List.flatten(data)
{:error, %LangChainError{} = error} ->
{:error, error}
{:error, %Req.TransportError{reason: :timeout} = err} ->
{:error,
LangChainError.exception(type: "timeout", message: "Request timed out", original: err)}
{:error, %Req.TransportError{reason: :closed}} ->
# Force a retry by making a recursive call decrementing the counter
Logger.debug(fn -> "Connection closed: retry count = #{inspect(retry_count)}" end)
do_api_request(openai, messages, tools, retry_count - 1)
other ->
Logger.warning(fn ->
"Unhandled and unexpected response from streamed post call. #{inspect(other)}"
end)
{:error,
LangChainError.exception(
type: "unexpected_response",
message: "Unexpected response",
original: other
)}
end
end
# The Responses API streams events in the form:
# event: <event_type>\ndata: { ...json... }
# We want to extract each pair and parse the JSON from the `data:` line.
# A list of all events can be found here: https://platform.openai.com/docs/api-reference/responses-streaming
# Unlike the Chat Completions API, we do not get a [DONE] token at the end of the stream.
@spec decode_stream({String.t(), String.t()}) :: {[map()], String.t()}
def decode_stream({raw_data, buffer}) do
combined = buffer <> raw_data
segments = String.split(combined, "\n\n")
{complete_segments, [maybe_incomplete]} = Enum.split(segments, -1)
parsed =
complete_segments
|> Enum.flat_map(fn segment ->
segment
|> String.split("\n")
|> Enum.find_value(fn
"data: " <> json -> json
"data:" <> json -> String.trim_leading(json)
_ -> nil
end)
|> case do
nil ->
[]
json ->
case Jason.decode(json) do
{:ok, parsed} -> [parsed]
{:error, _} -> []
end
end
end)
{parsed, maybe_incomplete}
end
# Parse a new message response
@doc false
@spec do_process_response(t(), data :: any()) ::
:skip
| TokenUsage.t()
| Message.t()
| [Message.t() | MessageDelta.t() | TokenUsage.t() | {:error, LangChainError.t()}]
| MessageDelta.t()
| [MessageDelta.t()]
| {:error, LangChainError.t()}
# Complete Response with output lists
def do_process_response(
_model,
%{"status" => "completed", "output" => content_items} = response
)
when is_list(content_items) do
{content_parts, tool_calls} = content_items_to_content_parts_and_tool_calls(content_items)
metadata =
case get_token_usage(response) do
nil -> %{}
%TokenUsage{} = usage -> %{usage: usage}
end
|> maybe_add_response_id(response)
Message.new!(%{
content: content_parts,
status: :complete,
role: :assistant,
tool_calls: tool_calls,
metadata: metadata
})
end
# Handle streaming events
# Streamed events are returned as a raw list of events
# Even if there is only one event, it is returned within a list.
# Open to feedback this should get moved up and down the pattern-matching
# priority here.
def do_process_response(model, list) when is_list(list) do
Enum.map(list, &do_process_response(model, &1))
end
# Deltas arrive in the following shape:
# %{
# "content_index" => 0,
# "delta" => "Hello",
# "item_id" => "msg_1234567890",
# "output_index" => 0,
# "sequence_number" => 4,
# "type" => "response.output_text.delta"
# }
def do_process_response(_model, %{
"type" => "response.reasoning.delta",
"output_index" => output_index,
"delta" => delta_text
}) do
data = %{
content: ContentPart.new!(%{type: :thinking, content: delta_text}),
status: :incomplete,
role: :assistant,
index: output_index
}
case MessageDelta.new(data) do
{:ok, message} ->
message
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
def do_process_response(
_model,
%{
"type" => "response.output_text.delta",
"output_index" => output_index,
"delta" => delta_text
}
) do
data = %{
content: delta_text,
# Will need to be updated to :complete when the response is complete
status: :incomplete,
role: :assistant,
index: output_index
}
case MessageDelta.new(data) do
{:ok, message} ->
message
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
def do_process_response(_model, %{"type" => "response.output_text.delta", "delta" => delta_text}) do
data = %{
content: delta_text,
# Will need to be updated to :complete when the response is complete
status: :incomplete,
role: :assistant
}
case MessageDelta.new(data) do
{:ok, message} ->
message
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# Annotation events arrive during streaming when the model cites sources
# (e.g., from web search). Each event carries one annotation that we convert
# to a Citation and attach to a ContentPart with nil content, which gets
# merged into the accumulated delta via ContentPart.merge_part/2.
def do_process_response(_model, %{
"type" => "response.output_text.annotation.added",
"annotation" => annotation_data,
"output_index" => output_index
}) do
citation = parse_openai_annotation(annotation_data)
content_part = %ContentPart{type: :text, content: nil, citations: [citation]}
data = %{
content: content_part,
status: :incomplete,
role: :assistant,
index: output_index
}
case MessageDelta.new(data) do
{:ok, message} ->
message
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# Open question: is it possible we get multiples of `response.output_text.done`?
# It precedes `response.content_part.done` and `response.output_item.done`
# and theoretically we could get multiple text content_parts and output_items.
# I believe, semantically, these deltas are "outside" the output item and content part
# and can be treated as a "global" stream of deltas -- meaning "done" is truly
# "done" -- but that remains unconfirmed.
def do_process_response(_model, %{"type" => "response.output_text.done"}) do
data = %{
content: "",
status: :complete,
role: :assistant
}
case MessageDelta.new(data) do
{:ok, message} ->
message
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# This is the first event we get for a reasoning/thinking block.
# It is followed by a series of `response.reasoning.delta` events.
# Finally, it is followed by a `response.output_item.done` event.
def do_process_response(_model, %{
"type" => "response.output_item.added",
"output_index" => output_index,
"item" => %{
"type" => "reasoning",
"id" => _reasoning_id
}
}) do
data = %{
content: ContentPart.new!(%{type: :thinking, content: ""}),
status: :incomplete,
role: :assistant,
index: output_index
}
case MessageDelta.new(data) do
{:ok, delta} ->
delta
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# This is the first event we get for a function call.
# It is followed by a series of `response.function_call_arguments.delta` events.
# It is followed by a `response.function_call_arguments.done` event. (which we skip)
# Finally, it is followed by a `response.output_item.done` event.
def do_process_response(_model, %{
"type" => "response.output_item.added",
"output_index" => output_index,
"item" => %{
"type" => "function_call",
"call_id" => call_id,
"name" => name,
"arguments" => args
}
}) do
data = %{
status: :incomplete,
type: :function,
call_id: call_id,
name: name,
arguments: args,
index: output_index
}
with {:ok, %ToolCall{} = call} <- ToolCall.new(data),
{:ok, delta} <-
MessageDelta.new(%{
content: "",
status: :incomplete,
role: :assistant,
tool_calls: [call]
}) do
delta
else
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
def do_process_response(_model, %{
"type" => "response.function_call_arguments.delta",
"output_index" => output_index,
"delta" => delta_text
}) do
data = %{
arguments: delta_text,
index: output_index
}
with {:ok, call} <- ToolCall.new(data),
{:ok, message} <-
MessageDelta.new(%{
content: "",
status: :incomplete,
role: :assistant,
tool_calls: [call]
}) do
message
else
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
def do_process_response(_model, %{
"type" => "response.output_item.done",
"output_index" => output_index,
"item" => %{"type" => "reasoning"}
}) do
data = %{
content: ContentPart.new!(%{type: :thinking, content: ""}),
status: :complete,
role: :assistant,
index: output_index
}
case MessageDelta.new(data) do
{:ok, delta} ->
delta
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
def do_process_response(_model, %{
"type" => "response.output_item.done",
"output_index" => output_index,
"item" => %{"type" => "function_call"} = item
}) do
data = %{
status: :complete,
index: output_index,
call_id: item["call_id"],
arguments: item["arguments"],
name: item["name"]
}
with {:ok, call} <- ToolCall.new(data),
{:ok, message} <-
MessageDelta.new(%{
status: :complete,
role: :assistant,
tool_calls: [call]
}) do
message
else
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
def do_process_response(_model, %{
"type" => "response.completed",
"response" => response
}) do
usage = get_token_usage(response)
metadata = %{usage: usage} |> maybe_add_response_id(response)
data = %{
content: "",
status: :complete,
role: :assistant,
metadata: metadata
}
case MessageDelta.new(data) do
{:ok, message} ->
message
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# Streaming events explicitly skipped
# Items we should come back and implement:
# - refusals
# - function_calls
# - error
@reasoning_summary_events [
"response.reasoning_summary_part.added",
"response.reasoning_summary_part.done",
"response.reasoning_summary_text.delta",
"response.reasoning_summary_text.done",
"response.reasoning_summary.delta",
"response.reasoning_summary.done"
]
@skippable_streaming_events [
"keepalive",
"response.created",
"response.in_progress",
"response.incomplete",
"response.output_item.added",
"response.output_item.done",
"response.content_part.added",
"response.content_part.done",
"response.refusal.delta",
"response.refusal.done",
"response.function_call_arguments.done",
"response.file_search_call.in_progress",
"response.file_search_call.searching",
"response.file_search_call.completed",
"response.web_search_call.in_progress",
"response.web_search_call.searching",
"response.web_search_call.completed",
"response.image_generation_call.completed",
"response.image_generation_call.generating",
"response.image_generation_call.in_progress",
"response.image_generation_call.partial_image",
"response.mcp_call.arguments.delta",
"response.mcp_call.arguments.done",
"response.mcp_call.completed",
"response.mcp_call.failed",
"response.mcp_call.in_progress",
"response.queued",
"error"
]
# Handle reasoning summary delta events - fire callback and return :skip
def do_process_response(model, %{
"type" => "response.reasoning_summary_text.delta",
"delta" => delta
}) do
if model.verbose_api do
Logger.debug("[LANGCHAIN] Reasoning text delta received")
end
Callbacks.fire(model.callbacks, :on_llm_reasoning_delta, [delta])
:skip
end
def do_process_response(model, %{
"type" => "response.reasoning_summary.delta",
"delta" => delta
}) do
if model.verbose_api do
Logger.debug("[LANGCHAIN] Reasoning summary delta received")
end
Callbacks.fire(model.callbacks, :on_llm_reasoning_delta, [delta])
:skip
end
def do_process_response(model, %{"type" => event} = _data)
when event in @reasoning_summary_events do
if model.verbose_api do
Logger.debug("[LANGCHAIN] Reasoning event: #{event}")
end
:skip
end
def do_process_response(model, %{"type" => event})
when event in @skippable_streaming_events do
if model.verbose_api do
Logger.debug("[LANGCHAIN] Skipping streaming event: #{event}")
end
:skip
end
def do_process_response(_model, %{"error" => %{"message" => reason}}) do
{:error, LangChainError.exception(message: reason)}
end
# Handle failed response status (streaming event with type "response.failed")
def do_process_response(_model, %{
"type" => "response.failed",
"response" => %{"status" => "failed"} = response
}) do
build_failed_response_error(response)
end
# Handle failed response status (non-streaming / full response object)
def do_process_response(_model, %{"response" => %{"status" => "failed"} = response}) do
build_failed_response_error(response)
end
def do_process_response(_model, {:error, %Jason.DecodeError{} = response}) do
error_message = "Received invalid JSON: #{inspect(response)}"
{:error,
LangChainError.exception(type: "invalid_json", message: error_message, original: response)}
end
def do_process_response(_model, other) do
{:error,
LangChainError.exception(
type: "unexpected_response",
message: "Unexpected response",
original: other
)}
end
# Extracts error details from a failed OpenAI response and builds an error tuple.
# Handles various error formats defensively:
# - %{"error" => %{"code" => "...", "message" => "..."}}
# - %{"error" => "string message"}
# - %{} (no error details)
@spec build_failed_response_error(map()) :: {:error, LangChainError.t()}
defp build_failed_response_error(response) do
{error_type, error_message} = extract_error_details(response)
{:error,
LangChainError.exception(
type: error_type,
message: "OpenAI request failed: #{error_message}",
original: response
)}
end
@spec extract_error_details(map()) :: {String.t(), String.t()}
defp extract_error_details(response) do
case Map.get(response, "error") do
%{} = error_info ->
message = Map.get(error_info, "message", "Request failed")
code = Map.get(error_info, "code", "api_error")
{code, message}
message when is_binary(message) ->
{"api_error", message}
_ ->
{"api_error", "Request failed"}
end
end
defp get_token_usage(%{"usage" => usage} = _response_body) when is_map(usage) do
# extract out the reported response token usage
#
# https://platform.openai.com/docs/api-reference/responses_streaming/response/completed#responses_streaming/response/completed-response-usage
TokenUsage.new!(%{
input: Map.get(usage, "input_tokens"),
output: Map.get(usage, "output_tokens"),
raw: usage
})
end
defp get_token_usage(_response_body), do: nil
defp maybe_add_response_id(metadata, %{"id" => id}) when is_binary(id) do
Map.put(metadata, :response_id, id)
end
defp maybe_add_response_id(metadata, _response), do: metadata
defp content_items_to_content_parts_and_tool_calls(content_items) do
Enum.reduce(content_items, {[], []}, fn content_item, {content_parts, tool_calls} ->
case content_item_to_content_part_or_tool_call(content_item) do
%ContentPart{} = cp ->
{content_parts ++ [cp], tool_calls}
%ToolCall{} = tc ->
{content_parts, tool_calls ++ [tc]}
end
end)
end
defp content_item_to_content_part_or_tool_call(%{
"type" => "message",
"content" => message_contents
}) do
{text_parts, all_citations} =
Enum.reduce(message_contents, {[], []}, fn
%{"type" => "output_text", "text" => text} = item, {texts, citations} ->
new_citations = parse_openai_annotations(Map.get(item, "annotations", []))
{texts ++ [text], citations ++ new_citations}
%{"type" => "refusal", "refusal" => refusal}, {texts, citations} ->
{texts ++ [refusal], citations}
end)
text = Enum.join(text_parts, " ")
part = ContentPart.text!(text)
%{part | citations: all_citations}
end
defp content_item_to_content_part_or_tool_call(%{
"type" => "function_call",
"call_id" => call_id,
"name" => name,
"arguments" => args
}) do
case ToolCall.new(%{
type: :function,
status: :complete,
name: name,
arguments: args,
call_id: call_id
}) do
{:ok, %ToolCall{} = call} ->
call
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# The Responses API returns web_search_call as a sibling of assistant messages, as
# in:
# %{
# ...,
# "output" => [
# %{"type" => "web_search_call", ...},
# %{"type" => "message", "content" => [...content_parts...]}
# ]
# }
# however we embed it within the message as an unsupported content part to maintain the
# idiom of returning a single %Message{} per API call.
defp content_item_to_content_part_or_tool_call(%{
"type" => "web_search_call",
"id" => web_search_call_id,
"status" => "completed"
}) do
case ContentPart.new(%{
type: :unsupported,
options: %{
id: web_search_call_id,
status: "completed",
type: "web_search_call"
},
call_id: web_search_call_id
}) do
{:ok, %ContentPart{} = call} ->
call
{:error, %Ecto.Changeset{} = changeset} ->
{:error, LangChainError.exception(changeset)}
end
end
# Handle reasoning output from gpt-5 and o-series models
# We can either ignore it or store it as metadata
defp content_item_to_content_part_or_tool_call(%{
"type" => "reasoning",
"id" => reasoning_id,
"summary" => summary
}) do
# Store reasoning as an unsupported content part for now
# This preserves the information without breaking the flow
case ContentPart.new(%{
type: :unsupported,
options: %{
id: reasoning_id,
summary: summary,
type: "reasoning"
}
}) do
{:ok, %ContentPart{} = part} ->
part
{:error, %Ecto.Changeset{} = changeset} ->
reason = Utils.changeset_error_to_string(changeset)
Logger.warning("Failed to process reasoning output. Reason: #{reason}")
# Return a minimal content part to avoid breaking the flow
ContentPart.text!("")
end
end
defp content_item_to_content_part_or_tool_call(
%{
"type" => "file_search_call"
} = part
) do
# Store reasoning as an unsupported content part for now
# This preserves the information without breaking the flow
case ContentPart.new(%{
type: :unsupported,
options: %{
id: part["id"],
type: "file_search_call",
queries: part["queries"],
results: part["results"]
}
}) do
{:ok, %ContentPart{} = part} ->
part
{:error, %Ecto.Changeset{} = changeset} ->
reason = Utils.changeset_error_to_string(changeset)
Logger.warning("Failed to process file_search output. Reason: #{reason}")
# Return a minimal content part to avoid breaking the flow
ContentPart.text!("")
end
end
# Catch-all for unknown content item types
defp content_item_to_content_part_or_tool_call(%{"type" => type} = item) do
Logger.warning("Unknown content item type: #{type}. Item: #{inspect(item)}")
# Return empty text content to avoid breaking the flow
ContentPart.text!("")
end
# -- OpenAI annotation parsing helpers --
defp parse_openai_annotations(annotations) when is_list(annotations) do
Enum.map(annotations, &parse_openai_annotation/1)
end
defp parse_openai_annotations(_), do: []
defp parse_openai_annotation(%{"type" => "url_citation"} = ann) do
Citation.new!(%{
source: %{
type: :web,
title: Map.get(ann, "title"),
url: Map.get(ann, "url")
},
start_index: Map.get(ann, "start_index"),
end_index: Map.get(ann, "end_index"),
metadata: %{"provider_type" => "url_citation"}
})
end
defp parse_openai_annotation(%{"type" => "file_citation"} = ann) do
Citation.new!(%{
source: %{
type: :document,
document_id: Map.get(ann, "file_id"),
title: Map.get(ann, "filename"),
metadata: %{"filename" => Map.get(ann, "filename")}
},
start_index: Map.get(ann, "start_index"),
end_index: Map.get(ann, "end_index"),
metadata: %{"provider_type" => "file_citation"}
})
end
defp parse_openai_annotation(%{"type" => type} = ann) do
Logger.warning("Unknown OpenAI annotation type: #{type}")
Citation.new!(%{
metadata: %{"provider_type" => type, "raw" => ann}
})
end
defp maybe_add_org_id_header(%Req.Request{} = req) do
org_id = get_org_id()
if org_id do
Req.Request.put_header(req, "OpenAI-Organization", org_id)
else
req
end
end
defp maybe_add_proj_id_header(%Req.Request{} = req) do
proj_id = get_proj_id()
if proj_id do
Req.Request.put_header(req, "OpenAI-Project", proj_id)
else
req
end
end
defp get_ratelimit_info(response_headers) do
# extract out all the ratelimit response headers
#
# https://platform.openai.com/docs/guides/rate-limits/rate-limits-in-headers
{return, _} =
Map.split(response_headers, [
"x-ratelimit-limit-requests",
"x-ratelimit-limit-tokens",
"x-ratelimit-remaining-requests",
"x-ratelimit-remaining-tokens",
"x-ratelimit-reset-requests",
"x-ratelimit-reset-tokens",
"x-request-id"
])
return
end
@doc """
Generate a config map that can later restore the model's configuration.
"""
@impl ChatModel
@spec serialize_config(t()) :: %{String.t() => any()}
def serialize_config(%ChatOpenAIResponses{} = model) do
Utils.to_serializable_map(
model,
[
:endpoint,
:model,
:temperature,
:frequency_penalty,
:reasoning,
:receive_timeout,
:seed,
:n,
:json_response,
:json_schema,
:stream,
:max_tokens,
:stream_options,
:verbosity
],
@current_config_version
)
end
@doc """
Restores the model from the config.
"""
@impl ChatModel
def restore_from_map(%{"version" => 1} = data) do
ChatOpenAIResponses.new(data)
end
@impl ChatModel
def provider, do: "openai_responses"
@doc """
Determine if an error should be retried with a fallback model.
Aligns with other providers.
"""
@impl ChatModel
@spec retry_on_fallback?(LangChainError.t()) :: boolean()
def retry_on_fallback?(%LangChainError{type: "rate_limited"}), do: true
def retry_on_fallback?(%LangChainError{type: "rate_limit_exceeded"}), do: true
def retry_on_fallback?(%LangChainError{type: "timeout"}), do: true
def retry_on_fallback?(%LangChainError{type: "too_many_requests"}), do: true
def retry_on_fallback?(_), do: false
end