1use crate::a2a::cli::A2aCommands;
10
11pub 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#[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
77async 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
133async 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 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
211async 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 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
263async 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
324async 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 #[tokio::test]
352 async fn test_discover_agent_display() {
353 }
355}