use crate::completion::history::Replay;
use crate::completion::options::{BaseInput, FinalBody, RawAt, Rewrite, request_params};
use crate::error::EncodeError;
use crate::json_utils;
use crate::json_utils::Lenient;
use crate::message::{
AssistantContent, Document, DocumentMediaType, DocumentSourceKind, Message, MimeType,
ToolResultContent, UserContent,
};
use crate::providers::internal::wire_ids::WireIds;
use crate::wire::Mode;
use crate::{completion, message};
use serde::{Deserialize, Serialize, Serializer};
use serde_json::{Map, Value, json};
pub mod streaming;
#[cfg(feature = "websocket")]
#[cfg_attr(docsrs, doc(cfg(feature = "websocket")))]
pub mod websocket;
pub mod wire;
fn image_part(image: &message::Image) -> Option<Value> {
let (key, source) = match &image.data {
DocumentSourceKind::Base64(data) => (
"image_url",
format!(
"data:{};base64,{data}",
image.media_type.as_ref()?.to_mime_type()
),
),
DocumentSourceKind::Url(url) => ("image_url", url.clone()),
DocumentSourceKind::FileId(file_id) => ("file_id", file_id.clone()),
_ => return None,
};
let mut part =
json!({"type": "input_image", "detail": image.detail.clone().unwrap_or_default()});
set(&mut part, key, Value::String(source));
Some(part)
}
fn document_part(document: &Document) -> Option<Value> {
Some(match &document.data {
DocumentSourceKind::FileId(file_id) => json!({"type": "input_file", "file_id": file_id}),
DocumentSourceKind::Url(url) => json!({"type": "input_file", "file_url": url}),
DocumentSourceKind::Base64(data) if document.media_type == Some(DocumentMediaType::PDF) => {
json!({
"type": "input_file",
"file_data": format!("data:application/pdf;base64,{data}"),
"filename": "document.pdf",
})
}
DocumentSourceKind::String(text) => json!({"type": "input_text", "text": text}),
_ => return None,
})
}
fn set(value: &mut Value, key: &str, field: Value) {
if let Some(fields) = value.as_object_mut() {
fields.insert(key.to_owned(), field);
}
}
fn unsendable(part: &str) -> EncodeError {
EncodeError::request(format!(
"the Responses API cannot carry this {part}; prepare the request first"
))
}
fn result_output(content: &[ToolResultContent]) -> Result<Value, EncodeError> {
let mut parts = content
.iter()
.map(|part| match part {
ToolResultContent::Text(text) => Ok(json!({"type": "input_text", "text": text.text})),
ToolResultContent::Json { value } => {
Ok(json!({"type": "input_text", "text": value.to_string()}))
}
ToolResultContent::Image(image) => {
image_part(image).ok_or_else(|| unsendable("tool-result image"))
}
})
.collect::<Result<Vec<_>, _>>()?;
Ok(match parts.as_mut_slice() {
[part] if part.str("type") == Some("input_text") => {
part.get_mut("text").map(Value::take).unwrap_or_default()
}
_ => Value::Array(parts),
})
}
struct Custom {
tools: std::collections::HashSet<String>,
calls: std::collections::HashMap<String, bool>,
}
fn input(
history: &[Message],
target: &wire::Responses,
model: &str,
custom: &mut Custom,
stateless: bool,
) -> Result<Vec<Value>, EncodeError> {
let ids = WireIds::for_target(history, target, model);
let mut items = Vec::new();
for (position, message) in history.iter().enumerate() {
match message {
Message::System { content } => items.push(json!({
"type": "message",
"role": "system",
"content": [{"type": "input_text", "text": content}],
})),
Message::User { content } => {
for part in content {
let part = match part {
UserContent::Text(text) if text.text.trim().is_empty() => continue,
UserContent::Text(text) => json!({"type": "input_text", "text": text.text}),
UserContent::ToolResult(result) => {
let call_id = ids.spell(&result.call);
let output = result_output(&result.content)?;
items.push(
if custom.calls.get(&call_id).copied().unwrap_or_else(|| {
custom.tools.contains(result.name.as_str())
}) {
json!({"type": "custom_tool_call_output", "call_id": call_id, "output": output})
} else {
json!({"type": "function_call_output", "call_id": call_id, "output": output, "status": "completed"})
},
);
continue;
}
UserContent::Image(image) => {
image_part(image).ok_or_else(|| unsendable("image"))?
}
UserContent::Document(document) => {
document_part(document).ok_or_else(|| unsendable("document"))?
}
UserContent::Audio(_) => return Err(unsendable("audio")),
UserContent::Video(_) => return Err(unsendable("video")),
};
items.push(json!({"type": "message", "role": "user", "content": [part]}));
}
}
Message::Assistant(turn) => {
let mut texts = 0usize;
let mut unpaired = false;
for block in &turn.content {
let mut replay = block.replay(target, &ids);
if unpaired && !matches!(block, AssistantContent::Reasoning(_)) {
unpaired = false;
replay = Replay::Rebuild;
}
let ciphertext = |item: &Map<String, Value>| {
item.get("encrypted_content")
.and_then(Value::as_str)
.is_some_and(|cipher| !cipher.is_empty())
};
let unresolvable = stateless
&& matches!(block, AssistantContent::Reasoning(_))
&& match &replay {
Replay::Item(item) => !item.as_object().is_some_and(ciphertext),
Replay::Identity(identity) => !ciphertext(identity),
Replay::Rebuild => false,
};
if unresolvable {
unpaired = true;
continue;
}
let identity = match replay {
Replay::Item(item) => {
if let Some(call_id) = item.str("call_id") {
let kind = item.str("type");
if matches!(kind, Some("custom_tool_call" | "function_call")) {
custom.calls.insert(
call_id.to_owned(),
kind == Some("custom_tool_call"),
);
}
}
items.push(item.into_owned());
continue;
}
Replay::Identity(identity) => identity,
Replay::Rebuild => Map::new(),
};
let id = |prefix: &str| {
identity
.get("id")
.and_then(Value::as_str)
.filter(|id| !id.is_empty() && id.starts_with(prefix) && id.len() <= 64)
.map(str::to_owned)
};
let item = match block {
AssistantContent::Text(text) => {
let synthetic = match texts {
0 => format!("msg_rig_{position}"),
n => format!("msg_rig_{position}_{n}"),
};
texts += 1;
let mut item = json!({
"type": "message",
"role": "assistant",
"content": [{"type": "output_text", "text": text.text, "annotations": []}],
"status": "completed",
"id": id("").unwrap_or(synthetic),
});
if let Some(phase) =
identity.get("phase").filter(|phase| phase.is_string())
{
set(&mut item, "phase", phase.clone());
}
items.push(item);
continue;
}
AssistantContent::ToolCall(call) => {
let call_id = ids.spell(&call.id);
let name = call.function.name.as_str();
let kind = identity.get("type").and_then(Value::as_str);
let arguments = call.function.arguments_value().to_string();
let (kind, key, payload, prefix) = if kind == Some("custom_tool_call")
|| (kind.is_none() && custom.tools.contains(name))
{
custom.calls.insert(call_id.clone(), true);
let input =
call.function.arguments.get("input").and_then(Value::as_str);
let input = input.map_or(arguments, str::to_owned);
("custom_tool_call", "input", input, "ctc_")
} else {
custom.calls.insert(call_id.clone(), false);
("function_call", "arguments", arguments, "fc_")
};
let mut item = json!({"type": kind, "call_id": call_id, "name": name});
set(&mut item, key, Value::String(payload));
if let Some(id) = id(prefix) {
set(&mut item, "id", json!(id));
}
item
}
AssistantContent::Reasoning(reasoning) => {
let Some(rs) = id("") else {
continue;
};
let summary: Vec<Value> = (!reasoning.text.is_empty())
.then(|| json!({"type": "summary_text", "text": reasoning.text}))
.into_iter()
.collect();
let mut item =
json!({"type": "reasoning", "id": rs, "summary": summary});
if let Some(ciphertext) = identity.get("encrypted_content") {
set(&mut item, "encrypted_content", ciphertext.clone());
}
item
}
AssistantContent::Opaque(opaque) if opaque.replay => opaque.item.clone(),
AssistantContent::Opaque(_) => continue,
AssistantContent::Image(_) => return Err(unsendable("assistant image")),
};
items.push(item);
}
}
}
}
Ok(items)
}
fn tool_choice(choice: message::ToolChoice) -> Result<Value, EncodeError> {
Ok(match choice {
message::ToolChoice::Auto => json!("auto"),
message::ToolChoice::None => json!("none"),
message::ToolChoice::Required => json!("required"),
message::ToolChoice::Specific { function_names } => match function_names.as_slice() {
[] => {
return Err(EncodeError::request(
"ToolChoice::Specific requires at least one function name",
));
}
[name] => json!({"type": "function", "name": name}),
names => json!({
"type": "allowed_tools",
"mode": "required",
"tools": names.iter().map(|name| json!({"type": "function", "name": name})).collect::<Vec<_>>(),
}),
},
})
}
pub(crate) fn include_ciphertext(body: &mut Map<String, Value>) {
const CIPHERTEXT: &str = "reasoning.encrypted_content";
let mut include = body
.get("include")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
if !include.iter().any(|item| item == CIPHERTEXT) {
include.push(json!(CIPHERTEXT));
}
body.insert("include".to_owned(), Value::Array(include));
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum Delivery {
Http(Mode),
WebSocket,
}
impl wire::Responses {
pub(crate) fn responses_request(
&self,
request: &completion::CompletionRequest,
delivery: Delivery,
) -> Result<FinalBody, EncodeError> {
let codex =
self.provider.dialect.quirks.responses.contract == wire::ResponsesContract::Codex;
let mut rewrites = Vec::new();
match delivery {
Delivery::Http(Mode::Streaming) => rewrites.push(Rewrite::Stream(true)),
Delivery::Http(Mode::Unary) | Delivery::WebSocket => rewrites.push(Rewrite::NoStream),
}
if delivery == Delivery::WebSocket {
rewrites.push(Rewrite::NoBackground);
}
if codex {
rewrites.push(Rewrite::CodexStore);
}
rewrites.push(Rewrite::ReasoningCiphertext(codex));
let body = request_params(
self,
request,
|layers| self.base(request, codex, layers),
RawAt::Top,
&rewrites,
)?;
crate::providers::openai::options::check_body(
self,
request,
body,
crate::providers::openai::options::Endpoint::Responses,
)
}
fn base(
&self,
request: &completion::CompletionRequest,
codex: bool,
layers: &mut BaseInput<'_>,
) -> Result<Map<String, Value>, EncodeError> {
let model = request.model.clone().unwrap_or_else(|| self.model.clone());
let mut tools: Vec<ResponsesToolDefinition> = request
.tools
.iter()
.cloned()
.map(ResponsesToolDefinition::from)
.collect();
let extra = layers.raw_tools()?;
if !extra.is_empty() {
tools.extend(
serde_json::from_value::<Vec<ResponsesToolDefinition>>(Value::Array(extra))
.map_err(|err| {
EncodeError::request(format!(
"Invalid OpenAI Responses tools payload in additional_params: {err}"
))
})?,
);
}
tools.extend(self.tools.iter().cloned());
if self.strict_tools {
tools = tools
.into_iter()
.map(ResponsesToolDefinition::with_strict)
.collect();
}
let mut custom = Custom {
tools: tools
.iter()
.filter(|tool| tool.kind == "custom")
.map(|tool| tool.name.clone())
.collect(),
calls: Default::default(),
};
let stateless = codex || layers.param("store") == Some(&json!(false));
let mut items = input(&request.chat_history, self, &model, &mut custom, stateless)?;
let system = |item: &Value| {
(item.str("role") == Some("system")).then(|| {
item.at("/content/0/text")
.and_then(Value::as_str)
.unwrap_or_default()
.to_owned()
})
};
let before = items.len();
let mut lifted = Vec::new();
match self.system_instructions {
SystemInstructionsPlacement::Instructions => {
let leading = items
.iter()
.take_while(|item| system(item).is_some())
.count();
if leading < items.len() {
lifted.extend(items.drain(..leading).filter_map(|item| system(&item)));
}
}
SystemInstructionsPlacement::AllInstructions => {
items.retain(|item| match system(item) {
Some(text) => {
lifted.push(text);
false
}
None => true,
})
}
SystemInstructionsPlacement::InputSystemMessages => {}
}
if items.is_empty() {
return Err(EncodeError::request(if items.len() < before {
"OpenAI Responses request input must contain at least one non-system item \
(system messages were lifted into the top-level `instructions` field)"
} else {
"OpenAI Responses request input must contain at least one item"
}));
}
let lifted: Vec<&str> = lifted
.iter()
.map(|text| text.trim())
.filter(|text| !text.is_empty())
.collect();
let lifted = lifted.join("\n\n");
let instructions = match &self.provider.instructions {
Some(gateway) if lifted.is_empty() => Some(gateway.clone()),
Some(gateway) if !lifted.contains(gateway.as_str()) => {
Some(format!("{gateway}\n\n{lifted}"))
}
_ => (!lifted.is_empty()).then_some(lifted),
};
let text = request
.output_schema
.clone()
.filter(|_| !codex)
.map(|schema| {
let (name, schema) = super::structured_output_schema(schema);
json!({"format": {"type": "json_schema", "name": name, "schema": schema, "strict": true}})
});
let fields = [
("model", Some(Value::String(model))),
("input", Some(Value::Array(items))),
("instructions", instructions.map(Value::from)),
(
"max_output_tokens",
request.max_tokens.filter(|_| !codex).map(Value::from),
),
(
"temperature",
request.temperature.filter(|_| !codex).map(Value::from),
),
(
"tool_choice",
request.tool_choice.clone().map(tool_choice).transpose()?,
),
("tools", (!tools.is_empty()).then(|| json!(tools))),
("text", text),
];
Ok(fields
.into_iter()
.filter_map(|(key, value)| Some((key.to_owned(), value?)))
.collect())
}
}
#[derive(Debug, Deserialize, Clone, PartialEq)]
pub struct ResponsesToolDefinition {
#[serde(rename = "type")]
pub kind: String,
#[serde(default)]
pub name: String,
#[serde(default)]
pub parameters: serde_json::Value,
#[serde(default, deserialize_with = "json_utils::null_or_default")]
pub strict: bool,
#[serde(default, deserialize_with = "json_utils::null_or_default")]
pub description: String,
#[serde(flatten, default)]
pub config: Map<String, Value>,
}
impl Serialize for ResponsesToolDefinition {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
use serde::ser::SerializeMap;
let mut map = serializer.serialize_map(None)?;
map.serialize_entry("type", &self.kind)?;
if !self.name.is_empty() {
map.serialize_entry("name", &self.name)?;
}
if !self.parameters.is_null() {
map.serialize_entry("parameters", &self.parameters)?;
}
if self.kind == "function" {
map.serialize_entry("strict", &self.strict)?;
}
if !self.description.is_empty() {
map.serialize_entry("description", &self.description)?;
}
for (key, value) in &self.config {
map.serialize_entry(key, value)?;
}
map.end()
}
}
impl ResponsesToolDefinition {
pub fn function(
name: impl Into<String>,
description: impl Into<String>,
parameters: serde_json::Value,
) -> Self {
Self {
kind: "function".to_string(),
name: name.into(),
parameters,
strict: false,
description: description.into(),
config: Map::new(),
}
}
pub fn strict_function(
name: impl Into<String>,
description: impl Into<String>,
parameters: serde_json::Value,
) -> Self {
Self::function(name, description, parameters).with_strict()
}
pub fn with_strict(mut self) -> Self {
if self.kind == "function" {
super::sanitize_schema(&mut self.parameters);
self.strict = true;
}
self
}
pub fn hosted(kind: impl Into<String>) -> Self {
Self {
kind: kind.into(),
name: String::new(),
parameters: Value::Null,
strict: false,
description: String::new(),
config: Map::new(),
}
}
pub fn web_search() -> Self {
Self::hosted("web_search")
}
pub fn file_search() -> Self {
Self::hosted("file_search")
}
pub fn computer_use() -> Self {
Self::hosted("computer_use")
}
pub fn with_config(mut self, key: impl Into<String>, value: Value) -> Self {
self.config.insert(key.into(), value);
self
}
}
impl From<completion::ToolDefinition> for ResponsesToolDefinition {
fn from(value: completion::ToolDefinition) -> Self {
let completion::ToolDefinition {
name,
parameters,
description,
} = value;
Self::function(name, description, parameters)
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SystemInstructionsPlacement {
#[default]
Instructions,
AllInstructions,
InputSystemMessages,
}
#[cfg(test)]
mod history_tests;
#[cfg(test)]
mod tests;