agentplane/peers/mod.rs
1//! Calling other agents.
2//!
3//! A peer hop is a tool call with two extra problems, and both are about
4//! identity rather than transport.
5//!
6//! # Token confusion
7//!
8//! To call a peer you hand it a credential. If that credential is not bound to
9//! *that peer*, the peer can replay it somewhere else — and it does not need to
10//! be malicious to do so, only compromised or confused. A bearer token sent to
11//! peer B and accepted by peer A is the whole vulnerability class, and it is why
12//! OAuth grew Resource Indicators (RFC 8707): the token says which audience it
13//! is for, and everyone else refuses it.
14//!
15//! This runtime cannot make peer A check that. What it *can* do — and what
16//! [`PeerRegistry`] enforces — is never send a credential to an audience it was
17//! not minted for. A credential for `settlement.example` is structurally
18//! unusable when calling `reviewer.example`; the call is refused before anything
19//! leaves.
20//!
21//! # Authority must narrow at the boundary
22//!
23//! A peer acts on our behalf, so the chain it receives is our chain with one
24//! more link. [`Delegation::delegate`] already refuses to widen and caps depth,
25//! so a hop cannot hand a peer more authority than the caller holds, and a
26//! request cannot wander arbitrarily far from the human who authorised it.
27//! "Our chain" is the **run's** — `StepCtx::acting_as`, which on a served
28//! plane is the caller's — never one a skill holds for itself, because a
29//! chain held by the skill is the same owner on every call whoever asked.
30//!
31//! # Who asked
32//!
33//! The chain narrows in-process; what reaches the peer is a credential. A
34//! credential held for the peer names nobody, so every run on the plane looks
35//! the same from the far side. A peer wired with a
36//! [`CredentialSource`] is called instead with a credential the source obtains
37//! for the run's owner — the person the run acts for — with the plane as actor,
38//! and a run that acts for nobody is refused rather than sent out under the
39//! plane's name. Which of the two a hop presented is on its announcement
40//! ([`CredentialBinding`]); the credential never is.
41//!
42//! The registry decides what each peer is *granted*, and that is an operator's
43//! declaration. It is not taken from the peer's agent card, for exactly the
44//! reason MCP annotations are not taken from a server: a party describing its own
45//! privileges is not a source of truth about them.
46
47use std::collections::BTreeMap;
48use std::fmt::Debug;
49use std::sync::Arc;
50
51use async_trait::async_trait;
52use serde::{Deserialize, Serialize};
53use serde_json::Value;
54
55/// The A2A protocol version this crate speaks, as client and as server.
56///
57/// One definition because it is one fact. The client sends it in `A2A-Version`,
58/// the published card names it on every interface, and the server refuses a
59/// request that asks for something else — three places that must agree, and a
60/// version that disagrees with the card is the kind of drift a caller finds
61/// before we do.
62pub const PROTOCOL_VERSION: &str = "1.0";
63
64/// The `google.rpc.ErrorInfo` domain under which this plane defines its own
65/// error reasons on the A2A surface.
66///
67/// One definition read by both halves, because the pair `(domain, reason)` is
68/// what makes a server-defined error *identifiable*: the numeric code alone
69/// sits in the range JSON-RPC gives implementations and A2A 1.0 reserves for
70/// its own table, so two parties can hold the same number for different
71/// facts. A domain this project controls cannot collide with either.
72pub const ERROR_DOMAIN: &str = "agentplane.hupe1980.github.io";
73
74/// The reason token a full quota answers with, inside [`ERROR_DOMAIN`].
75///
76/// The A2A client refuses-and-backs-off only on this exact pair, so the token
77/// is protocol surface, not a message string.
78pub const QUOTA_EXHAUSTED_REASON: &str = "QUOTA_EXHAUSTED";
79
80/// The reason token a halted agent answers with, inside [`ERROR_DOMAIN`].
81///
82/// Distinct from [`QUOTA_EXHAUSTED_REASON`] because the two ask opposite
83/// things of a caller: a ceiling says *come back*, and a halt says *somebody
84/// is dealing with an incident* — retrying is exactly what an operator pulling
85/// the switch is trying to stop. Answered under the quota token, a halt would
86/// teach every peer to hammer the one refusal that means stop.
87pub const HALTED_REASON: &str = "HALTED";
88
89/// The reason token a draining instance answers with, inside [`ERROR_DOMAIN`].
90///
91/// The third of the three admission refusals, and the only one a caller can act
92/// on *immediately*: a ceiling clears when a run finishes on that plane, a halt
93/// clears when a person lifts it, and this one clears the moment the caller
94/// reaches a different instance. Answered under the quota token it would teach a
95/// peer to wait out a back-off for a refusal that a retry now would pass.
96pub const DRAINING_REASON: &str = "DRAINING";
97
98/// Parse an A2A protocol version into the `Major.Minor` pair used for
99/// negotiation.
100///
101/// The specification requires decimal `Major.Minor`. A numeric patch is
102/// tolerated because patch releases MUST NOT affect compatibility, but an
103/// arbitrary suffix is not a patch version and must not turn `1.0.preview`
104/// into `1.0`. Keeping this in one place prevents card selection and server
105/// negotiation from accepting different version languages.
106#[cfg(any(feature = "a2a-server", all(feature = "a2a", feature = "manifest")))]
107pub(crate) fn protocol_major_minor(version: &str) -> Option<(u64, u64)> {
108 let mut parts = version.split('.');
109 let major = parts.next()?.parse().ok()?;
110 let minor = parts.next()?.parse().ok()?;
111 match parts.next() {
112 None => Some((major, minor)),
113 Some(patch) if !patch.is_empty() && patch.parse::<u64>().is_ok() => {
114 if parts.next().is_none() {
115 Some((major, minor))
116 } else {
117 None
118 }
119 }
120 Some(_) => None,
121 }
122}
123
124#[cfg(feature = "a2a")]
125pub mod a2a;
126#[cfg(feature = "manifest")]
127mod card;
128#[cfg(feature = "manifest")]
129mod card_sig;
130#[cfg(all(feature = "manifest", feature = "a2a"))]
131mod discovery;
132#[cfg(feature = "manifest")]
133pub use card::{
134 AgentCard, AgentExtension, CardCapabilities, CardInterface, CardSecurity,
135 CardSecurityRequirement, CardSecurityScheme, CardSkill, EXT_AGENT_DIRECTORY, EXT_GOVERNANCE,
136 EXT_MANIFEST_PROVENANCE, ExtendedAgentCard, ExtendedBudget, ExtendedTool,
137 HttpAuthSecurityScheme, SecurityScopeList, WELL_KNOWN_PATH, agent_card_path,
138};
139#[cfg(feature = "manifest")]
140pub use card_sig::{
141 ALG, CardSignature, CardSignatureError, CardSigner, CardVerifier, signing_input,
142};
143#[cfg(all(feature = "manifest", feature = "a2a"))]
144pub use discovery::{CardClient, DiscoveryError, JSONRPC};
145mod credentials;
146pub use credentials::{Cached, CredentialError, CredentialSource, TokenExchange};
147#[cfg(feature = "a2a")]
148mod exchange;
149#[cfg(feature = "a2a")]
150pub use exchange::{ACCESS_TOKEN_TYPE, TOKEN_EXCHANGE_GRANT, TokenEndpoint};
151
152use std::time::Duration;
153
154use crate::core::{
155 Capability, CredentialBinding, Delegation, DelegationError, Disposition, Effect,
156 EffectDescriptor, EffectError, Principal, ProtectedField, Recovery, RetryPolicy, Scope, Secret,
157 Sensitivity, SourceId, Timestamp, Trust,
158};
159
160/// Another agent, addressed by the name this plane knows it by.
161///
162/// Local, like a tool's server name. A peer that could choose its own identifier
163/// could step into another peer's grant.
164#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
165pub struct PeerId(pub String);
166
167impl PeerId {
168 pub fn new(s: impl Into<String>) -> Self {
169 Self(s.into())
170 }
171}
172
173impl std::fmt::Display for PeerId {
174 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
175 f.write_str(&self.0)
176 }
177}
178
179/// A credential, and the single audience it may be presented to.
180///
181/// The audience is not decoration. A bearer token with no stated audience is one
182/// that any recipient can replay at any other, and the whole point of RFC 8707 is
183/// that a token names where it is allowed to be spent.
184#[derive(Clone, PartialEq, Eq)]
185pub struct PeerCredential {
186 audience: PeerId,
187 /// The principal it was issued for, when the issuer named one.
188 subject: Option<String>,
189 /// Wiped when it drops, and compared in constant time.
190 ///
191 /// The redacting `Debug` and the absent `Serialize` stop it being *written*
192 /// somewhere. Neither stops it *staying* in freed heap after the credential
193 /// is gone, where a core dump or a swap file finds it — which is what
194 /// [`Secret`] is for.
195 secret: Secret,
196 /// When it stops being accepted, if the issuer said.
197 ///
198 /// `None` means the issuer gave no expiry — treated as usable, because
199 /// inventing one here would either reject working credentials or invent a
200 /// guarantee the issuer did not make.
201 expires_at: Option<Timestamp>,
202}
203
204impl PeerCredential {
205 /// Mint a credential for exactly one peer, naming nobody.
206 pub fn for_audience(audience: PeerId, secret: impl Into<String>) -> Self {
207 Self {
208 audience,
209 subject: None,
210 secret: Secret::new(secret),
211 expires_at: None,
212 }
213 }
214
215 /// Mint a credential for exactly one peer, naming the principal it was
216 /// issued for.
217 pub fn for_subject(
218 audience: PeerId,
219 subject: impl Into<String>,
220 secret: impl Into<String>,
221 ) -> Self {
222 Self {
223 subject: Some(subject.into()),
224 ..Self::for_audience(audience, secret)
225 }
226 }
227
228 /// Say when this credential stops being accepted.
229 #[must_use]
230 pub const fn expiring_at(mut self, at: Timestamp) -> Self {
231 self.expires_at = Some(at);
232 self
233 }
234
235 #[must_use]
236 pub const fn expires_at(&self) -> Option<Timestamp> {
237 self.expires_at
238 }
239
240 /// Whether this is still worth sending at `now`.
241 ///
242 /// `skew` is subtracted from the expiry, so a credential that expires in two
243 /// seconds is treated as already spent. Without that margin a token is sent
244 /// *just* before it lapses and is rejected in flight — which arrives as a
245 /// peer failure of unknown disposition, when it was really a refresh nobody
246 /// scheduled.
247 ///
248 /// `now` is a parameter rather than a clock read: expiry is transport
249 /// metadata and this keeps it testable at arbitrary instants.
250 #[must_use]
251 pub fn is_usable_at(&self, now: Timestamp, skew: Duration) -> bool {
252 let Some(expiry) = self.expires_at else {
253 return true;
254 };
255 let margin = i64::try_from(skew.as_secs()).unwrap_or(i64::MAX);
256 now.unix_timestamp().saturating_add(margin) < expiry.unix_timestamp()
257 }
258
259 #[must_use]
260 pub const fn audience(&self) -> &PeerId {
261 &self.audience
262 }
263
264 /// The principal this credential was issued for, if it names one.
265 #[must_use]
266 pub fn subject(&self) -> Option<&str> {
267 self.subject.as_deref()
268 }
269
270 /// The bearer value.
271 ///
272 /// Reachable only once a caller has already been past the audience check, so
273 /// there is no path that sends this to the wrong peer without going around
274 /// [`PeerRegistry::credential_for`] deliberately.
275 #[must_use]
276 pub fn expose(&self) -> &str {
277 self.secret.expose()
278 }
279}
280
281/// Never renders the secret.
282///
283/// A credential that prints itself ends up in a log, a span attribute, or an
284/// error message — and this crate writes all three. The audience is shown
285/// because that is the part worth debugging.
286impl Debug for PeerCredential {
287 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
288 f.debug_struct("PeerCredential")
289 .field("audience", &self.audience)
290 .field("subject", &self.subject)
291 .field("expires_at", &self.expires_at)
292 .field("secret", &"<redacted>")
293 .finish()
294 }
295}
296
297/// What the operator grants a peer.
298#[derive(Debug, Clone)]
299pub struct PeerGrant {
300 /// The authority this peer may hold. Narrowed further by the caller's own
301 /// scope at every hop — a grant is a ceiling, not an entitlement.
302 pub scope: Scope,
303 /// Whether a call to this peer changes the world.
304 pub mutates: bool,
305 pub recovery: Recovery,
306 pub max_sensitivity: Sensitivity,
307 pub output_sensitivity: Sensitivity,
308 pub retry: RetryPolicy,
309 /// The credential held for this peer, naming nobody. On a peer with a
310 /// [`source`](Self::with_source), the plane's own credential, presented
311 /// only for a run admitted as the plane.
312 credential: Option<PeerCredential>,
313 /// Where a credential naming the run's owner comes from, for a peer that
314 /// is told who asked.
315 source: Option<Arc<dyn CredentialSource>>,
316}
317
318impl PeerGrant {
319 /// A grant with the conservative posture: mutating, operator-resolved.
320 #[must_use]
321 pub fn new(scope: Scope) -> Self {
322 Self {
323 scope,
324 mutates: true,
325 recovery: Recovery::RequiresOperator,
326 max_sensitivity: Sensitivity::Public,
327 output_sensitivity: Sensitivity::Public,
328 retry: RetryPolicy::never(),
329 credential: None,
330 source: None,
331 }
332 }
333
334 /// Attach the credential this peer is called with.
335 ///
336 /// # Panics
337 ///
338 /// If the credential's audience is not this peer. That is a configuration
339 /// error the operator must see at startup rather than a refusal at 3am, and
340 /// there is no sensible way to continue: the alternative is holding a
341 /// credential that can only ever be sent to the wrong place.
342 #[must_use]
343 pub fn with_credential(mut self, peer: &PeerId, credential: PeerCredential) -> Self {
344 assert_eq!(
345 credential.audience(),
346 peer,
347 "a credential for '{}' was attached to peer '{peer}'. An audience-bound \
348 credential presented to the wrong peer is exactly what binding exists \
349 to prevent",
350 credential.audience()
351 );
352 self.credential = Some(credential);
353 self
354 }
355
356 /// Call this peer with a credential naming the person each run acts for,
357 /// obtained from `source` when the call is made.
358 ///
359 /// From then on a run that acts for nobody is refused this peer rather
360 /// than sent out under the plane's name, and the credential held through
361 /// [`with_credential`](Self::with_credential), if any, is presented only
362 /// for a run admitted as the plane.
363 #[must_use]
364 pub fn with_source(mut self, source: Arc<dyn CredentialSource>) -> Self {
365 self.source = Some(source);
366 self
367 }
368
369 /// Where this peer's subject-bound credentials come from, if it has one.
370 #[must_use]
371 pub fn source(&self) -> Option<&Arc<dyn CredentialSource>> {
372 self.source.as_ref()
373 }
374
375 #[must_use]
376 pub fn read_only(mut self) -> Self {
377 self.mutates = false;
378 self.recovery = Recovery::Retry;
379 self
380 }
381
382 #[must_use]
383 pub const fn output_sensitivity(mut self, s: Sensitivity) -> Self {
384 self.output_sensitivity = s;
385 self
386 }
387}
388
389/// Why a peer call could not be made, or did not work.
390#[derive(Debug, thiserror::Error)]
391pub enum PeerError {
392 /// This peer is not in the registry.
393 #[error(
394 "peer '{peer}' is not registered; a peer nobody declared is a peer nobody granted anything"
395 )]
396 Unknown { peer: PeerId },
397
398 /// The hop would hand the peer authority the caller does not hold, or would
399 /// take the request too far from the human who authorised it.
400 #[error("delegating to '{peer}' is refused: {source}")]
401 Delegation {
402 peer: PeerId,
403 /// Boxed: the refusal names both links' bounds, and every peer
404 /// call's `Result` would otherwise carry that width on its happy path.
405 #[source]
406 source: Box<DelegationError>,
407 },
408
409 /// The credential on hand is for someone else.
410 #[error(
411 "the credential held for this call is bound to '{held_for}', not '{peer}' — \
412 presenting it would let '{peer}' replay it at '{held_for}'"
413 )]
414 WrongAudience { peer: PeerId, held_for: PeerId },
415
416 /// The peer is told who each call is for, and this run acts for nobody.
417 ///
418 /// Refused rather than sent with the credential held for the peer: that
419 /// credential names nobody, and presenting it for a run with no chain is
420 /// the ambient authority a subject-bound peer was wired to rule out.
421 #[error(
422 "peer '{peer}' is called with a credential naming who asked, and this run acts \
423 for nobody — admit it under a chain"
424 )]
425 NoSubject { peer: PeerId },
426
427 /// A run admitted as the plane called a subject-bound peer, and the grant
428 /// holds no credential of the plane's own for it.
429 #[error(
430 "peer '{peer}' is called with a credential naming who asked; a run admitted as \
431 the plane presents the plane's own, and none is held for '{peer}'"
432 )]
433 NoPlaneCredential { peer: PeerId },
434
435 /// The capability asked for is outside what this peer is granted.
436 ///
437 /// The peer would refuse it at admission — the chain it receives permits
438 /// only the grant's scope — but a call that cannot succeed should not
439 /// leave: it costs a round trip, lands in the peer's journal as a refused
440 /// admission, and reads to the operator as the peer declining rather than
441 /// as this plane never having granted it.
442 #[error(
443 "peer '{peer}' is not granted '{capability}'; the registry's scope for it is the \
444 ceiling on what this plane may ask it to do"
445 )]
446 NotGranted { peer: PeerId, capability: String },
447
448 /// The call never left.
449 #[error("could not reach '{peer}': {detail}")]
450 Unreachable { peer: PeerId, detail: String },
451
452 /// The peer received it and declined without acting.
453 #[error("'{peer}' refused the request: {detail}")]
454 Refused { peer: PeerId, detail: String },
455
456 /// Sent, and no answer came within the deadline.
457 ///
458 /// Only for actual timeouts. Everything else that leaves the outcome
459 /// unknown — an answered fault, an unreadable response, a connection that
460 /// died mid-flight — is [`InDoubt`](Self::InDoubt), because "did not
461 /// answer in time" is a false diagnosis of a peer that answered HTTP 500
462 /// promptly, and a false diagnosis is what an operator debugs first.
463 #[error("'{peer}' did not answer in time: {detail}")]
464 TimedOut { peer: PeerId, detail: String },
465
466 /// The request may have reached the peer, and the outcome is unknown.
467 ///
468 /// The in-doubt bucket stated honestly: an HTTP 5xx, a JSON-RPC internal
469 /// error, a response that could not be read, a request that failed
470 /// mid-flight. Each says the peer may have acted; none says it was slow.
471 #[error("the outcome at '{peer}' is unknown: {detail}")]
472 InDoubt { peer: PeerId, detail: String },
473
474 /// Sent, but the peer's response did not conform to the negotiated protocol.
475 ///
476 /// The request may already have caused work, so this is in doubt rather
477 /// than a clean refusal.
478 #[error("'{peer}' returned an invalid response: {detail}")]
479 InvalidResponse { peer: PeerId, detail: String },
480
481 /// The peer acted and reported failure.
482 #[error("'{peer}' reported a failure: {detail}")]
483 Failed { peer: PeerId, detail: String },
484}
485
486impl PeerError {
487 /// What this failure says about whether the request reached the peer.
488 #[must_use]
489 pub const fn disposition(&self) -> Disposition {
490 match self {
491 // Nothing was sent: refused locally, or refused by the peer before
492 // it acted.
493 Self::Unknown { .. }
494 | Self::Delegation { .. }
495 | Self::WrongAudience { .. }
496 | Self::NoSubject { .. }
497 | Self::NoPlaneCredential { .. }
498 | Self::NotGranted { .. }
499 | Self::Unreachable { .. }
500 | Self::Refused { .. } => Disposition::DidNotHappen,
501 Self::TimedOut { .. } | Self::InDoubt { .. } | Self::InvalidResponse { .. } => {
502 Disposition::InDoubt
503 }
504 Self::Failed { .. } => Disposition::Landed,
505 }
506 }
507}
508
509/// Carries a request to a peer.
510#[async_trait]
511pub trait PeerClient: Send + Sync + Debug {
512 /// Send a request on behalf of a delegation chain.
513 ///
514 /// # Errors
515 ///
516 /// A [`PeerError`] whose variant states what is known about whether the
517 /// request reached the peer.
518 async fn send(
519 &self,
520 peer: &PeerId,
521 capability: &str,
522 payload: &Value,
523 acting_as: &Delegation,
524 credential: Option<&PeerCredential>,
525 provenance: Option<&crate::core::Provenance>,
526 ) -> Result<Value, PeerError>;
527
528 /// Read one previously accepted remote task.
529 ///
530 /// A default refusal keeps non-task peer transports honest. Implementors
531 /// must override this only when the wire has a stable task handle and an
532 /// idempotent read operation.
533 async fn get_task(
534 &self,
535 peer: &PeerId,
536 task_id: &str,
537 credential: Option<&PeerCredential>,
538 ) -> Result<Value, PeerError> {
539 let _ = (task_id, credential);
540 Err(PeerError::Refused {
541 peer: peer.clone(),
542 detail: "this peer transport does not support task lookup".to_owned(),
543 })
544 }
545
546 /// Ask the peer to stop one previously accepted remote task.
547 ///
548 /// Cooperative, exactly as this plane's own server implements it: the
549 /// answer means the request was durably recorded, and the task stops at
550 /// its next step boundary — polling is how the eventual state is
551 /// observed. The same default refusal as [`get_task`](Self::get_task),
552 /// for the same reason.
553 async fn cancel_task(
554 &self,
555 peer: &PeerId,
556 task_id: &str,
557 credential: Option<&PeerCredential>,
558 ) -> Result<Value, PeerError> {
559 let _ = (task_id, credential);
560 Err(PeerError::Refused {
561 peer: peer.clone(),
562 detail: "this peer transport does not support task cancellation".to_owned(),
563 })
564 }
565}
566
567/// One client per peer, resolved by the name a grant carries.
568///
569/// The peers' twin of `ToolRouter`, for the same reason: a single transport
570/// handed every peer id would have to hold one endpoint, and a plane that
571/// consults two peers has two. A peer nobody routed is unreachable — a
572/// refusal that never leaves — rather than a guess at the nearest endpoint.
573#[derive(Debug, Default)]
574pub struct PeerRouter {
575 routes: BTreeMap<PeerId, Arc<dyn PeerClient>>,
576}
577
578impl PeerRouter {
579 #[must_use]
580 pub fn new() -> Self {
581 Self::default()
582 }
583
584 /// Reach one peer through one client.
585 ///
586 /// # Panics
587 ///
588 /// If the peer is already routed: silently replacing would make
589 /// registration order decide which endpoint a call reaches.
590 #[must_use]
591 pub fn peer(mut self, peer: PeerId, client: Arc<dyn PeerClient>) -> Self {
592 assert!(
593 !self.routes.contains_key(&peer),
594 "peer '{peer}' is routed twice — one of the two endpoints would silently \
595 never be called"
596 );
597 self.routes.insert(peer, client);
598 self
599 }
600
601 /// The peers this router can reach.
602 pub fn peers(&self) -> impl Iterator<Item = &PeerId> {
603 self.routes.keys()
604 }
605
606 fn route(&self, peer: &PeerId) -> Result<&Arc<dyn PeerClient>, PeerError> {
607 self.routes.get(peer).ok_or_else(|| PeerError::Unreachable {
608 peer: peer.clone(),
609 detail: format!(
610 "no transport is wired for this peer; this plane routes {:?}",
611 self.routes
612 .keys()
613 .map(ToString::to_string)
614 .collect::<Vec<_>>()
615 ),
616 })
617 }
618}
619
620#[async_trait]
621impl PeerClient for PeerRouter {
622 async fn send(
623 &self,
624 peer: &PeerId,
625 capability: &str,
626 payload: &Value,
627 acting_as: &Delegation,
628 credential: Option<&PeerCredential>,
629 provenance: Option<&crate::core::Provenance>,
630 ) -> Result<Value, PeerError> {
631 self.route(peer)?
632 .send(peer, capability, payload, acting_as, credential, provenance)
633 .await
634 }
635
636 async fn get_task(
637 &self,
638 peer: &PeerId,
639 task_id: &str,
640 credential: Option<&PeerCredential>,
641 ) -> Result<Value, PeerError> {
642 self.route(peer)?.get_task(peer, task_id, credential).await
643 }
644
645 async fn cancel_task(
646 &self,
647 peer: &PeerId,
648 task_id: &str,
649 credential: Option<&PeerCredential>,
650 ) -> Result<Value, PeerError> {
651 self.route(peer)?
652 .cancel_task(peer, task_id, credential)
653 .await
654 }
655}
656
657/// A stable handle returned by a remote peer.
658#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
659pub struct PeerTask {
660 pub peer: PeerId,
661 pub id: String,
662 #[serde(default, skip_serializing_if = "Option::is_none")]
663 pub context_id: Option<String>,
664}
665
666impl PeerTask {
667 /// Extract a task handle from a peer response, or `None` for a direct
668 /// message response.
669 pub fn from_response(peer: PeerId, response: &Value) -> Result<Option<Self>, PeerError> {
670 if response.get("role").is_some() {
671 return Ok(None);
672 }
673 let id = response.get("id").and_then(Value::as_str).ok_or_else(|| {
674 PeerError::InvalidResponse {
675 peer: peer.clone(),
676 detail: "task response has no string id".to_owned(),
677 }
678 })?;
679 if response
680 .get("status")
681 .and_then(|status| status.get("state"))
682 .and_then(Value::as_str)
683 .is_none()
684 {
685 return Err(PeerError::InvalidResponse {
686 peer,
687 detail: "task response has no status.state".to_owned(),
688 });
689 }
690 Ok(Some(Self {
691 peer,
692 id: id.to_owned(),
693 context_id: response
694 .get("contextId")
695 .and_then(Value::as_str)
696 .map(ToOwned::to_owned),
697 }))
698 }
699}
700
701/// A2A's task lifecycle, normalized for callers that need to decide whether to
702/// poll, provide input, or consume the result.
703#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
704#[serde(rename_all = "snake_case")]
705#[non_exhaustive]
706pub enum PeerTaskState {
707 Submitted,
708 Working,
709 Completed,
710 Failed,
711 Canceled,
712 Rejected,
713 InputRequired,
714 AuthRequired,
715}
716
717impl PeerTaskState {
718 fn parse(peer: &PeerId, value: &Value) -> Result<Self, PeerError> {
719 match value
720 .get("status")
721 .and_then(|status| status.get("state"))
722 .and_then(Value::as_str)
723 {
724 Some("TASK_STATE_SUBMITTED") => Ok(Self::Submitted),
725 Some("TASK_STATE_WORKING") => Ok(Self::Working),
726 Some("TASK_STATE_COMPLETED") => Ok(Self::Completed),
727 Some("TASK_STATE_FAILED") => Ok(Self::Failed),
728 Some("TASK_STATE_CANCELED") => Ok(Self::Canceled),
729 Some("TASK_STATE_REJECTED") => Ok(Self::Rejected),
730 Some("TASK_STATE_INPUT_REQUIRED") => Ok(Self::InputRequired),
731 Some("TASK_STATE_AUTH_REQUIRED") => Ok(Self::AuthRequired),
732 Some(other) => Err(PeerError::InvalidResponse {
733 peer: peer.clone(),
734 detail: format!("task response has unknown state '{other}'"),
735 }),
736 None => Err(PeerError::InvalidResponse {
737 peer: peer.clone(),
738 detail: "task response has no status.state".to_owned(),
739 }),
740 }
741 }
742
743 #[must_use]
744 pub const fn is_terminal(self) -> bool {
745 matches!(
746 self,
747 Self::Completed | Self::Failed | Self::Canceled | Self::Rejected
748 )
749 }
750}
751
752/// One journaled observation of a remote task.
753#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
754pub struct PeerTaskSnapshot {
755 pub task: PeerTask,
756 pub state: PeerTaskState,
757 pub value: Value,
758}
759
760/// On whose behalf a hop presents its credential.
761#[derive(Debug, Clone, Copy)]
762pub enum Asker<'a> {
763 /// The person this chain acts for — its owner.
764 Owner(&'a Delegation),
765 /// The plane itself: a run admitted under the plane's own chain.
766 Plane,
767 /// Nobody: a run that acts under no chain.
768 Nobody,
769}
770
771/// Which credential a hop presents: decided when the call is prepared,
772/// obtained when it is performed.
773///
774/// Obtained at `perform` because that is the one place that runs live and
775/// only live: a call is constructed on every pass, a replay included, and a
776/// replay must reach no token endpoint.
777#[derive(Debug, Clone)]
778enum Presentation {
779 /// The credential held for the peer, or none; it names nobody.
780 Held(Option<PeerCredential>),
781 /// The plane's own held credential, for a run admitted as the plane.
782 Plane(PeerCredential),
783 /// Obtained from the peer's source for the run's owner.
784 Subject {
785 source: Arc<dyn CredentialSource>,
786 subject: String,
787 },
788}
789
790impl Presentation {
791 /// What the announcement records about this hop's credential.
792 fn binding(&self, audience: &PeerId) -> CredentialBinding {
793 let audience = audience.to_string();
794 match self {
795 Self::Held(_) => CredentialBinding::Unbound { audience },
796 Self::Plane(_) => CredentialBinding::Plane { audience },
797 Self::Subject { subject, .. } => CredentialBinding::Subject {
798 audience,
799 subject: subject.clone(),
800 },
801 }
802 }
803
804 /// The credential to put on the wire, checked for its audience and
805 /// subject whichever source produced it.
806 ///
807 /// A failure here is a refusal before anything was sent: an issuer that
808 /// could not be reached is [`EffectError::Unavailable`], a credential
809 /// for somebody else is [`EffectError::Refused`].
810 async fn obtain(&self, audience: &PeerId) -> Result<Option<PeerCredential>, EffectError> {
811 match self {
812 Self::Held(held) => Ok(held.clone()),
813 Self::Plane(own) => Ok(Some(own.clone())),
814 Self::Subject { source, subject } => {
815 let obtained = source
816 .credential(audience, subject, credentials::now())
817 .await
818 .and_then(|c| credentials::bound_to(c, audience, subject));
819 match obtained {
820 Ok(c) => Ok(Some(c)),
821 Err(CredentialError::Unavailable { detail, .. }) => {
822 Err(EffectError::Unavailable {
823 driver: audience.to_string(),
824 detail,
825 })
826 }
827 Err(refused) => Err(EffectError::Refused(refused.to_string())),
828 }
829 }
830 }
831 }
832}
833
834/// Map a peer's failure onto what it says about the world.
835fn effect_error(peer: &PeerId, error: &PeerError) -> EffectError {
836 let detail = error.to_string();
837 match error.disposition() {
838 Disposition::DidNotHappen => EffectError::Rejected(detail),
839 Disposition::InDoubt => EffectError::Interrupted {
840 driver: peer.to_string(),
841 detail,
842 },
843 Disposition::Landed => EffectError::Performed(detail),
844 }
845}
846
847/// A journaled, idempotent read of a remote task.
848#[derive(Debug)]
849pub struct PeerTaskCall {
850 task: PeerTask,
851 grant: PeerGrant,
852 presentation: Presentation,
853 client: Arc<dyn PeerClient>,
854}
855
856impl PeerTaskCall {
857 /// Prepare a task read under the same peer grant as the call that
858 /// created it, presenting the credential `asker` is owed.
859 ///
860 /// # Errors
861 ///
862 /// [`PeerError::Unknown`] for an unregistered peer,
863 /// [`PeerError::WrongAudience`] for a held credential bound elsewhere,
864 /// and, at a peer with a credential source, [`PeerError::NoSubject`] for
865 /// a run acting for nobody and [`PeerError::NoPlaneCredential`] for a run
866 /// admitted as the plane with no held credential.
867 pub fn prepare(
868 registry: &PeerRegistry,
869 client: Arc<dyn PeerClient>,
870 task: PeerTask,
871 asker: Asker<'_>,
872 ) -> Result<Self, PeerError> {
873 let Some(grant) = registry.grant(&task.peer).cloned() else {
874 return Err(PeerError::Unknown { peer: task.peer });
875 };
876 let presentation = registry.presentation(&task.peer, asker)?;
877 Ok(Self {
878 task,
879 grant,
880 presentation,
881 client,
882 })
883 }
884}
885
886#[async_trait]
887impl Effect for PeerTaskCall {
888 type Output = PeerTaskSnapshot;
889
890 fn descriptor(&self) -> EffectDescriptor {
891 EffectDescriptor::new(
892 "a2a.task/get",
893 serde_json::json!({
894 "peer": self.task.peer.0,
895 "task_id": self.task.id,
896 "context_id": self.task.context_id,
897 }),
898 )
899 }
900
901 fn mutates(&self) -> bool {
902 false
903 }
904
905 fn recovery(&self) -> Recovery {
906 Recovery::Retry
907 }
908
909 fn retry(&self) -> RetryPolicy {
910 self.grant.retry
911 }
912
913 fn output_sensitivity(&self) -> Sensitivity {
914 self.grant.output_sensitivity
915 }
916
917 fn trust(&self) -> Trust {
918 Trust::Untrusted
919 }
920
921 fn credential_binding(&self) -> Option<CredentialBinding> {
922 Some(self.presentation.binding(&self.task.peer))
923 }
924
925 async fn perform(&self) -> Result<Self::Output, EffectError> {
926 let credential = self.presentation.obtain(&self.task.peer).await?;
927 let value = self
928 .client
929 .get_task(&self.task.peer, &self.task.id, credential.as_ref())
930 .await
931 .map_err(|error| effect_error(&self.task.peer, &error))?;
932 let state = PeerTaskState::parse(&self.task.peer, &value).map_err(|error| {
933 EffectError::Interrupted {
934 driver: self.task.peer.to_string(),
935 detail: error.to_string(),
936 }
937 })?;
938 Ok(PeerTaskSnapshot {
939 task: self.task.clone(),
940 state,
941 value,
942 })
943 }
944}
945
946/// A journaled, cooperative cancellation of a remote task.
947///
948/// The A2A twin of [`McpTaskCancel`](crate::tools::McpTaskCancel), and it
949/// exists for the same lifecycle reason: a run that commissioned work at a
950/// peer and is itself cancelled or unwound must be able to tell the peer to
951/// stop, or the cancellation ends at this plane's edge while the peer keeps
952/// spending on an answer nobody will read.
953///
954/// # What the answer means
955///
956/// Acceptance of intent, never completion: this plane's own server records
957/// the request durably and the run stops at its next step boundary, so the
958/// snapshot that comes back is typically still `Working`. Polling through
959/// [`PeerTaskCall`] is how the eventual `Canceled` is observed. A task that
960/// already finished is *refused* by the far side (A2A's `TaskNotCancelable`,
961/// `-32002`) — a clean, pre-action decline, which is why retrying this
962/// effect is safe: a repeat of a cancel that landed meets that refusal
963/// rather than a second effect.
964#[derive(Debug)]
965pub struct PeerTaskCancel {
966 task: PeerTask,
967 grant: PeerGrant,
968 presentation: Presentation,
969 client: Arc<dyn PeerClient>,
970}
971
972impl PeerTaskCancel {
973 /// Prepare a cancellation under the same peer grant as the call that
974 /// created the task, presenting the credential `asker` is owed.
975 ///
976 /// # Errors
977 ///
978 /// As [`PeerTaskCall::prepare`].
979 pub fn prepare(
980 registry: &PeerRegistry,
981 client: Arc<dyn PeerClient>,
982 task: PeerTask,
983 asker: Asker<'_>,
984 ) -> Result<Self, PeerError> {
985 let Some(grant) = registry.grant(&task.peer).cloned() else {
986 return Err(PeerError::Unknown { peer: task.peer });
987 };
988 let presentation = registry.presentation(&task.peer, asker)?;
989 Ok(Self {
990 task,
991 grant,
992 presentation,
993 client,
994 })
995 }
996}
997
998#[async_trait]
999impl Effect for PeerTaskCancel {
1000 type Output = PeerTaskSnapshot;
1001
1002 fn descriptor(&self) -> EffectDescriptor {
1003 EffectDescriptor::new(
1004 "a2a.task/cancel",
1005 serde_json::json!({
1006 "peer": self.task.peer.0,
1007 "task_id": self.task.id,
1008 "context_id": self.task.context_id,
1009 }),
1010 )
1011 }
1012
1013 /// Asking a peer to stop is asking it to change what happens.
1014 fn mutates(&self) -> bool {
1015 true
1016 }
1017
1018 /// Retry is safe despite mutating: the far side answers a repeat of a
1019 /// cancel that landed with a clean refusal (`TaskNotCancelable`), never
1020 /// with a second act — the type-level statement of the module docs.
1021 fn recovery(&self) -> Recovery {
1022 Recovery::Retry
1023 }
1024
1025 fn retry(&self) -> RetryPolicy {
1026 self.grant.retry
1027 }
1028
1029 fn output_sensitivity(&self) -> Sensitivity {
1030 self.grant.output_sensitivity
1031 }
1032
1033 fn trust(&self) -> Trust {
1034 Trust::Untrusted
1035 }
1036
1037 fn credential_binding(&self) -> Option<CredentialBinding> {
1038 Some(self.presentation.binding(&self.task.peer))
1039 }
1040
1041 async fn perform(&self) -> Result<Self::Output, EffectError> {
1042 let credential = self.presentation.obtain(&self.task.peer).await?;
1043 let value = self
1044 .client
1045 .cancel_task(&self.task.peer, &self.task.id, credential.as_ref())
1046 .await
1047 .map_err(|error| effect_error(&self.task.peer, &error))?;
1048 let state = PeerTaskState::parse(&self.task.peer, &value).map_err(|error| {
1049 EffectError::Interrupted {
1050 driver: self.task.peer.to_string(),
1051 detail: error.to_string(),
1052 }
1053 })?;
1054 Ok(PeerTaskSnapshot {
1055 task: self.task.clone(),
1056 state,
1057 value,
1058 })
1059 }
1060}
1061
1062/// The peers this plane may call, and what each is granted.
1063#[derive(Debug, Default, Clone)]
1064pub struct PeerRegistry {
1065 peers: BTreeMap<PeerId, PeerGrant>,
1066}
1067
1068impl PeerRegistry {
1069 #[must_use]
1070 pub fn new() -> Self {
1071 Self::default()
1072 }
1073
1074 #[must_use]
1075 pub fn allow(mut self, peer: PeerId, grant: PeerGrant) -> Self {
1076 self.peers.insert(peer, grant);
1077 self
1078 }
1079
1080 #[must_use]
1081 pub fn grant(&self, peer: &PeerId) -> Option<&PeerGrant> {
1082 self.peers.get(peer)
1083 }
1084
1085 /// Every peer this registry names.
1086 pub fn peers(&self) -> impl Iterator<Item = &PeerId> {
1087 self.peers.keys()
1088 }
1089
1090 /// The credential for this peer, if one is held *and* bound to it.
1091 ///
1092 /// The audience check is here rather than at the call site so that no code
1093 /// path can reach a credential without passing it.
1094 ///
1095 /// # Errors
1096 ///
1097 /// [`PeerError::WrongAudience`] if a credential is held whose audience is a
1098 /// different peer.
1099 pub fn credential_for(&self, peer: &PeerId) -> Result<Option<&PeerCredential>, PeerError> {
1100 let Some(grant) = self.peers.get(peer) else {
1101 return Ok(None);
1102 };
1103 match grant.credential.as_ref() {
1104 None => Ok(None),
1105 Some(c) if c.audience() == peer => Ok(Some(c)),
1106 Some(c) => Err(PeerError::WrongAudience {
1107 peer: peer.clone(),
1108 held_for: c.audience().clone(),
1109 }),
1110 }
1111 }
1112
1113 /// Which credential a hop to `peer` presents for `asker`.
1114 ///
1115 /// A peer without a source presents the credential held for it, whoever
1116 /// asked. A peer with one presents a credential naming the chain's owner,
1117 /// the plane's own held credential for a run admitted as the plane, and
1118 /// nothing at all for a run that acts for nobody.
1119 ///
1120 /// # Errors
1121 ///
1122 /// * [`PeerError::Unknown`] if the peer is not registered.
1123 /// * [`PeerError::WrongAudience`] if the credential held is for someone else.
1124 /// * [`PeerError::NoSubject`] for a run acting for nobody at a peer with a
1125 /// source.
1126 /// * [`PeerError::NoPlaneCredential`] for a run admitted as the plane at a
1127 /// peer with a source and no held credential.
1128 fn presentation(&self, peer: &PeerId, asker: Asker<'_>) -> Result<Presentation, PeerError> {
1129 let Some(grant) = self.peers.get(peer) else {
1130 return Err(PeerError::Unknown { peer: peer.clone() });
1131 };
1132 let held = self.credential_for(peer)?.cloned();
1133 let Some(source) = grant.source.as_ref() else {
1134 return Ok(Presentation::Held(held));
1135 };
1136 match asker {
1137 Asker::Owner(chain) => Ok(Presentation::Subject {
1138 source: Arc::clone(source),
1139 subject: chain.owner().id.clone(),
1140 }),
1141 Asker::Plane => held
1142 .map(Presentation::Plane)
1143 .ok_or_else(|| PeerError::NoPlaneCredential { peer: peer.clone() }),
1144 Asker::Nobody => Err(PeerError::NoSubject { peer: peer.clone() }),
1145 }
1146 }
1147
1148 /// Drop every credential held for `subject`, at every peer's source.
1149 pub fn forget(&self, subject: &str) {
1150 for grant in self.peers.values() {
1151 if let Some(source) = grant.source.as_ref() {
1152 source.forget(subject);
1153 }
1154 }
1155 }
1156}
1157
1158/// One request to one peer.
1159#[derive(Debug)]
1160pub struct PeerCall {
1161 peer: PeerId,
1162 capability: String,
1163 payload: Value,
1164 grant: PeerGrant,
1165 /// The chain the *peer* acts under: ours, narrowed, with the peer appended.
1166 acting_as: Delegation,
1167 presentation: Presentation,
1168 client: Arc<dyn PeerClient>,
1169 /// Who is calling, sealed for this hop. Set by the runtime via
1170 /// [`Effect::attach`](crate::core::Effect::attach).
1171 provenance: Option<crate::core::Provenance>,
1172 /// Authority-bearing payload fields and their source rules, from the
1173 /// manifest grant that governs this call — see
1174 /// [`governed_by`](Self::governed_by). Empty for a call nothing governs.
1175 protected: Vec<ProtectedField>,
1176}
1177
1178impl PeerCall {
1179 /// Prepare a hop, attenuating the caller's authority onto the peer and
1180 /// presenting the credential the chain's owner is owed.
1181 ///
1182 /// # Errors
1183 ///
1184 /// * [`PeerError::Unknown`] if the peer is not registered — fail closed.
1185 /// * [`PeerError::NotGranted`] if the capability is outside the grant.
1186 /// * [`PeerError::Delegation`] if the grant would widen the caller's own
1187 /// authority, or if the chain is already at its depth limit.
1188 /// * [`PeerError::WrongAudience`] if the credential held is for someone else.
1189 pub fn prepare(
1190 registry: &PeerRegistry,
1191 client: Arc<dyn PeerClient>,
1192 caller: &Delegation,
1193 peer: PeerId,
1194 capability: impl Into<String>,
1195 payload: Value,
1196 ) -> Result<Self, PeerError> {
1197 Self::build(
1198 registry,
1199 client,
1200 caller,
1201 Asker::Owner(caller),
1202 peer,
1203 capability,
1204 payload,
1205 )
1206 }
1207
1208 /// Prepare a hop for a run admitted as the plane, under the plane's own
1209 /// chain.
1210 ///
1211 /// The plane is the party that asked, so a peer told who asked is shown
1212 /// the plane's own credential rather than one naming the chain's owner as
1213 /// if a person had asked.
1214 ///
1215 /// # Errors
1216 ///
1217 /// As [`prepare`](Self::prepare), plus [`PeerError::NoPlaneCredential`].
1218 pub fn prepare_as_plane(
1219 registry: &PeerRegistry,
1220 client: Arc<dyn PeerClient>,
1221 chain: &Delegation,
1222 peer: PeerId,
1223 capability: impl Into<String>,
1224 payload: Value,
1225 ) -> Result<Self, PeerError> {
1226 Self::build(
1227 registry,
1228 client,
1229 chain,
1230 Asker::Plane,
1231 peer,
1232 capability,
1233 payload,
1234 )
1235 }
1236
1237 fn build(
1238 registry: &PeerRegistry,
1239 client: Arc<dyn PeerClient>,
1240 caller: &Delegation,
1241 asker: Asker<'_>,
1242 peer: PeerId,
1243 capability: impl Into<String>,
1244 payload: Value,
1245 ) -> Result<Self, PeerError> {
1246 let Some(grant) = registry.grant(&peer) else {
1247 return Err(PeerError::Unknown { peer });
1248 };
1249 let capability = capability.into();
1250 // Refused here rather than by the peer: the chain it would receive
1251 // permits exactly `grant.scope`, so a capability outside it is a call
1252 // the far side's admission gate refuses after a round trip, journaled
1253 // there as *their* decline.
1254 if !grant.scope.permits(&Capability::new(capability.as_str())) {
1255 return Err(PeerError::NotGranted { peer, capability });
1256 }
1257
1258 // The grant is a ceiling. What the peer actually receives is bounded by
1259 // what *we* hold, which is what `delegate` enforces — a grant wider than
1260 // the caller's own authority is refused rather than silently clipped,
1261 // because silently clipping hides a misconfiguration that matters.
1262 let acting_as = caller
1263 .delegate(Principal::new(peer.to_string(), grant.scope.clone()))
1264 .map_err(|source| PeerError::Delegation {
1265 peer: peer.clone(),
1266 source: Box::new(source),
1267 })?;
1268
1269 let presentation = registry.presentation(&peer, asker)?;
1270
1271 Ok(Self {
1272 provenance: None,
1273 grant: grant.clone(),
1274 capability,
1275 payload,
1276 acting_as,
1277 presentation,
1278 peer,
1279 client,
1280 protected: Vec::new(),
1281 })
1282 }
1283
1284 /// Hold this call to a manifest grant's declaration.
1285 ///
1286 /// A registry entry is operator *wiring* — where the peer is, what it is
1287 /// granted, which credential reaches it. The reviewed `tool://<peer>/<capability>`
1288 /// grant in the calling agent's manifest is what says how the call may be
1289 /// *used*, and it governs the same way it governs a tool: its protected
1290 /// fields are checked at the sink, its ceiling bounds what may be sent, and
1291 /// its `mutates` can only make the call more cautious than the wiring did.
1292 #[must_use]
1293 pub fn governed_by(mut self, safety: &crate::tools::ToolSafety) -> Self {
1294 self.grant.mutates |= safety.mutates;
1295 self.grant.max_sensitivity = safety.max_sensitivity;
1296 self.grant.output_sensitivity =
1297 self.grant.output_sensitivity.max(safety.output_sensitivity);
1298 self.protected.clone_from(&safety.protected_fields);
1299 self
1300 }
1301
1302 /// The chain the peer will act under.
1303 #[must_use]
1304 pub const fn acting_as(&self) -> &Delegation {
1305 &self.acting_as
1306 }
1307}
1308
1309#[async_trait]
1310impl Effect for PeerCall {
1311 type Output = Value;
1312
1313 fn descriptor(&self) -> EffectDescriptor {
1314 EffectDescriptor::new(
1315 "a2a.peer/call",
1316 serde_json::json!({
1317 "peer": self.peer.0,
1318 "capability": self.capability,
1319 "payload": self.payload,
1320 }),
1321 )
1322 }
1323
1324 fn mutates(&self) -> bool {
1325 self.grant.mutates
1326 }
1327
1328 fn recovery(&self) -> Recovery {
1329 self.grant.recovery.clone()
1330 }
1331
1332 fn retry(&self) -> RetryPolicy {
1333 self.grant.retry
1334 }
1335
1336 fn max_sensitivity(&self) -> Sensitivity {
1337 self.grant.max_sensitivity
1338 }
1339
1340 fn delegation_depth(&self) -> Option<usize> {
1341 Some(self.acting_as.depth())
1342 }
1343
1344 /// The payload is what reaches the peer, so it is what the sink gate
1345 /// judges — the whole-value taint rule for a mutating grant, and the
1346 /// per-field rules a manifest declares.
1347 fn sink_arguments(&self) -> Option<&Value> {
1348 Some(&self.payload)
1349 }
1350
1351 fn protected_fields(&self) -> &[ProtectedField] {
1352 &self.protected
1353 }
1354
1355 /// The reference a manifest grants this call under, so a source rule can
1356 /// name *this* peer's answer — `tool://reviewer/audit.check` — rather than
1357 /// whichever peer an injected prompt reached first.
1358 fn source(&self) -> SourceId {
1359 SourceId::new(format!(
1360 "{}{}/{}",
1361 crate::tools::TOOL_SCHEME,
1362 self.peer,
1363 self.capability
1364 ))
1365 }
1366
1367 fn output_sensitivity(&self) -> Sensitivity {
1368 self.grant.output_sensitivity
1369 }
1370
1371 /// A peer's answer is another party's data.
1372 ///
1373 /// Stated rather than inherited because a peer feels more trusted than a
1374 /// tool — it is *our* agent, on our side. It is not: it runs somewhere else,
1375 /// under someone else's control, and it may itself have read the internet.
1376 fn trust(&self) -> Trust {
1377 Trust::Untrusted
1378 }
1379
1380 fn attach(&mut self, provenance: &crate::core::Provenance) {
1381 self.provenance = Some(provenance.clone());
1382 }
1383
1384 fn credential_binding(&self) -> Option<CredentialBinding> {
1385 Some(self.presentation.binding(&self.peer))
1386 }
1387
1388 async fn perform(&self) -> Result<Value, EffectError> {
1389 let credential = self.presentation.obtain(&self.peer).await?;
1390 self.client
1391 .send(
1392 &self.peer,
1393 &self.capability,
1394 &self.payload,
1395 &self.acting_as,
1396 credential.as_ref(),
1397 self.provenance.as_ref(),
1398 )
1399 .await
1400 .map_err(|e| effect_error(&self.peer, &e))
1401 }
1402}