use std::time::Duration;
use async_trait::async_trait;
use serde::Deserialize;
use serde_json::{Value, json};
use crate::core::Delegation;
use super::{PeerClient, PeerCredential, PeerError, PeerId};
#[derive(Debug, Clone)]
pub struct Endpoint {
pub url: String,
pub timeout: Duration,
}
impl Endpoint {
pub const DEFAULT_TIMEOUT: Duration = Duration::from_mins(2);
pub fn new(url: impl Into<String>) -> Self {
Self {
url: url.into(),
timeout: Self::DEFAULT_TIMEOUT,
}
}
#[must_use]
pub const fn timeout(mut self, d: Duration) -> Self {
self.timeout = d;
self
}
}
pub const EXTENSION_URI: &str = "https://hupe1980.github.io/agentplane/a2a/delegation/v1";
#[derive(Debug)]
pub struct A2aClient {
egress: Option<crate::core::Egress>,
http: reqwest::Client,
endpoint: Endpoint,
}
impl A2aClient {
#[must_use]
pub fn egress(mut self, egress: crate::core::Egress) -> Self {
self.egress = Some(egress);
self
}
pub fn new(endpoint: Endpoint) -> Result<Self, PeerError> {
let http = reqwest::Client::builder()
.timeout(endpoint.timeout)
.build()
.map_err(|e| PeerError::Unreachable {
peer: PeerId::new("<local>"),
detail: format!("could not build an HTTP client: {e}"),
})?;
Ok(Self {
egress: None,
http,
endpoint,
})
}
fn body(
capability: &str,
payload: &Value,
acting_as: &Delegation,
provenance: Option<&crate::core::Provenance>,
) -> Value {
let attested = provenance.map(|p| Value::Object(p.to_meta()));
json!({
"jsonrpc": "2.0",
"id": 1,
"method": "message/send",
"params": {
"message": {
"role": "user",
"messageId": capability,
"parts": [{ "kind": "data", "data": payload }],
"extensions": [EXTENSION_URI],
"metadata": {
EXTENSION_URI: {
"capability": capability,
"chain": acting_as,
"provenance": attested,
}
}
}
}
})
}
}
#[derive(Debug, Deserialize)]
struct RpcError {
code: i64,
message: String,
}
#[derive(Debug, Deserialize)]
struct RpcResponse {
#[serde(default)]
result: Option<Value>,
#[serde(default)]
error: Option<RpcError>,
}
fn classify_rpc(peer: &PeerId, e: &RpcError) -> PeerError {
let detail = format!("{} (code {})", e.message, e.code);
match e.code {
-32700 | -32600 | -32601 | -32602 | -32006..=-32001 => PeerError::Refused {
peer: peer.clone(),
detail,
},
_ => PeerError::TimedOut {
peer: peer.clone(),
detail: format!("{detail} — the peer did not say whether it acted"),
},
}
}
fn classify_status(peer: &PeerId, status: reqwest::StatusCode) -> PeerError {
if status.is_server_error() {
return PeerError::TimedOut {
peer: peer.clone(),
detail: format!("HTTP {status} — the peer did not say whether it acted"),
};
}
PeerError::Refused {
peer: peer.clone(),
detail: format!("HTTP {status}"),
}
}
fn classify_transport(peer: &PeerId, e: &reqwest::Error) -> PeerError {
if e.is_connect() {
return PeerError::Unreachable {
peer: peer.clone(),
detail: format!("could not connect: {e}"),
};
}
if e.is_timeout() {
return PeerError::TimedOut {
peer: peer.clone(),
detail: format!("timed out: {e}"),
};
}
if e.is_body() || e.is_decode() {
return PeerError::TimedOut {
peer: peer.clone(),
detail: format!("the response could not be read: {e}"),
};
}
if e.is_request() {
return PeerError::TimedOut {
peer: peer.clone(),
detail: format!("the request failed in flight: {e}"),
};
}
PeerError::TimedOut {
peer: peer.clone(),
detail: e.to_string(),
}
}
fn task_failure(peer: &PeerId, result: &Value) -> Option<PeerError> {
let state = result.get("status")?.get("state")?.as_str()?;
match state {
"failed" => Some(PeerError::Failed {
peer: peer.clone(),
detail: result
.get("status")
.and_then(|s| s.get("message"))
.map_or_else(
|| "the peer reported the task failed".to_owned(),
std::string::ToString::to_string,
),
}),
"rejected" => Some(PeerError::Refused {
peer: peer.clone(),
detail: "the peer rejected the task".to_owned(),
}),
_ => None,
}
}
#[async_trait]
impl PeerClient for A2aClient {
async fn send(
&self,
peer: &PeerId,
capability: &str,
payload: &Value,
acting_as: &Delegation,
credential: Option<&PeerCredential>,
provenance: Option<&crate::core::Provenance>,
) -> Result<Value, PeerError> {
if let Some(egress) = &self.egress {
let host = reqwest::Url::parse(&self.endpoint.url)
.ok()
.and_then(|u| u.host_str().map(ToOwned::to_owned));
if let Err(e) = egress.permits(host.as_deref()) {
return Err(PeerError::Refused {
peer: peer.clone(),
detail: e.to_string(),
});
}
}
let mut req = self
.http
.post(&self.endpoint.url)
.json(&Self::body(capability, payload, acting_as, provenance));
if let Some(c) = credential {
req = req.bearer_auth(c.expose());
}
let response = req.send().await.map_err(|e| classify_transport(peer, &e))?;
let status = response.status();
let body: Result<RpcResponse, _> = response.json().await;
let Ok(rpc) = body else {
return Err(classify_status(peer, status));
};
if let Some(e) = rpc.error {
return Err(classify_rpc(peer, &e));
}
if !status.is_success() {
return Err(classify_status(peer, status));
}
let result = rpc.result.unwrap_or(Value::Null);
if let Some(failure) = task_failure(peer, &result) {
return Err(failure);
}
Ok(result)
}
}