Skip to main content

rig_core/providers/openai/responses_api/
wire.rs

1//! Responses endpoint encoding and observation for configured OpenAI dialects.
2//!
3//! ```
4//! use rig_core::providers::openai::OpenAI;
5//! let wire = OpenAI::new("key").responses("gpt-5.2");
6//! ```
7
8use 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/// The Responses wire: `POST /responses`, SSE when streamed.
26#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
27pub struct Responses {
28    /// The provider this wire speaks to.
29    pub provider: OpenAIConfig,
30    /// The model to address.
31    pub model: String,
32    /// Tools added to every request from this wire.
33    pub tools: Vec<ResponsesToolDefinition>,
34    /// Whether Rig-generated tools request OpenAI's strict validation.
35    pub strict_tools: bool,
36    /// Where this wire puts Rig's system instructions. Defaults to the
37    /// dialect's placement.
38    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        // The codex gateway only ever answers with an event stream, and
54        // names no content type on it. It is asked for one whatever the
55        // caller wanted: the reply is framed the same way either way, and
56        // the driver folds it.
57        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    /// Create a wire with the provider's instruction placement and dialect's
96    /// strict-tool default, without additional tools.
97    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    /// Sanitize function schemas for strict mode and send `strict: true`.
108    pub fn with_strict_tools(mut self) -> Self {
109        self.strict_tools = true;
110        self
111    }
112
113    /// Add a tool to every request from this wire.
114    pub fn with_tool(mut self, tool: impl Into<ResponsesToolDefinition>) -> Self {
115        self.tools.push(tool.into());
116        self
117    }
118
119    /// Add tools to every request from this wire.
120    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    /// Put Rig's system instructions somewhere other than the dialect's
130    /// default placement.
131    pub fn with_system_instructions_placement(
132        mut self,
133        placement: SystemInstructionsPlacement,
134    ) -> Self {
135        self.system_instructions = placement;
136        self
137    }
138
139    /// Send Rig's system instructions as `system` messages in `input`, for a
140    /// backend that rejects or ignores top-level `instructions`.
141    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    /// The xAI contract does not compose native structured output with tools.
154    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    /// Section 6.3 of the typed-options design, by dialect.
180    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    /// The wire's own tools are declared beside the request's.
201    fn declares_tools(&self, request: &completion::CompletionRequest) -> bool {
202        !self.tools.is_empty() || crate::completion::history::declares_tools(request)
203    }
204
205    /// A call item names its id in `call_id`.
206    fn call_id_slot(&self) -> Option<&'static str> {
207        Some("/call_id")
208    }
209
210    /// An edited block keeps its item's `id` and `type`, a message its
211    /// `phase`, and reasoning its ciphertext, so the items after it stay
212    /// paired (pi's text signature is `{id, phase}`).
213    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    /// A reasoning item goes only with the item it preceded.
226    fn needs_next(&self, item: &Value) -> bool {
227        item.str("type") == Some("reasoning")
228    }
229
230    /// Responses reads images in user input and in function outputs, never
231    /// in assistant messages, and only on a model with vision input. Every
232    /// documented model calls tools except `o1-mini` and `o1-preview`.
233    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    /// Responses carries user and tool-result images as data URLs, URLs or
245    /// file ids, documents as files or text, and no audio, video or
246    /// assistant image. File ids need a dialect that resolves them, and
247    /// xAI's `input_image` takes none.
248    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    /// A request naming `previous_response_id` or a `conversation`
268    /// continues state the provider stores, which holds the calls its first
269    /// results answer.
270    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    /// pi's `normalizeIdPart`: characters outside `[a-zA-Z0-9_-]` become
278    /// `_`, the id is cut to 64 characters and loses its trailing `_`.
279    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            // An id with nothing legal in it still names its call.
289            "" => short_hash(id),
290            trimmed => trimmed.to_owned(),
291        }
292    }
293}
294
295/// Whether `model` reads images, past a `vendor/` prefix: its catalog
296/// entry's input, or for a model the catalog does not list its vendor's
297/// documented text-only models. An unknown model reads them.
298fn 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
314/// The facts a Responses payload carries before normalization discards
315/// them: the verdict, the model, the response id, the usage and any error
316/// envelope.
317///
318/// The unary reply is the response object itself; a stream event wraps that
319/// object under `response` (`response.created`, `.completed`, `.failed`,
320/// `.incomplete`) or, for `error`, carries the envelope's fields itself.
321pub(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        // The event carries its envelope either nested under `error` or as
335        // its own top-level fields; the nested form names the error type.
336        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    // `status` is the provider's verdict; `in_progress` on a stream's
364    // opening event is not one yet, so it is left out of the projection.
365    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;