use crate::completion::{self, ProviderCapabilities};
use crate::error::EncodeError;
use crate::json_utils::Lenient;
use crate::observe::ObservedError;
use crate::operation::Completion;
use crate::providers::openai::wire::OpenAIConfig;
pub(crate) use crate::providers::openai::wire::ResponsesContract;
use crate::wire::{
AdapterEvent, AdapterUsage, AdapterVerdict, Capabilities, Descriptor, Encoded, Framing, Mode,
ObservationSink, Wire,
};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use super::streaming::{ResponsesDecoder, usage_of};
use super::{ResponsesToolDefinition, SystemInstructionsPlacement};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Responses {
pub provider: OpenAIConfig,
pub model: String,
pub tools: Vec<ResponsesToolDefinition>,
pub strict_tools: bool,
pub system_instructions: SystemInstructionsPlacement,
}
impl Responses {
pub(crate) fn encode_with_headers(
&self,
request: completion::CompletionRequest,
mode: Mode,
headers: impl FnOnce(
&OpenAIConfig,
&completion::CompletionRequest,
http::request::Builder,
) -> http::request::Builder,
) -> Result<Encoded, EncodeError> {
let quirks = &self.provider.dialect.quirks.responses;
let codex = quirks.contract == ResponsesContract::Codex;
let streaming = matches!(mode, Mode::Streaming) || codex;
let builder = headers(
&self.provider,
&request,
http::Request::post(self.provider.uri(quirks.path, None)),
);
let mode = if streaming {
Mode::Streaming
} else {
Mode::Unary
};
let body = self.responses_request(&request, super::Delivery::Http(mode))?;
crate::providers::internal::trace_json(
crate::providers::internal::LogTarget::Completions,
"Responses completion request",
&body,
);
let request = builder
.header(http::header::CONTENT_TYPE, "application/json")
.body(body.into_body())?;
let framing = if streaming {
Framing::Sse
} else {
Framing::Whole
};
let encoded = Encoded::new(request, framing)
.with_request_id_header(self.provider.dialect.request_id_header)
.with_route(Some(self.provider.dialect.quirks.responses.path))
.with_projection(project_payload);
Ok(if codex {
encoded.with_relaxed_content_type()
} else {
encoded
})
}
pub fn new(provider: OpenAIConfig, model: impl Into<String>) -> Self {
Self {
system_instructions: provider.system_instructions_placement(),
strict_tools: provider.dialect.quirks.responses.strict_tools_by_default,
provider,
model: model.into(),
tools: Vec::new(),
}
}
pub fn with_strict_tools(mut self) -> Self {
self.strict_tools = true;
self
}
pub fn with_tool(mut self, tool: impl Into<ResponsesToolDefinition>) -> Self {
self.tools.push(tool.into());
self
}
pub fn with_tools<I, Tool>(mut self, tools: I) -> Self
where
I: IntoIterator<Item = Tool>,
Tool: Into<ResponsesToolDefinition>,
{
self.tools.extend(tools.into_iter().map(Into::into));
self
}
pub fn with_system_instructions_placement(
mut self,
placement: SystemInstructionsPlacement,
) -> Self {
self.system_instructions = placement;
self
}
pub fn with_system_instructions_as_messages(self) -> Self {
self.with_system_instructions_placement(SystemInstructionsPlacement::InputSystemMessages)
}
}
impl Wire for Responses {
type Op = Completion;
type Payload = crate::wire::Encoded;
type Frame = crate::wire::WireFrame;
type Decoder<'id> = ResponsesDecoder;
type Reassembler = super::streaming::document::Response;
fn describe(&self) -> Descriptor<'_> {
Descriptor::new(self.provider.dialect.name)
.model(self.model.as_str())
.capabilities(Capabilities::completion(
ProviderCapabilities::default().with_native_output_tool_composition(
self.provider.dialect.quirks.responses.contract != ResponsesContract::Xai,
),
))
.replay(self)
}
fn encode(
&self,
request: completion::CompletionRequest,
mode: Mode,
) -> Result<Encoded, EncodeError> {
self.encode_with_headers(request, mode, OpenAIConfig::completion_headers)
}
fn decoder<'id>(&self) -> Self::Decoder<'id> {
ResponsesDecoder::new()
}
}
impl crate::completion::ReplayTarget for Responses {
fn map_options(
&self,
request: &crate::completion::CompletionRequest,
fields: crate::completion::options::OptionFields<'_>,
) -> crate::completion::options::OptionMap {
crate::providers::openai::options::responses_options(self, request, fields)
}
fn api(&self) -> crate::message::Api {
crate::message::Api::from_static("openai.responses")
}
fn provider(&self) -> &str {
self.provider.dialect.name
}
fn model(&self) -> &str {
&self.model
}
fn declares_tools(&self, request: &completion::CompletionRequest) -> bool {
!self.tools.is_empty() || crate::completion::history::declares_tools(request)
}
fn call_id_slot(&self) -> Option<&'static str> {
Some("/call_id")
}
fn identity(&self, item: &Value) -> Map<String, Value> {
let keys: &[&str] = match item.str("type") {
Some("message") => &["type", "id", "phase"],
Some("reasoning") => &["type", "id", "encrypted_content"],
Some("function_call" | "custom_tool_call") => &["type", "id"],
_ => &[],
};
keys.iter()
.filter_map(|key| Some(((*key).to_owned(), item.get(*key)?.clone())))
.collect()
}
fn needs_next(&self, item: &Value) -> bool {
item.str("type") == Some("reasoning")
}
fn accepts(&self, model: &str) -> crate::completion::Accepts {
let images = reads_images(self.provider.dialect.quirks.responses.contract, model);
let model = model.to_ascii_lowercase();
crate::completion::Accepts {
user_images: images,
assistant_images: false,
tool_result_images: images,
tools: !(model.starts_with("o1-mini") || model.starts_with("o1-preview")),
}
}
fn encodes(&self, _model: &str, media: crate::completion::Media<'_>) -> bool {
use crate::completion::{Media, Place};
use crate::message::DocumentSourceKind;
let quirks = &self.provider.dialect.quirks;
let file_ids = quirks.accepts_file_ids;
match media {
Media::Image(_, Place::Assistant) | Media::Audio(_) | Media::Video(_) => false,
Media::Image(image, _) => {
super::image_part(image).is_some()
&& (!matches!(image.data, DocumentSourceKind::FileId(_))
|| file_ids && quirks.responses.contract != ResponsesContract::Xai)
}
Media::Document(document) => {
super::document_part(document).is_some()
&& (file_ids || !matches!(document.data, DocumentSourceKind::FileId(_)))
}
}
}
fn continues_stored(&self, request: &completion::CompletionRequest) -> bool {
["previous_response_id", "conversation"].iter().any(|key| {
crate::completion::options::param(self, request, key)
.is_some_and(|value| !value.is_null())
})
}
fn normalize_tool_call_id(
&self,
id: &str,
_model: &str,
_source: Option<&crate::message::Origin>,
) -> String {
use crate::providers::internal::wire_ids::{legal_call_id, short_hash};
let legal = legal_call_id(id, 64);
match legal.trim_end_matches('_') {
"" => short_hash(id),
trimmed => trimmed.to_owned(),
}
}
}
fn reads_images(contract: ResponsesContract, model: &str) -> bool {
let model = model.rsplit('/').next().unwrap_or_default();
match contract {
ResponsesContract::Xai => crate::catalog::reads_images_or(
crate::providers::xai::DIALECT.name,
model,
crate::providers::xai::reads_images,
),
ResponsesContract::OpenAi | ResponsesContract::Codex => crate::catalog::reads_images_or(
crate::providers::openai::wire::OPENAI.name,
model,
crate::providers::openai::reads_images,
),
}
}
pub(crate) fn project_payload(payload: &[u8], sink: &mut ObservationSink<'_>) {
let Ok(payload) = serde_json::from_slice::<Value>(payload) else {
return;
};
let envelope = |error: &Value| ObservedError {
code: error.get("code").filter(|code| !code.is_null()).cloned(),
kind: error
.str("type")
.or_else(|| error.str("status"))
.map(str::to_owned),
message: error.str("message").map(str::to_owned),
};
if payload.str("type") == Some("error") {
match payload.get("error").filter(|error| error.is_object()) {
Some(error) => envelope(error),
None => ObservedError {
kind: None,
..envelope(&payload)
},
}
.emit(sink);
return;
}
let object = payload
.get("response")
.filter(|response| response.is_object())
.unwrap_or(&payload);
if let Some(usage) = object.get("usage").filter(|usage| usage.is_object()) {
let usage = usage_of(usage);
sink.emit(AdapterEvent::Usage {
usage: AdapterUsage {
input_tokens: usage.input_tokens,
output_tokens: usage.output_tokens,
total_tokens: usage.total_tokens,
cached_input_tokens: usage.cached_input_tokens,
reasoning_tokens: usage.reasoning_tokens,
tool_input_tokens: None,
},
});
}
let verdict = AdapterVerdict {
finish_reason: object
.str("status")
.filter(|status| *status != "in_progress" && *status != "queued")
.map(|value| sink.scrub(value)),
block_reason: None,
detail: object
.at("/incomplete_details/reason")
.and_then(Value::as_str)
.map(|value| sink.scrub(value)),
model: object.str("model").map(|value| sink.scrub(value)),
};
let response_id = object.str("id").map(|value| sink.scrub(value));
sink.provider(verdict, response_id);
if let Some(error) = object.get("error").filter(|error| error.is_object()) {
envelope(error).emit(sink);
}
}
#[cfg(test)]
mod tests;