Skip to main content

vtcode_core/cli/
a2a.rs

1//! A2A Protocol CLI command handlers
2//!
3//! Implements the actual logic for A2A CLI commands including:
4//! - Starting the A2A server
5//! - Discovering remote agents
6//! - Sending tasks to agents
7//! - Managing A2A agent connections
8
9use crate::a2a::cli::A2aCommands;
10
11/// Execute an A2A CLI command
12pub async fn execute_a2a_command(command: A2aCommands) -> anyhow::Result<()> {
13    match command {
14        A2aCommands::Serve { host, port, base_url, enable_push } => {
15            serve_a2a_agent(host, port, base_url, enable_push).await
16        }
17        A2aCommands::Discover { agent_url } => discover_agent(agent_url).await,
18        A2aCommands::SendTask { agent_url, message, stream, context_id } => {
19            send_task_to_agent(agent_url, message, stream, context_id).await
20        }
21        A2aCommands::ListTasks { agent_url, context_id, limit } => list_agent_tasks(agent_url, context_id, limit).await,
22        A2aCommands::GetTask { agent_url, task_id } => get_agent_task(agent_url, task_id).await,
23        A2aCommands::CancelTask { agent_url, task_id } => cancel_agent_task(agent_url, task_id).await,
24    }
25}
26
27/// Serve VT Code as an A2A agent
28#[cfg(feature = "a2a-server")]
29async fn serve_a2a_agent(host: String, port: u16, base_url: Option<String>, _enable_push: bool) -> anyhow::Result<()> {
30    use crate::a2a::server::{A2aServerState, create_router};
31    use crate::a2a::{AgentCard, TaskManager};
32    use std::net::SocketAddr;
33
34    println!("Starting VT Code A2A Agent Server...");
35    println!("Feature: a2a-server enabled ✓");
36
37    let base_url = base_url.unwrap_or_else(|| format!("http://{host}:{port}"));
38    let agent_card = AgentCard::vtcode_default(&base_url);
39    let task_manager = TaskManager::new();
40
41    let server_state = match std::env::var(A2A_TOKEN_ENV_VAR) {
42        Ok(auth_token) => A2aServerState::new_with_auth_token(task_manager, agent_card, auth_token)?,
43        Err(std::env::VarError::NotPresent) => A2aServerState::new(task_manager, agent_card),
44        Err(error) => anyhow::bail!("Failed to read {A2A_TOKEN_ENV_VAR}: {error}"),
45    };
46    let auth_token = server_state.auth_token().to_string();
47    let router = create_router(server_state);
48
49    let addr = format!("{host}:{port}").parse::<SocketAddr>()?;
50    println!("Listening on http://{addr}");
51    println!("Agent Card: http://{addr}/.well-known/agent-card.json");
52    println!("JSON-RPC API: http://{addr}/a2a");
53    println!("Streaming API: http://{addr}/a2a/stream");
54    println!("Bearer token: {auth_token}");
55
56    let listener = tokio::net::TcpListener::bind(addr).await?;
57    axum::serve(listener, router)
58        .with_graceful_shutdown(crate::shutdown::shutdown_signal_logged("A2A"))
59        .await?;
60
61    Ok(())
62}
63
64#[cfg(not(feature = "a2a-server"))]
65async fn serve_a2a_agent(
66    _host: String,
67    _port: u16,
68    _base_url: Option<String>,
69    _enable_push: bool,
70) -> anyhow::Result<()> {
71    anyhow::bail!(
72        "A2A server is not enabled. Build with '--features a2a-server' to enable this feature.\n\
73         Example: cargo build --release --features a2a-server"
74    )
75}
76
77/// Discover and display information about a remote A2A agent
78async fn discover_agent(agent_url: String) -> anyhow::Result<()> {
79    use crate::a2a::A2aClient;
80
81    println!("Discovering A2A agent at: {agent_url}");
82
83    let client = A2aClient::new(&agent_url)?;
84    let agent_card = client.agent_card().await?;
85
86    println!("\n═══════════════════════════════════════════════════════════════");
87    println!("A2A Agent Discovery");
88    println!("═══════════════════════════════════════════════════════════════\n");
89
90    println!("Name: {}", agent_card.name);
91    println!("Description: {}", agent_card.description);
92    println!("Version: {}", agent_card.version);
93    println!("Protocol Version: {}", agent_card.protocol_version);
94    println!("URL: {}", agent_card.url);
95
96    if let Some(provider) = &agent_card.provider {
97        println!("\nProvider:");
98        println!("  Organization: {}", provider.organization);
99        if let Some(url) = &provider.url {
100            println!("  URL: {url}");
101        }
102    }
103
104    if let Some(capabilities) = &agent_card.capabilities {
105        println!("\nCapabilities:");
106        println!("  Streaming: {}", capabilities.streaming);
107        println!("  Push Notifications: {}", capabilities.push_notifications);
108        println!("  State Transition History: {}", capabilities.state_transition_history);
109        if !capabilities.extensions.is_empty() {
110            println!("  Extensions: {:?}", capabilities.extensions);
111        }
112    }
113
114    if !agent_card.skills.is_empty() {
115        println!("\nSkills:");
116        for skill in &agent_card.skills {
117            println!("  - {}", skill.name);
118            if let Some(desc) = &skill.description {
119                println!("    Description: {desc}");
120            }
121            if !skill.tags.is_empty() {
122                println!("    Tags: {:?}", skill.tags);
123            }
124        }
125    }
126
127    println!("\nInput Modes: {:?}", agent_card.default_input_modes);
128    println!("Output Modes: {:?}", agent_card.default_output_modes);
129
130    Ok(())
131}
132
133/// Send a task to a remote A2A agent
134async fn send_task_to_agent(
135    agent_url: String,
136    message: String,
137    stream: bool,
138    context_id: Option<String>,
139) -> anyhow::Result<()> {
140    use crate::a2a::{Message, rpc::MessageSendParams};
141    use futures::StreamExt;
142
143    println!("Connecting to A2A agent: {agent_url}");
144
145    let client = a2a_client(&agent_url)?;
146    let msg = Message::user_text(message);
147
148    let mut params = MessageSendParams::new(msg);
149    if let Some(ctx_id) = context_id {
150        params = params.with_context_id(ctx_id);
151    }
152
153    if stream {
154        println!("Streaming task execution...\n");
155        let stream = client.stream_message(params).await?;
156        futures::pin_mut!(stream);
157
158        while let Some(event) = stream.next().await {
159            match event {
160                Ok(event) => {
161                    // Handle different event types
162                    match event {
163                        crate::a2a::rpc::StreamingEvent::Message { message, .. } => {
164                            if let Some(text) = message.parts.iter().find_map(|p| p.as_text()) {
165                                println!("Agent: {text}");
166                            }
167                        }
168                        crate::a2a::rpc::StreamingEvent::TaskStatus { status, .. } => {
169                            println!("Status: {:?}", status.state);
170                            if let Some(msg) = status.message
171                                && let Some(text) = msg.parts.iter().find_map(|p| p.as_text())
172                            {
173                                println!("  Message: {text}");
174                            }
175                        }
176                        crate::a2a::rpc::StreamingEvent::TaskArtifact { artifact, .. } => {
177                            println!("Artifact: {}", artifact.id);
178                        }
179                        crate::a2a::rpc::StreamingEvent::Unknown => {}
180                    }
181                }
182                Err(e) => {
183                    eprintln!("Stream error: {e}");
184                    break;
185                }
186            }
187        }
188    } else {
189        println!("Sending task...\n");
190        let task = client.send_message(params).await?;
191        println!("Task created: {}", task.id);
192        println!("Status: {:?}\n", task.status.state);
193
194        if let Some(msg) = &task.status.message
195            && let Some(text) = msg.parts.iter().find_map(|p| p.as_text())
196        {
197            println!("Response: {text}");
198        }
199
200        if !task.artifacts.is_empty() {
201            println!("\nArtifacts:");
202            for artifact in &task.artifacts {
203                println!("  - {} ({} parts)", artifact.id, artifact.parts.len());
204            }
205        }
206    }
207
208    Ok(())
209}
210
211/// List tasks from a remote A2A agent
212async fn list_agent_tasks(agent_url: String, context_id: Option<String>, limit: u32) -> anyhow::Result<()> {
213    use crate::a2a::rpc::ListTasksParams;
214    use serde_json::Value;
215
216    println!("Fetching tasks from: {agent_url}");
217
218    let client = a2a_client(&agent_url)?;
219    let mut params = ListTasksParams::default();
220
221    if let Some(ctx_id) = context_id {
222        params.context_id = Some(ctx_id);
223    }
224    params.page_size = Some(limit);
225
226    let result_value: Value = client.list_tasks(Some(params)).await?;
227
228    // Parse the JSON result
229    let tasks_array = result_value
230        .get("tasks")
231        .and_then(|v| v.as_array())
232        .ok_or_else(|| anyhow::anyhow!("Invalid response format"))?;
233
234    let total_size = result_value
235        .get("totalSize")
236        .and_then(|v| v.as_u64())
237        .unwrap_or(tasks_array.len() as u64);
238
239    println!("\nTasks ({} total, showing {}):", total_size, tasks_array.len());
240    println!("═══════════════════════════════════════════════════════════════\n");
241
242    for task_value in tasks_array {
243        if let Some(task_id) = task_value.get("id").and_then(|v| v.as_str()) {
244            println!("Task: {task_id}");
245        }
246        if let Some(status) = task_value.get("status")
247            && let Some(state) = status.get("state").and_then(|v| v.as_str())
248        {
249            println!("  Status: {state}");
250        }
251        if let Some(ctx_id) = task_value.get("contextId").and_then(|v| v.as_str()) {
252            println!("  Context: {ctx_id}");
253        }
254        if let Some(artifacts) = task_value.get("artifacts").and_then(|v| v.as_array()) {
255            println!("  Artifacts: {}", artifacts.len());
256        }
257        println!();
258    }
259
260    Ok(())
261}
262
263/// Get details about a specific task
264async fn get_agent_task(agent_url: String, task_id: String) -> anyhow::Result<()> {
265    println!("Fetching task {task_id} from: {agent_url}");
266
267    let client = a2a_client(&agent_url)?;
268    let task = client.get_task(task_id.clone()).await?;
269
270    println!("\n═══════════════════════════════════════════════════════════════");
271    println!("Task: {}", task.id);
272    println!("═══════════════════════════════════════════════════════════════\n");
273
274    println!("Status: {:?}", task.status.state);
275    if let Some(ctx_id) = &task.context_id {
276        println!("Context: {ctx_id}");
277    }
278
279    if let Some(msg) = &task.status.message {
280        println!("\nLatest Message:");
281        println!("  Role: {:?}", msg.role);
282        for part in &msg.parts {
283            match part {
284                crate::a2a::types::Part::Text { text } => println!("  Text: {text}"),
285                crate::a2a::types::Part::File { file } => println!("  File: {file:?}"),
286                crate::a2a::types::Part::Data { data } => println!("  Data: {data}"),
287                crate::a2a::types::Part::Unknown => {}
288            }
289        }
290    }
291
292    if !task.artifacts.is_empty() {
293        println!("\nArtifacts:");
294        for artifact in &task.artifacts {
295            println!("  - {}:", artifact.id);
296            for part in &artifact.parts {
297                match part {
298                    crate::a2a::types::Part::Text { text } => {
299                        let preview = vtcode_commons::formatting::truncate_byte_budget(text, 60, "...");
300                        println!("    Text: {preview}");
301                    }
302                    crate::a2a::types::Part::File { file } => println!("    File: {file:?}"),
303                    crate::a2a::types::Part::Data { data } => {
304                        let s = data.to_string();
305                        let preview = vtcode_commons::formatting::truncate_byte_budget(&s, 60, "...");
306                        println!("    Data: {preview}");
307                    }
308                    crate::a2a::types::Part::Unknown => {}
309                }
310            }
311        }
312    }
313
314    if !task.history.is_empty() {
315        println!("\nHistory ({} messages):", task.history.len());
316        for (i, msg) in task.history.iter().enumerate() {
317            println!("  {}. {:?}: {} parts", i + 1, msg.role, msg.parts.len());
318        }
319    }
320
321    Ok(())
322}
323
324/// Cancel a running task
325async fn cancel_agent_task(agent_url: String, task_id: String) -> anyhow::Result<()> {
326    println!("Canceling task {task_id} at: {agent_url}");
327
328    let client = a2a_client(&agent_url)?;
329    client.cancel_task(task_id).await?;
330
331    println!("Task cancellation requested successfully.");
332
333    Ok(())
334}
335
336const A2A_TOKEN_ENV_VAR: &str = "VTCODE_A2A_TOKEN";
337
338fn a2a_client(agent_url: &str) -> anyhow::Result<crate::a2a::A2aClient> {
339    let client = crate::a2a::A2aClient::new(agent_url)?;
340    match std::env::var(A2A_TOKEN_ENV_VAR) {
341        Ok(auth_token) => Ok(client.with_bearer_token(auth_token)),
342        Err(std::env::VarError::NotPresent) => Ok(client),
343        Err(error) => anyhow::bail!("Failed to read {A2A_TOKEN_ENV_VAR}: {error}"),
344    }
345}
346
347#[cfg(test)]
348mod tests {
349    // Intentionally empty or specialized imports if needed later
350
351    #[tokio::test]
352    async fn test_discover_agent_display() {
353        // This is a simple display test - actual client functionality is tested in integration tests
354    }
355}