Skip to main content

vtcode_llm/providers/openresponses/
provider.rs

1use crate::error_display;
2use crate::provider::{
3    FinishReason, LLMError, LLMNormalizedStream, LLMProvider, LLMRequest, LLMResponse, LLMStream, LLMStreamEvent,
4    Message, ResponsesCompactionOptions, ToolCall,
5};
6use crate::providers::common::{
7    append_normalized_reasoning_detail_items, chat_completions_url, serialize_message_content_openai,
8};
9use crate::providers::shared::{
10    ResponsesNormalizedStreamOptions, Utf8StreamDecoder, collect_tool_references_from_tool_search_output,
11    create_responses_normalized_stream, function_output_value_from_message_content, parse_compacted_output_messages,
12};
13use anyhow::Result;
14use async_stream::try_stream;
15use async_trait::async_trait;
16use futures::StreamExt;
17use reqwest::Client as HttpClient;
18use serde::Deserialize;
19use serde_json::{Value, json};
20use vtcode_config::TimeoutsConfig;
21use vtcode_config::constants::{env_vars, models, urls};
22use vtcode_config::core::{AnthropicConfig, ModelConfig, PromptCachingConfig};
23
24use super::super::common::{override_base_url, resolve_model};
25use super::super::error_handling::{format_network_error, format_parse_error};
26
27pub struct OpenResponsesProvider {
28    http_client: HttpClient,
29    base_url: String,
30    model: String,
31    api_key: String,
32    model_behavior: Option<ModelConfig>,
33}
34
35/// Borrowed fields used by the native Responses SSE hot path.
36///
37/// Responses events carry a discriminator and a small set of fields consumed
38/// by this provider. Deserializing into `Value` first allocates the complete
39/// event tree, including fields that are ignored by the streaming loop.
40#[derive(Debug, Deserialize)]
41struct NativeStreamEventWire<'a> {
42    #[serde(rename = "type")]
43    event_type: Option<&'a str>,
44    #[serde(borrow)]
45    delta: Option<&'a str>,
46    #[serde(borrow)]
47    item_id: Option<&'a str>,
48    item: Option<Value>,
49}
50
51#[derive(Debug, Deserialize)]
52struct ChatCompletionStreamEventWire<'a> {
53    #[serde(borrow)]
54    choices: Option<Vec<ChatCompletionChoiceWire<'a>>>,
55}
56
57#[derive(Debug, Deserialize)]
58struct ChatCompletionChoiceWire<'a> {
59    #[serde(borrow)]
60    delta: Option<ChatCompletionDeltaWire<'a>>,
61}
62
63#[derive(Debug, Deserialize)]
64struct ChatCompletionDeltaWire<'a> {
65    #[serde(borrow)]
66    content: Option<&'a str>,
67    tool_calls: Option<Vec<Value>>,
68}
69
70impl OpenResponsesProvider {
71    fn parse_native_response_payload(json: Value, model: String) -> Result<LLMResponse, LLMError> {
72        let output = json
73            .get("output")
74            .and_then(|o| o.as_array())
75            .ok_or_else(|| LLMError::Provider {
76                message: "Invalid response from OpenResponses: missing output".to_string(),
77                metadata: None,
78            })?;
79
80        let mut content = String::new();
81        let mut tool_calls = Vec::new();
82        let mut reasoning = None;
83        let mut tool_references = Vec::new();
84        let mut replay_items = Vec::new();
85
86        for item_val in output {
87            let item_type = item_val.get("type").and_then(|t| t.as_str()).unwrap_or("");
88            match item_type {
89                "message" => {
90                    if let Some(content_parts) = item_val.get("content").and_then(|c| c.as_array()) {
91                        for part in content_parts {
92                            if let Some(text) = part.get("text").and_then(|t| t.as_str()) {
93                                content.push_str(text);
94                            }
95                        }
96                    }
97                }
98                "reasoning" => {
99                    replay_items.push(item_val.clone());
100                    if let Some(text) = item_val.get("content").and_then(|t| t.as_str()) {
101                        reasoning = Some(text.to_string());
102                    }
103                }
104                "function_call" => {
105                    let id = item_val.get("id").and_then(|v| v.as_str()).unwrap_or("").to_string();
106                    let name = item_val.get("name").and_then(|v| v.as_str()).unwrap_or("").to_string();
107                    let arguments = item_val
108                        .get("arguments")
109                        .map(|v| v.to_string())
110                        .unwrap_or_else(|| "{}".to_string());
111                    let namespace = item_val.get("namespace").and_then(|v| v.as_str()).map(ToOwned::to_owned);
112                    tool_calls.push(ToolCall::function_with_namespace(id, namespace, name, arguments));
113                }
114                "tool_search_output" => {
115                    collect_tool_references_from_tool_search_output(item_val, &mut tool_references);
116                }
117                _ if !item_type.is_empty() => replay_items.push(item_val.clone()),
118                _ => {}
119            }
120        }
121
122        let mut reasoning_details = if replay_items.is_empty() {
123            None
124        } else {
125            Some(replay_items.into_iter().map(|item| item.to_string()).collect())
126        };
127        let (final_reasoning, final_content) = if reasoning.is_none() && !content.is_empty() {
128            let (reasoning_parts, cleaned_content) = crate::utils::extract_reasoning_content(&content);
129            if reasoning_parts.is_empty() {
130                (None, Some(content))
131            } else {
132                crate::providers::common::preserve_interleaved_content_in_reasoning_details(
133                    &mut reasoning_details,
134                    &content,
135                );
136                (Some(reasoning_parts.join("\n\n")), cleaned_content.or(Some(content)))
137            }
138        } else {
139            (reasoning, Some(content))
140        };
141
142        let finish_reason = match json.get("status").and_then(|s| s.as_str()) {
143            Some("completed") => FinishReason::Stop,
144            Some("incomplete") => FinishReason::Length,
145            _ => FinishReason::Stop,
146        };
147
148        Ok(LLMResponse {
149            content: final_content.filter(|c| !c.is_empty()),
150            tool_calls: if tool_calls.is_empty() { None } else { Some(tool_calls) },
151            model,
152            usage: None,
153            finish_reason,
154            reasoning: final_reasoning,
155            reasoning_details,
156            tool_references,
157            request_id: json.get("id").and_then(|v| v.as_str()).map(|s| s.to_string()),
158            organization_id: None,
159            compaction: None,
160        })
161    }
162
163    fn output_item_to_value(item: crate::open_responses::OutputItem) -> Result<Value, LLMError> {
164        serde_json::to_value(item).map_err(|e| LLMError::Provider {
165            message: format!("Failed to serialize Open Responses input item: {e}"),
166            metadata: None,
167        })
168    }
169
170    pub fn new(api_key: String) -> Self {
171        Self::with_model(api_key, models::openresponses::DEFAULT_MODEL.to_string())
172    }
173
174    pub fn with_model(api_key: String, model: String) -> Self {
175        Self::with_model_internal(model, None, api_key, TimeoutsConfig::default(), None)
176    }
177
178    fn new_with_client(
179        api_key: String,
180        model: String,
181        http_client: reqwest::Client,
182        base_url: String,
183        _timeouts: TimeoutsConfig,
184    ) -> Self {
185        Self {
186            http_client,
187            base_url,
188            model,
189            api_key,
190            model_behavior: None,
191        }
192    }
193
194    pub fn from_config(
195        api_key: Option<String>,
196        model: Option<String>,
197        base_url: Option<String>,
198        _prompt_cache: Option<PromptCachingConfig>,
199        timeouts: Option<TimeoutsConfig>,
200        _anthropic: Option<AnthropicConfig>,
201        model_behavior: Option<ModelConfig>,
202    ) -> Self {
203        let api_key_value = api_key.unwrap_or_default();
204        let resolved_model = resolve_model(model, models::openresponses::DEFAULT_MODEL);
205        Self::with_model_internal(resolved_model, base_url, api_key_value, timeouts.unwrap_or_default(), model_behavior)
206    }
207
208    fn with_model_internal(
209        model: String,
210        base_url: Option<String>,
211        api_key: String,
212        timeouts: TimeoutsConfig,
213        model_behavior: Option<ModelConfig>,
214    ) -> Self {
215        use crate::http_client::HttpClientFactory;
216
217        Self {
218            http_client: HttpClientFactory::for_llm(&timeouts),
219            base_url: override_base_url(urls::OPENRESPONSES_API_BASE, base_url, Some(env_vars::OPENRESPONSES_BASE_URL)),
220            model,
221            api_key,
222            model_behavior,
223        }
224    }
225
226    fn responses_url(&self) -> String {
227        format!("{}/responses", self.base_url.trim_end_matches('/'))
228    }
229
230    fn responses_compact_url(&self) -> String {
231        format!("{}/responses/compact", self.base_url.trim_end_matches('/'))
232    }
233
234    /// Client for driving another endpoint's OpenAI-compatible
235    /// `/responses/compact` surface (e.g. Vercel AI Gateway OpenAI routes, xAI
236    /// Grok models) with this provider's transport and response parsing.
237    /// Capability gating stays with the caller: only construct this for routes
238    /// whose upstream documents the endpoint.
239    pub(crate) fn compact_endpoint_client(configured_model: &str, base_url: &str, api_key: &str, model: &str) -> Self {
240        let resolved = if model.trim().is_empty() {
241            configured_model.to_string()
242        } else {
243            model.to_string()
244        };
245        Self::from_config(Some(api_key.to_string()), Some(resolved), Some(base_url.to_string()), None, None, None, None)
246    }
247
248    fn supports_compaction_endpoint(&self) -> bool {
249        self.base_url.contains("api.openai.com") || self.base_url.contains("api.openresponses.com")
250    }
251
252    pub(crate) async fn compact_history_request(
253        &self,
254        model: &str,
255        history: &[Message],
256    ) -> Result<Vec<Message>, LLMError> {
257        self.compact_history_request_with_options(model, history, &ResponsesCompactionOptions::default())
258            .await
259    }
260
261    pub(crate) async fn compact_history_request_with_options(
262        &self,
263        model: &str,
264        history: &[Message],
265        options: &ResponsesCompactionOptions,
266    ) -> Result<Vec<Message>, LLMError> {
267        let resolved_model = if model.trim().is_empty() {
268            self.model.clone()
269        } else {
270            model.trim().to_string()
271        };
272        let request = LLMRequest {
273            model: resolved_model.clone(),
274            messages: std::sync::Arc::new(history.to_vec()),
275            ..Default::default()
276        };
277        let native_payload = self.build_native_payload(&request, false)?;
278        let input = native_payload.get("input").cloned().unwrap_or_else(|| json!([]));
279        let mut compact_payload = json!({
280            "model": resolved_model,
281            "input": input,
282        });
283        if let Some(map) = compact_payload.as_object_mut() {
284            if let Some(instructions) = options.instructions.as_deref().map(str::trim).filter(|value| !value.is_empty())
285            {
286                map.insert("instructions".to_string(), json!(instructions));
287            }
288            if let Some(service_tier) = options.service_tier.as_deref().map(str::trim).filter(|value| !value.is_empty())
289            {
290                map.insert("service_tier".to_string(), json!(service_tier));
291            }
292            if let Some(prompt_cache_key) = options
293                .prompt_cache_key
294                .as_deref()
295                .map(str::trim)
296                .filter(|value| !value.is_empty())
297            {
298                map.insert("prompt_cache_key".to_string(), json!(prompt_cache_key));
299            }
300        }
301
302        let response = self
303            .http_client
304            .post(self.responses_compact_url())
305            .bearer_auth(&self.api_key)
306            .json(&compact_payload)
307            .send()
308            .await
309            .map_err(|e| format_network_error("OpenResponses", &e))?;
310
311        if !response.status().is_success() {
312            let status = response.status();
313            let body = crate::providers::common::read_provider_error_body(response).await;
314            let formatted_error = error_display::format_llm_error(
315                "OpenResponses",
316                &format!("Compaction endpoint error (HTTP {status}): {body}"),
317            );
318            return Err(LLMError::Provider { message: formatted_error, metadata: None });
319        }
320
321        let json: Value = response.json().await.map_err(|e| format_parse_error("OpenResponses", &e))?;
322        let output = json
323            .get("output")
324            .and_then(|value| value.as_array())
325            .ok_or_else(|| LLMError::Provider {
326                message: "Invalid response from OpenResponses compact endpoint: missing output array".to_string(),
327                metadata: None,
328            })?;
329
330        let compacted = parse_compacted_output_messages(output);
331        if compacted.is_empty() {
332            return Err(LLMError::Provider {
333                message: "Compaction response contained no reusable messages".to_string(),
334                metadata: None,
335            });
336        }
337
338        Ok(compacted)
339    }
340
341    fn build_native_payload(&self, request: &LLMRequest, stream: bool) -> Result<Value, LLMError> {
342        use crate::open_responses::{
343            ContentPart, ImageDetail, InputFileContent, InputImageContent, MessageRole, OutputItem, Request,
344        };
345
346        let mut input: Vec<Value> = Vec::new();
347
348        if let Some(system) = &request.system_prompt {
349            input.push(Self::output_item_to_value(OutputItem::completed_message(
350                "msg_system",
351                MessageRole::System,
352                vec![ContentPart::input_text(system.as_ref())],
353            ))?);
354        }
355
356        for (i, message) in request.messages.iter().enumerate() {
357            if let Some(reasoning_details) = &message.reasoning_details {
358                append_normalized_reasoning_detail_items(&mut input, reasoning_details);
359            }
360
361            let role = match message.role.as_generic_str() {
362                "user" => Some(MessageRole::User),
363                "assistant" => Some(MessageRole::Assistant),
364                "system" => Some(MessageRole::System),
365                // Tool responses are represented by function_call_output items below.
366                "tool" => None,
367                _ => Some(MessageRole::User),
368            };
369
370            if let Some(role) = role {
371                let id = format!("msg_{i}");
372                let mut content = Vec::new();
373                match &message.content {
374                    crate::provider::MessageContent::Text(text) => {
375                        if !text.trim().is_empty() {
376                            content.push(ContentPart::input_text(text.as_str()));
377                        }
378                    }
379                    crate::provider::MessageContent::Parts(parts) => {
380                        for part in parts {
381                            match part {
382                                crate::provider::ContentPart::Text { text } => {
383                                    if !text.trim().is_empty() {
384                                        content.push(ContentPart::input_text(text.as_str()));
385                                    }
386                                }
387                                crate::provider::ContentPart::Image { data, mime_type, .. } => {
388                                    // Providers accept only JPEG/PNG/GIF/WebP. Anything else
389                                    // (notably SVG auto-attached from quoted paths in diffs)
390                                    // fails the whole request with 400, so drop it here
391                                    // instead of letting one bad part poison the turn.
392                                    if !vtcode_commons::image::is_supported_image_mime_type(mime_type) {
393                                        tracing::warn!(
394                                            mime_type = %mime_type,
395                                            "dropping unsupported image MIME type for Responses API"
396                                        );
397                                        continue;
398                                    }
399                                    content.push(ContentPart::InputImage(InputImageContent {
400                                        image_url: format!("data:{mime_type};base64,{data}"),
401                                        detail: Some(ImageDetail::Auto),
402                                    }));
403                                }
404                                crate::provider::ContentPart::File {
405                                    filename, file_id, file_data, file_url, ..
406                                } => {
407                                    content.push(ContentPart::InputFile(InputFileContent {
408                                        filename: filename.clone(),
409                                        file_id: file_id.clone(),
410                                        file_data: file_data.clone(),
411                                        file_url: file_url.clone(),
412                                    }));
413                                }
414                            }
415                        }
416                    }
417                }
418                if content.is_empty() {
419                    let content_text = message.content.as_text();
420                    if !content_text.trim().is_empty() {
421                        content.push(ContentPart::input_text(content_text.to_string()));
422                    }
423                }
424                if !content.is_empty() {
425                    input.push(Self::output_item_to_value(OutputItem::completed_message(id, role, content))?);
426                }
427            }
428
429            // Handle tool calls and outputs if present in message history
430            if let Some(tool_calls) = &message.tool_calls {
431                for (j, tc) in tool_calls.iter().enumerate() {
432                    if let Some(f) = &tc.function {
433                        input.push(Self::output_item_to_value(OutputItem::function_call(
434                            format!("fc_{i}_{j}"),
435                            &f.name,
436                            tc.parsed_arguments().unwrap_or(Value::Null),
437                        ))?);
438                    }
439                }
440            }
441
442            if let Some(tool_call_id) = &message.tool_call_id {
443                // If this message is a tool output, add it as FunctionCallOutput
444                input.push(json!({
445                    "type": "function_call_output",
446                    "id": format!("fco_{i}"),
447                    "status": "completed",
448                    "call_id": tool_call_id,
449                    "output": function_output_value_from_message_content(&message.content),
450                }));
451            }
452        }
453
454        let mut req = Request::new(&request.model, Vec::new());
455        req.stream = stream;
456        req.temperature = request.temperature.map(crate::providers::common::sampling_param_f64);
457        req.max_output_tokens = request.max_tokens.map(|t| t as u64);
458        req.previous_response_id = request
459            .previous_response_id
460            .as_ref()
461            .map(|value| value.trim().to_string())
462            .filter(|value| !value.is_empty());
463        req.store = request.response_store;
464        req.include = request.responses_include.as_ref().and_then(|fields| {
465            let values: Vec<String> = fields
466                .iter()
467                .map(|field| field.trim())
468                .filter(|field| !field.is_empty())
469                .map(ToOwned::to_owned)
470                .collect();
471            if values.is_empty() { None } else { Some(values) }
472        });
473
474        if let Some(tools) = &request.tools {
475            req.tools = Some(tools.as_ref().clone());
476        }
477
478        let mut payload = serde_json::to_value(req).map_err(|e| LLMError::Provider {
479            message: format!("Failed to serialize Open Responses request: {e}"),
480            metadata: None,
481        })?;
482        if let Some(map) = payload.as_object_mut() {
483            map.insert("input".to_string(), Value::Array(input));
484        }
485
486        if let Some(context_management) = &request.context_management
487            && let Some(map) = payload.as_object_mut()
488        {
489            map.insert("context_management".to_string(), context_management.clone());
490        }
491
492        Ok(payload)
493    }
494
495    fn build_payload(&self, request: &LLMRequest, stream: bool) -> Result<Value, LLMError> {
496        let mut messages = Vec::new();
497
498        if let Some(system) = &request.system_prompt {
499            messages.push(json!({
500                "role": "system",
501                "content": system
502            }));
503        }
504
505        for message in request.messages.iter() {
506            let role = message.role.as_generic_str();
507            let mut message_obj = json!({
508                "role": role,
509                "content": serialize_message_content_openai(&message.content)
510            });
511
512            if let Some(tool_calls) = &message.tool_calls {
513                let tool_calls_json: Vec<Value> = tool_calls
514                    .iter()
515                    .filter_map(|tc| {
516                        tc.function.as_ref().map(|f| {
517                            json!({
518                                "id": tc.id,
519                                "type": "function",
520                                "function": {
521                                    "name": f.name,
522                                    "arguments": f.arguments
523                                }
524                            })
525                        })
526                    })
527                    .collect();
528                message_obj["tool_calls"] = json!(tool_calls_json);
529            }
530
531            if let Some(tool_call_id) = &message.tool_call_id {
532                message_obj["tool_call_id"] = json!(tool_call_id);
533            }
534
535            messages.push(message_obj);
536        }
537
538        let mut payload = json!({
539            "model": request.model,
540            "messages": messages,
541            "stream": stream
542        });
543
544        if let Some(max_tokens) = request.max_tokens {
545            payload["max_tokens"] = json!(max_tokens);
546        }
547
548        if let Some(temp) = request.temperature {
549            payload["temperature"] = json!(crate::providers::common::sampling_param_f64(temp));
550        }
551
552        if let Some(tools) = &request.tools {
553            let tools_json: Vec<Value> = tools
554                .iter()
555                .filter_map(|t| {
556                    t.function.as_ref().map(|f| {
557                        json!({
558                            "type": "function",
559                            "function": {
560                                "name": f.name,
561                                "description": f.description,
562                                "parameters": f.parameters
563                            }
564                        })
565                    })
566                })
567                .collect();
568            payload["tools"] = json!(tools_json);
569        }
570
571        Ok(payload)
572    }
573
574    async fn generate_fallback(&self, request: LLMRequest) -> Result<LLMResponse, LLMError> {
575        let model = request.model.clone();
576        let payload = self.build_payload(&request, false)?;
577        let url = chat_completions_url(&self.base_url);
578
579        let response = self
580            .http_client
581            .post(url)
582            .bearer_auth(&self.api_key)
583            .json(&payload)
584            .send()
585            .await
586            .map_err(|e| format_network_error("OpenResponses", &e))?;
587
588        if !response.status().is_success() {
589            let status = response.status();
590            let body = crate::providers::common::read_provider_error_body(response).await;
591            let formatted_error = error_display::format_llm_error("OpenResponses", &format!("HTTP {status}: {body}"));
592            return Err(LLMError::Provider { message: formatted_error, metadata: None });
593        }
594
595        let json: Value = response.json().await.map_err(|e| format_parse_error("OpenResponses", &e))?;
596
597        let choice = json
598            .get("choices")
599            .and_then(|c| c.as_array())
600            .and_then(|c| c.first())
601            .ok_or_else(|| LLMError::Provider {
602                message: "Invalid response from OpenResponses: missing choices".to_string(),
603                metadata: None,
604            })?;
605
606        let message = choice.get("message").ok_or_else(|| LLMError::Provider {
607            message: "Invalid response from OpenResponses: missing message".to_string(),
608            metadata: None,
609        })?;
610
611        let content = message.get("content").and_then(|c| c.as_str()).map(|s| s.to_string());
612
613        let tool_calls = message
614            .get("tool_calls")
615            .and_then(|tc| tc.as_array())
616            .map(|calls| {
617                calls
618                    .iter()
619                    .filter_map(|call| {
620                        let id = call.get("id").and_then(|v| v.as_str())?;
621                        let function = call.get("function")?;
622                        let namespace = call
623                            .get("namespace")
624                            .and_then(|v| v.as_str())
625                            .or_else(|| function.get("namespace").and_then(|v| v.as_str()))
626                            .map(ToOwned::to_owned);
627                        let name = function.get("name").and_then(|v| v.as_str())?;
628                        let arguments = function.get("arguments").and_then(|v| v.as_str())?;
629                        Some(ToolCall::function_with_namespace(
630                            id.to_string(),
631                            namespace,
632                            name.to_string(),
633                            arguments.to_string(),
634                        ))
635                    })
636                    .collect::<Vec<_>>()
637            })
638            .filter(|calls| !calls.is_empty());
639
640        let finish_reason = choice
641            .get("finish_reason")
642            .and_then(|fr| fr.as_str())
643            .map(|fr| match fr {
644                "stop" => FinishReason::Stop,
645                "length" => FinishReason::Length,
646                "tool_calls" => FinishReason::ToolCalls,
647                other => FinishReason::Error(other.to_string()),
648            })
649            .unwrap_or(FinishReason::Stop);
650
651        Ok(LLMResponse {
652            content,
653            tool_calls,
654            model,
655            usage: None,
656            finish_reason,
657            reasoning: None,
658            reasoning_details: None,
659            tool_references: Vec::new(),
660            request_id: json.get("id").and_then(|v| v.as_str()).map(|s| s.to_string()),
661            organization_id: None,
662            compaction: None,
663        })
664    }
665
666    async fn stream_fallback(&self, request: LLMRequest) -> Result<LLMStream, LLMError> {
667        let model = request.model.clone();
668        let payload = self.build_payload(&request, true)?;
669        let url = chat_completions_url(&self.base_url);
670
671        let response = self
672            .http_client
673            .post(url)
674            .bearer_auth(&self.api_key)
675            .json(&payload)
676            .send()
677            .await
678            .map_err(|e| format_network_error("OpenResponses", &e))?;
679
680        if !response.status().is_success() {
681            let status = response.status();
682            let body = crate::providers::common::read_provider_error_body(response).await;
683            let formatted_error = error_display::format_llm_error("OpenResponses", &format!("HTTP {status}: {body}"));
684            return Err(LLMError::Provider { message: formatted_error, metadata: None });
685        }
686
687        let stream = try_stream! {
688            let mut body_stream = response.bytes_stream();
689            let mut buf: Vec<u8> = Vec::new();
690            let mut offset = 0usize;
691            let mut decoder = Utf8StreamDecoder::new();
692            let mut aggregator = crate::providers::shared::StreamAggregator::new(model);
693
694            while let Some(chunk_result) = body_stream.next().await {
695                let chunk = chunk_result.map_err(|e| format_network_error("OpenResponses", &e))?;
696                decoder.push_bytes(&chunk, &mut buf);
697
698                while let Some(event) =
699                    crate::providers::shared::next_sse_event(&buf, &mut offset).expect("valid utf-8 stream data")
700                {
701
702                    if let Some(data_payload) = crate::providers::shared::extract_data_payload(event) {
703                        let trimmed = data_payload.trim();
704                        if trimmed.is_empty() || trimmed == "[DONE]" {
705                            continue;
706                        }
707
708                        if let Ok(payload) = serde_json::from_str::<ChatCompletionStreamEventWire<'_>>(trimmed)
709                            && let Some(delta) = payload.choices.and_then(|choices| choices.into_iter().next()).and_then(|choice| choice.delta)
710                        {
711                            if let Some(content) = delta.content {
712                                for ev in aggregator.handle_content(content) {
713                                    yield ev;
714                                }
715                            }
716
717                            if let Some(tool_calls) = delta.tool_calls.as_deref() {
718                                aggregator.handle_tool_calls(tool_calls);
719                            }
720                        }
721                    }
722                }
723
724                // Keep `buf` bounded to the unprocessed tail rather than
725                // growing for the entire stream.
726                crate::providers::shared::drain_consumed_sse(&mut buf, &mut offset);
727            }
728
729            yield LLMStreamEvent::Completed { response: Box::new(aggregator.finalize()) };
730        };
731
732        Ok(Box::pin(stream))
733    }
734}
735
736#[async_trait]
737impl LLMProvider for OpenResponsesProvider {
738    fn name(&self) -> &str {
739        "openresponses"
740    }
741
742    fn supports_streaming(&self) -> bool {
743        true
744    }
745
746    fn supports_non_streaming(&self, _model: &str) -> bool {
747        // Pinned so the stream-timeout fallback cannot silently regress.
748        true
749    }
750
751    fn supports_reasoning(&self, _model: &str) -> bool {
752        self.model_behavior
753            .as_ref()
754            .and_then(|b| b.model_supports_reasoning)
755            .unwrap_or(true) // Open Responses usually implies reasoning support
756    }
757
758    fn supports_reasoning_effort(&self, _model: &str) -> bool {
759        self.model_behavior
760            .as_ref()
761            .and_then(|b| b.model_supports_reasoning_effort)
762            .unwrap_or(true)
763    }
764
765    fn supports_responses_compaction(&self, _model: &str) -> bool {
766        self.supports_compaction_endpoint()
767    }
768
769    // OpenResponses exposes the same standalone `/responses/compact` endpoint as
770    // the native OpenAI provider, so it opts into the NativeStandalone dispatch
771    // (`compact_history_with_options`). Without this, the unified dispatch would
772    // route it to NativeInline (Anthropic `compact_20260112` via `generate`),
773    // which OpenResponses cannot serve, silently downgrading auto/recovery
774    // compaction to the local summarization fallback.
775    fn supports_manual_openai_compaction(&self, _model: &str) -> bool {
776        self.supports_compaction_endpoint()
777    }
778
779    async fn compact_history(&self, model: &str, history: &[Message]) -> Result<Vec<Message>, LLMError> {
780        if !self.supports_compaction_endpoint() {
781            return Err(LLMError::Provider {
782                message: "OpenResponses compact endpoint is not supported for this configured base URL".to_string(),
783                metadata: None,
784            });
785        }
786
787        self.compact_history_request(model, history).await
788    }
789
790    // The compact endpoint accepts the common instruction and routing fields;
791    // output-shaping options are intentionally omitted because this endpoint
792    // does not expose them in its request schema.
793    async fn compact_history_with_options(
794        &self,
795        model: &str,
796        history: &[Message],
797        options: &ResponsesCompactionOptions,
798    ) -> Result<Vec<Message>, LLMError> {
799        if !self.supports_compaction_endpoint() {
800            return Err(LLMError::Provider {
801                message: "OpenResponses compact endpoint is not supported for this configured base URL".to_string(),
802                metadata: None,
803            });
804        }
805
806        self.compact_history_request_with_options(model, history, options).await
807    }
808
809    fn supported_models(&self) -> Vec<String> {
810        use vtcode_config::constants::models::openresponses::SUPPORTED_MODELS;
811        SUPPORTED_MODELS.iter().map(|s| s.to_string()).collect()
812    }
813
814    fn validate_request(&self, request: &LLMRequest) -> Result<(), LLMError> {
815        if request.model.is_empty() {
816            return Err(LLMError::Provider {
817                message: "Model is required for OpenResponses provider".to_string(),
818                metadata: None,
819            });
820        }
821
822        let supported = self.supported_models();
823        if !supported.contains(&request.model) {
824            return Err(LLMError::Provider {
825                message: format!(
826                    "Model '{}' is not supported by OpenResponses provider. Supported models: {}",
827                    request.model,
828                    supported.join(", ")
829                ),
830                metadata: None,
831            });
832        }
833
834        Ok(())
835    }
836
837    async fn generate(&self, mut request: LLMRequest) -> Result<LLMResponse, LLMError> {
838        if request.model.is_empty() {
839            request.model = self.model.clone();
840        }
841        let model = request.model.clone();
842
843        // Try native Open Responses endpoint first
844        let payload = self.build_native_payload(&request, false)?;
845        let url = self.responses_url();
846
847        let response = self
848            .http_client
849            .post(url)
850            .bearer_auth(&self.api_key)
851            .json(&payload)
852            .send()
853            .await
854            .map_err(|e| format_network_error("OpenResponses", &e))?;
855
856        // If native endpoint fails with 404, fallback to chat/completions
857        if response.status() == reqwest::StatusCode::NOT_FOUND {
858            return self.generate_fallback(request).await;
859        }
860
861        if !response.status().is_success() {
862            let status = response.status();
863            let body = crate::providers::common::read_provider_error_body(response).await;
864            let formatted_error = error_display::format_llm_error("OpenResponses", &format!("HTTP {status}: {body}"));
865            return Err(LLMError::Provider { message: formatted_error, metadata: None });
866        }
867
868        let json: Value = response.json().await.map_err(|e| format_parse_error("OpenResponses", &e))?;
869
870        Self::parse_native_response_payload(json, model)
871    }
872
873    async fn stream(&self, mut request: LLMRequest) -> Result<LLMStream, LLMError> {
874        if request.model.is_empty() {
875            request.model = self.model.clone();
876        }
877        let model = request.model.clone();
878
879        let payload = self.build_native_payload(&request, true)?;
880        let url = self.responses_url();
881
882        let response = self
883            .http_client
884            .post(url)
885            .bearer_auth(&self.api_key)
886            .json(&payload)
887            .send()
888            .await
889            .map_err(|e| format_network_error("OpenResponses", &e))?;
890
891        if response.status() == reqwest::StatusCode::NOT_FOUND {
892            return self.stream_fallback(request).await;
893        }
894
895        if !response.status().is_success() {
896            let status = response.status();
897            let body = crate::providers::common::read_provider_error_body(response).await;
898            let formatted_error = error_display::format_llm_error("OpenResponses", &format!("HTTP {status}: {body}"));
899            return Err(LLMError::Provider { message: formatted_error, metadata: None });
900        }
901
902        let stream = try_stream! {
903            let mut body_stream = response.bytes_stream();
904            let mut buf: Vec<u8> = Vec::new();
905            let mut offset = 0usize;
906            let mut decoder = Utf8StreamDecoder::new();
907            let mut aggregator = crate::providers::shared::StreamAggregator::new(model);
908
909            while let Some(chunk_result) = body_stream.next().await {
910                let chunk = chunk_result.map_err(|e| format_network_error("OpenResponses", &e))?;
911                decoder.push_bytes(&chunk, &mut buf);
912
913                while let Some(event) =
914                    crate::providers::shared::next_sse_event(&buf, &mut offset).expect("valid utf-8 stream data")
915                {
916
917                    if let Some(data_payload) = crate::providers::shared::extract_data_payload(event) {
918                        let trimmed = data_payload.trim();
919                        if trimmed.is_empty() || trimmed == "[DONE]" {
920                            continue;
921                        }
922
923                        if let Ok(event) = serde_json::from_str::<NativeStreamEventWire<'_>>(trimmed) {
924                            let event_type = event.event_type.unwrap_or("");
925
926                            match event_type {
927                                "response.output_text.delta" => {
928                                    if let Some(delta) = event.delta {
929                                        // Use aggregator's sanitizer to extract reasoning tags from content
930                                        for ev in aggregator.handle_content(delta) {
931                                            yield ev;
932                                        }
933                                    }
934                                }
935                                "response.function_call_arguments.delta" => {
936                                    if let Some(delta) = event.delta {
937                                        let tc_json = json!([{
938                                            "index": 0,
939                                            "id": event.item_id,
940                                            "function": { "arguments": delta }
941                                        }]);
942                                        if let Some(tool_calls) = tc_json.as_array() {
943                                            aggregator.handle_tool_calls(tool_calls);
944                                        }
945                                    }
946                                }
947                                "response.reasoning.delta" => {
948                                    // Legacy/simple reasoning event
949                                    if let Some(delta) = event.delta {
950                                        yield LLMStreamEvent::Reasoning { delta: delta.to_string() };
951                                    }
952                                }
953                                "response.reasoning_content.delta" => {
954                                    // Raw reasoning traces (preferred)
955                                    if let Some(delta) = event.delta {
956                                        yield LLMStreamEvent::Reasoning { delta: delta.to_string() };
957                                    }
958                                }
959                                "response.reasoning_summary_text.delta" => {
960                                    // Summary reasoning (fallback when raw not available)
961                                    if let Some(delta) = event.delta {
962                                        yield LLMStreamEvent::Reasoning { delta: delta.to_string() };
963                                    }
964                                }
965                                // The added event may contain an incomplete
966                                // opaque item. Only the done snapshot is safe
967                                // to replay on the next request.
968                                "response.output_item.done" => {
969                                    if let Some(item) = event.item.as_ref()
970                                        && item.get("type").and_then(Value::as_str) == Some("compaction")
971                                    {
972                                        aggregator.append_reasoning_detail(item);
973                                    }
974                                }
975                                _ => {}
976                            }
977                        }
978                    }
979                }
980
981                // Keep `buf` bounded to the unprocessed tail rather than
982                // growing for the entire stream.
983                crate::providers::shared::drain_consumed_sse(&mut buf, &mut offset);
984            }
985
986            yield LLMStreamEvent::Completed { response: Box::new(aggregator.finalize()) };
987        };
988
989        Ok(Box::pin(stream))
990    }
991
992    async fn stream_normalized(&self, mut request: LLMRequest) -> Result<LLMNormalizedStream, LLMError> {
993        if request.model.is_empty() {
994            request.model = self.model.clone();
995        }
996        let model = request.model.clone();
997
998        let payload = self.build_native_payload(&request, true)?;
999        let url = self.responses_url();
1000
1001        let response = self
1002            .http_client
1003            .post(url)
1004            .bearer_auth(&self.api_key)
1005            .json(&payload)
1006            .send()
1007            .await
1008            .map_err(|e| format_network_error("OpenResponses", &e))?;
1009
1010        if response.status() == reqwest::StatusCode::NOT_FOUND {
1011            let mut legacy_stream = self.stream_fallback(request).await?;
1012            let stream = try_stream! {
1013                while let Some(event) = legacy_stream.next().await {
1014                    for normalized in event?.into_normalized() {
1015                        yield normalized;
1016                    }
1017                }
1018            };
1019            return Ok(Box::pin(stream));
1020        }
1021
1022        if !response.status().is_success() {
1023            let status = response.status();
1024            let body = crate::providers::common::read_provider_error_body(response).await;
1025            let formatted_error = error_display::format_llm_error("OpenResponses", &format!("HTTP {status}: {body}"));
1026            return Err(LLMError::Provider { message: formatted_error, metadata: None });
1027        }
1028
1029        let emit_reasoning = self.supports_reasoning(&model);
1030        Ok(create_responses_normalized_stream(
1031            response,
1032            ResponsesNormalizedStreamOptions {
1033                provider_name: "OpenResponses",
1034                model: model.clone(),
1035                emit_reasoning,
1036                include_cached_prompt_metrics: false,
1037            },
1038            move |value| Self::parse_native_response_payload(value, model.clone()),
1039        ))
1040    }
1041}
1042
1043#[cfg(test)]
1044mod tests;