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(¶ms)?;
2056 if req.method == method::SEND_STREAMING
2057 && let Some(configuration) = ¶ms.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(¶ms)?;
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(¶ms)?;
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, ¶ms).await,
2207 method::CREATE_PUSH => push_create(server, headers, ¶ms).await,
2208 method::GET_PUSH => push_get(server, headers, ¶ms).await,
2209 method::LIST_PUSH => push_list(server, headers, ¶ms).await,
2210 method::DELETE_PUSH => push_delete(server, headers, ¶ms).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) = ¶ms.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, ¶ms)
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, ¶ms)
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}