use std::time::Duration;
use futures::StreamExt;
use a2a_rs::domain::{Message, SendCompletion};
use a2a_rs::{JsonRpcClient, StreamItem, Transport, connect, default_registry};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let base_url = std::env::args()
.nth(1)
.unwrap_or_else(|| "http://127.0.0.1:8137".to_string());
let transport: Box<dyn Transport> = match connect(&base_url, &default_registry()).await {
Ok(t) => {
println!("✅ negotiated transport: {}", t.protocol());
t
}
Err(e) => {
println!("⚠️ card negotiation failed ({e}); falling back to direct JSON-RPC client");
Box::new(JsonRpcClient::new(base_url.clone()))
}
};
let task = transport
.send_task_message(
"demo-task",
&Message::user_text("hello".to_string(), "m1".to_string()),
None,
None,
SendCompletion::WhenSettled,
)
.await?;
println!("📨 sent message; task id = {}", task.id);
let fetched = transport.get_task(&task.id, None).await?;
println!("📥 get_task → state {:?}", fetched.status.state);
let mut stream = transport.subscribe_to_task(&task.id, None, None).await?;
println!("📡 subscribing (up to 3 events / 5s)…");
for _ in 0..3 {
match tokio::time::timeout(Duration::from_secs(5), stream.next()).await {
Ok(Some(Ok(event))) => match &event.item {
StreamItem::Task(t) => {
println!(" • snapshot: task {} ({:?})", t.id, t.status.state)
}
StreamItem::StatusUpdate(u) => println!(" • status: {:?}", u.status.state),
StreamItem::ArtifactUpdate(_) => println!(" • artifact update"),
},
Ok(Some(Err(e))) => {
println!(" • stream error: {e}");
break;
}
Ok(None) => break, Err(_) => break, }
}
drop(stream);
let canceled = transport.cancel_task(&task.id).await?;
println!("🛑 canceled; final state {:?}", canceled.status.state);
Ok(())
}