Skip to main content

agentplane/api/
a2a.rs

1//! Serving A2A: this plane, as an agent other agents can call.
2//!
3//! The client side ([`crate::peers::a2a`]) lets this plane *call* peers. This is
4//! the other half — a peer calling us — and it is a different problem, because
5//! everything arriving here came from somebody else.
6//!
7//! # Why this is not a route on the operator API
8//!
9//! Every route on [`crate::api::Api`] authenticates and then authorizes. An
10//! Agent Card is public by design: it is what a caller reads *before* it has
11//! credentials, and a card behind authentication cannot be discovered. Bolting
12//! an unauthenticated path onto a surface whose invariant is "every route
13//! authenticates" would delete that invariant for the one route nobody would
14//! think to check.
15//!
16//! So this is its own router with its own rule: **the card is public, every
17//! method call is authenticated and authorized**.
18//!
19//! # What arrives here is untrusted
20//!
21//! A message from a peer is data written by a party this plane does not control,
22//! so it is admitted as `Tainted` with the sending
23//! peer's identity as its provenance source — never as trusted input. A skill
24//! that wants to act on it has to say so at a gate, and a protected sink field
25//! can name the one counterparty it will accept an amount from.
26//!
27//! This is the same reason the operator API takes an event's `source` from the
28//! authenticated caller rather than the body: a party describing itself is not
29//! evidence about itself.
30//!
31//! # The capability is named, never inferred
32//!
33//! A2A messages do not carry a "call this skill" field — the protocol assumes an
34//! agent works out what is being asked. This plane will not: choosing which
35//! capability to run on the strength of an untrusted message is a dispatch
36//! decision made by inference, and the thing doing the inferring would be a
37//! model reading attacker-controlled text.
38//!
39//! The skill is taken from `message.metadata.skill`, matched against the card's
40//! advertised skill ids. When the agent advertises exactly one there is nothing
41//! to infer and it is used; when it advertises several and none was named, the
42//! call is refused rather than guessed.
43//!
44//! # Optional capabilities say what is wired
45//!
46//! Streaming and non-terminal task subscription are durable journal views.
47//! Push is advertised only after `with_push` supplies durable registration
48//! storage and a retrying transport worker; without that wiring every push
49//! method uses `PushNotificationNotSupportedError` and the card advertises
50//! false.
51//!
52//! # A task belongs to the peer that admitted it
53//!
54//! Every read, write, stream, push configuration and `contextId` join is scoped
55//! to the admitting peer; another peer's task answers `TASK_NOT_FOUND`, never a
56//! refusal that confirms it exists.
57
58use std::sync::Arc;
59
60use axum::extract::State;
61use axum::http::{HeaderMap, StatusCode};
62use axum::response::{IntoResponse, Response};
63use axum::routing::{get, post};
64use axum::{Json, Router};
65use serde::{Deserialize, Serialize};
66use serde_json::{Value, json};
67
68use crate::core::{PolicyDecision, PolicyRequest, RunId, Seq, SourceId, Tainted};
69use crate::journal::RecordKind;
70use crate::manifest::Manifest;
71use crate::peers::{AgentCard, ExtendedAgentCard, WELL_KNOWN_PATH};
72use crate::runtime::Runtime;
73
74use super::{Authenticator, Caller};
75
76/// JSON-RPC method names, exactly as A2A 1.0 spells them.
77///
78/// 1.0 renamed these: `message/send` was the 0.3 spelling, and a server
79/// answering the old names would silently accept clients that have lost half the
80/// protocol. Constants rather than inline literals so the dispatch table and the
81/// refusal table cannot disagree.
82pub mod method {
83    pub const SEND_MESSAGE: &str = "SendMessage";
84    pub const GET_TASK: &str = "GetTask";
85    pub const CANCEL_TASK: &str = "CancelTask";
86    pub const GET_EXTENDED_CARD: &str = "GetExtendedAgentCard";
87
88    /// Optional operations implemented by this server.
89    pub const SEND_STREAMING: &str = "SendStreamingMessage";
90    pub const SUBSCRIBE: &str = "SubscribeToTask";
91    pub const LIST_TASKS: &str = "ListTasks";
92    pub const CREATE_PUSH: &str = "CreateTaskPushNotificationConfig";
93    pub const GET_PUSH: &str = "GetTaskPushNotificationConfig";
94    pub const LIST_PUSH: &str = "ListTaskPushNotificationConfigs";
95    pub const DELETE_PUSH: &str = "DeleteTaskPushNotificationConfig";
96}
97
98/// A2A-specific JSON-RPC error codes, from the spec's mapping table.
99pub mod code {
100    pub const PARSE_ERROR: i32 = -32700;
101    pub const INVALID_REQUEST: i32 = -32600;
102    pub const METHOD_NOT_FOUND: i32 = -32601;
103    pub const INVALID_PARAMS: i32 = -32602;
104    pub const INTERNAL_ERROR: i32 = -32603;
105
106    pub const TASK_NOT_FOUND: i32 = -32001;
107    pub const TASK_NOT_CANCELABLE: i32 = -32002;
108    pub const PUSH_NOT_SUPPORTED: i32 = -32003;
109    pub const UNSUPPORTED_OPERATION: i32 = -32004;
110    pub const CONTENT_TYPE_NOT_SUPPORTED: i32 = -32005;
111    pub const EXTENDED_CARD_NOT_CONFIGURED: i32 = -32007;
112    pub const VERSION_NOT_SUPPORTED: i32 = -32009;
113
114    /// Admission back-pressure: a quota ceiling refused the request before any
115    /// work was admitted.
116    ///
117    /// Not in A2A 1.0's error table — the spec defines only permanent
118    /// missing-capability codes and assigns back-pressure no code at all; its
119    /// own guidance stops at "return appropriate error responses when rate
120    /// limits are exceeded". Reaching into the table anyway would be worse
121    /// than inventing a code: `-32004 UnsupportedOperationError` for a full
122    /// quota teaches a compliant caller that the *operation does not exist
123    /// here* — the correct response to which is to abandon, never to retry —
124    /// when the truthful answer is "not right now". So the code is drawn from
125    /// JSON-RPC's implementation-defined server-error range.
126    ///
127    /// The **number is not the identity**. A2A 1.0 reserves `-32001..-32099`
128    /// for its own errors, the same band JSON-RPC gives implementations, so a
129    /// future spec revision could assign this numeral a meaning of its own.
130    /// What makes the refusal identifiable is the `google.rpc.ErrorInfo`
131    /// beside it — [`ERROR_DOMAIN`](crate::peers::ERROR_DOMAIN) +
132    /// [`QUOTA_EXHAUSTED_REASON`](crate::peers::QUOTA_EXHAUSTED_REASON), a
133    /// domain this project controls — and that pair is what this crate's own
134    /// client keys its refuse-and-come-back classification on. A bare
135    /// `-32029` from a foreign server proves nothing and is classified as an
136    /// unknown fault.
137    ///
138    /// What this code does NOT carry is the quota arithmetic: declines on this
139    /// surface are deliberately uniform, and counters or limits in the message
140    /// would hand an unauthenticated prober the tenant's ceilings. The message
141    /// beside this code is a fixed sentence with no numbers in it.
142    pub const QUOTA_EXHAUSTED: i32 = -32029;
143
144    /// This agent is halted by its operator.
145    ///
146    /// Its own code, beside [`QUOTA_EXHAUSTED`] and
147    /// under the same identification rule — the `(domain, reason)` pair, never
148    /// the numeral — because the two ask opposite things of a caller. A ceiling
149    /// says *come back*; a halt says *somebody is dealing with an incident*,
150    /// and a peer told to back off and retry keeps knocking on the one door
151    /// that means stop. The message is fixed and carries none of the
152    /// operator's reason: the counterparty gets the outcome, not the plane's
153    /// internals.
154    pub const HALTED: i32 = -32030;
155
156    /// This instance is shutting down and did not admit the request.
157    ///
158    /// The third admission refusal, under the same identification rule — the
159    /// `(domain, reason)` pair, never the numeral. It is its own code because
160    /// it is the only one of the three a caller clears by *moving*: a ceiling
161    /// and a halt are facts about the agent, and this is a fact about one
162    /// process serving it. A peer told to back off waits out a window it never
163    /// needed; a peer told to abandon gives up on an agent that is fine.
164    pub const DRAINING: i32 = -32031;
165}
166
167/// Actions this surface asks the policy engine about.
168pub mod action {
169    pub const MESSAGE_SEND: &str = "a2a:message.send";
170    pub const TASK_READ: &str = "a2a:task.read";
171    pub const TASK_CONTINUE: &str = "a2a:task.continue";
172    pub const TASK_CANCEL: &str = "a2a:task.cancel";
173    pub const CARD_EXTENDED: &str = "a2a:card.extended";
174    /// Registering, reading or removing a webhook for a task.
175    pub const TASK_PUSH: &str = "a2a:task.push";
176    /// Supplying an event of a kind, asked on the kind: a continuation is
177    /// delivered as the event its task awaits, so a peer may continue a task
178    /// only with input a rule lets it supply.
179    pub const EVENT_DELIVER: &str = "a2a:event.deliver";
180
181    /// Every action this surface can ask about, so a deployment can enumerate
182    /// what it must write rules for.
183    pub const ALL: &[&str] = &[
184        MESSAGE_SEND,
185        TASK_READ,
186        TASK_CONTINUE,
187        TASK_CANCEL,
188        CARD_EXTENDED,
189        TASK_PUSH,
190        EVENT_DELIVER,
191    ];
192}
193
194/// The policy context every A2A action is asked under: roles, the peer's own
195/// name, and tenant — no chain, no label, no depth. A task action adds the
196/// task's `owner` ([`A2aServer::context`]).
197fn peer_context(caller: &Caller) -> Value {
198    json!({
199        "roles": caller.roles,
200        "peer": caller.actor,
201        "tenant": caller.tenant.as_str(),
202    })
203}
204
205/// Whether `engine` can evaluate every request this surface puts to it.
206///
207/// Each [`action::ALL`] verb is probed in exactly the shapes it is asked in:
208/// with the task's `owner` where it acts on a task, and without it where it
209/// does not — sending, the extended card, the `tasks` listing, and an inline
210/// push registration for a task that does not exist yet. Probing a task action
211/// without its owner would refuse a rule set that reads `context.owner`, which
212/// every such request carries; probing the listing with one would miss a rule
213/// that breaks on the request that carries none. [`A2aServer::hosting`]
214/// refuses a runtime whose engine reports anything: a rule this surface
215/// cannot evaluate declines every peer, and the decline says nothing to the
216/// peer about why.
217#[must_use]
218pub fn policy_problems(engine: &dyn crate::core::PolicyEngine) -> Vec<String> {
219    let caller = Caller::new("preflight", vec!["preflight".to_owned()]);
220    let bare = peer_context(&caller);
221    let owned = A2aServer::context(&caller, Some(&caller.actor));
222    let shapes: &[(&str, &Value)] = &[
223        (action::MESSAGE_SEND, &bare),
224        (action::TASK_READ, &owned),
225        (action::TASK_READ, &bare),
226        (action::TASK_CONTINUE, &owned),
227        (action::TASK_CANCEL, &owned),
228        (action::CARD_EXTENDED, &bare),
229        (action::TASK_PUSH, &owned),
230        (action::TASK_PUSH, &bare),
231        (action::EVENT_DELIVER, &owned),
232    ];
233    let requests: Vec<PolicyRequest<'_>> = shapes
234        .iter()
235        .map(|(action, context)| PolicyRequest {
236            principal: &caller.actor,
237            principal_kind: crate::core::PrincipalKind::Subject,
238            action,
239            resource: "preflight.resource",
240            context,
241        })
242        .collect();
243    engine.preflight(&requests)
244}
245
246/// The `A2A-Version` service parameter.
247const VERSION_HEADER: &str = "a2a-version";
248
249/// The one sentence a full quota answers with, beside
250/// [`code::QUOTA_EXHAUSTED`].
251///
252/// A single constant rather than `why.to_string()` at each site, because the
253/// quota error's own rendering names the counter and its ceiling — numbers an
254/// external caller has no business learning from a decline. Deliberately
255/// digit-free; the test on this path asserts exactly that.
256const QUOTA_EXHAUSTED_MESSAGE: &str =
257    "this agent cannot take the request on right now; retry later";
258
259/// The one sentence a halt answers with, beside [`code::HALTED`]. No reason,
260/// no scope: an incident's description is for the operator's own worklist.
261const HALTED_MESSAGE: &str =
262    "this agent is halted by its operator; do not retry until it is lifted";
263
264/// The one sentence a draining instance answers with, beside [`code::DRAINING`].
265///
266/// No hostname, no instance id and no grace period: which process is going away
267/// is this deployment's topology, and a caller's correct behaviour does not
268/// depend on knowing it.
269const DRAINING_MESSAGE: &str =
270    "this instance is shutting down and did not take the request on; retry";
271
272/// Which of the two admission refusals a quota error is, on the wire.
273///
274/// One function for the three admission sites, because a halt that reached one
275/// of them as a ceiling would be the sibling-divergence shape: a peer told
276/// *retry* by `message/send` and *stop* by `message/stream`.
277fn quota_refusal(e: &crate::quota::QuotaError) -> RpcError {
278    match e {
279        crate::quota::QuotaError::Halted { .. } => RpcError::new(code::HALTED, HALTED_MESSAGE),
280        _ => RpcError::new(code::QUOTA_EXHAUSTED, QUOTA_EXHAUSTED_MESSAGE),
281    }
282}
283
284/// A message its agent cannot read a data subject from, on the wire: invalid
285/// params at every admission site. The error names the binding and why it
286/// did not resolve, never a subject.
287fn subject_unbound(e: &crate::core::RuntimeError) -> RpcError {
288    RpcError::new(code::INVALID_PARAMS, e.to_string())
289}
290
291/// A task's lifecycle state, as A2A names them.
292///
293/// `ProtoJSON` spelling — `TASK_STATE_WORKING`, not `working`. A client matching
294/// on the enum gets nothing from a friendlier spelling.
295///
296/// # Why this is the whole vocabulary, not the part this plane produces
297///
298/// It is **both** an output and an input: `ListTasks` takes a `status` filter,
299/// so a caller may name any state A2A defines and must get a valid, empty
300/// answer rather than `INVALID_PARAMS`. A subset would turn a legitimate query
301/// into a protocol error.
302///
303/// Three of these this plane never *produces*, and the reasons are worth
304/// keeping because they are what a client would otherwise have to guess.
305/// `TASK_STATE_SUBMITTED` means accepted but not yet started, and there is no
306/// such moment here: a run starts inside the call that admits it.
307/// `TASK_STATE_AUTH_REQUIRED` cannot arise either — an unauthenticated request
308/// never reaches dispatch, so there is no task to report it against.
309///
310/// `TASK_STATE_REJECTED` **is** produced, and not as a refusal: a sweep and a
311/// break-glass crossing take it, because
312/// neither is a task any peer submitted. A caller polling one of those is
313/// asking about this plane's record of itself, and gets the same answer as for
314/// a run id that does not exist at all.
315#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
316pub enum TaskState {
317    #[serde(rename = "TASK_STATE_UNSPECIFIED")]
318    Unspecified,
319    #[serde(rename = "TASK_STATE_SUBMITTED")]
320    Submitted,
321    #[serde(rename = "TASK_STATE_WORKING")]
322    Working,
323    #[serde(rename = "TASK_STATE_COMPLETED")]
324    Completed,
325    #[serde(rename = "TASK_STATE_FAILED")]
326    Failed,
327    #[serde(rename = "TASK_STATE_CANCELED")]
328    Canceled,
329    #[serde(rename = "TASK_STATE_INPUT_REQUIRED")]
330    InputRequired,
331    #[serde(rename = "TASK_STATE_REJECTED")]
332    Rejected,
333    #[serde(rename = "TASK_STATE_AUTH_REQUIRED")]
334    AuthRequired,
335}
336
337/// The A2A state a live run's status surfaces as.
338///
339/// Two of these are worth stating because the obvious mapping is wrong.
340///
341/// A **suspended** run is `INPUT_REQUIRED`, not `WORKING`: it is stopped and
342/// will not move until something external happens, and a caller polling a
343/// `WORKING` task waits forever for a task that is not running.
344///
345/// A **quarantined** run is `FAILED`, not `REJECTED`. `REJECTED` means the agent
346/// declined the work; quarantine means it accepted the work, started, and can no
347/// longer be trusted to describe what it did. Reporting that as a refusal tells
348/// the caller nothing happened, when something did.
349///
350/// This and [`sealed_state`] are the **same question asked on two paths**: this
351/// one answers the caller who receives the immediate `SendMessage` response, and
352/// `sealed_state` answers everyone who reads the task back — `GetTask`,
353/// `SubscribeToTask`, and every streamed status update. They must agree, and
354/// `a_live_status_and_its_sealed_outcome_agree` holds them to it over every
355/// variant.
356///
357/// **This is the exhaustive one, and that is the point.** The two were unified
358/// once by making this function delegate to `sealed_state`, which reads as the
359/// tidier direction and is the wrong one: `sealed_state` matches *strings*
360/// behind a `_ => Failed`, so delegating to it deleted the only compile-time
361/// check on the mapping while leaving a comment claiming the compiler still
362/// enforced it. Adding a `RunStatus` variant for something that is **not** a
363/// failure — an authorization wait, a rejection — would then have compiled
364/// cleanly and reported `Failed` on both paths. That is not agreement; it is
365/// the same wrong answer twice, which is strictly harder to notice than two
366/// different ones.
367///
368/// So the enum match is the definition and the string match is checked against
369/// it. A new variant fails to compile *here*, which is where the decision
370/// belongs.
371fn state_of(status: &crate::runtime::RunStatus) -> TaskState {
372    use crate::runtime::RunStatus;
373    match status {
374        RunStatus::Succeeded => TaskState::Completed,
375        RunStatus::Suspended(_) => TaskState::InputRequired,
376        RunStatus::Cancelled { .. } => TaskState::Canceled,
377        RunStatus::Failed(_)
378        | RunStatus::Exhausted(_)
379        | RunStatus::Quarantined(_)
380        // Not `Canceled`: A2A's cancellation means a client asked and the work
381        // stopped, which a peer reads as "nothing of mine is standing". An
382        // abandoned run is the opposite claim — something may well be standing
383        // and nobody could establish what — so it takes the state a peer
384        // investigates rather than the one they dismiss.
385        | RunStatus::Abandoned { .. }
386        // A withdrawal is a *pause* here and still `FAILED` on this wire, which
387        // is the honest mapping rather than a lazy one. A2A has no state for
388        // "stopped, intact, and resumable by somebody else's action": the peer
389        // cannot lift the halt, cannot wait for it, and has nothing to supply —
390        // so `INPUT_REQUIRED` would ask them for input that would never help.
391        // The run's own status keeps the distinction for the party who can act.
392        | RunStatus::Withheld { .. }
393        | RunStatus::Replanning(_) => TaskState::Failed,
394        // Neither is a task, and `REJECTED` is the state that says so: no peer
395        // submitted a sweep or a break-glass crossing, and a run id that
396        // resolves to one is a caller asking about this plane's record of
397        // itself. `FAILED` would claim the agent tried and could not; `CANCELED`
398        // would claim somebody's work was stopped. Both invent a task where
399        // there was none, and a peer reads either as being about *their*
400        // request. Nothing narrows further here on purpose — which of this
401        // plane's internal runs exists is not a peer's question, so the answer
402        // is the same one for a run that does not exist at all.
403        // An observed session joins them for the same reason and one more: it
404        // is not even this plane's work. No peer submitted it, nothing here
405        // executed it, and the only honest thing to tell a caller asking about
406        // somebody else's agent is the answer they get for a run that is none
407        // of their business.
408        RunStatus::Swept
409        | RunStatus::BrokeGlass { .. }
410        | RunStatus::HaltLifted { .. }
411        | RunStatus::HoldReleased { .. }
412        | RunStatus::Observed => TaskState::Rejected,
413    }
414}
415
416/// A2A's `TaskStatus`.
417#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
418pub struct TaskStatus {
419    pub state: TaskState,
420    #[serde(skip_serializing_if = "Option::is_none")]
421    pub message: Option<A2aMessage>,
422    #[serde(skip_serializing_if = "Option::is_none")]
423    pub timestamp: Option<String>,
424}
425
426/// A2A's `Task` — what this plane calls a run.
427#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
428#[serde(rename_all = "camelCase")]
429pub struct A2aTask {
430    pub id: String,
431    /// The case, when the run belongs to one.
432    ///
433    /// A2A's `contextId` is "the thing this conversation is part of", which is
434    /// exactly a case: several runs, one matter, shared state. Mapping it to
435    /// anything else would give a caller a correlation handle that does not
436    /// correlate.
437    #[serde(skip_serializing_if = "Option::is_none")]
438    pub context_id: Option<String>,
439    pub status: TaskStatus,
440    #[serde(skip_serializing_if = "Option::is_none")]
441    pub artifacts: Option<Vec<A2aArtifact>>,
442    #[serde(skip_serializing_if = "Option::is_none")]
443    pub history: Option<Vec<A2aMessage>>,
444    #[serde(skip_serializing_if = "Option::is_none")]
445    pub metadata: Option<Value>,
446}
447
448/// One output produced by an A2A task.
449#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
450#[serde(rename_all = "camelCase")]
451pub struct A2aArtifact {
452    pub artifact_id: String,
453    #[serde(skip_serializing_if = "Option::is_none")]
454    pub name: Option<String>,
455    #[serde(skip_serializing_if = "Option::is_none")]
456    pub description: Option<String>,
457    pub parts: Vec<Part>,
458    #[serde(skip_serializing_if = "Option::is_none")]
459    pub metadata: Option<Value>,
460    #[serde(default, skip_serializing_if = "Vec::is_empty")]
461    pub extensions: Vec<String>,
462}
463
464/// One piece of a message.
465#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
466#[serde(rename_all = "camelCase")]
467pub struct Part {
468    #[serde(skip_serializing_if = "Option::is_none")]
469    pub text: Option<String>,
470    #[serde(skip_serializing_if = "Option::is_none")]
471    pub data: Option<Value>,
472    #[serde(skip_serializing_if = "Option::is_none")]
473    pub raw: Option<String>,
474    #[serde(skip_serializing_if = "Option::is_none")]
475    pub url: Option<String>,
476    #[serde(skip_serializing_if = "Option::is_none")]
477    pub filename: Option<String>,
478    #[serde(skip_serializing_if = "Option::is_none")]
479    pub media_type: Option<String>,
480    #[serde(skip_serializing_if = "Option::is_none")]
481    pub metadata: Option<Value>,
482}
483
484impl Part {
485    /// A text part.
486    #[must_use]
487    pub fn text(text: impl Into<String>) -> Self {
488        Self {
489            text: Some(text.into()),
490            data: None,
491            raw: None,
492            url: None,
493            filename: None,
494            media_type: Some("text/plain".to_owned()),
495            metadata: None,
496        }
497    }
498
499    /// A structured-data part.
500    #[must_use]
501    pub fn data(data: Value) -> Self {
502        Self {
503            text: None,
504            data: Some(data),
505            raw: None,
506            url: None,
507            filename: None,
508            media_type: Some("application/json".to_owned()),
509            metadata: None,
510        }
511    }
512
513    /// A file part carrying its bytes inline, base64-encoded per `ProtoJSON`.
514    #[must_use]
515    pub fn file_raw(
516        raw_base64: impl Into<String>,
517        media_type: impl Into<String>,
518        filename: impl Into<String>,
519    ) -> Self {
520        Self {
521            text: None,
522            data: None,
523            raw: Some(raw_base64.into()),
524            url: None,
525            filename: Some(filename.into()),
526            media_type: Some(media_type.into()),
527            metadata: None,
528        }
529    }
530
531    /// A file part referring to bytes by URL.
532    #[must_use]
533    pub fn file_url(
534        url: impl Into<String>,
535        media_type: impl Into<String>,
536        filename: impl Into<String>,
537    ) -> Self {
538        Self {
539            text: None,
540            data: None,
541            raw: None,
542            url: Some(url.into()),
543            filename: Some(filename.into()),
544            media_type: Some(media_type.into()),
545            metadata: None,
546        }
547    }
548}
549
550/// How a skill shapes its A2A answer, when the default projection is not it.
551///
552/// The default stands for most skills: a string output becomes a text part and
553/// anything else a data part, in one artifact. What that projection cannot say
554/// is *file content* — inline bytes or a URL with a filename — or that the
555/// answer is a quick, stateless `Message` rather than a task's artifact. Both
556/// are ordinary A2A response shapes a peer may expect.
557///
558/// A skill opts in by returning this value (via [`A2aReply::into_value`]) as
559/// its outcome. The runtime journals it as it journals any output — the shape
560/// is a *projection instruction* read at the protocol boundary, not a second
561/// channel around the journal.
562///
563/// ```no_run
564/// # use agentplane::api::a2a::{A2aReply, Part};
565/// # use agentplane::core::{Outcome, Tainted};
566/// // A task whose artifact is a file reference:
567/// let reply = A2aReply::artifact(vec![Part::file_url(
568///     "https://example.com/report.pdf",
569///     "application/pdf",
570///     "report.pdf",
571/// )]);
572/// # let _ = Outcome::done(Tainted::trusted(reply.into_value()));
573/// ```
574#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
575pub struct A2aReply {
576    #[serde(skip_serializing_if = "Option::is_none")]
577    message: Option<Vec<Part>>,
578    #[serde(skip_serializing_if = "Option::is_none")]
579    artifacts: Option<Vec<Vec<Part>>>,
580}
581
582/// The output key an [`A2aReply`] travels under.
583///
584/// A `$`-prefixed marker, like `$a2a_message` on the inbound side: ordinary
585/// skill output is domain data this runtime never interprets, so the one shape
586/// it *does* interpret must be unmistakably deliberate rather than a field
587/// name a domain object happens to share.
588const REPLY_KEY: &str = "$a2a_reply";
589
590impl A2aReply {
591    /// Answer the blocking send with a direct `Message` instead of a task.
592    ///
593    /// Outside a blocking send — a spawned task, a stream, a later `GetTask` —
594    /// there is no message to deliver, and the parts become the task's
595    /// artifact instead: the content survives, only the envelope differs.
596    #[must_use]
597    pub const fn message(parts: Vec<Part>) -> Self {
598        Self {
599            message: Some(parts),
600            artifacts: None,
601        }
602    }
603
604    /// Answer with one artifact holding these parts.
605    #[must_use]
606    pub fn artifact(parts: Vec<Part>) -> Self {
607        Self {
608            message: None,
609            artifacts: Some(vec![parts]),
610        }
611    }
612
613    /// Answer with several artifacts.
614    #[must_use]
615    pub const fn artifacts(artifacts: Vec<Vec<Part>>) -> Self {
616        Self {
617            message: None,
618            artifacts: Some(artifacts),
619        }
620    }
621
622    /// The outcome value a skill returns.
623    #[must_use]
624    pub fn into_value(self) -> Value {
625        json!({ REPLY_KEY: self })
626    }
627
628    /// Read a reply back out of a run's output, **if the run may declare one**.
629    ///
630    /// A projection instruction is authority: it decides whether the answer is
631    /// a task artifact or a direct `Message`, and what parts it carries —
632    /// including a file part naming a URL. Skill output routinely *contains*
633    /// untrusted data: a summariser quotes its input, a declarative agent's
634    /// answer is a model's words, and an echoing skill returns a peer's own
635    /// bytes. So the marker is honoured only from a **trusted** output.
636    ///
637    /// Otherwise a peer could put the marker in its message, have an ordinary
638    /// echoing skill return it, and choose the envelope its own reply arrived
639    /// in — a file URL of the attacker's naming, presented as the agent's
640    /// answer. Model output is a proposal, never authority; so is a peer's
641    /// message.
642    ///
643    /// An untrusted answer still reaches the caller — as the artifact content
644    /// it is, rather than as an instruction about the envelope.
645    pub(super) fn of_output(output: &crate::core::Tainted<Value>) -> Option<Self> {
646        if output.label().is_untrusted() {
647            return None;
648        }
649        serde_json::from_value(output.peek().get(REPLY_KEY)?.clone()).ok()
650    }
651
652    /// The message parts, when this reply is a direct `Message`.
653    pub(super) fn message_parts(&self) -> Option<Vec<Part>> {
654        self.message.clone()
655    }
656
657    /// Every part, for contexts that can only carry artifacts.
658    pub(super) fn artifact_parts(&self) -> Vec<Vec<Part>> {
659        match (&self.artifacts, &self.message) {
660            (Some(artifacts), _) => artifacts.clone(),
661            (None, Some(message)) => vec![message.clone()],
662            (None, None) => Vec::new(),
663        }
664    }
665}
666
667/// A2A's `Message`.
668#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
669#[serde(rename_all = "camelCase")]
670pub struct A2aMessage {
671    pub message_id: String,
672    pub role: String,
673    pub parts: Vec<Part>,
674    #[serde(default, skip_serializing_if = "Option::is_none")]
675    pub context_id: Option<String>,
676    #[serde(default, skip_serializing_if = "Option::is_none")]
677    pub task_id: Option<String>,
678    #[serde(default, skip_serializing_if = "Option::is_none")]
679    pub metadata: Option<Value>,
680    #[serde(default, skip_serializing_if = "Vec::is_empty")]
681    pub extensions: Vec<String>,
682    #[serde(default, skip_serializing_if = "Vec::is_empty")]
683    pub reference_task_ids: Vec<String>,
684}
685
686impl A2aMessage {
687    fn validate_parts(&self) -> Result<(), RpcError> {
688        if self.role != "ROLE_USER" {
689            return Err(RpcError::new(
690                code::INVALID_PARAMS,
691                "an inbound Message role must be ROLE_USER",
692            ));
693        }
694        // The id is half of the admission key — an empty one would fold every
695        // caller's unnamed messages into a single "retry" of the first.
696        if self.message_id.is_empty() {
697            return Err(RpcError::new(
698                code::INVALID_PARAMS,
699                "messageId must be non-empty: it is what lets a retry of this \
700                 message be recognised as the same message",
701            ));
702        }
703        if self.parts.is_empty() {
704            return Err(RpcError::new(
705                code::CONTENT_TYPE_NOT_SUPPORTED,
706                "the message has no parts this agent can read; it accepts text and data parts",
707            ));
708        }
709        for (index, part) in self.parts.iter().enumerate() {
710            // Only the file modalities are refused. A `mediaType` on a text or
711            // data part is a *label* the spec allows on every part type — a
712            // peer sending `text/markdown` text is conformant, and the card's
713            // input modes are a tailoring hint, not a per-part contract the
714            // server may bounce messages on.
715            if part.raw.is_some() || part.url.is_some() {
716                return Err(RpcError::new(
717                    code::CONTENT_TYPE_NOT_SUPPORTED,
718                    format!(
719                        "message.parts[{index}] is file content; this agent card advertises only text/plain and application/json"
720                    ),
721                ));
722            }
723            let variants = usize::from(part.text.is_some()) + usize::from(part.data.is_some());
724            if variants != 1 {
725                return Err(RpcError::new(
726                    code::INVALID_PARAMS,
727                    format!("message.parts[{index}] must contain exactly one of text or data"),
728                ));
729            }
730        }
731        Ok(())
732    }
733
734    /// What the runtime receives as input.
735    ///
736    /// A stable shape, always the same three keys, rather than a clever unwrapping
737    /// that hands a skill a bare string sometimes and an object other times. A
738    /// skill parsing its own input should not have to branch on how many parts
739    /// the caller happened to send. `$a2a_message` preserves the exact protocol
740    /// object for task-history reconstruction; `text` and `data` remain the
741    /// ergonomic projections skills normally consume.
742    fn to_input(&self) -> Value {
743        let text: Vec<&str> = self
744            .parts
745            .iter()
746            .filter_map(|p| p.text.as_deref())
747            .collect();
748        let data: Vec<Value> = self.parts.iter().filter_map(|p| p.data.clone()).collect();
749        json!({
750            "text": text.join("\n"),
751            "data": data,
752            "$a2a_message": self,
753        })
754    }
755
756    /// The skill this message asks for, if it named one.
757    fn requested_skill(&self) -> Option<&str> {
758        self.metadata.as_ref()?.get("skill")?.as_str()
759    }
760}
761
762/// A JSON-RPC 2.0 request.
763#[derive(Debug, Clone, Deserialize)]
764struct RpcRequest {
765    #[serde(default)]
766    jsonrpc: String,
767    #[serde(default)]
768    id: Value,
769    method: String,
770    #[serde(default)]
771    params: Value,
772}
773
774/// The parts of `SendMessageConfiguration` this server acts on.
775#[derive(Debug, Clone, Default, Deserialize)]
776#[serde(rename_all = "camelCase")]
777struct SendConfiguration {
778    /// Return as soon as the task exists, rather than when it finishes.
779    ///
780    /// Blocking is the spec's default and the default here: unset means wait.
781    #[serde(default)]
782    return_immediately: bool,
783    #[serde(default)]
784    accepted_output_modes: Vec<String>,
785    #[serde(default)]
786    history_length: Option<usize>,
787    #[serde(default)]
788    task_push_notification_config: Option<PushRequest>,
789}
790
791impl SendConfiguration {
792    fn validate(&self) -> Result<(), RpcError> {
793        if !self.accepted_output_modes.is_empty()
794            && !self
795                .accepted_output_modes
796                .iter()
797                .any(|mode| matches!(mode.as_str(), "text/plain" | "application/json"))
798        {
799            return Err(RpcError::new(
800                code::CONTENT_TYPE_NOT_SUPPORTED,
801                "acceptedOutputModes contains no mode this agent can produce",
802            ));
803        }
804        Ok(())
805    }
806}
807
808#[derive(Debug, Clone, Deserialize)]
809#[serde(rename_all = "camelCase")]
810struct PushAuthenticationRequest {
811    scheme: String,
812    #[serde(default)]
813    credentials: Option<String>,
814}
815
816#[derive(Debug, Clone, Deserialize)]
817#[serde(rename_all = "camelCase")]
818struct PushRequest {
819    #[serde(default)]
820    id: Option<String>,
821    #[serde(default)]
822    task_id: Option<String>,
823    url: String,
824    #[serde(default)]
825    token: Option<String>,
826    #[serde(default)]
827    authentication: Option<PushAuthenticationRequest>,
828}
829
830/// The longest push configuration id a caller may choose, in bytes.
831const MAX_PUSH_ID_LEN: usize = 128;
832
833/// The most push configurations one task holds.
834///
835/// Every one is a delivery per record, so without a ceiling a peer multiplies
836/// the plane's outbound traffic and store rows by however many it registers.
837/// Checked before the write rather than inside it, so two registrations racing
838/// for the last slot can both land; the bound is on what a peer can build up,
839/// not an exact count.
840const MAX_PUSH_CONFIGS_PER_TASK: usize = 10;
841
842impl PushRequest {
843    /// Refuse an id in the namespace an operator destination owns, or one
844    /// longer than [`MAX_PUSH_ID_LEN`].
845    ///
846    /// A caller and the deployment share one push store, and the two are told
847    /// apart by an id prefix. A caller allowed to write into that namespace could
848    /// point one of the deployment's own destinations at an address it chose —
849    /// and operator destinations are deliberately exempt from the host
850    /// allowlist, HTTPS and the public-address check, because there is supposed
851    /// to be no caller involved. This is the check that keeps that supposition
852    /// true.
853    fn validate(&self) -> Result<(), RpcError> {
854        if self
855            .id
856            .as_deref()
857            .is_some_and(|id| id.len() > MAX_PUSH_ID_LEN)
858        {
859            return Err(RpcError::new(
860                code::INVALID_PARAMS,
861                format!("a pushNotificationConfig id is at most {MAX_PUSH_ID_LEN} bytes"),
862            ));
863        }
864        if self.id.as_deref().is_some_and(crate::push::is_operator_id) {
865            return Err(RpcError::new(
866                code::INVALID_PARAMS,
867                format!(
868                    "a pushNotificationConfig id may not begin with '{}': that namespace \
869                     belongs to destinations this deployment configured for itself",
870                    crate::push::OPERATOR_PREFIX
871                ),
872            ));
873        }
874        Ok(())
875    }
876
877    fn config(&self, task: RunId) -> crate::push::PushConfig {
878        crate::push::PushConfig {
879            id: self.id.clone().unwrap_or_else(|| format!("push-{task}")),
880            task,
881            url: self.url.clone(),
882            token: self.token.clone().map(crate::core::Secret::new),
883            authentication: self.authentication.as_ref().map(|authentication| {
884                crate::push::PushAuthentication {
885                    scheme: authentication.scheme.clone(),
886                    credentials: crate::core::Secret::new(
887                        authentication.credentials.clone().unwrap_or_default(),
888                    ),
889                }
890            }),
891        }
892    }
893}
894
895/// What every method's params may carry.
896///
897/// One struct for every method, which is the same shape the binary's arguments
898/// had before they moved onto per-verb structs: a field belonging to one method
899/// was **silently accepted** by another and did nothing. On the wire that reads
900/// worse than at a command line, because the caller is a stranger who cannot
901/// see the source. `ListTasks` was the case that mattered — a request naming
902/// `contxtId`, or the `context_id` the protocol's own conformance kit sends,
903/// parsed cleanly, dropped the filter, and answered with **every** task the
904/// caller may see, shaped exactly like the scoped list that was asked for.
905///
906/// Two mechanisms close it, and they are different questions.
907/// [`deny_unknown_fields`] refuses a name this surface does not know at all
908/// (`contxtId`, `task_id`), which the A2A specification licenses
909/// outright: A2A §5.5 says JSON field names **MUST** be camelCase, so `context_id` is not an
910/// alternative spelling but a violation. [`FIELDS_BY_METHOD`] refuses a name
911/// this surface knows and *this method* does not (`pageSize` on `CancelTask`).
912///
913/// Neither subsumes the other.
914///
915/// [`deny_unknown_fields`]: https://serde.rs/container-attrs.html
916#[derive(Debug, Clone, Default, Deserialize)]
917#[serde(deny_unknown_fields)]
918struct CommonParams {
919    /// A2A's opaque routing identifier.
920    #[serde(default)]
921    tenant: Option<String>,
922    #[serde(default)]
923    message: Option<A2aMessage>,
924    #[serde(default)]
925    id: Option<String>,
926    #[serde(default)]
927    configuration: Option<SendConfiguration>,
928    /// `SendMessageRequest.metadata`, accepted and not interpreted.
929    ///
930    /// Modelled rather than ignored because unknown fields are refused, and
931    /// the two are different answers: the specification defines this field, so
932    /// a conforming client may send it and must not meet `-32602`. It is
933    /// deliberately not read — opaque caller data has no governed meaning here,
934    /// and inventing one would be a control nothing enforces. What the runtime
935    /// does record about an inbound message is its authenticated sender, which
936    /// is provenance rather than a field the sender chose.
937    ///
938    /// `dead_code` is allowed here for the one case where it is the point: the
939    /// field exists so that deserialization *accepts* the name, and reading it
940    /// is what would be wrong.
941    #[serde(default)]
942    #[allow(dead_code)]
943    metadata: Option<Value>,
944    #[serde(default, rename = "taskId")]
945    push_task: Option<String>,
946    #[serde(default)]
947    url: Option<String>,
948    #[serde(default)]
949    token: Option<String>,
950    #[serde(default)]
951    authentication: Option<PushAuthenticationRequest>,
952    #[serde(default, rename = "contextId")]
953    context_id: Option<String>,
954    #[serde(default)]
955    status: Option<TaskState>,
956    #[serde(default, rename = "pageSize")]
957    page_size: Option<usize>,
958    #[serde(default, rename = "pageToken")]
959    page_token: Option<String>,
960    #[serde(default, rename = "historyLength")]
961    history_length: Option<usize>,
962    #[serde(default, rename = "statusTimestampAfter")]
963    status_timestamp_after: Option<String>,
964    #[serde(default, rename = "includeArtifacts")]
965    include_artifacts: bool,
966}
967
968/// A JSON-RPC error, carrying the HTTP status it should be served with.
969#[derive(Debug, Clone)]
970pub struct RpcError {
971    code: i32,
972    message: String,
973    /// The id of the request being answered — `null` when it could not be read.
974    id: Value,
975    /// Served as HTTP 401 with a bearer challenge rather than 200.
976    unauthenticated: bool,
977}
978
979impl RpcError {
980    fn new(code: i32, message: impl Into<String>) -> Self {
981        Self {
982            code,
983            message: message.into(),
984            id: Value::Null,
985            unauthenticated: false,
986        }
987    }
988
989    /// A caller whose credential was refused: HTTP 401 with
990    /// `WWW-Authenticate`, as A2A requires, so the status says *authenticate*
991    /// rather than the body saying *malformed*.
992    fn unauthenticated(message: impl Into<String>) -> Self {
993        Self {
994            unauthenticated: true,
995            ..Self::new(code::INVALID_REQUEST, message)
996        }
997    }
998
999    fn with_id(mut self, id: Value) -> Self {
1000        self.id = id;
1001        self
1002    }
1003
1004    /// The machine-readable reason A2A 1.0's error-handling rules require
1005    /// beside an A2A-specific code, as a `google.rpc.ErrorInfo` in
1006    /// `error.data` — `(domain, reason)`.
1007    ///
1008    /// Derived from the code rather than declared per call site, because the
1009    /// two are one fact: the spec's table maps each code to exactly one reason
1010    /// token, and a site that could set them independently is a site that can
1011    /// disagree with itself. Standard JSON-RPC codes carry no reason — the
1012    /// spec assigns them none — so they return `None` and the error omits
1013    /// `data` rather than inventing a token.
1014    ///
1015    /// The spec's codes carry the spec's domain. [`code::QUOTA_EXHAUSTED`]
1016    /// carries this plane's own — [`ERROR_DOMAIN`](crate::peers::ERROR_DOMAIN)
1017    /// — and that domain is load-bearing rather than decorative: the numeric
1018    /// code sits in space A2A 1.0 reserves for its own future errors, so the
1019    /// code alone is not proof of what it means. The `(domain, reason)` pair
1020    /// is, and it is what this crate's client keys its back-off
1021    /// classification on.
1022    ///
1023    /// This adds no information a prober does not already have: the reason is
1024    /// a restatement of the code in the same response. The uniform-refusal
1025    /// rule governs what a *model* is told; this is the protocol channel to an
1026    /// authenticated peer.
1027    const fn reason(&self) -> Option<(&'static str, &'static str)> {
1028        const A2A: &str = "a2a-protocol.org";
1029        const SELF_DOMAIN: &str = crate::peers::ERROR_DOMAIN;
1030        const QUOTA_EXHAUSTED_REASON: &str = crate::peers::QUOTA_EXHAUSTED_REASON;
1031        const HALTED_REASON: &str = crate::peers::HALTED_REASON;
1032        const DRAINING_REASON: &str = crate::peers::DRAINING_REASON;
1033        match self.code {
1034            code::TASK_NOT_FOUND => Some((A2A, "TASK_NOT_FOUND")),
1035            code::TASK_NOT_CANCELABLE => Some((A2A, "TASK_NOT_CANCELABLE")),
1036            code::PUSH_NOT_SUPPORTED => Some((A2A, "PUSH_NOTIFICATION_NOT_SUPPORTED")),
1037            code::UNSUPPORTED_OPERATION => Some((A2A, "UNSUPPORTED_OPERATION")),
1038            code::CONTENT_TYPE_NOT_SUPPORTED => Some((A2A, "CONTENT_TYPE_NOT_SUPPORTED")),
1039            code::EXTENDED_CARD_NOT_CONFIGURED => Some((A2A, "EXTENDED_AGENT_CARD_NOT_CONFIGURED")),
1040            code::VERSION_NOT_SUPPORTED => Some((A2A, "VERSION_NOT_SUPPORTED")),
1041            code::QUOTA_EXHAUSTED => Some((SELF_DOMAIN, QUOTA_EXHAUSTED_REASON)),
1042            code::HALTED => Some((SELF_DOMAIN, HALTED_REASON)),
1043            code::DRAINING => Some((SELF_DOMAIN, DRAINING_REASON)),
1044            _ => None,
1045        }
1046    }
1047
1048    /// The JSON-RPC error member, with the required `ErrorInfo` when one
1049    /// applies.
1050    fn body(&self) -> Value {
1051        match self.reason() {
1052            Some((domain, reason)) => json!({
1053                "code": self.code,
1054                "message": self.message,
1055                "data": [{
1056                    "@type": "type.googleapis.com/google.rpc.ErrorInfo",
1057                    "domain": domain,
1058                    "reason": reason,
1059                }],
1060            }),
1061            None => json!({ "code": self.code, "message": self.message }),
1062        }
1063    }
1064}
1065
1066impl IntoResponse for RpcError {
1067    fn into_response(self) -> Response {
1068        // HTTP 200 with a JSON-RPC error body, except for a refused credential.
1069        // JSON-RPC carries its own error channel, and a transport-level status
1070        // for an application-level refusal is how a client ends up retrying a
1071        // permanent decline: an A2A client reads `error.code`, and many treat
1072        // a 5xx as retryable without ever parsing the body. Authentication is
1073        // the transport's question, and A2A answers it with 401.
1074        let body = Json(json!({
1075            "jsonrpc": "2.0",
1076            "id": self.id,
1077            "error": self.body(),
1078        }));
1079        if self.unauthenticated {
1080            return (
1081                StatusCode::UNAUTHORIZED,
1082                [(axum::http::header::WWW_AUTHENTICATE, "Bearer")],
1083                body,
1084            )
1085                .into_response();
1086        }
1087        (StatusCode::OK, body).into_response()
1088    }
1089}
1090
1091/// This plane, served as an A2A agent.
1092#[derive(Clone)]
1093pub struct A2aServer {
1094    runtime: Arc<Runtime>,
1095    auth: Arc<dyn Authenticator>,
1096    policy: Arc<dyn crate::core::PolicyEngine>,
1097    card: AgentCard,
1098    extended: ExtendedAgentCard,
1099    /// Each agent's own card, by agent name.
1100    per_agent: std::collections::BTreeMap<String, crate::peers::AgentCard>,
1101    /// The card's advertised skill ids — what a caller may ask for.
1102    skills: Vec<String>,
1103    push: Option<PushRuntime>,
1104    /// How many of the caller's own tasks one `ListTasks` may examine.
1105    ///
1106    /// See [`Self::filter_scan_budget`] for why this is a ceiling and not a
1107    /// page size.
1108    filter_scan_budget: usize,
1109    /// How many strict replays one `ListTasks` with `includeArtifacts` may
1110    /// perform. See [`Self::artifact_replay_budget`].
1111    artifact_replay_budget: usize,
1112    /// Artifact projections of **sealed** runs, so a poll loop does not buy a
1113    /// strict replay per read. See [`ArtifactCache`].
1114    artifact_cache: Arc<std::sync::Mutex<ArtifactCache>>,
1115    /// Open streams, per caller. See [`Self::streams_per_caller`].
1116    streams: StreamSlots,
1117}
1118
1119/// Default for [`A2aServer::streams_per_caller`].
1120const STREAMS_PER_CALLER: usize = 8;
1121
1122/// How many streams each caller holds open on this server.
1123#[derive(Debug, Clone)]
1124struct StreamSlots {
1125    limit: usize,
1126    open: Arc<std::sync::Mutex<std::collections::HashMap<String, usize>>>,
1127}
1128
1129impl StreamSlots {
1130    fn new(limit: usize) -> Self {
1131        Self {
1132            limit,
1133            open: Arc::default(),
1134        }
1135    }
1136
1137    /// One more stream for `actor`, or `None` at its ceiling.
1138    fn claim(&self, actor: &str) -> Option<StreamSlot> {
1139        let mut open = crate::core::poison::recover(&self.open);
1140        let held = open.entry(actor.to_owned()).or_default();
1141        if *held >= self.limit {
1142            return None;
1143        }
1144        *held += 1;
1145        Some(StreamSlot {
1146            open: Arc::clone(&self.open),
1147            actor: actor.to_owned(),
1148        })
1149    }
1150}
1151
1152/// One open stream, given back when the stream is dropped — whether it ended
1153/// or its client went away.
1154#[derive(Debug)]
1155pub(super) struct StreamSlot {
1156    open: Arc<std::sync::Mutex<std::collections::HashMap<String, usize>>>,
1157    actor: String,
1158}
1159
1160impl Drop for StreamSlot {
1161    fn drop(&mut self) {
1162        let mut open = crate::core::poison::recover(&self.open);
1163        if let Some(held) = open.get_mut(&self.actor) {
1164            *held = held.saturating_sub(1);
1165            if *held == 0 {
1166                open.remove(&self.actor);
1167            }
1168        }
1169    }
1170}
1171
1172/// A bounded in-process cache of sealed runs' artifact projections.
1173///
1174/// A completed task's artifacts exist only as a projection of its journal, so
1175/// every read recovers them by **strict replay** — the most expensive read
1176/// this surface performs. The cache is sound because it holds only sealed
1177/// runs: a seal is the last record a run ever gets, so the projection is a
1178/// pure function of an immutable history and can never go stale. Bounded and
1179/// evicting oldest-first, because an unbounded cache is a leak with a
1180/// performance story; insertion order is enough where every entry is equally
1181/// immutable.
1182#[derive(Debug, Default)]
1183pub(super) struct ArtifactCache {
1184    map: std::collections::HashMap<RunId, Option<Vec<A2aArtifact>>>,
1185    order: std::collections::VecDeque<RunId>,
1186}
1187
1188/// How many sealed runs' artifact projections stay cached per server.
1189const ARTIFACT_CACHE_CAPACITY: usize = 256;
1190
1191/// What one bounded artifact read produced.
1192enum ArtifactRead {
1193    /// The projection — from the cache or from one counted replay.
1194    Artifacts(Option<Vec<A2aArtifact>>),
1195    /// The request's replay budget is spent; the caller must say so.
1196    OverBudget,
1197}
1198
1199impl ArtifactCache {
1200    /// A hit, already wrapped as the read it answers.
1201    fn get(&self, run: RunId) -> Option<ArtifactRead> {
1202        self.map
1203            .get(&run)
1204            .map(|artifacts| ArtifactRead::Artifacts(artifacts.clone()))
1205    }
1206
1207    fn insert(&mut self, run: RunId, artifacts: Option<Vec<A2aArtifact>>) {
1208        if self.map.insert(run, artifacts).is_none() {
1209            self.order.push_back(run);
1210        }
1211        while self.map.len() > ARTIFACT_CACHE_CAPACITY {
1212            let Some(evicted) = self.order.pop_front() else {
1213                break;
1214            };
1215            self.map.remove(&evicted);
1216        }
1217    }
1218}
1219
1220#[derive(Debug, Clone)]
1221struct PushRuntime {
1222    store: Arc<dyn crate::push::PushStore>,
1223    transport: Arc<dyn crate::push::PushTransport>,
1224}
1225
1226/// Durable A2A webhook delivery, driven by an operator scheduler.
1227///
1228/// A name, not a type of its own: [`A2aServer::push_worker`] hands back the
1229/// [`DeliveryWorker`](crate::push::DeliveryWorker) itself, bound to the A2A
1230/// projection. The struct this alias replaced was a pure delegation shell — its
1231/// own docs said so — that re-stated the loop's retry ceiling and its default
1232/// in a second place, which is two copies of one contract with nothing holding
1233/// them together. The cursor discipline lives in `push` because it has nothing
1234/// to do with A2A; what is A2A about this worker is only the projection it was
1235/// constructed with.
1236pub type A2aPushWorker = crate::push::DeliveryWorker;
1237
1238/// Outcome of one bounded push sweep.
1239///
1240/// Re-exported from [`crate::push`], which owns the delivery loop.
1241pub use crate::push::PushSweepReport;
1242
1243/// `StreamResponse` payloads, for **caller-registered** webhooks only.
1244///
1245/// It claims exactly the registrations an operator destination does not — see
1246/// [`crate::push::Outbox`] on why the two share one store and must not serve
1247/// each other's rows.
1248#[derive(Clone)]
1249struct A2aProjection {
1250    runtime: Arc<Runtime>,
1251}
1252
1253impl std::fmt::Debug for A2aProjection {
1254    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1255        f.debug_struct("A2aProjection").finish_non_exhaustive()
1256    }
1257}
1258
1259#[async_trait::async_trait]
1260impl crate::push::Projection for A2aProjection {
1261    async fn messages(
1262        &self,
1263        record: &crate::journal::Record,
1264    ) -> Result<Vec<crate::push::PushMessage>, crate::core::StoreError> {
1265        let case = record.body.case.map(|case| case.to_string());
1266        let payloads =
1267            super::a2a_stream::payloads_for_record(&self.runtime, record, case.as_deref())
1268                .await
1269                // A projection failure is this plane's own bug and always transient
1270                // to the worker; the shape it travels in is the store's error type
1271                // because that is what the seam speaks.
1272                .map_err(|error| crate::core::StoreError::Backend(error.to_string()))?;
1273        // One record can produce a status event and an artifact event, so the
1274        // record's sequence alone is not an identity. The id a receiver
1275        // deduplicates on is the task, the sequence and the position within
1276        // the record — stable across every retry of the same message, distinct
1277        // between the two.
1278        Ok(payloads
1279            .into_iter()
1280            .enumerate()
1281            .map(|(index, payload)| {
1282                crate::push::PushMessage::json(
1283                    format!("{}/{}/{index}", record.body.run, record.body.seq),
1284                    payload,
1285                )
1286                // A2A's media type, so a receiver routes push and streaming
1287                // through the same parser.
1288                .typed("application/a2a+json")
1289            })
1290            .collect())
1291    }
1292
1293    fn namespace(&self) -> crate::push::PushNamespace {
1294        crate::push::PushNamespace::Caller
1295    }
1296}
1297
1298impl std::fmt::Debug for A2aServer {
1299    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1300        f.debug_struct("A2aServer")
1301            .field("agent", &self.card.name)
1302            .field("skills", &self.skills)
1303            .finish_non_exhaustive()
1304    }
1305}
1306
1307/// Why a server could not be built.
1308#[derive(Debug, thiserror::Error)]
1309pub enum ServerSetupError {
1310    #[error(
1311        "the runtime has no policy engine, so every A2A method would be \
1312         unauthorized. A surface reachable by other agents cannot be the one \
1313         place that skips the gate"
1314    )]
1315    NoPolicy,
1316    /// The policy set cannot evaluate a request this surface asks. A2A asks
1317    /// under `{roles, peer, tenant}` alone, so an unscoped rule reading
1318    /// `context.delegation_depth` or `context.label` unguarded declines every
1319    /// peer.
1320    #[error(
1321        "the policy set cannot evaluate the A2A surface's requests: {problems} — \
1322         they carry only `roles`, `peer` and `tenant`; scope the rule to the \
1323         actions it is about, or guard the read with `context has …`"
1324    )]
1325    PolicyUnevaluable { problems: String },
1326    #[error(
1327        "the runtime has no case layer, so this server cannot mint the \
1328         contextId A2A 1.0 requires on every task — a generated contextId \
1329         must be continuable, and continuation here is a case. Build the \
1330         runtime with `.cases(store)`"
1331    )]
1332    NoCases,
1333    #[error("the agent card could not be derived: {0}")]
1334    Card(#[from] crate::manifest::ManifestError),
1335    #[error("push changes the signed Agent Card; configure it before calling signing_cards_with")]
1336    CardAlreadySigned,
1337    #[error(
1338        "an A2A server serves at least one agent, and no manifest was given — the well-known card path must answer with a card describing something"
1339    )]
1340    NoAgents,
1341    #[error(
1342        "agents '{first}' and '{second}' both advertise skill '{skill}', so a request naming it would be a routing decision the caller did not make. A2A dispatch is named, never inferred: give the skill distinct capability names, or serve the two agents from separate planes"
1343    )]
1344    AmbiguousSkill {
1345        skill: String,
1346        first: String,
1347        second: String,
1348    },
1349}
1350
1351impl A2aServer {
1352    /// Serve `manifest`'s agent at `url`.
1353    ///
1354    /// `url` is where callers reach this plane, and it goes on the card — the
1355    /// same deployment-wiring split as everywhere else: an agent's declaration
1356    /// must not change when its address does.
1357    ///
1358    /// # Errors
1359    ///
1360    /// [`ServerSetupError::NoPolicy`] when the runtime has no policy engine, and
1361    /// [`ServerSetupError::Card`] when the manifest's digest cannot be computed.
1362    pub fn new(
1363        runtime: Arc<Runtime>,
1364        auth: Arc<dyn Authenticator>,
1365        security: &crate::peers::CardSecurity,
1366        manifest: &Manifest,
1367        url: impl Into<String>,
1368    ) -> Result<Self, ServerSetupError> {
1369        Self::hosting(runtime, auth, security, &[manifest], url)
1370    }
1371
1372    /// Serve several declared agents from one plane.
1373    ///
1374    /// A2A's well-known card path is singular per host, so the alternative to
1375    /// this is a server per agent. The first manifest is the one the **well-known card
1376    /// describes** — a room's orchestrator, in the shape the CLI already uses —
1377    /// and every agent additionally gets its own full card at
1378    /// [`agent_card_path`](crate::peers::agent_card_path), listed in the
1379    /// [`EXT_AGENT_DIRECTORY`](crate::peers::EXT_AGENT_DIRECTORY) extension so
1380    /// a caller can find them.
1381    ///
1382    /// Skill dispatch spans every agent, because they are all on the runtime
1383    /// already — what was missing was only discovery. Two agents advertising
1384    /// one skill id is **refused**: dispatch names a skill, and a name that
1385    /// resolves to two agents is a routing decision the caller did not make.
1386    ///
1387    /// # Errors
1388    ///
1389    /// [`ServerSetupError::NoPolicy`], [`ServerSetupError::NoCases`],
1390    /// [`ServerSetupError::NoAgents`] for an empty slice, or
1391    /// [`ServerSetupError::AmbiguousSkill`].
1392    pub fn hosting(
1393        runtime: Arc<Runtime>,
1394        auth: Arc<dyn Authenticator>,
1395        security: &crate::peers::CardSecurity,
1396        manifests: &[&Manifest],
1397        url: impl Into<String>,
1398    ) -> Result<Self, ServerSetupError> {
1399        let policy = runtime.policy().ok_or(ServerSetupError::NoPolicy)?.clone();
1400        let mut problems = policy_problems(policy.as_ref());
1401        // A peer with no chain of its own acts under none on this plane, so
1402        // the runtime's requests for it take a shape the build did not probe.
1403        problems.extend(runtime.served_policy_problems());
1404        if !problems.is_empty() {
1405            return Err(ServerSetupError::PolicyUnevaluable {
1406                problems: problems.join("; "),
1407            });
1408        }
1409        if runtime.cases().is_none() {
1410            return Err(ServerSetupError::NoCases);
1411        }
1412        let [primary, rest @ ..] = manifests else {
1413            return Err(ServerSetupError::NoAgents);
1414        };
1415        let url = url.into();
1416
1417        // Every agent's own card, derived exactly as a lone agent's would be —
1418        // same digest, same skills, same ceilings. A card that differed because
1419        // its agent shared a plane would make the plane part of the identity a
1420        // consumer pins, which the room work already refused for manifests.
1421        let mut directory = Vec::new();
1422        let mut per_agent = std::collections::BTreeMap::new();
1423        let mut owner_of_skill: std::collections::BTreeMap<String, String> =
1424            std::collections::BTreeMap::new();
1425        for m in manifests {
1426            let name = m.metadata.name.clone();
1427            let mut agent_card = AgentCard::derive(m, url.clone())?;
1428            security.apply(&mut agent_card);
1429            for skill in &agent_card.skills {
1430                if let Some(other) = owner_of_skill.insert(skill.id.clone(), name.clone())
1431                    && other != name
1432                {
1433                    return Err(ServerSetupError::AmbiguousSkill {
1434                        skill: skill.id.clone(),
1435                        first: other,
1436                        second: name,
1437                    });
1438                }
1439            }
1440            directory.push(serde_json::json!({
1441                "name": name,
1442                "version": m.metadata.version,
1443                "cardPath": crate::peers::agent_card_path(&name),
1444                "manifestDigest": m.digest()?.to_hex(),
1445                "skills": agent_card.skills.iter().map(|s| s.id.clone()).collect::<Vec<_>>(),
1446            }));
1447            per_agent.insert(name, agent_card);
1448        }
1449
1450        let mut card = AgentCard::derive(primary, url.clone())?;
1451        let mut extended = ExtendedAgentCard::derive(primary, url)?;
1452        security.apply(&mut card);
1453        security.apply(&mut extended.public);
1454
1455        // Only when there is more than one, so a single-agent plane's card is
1456        // byte-for-byte what it was: an extension nobody needs is a claim a
1457        // verifier has to understand for nothing.
1458        if !rest.is_empty() {
1459            let ext = crate::peers::AgentExtension {
1460                uri: crate::peers::EXT_AGENT_DIRECTORY.to_owned(),
1461                description: Some(
1462                    "Every agent this plane serves, and where each agent's own card is.".to_owned(),
1463                ),
1464                required: false,
1465                params: Some(serde_json::json!({ "agents": directory })),
1466            };
1467            card.capabilities.extensions.push(ext.clone());
1468            extended.public.capabilities.extensions.push(ext);
1469        }
1470
1471        // The card names the tenant only when there is one to route on. A2A's
1472        // rule is that a client echoes this value back in every request, so
1473        // advertising `default` would make every caller send a routing
1474        // identifier that routes nowhere.
1475        let tenant = runtime.tenant();
1476        if tenant.as_str() != crate::core::TenantId::DEFAULT {
1477            for iface in &mut card.supported_interfaces {
1478                iface.tenant = Some(tenant.to_string());
1479            }
1480            for iface in &mut extended.public.supported_interfaces {
1481                iface.tenant = Some(tenant.to_string());
1482            }
1483        }
1484
1485        // The union, because every agent is already on the runtime and a skill
1486        // that dispatches but is not accepted here would be a refusal the plane
1487        // could have answered.
1488        let skills = owner_of_skill.keys().cloned().collect();
1489        Ok(Self {
1490            runtime,
1491            auth,
1492            policy,
1493            card,
1494            extended,
1495            per_agent,
1496            skills,
1497            push: None,
1498            filter_scan_budget: FILTER_SCAN_BUDGET,
1499            artifact_replay_budget: ARTIFACT_REPLAY_BUDGET,
1500            artifact_cache: Arc::new(std::sync::Mutex::new(ArtifactCache::default())),
1501            streams: StreamSlots::new(STREAMS_PER_CALLER),
1502        })
1503    }
1504
1505    /// Bound how many streams one caller may hold open at once.
1506    ///
1507    /// A `SendStreamingMessage` or `SubscribeToTask` is a connection held for
1508    /// as long as the run lives and a journal read per poll interval — so an
1509    /// unbounded count is a cost any authenticated peer chooses for this
1510    /// plane. Past the ceiling a new stream is refused as back-pressure
1511    /// ([`code::QUOTA_EXHAUSTED`]), before anything is admitted; one ending,
1512    /// or its client going away, frees the slot.
1513    #[must_use]
1514    pub fn streams_per_caller(mut self, limit: usize) -> Self {
1515        self.streams = StreamSlots::new(limit.max(1));
1516        self
1517    }
1518
1519    /// Bound what one `ListTasks` may cost.
1520    ///
1521    /// The spec's `totalSize` is the exact total, and a `status` or `contextId`
1522    /// filter is decided only by reading each candidate's journal — so the
1523    /// cost of a listing is one the *caller* chooses.
1524    ///
1525    /// The budget counts the caller's own tasks. The index a listing reads is
1526    /// narrowed to them, so another caller's volume — or the embedder's — does
1527    /// not refuse this one's listing.
1528    ///
1529    /// Over budget, the request is refused with the narrowing lever named —
1530    /// `statusTimestampAfter` is answered from the index, so tightening it
1531    /// shrinks the candidate set without reading anything. A refusal is honest
1532    /// where a truncated total would be a lie shaped like an answer: the spec
1533    /// requires the exact count, and a bound that quietly stopped counting
1534    /// would report a smaller tenant, not a bounded scan.
1535    #[must_use]
1536    pub fn filter_scan_budget(mut self, budget: usize) -> Self {
1537        self.filter_scan_budget = budget.max(1);
1538        self
1539    }
1540
1541    /// Bound what one `ListTasks` with `includeArtifacts` may cost.
1542    ///
1543    /// A completed task's artifacts are recovered by strict replay, and a page
1544    /// holds up to a hundred tasks — so without a ceiling, one authenticated
1545    /// request buys a hundred full replays, per request, forever. Replays
1546    /// served from the sealed-run cache cost nothing and do not count; past
1547    /// the budget a task's artifacts are **omitted and marked** (see
1548    /// [`ARTIFACTS_OMITTED_KEY`]) rather than the request refused, because
1549    /// unlike `totalSize` an artifact list is per-task enrichment: omitting
1550    /// one is honest as long as it does not look complete, and `GetTask` on
1551    /// the marked id recovers exactly that task's artifacts.
1552    #[must_use]
1553    pub fn artifact_replay_budget(mut self, budget: usize) -> Self {
1554        self.artifact_replay_budget = budget.max(1);
1555        self
1556    }
1557
1558    /// The single-task read, where one replay is the request's stated cost.
1559    ///
1560    /// A separate entry point rather than `artifacts_bounded(.., None)`,
1561    /// because the caller that passes no budget cannot be handed
1562    /// [`ArtifactRead::OverBudget`] — and saying that with `unreachable!` in a
1563    /// request handler makes a caller's argument the only thing between a
1564    /// JSON-RPC method and a panic. Here the type says it instead.
1565    async fn artifacts_unmetered(
1566        &self,
1567        run: RunId,
1568        state: TaskState,
1569    ) -> Result<Option<Vec<A2aArtifact>>, crate::core::RuntimeError> {
1570        match self.artifacts_bounded(run, state, None).await? {
1571            ArtifactRead::Artifacts(artifacts) => Ok(artifacts),
1572            // Unreachable by construction — no budget, nothing to exceed — and
1573            // answered rather than asserted: an empty artifact list is what a
1574            // task with none looks like, which is the honest reading of "this
1575            // read produced none".
1576            ArtifactRead::OverBudget => Ok(None),
1577        }
1578    }
1579
1580    /// A task's artifacts, replaying at most once per sealed run.
1581    ///
1582    /// Consults [`ArtifactCache`] first and fills it for completed runs; the
1583    /// uncached fall-through is [`task_artifacts`]. `budget` is decremented on
1584    /// every actual replay; `None` means unmetered, which reaches here only
1585    /// through [`artifacts_unmetered`](Self::artifacts_unmetered).
1586    ///
1587    /// A budget hit comes back as [`ArtifactRead::OverBudget`], never as an
1588    /// absent list, so the caller must decide how to mark it — silence here
1589    /// would be a bounded result shaped exactly like a complete one.
1590    async fn artifacts_bounded(
1591        &self,
1592        run: RunId,
1593        state: TaskState,
1594        budget: Option<&mut usize>,
1595    ) -> Result<ArtifactRead, crate::core::RuntimeError> {
1596        if state != TaskState::Completed {
1597            return Ok(ArtifactRead::Artifacts(None));
1598        }
1599        if let Some(cached) = crate::core::poison::recover(&self.artifact_cache).get(run) {
1600            return Ok(cached);
1601        }
1602        if let Some(budget) = &budget
1603            && **budget == 0
1604        {
1605            return Ok(ArtifactRead::OverBudget);
1606        }
1607        let artifacts = task_artifacts(&self.runtime, run, state).await?;
1608        if let Some(budget) = budget {
1609            *budget -= 1;
1610        }
1611        crate::core::poison::recover(&self.artifact_cache).insert(run, artifacts.clone());
1612        Ok(ArtifactRead::Artifacts(artifacts))
1613    }
1614
1615    /// Publish a **signed** card.
1616    ///
1617    /// The card is served unauthenticated from a host a caller may not control.
1618    /// TLS says the bytes came from that host; it says nothing about whether the
1619    /// host is the party whose capabilities the card describes. A signature says
1620    /// that, and it keeps saying it after the card has been copied into a
1621    /// registry, a cache, or somebody's repository.
1622    ///
1623    /// Signed here rather than at derivation because the signature covers the
1624    /// **published** card, interface URL and tenant included — those are
1625    /// deployment facts, and a signature taken before they were set would cover
1626    /// a document nobody serves.
1627    ///
1628    /// # Errors
1629    ///
1630    /// If the card cannot be canonicalized.
1631    pub fn signing_cards_with(
1632        mut self,
1633        signer: &dyn crate::peers::CardSigner,
1634    ) -> Result<Self, crate::peers::CardSignatureError> {
1635        self.card.sign(signer)?;
1636        self.extended.public.sign(signer)?;
1637        Ok(self)
1638    }
1639
1640    /// Enable A2A push configuration and expose a durable delivery worker.
1641    ///
1642    /// The worker uses the task journal itself as its outbox. A registration's
1643    /// cursor advances only after a receiver returns 2xx, so a crash between
1644    /// POST and acknowledgement persistence repeats an event and never loses
1645    /// one. Call [`A2aServer::push_worker`] before consuming the server into its
1646    /// router and schedule
1647    /// [`DeliveryWorker::run_once`](crate::push::DeliveryWorker::run_once) from
1648    /// every instance.
1649    pub fn with_push(
1650        mut self,
1651        store: Arc<dyn crate::push::PushStore>,
1652        transport: Arc<dyn crate::push::PushTransport>,
1653    ) -> Result<Self, ServerSetupError> {
1654        if !self.card.signatures.is_empty() || !self.extended.public.signatures.is_empty() {
1655            return Err(ServerSetupError::CardAlreadySigned);
1656        }
1657        self.card.capabilities.push_notifications = true;
1658        self.extended.public.capabilities.push_notifications = true;
1659        self.push = Some(PushRuntime { store, transport });
1660        Ok(self)
1661    }
1662
1663    /// A cloneable worker handle, when push is configured.
1664    ///
1665    /// The worker is [`crate::push::DeliveryWorker`] itself, bound to the A2A
1666    /// projection — rows belonging to an operator
1667    /// [`Outbox`](crate::push::Outbox) are left alone, because the projection
1668    /// declares the caller namespace and the store's due query filters on it.
1669    /// Retry ceiling, backoff and abandonment reporting are the loop's own;
1670    /// see its docs rather than a restatement here.
1671    #[must_use]
1672    pub fn push_worker(&self) -> Option<A2aPushWorker> {
1673        self.push.clone().map(|push| {
1674            crate::push::DeliveryWorker::new(
1675                Arc::clone(self.runtime.journal()),
1676                push.store,
1677                push.transport,
1678                Arc::new(A2aProjection {
1679                    runtime: Arc::clone(&self.runtime),
1680                }),
1681            )
1682        })
1683    }
1684
1685    /// The router.
1686    ///
1687    /// The card path is unauthenticated by design and everything else is not.
1688    ///
1689    /// The RPC endpoint answers with and without a trailing slash, and that is
1690    /// an interoperability fact rather than a courtesy: mainstream HTTP clients
1691    /// that take a base URL — httpx among them, which is what the official A2A
1692    /// conformance kit and the reference Python SDK are built on — resolve a
1693    /// request for `"/"` against the card's interface URL per RFC 3986, so the
1694    /// wire carries `POST {interface}/`. A router that 404s the slash form
1695    /// passes every test written here and fails the first real peer. Exactness
1696    /// elsewhere (method names, version headers) guards protocol *semantics*;
1697    /// a trailing slash is a client-side join artifact with none.
1698    pub fn router(self) -> Router {
1699        Router::new()
1700            .route(WELL_KNOWN_PATH, get(agent_card))
1701            // Unauthenticated like the well-known card, and for the same
1702            // reason: a card a caller must already be trusted to read is a card
1703            // that cannot be discovered.
1704            .route("/agents/{agent}/agent-card.json", get(one_agent_card))
1705            .route("/a2a", post(rpc))
1706            .route("/a2a/", post(rpc))
1707            .with_state(self)
1708    }
1709
1710    /// Authenticate, then authorize.
1711    async fn gate(
1712        &self,
1713        headers: &HeaderMap,
1714        action: &str,
1715        resource: &str,
1716    ) -> Result<Caller, RpcError> {
1717        let caller = self.authenticate(headers).await?;
1718        self.authorize(&caller, action, resource, None)?;
1719        Ok(caller)
1720    }
1721
1722    /// Authenticate, find the task, and refuse it unless this caller admitted
1723    /// it — then authorize, with the owner in the rule's context.
1724    ///
1725    /// A task id is a bearer of nothing. Policy decides what a caller may do
1726    /// with *its* tasks; whose task an id names is decided here, before policy
1727    /// is asked, and a task another peer admitted — or one nobody admitted over
1728    /// this surface — answers exactly as a task that does not exist. Anything
1729    /// else tells a peer holding a guessed or leaked id that it is real.
1730    ///
1731    /// The id is read only once the caller is authenticated, so a malformed
1732    /// one earns an unauthenticated caller nothing but the challenge.
1733    ///
1734    /// Hands back the task's journal, read once, because every caller of this
1735    /// was about to read it anyway.
1736    async fn gate_task(
1737        &self,
1738        headers: &HeaderMap,
1739        action: &str,
1740        params: &CommonParams,
1741    ) -> Result<(Caller, RunId, Vec<crate::journal::Record>), RpcError> {
1742        let caller = self.authenticate(headers).await?;
1743        let run = task_id(params)?;
1744        let records = self.authorized_task(&caller, action, run).await?;
1745        Ok((caller, run, records))
1746    }
1747
1748    /// [`gate_task`](Self::gate_task) for a caller already authenticated, so
1749    /// a request that needs the caller earlier authenticates once.
1750    async fn authorized_task(
1751        &self,
1752        caller: &Caller,
1753        action: &str,
1754        run: RunId,
1755    ) -> Result<Vec<crate::journal::Record>, RpcError> {
1756        let records = self
1757            .runtime
1758            .journal()
1759            .read(run, 1)
1760            .await
1761            .map_err(|e| internal("reading a task's journal", &e))?;
1762        let owner = task_owner(&records);
1763        if owner != Some(caller.actor.as_str()) {
1764            return Err(task_not_found(run));
1765        }
1766        self.authorize(caller, action, &run.to_string(), owner)?;
1767        Ok(records)
1768    }
1769
1770    /// Who is calling, from the credential — and only if it is this tenant.
1771    async fn authenticate(&self, headers: &HeaderMap) -> Result<Caller, RpcError> {
1772        let caller = self
1773            .auth
1774            .authenticate(headers)
1775            .await
1776            .map_err(|e| RpcError::unauthenticated(e.to_string()))?;
1777        // The caller's tenant is the second half of the check `check_tenant`
1778        // starts. That one compares what the *request* asked for against what
1779        // the card advertises; this compares it against what the *credential*
1780        // says. A peer that authenticates into one tenant must not be served
1781        // from another's runs by naming it in a field, and a peer holding a
1782        // valid credential for a different tenant is exactly who would try.
1783        // Refused as a credential, in the same sentence as any other this
1784        // endpoint does not accept: naming the tenant would tell a prober the
1785        // token is valid somewhere else.
1786        if caller.tenant != *self.runtime.tenant() {
1787            return Err(RpcError::unauthenticated(
1788                crate::api::AuthError::Rejected.to_string(),
1789            ));
1790        }
1791        Ok(caller)
1792    }
1793
1794    /// What a rule on this surface can key on.
1795    ///
1796    /// `owner` is present on a task action: the peer that admitted the task,
1797    /// which by the time a rule is asked is always the caller — so a rule set
1798    /// can say "a peer reads its own tasks" in its own words rather than trust
1799    /// this surface to have said it.
1800    fn context(caller: &Caller, owner: Option<&str>) -> Value {
1801        let mut context = peer_context(caller);
1802        if let Some(owner) = owner {
1803            context["owner"] = json!(owner);
1804        }
1805        context
1806    }
1807
1808    fn authorize(
1809        &self,
1810        caller: &Caller,
1811        action: &str,
1812        resource: &str,
1813        owner: Option<&str>,
1814    ) -> Result<(), RpcError> {
1815        let context = Self::context(caller, owner);
1816        match self.policy.authorize(&PolicyRequest {
1817            principal: &caller.actor,
1818            principal_kind: crate::core::PrincipalKind::Subject,
1819            action,
1820            resource,
1821            context: &context,
1822        }) {
1823            PolicyDecision::Permit => Ok(()),
1824            PolicyDecision::Deny { reason } => {
1825                // The determining policy and its reason stay operator-side. A
1826                // Cedar denial names the action, the resource, and the policy
1827                // ids that fired; returned to an external caller that is a
1828                // probe-able map of the authorization vocabulary. A decline
1829                // carries no reason on the wire — the runtime's own denial
1830                // already names the action and resource the gate keyed on. The
1831                // caller learns only that it was declined; the reason reaches
1832                // whoever runs the plane.
1833                tracing::warn!(
1834                    target: "agentplane::a2a",
1835                    action,
1836                    resource,
1837                    reason,
1838                    "A2A request denied at admission"
1839                );
1840                Err(RpcError::new(
1841                    code::INVALID_REQUEST,
1842                    "this request was not permitted",
1843                ))
1844            }
1845            // The rules could not be evaluated. The caller is declined
1846            // exactly as for a denial — telling a peer that the plane's
1847            // policy set is broken is a fact about our operations, and a
1848            // probe would learn which shapes break it — but the operator's
1849            // side says defect rather than refusal, because the fix is the
1850            // policy set and nobody should be reading rules looking for the
1851            // one that fired.
1852            PolicyDecision::Malformed { reason } => {
1853                tracing::error!(
1854                    target: "agentplane::a2a",
1855                    action,
1856                    resource,
1857                    policy_error = true,
1858                    reason,
1859                    "A2A request refused because the policy set could not be evaluated"
1860                );
1861                Err(RpcError::new(
1862                    code::INVALID_REQUEST,
1863                    "this request was not permitted",
1864                ))
1865            }
1866        }
1867    }
1868
1869    fn permits(&self, caller: &Caller, action: &str, resource: &str, owner: Option<&str>) -> bool {
1870        let context = Self::context(caller, owner);
1871        matches!(
1872            self.policy.authorize(&PolicyRequest {
1873                principal: &caller.actor,
1874                principal_kind: crate::core::PrincipalKind::Subject,
1875                action,
1876                resource,
1877                context: &context,
1878            }),
1879            PolicyDecision::Permit
1880        )
1881    }
1882
1883    /// Refuse a version this server does not speak.
1884    ///
1885    /// An **absent** header is a refusal, not a default. The spec says an empty
1886    /// value means 0.3, so a missing header is a 0.3 client — and answering it
1887    /// with 1.0 semantics is how a caller silently loses half the protocol.
1888    /// Matching is on `Major.Minor`, as the spec requires.
1889    fn check_version(headers: &HeaderMap) -> Result<(), RpcError> {
1890        let claimed = headers
1891            .get(VERSION_HEADER)
1892            .and_then(|v| v.to_str().ok())
1893            .unwrap_or("");
1894        let claimed_version = crate::peers::protocol_major_minor(claimed);
1895        if claimed_version == crate::peers::protocol_major_minor(crate::peers::PROTOCOL_VERSION)
1896            && claimed_version.is_some()
1897        {
1898            return Ok(());
1899        }
1900        let seen = if claimed.is_empty() {
1901            "0.3 (no A2A-Version header, which the spec reads as 0.3)".to_owned()
1902        } else {
1903            claimed.to_owned()
1904        };
1905        Err(RpcError::new(
1906            code::VERSION_NOT_SUPPORTED,
1907            format!(
1908                "this agent speaks A2A {}, and the request asked for {seen}",
1909                crate::peers::PROTOCOL_VERSION
1910            ),
1911        ))
1912    }
1913
1914    /// Refuse a request routed to a different tenant.
1915    ///
1916    /// A2A's rule is that the client echoes the `tenant` from the interface it
1917    /// selected on the card, omitting it when the card omits it. So a value that
1918    /// does not match ours is a request meant for somebody else — plausibly
1919    /// another plane behind the same address — and answering it would serve one
1920    /// tenant's caller from another tenant's runs.
1921    fn check_tenant(&self, params: &CommonParams) -> Result<(), RpcError> {
1922        let ours = self.runtime.tenant().as_str();
1923        let advertised = if ours == crate::core::TenantId::DEFAULT {
1924            ""
1925        } else {
1926            ours
1927        };
1928        let sent = params.tenant.as_deref().unwrap_or("");
1929        if sent == advertised {
1930            return Ok(());
1931        }
1932        Err(RpcError::new(
1933            code::INVALID_PARAMS,
1934            format!(
1935                "this endpoint serves the tenant advertised on its card, and \
1936                 the request named '{sent}'. A2A clients echo the `tenant` from \
1937                 the interface they selected; a different value is a request \
1938                 for a different agent"
1939            ),
1940        ))
1941    }
1942}
1943
1944/// The public Agent Card.
1945///
1946/// Unauthenticated on purpose — see the module docs. It is derived from the
1947/// manifest, so it cannot describe a capability the plane would refuse to
1948/// dispatch, which is what makes publishing it safe.
1949async fn agent_card(State(server): State<A2aServer>) -> Json<AgentCard> {
1950    Json(server.card)
1951}
1952
1953/// One agent's own card.
1954///
1955/// A 404 for an unknown name, deliberately naming nothing: the directory
1956/// extension on the well-known card is how a caller learns which names exist,
1957/// and answering an unknown one with the list would make this path a way to
1958/// enumerate a plane without reading the card that governs it.
1959async fn one_agent_card(
1960    State(server): State<A2aServer>,
1961    axum::extract::Path(agent): axum::extract::Path<String>,
1962) -> Response {
1963    server.per_agent.get(&agent).map_or_else(
1964        || (StatusCode::NOT_FOUND, "no such agent on this plane").into_response(),
1965        |card| Json(card.clone()).into_response(),
1966    )
1967}
1968
1969async fn rpc(
1970    State(server): State<A2aServer>,
1971    headers: HeaderMap,
1972    body: Result<Json<RpcRequest>, axum::extract::rejection::JsonRejection>,
1973) -> Response {
1974    let req = match body {
1975        Ok(Json(req)) => req,
1976        // The two refusals are different spec rows and must not share a code:
1977        // a body that is not `application/json` is `ContentTypeNotSupported`
1978        // (-32005) in A2A's own mapping table, while a body that *claims* JSON
1979        // and is not is JSON-RPC's ParseError. Collapsing them tells a caller
1980        // with a wrong header that its serializer is broken.
1981        Err(rejection) => {
1982            let error = match &rejection {
1983                axum::extract::rejection::JsonRejection::MissingJsonContentType(_) => {
1984                    RpcError::new(
1985                        code::CONTENT_TYPE_NOT_SUPPORTED,
1986                        "the request body must be application/json",
1987                    )
1988                }
1989                _ => RpcError::new(code::PARSE_ERROR, "the request body is not valid JSON-RPC"),
1990            };
1991            return error.into_response();
1992        }
1993    };
1994    let id = req.id.clone();
1995
1996    if req.jsonrpc != "2.0" {
1997        return RpcError::new(
1998            code::INVALID_REQUEST,
1999            format!("`jsonrpc` must be \"2.0\", not {:?}", req.jsonrpc),
2000        )
2001        .with_id(id)
2002        .into_response();
2003    }
2004    if let Err(e) = A2aServer::check_version(&headers) {
2005        return e.with_id(id).into_response();
2006    }
2007
2008    // Streaming methods are dispatched first because they answer with a
2009    // different *kind* of response: an SSE body, not a JSON-RPC envelope. Folding
2010    // them into `dispatch` would mean a function whose return type is "a value
2011    // or an entire HTTP response", which is how one of the two paths quietly
2012    // stops setting its content type.
2013    if matches!(
2014        req.method.as_str(),
2015        method::SEND_STREAMING | method::SUBSCRIBE
2016    ) {
2017        return match stream_method(server, headers, req).await {
2018            Ok(sse) => sse.into_response(),
2019            Err(e) => e.with_id(id).into_response(),
2020        };
2021    }
2022
2023    // **Boxed, because this future is the whole runtime.** `dispatch` fans out
2024    // to every A2A method and one of them admits and executes a run, so the
2025    // state machine inlined here carries the executor's own locals — sixteen
2026    // kilobytes of them. Held on the stack it is sixteen kilobytes *per
2027    // concurrent request*, paid by every caller including the ones asking for a
2028    // task's status. One heap allocation per request buys that back, and a
2029    // request that reaches this line is already doing far more work than an
2030    // allocation.
2031    match Box::pin(dispatch(&server, &headers, &req)).await {
2032        Ok(result) => Json(json!({ "jsonrpc": "2.0", "id": id, "result": result })).into_response(),
2033        Err(e) => e.with_id(id).into_response(),
2034    }
2035}
2036
2037/// `SendStreamingMessage` and `SubscribeToTask`.
2038///
2039/// Both are the same thing once the run exists: a view of the journal from a
2040/// point onward. The only difference is whether this call is what created it.
2041#[allow(clippy::too_many_lines)]
2042async fn stream_method(
2043    server: A2aServer,
2044    headers: HeaderMap,
2045    req: RpcRequest,
2046) -> Result<
2047    axum::response::sse::Sse<
2048        impl futures_util::stream::Stream<
2049            Item = Result<axum::response::sse::Event, std::convert::Infallible>,
2050        >,
2051    >,
2052    RpcError,
2053> {
2054    let params = parse_params(&req.method, &req.params)?;
2055    server.check_tenant(&params)?;
2056    if req.method == method::SEND_STREAMING
2057        && let Some(configuration) = &params.configuration
2058    {
2059        configuration.validate()?;
2060    }
2061
2062    // A slot before anything else is done for this caller: a stream is a
2063    // connection held open and a journal polled for as long as the run lives,
2064    // so how many one peer may hold is a bound this surface states. Taken
2065    // before admission, so a caller at its ceiling starts no work it then
2066    // cannot watch; released when the stream ends or the client goes away.
2067    let caller = server.authenticate(&headers).await?;
2068    let slot = server
2069        .streams
2070        .claim(&caller.actor)
2071        .ok_or_else(|| RpcError::new(code::QUOTA_EXHAUSTED, QUOTA_EXHAUSTED_MESSAGE))?;
2072
2073    let (run, records) = if req.method == method::SUBSCRIBE {
2074        let id = task_id(&params)?;
2075        let records = server
2076            .authorized_task(&caller, action::TASK_READ, id)
2077            .await?;
2078        (id, records)
2079    } else {
2080        let Some(message) = params.message.clone() else {
2081            return Err(RpcError::new(
2082                code::INVALID_PARAMS,
2083                "`message` is required by SendStreamingMessage",
2084            ));
2085        };
2086        message.validate_parts()?;
2087        let run = if message.task_id.is_some() {
2088            continue_task(&server, &caller, &message).await?
2089        } else {
2090            let skill = resolve_skill(&server, &message)?;
2091            server.authorize(&caller, action::MESSAGE_SEND, &skill, None)?;
2092            check_context(&server, &message, &caller).await?;
2093            if let Some(push) = params
2094                .configuration
2095                .as_ref()
2096                .and_then(|configuration| configuration.task_push_notification_config.as_ref())
2097            {
2098                validate_inline_push(&server, &caller, &skill, push)?;
2099            }
2100            // Admitted before the stream opens. A stream that begins and *then*
2101            // reports a refusal has already told the client the work started —
2102            // and an SSE body cannot carry a JSON-RPC error the client is looking
2103            // for at that point.
2104            let source = super::peer_source(&caller.actor);
2105            let input = Tainted::from_source(message.to_input(), SourceId::new(&source));
2106            match spawn_a2a(&server, &skill, input, &message, &caller).await {
2107                Ok(run) => {
2108                    if let Some(push) = params.configuration.as_ref().and_then(|configuration| {
2109                        configuration.task_push_notification_config.as_ref()
2110                    }) {
2111                        register_push(&server, push, run, 1).await?;
2112                    }
2113                    run
2114                }
2115                Err(
2116                    crate::core::RuntimeError::PolicyDenied(_)
2117                    | crate::core::RuntimeError::Delegation(_),
2118                ) => {
2119                    return Err(RpcError::new(
2120                        code::UNSUPPORTED_OPERATION,
2121                        "this agent declined the request",
2122                    ));
2123                }
2124                Err(crate::core::RuntimeError::QuotaExceeded(e)) => {
2125                    return Err(quota_refusal(&e));
2126                }
2127                Err(crate::core::RuntimeError::Draining) => {
2128                    return Err(RpcError::new(code::DRAINING, DRAINING_MESSAGE));
2129                }
2130                Err(e @ crate::core::RuntimeError::SubjectUnbound { .. }) => {
2131                    return Err(subject_unbound(&e));
2132                }
2133                Err(e) => return Err(internal("admitting a streamed message", &e)),
2134            }
2135        };
2136        let records = server
2137            .runtime
2138            .journal()
2139            .read(run, 1)
2140            .await
2141            .map_err(|e| internal("reading a task's journal", &e))?;
2142        (run, records)
2143    };
2144
2145    // Decided from the records before anything is replayed: a subscription
2146    // to a finished task is refused, and refusing it after a strict replay to
2147    // build artifacts nobody will receive is a cost the caller chose for us.
2148    let Some((state, detail)) = state_from_history(&records) else {
2149        return Err(task_not_found(run));
2150    };
2151    // `closes` and not a second `matches!` with the same four states: this
2152    // decides whether a subscription is *refused*, and `closes` decides whether
2153    // a stream *ends*. Two spellings of one rule disagree the day somebody adds
2154    // a terminal state to one of them, and the result is either a subscription
2155    // accepted that shuts immediately or one refused that would have streamed.
2156    if req.method == method::SUBSCRIBE && super::a2a_stream::closes(state) {
2157        return Err(RpcError::new(
2158            code::UNSUPPORTED_OPERATION,
2159            "SubscribeToTask requires a non-terminal task",
2160        ));
2161    }
2162    let case = records
2163        .iter()
2164        .find_map(|r| r.body.case.map(|c| c.to_string()));
2165    let from = records.last().map_or(1, |last| last.body.seq + 1);
2166    let mut task = task_of(run, state, &detail, case.clone());
2167    task.artifacts = server
2168        .artifacts_unmetered(run, state)
2169        .await
2170        .map_err(|e| internal("projecting a task's artifacts", &e))?;
2171
2172    Ok(super::a2a_stream::tail(
2173        Arc::clone(&server.runtime),
2174        Arc::clone(&server.artifact_cache),
2175        slot,
2176        run,
2177        case,
2178        req.id,
2179        task,
2180        from,
2181    ))
2182}
2183
2184async fn dispatch(
2185    server: &A2aServer,
2186    headers: &HeaderMap,
2187    req: &RpcRequest,
2188) -> Result<Value, RpcError> {
2189    let params = parse_params(&req.method, &req.params)?;
2190    server.check_tenant(&params)?;
2191
2192    match req.method.as_str() {
2193        // Boxed: a send admits and may run a whole run inline, and that future
2194        // is larger than the others this match holds.
2195        method::SEND_MESSAGE => Box::pin(send_message(server, headers, params)).await,
2196        method::GET_TASK => get_task(server, headers, params).await,
2197        method::CANCEL_TASK => cancel_task(server, headers, params).await,
2198        method::GET_EXTENDED_CARD => get_extended_card(server, headers).await,
2199
2200        // Streaming is handled before dispatch; reaching here means the router
2201        // changed and this arm did not.
2202        method::SEND_STREAMING | method::SUBSCRIBE => Err(RpcError::new(
2203            code::INTERNAL_ERROR,
2204            "a streaming method reached the non-streaming dispatcher",
2205        )),
2206        method::LIST_TASKS => list_tasks(server, headers, &params).await,
2207        method::CREATE_PUSH => push_create(server, headers, &params).await,
2208        method::GET_PUSH => push_get(server, headers, &params).await,
2209        method::LIST_PUSH => push_list(server, headers, &params).await,
2210        method::DELETE_PUSH => push_delete(server, headers, &params).await,
2211        other => Err(RpcError::new(
2212            code::METHOD_NOT_FOUND,
2213            format!("no such A2A method: {other}"),
2214        )),
2215    }
2216}
2217
2218/// Which parameter names each method actually reads.
2219///
2220/// `tenant` is on every row rather than special-cased: A2A's routing identifier
2221/// is orthogonal to the method, and leaving it out of one row would refuse a
2222/// correctly-routed request for that method alone — the kind of hole a table
2223/// exists to make visible.
2224const FIELDS_BY_METHOD: &[(&str, &[&str])] = &[
2225    (
2226        method::SEND_MESSAGE,
2227        &["tenant", "message", "configuration", "metadata"],
2228    ),
2229    (
2230        method::SEND_STREAMING,
2231        &["tenant", "message", "configuration", "metadata"],
2232    ),
2233    (method::GET_TASK, &["tenant", "id", "historyLength"]),
2234    (method::CANCEL_TASK, &["tenant", "id"]),
2235    (method::SUBSCRIBE, &["tenant", "id"]),
2236    (method::GET_EXTENDED_CARD, &["tenant"]),
2237    (
2238        method::LIST_TASKS,
2239        &[
2240            "tenant",
2241            "contextId",
2242            "status",
2243            "pageSize",
2244            "pageToken",
2245            "historyLength",
2246            "statusTimestampAfter",
2247            "includeArtifacts",
2248        ],
2249    ),
2250    (
2251        method::CREATE_PUSH,
2252        &["tenant", "taskId", "id", "url", "token", "authentication"],
2253    ),
2254    (method::GET_PUSH, &["tenant", "taskId", "id"]),
2255    (
2256        method::LIST_PUSH,
2257        &["tenant", "taskId", "pageSize", "pageToken"],
2258    ),
2259    (method::DELETE_PUSH, &["tenant", "taskId", "id"]),
2260];
2261
2262fn parse_params(method: &str, value: &Value) -> Result<CommonParams, RpcError> {
2263    if value.is_null() {
2264        return Ok(CommonParams::default());
2265    }
2266    let Some(object) = value.as_object() else {
2267        return Err(RpcError::new(
2268            code::INVALID_PARAMS,
2269            "A2A method parameters must be a JSON object",
2270        ));
2271    };
2272
2273    // Before deserializing, because the union struct would accept the field and
2274    // the method would then ignore it. An unknown method falls through to the
2275    // dispatcher's own `METHOD_NOT_FOUND`, which is a better answer than a
2276    // parameter complaint about a method that does not exist.
2277    if let Some((_, allowed)) = FIELDS_BY_METHOD.iter().find(|(m, _)| *m == method)
2278        && let Some(stray) = object.keys().find(|k| !allowed.contains(&k.as_str()))
2279    {
2280        return Err(RpcError::new(
2281            code::INVALID_PARAMS,
2282            format!(
2283                "'{stray}' is not a parameter of {method}; it takes {}",
2284                allowed.join(", ")
2285            ),
2286        ));
2287    }
2288
2289    serde_json::from_value(value.clone()).map_err(|error| {
2290        RpcError::new(
2291            code::INVALID_PARAMS,
2292            format!("request parameters do not match the A2A method schema: {error}"),
2293        )
2294    })
2295}
2296
2297/// Answer a message aimed at a task that has already sealed.
2298///
2299/// Decided from the task's own journal: the event store's subscription may
2300/// already be retired, and routing a retransmit through `deliver_to` would
2301/// answer a legitimate retry with a terminal-state error. The journal records
2302/// every delivered continuation message verbatim, so the retry of the message
2303/// that completed the task is recognisable there — and a *new* message aimed
2304/// at a finished task stays refused.
2305fn continue_sealed_task(
2306    run: RunId,
2307    records: &[crate::journal::Record],
2308    message: &A2aMessage,
2309) -> Result<RunId, RpcError> {
2310    let delivered_before = records.iter().any(|record| {
2311        matches!(
2312            record.kind(),
2313            RecordKind::EffectDone { output, .. }
2314                if output
2315                    .get("$a2a_message")
2316                    .and_then(|m| m.get("messageId"))
2317                    .and_then(Value::as_str)
2318                    == Some(message.message_id.as_str())
2319        )
2320    });
2321    if delivered_before {
2322        return Ok(run);
2323    }
2324    Err(RpcError::new(
2325        code::UNSUPPORTED_OPERATION,
2326        "this task has completed and is not waiting for input",
2327    ))
2328}
2329
2330/// Continue one interrupted A2A task with a client message.
2331///
2332/// The run remains append-only: the message becomes the output of the exact
2333/// `event.await` effect on which the task stopped, then ordinary resume replay
2334/// carries execution forward. `EventStore::deliver_to` is task-addressed and
2335/// atomic, so another run sharing the same business correlation key cannot
2336/// consume this message.
2337async fn continue_task(
2338    server: &A2aServer,
2339    caller: &Caller,
2340    message: &A2aMessage,
2341) -> Result<RunId, RpcError> {
2342    let raw = message
2343        .task_id
2344        .as_deref()
2345        .ok_or_else(|| RpcError::new(code::INVALID_PARAMS, "`taskId` is required"))?;
2346    let run = RunId::parse(raw)
2347        .map_err(|_| RpcError::new(code::TASK_NOT_FOUND, format!("no such task: {raw}")))?;
2348    let records = server
2349        .authorized_task(caller, action::TASK_CONTINUE, run)
2350        .await?;
2351    let Some(last) = records.last() else {
2352        return Err(task_not_found(run));
2353    };
2354    if matches!(
2355        last.kind(),
2356        RecordKind::RunSuspended {
2357            reason: crate::core::SuspendReason::AwaitingTime { .. }
2358        }
2359    ) {
2360        return Err(RpcError::new(
2361            code::UNSUPPORTED_OPERATION,
2362            "this task is sleeping until a timer fires and cannot accept input",
2363        ));
2364    }
2365    if !matches!(
2366        last.kind(),
2367        RecordKind::RunSuspended { .. } | RecordKind::RunConcluded { .. }
2368    ) {
2369        return Err(RpcError::new(
2370            code::UNSUPPORTED_OPERATION,
2371            "this task is not waiting for input",
2372        ));
2373    }
2374    if matches!(last.kind(), RecordKind::RunConcluded { .. }) {
2375        return continue_sealed_task(run, &records, message);
2376    }
2377    // Search history rather than only the last record so a transport retry of
2378    // the message that completed a task can still be recognized as the same
2379    // `(source, messageId)` and return the current task instead of a spurious
2380    // terminal-task error. A different message id finds no live subscription
2381    // and remains refused.
2382    let (kind, correlation) = records
2383        .iter()
2384        .rev()
2385        .find_map(|record| match record.kind() {
2386            RecordKind::RunSuspended {
2387                reason:
2388                    crate::core::SuspendReason::AwaitingEvent {
2389                        kind, correlation, ..
2390                    },
2391            } => Some((kind.clone(), correlation.clone())),
2392            _ => None,
2393        })
2394        .ok_or_else(|| {
2395            RpcError::new(
2396                code::UNSUPPORTED_OPERATION,
2397                "this task has no input wait to continue",
2398            )
2399        })?;
2400    let context = records
2401        .iter()
2402        .find_map(|record| record.body.case.map(|case| case.to_string()));
2403    if let Some(sent) = message.context_id.as_deref()
2404        && context.as_deref() != Some(sent)
2405    {
2406        return Err(RpcError::new(
2407            code::INVALID_PARAMS,
2408            "message.contextId does not match the referenced task",
2409        ));
2410    }
2411    // The task gate says this peer may continue its own task; this one says it
2412    // may supply the kind the task awaits.
2413    server.authorize(caller, action::EVENT_DELIVER, &kind, Some(&caller.actor))?;
2414
2415    let event = crate::core::InboundEvent {
2416        // `peer:{actor}`, the one spelling a counterparty's provenance has on
2417        // every transport — see [`super::peer_source`]. This source becomes the
2418        // delivered value's provenance, so a transport-qualified variant here
2419        // would give one counterparty two names.
2420        source: super::peer_source(&caller.actor),
2421        id: message.message_id.clone(),
2422        kind,
2423        correlation,
2424        payload: message.to_input(),
2425        // A counterparty's message, so nobody on this plane minted it. The
2426        // sender is provenance, not an operator, and the one thing this field
2427        // must not become is somewhere a peer can claim to be a person here.
2428        by: None,
2429    };
2430    match server.runtime.deliver_to(run, &event).await {
2431        // `Buffered` is a success too: when the owner's lease is stalled the
2432        // input is held durably and the sweep resumes the run — an internal
2433        // error here would tell the caller a continuation failed that
2434        // actually succeeded-pending.
2435        Ok(
2436            crate::core::Delivery::Resumed { .. }
2437            | crate::core::Delivery::Duplicate
2438            | crate::core::Delivery::Buffered,
2439        ) => Ok(run),
2440        Err(crate::core::RuntimeError::PlanContract(_)) => Err(RpcError::new(
2441            code::UNSUPPORTED_OPERATION,
2442            "this task is no longer waiting for input",
2443        )),
2444        // A run waiting on a human task waits for the worklist's answer, and a
2445        // message addressed to it by task id is not one.
2446        Err(crate::core::RuntimeError::ReservedEventKind { .. }) => Err(RpcError::new(
2447            code::UNSUPPORTED_OPERATION,
2448            "this task is waiting for a decision on this plane's worklist, not for input",
2449        )),
2450        Err(error) => Err(internal("delivering a continuation", &error)),
2451    }
2452}
2453
2454#[allow(clippy::too_many_lines)]
2455async fn send_message(
2456    server: &A2aServer,
2457    headers: &HeaderMap,
2458    params: CommonParams,
2459) -> Result<Value, RpcError> {
2460    if let Some(configuration) = &params.configuration {
2461        configuration.validate()?;
2462    }
2463    let inline_push = params
2464        .configuration
2465        .as_ref()
2466        .and_then(|configuration| configuration.task_push_notification_config.clone());
2467    let Some(message) = params.message else {
2468        return Err(RpcError::new(
2469            code::INVALID_PARAMS,
2470            "`message` is required by SendMessage",
2471        ));
2472    };
2473    message.validate_parts()?;
2474    if message.task_id.is_some() {
2475        let history_length = params
2476            .configuration
2477            .as_ref()
2478            .and_then(|configuration| configuration.history_length);
2479        let caller = server.authenticate(headers).await?;
2480        let run = continue_task(server, &caller, &message).await?;
2481        return get_task(
2482            server,
2483            headers,
2484            CommonParams {
2485                id: Some(run.to_string()),
2486                history_length,
2487                ..CommonParams::default()
2488            },
2489        )
2490        .await
2491        .map(|task| json!({ "task": task }));
2492    }
2493    let skill = resolve_skill(server, &message)?;
2494    let caller = server.gate(headers, action::MESSAGE_SEND, &skill).await?;
2495    check_context(server, &message, &caller).await?;
2496    if let Some(push) = &inline_push {
2497        validate_inline_push(server, &caller, &skill, push)?;
2498    }
2499
2500    // Untrusted, and provenanced to the peer that sent it. A protected sink
2501    // field can then name the one counterparty it will take an amount from, and
2502    // a skill that wants to act on this has to pass a gate to do it.
2503    let source = super::peer_source(&caller.actor);
2504    let input = Tainted::from_source(message.to_input(), SourceId::new(&source));
2505
2506    // Non-blocking: the spec requires returning as soon as the task exists,
2507    // with an in-progress state, leaving the caller to poll `GetTask`. Admission
2508    // still happens synchronously, so a refusal is still an immediate answer —
2509    // returning a task id for a run the gate rejected would hand the caller a
2510    // handle to nothing and turn a decline into a task that never appears.
2511    if params
2512        .configuration
2513        .as_ref()
2514        .is_some_and(|c| c.return_immediately)
2515    {
2516        return match spawn_a2a(server, &skill, input, &message, &caller).await {
2517            Ok(run) => {
2518                if let Some(push) = &inline_push {
2519                    register_push(server, push, run, 1).await?;
2520                }
2521                let case = task_context(server, run).await?;
2522                Ok(json!({
2523                    "task": task_of(run, TaskState::Working, "accepted", case)
2524                }))
2525            }
2526            Err(
2527                crate::core::RuntimeError::PolicyDenied(_)
2528                | crate::core::RuntimeError::Delegation(_),
2529            ) => Ok(json!({ "message": declined(&skill) })),
2530            Err(crate::core::RuntimeError::QuotaExceeded(e)) => Err(quota_refusal(&e)),
2531            Err(crate::core::RuntimeError::Draining) => {
2532                Err(RpcError::new(code::DRAINING, DRAINING_MESSAGE))
2533            }
2534            Err(crate::core::RuntimeError::PlanContract(why)) if message.context_id.is_some() => {
2535                Err(RpcError::new(code::TASK_NOT_FOUND, why))
2536            }
2537            Err(e @ crate::core::RuntimeError::SubjectUnbound { .. }) => Err(subject_unbound(&e)),
2538            Err(e) => Err(internal("admitting a message", &e)),
2539        };
2540    }
2541
2542    let admission = match run_a2a(server, &skill, input, &message, &caller).await {
2543        Ok(admission) => admission,
2544        // A policy denial is the agent *declining*, not the agent breaking.
2545        // Reported as `-32603 Internal error` it reads as "this server is
2546        // faulty, retry later", and the caller retries a decision that will
2547        // never change.
2548        //
2549        // Answered as a `Message` rather than a rejected `Task` because no task
2550        // exists: nothing was admitted, so there is no id to poll and no
2551        // history to fetch. A2A's response is a oneof for exactly this.
2552        Err(
2553            crate::core::RuntimeError::PolicyDenied(_) | crate::core::RuntimeError::Delegation(_),
2554        ) => {
2555            return Ok(json!({ "message": declined(&skill) }));
2556        }
2557        // Back-pressure, not a fault. `-32603` reads as "the far side is
2558        // broken, retry later" and a caller may well retry the same second;
2559        // this says *the agent cannot take this on right now*, which is what a
2560        // ceiling means and what a caller should back off from. The code and
2561        // message are fixed — see `code::QUOTA_EXHAUSTED` for why neither is
2562        // the spec's `-32004` and why the quota arithmetic stays out of it.
2563        Err(crate::core::RuntimeError::QuotaExceeded(e)) => {
2564            return Err(quota_refusal(&e));
2565        }
2566        // Not a fault and not back-pressure from the agent: this process is
2567        // going away and admitted nothing. A caller retrying now reaches a
2568        // different instance and is served.
2569        Err(crate::core::RuntimeError::Draining) => {
2570            return Err(RpcError::new(code::DRAINING, DRAINING_MESSAGE));
2571        }
2572        Err(crate::core::RuntimeError::PlanContract(why)) if message.context_id.is_some() => {
2573            return Err(RpcError::new(code::TASK_NOT_FOUND, why));
2574        }
2575        // The message, not the server: it lacks what the agent's declaration
2576        // reads its data subject from, and the same message is refused again.
2577        Err(e @ crate::core::RuntimeError::SubjectUnbound { .. }) => {
2578            return Err(subject_unbound(&e));
2579        }
2580        Err(e) => return Err(internal("admitting a message", &e)),
2581    };
2582    let outcome = match admission {
2583        crate::runtime::Admission::Fresh(outcome)
2584        | crate::runtime::Admission::Replayed(outcome) => outcome,
2585        // Another instance is executing this message right now. The honest
2586        // answer to a retry is *accepted, already in progress* — a Working
2587        // task naming the run — never a retry-provoking failure.
2588        crate::runtime::Admission::InFlight(run) => {
2589            let case = task_context(server, run).await?;
2590            return Ok(json!({
2591                "task": task_of(run, TaskState::Working, "accepted", case)
2592            }));
2593        }
2594    };
2595
2596    let case = task_context(server, outcome.run_id).await?;
2597    // A declared `Message` reply answers the blocking send directly — A2A's
2598    // response is a oneof for exactly this, and it is the only path with a
2599    // caller still waiting to hand a message to. Returned before push
2600    // registration, because a response with no task has no task to push about;
2601    // the run and its journal exist either way.
2602    if matches!(outcome.status, crate::runtime::RunStatus::Succeeded)
2603        && let Some(reply) = outcome.output.as_ref().and_then(A2aReply::of_output)
2604        && let Some(parts) = reply.message_parts()
2605    {
2606        return Ok(json!({
2607            "message": A2aMessage {
2608                message_id: format!("reply-{}", outcome.run_id),
2609                role: "ROLE_AGENT".to_owned(),
2610                parts,
2611                context_id: case,
2612                task_id: None,
2613                metadata: None,
2614                extensions: Vec::new(),
2615                reference_task_ids: Vec::new(),
2616            }
2617        }));
2618    }
2619    if let Some(push) = &inline_push {
2620        register_push(server, push, outcome.run_id, 1).await?;
2621    }
2622    let mut task = task_of_outcome(&outcome);
2623    task.context_id = case;
2624    // Unset means the full (capped) history, per the protocol's default —
2625    // only an explicit `0` skips the read.
2626    let history_length = params
2627        .configuration
2628        .as_ref()
2629        .and_then(|configuration| configuration.history_length);
2630    if history_length != Some(0) {
2631        let records = server
2632            .runtime
2633            .journal()
2634            .read(outcome.run_id, 1)
2635            .await
2636            .map_err(|_| {
2637                RpcError::new(code::INTERNAL_ERROR, "the task journal could not be read")
2638            })?;
2639        task.history = task_history(
2640            outcome.run_id,
2641            &records,
2642            history_length,
2643            task.context_id.as_deref(),
2644        );
2645    }
2646    Ok(json!({ "task": task }))
2647}
2648
2649async fn run_a2a(
2650    server: &A2aServer,
2651    skill: &str,
2652    input: Tainted<Value>,
2653    message: &A2aMessage,
2654    caller: &Caller,
2655) -> Result<crate::runtime::Admission, crate::core::RuntimeError> {
2656    server
2657        .runtime
2658        .run_under(skill, input, run_terms(message, caller)?)
2659        .await
2660}
2661
2662async fn spawn_a2a(
2663    server: &A2aServer,
2664    skill: &str,
2665    input: Tainted<Value>,
2666    message: &A2aMessage,
2667    caller: &Caller,
2668) -> Result<RunId, crate::core::RuntimeError> {
2669    // Producer-scoped and deduplicating, for the reasons `run_terms` states.
2670    // A duplicate answers with the run the key already admitted, which is
2671    // exactly what a retrying caller needs back.
2672    Ok(server
2673        .runtime
2674        .spawn_under(skill, input, run_terms(message, caller)?)
2675        .await?
2676        .run)
2677}
2678
2679/// The terms every A2A admission runs under: keyed by the producer-scoped
2680/// message id, inside the case the caller named or a fresh one, and **acting
2681/// as the caller** — under the chain its credential carried, or under none.
2682///
2683/// The key is also what makes the run this caller's: its source half is the
2684/// authenticated sender, so every later read or write of the task can ask
2685/// whose it is without a second record — see [`task_owner`].
2686///
2687/// The key carries its producer, in
2688/// [`origin_key`](crate::core::origin_key)'s spelling: a bare `messageId` would
2689/// let two counterparties swallow each other's messages as apparent retries —
2690/// or *join* each other's cases, since correlation matches any open case in the
2691/// tenant. And it is an **admission** key, not only a correlation key: this
2692/// crate's own client keeps `messageId` stable across retries precisely so a
2693/// peer can deduplicate, and a server that correlates without deduplicating
2694/// starts a second run inside the right case. A different key space from an
2695/// event's, and the same construction, because the forgery it refuses is the
2696/// same one.
2697///
2698/// The chain is the caller's, never the plane's: a plane that admitted every
2699/// peer's run under its own chain would answer "on whose behalf" with the same
2700/// name for all of them, and act for a caller whose credential permits less. A
2701/// caller that presented no chain acts under **none**: the plane's chain
2702/// bounds what it may be admitted for, and lends it none of its authority —
2703/// the fallback would make the operator's authority ambient.
2704fn run_terms(
2705    message: &A2aMessage,
2706    caller: &Caller,
2707) -> Result<crate::runtime::RunTerms, crate::core::RuntimeError> {
2708    let keyed = crate::core::origin_key(&admission_source(&caller.actor), &message.message_id);
2709    // The authenticated peer, never a body field: the four-eyes exclusion
2710    // reads it, and a caller who could name it could name somebody else.
2711    let terms = crate::runtime::RunTerms::default()
2712        .once(&keyed)
2713        .served(caller.acting_as.clone())
2714        .admitted_by(&caller.actor);
2715    Ok(match message.context_id.as_deref() {
2716        Some(context) => {
2717            let case = crate::core::CaseId::parse(context).map_err(|_| {
2718                crate::core::RuntimeError::PlanContract(
2719                    "contextId is not a case issued here".into(),
2720                )
2721            })?;
2722            terms.in_case(case)
2723        }
2724        // Always correlated: `A2aServer::new` refuses a runtime without a case
2725        // layer, so every task gets a real, continuable context — the
2726        // contextId returned is one a client can send back, not a string that
2727        // satisfies a schema and continues nothing.
2728        None => terms.correlated(
2729            "a2a.context",
2730            &[crate::core::CorrelationKey::new(
2731                CONTEXT_KEY_NAMESPACE,
2732                &keyed,
2733            )],
2734        ),
2735    })
2736}
2737
2738/// The `source` half of every A2A admission key: this surface, then the
2739/// authenticated sender.
2740///
2741/// Its own namespace rather than the bare [`peer_source`](super::peer_source)
2742/// an inbound event's key carries, because ownership is read back from it: an
2743/// embedder that keys a run with an event's `dedup_key` — the key this crate
2744/// recommends — must not thereby hand that run to the event's sender as an A2A
2745/// task.
2746fn admission_source(actor: &str) -> String {
2747    format!("{ADMISSION_NAMESPACE}{}", super::peer_source(actor))
2748}
2749
2750const ADMISSION_NAMESPACE: &str = "a2a/";
2751
2752/// The peer that admitted a task over this surface, read from the run's
2753/// admission key. `None` for a run nobody admitted here.
2754fn task_owner(records: &[crate::journal::Record]) -> Option<&str> {
2755    records
2756        .first()?
2757        .admission_source()?
2758        .strip_prefix(ADMISSION_NAMESPACE)?
2759        .strip_prefix("peer:")
2760}
2761
2762/// The one answer for a task this caller may not address, whether it exists
2763/// or not.
2764fn task_not_found(run: RunId) -> RpcError {
2765    RpcError::new(code::TASK_NOT_FOUND, format!("no such task: {run}"))
2766}
2767
2768/// A fault on this plane's side, told to the peer in one fixed sentence —
2769/// see [`withheld_fault`](crate::core::withheld_fault).
2770fn internal(doing: &str, error: &dyn std::fmt::Display) -> RpcError {
2771    RpcError::new(
2772        code::INTERNAL_ERROR,
2773        crate::core::withheld_fault("a2a", doing, error),
2774    )
2775}
2776
2777/// Refuse a `contextId` this caller did not open.
2778///
2779/// A context is a case, and joining one puts this caller's run beside another
2780/// peer's work and its case state in reach. So the case must be one a message
2781/// from this caller opened — its correlation key carries the admitting
2782/// sender — and anything else answers as a context that does not exist.
2783async fn check_context(
2784    server: &A2aServer,
2785    message: &A2aMessage,
2786    caller: &Caller,
2787) -> Result<(), RpcError> {
2788    let Some(context) = message.context_id.as_deref() else {
2789        return Ok(());
2790    };
2791    let not_found = || RpcError::new(code::TASK_NOT_FOUND, format!("no such context: {context}"));
2792    let case = crate::core::CaseId::parse(context).map_err(|_| not_found())?;
2793    let Some(cases) = server.runtime.cases() else {
2794        return Err(not_found());
2795    };
2796    let Some(case) = cases
2797        .case(case)
2798        .await
2799        .map_err(|e| internal("reading a context's case", &e))?
2800    else {
2801        return Err(not_found());
2802    };
2803    let ours = admission_source(&caller.actor);
2804    let opened_by_caller = case.correlation.iter().any(|key| {
2805        key.namespace == CONTEXT_KEY_NAMESPACE
2806            && crate::core::origin_source(&key.value) == Some(ours.as_str())
2807    });
2808    if opened_by_caller {
2809        Ok(())
2810    } else {
2811        Err(not_found())
2812    }
2813}
2814
2815/// The correlation namespace a fresh A2A context is opened under.
2816const CONTEXT_KEY_NAMESPACE: &str = "a2a-message";
2817
2818async fn task_context(server: &A2aServer, run: RunId) -> Result<Option<String>, RpcError> {
2819    server
2820        .runtime
2821        .journal()
2822        .read(run, 1)
2823        .await
2824        .map_err(|_| RpcError::new(code::INTERNAL_ERROR, "the task journal could not be read"))
2825        .map(|records| {
2826            records
2827                .iter()
2828                .find_map(|record| record.body.case.map(|case| case.to_string()))
2829        })
2830}
2831
2832pub(super) fn task_of_outcome(outcome: &crate::runtime::RunOutcome) -> A2aTask {
2833    // A declared reply carries its own parts; otherwise the default projection
2834    // stands — a string is a text part, anything else a data part.
2835    let artifacts_parts: Vec<Vec<Part>> =
2836        match outcome.output.as_ref().and_then(A2aReply::of_output) {
2837            Some(reply) => reply.artifact_parts(),
2838            None => vec![vec![match outcome
2839                .output
2840                .as_ref()
2841                .map_or(Value::Null, |o| o.peek().clone())
2842            {
2843                Value::String(text) => Part::text(text),
2844                data => Part::data(data),
2845            }]],
2846        };
2847    A2aTask {
2848        id: outcome.run_id.to_string(),
2849        context_id: None,
2850        status: TaskStatus {
2851            state: state_of(&outcome.status),
2852            message: None,
2853            timestamp: None,
2854        },
2855        artifacts: matches!(outcome.status, crate::runtime::RunStatus::Succeeded).then(|| {
2856            artifacts_parts
2857                .into_iter()
2858                .enumerate()
2859                .map(|(i, parts)| A2aArtifact {
2860                    artifact_id: format!("{}-result-{i}", outcome.run_id),
2861                    name: None,
2862                    description: None,
2863                    parts,
2864                    metadata: None,
2865                    extensions: Vec::new(),
2866                })
2867                .collect()
2868        }),
2869        history: None,
2870        metadata: None,
2871    }
2872}
2873
2874/// Which capability this message asks for.
2875///
2876/// Named or unambiguous — never inferred. See the module docs: picking a
2877/// capability from the content of an untrusted message is a dispatch decision
2878/// made by reading attacker-controlled text.
2879fn resolve_skill(server: &A2aServer, message: &A2aMessage) -> Result<String, RpcError> {
2880    if let Some(asked) = message.requested_skill() {
2881        if server.skills.iter().any(|s| s == asked) {
2882            return Ok(asked.to_owned());
2883        }
2884        return Err(RpcError::new(
2885            code::INVALID_PARAMS,
2886            format!(
2887                "this agent has no skill '{asked}'. Its card advertises: {}",
2888                server.skills.join(", ")
2889            ),
2890        ));
2891    }
2892    match server.skills.as_slice() {
2893        [only] => Ok(only.clone()),
2894        [] => Err(RpcError::new(
2895            code::UNSUPPORTED_OPERATION,
2896            "this agent advertises no skills, so there is nothing to send a \
2897             message to",
2898        )),
2899        many => Err(RpcError::new(
2900            code::INVALID_PARAMS,
2901            format!(
2902                "this agent advertises {} skills, so `message.metadata.skill` \
2903                 must name one of: {}. It is not inferred from the message — \
2904                 choosing what to run by reading the text would let the sender \
2905                 pick the capability",
2906                many.len(),
2907                many.join(", ")
2908            ),
2909        )),
2910    }
2911}
2912
2913async fn get_task(
2914    server: &A2aServer,
2915    headers: &HeaderMap,
2916    params: CommonParams,
2917) -> Result<Value, RpcError> {
2918    let (_, id, records) = server
2919        .gate_task(headers, action::TASK_READ, &params)
2920        .await?;
2921    load_task(server, id, &records, params.history_length).await
2922}
2923
2924async fn load_task(
2925    server: &A2aServer,
2926    id: RunId,
2927    records: &[crate::journal::Record],
2928    history_length: Option<usize>,
2929) -> Result<Value, RpcError> {
2930    // No records is no such task: a run this plane never admitted and a task id
2931    // a caller invented are the same fact from here.
2932    let (state, detail) = state_from_history(records).ok_or_else(|| task_not_found(id))?;
2933    let case = records
2934        .iter()
2935        .find_map(|r| r.body.case.map(|c| c.to_string()));
2936
2937    let mut task = task_of(id, state, &detail, case.clone());
2938    task.history = task_history(id, records, history_length, case.as_deref());
2939    // Unmetered: one task, one replay, and the sealed-run cache spares the
2940    // poll loop that reads the same finished task every few seconds.
2941    task.artifacts = server
2942        .artifacts_unmetered(id, state)
2943        .await
2944        .map_err(|error| internal("projecting a task's artifacts", &error))?;
2945    serde_json::to_value(task).map_err(|error| internal("encoding a task", &error))
2946}
2947
2948#[derive(Debug, Serialize, Deserialize)]
2949#[serde(rename_all = "camelCase")]
2950struct TaskCursor {
2951    updated_at: u64,
2952    run: String,
2953    context_id: Option<String>,
2954    status: Option<TaskState>,
2955    status_timestamp_after: Option<String>,
2956}
2957
2958/// How many index rows `list_tasks` pulls per store round trip.
2959///
2960/// Bounded so no path holds the tenant's whole index in memory — which is what
2961/// the unbounded read this replaced did on *every* call, before going on to
2962/// read the complete journal of every run it returned.
2963const TASK_SCAN: usize = 256;
2964
2965/// Default for [`A2aServer::filter_scan_budget`]: the most candidates one
2966/// `ListTasks` may examine before being refused as too broad.
2967///
2968/// The number is a cost ceiling, not a result limit — results are bounded by
2969/// `pageSize` already. At the default, the worst request a peer can make costs
2970/// on the order of a thousand journal reads, once, and is told how to narrow;
2971/// without it the cost was every run the tenant ever wrote, per request,
2972/// forever.
2973const FILTER_SCAN_BUDGET: usize = 1024;
2974
2975/// Default for [`A2aServer::artifact_replay_budget`]: the most strict replays
2976/// one `ListTasks` with `includeArtifacts` may perform.
2977///
2978/// A cost ceiling in the same spirit as [`FILTER_SCAN_BUDGET`], for the other
2979/// caller-priced read: each completed task's artifacts are a full strict
2980/// replay, a page holds up to a hundred tasks, and the sealed-run cache means
2981/// only first sight of a run pays it. Sixteen replays is a bounded worst case
2982/// per request; a page that wanted more marks the remainder omitted, and a
2983/// second request finds the first sixteen cached.
2984const ARTIFACT_REPLAY_BUDGET: usize = 16;
2985
2986/// The task-metadata key marking artifacts withheld by the replay budget.
2987///
2988/// A bounded result must not be shaped like a complete one — a silent
2989/// truncation reads as an answer. A
2990/// task past the budget is returned without artifacts — indistinguishable from
2991/// several honest states — so the omission says its own name, in `Task.metadata`
2992/// because that is the field A2A gives a server for exactly this kind of
2993/// annotation; a new top-level response member would be a schema no client
2994/// expects. `GetTask` on the marked id recovers the artifacts, and warms the
2995/// cache doing it.
2996pub const ARTIFACTS_OMITTED_KEY: &str = "io.agentplane.a2a/artifactsOmitted";
2997
2998#[allow(clippy::too_many_lines)]
2999async fn list_tasks(
3000    server: &A2aServer,
3001    headers: &HeaderMap,
3002    params: &CommonParams,
3003) -> Result<Value, RpcError> {
3004    let caller = server.gate(headers, action::TASK_READ, "tasks").await?;
3005    let page_size = params.page_size.unwrap_or(50);
3006    if !(1..=100).contains(&page_size) {
3007        return Err(RpcError::new(
3008            code::INVALID_PARAMS,
3009            "pageSize must be between 1 and 100",
3010        ));
3011    }
3012    let after = params
3013        .status_timestamp_after
3014        .as_deref()
3015        .map(|value| {
3016            time::OffsetDateTime::parse(value, &time::format_description::well_known::Rfc3339)
3017                .map_err(|_| {
3018                    RpcError::new(
3019                        code::INVALID_PARAMS,
3020                        "statusTimestampAfter must be an RFC 3339 timestamp",
3021                    )
3022                })
3023        })
3024        .transpose()?;
3025    let cursor = params
3026        .page_token
3027        .as_deref()
3028        .map(decode_task_cursor)
3029        .transpose()?;
3030    if let Some(cursor) = &cursor
3031        && (cursor.context_id != params.context_id
3032            || cursor.status != params.status
3033            || cursor.status_timestamp_after != params.status_timestamp_after)
3034    {
3035        return Err(RpcError::new(
3036            code::INVALID_PARAMS,
3037            "pageToken was issued for different ListTasks filters",
3038        ));
3039    }
3040
3041    // Whether the caller asked anything that can only be answered by reading a
3042    // run's whole journal. `statusTimestampAfter` is not one of those — it
3043    // compares against the activity index's own timestamp.
3044    //
3045    // The candidates are the caller's own runs and nobody else's: the index
3046    // is narrowed to the admission source this surface keys the caller's runs
3047    // with, so another peer's runs and the embedder's cost this listing
3048    // nothing. An unfiltered listing reads the whole journal only of the
3049    // tasks on the page; a content-filtered one reads the whole journal of
3050    // every task it must examine.
3051    let content_filtered = params.context_id.is_some() || params.status.is_some();
3052    let source = admission_source(&caller.actor);
3053
3054    let mut cursor_pos = cursor
3055        .as_ref()
3056        .and_then(|c| RunId::parse(&c.run).ok().map(|run| (c.updated_at, run)));
3057    let mut page: Vec<(RunId, u64, A2aTask)> = Vec::new();
3058    let mut matched: u64 = 0;
3059    let mut has_more = false;
3060    let mut examined: usize = 0;
3061
3062    'scan: loop {
3063        let batch = server
3064            .runtime
3065            .journal()
3066            .recent_runs_from(&source, cursor_pos, TASK_SCAN)
3067            .await
3068            .map_err(|e| internal("reading the task index", &e))?;
3069        if batch.is_empty() {
3070            break;
3071        }
3072        cursor_pos = batch.last().map(|(run, updated)| (*updated, *run));
3073
3074        for (run, updated) in batch {
3075            if after.is_some_and(|cutoff| {
3076                i64::try_from(updated)
3077                    .ok()
3078                    .and_then(|seconds| time::OffsetDateTime::from_unix_timestamp(seconds).ok())
3079                    .is_none_or(|value| value < cutoff)
3080            }) {
3081                // The index is newest first, so every row after this one is
3082                // older still: the cutoff ends the scan, not just this row.
3083                break 'scan;
3084            }
3085
3086            // Whose task it is was settled by the index: every run in this
3087            // range was admitted under the caller's source.
3088            let owner = Some(caller.actor.as_str());
3089            // The ceiling on what one listing may cost, refused rather than
3090            // truncated — see `A2aServer::filter_scan_budget`.
3091            examined += 1;
3092            if examined > server.filter_scan_budget {
3093                return Err(RpcError::new(
3094                    code::INVALID_PARAMS,
3095                    format!(
3096                        "this listing would require examining more than {} tasks to answer \
3097                         exactly — narrow it with statusTimestampAfter and page from there",
3098                        server.filter_scan_budget
3099                    ),
3100                ));
3101            }
3102            // Load-bearing for more than cost: a run the caller may not read
3103            // must not reach the total either, or `totalSize` discloses the
3104            // existence of tasks it cannot see.
3105            if !server.permits(&caller, action::TASK_READ, &run.to_string(), owner) {
3106                continue;
3107            }
3108
3109            let wanted = page.len() < page_size;
3110            if !content_filtered && !wanted {
3111                matched += 1;
3112                has_more = true;
3113                continue;
3114            }
3115            let records = server
3116                .runtime
3117                .journal()
3118                .read(run, 1)
3119                .await
3120                .map_err(|e| internal("reading a task's journal", &e))?;
3121            let task = task_from_records(run, &records, updated, params.history_length);
3122            if params
3123                .context_id
3124                .as_ref()
3125                .is_some_and(|context| task.context_id.as_ref() != Some(context))
3126                || params
3127                    .status
3128                    .as_ref()
3129                    .is_some_and(|status| &task.status.state != status)
3130            {
3131                continue;
3132            }
3133            matched += 1;
3134            if wanted {
3135                page.push((run, updated, task));
3136            } else {
3137                has_more = true;
3138            }
3139        }
3140    }
3141
3142    // Exact, and counted over the same rows the caller was allowed to see. The
3143    // easy mistake is reporting the page's length, which tells every caller the
3144    // total is whatever fits on a screen; the dangerous one is counting the
3145    // store's index directly, which is cheaper and reveals the tasks policy
3146    // just hid. Saturated at the wire type's ceiling: the proto field is an
3147    // `int32`, and a count past it must not overflow a conformant client's
3148    // deserializer.
3149    let total_size = matched.min(u64::try_from(i32::MAX).unwrap_or(u64::MAX));
3150
3151    let visible = &page[..];
3152    let next_page_token = if has_more {
3153        let (run, updated, _) = visible.last().expect("a page with more has a last item");
3154        encode_task_cursor(&TaskCursor {
3155            updated_at: *updated,
3156            run: run.to_string(),
3157            context_id: params.context_id.clone(),
3158            status: params.status,
3159            status_timestamp_after: params.status_timestamp_after.clone(),
3160        })?
3161    } else {
3162        String::new()
3163    };
3164    let mut tasks = Vec::with_capacity(visible.len());
3165    // The replay budget for this whole request. Cache hits are free; each
3166    // cache miss is a strict replay and spends one. Past it a task's artifacts
3167    // are omitted **and marked** — see [`ARTIFACTS_OMITTED_KEY`] for why the
3168    // omission must say its own name.
3169    let mut replays_left = server.artifact_replay_budget;
3170    for (run, _, task) in visible {
3171        let mut task = task.clone();
3172        if params.include_artifacts {
3173            match server
3174                .artifacts_bounded(*run, task.status.state, Some(&mut replays_left))
3175                .await
3176                .map_err(|error| internal("projecting a task's artifacts", &error))?
3177            {
3178                ArtifactRead::Artifacts(artifacts) => task.artifacts = artifacts,
3179                ArtifactRead::OverBudget => {
3180                    task.metadata = Some(json!({ ARTIFACTS_OMITTED_KEY: true }));
3181                }
3182            }
3183        }
3184        tasks.push(task);
3185    }
3186    Ok(json!({
3187        "tasks": tasks,
3188        "nextPageToken": next_page_token,
3189        "pageSize": page_size,
3190        "totalSize": total_size,
3191    }))
3192}
3193
3194fn encode_task_cursor(cursor: &TaskCursor) -> Result<String, RpcError> {
3195    crate::core::canon::to_bytes(cursor)
3196        .map(crate::core::b64::encode_url)
3197        .map_err(|_| RpcError::new(code::INTERNAL_ERROR, "the task cursor could not be encoded"))
3198}
3199
3200fn decode_task_cursor(token: &str) -> Result<TaskCursor, RpcError> {
3201    crate::core::b64::decode_url(token)
3202        .and_then(|bytes| serde_json::from_slice(&bytes).ok())
3203        .ok_or_else(|| RpcError::new(code::INVALID_PARAMS, "pageToken is not a valid task cursor"))
3204}
3205
3206fn task_from_records(
3207    run: RunId,
3208    records: &[crate::journal::Record],
3209    updated: u64,
3210    history_length: Option<usize>,
3211) -> A2aTask {
3212    // A run with no records is one whose first append has not landed. It is
3213    // reported as working rather than refused, because this builds a row in a
3214    // listing where a single unreadable run must not fail the page.
3215    let (state, detail) =
3216        state_from_history(records).unwrap_or_else(|| (TaskState::Working, "running".to_owned()));
3217    let case = records
3218        .iter()
3219        .find_map(|record| record.body.case.map(|case| case.to_string()));
3220    let timestamp = i64::try_from(updated)
3221        .ok()
3222        .and_then(|seconds| time::OffsetDateTime::from_unix_timestamp(seconds).ok())
3223        .and_then(|value| {
3224            value
3225                .format(&time::format_description::well_known::Rfc3339)
3226                .ok()
3227        });
3228    let history = task_history(run, records, history_length, case.as_deref());
3229    let mut task = task_of(run, state, &detail, case);
3230    task.status.timestamp = timestamp;
3231    task.history = history;
3232    task
3233}
3234
3235/// The most messages an unbounded history request returns.
3236///
3237/// The protocol reads an **unset** `historyLength` as "the full history", so
3238/// the bound is the server's, stated here rather than left to whatever the
3239/// journal happens to hold — the same shape as the artifact replay budget.
3240const HISTORY_CAP: usize = 128;
3241
3242fn task_history(
3243    run: RunId,
3244    records: &[crate::journal::Record],
3245    history_length: Option<usize>,
3246    case: Option<&str>,
3247) -> Option<Vec<A2aMessage>> {
3248    // Unset means full (capped); `0` is the explicit request for none —
3249    // the protocol's own default, which a conformant client expecting its
3250    // conversation back depends on.
3251    let limit = history_length.unwrap_or(HISTORY_CAP);
3252    {
3253        if limit == 0 {
3254            return None;
3255        }
3256        let mut history = Vec::new();
3257        for record in records {
3258            let input = match record.kind() {
3259                RecordKind::RunAdmitted { input, .. } => Some(input),
3260                RecordKind::EffectDone { output, .. } if output.get("$a2a_message").is_some() => {
3261                    Some(output)
3262                }
3263                _ => None,
3264            };
3265            let Some(input) = input else { continue };
3266            if let Some(message) = input.get("$a2a_message")
3267                && let Ok(mut message) = serde_json::from_value::<A2aMessage>(message.clone())
3268            {
3269                message.task_id = Some(run.to_string());
3270                if message.context_id.is_none() {
3271                    message.context_id = case.map(ToOwned::to_owned);
3272                }
3273                history.push(message);
3274                continue;
3275            }
3276            // Non-A2A admission retained for old/directly-created tasks.
3277            if history.is_empty() {
3278                let text = input
3279                    .get("text")
3280                    .and_then(Value::as_str)
3281                    .map(ToOwned::to_owned);
3282                let media_type = if text.is_some() {
3283                    "text/plain"
3284                } else {
3285                    "application/json"
3286                };
3287                history.push(A2aMessage {
3288                    message_id: format!("{run}-input"),
3289                    role: "ROLE_USER".to_owned(),
3290                    parts: vec![Part {
3291                        data: text.is_none().then(|| input.clone()),
3292                        text,
3293                        raw: None,
3294                        url: None,
3295                        filename: None,
3296                        media_type: Some(media_type.to_owned()),
3297                        metadata: None,
3298                    }],
3299                    context_id: case.map(ToOwned::to_owned),
3300                    task_id: Some(run.to_string()),
3301                    metadata: None,
3302                    extensions: Vec::new(),
3303                    reference_task_ids: Vec::new(),
3304                });
3305            }
3306        }
3307        let keep_from = history.len().saturating_sub(limit);
3308        (!history.is_empty()).then(|| history.split_off(keep_from))
3309    }
3310}
3311
3312/// The uncached artifact projection: one strict replay of a completed run.
3313///
3314/// Server request paths go through [`A2aServer::artifacts_bounded`], which
3315/// caches sealed runs and meters `ListTasks`, and a stream through
3316/// [`cached_artifacts`]; this stays the raw read for the push projection.
3317pub(super) async fn task_artifacts(
3318    runtime: &Runtime,
3319    run: RunId,
3320    state: TaskState,
3321) -> Result<Option<Vec<A2aArtifact>>, crate::core::RuntimeError> {
3322    if state != TaskState::Completed {
3323        return Ok(None);
3324    }
3325    runtime
3326        .replay(run, crate::runtime::Mode::Strict)
3327        .await
3328        .map(|outcome| task_of_outcome(&outcome).artifacts)
3329}
3330
3331/// A task's artifacts through the sealed-run cache, unmetered.
3332///
3333/// What a stream reads when the run it watches concludes: every subscriber to
3334/// one task would otherwise buy its own strict replay of the same immutable
3335/// history.
3336pub(super) async fn cached_artifacts(
3337    runtime: &Runtime,
3338    cache: &std::sync::Mutex<ArtifactCache>,
3339    run: RunId,
3340    state: TaskState,
3341) -> Result<Option<Vec<A2aArtifact>>, crate::core::RuntimeError> {
3342    if state != TaskState::Completed {
3343        return Ok(None);
3344    }
3345    if let Some(ArtifactRead::Artifacts(artifacts)) = crate::core::poison::recover(cache).get(run) {
3346        return Ok(artifacts);
3347    }
3348    let artifacts = task_artifacts(runtime, run, state).await?;
3349    crate::core::poison::recover(cache).insert(run, artifacts.clone());
3350    Ok(artifacts)
3351}
3352
3353/// A sealed run's outcome word, as an A2A state.
3354///
3355/// The outcome is the same string [`RunStatus::as_str`](crate::runtime::RunStatus::as_str)
3356/// produces, which is what the executor seals with — so this is [`state_of`]
3357/// reached through a string, and the two are held to agreement by
3358/// `a_live_status_and_its_sealed_outcome_agree`.
3359///
3360/// The `_` arm is a *wire* decision rather than a modelling shortcut: an A2A
3361/// client must be told some state, and a word this build cannot interpret is not
3362/// something it may describe as completed or waiting. It is deliberately **not**
3363/// how the runtime itself treats an unrecognised outcome — `resume_is_closed`
3364/// quarantines on one, because refusing to guess is available there and is not
3365/// available here. The agreement test is what keeps this arm from quietly
3366/// swallowing a variant somebody added and forgot to map.
3367pub(super) fn sealed_state(outcome: &str) -> TaskState {
3368    match outcome {
3369        "succeeded" => TaskState::Completed,
3370        "cancelled" => TaskState::Canceled,
3371        "suspended" => TaskState::InputRequired,
3372        // Not tasks at all: nobody submitted this plane's record of a sweep, of
3373        // an operator crossing a tenant boundary or lifting a control, or of a session it merely
3374        // watched somebody else's agent run. `state_of` decides that and this
3375        // agrees, which is the direction the test beside it enforces.
3376        "swept"
3377        | "broke-glass"
3378        | crate::runtime::HALT_LIFTED_OUTCOME
3379        | crate::runtime::HOLD_RELEASED_OUTCOME
3380        | crate::runtime::OBSERVED_OUTCOME => TaskState::Rejected,
3381        _ => TaskState::Failed,
3382    }
3383}
3384
3385/// The A2A state of a run, read from the last record of its history.
3386///
3387/// Every surface that reports a task's state answers here: `tasks/get`,
3388/// `tasks/list`, and the event stream a client subscribes to. They are three
3389/// views of one fact, and three copies of this match would agree until somebody
3390/// added a record kind or reworded a suspension — after which a client polling
3391/// and the same client streaming would be told different things about the same
3392/// run, which is worse than either answer being wrong.
3393///
3394/// **The last record, not any record.** A run that waited, was resumed and
3395/// carried on has a suspension in its history and is not suspended now.
3396///
3397/// An empty history has no state to report and is not this function's to
3398/// invent: the caller decides whether that is a task that does not exist or one
3399/// whose first record has not landed yet, and those get different answers.
3400pub(super) fn state_from_history(
3401    records: &[crate::journal::Record],
3402) -> Option<(TaskState, String)> {
3403    let last = records.last()?;
3404    Some(match last.kind() {
3405        RecordKind::RunSuspended { reason } => (TaskState::InputRequired, reason.to_string()),
3406        RecordKind::RunConcluded { outcome, .. } => (sealed_state(outcome), outcome.clone()),
3407        _ => (TaskState::Working, "running".to_owned()),
3408    })
3409}
3410
3411async fn cancel_task(
3412    server: &A2aServer,
3413    headers: &HeaderMap,
3414    params: CommonParams,
3415) -> Result<Value, RpcError> {
3416    let (caller, id, records) = server
3417        .gate_task(headers, action::TASK_CANCEL, &params)
3418        .await?;
3419    let Some(last) = records.last() else {
3420        return Err(task_not_found(id));
3421    };
3422    // A sealed run is finished, and A2A has a code for exactly this. Accepting
3423    // the request and reporting success would tell the caller a completed run
3424    // is about to stop.
3425    if let RecordKind::RunConcluded { outcome, .. } = last.kind() {
3426        return Err(RpcError::new(
3427            code::TASK_NOT_CANCELABLE,
3428            format!("this task already finished as '{outcome}'"),
3429        ));
3430    }
3431
3432    server
3433        .runtime
3434        .request_cancel(
3435            id,
3436            &crate::core::Operator::authenticated(caller.actor.clone())
3437                .map_err(|e| internal("naming the authenticated caller", &e))?,
3438            "cancelled over A2A",
3439        )
3440        .await
3441        .map_err(|e| match e {
3442            // Concluded between the read above and the request: the same
3443            // answer the read would have given.
3444            crate::core::RuntimeError::AlreadyConcluded { outcome, .. } => RpcError::new(
3445                code::TASK_NOT_CANCELABLE,
3446                format!("this task already finished as '{outcome}'"),
3447            ),
3448            other => internal("requesting a cancellation", &other),
3449        })?;
3450
3451    // The same resolution every other path uses, from the records already in
3452    // hand: A2A 1.0 puts `contextId` on every task, and the one place that
3453    // hard-coded it absent handed a canceling caller a task it could not
3454    // continue or correlate.
3455    let case = records
3456        .iter()
3457        .find_map(|record| record.body.case.map(|case| case.to_string()));
3458
3459    // Still `WORKING`, deliberately. The request is durable, and the run stops
3460    // at its next step boundary — reporting `CANCELED` here would claim it had
3461    // already stopped and unwound, which is exactly what has not happened yet.
3462    serde_json::to_value(task_of(
3463        id,
3464        TaskState::Working,
3465        "cancellation requested",
3466        case,
3467    ))
3468    .map_err(|error| internal("encoding a task", &error))
3469}
3470
3471async fn get_extended_card(server: &A2aServer, headers: &HeaderMap) -> Result<Value, RpcError> {
3472    server
3473        .gate(headers, action::CARD_EXTENDED, &server.card.name)
3474        .await?;
3475    serde_json::to_value(&server.extended).map_err(|e| internal("encoding the extended card", &e))
3476}
3477
3478fn task_id(params: &CommonParams) -> Result<RunId, RpcError> {
3479    let Some(raw) = params.id.as_deref() else {
3480        return Err(RpcError::new(code::INVALID_PARAMS, "`id` is required"));
3481    };
3482    RunId::parse(raw).map_err(|_| {
3483        // Not found rather than invalid params: whether a string is a run id
3484        // this plane issued is not something a caller should learn from the
3485        // shape of the refusal.
3486        RpcError::new(code::TASK_NOT_FOUND, format!("no such task: {raw}"))
3487    })
3488}
3489
3490/// The agent declining, with nothing about *why*.
3491///
3492/// The reason is deliberately absent. A denial reason describes the rule that
3493/// fired, and a rule describes the classification it protects — so a caller who
3494/// can send messages and read refusals can map the policy by probing it. The
3495/// operator sees the reason in the journal, where it belongs; the peer sees that
3496/// it was declined, which is all it can act on anyway.
3497fn declined(skill: &str) -> A2aMessage {
3498    A2aMessage {
3499        message_id: format!("declined-{skill}"),
3500        role: "ROLE_AGENT".to_owned(),
3501        parts: vec![Part {
3502            text: Some("this agent declined the request".to_owned()),
3503            data: None,
3504            raw: None,
3505            url: None,
3506            filename: None,
3507            media_type: Some("text/plain".to_owned()),
3508            metadata: None,
3509        }],
3510        context_id: None,
3511        task_id: None,
3512        metadata: None,
3513        extensions: Vec::new(),
3514        reference_task_ids: Vec::new(),
3515    }
3516}
3517
3518pub(super) fn task_of(run: RunId, state: TaskState, detail: &str, case: Option<String>) -> A2aTask {
3519    A2aTask {
3520        id: run.to_string(),
3521        context_id: case,
3522        status: TaskStatus {
3523            state,
3524            message: Some(A2aMessage {
3525                message_id: format!("{run}-status"),
3526                role: "ROLE_AGENT".to_owned(),
3527                parts: vec![Part {
3528                    text: Some(detail.to_owned()),
3529                    data: None,
3530                    raw: None,
3531                    url: None,
3532                    filename: None,
3533                    media_type: Some("text/plain".to_owned()),
3534                    metadata: None,
3535                }],
3536                context_id: None,
3537                task_id: Some(run.to_string()),
3538                metadata: None,
3539                extensions: Vec::new(),
3540                reference_task_ids: Vec::new(),
3541            }),
3542            timestamp: None,
3543        },
3544        artifacts: None,
3545        history: None,
3546        metadata: None,
3547    }
3548}
3549
3550fn push_runtime(server: &A2aServer) -> Result<&PushRuntime, RpcError> {
3551    server.push.as_ref().ok_or_else(push_not_supported_error)
3552}
3553
3554fn push_not_supported_error() -> RpcError {
3555    RpcError::new(
3556        code::PUSH_NOT_SUPPORTED,
3557        "this agent does not implement push notifications; its card advertises pushNotifications as false",
3558    )
3559}
3560
3561/// The task a push method addresses, refused unless this caller admitted it.
3562///
3563/// Authenticated before the wiring is disclosed, and owned before anything is
3564/// read or written: a webhook registration is a standing instruction to send
3565/// a task's history somewhere, so one peer registering on — or reading,
3566/// listing, deleting the registrations of — another's task would be a read
3567/// of that task with a delivery address attached.
3568async fn push_task<'a>(
3569    server: &'a A2aServer,
3570    headers: &HeaderMap,
3571    raw: Option<&str>,
3572) -> Result<(RunId, &'a PushRuntime), RpcError> {
3573    let caller = server.authenticate(headers).await?;
3574    let push = push_runtime(server)?;
3575    let raw = raw.ok_or_else(|| RpcError::new(code::INVALID_PARAMS, "`taskId` is required"))?;
3576    let task = RunId::parse(raw)
3577        .map_err(|_| RpcError::new(code::TASK_NOT_FOUND, format!("no such task: {raw}")))?;
3578    server
3579        .authorized_task(&caller, action::TASK_PUSH, task)
3580        .await?;
3581    Ok((task, push))
3582}
3583
3584fn push_request(params: &CommonParams) -> Result<PushRequest, RpcError> {
3585    let url = params
3586        .url
3587        .clone()
3588        .ok_or_else(|| RpcError::new(code::INVALID_PARAMS, "`url` is required"))?;
3589    Ok(PushRequest {
3590        id: params.id.clone(),
3591        task_id: params.push_task.clone(),
3592        url,
3593        token: params.token.clone(),
3594        authentication: params.authentication.clone(),
3595    })
3596}
3597
3598fn validate_push_request(server: &A2aServer, request: &PushRequest) -> Result<(), RpcError> {
3599    let push = push_runtime(server)?;
3600    request.validate()?;
3601    let config = request.config(RunId::generate());
3602    if let Some(authentication) = &config.authentication {
3603        authentication
3604            .validate()
3605            .map_err(|error| RpcError::new(code::INVALID_PARAMS, error.to_string()))?;
3606    }
3607    push.transport
3608        .validate(&config)
3609        .map_err(|error| RpcError::new(code::INVALID_PARAMS, error.to_string()))
3610}
3611
3612fn validate_inline_push(
3613    server: &A2aServer,
3614    caller: &Caller,
3615    skill: &str,
3616    request: &PushRequest,
3617) -> Result<(), RpcError> {
3618    if request
3619        .task_id
3620        .as_deref()
3621        .is_some_and(|task| !task.is_empty())
3622    {
3623        return Err(RpcError::new(
3624            code::INVALID_PARAMS,
3625            "taskPushNotificationConfig.taskId must be empty in SendMessage",
3626        ));
3627    }
3628    server.authorize(caller, action::TASK_PUSH, &format!("new:{skill}"), None)?;
3629    validate_push_request(server, request)
3630}
3631
3632async fn register_push(
3633    server: &A2aServer,
3634    request: &PushRequest,
3635    task: RunId,
3636    next_seq: Seq,
3637) -> Result<crate::push::PushConfig, RpcError> {
3638    let push = push_runtime(server)?;
3639    if request
3640        .task_id
3641        .as_deref()
3642        .is_some_and(|configured| !configured.is_empty() && configured != task.to_string())
3643    {
3644        return Err(RpcError::new(
3645            code::INVALID_PARAMS,
3646            "push configuration taskId does not match its task",
3647        ));
3648    }
3649    request.validate()?;
3650    let config = request.config(task);
3651    if let Some(authentication) = &config.authentication {
3652        authentication
3653            .validate()
3654            .map_err(|error| RpcError::new(code::INVALID_PARAMS, error.to_string()))?;
3655    }
3656    push.transport
3657        .validate(&config)
3658        .map_err(|error| RpcError::new(code::INVALID_PARAMS, error.to_string()))?;
3659    // The deployment's own destinations are not the caller's to spend.
3660    let held: Vec<String> = push
3661        .store
3662        .list(task)
3663        .await
3664        .map_err(|error| internal("listing push configurations", &error))?
3665        .into_iter()
3666        .map(|c| c.id)
3667        .filter(|id| !crate::push::is_operator_id(id))
3668        .collect();
3669    if held.len() >= MAX_PUSH_CONFIGS_PER_TASK && !held.contains(&config.id) {
3670        return Err(RpcError::new(
3671            code::INVALID_PARAMS,
3672            format!(
3673                "a task holds at most {MAX_PUSH_CONFIGS_PER_TASK} push configurations; \
3674                 delete one first"
3675            ),
3676        ));
3677    }
3678    push.store
3679        .put(&config, next_seq)
3680        .await
3681        .map_err(|error| internal("storing a push configuration", &error))?;
3682    Ok(config)
3683}
3684
3685async fn push_create(
3686    server: &A2aServer,
3687    headers: &HeaderMap,
3688    params: &CommonParams,
3689) -> Result<Value, RpcError> {
3690    let (task, _) = push_task(server, headers, params.push_task.as_deref()).await?;
3691    let request = push_request(params)?;
3692    let head = server
3693        .runtime
3694        .journal()
3695        .head(task)
3696        .await
3697        .map_err(|error| internal("reading a task's head", &error))?;
3698    let tail = server
3699        .runtime
3700        .journal()
3701        .read(task, head.seq)
3702        .await
3703        .map_err(|error| internal("reading a task's journal", &error))?;
3704    let next_seq = if tail
3705        .last()
3706        .is_some_and(|record| matches!(record.kind(), RecordKind::RunConcluded { .. }))
3707    {
3708        head.seq
3709    } else {
3710        head.seq.saturating_add(1)
3711    };
3712    Ok(register_push(server, &request, task, next_seq)
3713        .await?
3714        .redacted())
3715}
3716
3717async fn push_get(
3718    server: &A2aServer,
3719    headers: &HeaderMap,
3720    params: &CommonParams,
3721) -> Result<Value, RpcError> {
3722    let (task, push) = push_task(server, headers, params.push_task.as_deref()).await?;
3723    let id = params
3724        .id
3725        .as_deref()
3726        .ok_or_else(|| RpcError::new(code::INVALID_PARAMS, "`id` is required"))?;
3727    push.store
3728        .get(task, id)
3729        .await
3730        .map_err(|error| internal("reading a push configuration", &error))?
3731        .map(|config| config.redacted())
3732        .ok_or_else(|| {
3733            RpcError::new(
3734                code::TASK_NOT_FOUND,
3735                format!("no push configuration '{id}' for task {task}"),
3736            )
3737        })
3738}
3739
3740async fn push_list(
3741    server: &A2aServer,
3742    headers: &HeaderMap,
3743    params: &CommonParams,
3744) -> Result<Value, RpcError> {
3745    let (task, push) = push_task(server, headers, params.push_task.as_deref()).await?;
3746    let configs = push
3747        .store
3748        .list(task)
3749        .await
3750        .map_err(|error| internal("listing push configurations", &error))?;
3751    Ok(json!({
3752        "configs": configs.iter().map(crate::push::PushConfig::redacted).collect::<Vec<_>>(),
3753        "nextPageToken": "",
3754    }))
3755}
3756
3757async fn push_delete(
3758    server: &A2aServer,
3759    headers: &HeaderMap,
3760    params: &CommonParams,
3761) -> Result<Value, RpcError> {
3762    let (task, push) = push_task(server, headers, params.push_task.as_deref()).await?;
3763    let id = params
3764        .id
3765        .as_deref()
3766        .ok_or_else(|| RpcError::new(code::INVALID_PARAMS, "`id` is required"))?;
3767    push.store
3768        .delete(task, id)
3769        .await
3770        .map_err(|error| internal("deleting a push configuration", &error))?;
3771    Ok(json!({}))
3772}
3773
3774#[cfg(test)]
3775mod state_agreement_tests {
3776    use super::{TaskState, sealed_state, state_of};
3777
3778    /// The same task must not have two states depending on which path a client took.
3779    ///
3780    /// `state_of` answers the caller holding the immediate `SendMessage`
3781    /// response; `sealed_state` answers `GetTask`, `SubscribeToTask` and every
3782    /// streamed status update. They read the same run.
3783    ///
3784    /// The check runs over the crate's one `RunStatus` list, so adding a variant
3785    /// fails to compile in `state_of` (the match is exhaustive) and, once
3786    /// somebody maps it there, fails *here* until `sealed_state` is taught the
3787    /// same answer. That is the ordering the defect needs: the enum decides, the
3788    /// string agrees. Unifying them the other way — `state_of` delegating to
3789    /// `sealed_state` — makes both paths return the `_ => Failed` fallback for a
3790    /// new variant, with nothing to notice, and that is what this test exists to
3791    /// stop being reintroduced as a tidy-up.
3792    #[test]
3793    fn a_live_status_and_its_sealed_outcome_agree() {
3794        let statuses = crate::runtime::every_status();
3795        assert_eq!(
3796            statuses.len(),
3797            14,
3798            "a RunStatus variant was added or removed — decide which A2A state it \
3799             surfaces as, in `state_of` and in `sealed_state` both"
3800        );
3801        for status in &statuses {
3802            assert_eq!(
3803                state_of(status),
3804                sealed_state(status.as_str()),
3805                "'{}' surfaces as one state live and another once sealed",
3806                status.as_str()
3807            );
3808        }
3809    }
3810
3811    /// And the mapping is not one answer for everything.
3812    ///
3813    /// Without this, `sealed_state` and `state_of` could both be changed to
3814    /// return `Failed` unconditionally and the agreement test above would pass
3815    /// perfectly — every A2A client would see every task as failed.
3816    #[test]
3817    fn the_mapping_distinguishes_more_than_failure() {
3818        let seen: std::collections::BTreeSet<_> = crate::runtime::every_status()
3819            .iter()
3820            .map(|status| format!("{:?}", state_of(status)))
3821            .collect();
3822        assert!(
3823            seen.len() >= 4,
3824            "the run-status mapping collapsed to {seen:?}; a client cannot tell \
3825             completed from cancelled from waiting"
3826        );
3827        assert_eq!(
3828            state_of(&crate::runtime::RunStatus::Succeeded),
3829            TaskState::Completed
3830        );
3831    }
3832}