use std::collections::BTreeMap;
use std::fmt::Debug;
use std::sync::Arc;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;
pub const PROTOCOL_VERSION: &str = "1.0";
#[cfg(any(feature = "a2a-server", all(feature = "a2a", feature = "manifest")))]
pub(crate) fn protocol_major_minor(version: &str) -> Option<(u64, u64)> {
let mut parts = version.split('.');
let major = parts.next()?.parse().ok()?;
let minor = parts.next()?.parse().ok()?;
match parts.next() {
None => Some((major, minor)),
Some(patch) if !patch.is_empty() && patch.parse::<u64>().is_ok() => {
if parts.next().is_none() {
Some((major, minor))
} else {
None
}
}
Some(_) => None,
}
}
#[cfg(feature = "a2a")]
pub mod a2a;
#[cfg(feature = "manifest")]
mod card;
#[cfg(feature = "manifest")]
mod card_sig;
#[cfg(all(feature = "manifest", feature = "a2a"))]
mod discovery;
#[cfg(feature = "manifest")]
pub use card::{
AgentCard, AgentExtension, CardCapabilities, CardInterface, CardSecurity,
CardSecurityRequirement, CardSecurityScheme, CardSkill, EXT_AGENT_DIRECTORY, EXT_GOVERNANCE,
EXT_MANIFEST_PROVENANCE, ExtendedAgentCard, ExtendedBudget, ExtendedTool,
HttpAuthSecurityScheme, SecurityScopeList, WELL_KNOWN_PATH, agent_card_path,
};
#[cfg(feature = "manifest")]
pub use card_sig::{
ALG, CardSignature, CardSignatureError, CardSigner, CardVerifier, signing_input,
};
#[cfg(all(feature = "manifest", feature = "a2a"))]
pub use discovery::{CardClient, DiscoveryError, JSONRPC};
mod credentials;
pub use credentials::{Cached, CredentialError, CredentialSource, Fixed, TokenExchange};
use std::time::Duration;
use crate::core::{
Delegation, DelegationError, Disposition, Effect, EffectDescriptor, EffectError, Principal,
Recovery, RetryPolicy, Scope, Secret, Sensitivity, Timestamp, Trust,
};
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub struct PeerId(pub String);
impl PeerId {
pub fn new(s: impl Into<String>) -> Self {
Self(s.into())
}
}
impl std::fmt::Display for PeerId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
#[derive(Clone, PartialEq, Eq)]
pub struct PeerCredential {
audience: PeerId,
secret: Secret,
expires_at: Option<Timestamp>,
}
impl PeerCredential {
pub fn for_audience(audience: PeerId, secret: impl Into<String>) -> Self {
Self {
audience,
secret: Secret::new(secret),
expires_at: None,
}
}
#[must_use]
pub const fn expiring_at(mut self, at: Timestamp) -> Self {
self.expires_at = Some(at);
self
}
#[must_use]
pub const fn expires_at(&self) -> Option<Timestamp> {
self.expires_at
}
#[must_use]
pub fn is_usable_at(&self, now: Timestamp, skew: Duration) -> bool {
let Some(expiry) = self.expires_at else {
return true;
};
let margin = i64::try_from(skew.as_secs()).unwrap_or(i64::MAX);
now.unix_timestamp().saturating_add(margin) < expiry.unix_timestamp()
}
#[must_use]
pub const fn audience(&self) -> &PeerId {
&self.audience
}
#[must_use]
pub fn expose(&self) -> &str {
self.secret.expose()
}
}
impl Debug for PeerCredential {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PeerCredential")
.field("audience", &self.audience)
.field("expires_at", &self.expires_at)
.field("secret", &"<redacted>")
.finish()
}
}
#[derive(Debug, Clone)]
pub struct PeerGrant {
pub scope: Scope,
pub mutates: bool,
pub recovery: Recovery,
pub max_sensitivity: Sensitivity,
pub output_sensitivity: Sensitivity,
pub retry: RetryPolicy,
credential: Option<PeerCredential>,
}
impl PeerGrant {
#[must_use]
pub fn new(scope: Scope) -> Self {
Self {
scope,
mutates: true,
recovery: Recovery::RequiresOperator,
max_sensitivity: Sensitivity::Public,
output_sensitivity: Sensitivity::Public,
retry: RetryPolicy::never(),
credential: None,
}
}
#[must_use]
pub fn with_credential(mut self, peer: &PeerId, credential: PeerCredential) -> Self {
assert_eq!(
credential.audience(),
peer,
"a credential for '{}' was attached to peer '{peer}'. An audience-bound \
credential presented to the wrong peer is exactly what binding exists \
to prevent",
credential.audience()
);
self.credential = Some(credential);
self
}
#[must_use]
pub fn read_only(mut self) -> Self {
self.mutates = false;
self.recovery = Recovery::Retry;
self
}
#[must_use]
pub const fn output_sensitivity(mut self, s: Sensitivity) -> Self {
self.output_sensitivity = s;
self
}
}
#[derive(Debug, thiserror::Error)]
pub enum PeerError {
#[error(
"peer '{peer}' is not registered; a peer nobody declared is a peer nobody granted anything"
)]
Unknown { peer: PeerId },
#[error("delegating to '{peer}' is refused: {source}")]
Delegation {
peer: PeerId,
#[source]
source: DelegationError,
},
#[error(
"the credential held for this call is bound to '{held_for}', not '{peer}' — \
presenting it would let '{peer}' replay it at '{held_for}'"
)]
WrongAudience { peer: PeerId, held_for: PeerId },
#[error("could not reach '{peer}': {detail}")]
Unreachable { peer: PeerId, detail: String },
#[error("'{peer}' refused the request: {detail}")]
Refused { peer: PeerId, detail: String },
#[error("'{peer}' did not answer in time: {detail}")]
TimedOut { peer: PeerId, detail: String },
#[error("'{peer}' returned an invalid response: {detail}")]
InvalidResponse { peer: PeerId, detail: String },
#[error("'{peer}' reported a failure: {detail}")]
Failed { peer: PeerId, detail: String },
}
impl PeerError {
#[must_use]
pub const fn disposition(&self) -> Disposition {
match self {
Self::Unknown { .. }
| Self::Delegation { .. }
| Self::WrongAudience { .. }
| Self::Unreachable { .. }
| Self::Refused { .. } => Disposition::DidNotHappen,
Self::TimedOut { .. } | Self::InvalidResponse { .. } => Disposition::InDoubt,
Self::Failed { .. } => Disposition::Landed,
}
}
}
#[async_trait]
pub trait PeerClient: Send + Sync + Debug {
async fn send(
&self,
peer: &PeerId,
capability: &str,
payload: &Value,
acting_as: &Delegation,
credential: Option<&PeerCredential>,
provenance: Option<&crate::core::Provenance>,
) -> Result<Value, PeerError>;
async fn get_task(
&self,
peer: &PeerId,
task_id: &str,
credential: Option<&PeerCredential>,
) -> Result<Value, PeerError> {
let _ = (task_id, credential);
Err(PeerError::Refused {
peer: peer.clone(),
detail: "this peer transport does not support task lookup".to_owned(),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PeerTask {
pub peer: PeerId,
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub context_id: Option<String>,
}
impl PeerTask {
pub fn from_response(peer: PeerId, response: &Value) -> Result<Option<Self>, PeerError> {
if response.get("role").is_some() {
return Ok(None);
}
let id = response.get("id").and_then(Value::as_str).ok_or_else(|| {
PeerError::InvalidResponse {
peer: peer.clone(),
detail: "task response has no string id".to_owned(),
}
})?;
if response
.get("status")
.and_then(|status| status.get("state"))
.and_then(Value::as_str)
.is_none()
{
return Err(PeerError::InvalidResponse {
peer,
detail: "task response has no status.state".to_owned(),
});
}
Ok(Some(Self {
peer,
id: id.to_owned(),
context_id: response
.get("contextId")
.and_then(Value::as_str)
.map(ToOwned::to_owned),
}))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum PeerTaskState {
Submitted,
Working,
Completed,
Failed,
Canceled,
Rejected,
InputRequired,
AuthRequired,
}
impl PeerTaskState {
fn parse(peer: &PeerId, value: &Value) -> Result<Self, PeerError> {
match value
.get("status")
.and_then(|status| status.get("state"))
.and_then(Value::as_str)
{
Some("TASK_STATE_SUBMITTED") => Ok(Self::Submitted),
Some("TASK_STATE_WORKING") => Ok(Self::Working),
Some("TASK_STATE_COMPLETED") => Ok(Self::Completed),
Some("TASK_STATE_FAILED") => Ok(Self::Failed),
Some("TASK_STATE_CANCELED") => Ok(Self::Canceled),
Some("TASK_STATE_REJECTED") => Ok(Self::Rejected),
Some("TASK_STATE_INPUT_REQUIRED") => Ok(Self::InputRequired),
Some("TASK_STATE_AUTH_REQUIRED") => Ok(Self::AuthRequired),
Some(other) => Err(PeerError::InvalidResponse {
peer: peer.clone(),
detail: format!("task response has unknown state '{other}'"),
}),
None => Err(PeerError::InvalidResponse {
peer: peer.clone(),
detail: "task response has no status.state".to_owned(),
}),
}
}
#[must_use]
pub const fn is_terminal(self) -> bool {
matches!(
self,
Self::Completed | Self::Failed | Self::Canceled | Self::Rejected
)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PeerTaskSnapshot {
pub task: PeerTask,
pub state: PeerTaskState,
pub value: Value,
}
#[derive(Debug)]
pub struct PeerTaskCall {
task: PeerTask,
grant: PeerGrant,
credential: Option<PeerCredential>,
client: Arc<dyn PeerClient>,
}
impl PeerTaskCall {
pub fn prepare(
registry: &PeerRegistry,
client: Arc<dyn PeerClient>,
task: PeerTask,
) -> Result<Self, PeerError> {
let Some(grant) = registry.grant(&task.peer).cloned() else {
return Err(PeerError::Unknown {
peer: task.peer.clone(),
});
};
let credential = registry.credential_for(&task.peer)?.cloned();
Ok(Self {
task,
grant,
credential,
client,
})
}
}
#[async_trait]
impl Effect for PeerTaskCall {
type Output = PeerTaskSnapshot;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"a2a.task/get",
serde_json::json!({
"peer": self.task.peer.0,
"task_id": self.task.id,
"context_id": self.task.context_id,
}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn retry(&self) -> RetryPolicy {
self.grant.retry
}
fn output_sensitivity(&self) -> Sensitivity {
self.grant.output_sensitivity
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
let value = self
.client
.get_task(&self.task.peer, &self.task.id, self.credential.as_ref())
.await
.map_err(|error| {
let detail = error.to_string();
match error.disposition() {
Disposition::DidNotHappen => EffectError::Rejected(detail),
Disposition::InDoubt => EffectError::Interrupted {
driver: self.task.peer.to_string(),
detail,
},
Disposition::Landed => EffectError::Performed(detail),
}
})?;
let state = PeerTaskState::parse(&self.task.peer, &value).map_err(|error| {
EffectError::Interrupted {
driver: self.task.peer.to_string(),
detail: error.to_string(),
}
})?;
Ok(PeerTaskSnapshot {
task: self.task.clone(),
state,
value,
})
}
}
#[derive(Debug, Default, Clone)]
pub struct PeerRegistry {
peers: BTreeMap<PeerId, PeerGrant>,
}
impl PeerRegistry {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn allow(mut self, peer: PeerId, grant: PeerGrant) -> Self {
self.peers.insert(peer, grant);
self
}
#[must_use]
pub fn grant(&self, peer: &PeerId) -> Option<&PeerGrant> {
self.peers.get(peer)
}
pub fn credential_for(&self, peer: &PeerId) -> Result<Option<&PeerCredential>, PeerError> {
let Some(grant) = self.peers.get(peer) else {
return Ok(None);
};
match grant.credential.as_ref() {
None => Ok(None),
Some(c) if c.audience() == peer => Ok(Some(c)),
Some(c) => Err(PeerError::WrongAudience {
peer: peer.clone(),
held_for: c.audience().clone(),
}),
}
}
}
#[derive(Debug)]
pub struct PeerCall {
peer: PeerId,
capability: String,
payload: Value,
grant: PeerGrant,
acting_as: Delegation,
credential: Option<PeerCredential>,
client: Arc<dyn PeerClient>,
provenance: Option<crate::core::Provenance>,
}
impl PeerCall {
pub fn prepare(
registry: &PeerRegistry,
client: Arc<dyn PeerClient>,
caller: &Delegation,
peer: PeerId,
capability: impl Into<String>,
payload: Value,
) -> Result<Self, PeerError> {
let held = registry.credential_for(&peer)?.cloned();
Self::prepare_with_credential(registry, client, caller, peer, capability, payload, held)
}
#[allow(clippy::too_many_arguments)]
pub fn prepare_with_credential(
registry: &PeerRegistry,
client: Arc<dyn PeerClient>,
caller: &Delegation,
peer: PeerId,
capability: impl Into<String>,
payload: Value,
credential: Option<PeerCredential>,
) -> Result<Self, PeerError> {
if let Some(c) = credential.as_ref()
&& c.audience() != &peer
{
return Err(PeerError::WrongAudience {
held_for: c.audience().clone(),
peer,
});
}
Self::build(
registry, client, caller, peer, capability, payload, credential,
)
}
#[allow(clippy::too_many_arguments)]
fn build(
registry: &PeerRegistry,
client: Arc<dyn PeerClient>,
caller: &Delegation,
peer: PeerId,
capability: impl Into<String>,
payload: Value,
credential: Option<PeerCredential>,
) -> Result<Self, PeerError> {
let Some(grant) = registry.grant(&peer) else {
return Err(PeerError::Unknown { peer });
};
let acting_as = caller
.delegate(Principal::new(peer.to_string(), grant.scope.clone()))
.map_err(|source| PeerError::Delegation {
peer: peer.clone(),
source,
})?;
Ok(Self {
provenance: None,
grant: grant.clone(),
capability: capability.into(),
payload,
acting_as,
credential,
peer,
client,
})
}
#[must_use]
pub const fn acting_as(&self) -> &Delegation {
&self.acting_as
}
}
#[async_trait]
impl Effect for PeerCall {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"a2a.peer/call",
serde_json::json!({
"peer": self.peer.0,
"capability": self.capability,
"payload": self.payload,
}),
)
}
fn mutates(&self) -> bool {
self.grant.mutates
}
fn recovery(&self) -> Recovery {
self.grant.recovery.clone()
}
fn retry(&self) -> RetryPolicy {
self.grant.retry
}
fn max_sensitivity(&self) -> Sensitivity {
self.grant.max_sensitivity
}
fn delegation_depth(&self) -> Option<usize> {
Some(self.acting_as.depth())
}
fn output_sensitivity(&self) -> Sensitivity {
self.grant.output_sensitivity
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
fn attach(&mut self, provenance: &crate::core::Provenance) {
self.provenance = Some(provenance.clone());
}
async fn perform(&self) -> Result<Value, EffectError> {
self.client
.send(
&self.peer,
&self.capability,
&self.payload,
&self.acting_as,
self.credential.as_ref(),
self.provenance.as_ref(),
)
.await
.map_err(|e| {
let detail = e.to_string();
match e.disposition() {
Disposition::DidNotHappen => EffectError::Rejected(detail),
Disposition::InDoubt => EffectError::Interrupted {
driver: self.peer.to_string(),
detail,
},
Disposition::Landed => EffectError::Performed(detail),
}
})
}
}