rig_core/providers/openai/responses_api/
wire.rs1use crate::completion::{self, ProviderCapabilities};
9use crate::error::EncodeError;
10use crate::json_utils::Lenient;
11use crate::observe::ObservedError;
12use crate::operation::Completion;
13use crate::providers::openai::wire::OpenAIConfig;
14pub(crate) use crate::providers::openai::wire::ResponsesContract;
15use crate::wire::{
16 AdapterEvent, AdapterUsage, AdapterVerdict, Capabilities, Descriptor, Encoded, Framing, Mode,
17 ObservationSink, Wire,
18};
19use serde::{Deserialize, Serialize};
20use serde_json::{Map, Value};
21
22use super::streaming::{ResponsesDecoder, usage_of};
23use super::{ResponsesToolDefinition, SystemInstructionsPlacement};
24
25#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
27pub struct Responses {
28 pub provider: OpenAIConfig,
30 pub model: String,
32 pub tools: Vec<ResponsesToolDefinition>,
34 pub strict_tools: bool,
36 pub system_instructions: SystemInstructionsPlacement,
39}
40
41impl Responses {
42 pub(crate) fn encode_with_headers(
43 &self,
44 request: completion::CompletionRequest,
45 mode: Mode,
46 headers: impl FnOnce(
47 &OpenAIConfig,
48 &completion::CompletionRequest,
49 http::request::Builder,
50 ) -> http::request::Builder,
51 ) -> Result<Encoded, EncodeError> {
52 let quirks = &self.provider.dialect.quirks.responses;
53 let codex = quirks.contract == ResponsesContract::Codex;
58 let streaming = matches!(mode, Mode::Streaming) || codex;
59 let builder = headers(
60 &self.provider,
61 &request,
62 http::Request::post(self.provider.uri(quirks.path, None)),
63 );
64 let mode = if streaming {
65 Mode::Streaming
66 } else {
67 Mode::Unary
68 };
69 let body = self.responses_request(&request, super::Delivery::Http(mode))?;
70 crate::providers::internal::trace_json(
71 crate::providers::internal::LogTarget::Completions,
72 "Responses completion request",
73 &body,
74 );
75 let request = builder
76 .header(http::header::CONTENT_TYPE, "application/json")
77 .body(body.into_body())?;
78
79 let framing = if streaming {
80 Framing::Sse
81 } else {
82 Framing::Whole
83 };
84 let encoded = Encoded::new(request, framing)
85 .with_request_id_header(self.provider.dialect.request_id_header)
86 .with_route(Some(self.provider.dialect.quirks.responses.path))
87 .with_projection(project_payload);
88 Ok(if codex {
89 encoded.with_relaxed_content_type()
90 } else {
91 encoded
92 })
93 }
94
95 pub fn new(provider: OpenAIConfig, model: impl Into<String>) -> Self {
98 Self {
99 system_instructions: provider.system_instructions_placement(),
100 strict_tools: provider.dialect.quirks.responses.strict_tools_by_default,
101 provider,
102 model: model.into(),
103 tools: Vec::new(),
104 }
105 }
106
107 pub fn with_strict_tools(mut self) -> Self {
109 self.strict_tools = true;
110 self
111 }
112
113 pub fn with_tool(mut self, tool: impl Into<ResponsesToolDefinition>) -> Self {
115 self.tools.push(tool.into());
116 self
117 }
118
119 pub fn with_tools<I, Tool>(mut self, tools: I) -> Self
121 where
122 I: IntoIterator<Item = Tool>,
123 Tool: Into<ResponsesToolDefinition>,
124 {
125 self.tools.extend(tools.into_iter().map(Into::into));
126 self
127 }
128
129 pub fn with_system_instructions_placement(
132 mut self,
133 placement: SystemInstructionsPlacement,
134 ) -> Self {
135 self.system_instructions = placement;
136 self
137 }
138
139 pub fn with_system_instructions_as_messages(self) -> Self {
142 self.with_system_instructions_placement(SystemInstructionsPlacement::InputSystemMessages)
143 }
144}
145
146impl Wire for Responses {
147 type Op = Completion;
148 type Payload = crate::wire::Encoded;
149 type Frame = crate::wire::WireFrame;
150 type Decoder<'id> = ResponsesDecoder;
151 type Reassembler = super::streaming::document::Response;
152
153 fn describe(&self) -> Descriptor<'_> {
155 Descriptor::new(self.provider.dialect.name)
156 .model(self.model.as_str())
157 .capabilities(Capabilities::completion(
158 ProviderCapabilities::default().with_native_output_tool_composition(
159 self.provider.dialect.quirks.responses.contract != ResponsesContract::Xai,
160 ),
161 ))
162 .replay(self)
163 }
164
165 fn encode(
166 &self,
167 request: completion::CompletionRequest,
168 mode: Mode,
169 ) -> Result<Encoded, EncodeError> {
170 self.encode_with_headers(request, mode, OpenAIConfig::completion_headers)
171 }
172
173 fn decoder<'id>(&self) -> Self::Decoder<'id> {
174 ResponsesDecoder::new()
175 }
176}
177
178impl crate::completion::ReplayTarget for Responses {
179 fn map_options(
181 &self,
182 request: &crate::completion::CompletionRequest,
183 fields: crate::completion::options::OptionFields<'_>,
184 ) -> crate::completion::options::OptionMap {
185 crate::providers::openai::options::responses_options(self, request, fields)
186 }
187
188 fn api(&self) -> crate::message::Api {
189 crate::message::Api::from_static("openai.responses")
190 }
191
192 fn provider(&self) -> &str {
193 self.provider.dialect.name
194 }
195
196 fn model(&self) -> &str {
197 &self.model
198 }
199
200 fn declares_tools(&self, request: &completion::CompletionRequest) -> bool {
202 !self.tools.is_empty() || crate::completion::history::declares_tools(request)
203 }
204
205 fn call_id_slot(&self) -> Option<&'static str> {
207 Some("/call_id")
208 }
209
210 fn identity(&self, item: &Value) -> Map<String, Value> {
214 let keys: &[&str] = match item.str("type") {
215 Some("message") => &["type", "id", "phase"],
216 Some("reasoning") => &["type", "id", "encrypted_content"],
217 Some("function_call" | "custom_tool_call") => &["type", "id"],
218 _ => &[],
219 };
220 keys.iter()
221 .filter_map(|key| Some(((*key).to_owned(), item.get(*key)?.clone())))
222 .collect()
223 }
224
225 fn needs_next(&self, item: &Value) -> bool {
227 item.str("type") == Some("reasoning")
228 }
229
230 fn accepts(&self, model: &str) -> crate::completion::Accepts {
234 let images = reads_images(self.provider.dialect.quirks.responses.contract, model);
235 let model = model.to_ascii_lowercase();
236 crate::completion::Accepts {
237 user_images: images,
238 assistant_images: false,
239 tool_result_images: images,
240 tools: !(model.starts_with("o1-mini") || model.starts_with("o1-preview")),
241 }
242 }
243
244 fn encodes(&self, _model: &str, media: crate::completion::Media<'_>) -> bool {
249 use crate::completion::{Media, Place};
250 use crate::message::DocumentSourceKind;
251 let quirks = &self.provider.dialect.quirks;
252 let file_ids = quirks.accepts_file_ids;
253 match media {
254 Media::Image(_, Place::Assistant) | Media::Audio(_) | Media::Video(_) => false,
255 Media::Image(image, _) => {
256 super::image_part(image).is_some()
257 && (!matches!(image.data, DocumentSourceKind::FileId(_))
258 || file_ids && quirks.responses.contract != ResponsesContract::Xai)
259 }
260 Media::Document(document) => {
261 super::document_part(document).is_some()
262 && (file_ids || !matches!(document.data, DocumentSourceKind::FileId(_)))
263 }
264 }
265 }
266
267 fn continues_stored(&self, request: &completion::CompletionRequest) -> bool {
271 ["previous_response_id", "conversation"].iter().any(|key| {
272 crate::completion::options::param(self, request, key)
273 .is_some_and(|value| !value.is_null())
274 })
275 }
276
277 fn normalize_tool_call_id(
280 &self,
281 id: &str,
282 _model: &str,
283 _source: Option<&crate::message::Origin>,
284 ) -> String {
285 use crate::providers::internal::wire_ids::{legal_call_id, short_hash};
286 let legal = legal_call_id(id, 64);
287 match legal.trim_end_matches('_') {
288 "" => short_hash(id),
290 trimmed => trimmed.to_owned(),
291 }
292 }
293}
294
295fn reads_images(contract: ResponsesContract, model: &str) -> bool {
299 let model = model.rsplit('/').next().unwrap_or_default();
300 match contract {
301 ResponsesContract::Xai => crate::catalog::reads_images_or(
302 crate::providers::xai::DIALECT.name,
303 model,
304 crate::providers::xai::reads_images,
305 ),
306 ResponsesContract::OpenAi | ResponsesContract::Codex => crate::catalog::reads_images_or(
307 crate::providers::openai::wire::OPENAI.name,
308 model,
309 crate::providers::openai::reads_images,
310 ),
311 }
312}
313
314pub(crate) fn project_payload(payload: &[u8], sink: &mut ObservationSink<'_>) {
322 let Ok(payload) = serde_json::from_slice::<Value>(payload) else {
323 return;
324 };
325 let envelope = |error: &Value| ObservedError {
326 code: error.get("code").filter(|code| !code.is_null()).cloned(),
327 kind: error
328 .str("type")
329 .or_else(|| error.str("status"))
330 .map(str::to_owned),
331 message: error.str("message").map(str::to_owned),
332 };
333 if payload.str("type") == Some("error") {
334 match payload.get("error").filter(|error| error.is_object()) {
337 Some(error) => envelope(error),
338 None => ObservedError {
339 kind: None,
340 ..envelope(&payload)
341 },
342 }
343 .emit(sink);
344 return;
345 }
346 let object = payload
347 .get("response")
348 .filter(|response| response.is_object())
349 .unwrap_or(&payload);
350 if let Some(usage) = object.get("usage").filter(|usage| usage.is_object()) {
351 let usage = usage_of(usage);
352 sink.emit(AdapterEvent::Usage {
353 usage: AdapterUsage {
354 input_tokens: usage.input_tokens,
355 output_tokens: usage.output_tokens,
356 total_tokens: usage.total_tokens,
357 cached_input_tokens: usage.cached_input_tokens,
358 reasoning_tokens: usage.reasoning_tokens,
359 tool_input_tokens: None,
360 },
361 });
362 }
363 let verdict = AdapterVerdict {
366 finish_reason: object
367 .str("status")
368 .filter(|status| *status != "in_progress" && *status != "queued")
369 .map(|value| sink.scrub(value)),
370 block_reason: None,
371 detail: object
372 .at("/incomplete_details/reason")
373 .and_then(Value::as_str)
374 .map(|value| sink.scrub(value)),
375 model: object.str("model").map(|value| sink.scrub(value)),
376 };
377 let response_id = object.str("id").map(|value| sink.scrub(value));
378 sink.provider(verdict, response_id);
379 if let Some(error) = object.get("error").filter(|error| error.is_object()) {
380 envelope(error).emit(sink);
381 }
382}
383
384#[cfg(test)]
385mod tests;