1use 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 }
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 }
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 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 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 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 let state = if let Some(state) = self.active_items.get_mut(&item.id) {
227 state
228 } else {
229 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 let delta = if let Some(suffix) = msg.text.strip_prefix(&state.prev_text) {
241 suffix
242 } else {
243 &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 let delta = if let Some(suffix) = r.text.strip_prefix(&state.prev_text) {
262 suffix
263 } else {
264 &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 }
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 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 let status = self.determine_item_status(&item.details);
314 let output_item = self.convert_thread_item(item, status);
315
316 if !was_started {
318 emitter.output_item_added(&self.response.id, output_index, output_item.clone());
319
320 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 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 match &output_item {
356 OutputItem::Message(msg) => {
357 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 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 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 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 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}