Skip to main content

vtcode_llm/open_responses/bridge/
items.rs

1//! ThreadEvent item handling: started/updated/completed item conversion.
2
3use super::*;
4
5fn tool_output_text(output: &ToolOutputItem) -> String {
6    if !output.output.is_empty() {
7        return output.output.clone();
8    }
9
10    output
11        .spool_path
12        .as_deref()
13        .map(|path| format!("Output saved to {path}"))
14        .unwrap_or_default()
15}
16
17impl ResponseBuilder {
18    pub fn process_event<E: StreamEventEmitter>(&mut self, event: &ThreadEvent, emitter: &mut E) {
19        match event {
20            ThreadEvent::ThreadStarted(_) => {
21                emitter.response_created(self.response.clone());
22                self.response.status = ResponseStatus::InProgress;
23                emitter.response_in_progress(self.response.clone());
24                self.normalized.response_started = true;
25            }
26
27            ThreadEvent::TurnStarted(_) => {
28                // Turn started is internal to VT Code; no direct Open Responses equivalent
29                // The response is already in progress from ThreadStarted
30            }
31
32            ThreadEvent::TurnCompleted(evt) => {
33                if self.response.status.is_terminal() {
34                    return;
35                }
36                self.response.usage = Some(OpenUsage::from_exec_usage(&evt.usage).into());
37                self.response.status = ResponseStatus::Completed;
38                self.response.complete();
39                emitter.response_completed(self.response.clone());
40            }
41
42            ThreadEvent::TurnFailed(evt) => {
43                if self.response.status.is_terminal() {
44                    return;
45                }
46                self.response.fail(OpenResponseError::model_error(&evt.message));
47                emitter.response_failed(self.response.clone());
48            }
49
50            ThreadEvent::TurnBlocked(evt) => {
51                self.emit_custom_event(
52                    emitter,
53                    "vtcode.turn_blocked",
54                    json!({
55                        "completed_at": evt.completed_at,
56                        "message": evt.message,
57                        "last_tool": evt.last_tool,
58                        "blocked_streak": evt.blocked_streak,
59                        "blocked_total": evt.blocked_total,
60                        "consecutive_cap": evt.consecutive_cap,
61                        "total_cap": evt.total_cap,
62                        "recovery_active": evt.recovery_active,
63                    }),
64                );
65            }
66
67            ThreadEvent::ThreadCompleted(evt) => {
68                self.emit_custom_event(
69                    emitter,
70                    "vtcode.thread_completed",
71                    json!({
72                        "completed_at": evt.completed_at,
73                        "thread_id": evt.thread_id,
74                        "session_id": evt.session_id,
75                        "subtype": evt.subtype.as_str(),
76                        "outcome_code": evt.outcome_code,
77                        "result": evt.result,
78                        "stop_reason": evt.stop_reason,
79                        "usage": evt.usage,
80                        "total_cost_usd": evt.total_cost_usd,
81                        "num_turns": evt.num_turns,
82                    }),
83                );
84            }
85
86            ThreadEvent::ThreadCompactBoundary(evt) => {
87                self.emit_custom_event(
88                    emitter,
89                    "vtcode.thread_compact_boundary",
90                    json!({
91                        "thread_id": evt.thread_id,
92                        "trigger": evt.trigger.as_str(),
93                        "mode": evt.mode.as_str(),
94                        "original_message_count": evt.original_message_count,
95                        "compacted_message_count": evt.compacted_message_count,
96                        "history_artifact_path": evt.history_artifact_path,
97                    }),
98                );
99            }
100
101            ThreadEvent::ContextReset(evt) => {
102                self.emit_custom_event(
103                    emitter,
104                    "vtcode.context_reset",
105                    json!({
106                        "thread_id": evt.thread_id,
107                        "turn_id": evt.turn_id,
108                        "trigger": evt.trigger,
109                        "plan_preserved": evt.plan_preserved,
110                        "previous_context_usage_percent": evt.previous_context_usage_percent,
111                        "tool_budget_reset": evt.tool_budget_reset,
112                    }),
113                );
114            }
115
116            ThreadEvent::ItemStarted(evt) => {
117                self.handle_item_started(&evt.item, emitter);
118            }
119
120            ThreadEvent::ItemUpdated(evt) => {
121                self.handle_item_updated(&evt.item, emitter);
122            }
123
124            ThreadEvent::ItemCompleted(evt) => {
125                self.handle_item_completed(&evt.item, emitter);
126            }
127            ThreadEvent::PlanDelta(_) => {
128                // Plan deltas are VT Code-specific extension events and are intentionally
129                // ignored by the Open Responses bridge. The completed Plan item carries
130                // the full final plan content.
131            }
132
133            ThreadEvent::PlanApprovalRequested(evt) => {
134                self.emit_custom_event(
135                    emitter,
136                    "vtcode.plan_approval_requested",
137                    json!({
138                        "thread_id": evt.thread_id,
139                        "turn_id": evt.turn_id,
140                        "plan_file": evt.plan_file,
141                    }),
142                );
143            }
144
145            ThreadEvent::PlanApprovalResolved(evt) => {
146                self.emit_custom_event(
147                    emitter,
148                    "vtcode.plan_approval_resolved",
149                    json!({
150                        "thread_id": evt.thread_id,
151                        "turn_id": evt.turn_id,
152                        "decision": evt.decision,
153                        "automatic": evt.automatic,
154                    }),
155                );
156            }
157
158            ThreadEvent::Error(evt) => {
159                if self.response.status.is_terminal() {
160                    return;
161                }
162                self.response.fail(OpenResponseError::server_error(&evt.message));
163                emitter.response_failed(self.response.clone());
164            }
165
166            // Unknown events from newer schema versions are silently skipped.
167            ThreadEvent::Unknown
168            | ThreadEvent::PermissionRequested(_)
169            | ThreadEvent::PermissionResolved(_)
170            | ThreadEvent::Interjected(_)
171            | ThreadEvent::MatrixUpdated(_) => {}
172        }
173    }
174
175    fn handle_item_started<E: StreamEventEmitter>(&mut self, item: &ThreadItem, emitter: &mut E) {
176        let output_index = self.next_output_index;
177        self.next_output_index += 1;
178        self.item_id_to_index.insert(item.id.clone(), output_index);
179
180        let output_item = self.convert_thread_item(item, ItemStatus::InProgress);
181
182        // Track active item state for streaming
183        // Initialize prev_text from initial content to prevent duplicate deltas
184        let initial_text = match &item.details {
185            ThreadItemDetails::AgentMessage(msg) => msg.text.clone(),
186            ThreadItemDetails::Plan(plan) => plan.text.clone(),
187            ThreadItemDetails::Reasoning(r) => r.text.clone(),
188            ThreadItemDetails::ToolOutput(output) => tool_output_text(output),
189            _ => String::new(),
190        };
191        let active_state = ActiveItemState {
192            output_index,
193            content_index: 0,
194            prev_text: initial_text,
195        };
196        self.active_items.insert(item.id.clone(), active_state);
197
198        self.response.add_output(output_item.clone());
199        emitter.output_item_added(&self.response.id, output_index, output_item.clone());
200
201        // Emit ContentPartAdded for items with text content
202        if let OutputItem::Message(ref msg) = output_item
203            && !msg.content.is_empty()
204        {
205            emitter.emit(ResponseStreamEvent::ContentPartAdded {
206                response_id: self.response.id.clone(),
207                item_id: item.id.clone(),
208                output_index,
209                content_index: 0,
210                part: msg.content[0].clone(),
211            });
212        }
213    }
214
215    fn emit_custom_event<E: StreamEventEmitter>(&self, emitter: &mut E, event_type: &str, data: serde_json::Value) {
216        emitter.emit(ResponseStreamEvent::CustomEvent {
217            response_id: self.response.id.clone(),
218            event_type: event_type.to_string(),
219            sequence_number: self.next_output_index as u64,
220            data,
221        });
222    }
223
224    fn handle_item_updated<E: StreamEventEmitter>(&mut self, item: &ThreadItem, emitter: &mut E) {
225        // Handle updates for items not yet started (implicit start)
226        let state = if let Some(state) = self.active_items.get_mut(&item.id) {
227            state
228        } else {
229            // Implicit start: create item and emit Added event
230            self.handle_item_started(item, emitter);
231            match self.active_items.get_mut(&item.id) {
232                Some(s) => s,
233                None => return,
234            }
235        };
236
237        match &item.details {
238            ThreadItemDetails::AgentMessage(msg) => {
239                // Use strip_prefix for safe UTF-8 delta computation
240                let delta = if let Some(suffix) = msg.text.strip_prefix(&state.prev_text) {
241                    suffix
242                } else {
243                    // Non-append update: emit full text as delta (fallback)
244                    &msg.text
245                };
246
247                if !delta.is_empty() {
248                    emitter.output_text_delta(
249                        &self.response.id,
250                        &item.id,
251                        state.output_index,
252                        state.content_index,
253                        delta,
254                    );
255                    state.prev_text = msg.text.clone();
256                }
257            }
258
259            ThreadItemDetails::Reasoning(r) => {
260                // Use strip_prefix for safe UTF-8 delta computation
261                let delta = if let Some(suffix) = r.text.strip_prefix(&state.prev_text) {
262                    suffix
263                } else {
264                    // Non-append update: emit full text as delta (fallback)
265                    &r.text
266                };
267
268                if !delta.is_empty() {
269                    emitter.reasoning_delta(&self.response.id, &item.id, state.output_index, delta);
270                    state.prev_text = r.text.clone();
271                }
272            }
273
274            ThreadItemDetails::ToolOutput(output) => {
275                let current_text = tool_output_text(output);
276                let delta = if let Some(suffix) = current_text.strip_prefix(&state.prev_text) {
277                    suffix
278                } else {
279                    current_text.as_str()
280                };
281
282                if !delta.is_empty() {
283                    emitter.output_text_delta(
284                        &self.response.id,
285                        &item.id,
286                        state.output_index,
287                        state.content_index,
288                        delta,
289                    );
290                    state.prev_text = current_text;
291                }
292            }
293
294            _ => {
295                // Other item types don't have incremental updates
296            }
297        }
298    }
299
300    fn handle_item_completed<E: StreamEventEmitter>(&mut self, item: &ThreadItem, emitter: &mut E) {
301        let (was_started, output_index) = match self.item_id_to_index.get(&item.id) {
302            Some(&idx) => (true, idx),
303            None => {
304                // Item was completed without being started (atomic item)
305                let idx = self.next_output_index;
306                self.next_output_index += 1;
307                self.item_id_to_index.insert(item.id.clone(), idx);
308                (false, idx)
309            }
310        };
311
312        // Determine final status
313        let status = self.determine_item_status(&item.details);
314        let output_item = self.convert_thread_item(item, status);
315
316        // For atomic completions (never started), emit Added first, then ContentPartAdded
317        if !was_started {
318            emitter.output_item_added(&self.response.id, output_index, output_item.clone());
319
320            // Emit ContentPartAdded for Message and Reasoning items
321            match &output_item {
322                OutputItem::Message(msg) => {
323                    if !msg.content.is_empty() {
324                        emitter.emit(ResponseStreamEvent::ContentPartAdded {
325                            response_id: self.response.id.clone(),
326                            item_id: item.id.clone(),
327                            output_index,
328                            content_index: 0,
329                            part: msg.content[0].clone(),
330                        });
331                    }
332                }
333                OutputItem::Reasoning(r) => {
334                    let text = r.content.clone().unwrap_or_default();
335                    emitter.emit(ResponseStreamEvent::ContentPartAdded {
336                        response_id: self.response.id.clone(),
337                        item_id: item.id.clone(),
338                        output_index,
339                        content_index: 0,
340                        part: ContentPart::output_text(text),
341                    });
342                }
343                _ => {}
344            }
345        }
346
347        // Update the response output
348        if output_index < self.response.output.len() {
349            self.response.output[output_index] = output_item.clone();
350        } else {
351            self.response.add_output(output_item.clone());
352        }
353
354        // Emit content-specific "done" events based on item type
355        match &output_item {
356            OutputItem::Message(msg) => {
357                // Emit OutputTextDone for text content
358                if let Some(ContentPart::OutputText(text_content)) = msg.content.first() {
359                    emitter.emit(ResponseStreamEvent::OutputTextDone {
360                        response_id: self.response.id.clone(),
361                        item_id: item.id.clone(),
362                        output_index,
363                        content_index: 0,
364                        text: text_content.text.clone(),
365                    });
366                    emitter.emit(ResponseStreamEvent::ContentPartDone {
367                        response_id: self.response.id.clone(),
368                        item_id: item.id.clone(),
369                        output_index,
370                        content_index: 0,
371                        part: msg.content[0].clone(),
372                    });
373                }
374            }
375            OutputItem::Reasoning(r) => {
376                // Emit ReasoningDone then ContentPartDone
377                emitter.emit(ResponseStreamEvent::ReasoningDone {
378                    response_id: self.response.id.clone(),
379                    item_id: item.id.clone(),
380                    output_index,
381                    item: output_item.clone(),
382                });
383                let text = r.content.clone().unwrap_or_default();
384                emitter.emit(ResponseStreamEvent::ContentPartDone {
385                    response_id: self.response.id.clone(),
386                    item_id: item.id.clone(),
387                    output_index,
388                    content_index: 0,
389                    part: ContentPart::output_text(text),
390                });
391            }
392            OutputItem::FunctionCall(fc) => {
393                // Emit FunctionCallArgumentsDone
394                if let Ok(args_str) = serde_json::to_string(&fc.arguments) {
395                    emitter.emit(ResponseStreamEvent::FunctionCallArgumentsDone {
396                        response_id: self.response.id.clone(),
397                        item_id: item.id.clone(),
398                        output_index,
399                        arguments: args_str,
400                    });
401                }
402            }
403            OutputItem::FunctionCallOutput(fco) if !fco.output.is_empty() => {
404                emitter.emit(ResponseStreamEvent::OutputTextDone {
405                    response_id: self.response.id.clone(),
406                    item_id: item.id.clone(),
407                    output_index,
408                    content_index: 0,
409                    text: fco.output.clone(),
410                });
411            }
412            _ => {}
413        }
414
415        // Clean up active state
416        self.active_items.remove(&item.id);
417
418        emitter.output_item_done(&self.response.id, output_index, output_item);
419    }
420
421    fn determine_item_status(&self, details: &ThreadItemDetails) -> ItemStatus {
422        match details {
423            ThreadItemDetails::CommandExecution(cmd) => match cmd.status {
424                CommandExecutionStatus::Completed => ItemStatus::Completed,
425                CommandExecutionStatus::Failed => ItemStatus::Failed,
426                CommandExecutionStatus::InProgress => ItemStatus::InProgress,
427            },
428            ThreadItemDetails::ToolInvocation(invocation) => match invocation.status {
429                vtcode_exec_events::ToolCallStatus::Completed => ItemStatus::Completed,
430                vtcode_exec_events::ToolCallStatus::Failed => ItemStatus::Failed,
431                vtcode_exec_events::ToolCallStatus::InProgress => ItemStatus::InProgress,
432            },
433            ThreadItemDetails::ToolOutput(output) => match output.status {
434                vtcode_exec_events::ToolCallStatus::Completed => ItemStatus::Completed,
435                vtcode_exec_events::ToolCallStatus::Failed => ItemStatus::Failed,
436                vtcode_exec_events::ToolCallStatus::InProgress => ItemStatus::InProgress,
437            },
438            ThreadItemDetails::FileChange(fc) => match fc.status {
439                PatchApplyStatus::Completed => ItemStatus::Completed,
440                PatchApplyStatus::Failed => ItemStatus::Failed,
441            },
442            ThreadItemDetails::McpToolCall(tc) => match tc.status {
443                Some(McpToolCallStatus::Completed) => ItemStatus::Completed,
444                Some(McpToolCallStatus::Failed) => ItemStatus::Failed,
445                Some(McpToolCallStatus::Started) | None => ItemStatus::InProgress,
446            },
447            ThreadItemDetails::Error(_) => ItemStatus::Failed,
448            _ => ItemStatus::Completed,
449        }
450    }
451
452    fn resolve_tool_call_correlation_id(&mut self, harness_call_id: &str, raw_tool_call_id: Option<&str>) -> String {
453        if let Some(existing) = self.tool_call_correlation_ids.get(harness_call_id) {
454            return existing.clone();
455        }
456
457        let correlation_id = match raw_tool_call_id {
458            Some(raw_id) if self.used_tool_call_ids.insert(raw_id.to_string()) => raw_id.to_string(),
459            _ => harness_call_id.to_string(),
460        };
461        self.tool_call_correlation_ids
462            .insert(harness_call_id.to_string(), correlation_id.clone());
463        correlation_id
464    }
465
466    fn convert_thread_item(&mut self, item: &ThreadItem, status: ItemStatus) -> OutputItem {
467        match &item.details {
468            ThreadItemDetails::Decision(decision) => OutputItem::Custom(CustomItem {
469                id: item.id.clone().into(),
470                status,
471                custom_type: "vtcode:decision".into(),
472                data: json!({"decision": decision, "context": item.context}),
473            }),
474            ThreadItemDetails::AgentMessage(msg) => OutputItem::Message(MessageItem {
475                id: item.id.clone().into(),
476                status,
477                role: MessageRole::Assistant,
478                content: vec![ContentPart::output_text(&msg.text)],
479            }),
480
481            ThreadItemDetails::Reasoning(r) => OutputItem::Reasoning(ReasoningItem {
482                id: item.id.clone().into(),
483                status,
484                summary: None,
485                content: Some(r.text.clone()),
486                encrypted_content: None,
487            }),
488
489            ThreadItemDetails::Plan(plan) => OutputItem::Custom(CustomItem {
490                id: item.id.clone().into(),
491                status,
492                custom_type: "vtcode:plan".to_string(),
493                data: json!({
494                    "text": plan.text,
495                }),
496            }),
497
498            ThreadItemDetails::CommandExecution(cmd) => OutputItem::Custom(CustomItem {
499                id: item.id.clone().into(),
500                status,
501                custom_type: "vtcode:command_execution".to_string(),
502                data: json!({
503                    "command": cmd.command,
504                    "arguments": cmd.arguments,
505                    "aggregated_output": cmd.aggregated_output,
506                    "exit_code": cmd.exit_code,
507                    "status": serde_json::to_value(&cmd.status).unwrap_or(serde_json::Value::Null),
508                }),
509            }),
510
511            ThreadItemDetails::ToolInvocation(invocation) => OutputItem::FunctionCall(FunctionCallItem {
512                id: item.id.clone().into(),
513                status,
514                name: invocation.tool_name.clone(),
515                arguments: invocation.arguments.clone().unwrap_or(json!({})),
516                call_id: Some(self.resolve_tool_call_correlation_id(&item.id, invocation.tool_call_id.as_deref())),
517            }),
518
519            ThreadItemDetails::ToolOutput(output) => {
520                OutputItem::FunctionCallOutput(crate::open_responses::FunctionCallOutputItem {
521                    id: item.id.clone().into(),
522                    status,
523                    call_id: Some(
524                        self.resolve_tool_call_correlation_id(&output.call_id, output.tool_call_id.as_deref()),
525                    ),
526                    output: tool_output_text(output),
527                })
528            }
529
530            ThreadItemDetails::FileChange(fc) => {
531                let changes: Vec<_> = fc
532                    .changes
533                    .iter()
534                    .map(|c| {
535                        json!({
536                            "path": c.path,
537                            "kind": format!("{:?}", c.kind).to_lowercase(),
538                        })
539                    })
540                    .collect();
541
542                OutputItem::Custom(CustomItem {
543                    id: item.id.clone().into(),
544                    status,
545                    custom_type: "vtcode:file_change".to_string(),
546                    data: json!({
547                        "changes": changes,
548                        "status": format!("{:?}", fc.status).to_lowercase(),
549                    }),
550                })
551            }
552
553            ThreadItemDetails::McpToolCall(tc) => OutputItem::FunctionCall(FunctionCallItem {
554                id: item.id.clone().into(),
555                status,
556                name: tc.tool_name.clone(),
557                arguments: tc.arguments.clone().unwrap_or(json!({})),
558                call_id: Some(item.id.clone()),
559            }),
560
561            ThreadItemDetails::WebSearch(ws) => OutputItem::Custom(CustomItem {
562                id: item.id.clone().into(),
563                status,
564                custom_type: "vtcode:web_search".to_string(),
565                data: json!({
566                    "query": ws.query,
567                    "provider": ws.provider,
568                    "results": ws.results,
569                }),
570            }),
571
572            ThreadItemDetails::Harness(event) => OutputItem::Custom(CustomItem {
573                id: item.id.clone().into(),
574                status,
575                custom_type: "vtcode:harness_event".to_string(),
576                data: json!({
577                    "event": serde_json::to_value(&event.event).unwrap_or(serde_json::Value::Null),
578                    "message": event.message,
579                    "command": event.command,
580                    "path": event.path,
581                    "exit_code": event.exit_code,
582                }),
583            }),
584
585            ThreadItemDetails::Error(err) => {
586                // Errors are represented as custom items
587                OutputItem::Custom(CustomItem {
588                    id: item.id.clone().into(),
589                    status: ItemStatus::Failed,
590                    custom_type: "vtcode:error".to_string(),
591                    data: json!({
592                        "message": err.message,
593                    }),
594                })
595            }
596        }
597    }
598}