vtcode_llm/open_responses/
integration.rs1use 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
15pub type OpenResponsesCallback = Arc<Mutex<Box<dyn FnMut(ResponseStreamEvent) + Send>>>;
17
18pub 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 pub fn new(config: OpenResponsesConfig) -> Self {
44 Self {
45 config,
46 builder: None,
47 events: Vec::new(),
48 callback: None,
49 }
50 }
51
52 fn disabled() -> Self {
54 Self::new(OpenResponsesConfig::default())
55 }
56
57 fn is_enabled(&self) -> bool {
59 self.config.enabled
60 }
61
62 pub fn set_callback(&mut self, callback: OpenResponsesCallback) {
64 self.callback = Some(callback);
65 }
66
67 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 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 let mut collector = VecStreamEmitter::new();
91 builder.process_event(event, &mut collector);
92
93 for stream_event in collector.into_events() {
95 if self.should_emit_event(&stream_event) {
97 self.emit_event(stream_event);
98 }
99 }
100 }
101
102 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 fn current_response(&self) -> Option<&Response> {
124 self.builder.as_ref().map(|b| b.response())
125 }
126
127 pub fn finish_response(&mut self) -> Option<Response> {
129 self.builder.take().map(|b| b.build())
130 }
131
132 fn events(&self) -> &[ResponseStreamEvent] {
134 &self.events
135 }
136
137 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 ResponseStreamEvent::ResponseCreated { .. }
146 | ResponseStreamEvent::ResponseInProgress { .. }
147 | ResponseStreamEvent::ResponseCompleted { .. }
148 | ResponseStreamEvent::ResponseFailed { .. }
149 | ResponseStreamEvent::ResponseIncomplete { .. } => true,
150
151 ResponseStreamEvent::OutputItemAdded { item, .. } | ResponseStreamEvent::OutputItemDone { item, .. } => {
153 self.should_include_item(item)
154 }
155
156 ResponseStreamEvent::ReasoningDelta { .. } | ResponseStreamEvent::ReasoningDone { .. } => {
158 self.config.include_reasoning
159 }
160
161 ResponseStreamEvent::FunctionCallArgumentsDelta { .. }
163 | ResponseStreamEvent::FunctionCallArgumentsDone { .. } => self.config.map_tool_calls,
164
165 ResponseStreamEvent::CustomEvent { .. } => self.config.include_extensions,
167
168 _ => 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 self.events.push(event.clone());
185
186 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
201pub trait OpenResponsesProvider {
203 fn open_responses(&self) -> Option<&OpenResponsesIntegration>;
205
206 fn open_responses_mut(&mut self) -> Option<&mut OpenResponsesIntegration>;
208}
209
210pub trait ToOpenResponse {
212 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 if let Some(usage) = &self.usage {
222 response.usage = Some(OpenUsage::from_llm_usage(usage).into());
223 }
224
225 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 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 if let Some(tool_calls) = &self.tool_calls {
253 for tc in tool_calls {
254 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 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}