use crate::a2a::peer::{self as a2a};
use crate::config::A2aEndpoint;
use crate::json::{Id, Request};
use a2a_rs::domain::TaskState;
use serde_json::{Value, json};
use std::time::{Duration, Instant};
const POLL_INTERVAL: Duration = Duration::from_millis(100);
const REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
pub enum DelegateOutcome {
Distillate(String),
Error(String),
}
pub fn delegate(
endpoint: &A2aEndpoint,
auth: PeerAuth,
objective: &str,
command: Option<&Value>,
output_contract: Option<&str>,
message_id: Option<&str>,
deadline: Instant,
) -> DelegateOutcome {
match endpoint {
A2aEndpoint::Https(url) => match HttpEp::parse(url) {
Ok(ep) => {
let mut conn = HttpConn::new(ep, auth);
match conn.call_streaming(objective, command, output_contract, message_id, deadline)
{
Err(e) => DelegateOutcome::Error(e),
Ok(StreamOutcome::Done(outcome)) => outcome,
Ok(StreamOutcome::Recover(task_id)) => poll_task(&mut conn, &task_id, deadline),
}
}
Err(e) => DelegateOutcome::Error(e),
},
}
}
pub fn send(
endpoint: &A2aEndpoint,
auth: PeerAuth,
parts: &Value,
context: Option<&str>,
message_id: Option<&str>,
deadline: Instant,
) -> Result<Value, String> {
match endpoint {
A2aEndpoint::Https(url) => {
let ep = HttpEp::parse(url)?;
let mut conn = HttpConn::new(ep, auth);
let mut message = serde_json::Map::new();
message.insert(
"messageId".into(),
Value::String(
message_id
.map(str::to_string)
.unwrap_or_else(mint_message_id),
),
);
message.insert("role".into(), Value::String("user".into()));
message.insert(
"parts".into(),
match parts {
Value::Array(_) => parts.clone(),
Value::String(t) => serde_json::json!([{"text": t}]),
Value::Null => serde_json::json!([]),
other => serde_json::json!([{"data": other}]),
},
);
if let Some(ctx) = context.filter(|c| !c.is_empty()) {
message.insert("contextId".into(), Value::String(ctx.to_string()));
}
conn.call(
"SendMessage",
serde_json::json!({"message": Value::Object(message)}),
deadline,
)
}
}
}
enum StreamOutcome {
Done(DelegateOutcome),
Recover(String),
}
fn poll_task<C: Caller>(conn: &mut C, task_id: &str, deadline: Instant) -> DelegateOutcome {
let get_params = json!({ "id": task_id });
loop {
if Instant::now() >= deadline {
return DelegateOutcome::Error(format!(
"a2a: delegation to peer timed out (task {task_id} still running)"
));
}
std::thread::sleep(POLL_INTERVAL);
let task = match conn.call("GetTask", get_params.clone(), deadline) {
Ok(t) => t,
Err(e) => return DelegateOutcome::Error(e),
};
if let Some(outcome) = terminal_outcome(&task) {
return outcome;
}
}
}
#[derive(Default)]
pub struct PeerAuth {
pub headers: Vec<(String, String)>,
#[cfg(feature = "tls")]
pub identity: Option<crate::net::tls::ClientIdentity>,
pub signer: Option<std::sync::Arc<dyn ::mcp::http::RequestSigner>>,
}
trait Caller {
fn call(&mut self, method: &str, params: Value, deadline: Instant) -> Result<Value, String>;
}
fn terminal_outcome(task: &Value) -> Option<DelegateOutcome> {
let state = a2a::task_state_of(task);
if !a2a::is_terminal(state) {
return None;
}
Some(match state {
TaskState::TASK_STATE_COMPLETED => DelegateOutcome::Distillate(a2a::artifact_text_of(task)),
TaskState::TASK_STATE_REJECTED => {
DelegateOutcome::Error("a2a: remote agent rejected the objective".into())
}
TaskState::TASK_STATE_CANCELED => {
DelegateOutcome::Error("a2a: remote task was canceled".into())
}
_ => DelegateOutcome::Error("a2a: remote task failed".into()),
})
}
struct HttpEp {
host: String,
port: u16,
path: String,
host_header: String,
tls: bool,
socket: Option<String>,
}
impl HttpEp {
fn parse(url: &str) -> Result<HttpEp, String> {
if let Some(path) = url
.strip_prefix("unix://")
.or_else(|| url.strip_prefix("unix:"))
{
if path.is_empty() {
return Err(format!("a2a: unix peer needs a socket path: {url}"));
}
return Ok(HttpEp {
host: String::new(),
port: 0,
path: "/".to_string(),
host_header: "localhost".to_string(),
tls: false,
socket: Some(path.to_string()),
});
}
let u = crate::net::http::Url::parse(url)
.map_err(|e| format!("a2a: bad peer url {url}: {e}"))?;
let path = if u.path.is_empty() || u.path == "/" {
"/".to_string()
} else {
u.path.clone()
};
Ok(HttpEp {
host_header: u.host_header(),
tls: u.is_tls(),
host: u.host,
port: u.port,
path,
socket: None,
})
}
}
struct HttpConn {
ep: HttpEp,
auth: PeerAuth,
next_id: i64,
}
impl HttpConn {
fn new(ep: HttpEp, auth: PeerAuth) -> HttpConn {
HttpConn {
ep,
auth,
next_id: 1,
}
}
fn signature_headers(&self, body: &[u8]) -> Vec<(String, String)> {
let mut out = Vec::new();
#[cfg(feature = "aauth")]
if let Some(signer) = crate::aauth::signer() {
out.extend(signer.sign("POST", &self.ep.host_header, &self.ep.path, body));
}
if let Some(signer) = &self.auth.signer {
out.extend(signer.sign("POST", &self.ep.host_header, &self.ep.path, body));
}
out
}
fn connect(&self, timeout: Duration) -> Result<Box<dyn crate::net::http::Stream>, String> {
if let Some(socket) = &self.ep.socket {
let s = crate::net::unixsock::connect(socket, timeout)
.map_err(|e| format!("a2a: cannot reach peer socket {socket}: {e}"))?;
return Ok(Box::new(s));
}
let tcp = crate::net::http::connect_tcp(&self.ep.host, self.ep.port, timeout)
.map_err(|e| format!("a2a: cannot reach peer {}: {e}", self.ep.host))?;
if self.ep.tls {
#[cfg(feature = "tls")]
{
let tls = crate::net::tls::connect(tcp, &self.ep.host, self.auth.identity.as_ref())
.map_err(|e| format!("a2a: tls to peer {}: {e}", self.ep.host))?;
Ok(Box::new(tls))
}
#[cfg(not(feature = "tls"))]
{
Err("a2a: https peer requires the 'tls' build feature".to_string())
}
} else {
Ok(Box::new(tcp))
}
}
}
impl HttpConn {
fn call_streaming(
&mut self,
objective: &str,
command: Option<&Value>,
output_contract: Option<&str>,
explicit_message_id: Option<&str>,
deadline: Instant,
) -> Result<StreamOutcome, String> {
let id = self.next_id;
self.next_id += 1;
let message_id = explicit_message_id
.map(str::to_string)
.unwrap_or_else(mint_message_id);
let params = a2a::send_message_params_cmd(objective, command, output_contract, &message_id);
let req = Request::new(Id::Num(id), "SendStreamingMessage", Some(params));
let body =
serde_json::to_vec(&req).map_err(|e| format!("a2a: encode streaming send: {e}"))?;
let remaining = deadline.saturating_duration_since(Instant::now());
let timeout = remaining
.min(Duration::from_secs(45))
.max(Duration::from_millis(1));
let stream = self.connect(timeout)?;
let sig = self.signature_headers(&body);
let mut headers: Vec<(&str, &str)> = vec![
("Content-Type", "application/json"),
("Accept", "text/event-stream"),
];
for (name, value) in &self.auth.headers {
headers.push((name.as_str(), value.as_str()));
}
for (name, value) in &sig {
headers.push((name.as_str(), value.as_str()));
}
let resp = crate::net::http::send_streaming(
stream,
&self.ep.host_header,
"POST",
&self.ep.path,
&headers,
&body,
)
.map_err(|e| format!("a2a: streaming send: {e}"))?;
if resp.status != 200 {
return Err(format!("a2a: SendStreamingMessage HTTP {}", resp.status));
}
let sse = resp
.header("content-type")
.is_some_and(|ct| ct.to_ascii_lowercase().contains("text/event-stream"));
if !sse {
use std::io::Read as _;
let mut text = String::new();
let _ = resp.into_reader().take(1 << 20).read_to_string(&mut text);
let frame: crate::json::Response = serde_json::from_str(text.trim())
.map_err(|e| format!("a2a: bad unary streaming reply: {e}"))?;
if let Some(err) = frame.error {
return Err(format!(
"a2a: streaming rpc error {}: {}",
err.code, err.message
));
}
let result = frame.result.unwrap_or(Value::Null);
if let Some(outcome) = terminal_outcome(&result) {
return Ok(StreamOutcome::Done(outcome));
}
let task_id = result
.pointer("/statusUpdate/taskId")
.and_then(Value::as_str)
.map(str::to_string)
.or_else(|| Some(a2a::task_id_of(&result)).filter(|s| !s.is_empty()));
return match task_id {
Some(tid) => Ok(StreamOutcome::Recover(tid)),
None => Err("a2a: unary streaming reply named no task".into()),
};
}
let mut events = resp.sse();
let mut task_id: Option<String> = None;
let mut distillate: Option<String> = None;
loop {
if Instant::now() >= deadline {
return match task_id {
Some(tid) => Ok(StreamOutcome::Recover(tid)),
None => Err("a2a: deadline while streaming".into()),
};
}
let ev = match events.next_event() {
Ok(Some(ev)) => ev,
Ok(None) => {
return match task_id {
Some(tid) => Ok(StreamOutcome::Recover(tid)),
None => Err("a2a: stream ended before any frame".into()),
};
}
Err(e) => {
return match task_id {
Some(tid) => Ok(StreamOutcome::Recover(tid)),
None => Err(format!("a2a: stream read: {e}")),
};
}
};
if ev.data.trim().is_empty() {
continue;
}
let Ok(frame) = serde_json::from_str::<crate::json::Response>(ev.data.trim()) else {
continue; };
if let Some(err) = frame.error {
return Ok(StreamOutcome::Done(DelegateOutcome::Error(format!(
"a2a: streaming rpc error {}: {}",
err.code, err.message
))));
}
let result = frame.result.unwrap_or(Value::Null);
if let Some(update) = result.get("statusUpdate") {
if let Some(tid) = update.get("taskId").and_then(Value::as_str) {
task_id = Some(tid.to_string());
}
let state = update
.pointer("/status/state")
.and_then(Value::as_str)
.unwrap_or("");
if a2a::is_terminal(
serde_json::from_value(json!(state))
.unwrap_or(TaskState::TASK_STATE_UNSPECIFIED),
) {
let outcome = match state {
"TASK_STATE_COMPLETED" => match distillate {
Some(text) => DelegateOutcome::Distillate(text),
None => {
return Ok(match &task_id {
Some(tid) => StreamOutcome::Recover(tid.clone()),
None => StreamOutcome::Done(DelegateOutcome::Error(
"a2a: remote completed without a distillate artifact"
.into(),
)),
});
}
},
"TASK_STATE_REJECTED" => DelegateOutcome::Error(
"a2a: remote agent rejected the objective".into(),
),
"TASK_STATE_CANCELLED" | "TASK_STATE_CANCELED" => {
DelegateOutcome::Error("a2a: remote task was canceled".into())
}
other => DelegateOutcome::Error(format!("a2a: remote task ended {other}")),
};
return Ok(StreamOutcome::Done(outcome));
}
} else if let Some(update) = result.get("artifactUpdate") {
if let Some(text) = update
.pointer("/artifact/parts/0/text")
.and_then(Value::as_str)
{
distillate = Some(text.to_string());
}
}
else if !a2a::task_id_of(&result).is_empty() {
task_id = Some(a2a::task_id_of(&result));
}
}
}
}
impl Caller for HttpConn {
fn call(&mut self, method: &str, params: Value, deadline: Instant) -> Result<Value, String> {
let id = self.next_id;
self.next_id += 1;
let timeout = request_timeout(deadline);
let req = Request::new(Id::Num(id), method, Some(params));
let body = serde_json::to_vec(&req).map_err(|e| format!("a2a: encode {method}: {e}"))?;
let mut stream = self.connect(timeout)?;
let sig = self.signature_headers(&body);
let mut headers: Vec<(&str, &str)> = vec![("Content-Type", "application/json")];
for (name, value) in &self.auth.headers {
headers.push((name.as_str(), value.as_str()));
}
for (name, value) in &sig {
headers.push((name.as_str(), value.as_str()));
}
let resp = crate::net::http::send(
&mut *stream,
&self.ep.host_header,
"POST",
&self.ep.path,
&headers,
&body,
)
.map_err(|e| format!("a2a: {method}: {e}"))?;
if !resp.is_success() {
return Err(format!("a2a: {method} HTTP {}", resp.status));
}
let response: crate::json::Response = serde_json::from_slice(&resp.body)
.map_err(|e| format!("a2a: {method} bad reply: {e}"))?;
if let Some(err) = response.error {
return Err(format!(
"a2a: {method} rpc error {}: {}",
err.code, err.message
));
}
Ok(response.result.unwrap_or(Value::Null))
}
}
fn mint_message_id() -> String {
use std::sync::atomic::{AtomicU64, Ordering};
static SEQ: AtomicU64 = AtomicU64::new(0);
let n = SEQ.fetch_add(1, Ordering::Relaxed);
let millis = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis();
format!("a2a-msg-{millis:x}-{n:x}")
}
fn request_timeout(deadline: Instant) -> Duration {
let remaining = deadline.saturating_duration_since(Instant::now());
remaining.min(REQUEST_TIMEOUT).max(Duration::from_millis(1))
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::{BufRead, BufReader, Write};
use std::thread;
fn task(id: &str, state: TaskState, artifact: Option<&str>) -> Value {
let mut t = json!({
"id": id,
"contextId": format!("ctx-{id}"),
"status": { "state": state, "timestamp": "1970-01-01T00:00:00.000Z" },
});
if let Some(text) = artifact {
t["artifacts"] =
json!([{ "artifactId": format!("{id}.distillate"), "parts": [{ "text": text }] }]);
}
t
}
fn serve_http_fixture(replies: Vec<Value>) -> String {
serve_http_impl(replies, false)
}
fn serve_http_error_fixture() -> String {
serve_http_impl(vec![json!({"__rpc_error__": true})], true)
}
fn serve_http_impl(replies: Vec<Value>, error: bool) -> String {
use std::io::Read as _;
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
thread::spawn(move || {
let mut idx = 0usize;
for conn in listener.incoming() {
let Ok(mut stream) = conn else { continue };
stream.set_read_timeout(Some(Duration::from_secs(2))).ok();
let mut reader = BufReader::new(stream.try_clone().unwrap());
let mut len = 0usize;
let mut req_id = json!(1);
loop {
let mut line = String::new();
if reader.read_line(&mut line).unwrap_or(0) == 0 || line.trim().is_empty() {
break;
}
if let Some((k, v)) = line.split_once(':')
&& k.trim().eq_ignore_ascii_case("content-length")
{
len = v.trim().parse().unwrap_or(0);
}
}
let mut body = vec![0u8; len];
if reader.read_exact(&mut body).is_ok()
&& let Ok(rpc) = serde_json::from_slice::<Value>(&body)
{
req_id = rpc["id"].clone();
}
if replies.is_empty() {
break;
}
let payload = if error {
json!({"jsonrpc": "2.0", "id": req_id, "error": {"code": -32602, "message": "bad params"}})
} else {
let result = replies[idx.min(replies.len() - 1)].clone();
json!({"jsonrpc": "2.0", "id": req_id, "result": result})
};
idx += 1;
let text = serde_json::to_vec(&payload).unwrap();
let head = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
text.len()
);
let _ = stream.write_all(head.as_bytes());
let _ = stream.write_all(&text);
let _ = stream.flush();
}
});
format!("http://{addr}")
}
fn serve_sse_fixture(frames: Vec<Value>, unary: Vec<Value>) -> String {
use std::io::Read as _;
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
thread::spawn(move || {
let mut first = true;
let mut uidx = 0usize;
for conn in listener.incoming() {
let Ok(mut stream) = conn else { continue };
stream.set_read_timeout(Some(Duration::from_secs(2))).ok();
let mut reader = BufReader::new(stream.try_clone().unwrap());
let mut len = 0usize;
let mut req_id = json!(1);
loop {
let mut line = String::new();
if reader.read_line(&mut line).unwrap_or(0) == 0 || line.trim().is_empty() {
break;
}
if let Some((k, v)) = line.split_once(':')
&& k.trim().eq_ignore_ascii_case("content-length")
{
len = v.trim().parse().unwrap_or(0);
}
}
let mut body = vec![0u8; len];
if reader.read_exact(&mut body).is_ok()
&& let Ok(rpc) = serde_json::from_slice::<Value>(&body)
{
req_id = rpc["id"].clone();
}
if first {
first = false;
let _ = stream.write_all(
b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nConnection: close\r\n\r\n",
);
let _ = stream.write_all(b": keep-alive\n\n");
for f in &frames {
let env = json!({"jsonrpc": "2.0", "id": req_id, "result": f});
let _ = stream.write_all(format!("data: {env}\n\n").as_bytes());
}
let _ = stream.flush();
continue; }
if unary.is_empty() {
break;
}
let result = unary[uidx.min(unary.len() - 1)].clone();
uidx += 1;
let payload = json!({"jsonrpc": "2.0", "id": req_id, "result": result}).to_string();
let head = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
payload.len()
);
let _ = stream.write_all(head.as_bytes());
let _ = stream.write_all(payload.as_bytes());
let _ = stream.flush();
}
});
format!("http://{addr}")
}
fn status_frame(id: &str, state: &str, _is_final: bool) -> Value {
json!({"statusUpdate": {"taskId": id, "contextId": "ctx", "status": {"state": state}}})
}
#[test]
fn delegate_consumes_an_sse_stream_to_the_distillate_without_polling() {
let url = serve_sse_fixture(
vec![
status_frame("s-1", "TASK_STATE_WORKING", false),
json!({"artifactUpdate": {"taskId": "s-1", "contextId": "ctx", "artifact": {"artifactId": "s-1.distillate", "parts": [{"text": "streamed answer"}]}, "lastChunk": true}}),
status_frame("s-1", "TASK_STATE_COMPLETED", true),
],
Vec::new(), );
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
match delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline) {
DelegateOutcome::Distillate(s) => assert_eq!(s, "streamed answer"),
DelegateOutcome::Error(e) => panic!("expected streamed distillate: {e}"),
}
}
#[test]
fn a_broken_stream_recovers_over_get_task() {
let url = serve_sse_fixture(
vec![status_frame("s-2", "TASK_STATE_WORKING", false)],
vec![task(
"s-2",
TaskState::TASK_STATE_COMPLETED,
Some("recovered answer"),
)],
);
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
match delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline) {
DelegateOutcome::Distillate(s) => assert_eq!(s, "recovered answer"),
DelegateOutcome::Error(e) => panic!("expected recovery: {e}"),
}
}
#[test]
fn a_final_failed_stream_frame_is_a_terminal_error() {
let url = serve_sse_fixture(
vec![
status_frame("s-3", "TASK_STATE_WORKING", false),
status_frame("s-3", "TASK_STATE_FAILED", true),
],
Vec::new(),
);
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
match delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline) {
DelegateOutcome::Error(e) => assert!(e.contains("FAILED"), "{e}"),
DelegateOutcome::Distillate(s) => panic!("expected error, got: {s}"),
}
}
#[test]
fn delegate_over_http_send_then_poll_returns_the_distillate() {
let url = serve_http_fixture(vec![
task("h-1", TaskState::TASK_STATE_WORKING, None),
task("h-1", TaskState::TASK_STATE_WORKING, None),
task(
"h-1",
TaskState::TASK_STATE_COMPLETED,
Some("http distilled answer"),
),
]);
let ep = A2aEndpoint::parse(&url).expect("parse https endpoint");
assert!(matches!(ep, A2aEndpoint::Https(_)));
let deadline = Instant::now() + Duration::from_secs(5);
match delegate(
&ep,
PeerAuth::default(),
"do the work",
None,
Some("one line"),
None,
deadline,
) {
DelegateOutcome::Distillate(s) => assert_eq!(s, "http distilled answer"),
DelegateOutcome::Error(e) => panic!("expected distillate, got error: {e}"),
}
}
#[test]
fn delegate_presents_the_peer_auth_headers() {
use std::sync::{Arc, Mutex};
let captured: Arc<Mutex<String>> = Arc::default();
let cap = Arc::clone(&captured);
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
thread::spawn(move || {
if let Some(Ok(mut stream)) = listener.incoming().next() {
stream.set_read_timeout(Some(Duration::from_secs(2))).ok();
let mut reader = BufReader::new(stream.try_clone().unwrap());
let mut head = String::new();
loop {
let mut line = String::new();
if reader.read_line(&mut line).unwrap_or(0) == 0 || line.trim().is_empty() {
break;
}
head.push_str(&line);
}
*cap.lock().unwrap() = head;
let result = task("h-a", TaskState::TASK_STATE_COMPLETED, Some("authed"));
let payload = json!({"jsonrpc": "2.0", "id": 1, "result": result}).to_string();
let head = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
payload.len()
);
let _ = stream.write_all(head.as_bytes());
let _ = stream.write_all(payload.as_bytes());
let _ = stream.flush();
}
});
let ep = A2aEndpoint::parse(&format!("http://{addr}")).unwrap();
let auth = PeerAuth {
headers: vec![("authorization".into(), "Bearer sekrit-token".into())],
..Default::default()
};
let deadline = Instant::now() + Duration::from_secs(5);
match delegate(&ep, auth, "obj", None, None, None, deadline) {
DelegateOutcome::Distillate(s) => assert_eq!(s, "authed"),
DelegateOutcome::Error(e) => panic!("unexpected error: {e}"),
}
let head = captured.lock().unwrap().clone();
assert!(
head.to_lowercase()
.contains("authorization: bearer sekrit-token"),
"the bearer header was presented to the peer:\n{head}"
);
}
#[test]
fn delegate_over_http_send_message_already_terminal_skips_polling() {
let url = serve_http_fixture(vec![task(
"h-t",
TaskState::TASK_STATE_COMPLETED,
Some("immediate"),
)]);
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
match delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline) {
DelegateOutcome::Distillate(s) => assert_eq!(s, "immediate"),
DelegateOutcome::Error(e) => panic!("unexpected error: {e}"),
}
}
#[test]
fn delegate_over_http_surfaces_a_failed_task() {
let url = serve_http_fixture(vec![task("h-2", TaskState::TASK_STATE_FAILED, None)]);
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_secs(5);
assert!(matches!(
delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline),
DelegateOutcome::Error(_)
));
}
#[test]
fn delegate_over_http_deadline_while_polling_is_a_timeout() {
let url = serve_http_fixture(vec![task("h-w", TaskState::TASK_STATE_WORKING, None)]);
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_millis(300);
match delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline) {
DelegateOutcome::Error(e) => assert!(e.contains("timed out"), "got: {e}"),
DelegateOutcome::Distillate(s) => panic!("expected timeout, got: {s}"),
}
}
#[test]
fn delegate_over_http_surfaces_a_peer_rpc_error() {
let url = serve_http_error_fixture();
let ep = A2aEndpoint::parse(&url).unwrap();
let deadline = Instant::now() + Duration::from_secs(2);
match delegate(&ep, PeerAuth::default(), "obj", None, None, None, deadline) {
DelegateOutcome::Error(e) => assert!(e.contains("rpc error"), "got: {e}"),
DelegateOutcome::Distillate(s) => panic!("expected error, got: {s}"),
}
}
}