use std::sync::{Arc, Mutex};
use agentplane::api::a2a::{A2aServer, method};
use agentplane::api::{AuthError, Authenticator, Caller};
use agentplane::core::{
Outcome, PolicyBundleIdentity, PolicyDecision, PolicyEngine, PolicyRequest, Skill,
SkillDescriptor, SkillError, Tainted,
};
use agentplane::manifest::Manifest;
use agentplane::peers::CardSecurity;
use agentplane::runtime::{Agent, Runtime, StepCtx};
use agentplane::store::RedbStore;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use serde_json::{Value, json};
use tower::ServiceExt as _;
const CHECKER: &str = r#"
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata:
name: settlement-checker
version: "1.0.0"
spec:
identity:
role: "Decide whether a settlement instruction may proceed."
constraints: "Answer only from the instruction. Do not invent counterparties."
budgets:
max_steps: 5
capabilities:
provides: [settlement.check]
"#;
#[derive(Debug, Clone, PartialEq, Eq)]
struct Seen {
untrusted: bool,
provenance: Vec<String>,
}
#[derive(Debug, Clone)]
struct Checks(Arc<Mutex<Vec<Seen>>>);
#[async_trait::async_trait]
impl Skill for Checks {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("settlement.check").provides("settlement.check")
}
async fn invoke(
&self,
_cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let label = input.label();
self.0.lock().expect("not poisoned").push(Seen {
untrusted: label.is_untrusted(),
provenance: label.provenance.iter().map(ToString::to_string).collect(),
});
Ok(Outcome::done(Tainted::trusted(json!({"cleared": true}))))
}
}
#[derive(Debug)]
struct BearerAuth;
#[async_trait::async_trait]
impl Authenticator for BearerAuth {
async fn authenticate(&self, headers: &axum::http::HeaderMap) -> Result<Caller, AuthError> {
let token = headers
.get("authorization")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.strip_prefix("Bearer "))
.ok_or(AuthError::Missing)?;
Ok(Caller::new(token, vec!["peer".to_owned()]))
}
}
#[derive(Debug)]
struct PermitPeers;
impl PolicyEngine for PermitPeers {
fn authorize(&self, _request: &PolicyRequest<'_>) -> PolicyDecision {
PolicyDecision::Permit
}
fn bundle(&self) -> PolicyBundleIdentity {
PolicyBundleIdentity::new(
agentplane::core::Digest::of(b"example.permit-peers"),
"example/permit-peers-v1",
)
}
}
async fn rpc(router: &axum::Router, token: Option<&str>, body: Value) -> (StatusCode, Value) {
let mut request = Request::builder()
.method("POST")
.uri("/a2a")
.header("a2a-version", "1.0");
if let Some(token) = token {
request = request.header("authorization", format!("Bearer {token}"));
}
let response = router
.clone()
.oneshot(
request
.header("content-type", "application/json")
.body(Body::from(body.to_string()))
.expect("request"),
)
.await
.expect("router answered");
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), 1 << 20)
.await
.expect("body");
let value = serde_json::from_slice(&bytes).unwrap_or(Value::Null);
(status, value)
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let manifest = Manifest::parse(CHECKER)?;
let seen = Arc::new(Mutex::new(Vec::new()));
let store = Arc::new(RedbStore::open_in_memory()?);
let runtime = Runtime::builder_on(store)
.policy(Arc::new(PermitPeers))
.agent(Agent::new(&manifest).skill(Checks(Arc::clone(&seen))))
.build();
let server = A2aServer::new(
runtime,
Arc::new(BearerAuth),
&CardSecurity::bearer("peer-token", ["a2a:invoke"]),
&manifest,
"https://settlement.example/a2a",
)?;
let router = server.router();
let card = router
.clone()
.oneshot(
Request::builder()
.uri("/.well-known/agent-card.json")
.body(Body::empty())?,
)
.await?;
println!("1. the card is public");
println!(" GET /.well-known/agent-card.json → {}", card.status());
let card: Value =
serde_json::from_slice(&axum::body::to_bytes(card.into_body(), 1 << 20).await?)?;
println!(
" protocolVersion: {} (on the interface — a card may serve several)",
card["supportedInterfaces"][0]["protocolVersion"]
);
println!(" name/version: {} {}", card["name"], card["version"]);
println!(
" skills: {} — derived from `capabilities.provides`, not written",
card["skills"][0]["id"]
);
let send = json!({
"jsonrpc": "2.0", "id": 1, "method": method::SEND_MESSAGE,
"params": { "message": {
"messageId": "m-1", "role": "ROLE_USER",
"parts": [{ "text": "settle GB-4471 for 12,000" }],
}},
});
let (status, body) = rpc(&router, None, send.clone()).await;
println!("\n2. every method is not");
println!(" POST /a2a with no credential → {status}");
println!(" body: {}", body["error"]["message"]);
let (status, body) = rpc(
&router,
Some("peer-alpha"),
json!({"jsonrpc": "2.0", "id": 2, "method": "message/send", "params": {}}),
)
.await;
println!("\n3. the 0.3 method spelling is refused, as an error a client can read");
println!(
" 'message/send' → HTTP {status}, code {}",
body["error"]["code"]
);
println!(" {}", body["error"]["message"]);
let unversioned = router
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/a2a")
.header("authorization", "Bearer peer-alpha")
.header("content-type", "application/json")
.body(Body::from(
json!({"jsonrpc": "2.0", "id": 3, "method": method::GET_TASK, "params": {}})
.to_string(),
))?,
)
.await?;
let unversioned: Value =
serde_json::from_slice(&axum::body::to_bytes(unversioned.into_body(), 1 << 20).await?)?;
println!(
" no A2A-Version header → code {}: {}",
unversioned["error"]["code"], unversioned["error"]["message"]
);
let (_, body) = rpc(&router, Some("peer-alpha"), send).await;
let task = &body["result"]["task"];
println!("\n4. an authenticated SendMessage");
println!(" task id: {}", task["id"]);
println!(" state: {}", task["status"]["state"]);
println!(
" artifact: {}",
task["artifacts"][0]["parts"][0]["data"]
);
let seen = seen.lock().expect("not poisoned");
let first = seen.first().expect("the skill ran");
println!("\n what the skill received:");
println!(" untrusted: {}", first.untrusted);
println!(" provenance: {:?}", first.provenance);
println!(
"\n A peer's message crossed a trust boundary, so it is labelled at the\n \
source. Reaching a mutating sink with it takes a journaled release —\n \
the peer being authenticated is not the same as its content being true."
);
Ok(())
}