Skip to main content

opendev_tools_impl/agents/
events.rs

1use std::collections::HashMap;
2
3use tokio::sync::mpsc;
4
5/// Events emitted by a running subagent, consumed by the parent agent or TUI.
6#[derive(Debug, Clone)]
7pub enum SubagentEvent {
8    /// Subagent started.
9    Started {
10        subagent_id: String,
11        subagent_name: String,
12        task: String,
13        cancel_token: Option<tokio_util::sync::CancellationToken>,
14    },
15    /// Subagent made a tool call.
16    ToolCall {
17        subagent_id: String,
18        subagent_name: String,
19        tool_name: String,
20        tool_id: String,
21        args: HashMap<String, serde_json::Value>,
22    },
23    /// A subagent tool call completed.
24    ToolComplete {
25        subagent_id: String,
26        subagent_name: String,
27        tool_name: String,
28        tool_id: String,
29        success: bool,
30    },
31    /// Subagent finished.
32    Finished {
33        subagent_id: String,
34        subagent_name: String,
35        success: bool,
36        result_summary: String,
37        tool_call_count: usize,
38        shallow_warning: Option<String>,
39    },
40    /// Token usage update from a subagent's LLM call.
41    TokenUpdate {
42        subagent_id: String,
43        subagent_name: String,
44        input_tokens: u64,
45        output_tokens: u64,
46    },
47}
48
49/// Progress callback that sends events through an mpsc channel.
50///
51/// Used to bridge subagent execution progress back to the TUI event loop.
52pub struct ChannelProgressCallback {
53    tx: mpsc::UnboundedSender<SubagentEvent>,
54    /// Unique identifier for this subagent instance (disambiguates parallel subagents).
55    subagent_id: String,
56    /// Per-subagent cancellation token (child of parent's token).
57    cancel_token: Option<tokio_util::sync::CancellationToken>,
58}
59
60impl ChannelProgressCallback {
61    /// Create a new channel-based progress callback with a unique subagent ID.
62    pub fn new(
63        tx: mpsc::UnboundedSender<SubagentEvent>,
64        subagent_id: String,
65        cancel_token: Option<tokio_util::sync::CancellationToken>,
66    ) -> Self {
67        Self {
68            tx,
69            subagent_id,
70            cancel_token,
71        }
72    }
73}
74
75impl std::fmt::Debug for ChannelProgressCallback {
76    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
77        f.debug_struct("ChannelProgressCallback").finish()
78    }
79}
80
81impl opendev_agents::SubagentProgressCallback for ChannelProgressCallback {
82    fn on_started(&self, subagent_name: &str, task: &str) {
83        let _ = self.tx.send(SubagentEvent::Started {
84            subagent_id: self.subagent_id.clone(),
85            subagent_name: subagent_name.to_string(),
86            task: task.to_string(),
87            cancel_token: self.cancel_token.clone(),
88        });
89    }
90
91    fn on_tool_call(
92        &self,
93        subagent_name: &str,
94        tool_name: &str,
95        tool_id: &str,
96        args: &HashMap<String, serde_json::Value>,
97    ) {
98        let _ = self.tx.send(SubagentEvent::ToolCall {
99            subagent_id: self.subagent_id.clone(),
100            subagent_name: subagent_name.to_string(),
101            tool_name: tool_name.to_string(),
102            tool_id: tool_id.to_string(),
103            args: args.clone(),
104        });
105    }
106
107    fn on_tool_complete(&self, subagent_name: &str, tool_name: &str, tool_id: &str, success: bool) {
108        let _ = self.tx.send(SubagentEvent::ToolComplete {
109            subagent_id: self.subagent_id.clone(),
110            subagent_name: subagent_name.to_string(),
111            tool_name: tool_name.to_string(),
112            tool_id: tool_id.to_string(),
113            success,
114        });
115    }
116
117    fn on_finished(&self, _subagent_name: &str, _success: bool, _result_summary: &str) {
118        // Don't emit Finished here — SpawnSubagentTool::execute() sends the
119        // authoritative Finished event with correct tool_call_count and shallow_warning.
120    }
121
122    fn on_token_usage(&self, subagent_name: &str, input_tokens: u64, output_tokens: u64) {
123        let _ = self.tx.send(SubagentEvent::TokenUpdate {
124            subagent_id: self.subagent_id.clone(),
125            subagent_name: subagent_name.to_string(),
126            input_tokens,
127            output_tokens,
128        });
129    }
130}
131
132#[cfg(test)]
133#[path = "events_tests.rs"]
134mod tests;