#![cfg(all(unix, feature = "a2a"))]
mod common;
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::process::{Child, Command, Stdio};
use std::time::{Duration, Instant};
use serde_json::{Value, json};
const TOKEN_A: &str = "authz-token-for-principal-a";
const TOKEN_B: &str = "authz-token-for-principal-b";
fn sigterm(pid: u32) {
unsafe {
libc::kill(pid as i32, libc::SIGTERM);
}
}
fn free_port() -> u16 {
TcpListener::bind("127.0.0.1:0")
.unwrap()
.local_addr()
.unwrap()
.port()
}
fn post(
addr: &str,
body: &str,
extra: &[(&str, &str)],
budget: Duration,
) -> (String, String, String) {
let mut s = TcpStream::connect(addr).expect("connect a2a http");
s.set_read_timeout(Some(Duration::from_millis(200))).ok();
let mut head = format!(
"POST / HTTP/1.1\r\nHost: x\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n",
body.len()
);
for (k, v) in extra {
head.push_str(&format!("{k}: {v}\r\n"));
}
head.push_str("\r\n");
s.write_all(head.as_bytes()).unwrap();
s.write_all(body.as_bytes()).unwrap();
s.flush().unwrap();
let deadline = Instant::now() + budget;
let mut raw = Vec::new();
let mut buf = [0u8; 8192];
while Instant::now() < deadline {
match s.read(&mut buf) {
Ok(0) => break,
Ok(n) => {
raw.extend_from_slice(&buf[..n]);
if raw.windows(6).any(|w| w == b"\ndata:") {
break;
}
}
Err(e)
if matches!(
e.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) => {}
Err(e) => panic!("read a2a response: {e}"),
}
}
let text = String::from_utf8_lossy(&raw).into_owned();
let (head, body) = match text.find("\r\n\r\n") {
Some(i) => (text[..i].to_string(), text[i + 4..].to_string()),
None => (text.clone(), String::new()),
};
let status = head.lines().next().unwrap_or("").to_string();
(status, head, body)
}
fn rpc_as(addr: &str, bearer: &str, id: i64, method: &str, params: Value) -> Value {
let body = json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params}).to_string();
let auth = format!("Bearer {bearer}");
let (status, _, body) = post(
addr,
&body,
&[("Authorization", &auth)],
Duration::from_secs(60),
);
serde_json::from_str(&body)
.unwrap_or_else(|e| panic!("non-JSON A2A response ({e}) to {method}: {status} {body:?}"))
}
fn wait_ready(addr: &str) {
let deadline = Instant::now() + Duration::from_secs(10);
loop {
if TcpStream::connect(addr).is_ok() {
return;
}
assert!(
Instant::now() < deadline,
"a2a listener never became connectable"
);
std::thread::sleep(Duration::from_millis(25));
}
}
struct MockLlm {
child: Child,
addr_file: String,
uri: String,
}
impl Drop for MockLlm {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
let _ = std::fs::remove_file(&self.addr_file);
}
}
fn spawn_mock_llm(playbook: &Value) -> MockLlm {
let pb = common::unique_path("authz-playbook", "json");
std::fs::write(&pb, playbook.to_string()).unwrap();
let addr_file = common::unique_path("authz-mock-llm", "addr");
let _ = std::fs::remove_file(&addr_file);
let child = Command::new(env!("CARGO_BIN_EXE_agentd"))
.args(["--internal-mock-llm", &addr_file, &format!("file:{pb}")])
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.expect("spawn mock llm");
let addr = common::read_addr_file(&addr_file);
MockLlm {
child,
addr_file,
uri: format!("http://{addr}"),
}
}
struct Daemon {
child: Child,
stderr_path: String,
}
impl Daemon {
fn pid(&self) -> u32 {
self.child.id()
}
fn alive(&mut self) -> bool {
matches!(self.child.try_wait(), Ok(None))
}
fn stderr(&self) -> String {
std::fs::read_to_string(&self.stderr_path).unwrap_or_default()
}
}
impl Drop for Daemon {
fn drop(&mut self) {
sigterm(self.child.id());
let deadline = Instant::now() + Duration::from_secs(3);
while Instant::now() < deadline {
if matches!(self.child.try_wait(), Ok(Some(_))) {
break;
}
std::thread::sleep(Duration::from_millis(20));
}
let _ = self.child.kill();
let _ = self.child.wait();
let _ = std::fs::remove_file(&self.stderr_path);
}
}
fn spawn_daemon(config: &str) -> Daemon {
let stderr_path = common::unique_path("a2a-authz-daemon", "log");
let errf = std::fs::File::create(&stderr_path).unwrap();
let child = Command::new(env!("CARGO_BIN_EXE_agentd"))
.args(["--config", config])
.env("AGENTD_AUTHZ_TOKEN_A", TOKEN_A)
.env("AGENTD_AUTHZ_TOKEN_B", TOKEN_B)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::from(errf))
.spawn()
.expect("spawn agentd a2a daemon");
Daemon { child, stderr_path }
}
fn write_config(yaml: &str) -> String {
let path = common::unique_path("a2a-authz", "yaml");
std::fs::write(&path, yaml).unwrap();
path
}
fn two_principal_config(llm: &str, port: u16) -> String {
format!(
"config_version: \"2\"\n\
agent:\n name: a2a-authz\n instruction: You are a helpful test agent.\n preflight: never\n\
intelligence:\n endpoints: {llm}\n model: mock\n\
store:\n kind: memory\n\
a2a:\n listen: http://127.0.0.1:{port}\n\
\x20 principals:\n\
\x20 - match: {{ bearer_ref: \"{{{{secret:AGENTD_AUTHZ_TOKEN_A}}}}\" }}\n\
\x20 role: user\n\
\x20 - match: {{ bearer_ref: \"{{{{secret:AGENTD_AUTHZ_TOKEN_B}}}}\" }}\n\
\x20 role: agent\n\
lifecycle:\n run_until: drained\n\
observability:\n log_level: info\n"
)
}
#[test]
fn one_principals_task_stream_is_not_readable_by_another() {
let llm = spawn_mock_llm(&json!({"turns": [{"content": "the private answer"}]}));
let port = free_port();
let addr = format!("127.0.0.1:{port}");
let cfg = write_config(&two_principal_config(&llm.uri, port));
let mut daemon = spawn_daemon(&cfg);
wait_ready(&addr);
let send = rpc_as(
&addr,
TOKEN_A,
1,
"SendMessage",
json!({"message": {"messageId": "m-a", "parts": [{"text": "hello"}]}}),
);
assert!(
send.get("error").is_none(),
"A's send should succeed: {send}"
);
let task = &send["result"]["task"];
let task_id = task["id"]
.as_str()
.unwrap_or_else(|| panic!("no task id in {send}"))
.to_string();
assert_eq!(
task["status"]["state"], "TASK_STATE_COMPLETED",
"A's task settled: {task}"
);
let got = rpc_as(&addr, TOKEN_B, 2, "GetTask", json!({"id": task_id}));
assert_eq!(got["error"]["code"], -32001, "B's GetTask: {got}");
let cancelled = rpc_as(&addr, TOKEN_B, 3, "CancelTask", json!({"id": task_id}));
assert_eq!(
cancelled["error"]["code"], -32001,
"B's cancel: {cancelled}"
);
let subscribe = json!({"jsonrpc": "2.0", "id": 4, "method": "SubscribeToTask",
"params": {"id": task_id}})
.to_string();
let auth_b = format!("Bearer {TOKEN_B}");
let (status, head, body) = post(
&addr,
&subscribe,
&[("Authorization", &auth_b), ("Last-Event-ID", "0")],
Duration::from_secs(10),
);
assert!(
!body.contains("data:"),
"B received stream frames for A's task: {status} / {body}"
);
assert!(
!body.contains(&task_id) || body.contains("error"),
"B's response carried A's task id outside an error: {body}"
);
assert!(
!body.contains("the private answer"),
"B replayed A's result artifact: {body}"
);
assert!(
head.contains("application/json"),
"a refused subscribe answers with an error, not a stream: {head}"
);
let refused: Value = serde_json::from_str(&body)
.unwrap_or_else(|e| panic!("non-JSON refusal ({e}): {status} {body:?}"));
assert_eq!(
refused["error"]["code"], -32001,
"B's subscribe is refused as not-found: {refused}"
);
let auth_a = format!("Bearer {TOKEN_A}");
let (a_status, _, a_body) = post(
&addr,
&json!({"jsonrpc": "2.0", "id": 5, "method": "SubscribeToTask",
"params": {"id": task_id}})
.to_string(),
&[("Authorization", &auth_a), ("Last-Event-ID", "0")],
Duration::from_secs(10),
);
assert!(
a_body.contains("data:") && a_body.contains(&task_id),
"the owner still receives its own task's events: {a_status} / {a_body}"
);
assert!(daemon.alive(), "daemon still serving: {}", daemon.stderr());
std::fs::remove_file(&cfg).ok();
}
fn rss_kb(pid: u32) -> u64 {
let status = std::fs::read_to_string(format!("/proc/{pid}/status"))
.unwrap_or_else(|e| panic!("read /proc/{pid}/status: {e}"));
status
.lines()
.find_map(|l| l.strip_prefix("VmRSS:"))
.and_then(|v| v.split_whitespace().next())
.and_then(|n| n.parse().ok())
.unwrap_or_else(|| panic!("no VmRSS in /proc/{pid}/status"))
}
fn loopback_config(port: u16) -> String {
format!(
"config_version: \"2\"\n\
agent:\n name: a2a-leak\n instruction: You are a test agent.\n preflight: never\n\
intelligence:\n endpoints: https://127.0.0.1:9\n model: mock\n\
store:\n kind: memory\n\
a2a:\n listen: http://127.0.0.1:{port}\n\
lifecycle:\n run_until: drained\n\
observability:\n log_level: info\n"
)
}
#[test]
fn a_flood_of_distinct_method_names_does_not_grow_the_daemon() {
const NAME: usize = 64_000;
const FLOOD: u64 = 250;
let port = free_port();
let addr = format!("127.0.0.1:{port}");
let cfg = write_config(&loopback_config(port));
let mut daemon = spawn_daemon(&cfg);
wait_ready(&addr);
let call = |n: u64| {
let method = format!("{n:08}{}", "m".repeat(NAME));
let body = json!({"jsonrpc": "2.0", "id": n, "method": method, "params": {}}).to_string();
let (_, _, body) = post(&addr, &body, &[], Duration::from_secs(30));
body
};
for n in 0..40 {
let answer = call(n);
assert!(
answer.contains("error"),
"an unknown method is refused, not served: {answer}"
);
}
let before = rss_kb(daemon.pid());
for n in 40..40 + FLOOD {
call(n);
}
let after = rss_kb(daemon.pid());
let growth = after.saturating_sub(before);
let leaked = FLOOD * NAME as u64 / 1024;
assert!(
growth < leaked / 3,
"RSS grew {growth} KiB over {FLOOD} requests \
(a leaked copy of every method name would be ~{leaked} KiB): \
{before} KiB → {after} KiB"
);
assert!(daemon.alive(), "daemon still serving: {}", daemon.stderr());
std::fs::remove_file(&cfg).ok();
}
#[test]
fn an_unauthenticated_caller_still_reaches_the_admin_check() {
let llm = spawn_mock_llm(&json!({"turns": [{"content": "unused"}]}));
let port = free_port();
let addr = format!("127.0.0.1:{port}");
let cfg = write_config(&two_principal_config(&llm.uri, port));
let mut daemon = spawn_daemon(&cfg);
wait_ready(&addr);
let refused = rpc_as(&addr, "not-a-real-token", 1, "GetTask", json!({"id": "t1"}));
assert_eq!(
refused["error"]["code"], -32003,
"a junk bearer is refused by the matrix, not by a 401: {refused}"
);
let (status, _, _) = post(
&addr,
&json!({"jsonrpc": "2.0", "id": 2, "method": "a2a.drainX", "params": {}}).to_string(),
&[("Authorization", "Bearer also-junk")],
Duration::from_secs(30),
);
assert!(
status.contains("200"),
"the request is dispatched and answered, not rejected at the door: {status}"
);
assert!(daemon.alive(), "daemon still serving: {}", daemon.stderr());
std::fs::remove_file(&cfg).ok();
}