Skip to main content

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}