Skip to main content

vtcode_llm/open_responses/
integration.rs

1//! Integration layer for Open Responses with VT Code agent infrastructure.
2//!
3//! This module provides the glue between VT Code's internal event system
4//! and the Open Responses specification. It uses configuration from
5//! `vtcode.toml` to control when and how Open Responses events are emitted.
6
7use std::sync::{Arc, Mutex};
8
9use crate::provider::NormalizedStreamEvent;
10use vtcode_config::OpenResponsesConfig;
11use vtcode_exec_events::ThreadEvent;
12
13use super::{OpenUsage, OutputItem, Response, ResponseBuilder, ResponseStreamEvent, VecStreamEmitter};
14
15/// Callback type for receiving Open Responses streaming events.
16pub type OpenResponsesCallback = Arc<Mutex<Box<dyn FnMut(ResponseStreamEvent) + Send>>>;
17
18/// Open Responses integration manager.
19///
20/// This struct manages the integration between VT Code's internal event system
21/// and the Open Responses specification. It respects the configuration flags
22/// to control event emission and item mapping.
23pub struct OpenResponsesIntegration {
24    config: OpenResponsesConfig,
25    builder: Option<ResponseBuilder>,
26    events: Vec<ResponseStreamEvent>,
27    callback: Option<OpenResponsesCallback>,
28}
29
30impl std::fmt::Debug for OpenResponsesIntegration {
31    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
32        f.debug_struct("OpenResponsesIntegration")
33            .field("config", &self.config)
34            .field("builder", &self.builder)
35            .field("events_count", &self.events.len())
36            .field("callback_set", &self.callback.is_some())
37            .finish()
38    }
39}
40
41impl OpenResponsesIntegration {
42    /// Creates a new integration manager with the given configuration.
43    pub fn new(config: OpenResponsesConfig) -> Self {
44        Self {
45            config,
46            builder: None,
47            events: Vec::new(),
48            callback: None,
49        }
50    }
51
52    /// Creates a new integration manager that is disabled.
53    fn disabled() -> Self {
54        Self::new(OpenResponsesConfig::default())
55    }
56
57    /// Returns true if Open Responses integration is enabled.
58    fn is_enabled(&self) -> bool {
59        self.config.enabled
60    }
61
62    /// Sets a callback for receiving Open Responses events.
63    pub fn set_callback(&mut self, callback: OpenResponsesCallback) {
64        self.callback = Some(callback);
65    }
66
67    /// Starts a new response session with the given model.
68    ///
69    /// This should be called when a new agent turn begins.
70    pub fn start_response(&mut self, model: &str) {
71        if !self.config.enabled {
72            return;
73        }
74
75        self.builder = Some(ResponseBuilder::new(model));
76        self.events.clear();
77    }
78
79    /// Processes a VT Code ThreadEvent and emits corresponding Open Responses events.
80    pub fn process_event(&mut self, event: &ThreadEvent) {
81        if !self.config.enabled || !self.config.emit_events {
82            return;
83        }
84
85        let Some(builder) = self.builder.as_mut() else {
86            return;
87        };
88
89        // Use a collecting emitter first
90        let mut collector = VecStreamEmitter::new();
91        builder.process_event(event, &mut collector);
92
93        // Process collected events
94        for stream_event in collector.into_events() {
95            // Apply filtering based on config
96            if self.should_emit_event(&stream_event) {
97                self.emit_event(stream_event);
98            }
99        }
100    }
101
102    /// Processes a normalized provider stream event and emits corresponding Open Responses events.
103    fn process_normalized_event(&mut self, event: &NormalizedStreamEvent) {
104        if !self.config.enabled || !self.config.emit_events {
105            return;
106        }
107
108        let Some(builder) = self.builder.as_mut() else {
109            return;
110        };
111
112        let mut collector = VecStreamEmitter::new();
113        builder.process_normalized_event(event, &mut collector);
114
115        for stream_event in collector.into_events() {
116            if self.should_emit_event(&stream_event) {
117                self.emit_event(stream_event);
118            }
119        }
120    }
121
122    /// Returns the current response, if any.
123    fn current_response(&self) -> Option<&Response> {
124        self.builder.as_ref().map(|b| b.response())
125    }
126
127    /// Finishes the current response and returns it.
128    pub fn finish_response(&mut self) -> Option<Response> {
129        self.builder.take().map(|b| b.build())
130    }
131
132    /// Returns all collected events.
133    fn events(&self) -> &[ResponseStreamEvent] {
134        &self.events
135    }
136
137    /// Takes all collected events, leaving the internal buffer empty.
138    pub fn take_events(&mut self) -> Vec<ResponseStreamEvent> {
139        std::mem::take(&mut self.events)
140    }
141
142    fn should_emit_event(&self, event: &ResponseStreamEvent) -> bool {
143        match event {
144            // Always emit response lifecycle events
145            ResponseStreamEvent::ResponseCreated { .. }
146            | ResponseStreamEvent::ResponseInProgress { .. }
147            | ResponseStreamEvent::ResponseCompleted { .. }
148            | ResponseStreamEvent::ResponseFailed { .. }
149            | ResponseStreamEvent::ResponseIncomplete { .. } => true,
150
151            // Filter output items based on config
152            ResponseStreamEvent::OutputItemAdded { item, .. } | ResponseStreamEvent::OutputItemDone { item, .. } => {
153                self.should_include_item(item)
154            }
155
156            // Reasoning events
157            ResponseStreamEvent::ReasoningDelta { .. } | ResponseStreamEvent::ReasoningDone { .. } => {
158                self.config.include_reasoning
159            }
160
161            // Function call events
162            ResponseStreamEvent::FunctionCallArgumentsDelta { .. }
163            | ResponseStreamEvent::FunctionCallArgumentsDone { .. } => self.config.map_tool_calls,
164
165            // Extension events
166            ResponseStreamEvent::CustomEvent { .. } => self.config.include_extensions,
167
168            // Text and content events are always emitted
169            _ => true,
170        }
171    }
172
173    fn should_include_item(&self, item: &OutputItem) -> bool {
174        match item {
175            OutputItem::Reasoning(_) => self.config.include_reasoning,
176            OutputItem::FunctionCall(_) | OutputItem::FunctionCallOutput(_) => self.config.map_tool_calls,
177            OutputItem::Custom(_) => self.config.include_extensions,
178            OutputItem::Message(_) => true,
179        }
180    }
181
182    fn emit_event(&mut self, event: ResponseStreamEvent) {
183        // Store in local buffer
184        self.events.push(event.clone());
185
186        // Send to callback if registered
187        if let Some(callback) = &self.callback
188            && let Ok(mut cb) = callback.lock()
189        {
190            cb(event);
191        }
192    }
193}
194
195impl Default for OpenResponsesIntegration {
196    fn default() -> Self {
197        Self::disabled()
198    }
199}
200
201/// Trait for types that can provide Open Responses integration.
202pub trait OpenResponsesProvider {
203    /// Returns a reference to the Open Responses integration, if enabled.
204    fn open_responses(&self) -> Option<&OpenResponsesIntegration>;
205
206    /// Returns a mutable reference to the Open Responses integration, if enabled.
207    fn open_responses_mut(&mut self) -> Option<&mut OpenResponsesIntegration>;
208}
209
210/// Extension trait for converting VT Code LLM responses to Open Responses format.
211pub trait ToOpenResponse {
212    /// Converts to an Open Responses Response object.
213    fn to_open_response(&self, response_id: &str, model: &str) -> Response;
214}
215
216impl ToOpenResponse for crate::provider::LLMResponse {
217    fn to_open_response(&self, response_id: &str, model: &str) -> Response {
218        let mut response = Response::new(response_id, model);
219
220        // Add usage if available
221        if let Some(usage) = &self.usage {
222            response.usage = Some(OpenUsage::from_llm_usage(usage).into());
223        }
224
225        // Add content as message item if present
226        if let Some(content) = &self.content
227            && !content.is_empty()
228        {
229            let item = OutputItem::completed_message(
230                super::response::generate_item_id(),
231                super::items::MessageRole::Assistant,
232                vec![super::ContentPart::output_text(content)],
233            );
234            response.add_output(item);
235        }
236
237        // Add reasoning if present
238        if let Some(reasoning) = &self.reasoning
239            && !reasoning.is_empty()
240        {
241            let item = OutputItem::Reasoning(super::items::ReasoningItem {
242                id: super::response::generate_item_id().into(),
243                status: super::ItemStatus::Completed,
244                summary: None,
245                content: Some(reasoning.clone()),
246                encrypted_content: None,
247            });
248            response.add_output(item);
249        }
250
251        // Add tool calls as function call items
252        if let Some(tool_calls) = &self.tool_calls {
253            for tc in tool_calls {
254                // Extract name and arguments from the function field
255                let (name, arguments) = if let Some(ref func) = tc.function {
256                    (func.name.clone(), serde_json::from_str(&func.arguments).unwrap_or(serde_json::json!({})))
257                } else {
258                    (tc.call_type.clone(), serde_json::json!({}))
259                };
260
261                let item = OutputItem::FunctionCall(super::items::FunctionCallItem {
262                    id: tc.id.clone().into(),
263                    status: super::ItemStatus::Completed,
264                    name,
265                    arguments,
266                    call_id: Some(tc.id.clone()),
267                });
268                response.add_output(item);
269            }
270        }
271
272        response.complete();
273        response
274    }
275}
276
277#[cfg(test)]
278mod tests {
279    use super::*;
280    use crate::provider::{FinishReason, LLMResponse, NormalizedStreamEvent};
281
282    #[test]
283    fn test_integration_disabled_by_default() {
284        let integration = OpenResponsesIntegration::disabled();
285        assert!(!integration.is_enabled());
286    }
287
288    #[test]
289    fn test_integration_enabled() {
290        let config = OpenResponsesConfig { enabled: true, ..Default::default() };
291        let integration = OpenResponsesIntegration::new(config);
292        assert!(integration.is_enabled());
293    }
294
295    #[test]
296    fn test_start_response() {
297        let config = OpenResponsesConfig { enabled: true, ..Default::default() };
298        let mut integration = OpenResponsesIntegration::new(config);
299        integration.start_response("gpt-5");
300        assert!(integration.current_response().is_some());
301    }
302
303    #[test]
304    fn test_disabled_skips_events() {
305        let mut integration = OpenResponsesIntegration::disabled();
306        integration.start_response("gpt-5");
307        // Should not create a builder when disabled
308        assert!(integration.current_response().is_none());
309    }
310
311    #[test]
312    fn integration_processes_normalized_events() {
313        let mut integration = OpenResponsesIntegration::new(OpenResponsesConfig {
314            enabled: true,
315            emit_events: true,
316            ..Default::default()
317        });
318        integration.start_response("gpt-5");
319
320        integration.process_normalized_event(&NormalizedStreamEvent::TextDelta { delta: "hello".to_string() });
321        integration.process_normalized_event(&NormalizedStreamEvent::Done {
322            response: Box::new(LLMResponse {
323                content: Some("hello".to_string()),
324                model: "gpt-5".to_string(),
325                tool_calls: None,
326                usage: None,
327                finish_reason: FinishReason::Stop,
328                reasoning: None,
329                reasoning_details: None,
330                organization_id: None,
331                request_id: None,
332                tool_references: Vec::new(),
333                compaction: None,
334            }),
335        });
336
337        assert!(integration.events().iter().any(|event| matches!(
338            event,
339            ResponseStreamEvent::OutputTextDelta { delta, .. } if delta == "hello"
340        )));
341        assert!(
342            integration
343                .events()
344                .iter()
345                .any(|event| matches!(event, ResponseStreamEvent::ResponseCompleted { .. }))
346        );
347    }
348}