use std::sync::Arc;
use agentplane::api::a2a::{A2aReply, A2aServer, Part};
use agentplane::api::{AuthError, Authenticator, Caller};
use agentplane::core::{
AwaitSpec, CorrelationKey, DeadlineSpec, PolicyBundleIdentity, PolicyDecision, PolicyEngine,
PolicyRequest,
};
use agentplane::manifest::Manifest;
use agentplane::peers::CardSecurity;
use agentplane::prelude::*;
use agentplane::runtime::Agent;
use serde_json::{Value, json};
const ECHO: &str = r#"
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata:
name: tck-echo
version: "1.0.0"
spec:
identity:
role: "Answer conformance-kit messages per the TCK's SUT contract."
constraints: "Dispatch on the messageId prefix; echo everything else."
budgets:
max_steps: 5
capabilities:
provides: [tck.echo]
"#;
#[derive(Debug)]
struct TckAgent;
#[async_trait::async_trait]
impl Skill for TckAgent {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("tck.echo")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let message_id = input.peek()["$a2a_message"]["messageId"]
.as_str()
.unwrap_or_default()
.to_owned();
let reply = if message_id.starts_with("tck-artifact-file-url") {
A2aReply::artifact(vec![Part::file_url(
"https://example.com/output.txt",
"text/plain",
"output.txt",
)])
} else if message_id.starts_with("tck-artifact-file") {
A2aReply::artifact(vec![Part::file_raw("dGNr", "text/plain", "output.txt")])
} else if message_id.starts_with("tck-artifact-text") {
A2aReply::artifact(vec![Part::text("Generated text content")])
} else if message_id.starts_with("tck-artifact-data") {
A2aReply::artifact(vec![Part::data(json!({"key": "value", "count": 42}))])
} else if message_id.starts_with("tck-message-response") {
A2aReply::message(vec![Part::text("Direct message response")])
} else if message_id.starts_with("tck-stream-artifact-text") {
A2aReply::artifact(vec![Part::text("Streamed text content")])
} else if message_id.starts_with("tck-stream-artifact-file") {
A2aReply::artifact(vec![Part::file_raw("dGNr", "text/plain", "output.txt")])
} else if message_id.starts_with("tck-stream-ordering-001") {
A2aReply::artifact(vec![Part::text("Ordered output")])
} else if message_id.starts_with("tck-stream-001") {
A2aReply::artifact(vec![Part::text("Stream hello from TCK")])
} else if message_id.starts_with("tck-stream-003") {
A2aReply::artifact(vec![Part::text("Stream task lifecycle")])
} else if message_id.starts_with("tck-complete-task") {
A2aReply::artifact(vec![Part::text("Hello from TCK")])
} else if message_id.starts_with("tck-input-required") {
cx.deadline("reply", &DeadlineSpec::days(1), None).await?;
let reply = cx
.await_event(
&AwaitSpec::new("a2a.task.input", "reply")
.correlate(CorrelationKey::new("task", "tck")),
)
.await?;
cx.meet_deadline("reply").await?;
return Ok(Outcome::done(reply));
} else if message_id.starts_with("tck-reject-task") {
return Ok(Outcome::fail("rejected, as the TCK asked"));
} else {
return Ok(Outcome::done(input.map(
|v| json!({ "echo": v["text"].as_str().unwrap_or_default() }),
)));
};
Ok(Outcome::done(Tainted::trusted(reply.into_value())))
}
}
#[derive(Debug)]
struct AnyoneIsTck;
#[async_trait::async_trait]
impl Authenticator for AnyoneIsTck {
async fn authenticate(&self, headers: &axum::http::HeaderMap) -> Result<Caller, AuthError> {
let actor = headers
.get("authorization")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.strip_prefix("Bearer "))
.unwrap_or("tck");
Ok(Caller::new(actor, vec!["peer".to_owned()]))
}
}
#[derive(Debug)]
struct PermitEverything;
impl PolicyEngine for PermitEverything {
fn authorize(&self, _request: &PolicyRequest<'_>) -> PolicyDecision {
PolicyDecision::Permit
}
fn bundle(&self) -> PolicyBundleIdentity {
PolicyBundleIdentity::new(
agentplane::core::Digest::of(b"tck.permit-everything"),
"tck/permit-everything-v1",
)
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let bind = std::env::var("A2A_TCK_ADDR").unwrap_or_else(|_| "127.0.0.1:9999".to_owned());
let public = std::env::var("A2A_TCK_URL").unwrap_or_else(|_| format!("http://{bind}"));
let manifest = Manifest::parse(ECHO)?;
let store = Arc::new(RedbStore::open_in_memory()?);
let runtime = Runtime::builder_on(Arc::clone(&store))
.policy(Arc::new(PermitEverything))
.agent(Agent::new(&manifest).skill(TckAgent))
.build();
let server = A2aServer::new(
runtime,
Arc::new(AnyoneIsTck),
&CardSecurity::bearer("tck", ["a2a:invoke"]),
&manifest,
format!("{public}/a2a"),
)?
.with_push(
Arc::clone(&store) as Arc<dyn agentplane::push::PushStore>,
Arc::new(
agentplane::push::PushSender::new(
agentplane::push::PushPolicy::new()
.allow_host("localhost")
.allow_host("127.0.0.1"),
)
.allow_plaintext_loopback(),
) as Arc<dyn agentplane::push::PushTransport>,
)?;
if let Some(worker) = server.push_worker() {
tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_millis(250));
loop {
tick.tick().await;
#[allow(clippy::disallowed_methods)]
let at = time::OffsetDateTime::now_utc().unix_timestamp();
let Ok(at) = u64::try_from(at) else { continue };
if let Err(error) = worker.run_once(at, 32).await {
eprintln!("push worker: {error}");
}
}
});
}
let listener = tokio::net::TcpListener::bind(&bind).await?;
eprintln!("a2a-tck fixture serving on http://{bind}");
eprintln!(" card: {public}/.well-known/agent-card.json");
eprintln!(" json-rpc: {public}/a2a");
eprintln!(" push: http://localhost:* (testkit loopback exception)");
axum::serve(listener, server.router()).await?;
Ok(())
}