use crate::session::{ClientSession, ServerState};
use car_peers::{
DeliveryGuard, DeliveryOutcome, GuardVerdict, PeerAddress, PeerDescriptor, PeerDirectory,
PeerKind, PeerMessage, PeerSource, StaticProvider,
};
use futures::SinkExt;
use serde_json::Value;
use tokio::sync::oneshot;
use tokio_tungstenite::tungstenite::Message;
#[derive(Debug, Clone, serde::Serialize)]
pub struct HeldPeerMessage {
pub message: PeerMessage,
pub target: PeerDescriptor,
pub held_at_ms: u64,
pub reason: String,
}
const PEER_ACK_TIMEOUT_SECS: u64 = 5;
pub async fn snapshot_attached(state: &ServerState) -> Vec<PeerDescriptor> {
let attached = state.attached_agents.lock().await.clone();
attached
.into_keys()
.filter(|id| car_peers::is_valid_peer_name(id))
.map(|agent_id| PeerDescriptor {
name: agent_id.clone(),
reference: None,
kind: PeerKind::CarAgent,
source: PeerSource::Attached,
address: PeerAddress::AttachedAgent { agent_id },
display_name: None,
capability: None,
last_seen_ms: Some(car_peers::now_ms()),
})
.collect()
}
async fn directory_for(state: &ServerState, session: &ClientSession) -> PeerDirectory {
let self_name = session.agent_id.lock().await.clone().unwrap_or_default();
let mut dir = PeerDirectory::new(self_name).with_provider(Box::new(StaticProvider::new(
"attached",
snapshot_attached(state).await,
)));
for (label, peers) in [
("parslee", snapshot_parslee(state).await),
("lan", snapshot_lan(state)),
] {
if !peers.is_empty() {
dir = dir.with_provider(Box::new(StaticProvider::new(label, peers)));
}
}
dir
}
pub async fn snapshot_parslee(state: &ServerState) -> Vec<PeerDescriptor> {
let handle = { state.sync.lock().unwrap_or_else(|e| e.into_inner()).clone() };
let Some(sync) = handle else {
return Vec::new();
};
let endpoints = { sync.lock().await.host_endpoints() };
endpoints
.into_iter()
.filter(|e| car_peers::is_valid_peer_name(&e.name))
.map(|e| PeerDescriptor {
name: e.name,
reference: None,
kind: PeerKind::RemoteCar,
source: PeerSource::Parslee,
address: PeerAddress::A2a { base_url: e.url },
display_name: Some(e.device_id),
capability: None,
last_seen_ms: None,
})
.collect()
}
pub async fn refresh_peer_trust(state: &ServerState) {
let handle = { state.sync.lock().unwrap_or_else(|e| e.into_inner()).clone() };
let Some(sync) = handle else {
state.peer_trust.set_trusted(Vec::<String>::new());
return;
};
let keys: Vec<String> = sync
.lock()
.await
.host_endpoints()
.into_iter()
.map(|e| e.pubkey)
.filter(|k| !k.trim().is_empty())
.collect();
let n = keys.len();
state.peer_trust.set_trusted(keys);
tracing::debug!(trusted_peers = n, "refreshed CAR peer trust set");
}
pub fn snapshot_lan(state: &ServerState) -> Vec<PeerDescriptor> {
let guard = state
.lan_discovery
.lock()
.unwrap_or_else(|e| e.into_inner());
let Some(dir) = guard.as_ref() else {
return Vec::new();
};
let trusted: std::collections::HashSet<String> = car_a2a::peers::PeerRegistry::user_default()
.map(|r| r.list().into_iter().map(|p| p.url).collect())
.unwrap_or_default();
dir.peers()
.into_iter()
.filter(|p| car_peers::is_valid_peer_name(&p.name))
.filter(|p| !trusted.contains(&p.url))
.map(|p| PeerDescriptor {
name: p.name,
reference: None,
kind: PeerKind::RemoteCar,
source: PeerSource::Lan,
address: PeerAddress::A2a { base_url: p.url },
display_name: None,
capability: None,
last_seen_ms: None,
})
.collect()
}
pub async fn handle_agents_peers(
state: &ServerState,
session: &ClientSession,
) -> Result<Value, String> {
let dir = directory_for(state, session).await;
let peers = dir.list();
let discoverable = state
.lan_discovery
.lock()
.unwrap_or_else(|e| e.into_inner())
.is_some();
Ok(serde_json::json!({
"self": if dir.self_name().is_empty() { Value::Null } else { Value::from(dir.self_name()) },
"lan_browsing": discoverable,
"peers": peers.iter().map(|p| serde_json::json!({
"name": p.name,
"address": p.address_form(),
"reference": p.reference,
"kind": p.kind.as_str(),
"source": p.source.as_str(),
"can_receive": p.kind.can_receive(),
"display_name": p.display_name,
"capability": p.capability,
"last_seen_ms": p.last_seen_ms,
})).collect::<Vec<_>>(),
"count": peers.len(),
}))
}
pub async fn handle_agents_message(
req: &crate::handler::JsonRpcMessage,
state: &ServerState,
session: &ClientSession,
) -> Result<Value, String> {
let to = req
.params
.get("to")
.and_then(|v| v.as_str())
.ok_or("missing `to`")?
.to_string();
let body = req
.params
.get("body")
.and_then(|v| v.as_str())
.ok_or("missing `body`")?
.to_string();
let summary = req
.params
.get("summary")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
let from = crate::handler::session_principal_for_peers(session).await;
let dir = directory_for(state, session).await;
let target = dir.resolve(&to).map_err(|e| e.to_string())?;
if !target.source.is_trusted_by_default() {
return Err(format!(
"`{}` was discovered on the local network and is not a trusted peer. Anyone on \
this network can advertise any name, so discovery makes a peer visible, not \
reachable. Promote it with `a2a.peers.add` first.",
target.name
));
}
if !target.kind.can_receive() {
return Err(format!(
"`{}` is a {} — it can message CAR while it runs but has no inbox to deliver into",
target.name,
target.kind.as_str()
));
}
let mut msg = PeerMessage::new(&from, &target.name, &body);
msg.summary = summary;
let verdict = {
let mut guards = state.peer_guards.lock().await;
let guard = guards
.entry(target.name.clone())
.or_insert_with(DeliveryGuard::new);
guard.admit(&msg, car_peers::now_ms())
};
if !verdict.is_accept() {
let outcome = DeliveryOutcome::Refused {
reason: verdict.reason(),
};
append_peer_audit(&msg, &target, &outcome);
return Err(guard_error(&verdict));
}
let outcome = admit(state, session, &msg, &target).await;
append_peer_audit(&msg, &target, &outcome);
match &outcome {
DeliveryOutcome::Refused { reason } => {
release(state, &target.name).await;
return Err(format!("message refused: {reason}"));
}
DeliveryOutcome::Held { reason } => {
release(state, &target.name).await;
let held = HeldPeerMessage {
message: msg.clone(),
target: target.clone(),
held_at_ms: car_peers::now_ms(),
reason: reason.clone(),
};
let dropped = {
let mut q = state.held_peer_messages.lock().await;
q.push_back(held);
if q.len() > car_peers::HOLD_CAP {
q.pop_front()
} else {
None
}
};
if let Some(evicted) = dropped {
tracing::warn!(
id = %evicted.message.id,
from = %evicted.message.from,
to = %evicted.target.name,
"hold queue full; dropped the oldest undecided peer message"
);
append_peer_audit(
&evicted.message,
&evicted.target,
&DeliveryOutcome::Refused {
reason: format!(
"evicted from the hold queue at {} undecided messages",
car_peers::HOLD_CAP
),
},
);
}
return Ok(serde_json::json!({
"id": msg.id,
"to": target.name,
"outcome": "held",
"retained": true,
"reason": reason,
}));
}
DeliveryOutcome::Delivered => {}
}
let result = deliver(state, &target, &msg).await;
release(state, &target.name).await;
result
}
pub async fn handle_agents_message_pending(
state: &ServerState,
session: &ClientSession,
) -> Result<Value, String> {
require_host(session, "agents.message.pending")?;
Ok(pending_snapshot(state).await)
}
pub async fn pending_snapshot(state: &ServerState) -> Value {
let q = state.held_peer_messages.lock().await;
serde_json::json!({
"held": q.iter().map(|h| serde_json::json!({
"id": h.message.id,
"from": h.message.from,
"to": h.target.name,
"body": h.message.body,
"held_at_ms": h.held_at_ms,
"reason": h.reason,
})).collect::<Vec<_>>(),
"count": q.len(),
"cap": car_peers::HOLD_CAP,
})
}
pub async fn handle_agents_message_approve(
req: &crate::handler::JsonRpcMessage,
state: &ServerState,
session: &ClientSession,
) -> Result<Value, String> {
require_host(session, "agents.message.approve")?;
let id = req
.params
.get("id")
.and_then(|v| v.as_str())
.ok_or("missing `id`")?
.to_string();
let approved = match req.params.get("decision") {
Some(Value::Bool(b)) => *b,
Some(Value::String(sv)) => {
matches!(
sv.to_ascii_lowercase().as_str(),
"approve" | "approved" | "yes"
)
}
_ => false,
};
decide_held(state, &id, approved).await
}
pub async fn decide_held(state: &ServerState, id: &str, approved: bool) -> Result<Value, String> {
let held = {
let mut q = state.held_peer_messages.lock().await;
let pos = q.iter().position(|h| h.message.id == id);
match pos {
Some(i) => q.remove(i).expect("position just found"),
None => return Err(format!("no held message with id `{id}`")),
}
};
if !approved {
append_peer_audit(
&held.message,
&held.target,
&DeliveryOutcome::Refused {
reason: "denied by the operator".into(),
},
);
return Ok(serde_json::json!({
"id": id,
"outcome": "denied",
}));
}
let verdict = {
let mut guards = state.peer_guards.lock().await;
guards
.entry(held.target.name.clone())
.or_insert_with(DeliveryGuard::new)
.admit(&held.message, car_peers::now_ms())
};
if !verdict.is_accept() {
append_peer_audit(
&held.message,
&held.target,
&DeliveryOutcome::Refused {
reason: verdict.reason(),
},
);
return Err(guard_error(&verdict));
}
append_peer_audit(&held.message, &held.target, &DeliveryOutcome::Delivered);
let result = deliver(state, &held.target, &held.message).await;
release(state, &held.target.name).await;
result
}
fn require_host(session: &ClientSession, method: &str) -> Result<(), String> {
if session.is_host.load(std::sync::atomic::Ordering::Acquire) {
return Ok(());
}
Err(require_host_message(method))
}
fn require_host_message(method: &str) -> String {
format!("`{method}` is host-only; an agent cannot approve the messages its own posture held")
}
async fn release(state: &ServerState, recipient: &str) {
if let Some(g) = state.peer_guards.lock().await.get_mut(recipient) {
g.consumed();
}
}
fn guard_error(v: &GuardVerdict) -> String {
match v {
GuardVerdict::Accept => "accepted".into(),
GuardVerdict::TooLarge { .. } => {
format!("{} — send a path or a state handle instead", v.reason())
}
GuardVerdict::RateLimited { .. } => {
format!("{} — batch the rest into one message", v.reason())
}
GuardVerdict::DuplicateWithinWindow => {
format!("{} — it was already delivered; do not resend", v.reason())
}
GuardVerdict::QueueFull { .. } => {
format!("{} — wait for it to drain", v.reason())
}
GuardVerdict::InvalidName { .. } => v.reason(),
}
}
async fn admit(
_state: &ServerState,
session: &ClientSession,
msg: &PeerMessage,
_target: &PeerDescriptor,
) -> DeliveryOutcome {
let sender_agent = session.agent_id.lock().await.clone();
let is_host = session.is_host.load(std::sync::atomic::Ordering::Acquire);
admit_with(
&crate::agent_permissions::load_policy(),
sender_agent,
is_host,
&msg.from,
)
}
fn admit_with(
policy: &car_policy::AgentPermissionPolicy,
sender_agent: Option<String>,
is_host: bool,
from: &str,
) -> DeliveryOutcome {
let Some(agent_id) = sender_agent else {
if is_host {
return DeliveryOutcome::Delivered;
}
return DeliveryOutcome::Refused {
reason: format!(
"sender `{from}` is neither a bound agent nor the host; a peer message needs an authenticated principal"
),
};
};
match policy.resolve(&agent_id, car_policy::PermissionTier::ReadOnly) {
car_policy::agent_permissions::ApprovalMode::AlwaysAllow => DeliveryOutcome::Delivered,
car_policy::agent_permissions::ApprovalMode::RequireApproval => DeliveryOutcome::Held {
reason: format!(
"`{agent_id}` is set to require approval; the message is held rather than dropped"
),
},
car_policy::agent_permissions::ApprovalMode::Deny => DeliveryOutcome::Refused {
reason: format!("`{agent_id}` is denied at the read_only tier"),
},
}
}
async fn deliver(
state: &ServerState,
target: &PeerDescriptor,
msg: &PeerMessage,
) -> Result<Value, String> {
let agent_id = match &target.address {
PeerAddress::AttachedAgent { agent_id } => agent_id,
PeerAddress::A2a { base_url } => return deliver_remote(state, base_url, target, msg).await,
};
let agent_client_id = state
.attached_agents
.lock()
.await
.get(agent_id)
.cloned()
.ok_or_else(|| format!("agent `{agent_id}` detached before the message could be sent"))?;
let channel = {
let sessions = state.sessions.lock().await;
sessions
.get(&agent_client_id)
.map(|s| s.channel.clone())
.ok_or_else(|| format!("agent `{agent_id}` raced with disconnect"))?
};
let request_id = channel.next_request_id();
let (tx, rx) = oneshot::channel();
channel.pending.lock().await.insert(request_id.clone(), tx);
let rpc = serde_json::json!({
"jsonrpc": "2.0",
"method": "agent.peer_message",
"params": {
"id": msg.id,
"from": msg.from,
"body": msg.body,
"sent_at_ms": msg.sent_at_ms,
"no_reply": msg.no_reply,
},
"id": request_id,
});
let frame = Message::Text(
serde_json::to_string(&rpc)
.map_err(|e| e.to_string())?
.into(),
);
if let Err(e) = channel.write.lock().await.send(frame).await {
channel.pending.lock().await.remove(&request_id);
return Err(format!("failed to deliver to `{agent_id}`: {e}"));
}
match tokio::time::timeout(std::time::Duration::from_secs(PEER_ACK_TIMEOUT_SECS), rx).await {
Ok(Ok(_)) => Ok(serde_json::json!({
"id": msg.id,
"to": target.name,
"outcome": "delivered",
})),
Ok(Err(_)) => Err(format!("agent `{agent_id}` closed before acknowledging")),
Err(_) => {
channel.pending.lock().await.remove(&request_id);
Ok(serde_json::json!({
"id": msg.id,
"to": target.name,
"outcome": "unacknowledged",
"detail": format!(
"written to `{agent_id}` but not acknowledged within {PEER_ACK_TIMEOUT_SECS}s"
),
}))
}
}
}
async fn deliver_remote(
state: &ServerState,
base_url: &str,
target: &PeerDescriptor,
msg: &PeerMessage,
) -> Result<Value, String> {
use car_a2a::types::{Message as A2aMessage, MessageRole, Part, TextPart};
let identity = {
state
.peer_identity
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
};
let Some(identity) = identity else {
return Err(format!(
"cannot reach `{}`: this daemon has no peer identity, so a remote CAR would \
refuse it. The identity is created when the A2A surface starts.",
target.name
));
};
let client = car_a2a::client::A2aClient::new(base_url).with_peer_identity(identity);
let mut metadata = std::collections::HashMap::new();
metadata.insert("carPeerFrom".to_string(), Value::from(msg.from.clone()));
metadata.insert("carPeerTo".to_string(), Value::from(target.name.clone()));
let a2a_msg = A2aMessage {
message_id: msg.id.clone(),
role: MessageRole::User,
parts: vec![Part::Text(TextPart {
text: msg.body.clone(),
metadata: std::collections::HashMap::new(),
})],
task_id: None,
context_id: None,
metadata,
};
match client.send_message(a2a_msg, true).await {
Ok(_) => Ok(serde_json::json!({
"id": msg.id,
"to": target.name,
"outcome": "delivered",
"transport": "a2a",
"url": base_url,
})),
Err(e) => Err(format!(
"failed to deliver to `{}` at {base_url}: {e}",
target.name
)),
}
}
pub fn append_peer_audit(msg: &PeerMessage, target: &PeerDescriptor, outcome: &DeliveryOutcome) {
let Some(car_dir) = car_home::root() else {
return;
};
if std::fs::create_dir_all(&car_dir).is_err() {
return;
}
append_peer_audit_at(&car_dir.join("peer-messages.jsonl"), msg, target, outcome);
}
pub fn append_peer_audit_at(
path: &std::path::Path,
msg: &PeerMessage,
target: &PeerDescriptor,
outcome: &DeliveryOutcome,
) {
use std::io::Write;
let record = serde_json::json!({
"ts": chrono::Utc::now().to_rfc3339(),
"id": msg.id,
"from": msg.from,
"to": target.name,
"kind": target.kind.as_str(),
"source": target.source.as_str(),
"bytes": msg.body.len(),
"outcome": outcome,
});
let Ok(line) = serde_json::to_string(&record) else {
return;
};
if let Ok(mut f) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(path)
{
let _ = writeln!(f, "{line}");
} else {
tracing::warn!(path = %path.display(), "failed to append peer-message audit record");
}
}
pub fn append_agent_chat_audit(principal: &str, agent_id: &str, session_id: &str) {
use std::io::Write;
let Some(car_dir) = car_home::root() else {
return;
};
if std::fs::create_dir_all(&car_dir).is_err() {
return;
}
let path = car_dir.join("peer-messages.jsonl");
let record = serde_json::json!({
"ts": chrono::Utc::now().to_rfc3339(),
"surface": "agents.chat",
"from": principal,
"to": agent_id,
"session_id": session_id,
});
let Ok(line) = serde_json::to_string(&record) else {
return;
};
if let Ok(mut f) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)
{
let _ = writeln!(f, "{line}");
}
}
pub fn register_peer_tools(
server: &mut car_mcp::Server,
state: std::sync::Arc<ServerState>,
) -> Result<(), car_mcp::RegisterError> {
server.register_tool(
peer_list_schema(),
std::sync::Arc::new(PeerListTool(state.clone())),
)?;
server.register_tool(
peer_message_schema(),
std::sync::Arc::new(PeerMessageTool(state)),
)?;
Ok(())
}
const MCP_PRINCIPAL: &str = "mcp:external-cli";
fn peer_list_schema() -> Value {
serde_json::json!({
"name": "peer_list",
"description": "List the CAR agents you can send a message to. Returns each peer's \
address, kind, and whether it can receive. Use the returned `address` \
verbatim as peer_message's `to` — it carries a disambiguating suffix \
when two live agents share a name.",
"inputSchema": { "type": "object", "properties": {} },
"annotations": {
"readOnlyHint": true,
"destructiveHint": false,
"idempotentHint": true,
"openWorldHint": false,
},
})
}
fn peer_message_schema() -> Value {
serde_json::json!({
"name": "peer_message",
"description": "Send a short plain-text message to one CAR agent — a finding, a status, \
a decision it is blocked on. The message is text only: it cannot run a \
command, approve anything, or change the recipient's configuration, and \
whatever the recipient does about it goes through its own permissions. \
Get `to` from peer_list. Keep it to one self-contained first line; \
identical repeats within 10s are dropped.",
"inputSchema": {
"type": "object",
"properties": {
"to": { "type": "string", "description": "An `address` from peer_list." },
"body": { "type": "string", "description": "Plain text. First line should stand alone." },
},
"required": ["to", "body"],
},
"annotations": {
"readOnlyHint": false,
"destructiveHint": false,
"idempotentHint": false,
"openWorldHint": true,
},
})
}
struct PeerListTool(std::sync::Arc<ServerState>);
#[async_trait::async_trait]
impl car_mcp::ToolHandler for PeerListTool {
async fn call(&self, _args: Value) -> Result<String, car_mcp::ToolError> {
let peers = snapshot_attached(&self.0).await;
let rows: Vec<Value> = peers
.iter()
.map(|p| {
serde_json::json!({
"address": p.address_form(),
"kind": p.kind.as_str(),
"can_receive": p.kind.can_receive(),
})
})
.collect();
serde_json::to_string(&serde_json::json!({ "peers": rows, "count": rows.len() }))
.map_err(|e| car_mcp::ToolError::Internal(e.to_string()))
}
}
struct PeerMessageTool(std::sync::Arc<ServerState>);
#[async_trait::async_trait]
impl car_mcp::ToolHandler for PeerMessageTool {
async fn call(&self, args: Value) -> Result<String, car_mcp::ToolError> {
let to = args
.get("to")
.and_then(|v| v.as_str())
.ok_or_else(|| car_mcp::ToolError::InvalidParams("missing `to`".into()))?;
let body = args
.get("body")
.and_then(|v| v.as_str())
.ok_or_else(|| car_mcp::ToolError::InvalidParams("missing `body`".into()))?;
let dir = PeerDirectory::new(MCP_PRINCIPAL).with_provider(Box::new(StaticProvider::new(
"attached",
snapshot_attached(&self.0).await,
)));
let target = dir
.resolve(to)
.map_err(|e| car_mcp::ToolError::Internal(e.to_string()))?;
if !target.kind.can_receive() {
return Err(car_mcp::ToolError::Internal(format!(
"`{}` has no inbox to deliver into",
target.name
)));
}
let msg = PeerMessage::new(MCP_PRINCIPAL, &target.name, body);
let verdict = {
let mut guards = self.0.peer_guards.lock().await;
let guard = guards
.entry(target.name.clone())
.or_insert_with(DeliveryGuard::new);
guard.admit(&msg, car_peers::now_ms())
};
if !verdict.is_accept() {
append_peer_audit(
&msg,
&target,
&DeliveryOutcome::Refused {
reason: verdict.reason(),
},
);
return Err(car_mcp::ToolError::Internal(guard_error(&verdict)));
}
append_peer_audit(&msg, &target, &DeliveryOutcome::Delivered);
let result = deliver(&self.0, &target, &msg).await;
release(&self.0, &target.name).await;
match result {
Ok(v) => {
serde_json::to_string(&v).map_err(|e| car_mcp::ToolError::Internal(e.to_string()))
}
Err(e) => Err(car_mcp::ToolError::Internal(e)),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
async fn test_state() -> (Arc<ServerState>, tempfile::TempDir) {
let temp = tempfile::tempdir().unwrap();
let state = Arc::new(ServerState::with_config(
crate::session::ServerStateConfig::new(temp.path().to_path_buf()),
));
(state, temp)
}
async fn attach(state: &ServerState, agent_id: &str) {
state
.attached_agents
.lock()
.await
.insert(agent_id.to_string(), format!("client-{agent_id}"));
}
fn held_fixture(id: &str, to: &str) -> HeldPeerMessage {
let mut m = PeerMessage::new("agent:sender", to, format!("body-{id}"));
m.id = id.to_string();
HeldPeerMessage {
message: m,
target: PeerDescriptor {
name: to.into(),
reference: None,
kind: PeerKind::CarAgent,
source: PeerSource::Attached,
address: PeerAddress::AttachedAgent {
agent_id: to.into(),
},
display_name: None,
capability: None,
last_seen_ms: None,
},
held_at_ms: 1_000,
reason: "requires approval".into(),
}
}
#[tokio::test]
async fn pending_lists_held_messages_oldest_first() {
let (state, _t) = test_state().await;
{
let mut q = state.held_peer_messages.lock().await;
q.push_back(held_fixture("first", "milo"));
q.push_back(held_fixture("second", "milo"));
}
let snap = pending_snapshot(&state).await;
assert_eq!(snap["count"], 2);
assert_eq!(snap["cap"], car_peers::HOLD_CAP);
assert_eq!(snap["held"][0]["id"], "first");
assert_eq!(snap["held"][1]["id"], "second");
}
#[tokio::test]
async fn denying_a_held_message_removes_it() {
let (state, _t) = test_state().await;
state
.held_peer_messages
.lock()
.await
.push_back(held_fixture("m1", "milo"));
let out = decide_held(&state, "m1", false).await.unwrap();
assert_eq!(out["outcome"], "denied");
assert_eq!(
pending_snapshot(&state).await["count"],
0,
"a decided message must leave the queue"
);
}
#[tokio::test]
async fn deciding_an_unknown_id_is_a_named_error() {
let (state, _t) = test_state().await;
let err = decide_held(&state, "ghost", true).await.unwrap_err();
assert!(err.contains("ghost"), "error should name the id: {err}");
}
#[tokio::test]
async fn a_held_message_cannot_be_decided_twice() {
let (state, _t) = test_state().await;
state
.held_peer_messages
.lock()
.await
.push_back(held_fixture("m1", "milo"));
assert!(decide_held(&state, "m1", false).await.is_ok());
assert!(decide_held(&state, "m1", true).await.is_err());
}
#[tokio::test]
async fn approving_a_detached_recipient_fails_loudly() {
let (state, _t) = test_state().await;
state
.held_peer_messages
.lock()
.await
.push_back(held_fixture("m1", "ghost"));
let err = decide_held(&state, "m1", true).await.unwrap_err();
assert!(
err.contains("ghost"),
"approval of a vanished recipient must name it: {err}"
);
assert_eq!(pending_snapshot(&state).await["count"], 0);
}
#[test]
fn only_the_host_may_read_or_decide_the_hold_queue() {
let msg = require_host_message("agents.message.approve");
assert!(msg.contains("host-only"), "{msg}");
assert!(msg.contains("its own posture held"), "{msg}");
}
#[test]
fn an_unauthenticated_sender_is_refused() {
let policy = car_policy::AgentPermissionPolicy::default();
let out = admit_with(&policy, None, false, "conn:abc");
assert!(
matches!(out, DeliveryOutcome::Refused { .. }),
"got {out:?}"
);
}
#[test]
fn the_host_needs_no_agent_posture() {
let policy = car_policy::AgentPermissionPolicy::default();
assert_eq!(
admit_with(&policy, None, true, "conn:host"),
DeliveryOutcome::Delivered
);
}
#[test]
fn a_bound_agent_is_allowed_by_default() {
let policy = car_policy::AgentPermissionPolicy::default();
assert_eq!(
admit_with(&policy, Some("milo".into()), false, "agent:milo"),
DeliveryOutcome::Delivered
);
}
#[test]
fn denying_an_agent_at_read_only_actually_stops_its_messages() {
let mut policy = car_policy::AgentPermissionPolicy::default();
policy.set_agent(
"milo",
car_policy::PermissionTier::ReadOnly,
car_policy::agent_permissions::ApprovalMode::Deny,
);
let out = admit_with(&policy, Some("milo".into()), false, "agent:milo");
assert!(
matches!(out, DeliveryOutcome::Refused { .. }),
"got {out:?}"
);
assert_eq!(
admit_with(&policy, Some("trader".into()), false, "agent:trader"),
DeliveryOutcome::Delivered
);
}
#[test]
fn require_approval_holds_rather_than_drops() {
let mut policy = car_policy::AgentPermissionPolicy::default();
policy.set_agent(
"milo",
car_policy::PermissionTier::ReadOnly,
car_policy::agent_permissions::ApprovalMode::RequireApproval,
);
assert!(matches!(
admit_with(&policy, Some("milo".into()), false, "agent:milo"),
DeliveryOutcome::Held { .. }
));
}
#[tokio::test]
async fn snapshot_lists_attached_agents() {
let (state, _t) = test_state().await;
attach(&state, "milo").await;
attach(&state, "trader").await;
let peers = snapshot_attached(&state).await;
assert_eq!(peers.len(), 2);
assert!(peers.iter().all(|p| p.kind == PeerKind::CarAgent));
assert!(peers.iter().all(|p| p.source == PeerSource::Attached));
}
#[tokio::test]
async fn snapshot_drops_names_that_are_not_addressable() {
let (state, _t) = test_state().await;
attach(&state, "milo").await;
attach(&state, "../escape").await;
let peers = snapshot_attached(&state).await;
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].name, "milo");
}
#[tokio::test]
async fn an_oversized_message_is_refused_before_delivery() {
let (state, _t) = test_state().await;
attach(&state, "milo").await;
let msg = PeerMessage::new("agent:sender", "milo", "x".repeat(2_000_000));
let verdict = {
let mut guards = state.peer_guards.lock().await;
guards
.entry("milo".to_string())
.or_insert_with(DeliveryGuard::new)
.admit(&msg, car_peers::now_ms())
};
assert!(matches!(verdict, GuardVerdict::TooLarge { .. }));
assert!(guard_error(&verdict).contains("state handle"));
}
#[tokio::test]
async fn a_detached_agent_yields_a_structured_error_not_a_hang() {
let (state, _t) = test_state().await;
attach(&state, "ghost").await;
let target = snapshot_attached(&state).await.remove(0);
let msg = PeerMessage::new("agent:sender", "ghost", "hello");
let err = deliver(&state, &target, &msg).await.unwrap_err();
assert!(
err.contains("ghost") && err.contains("disconnect"),
"error should name the agent and the cause, got: {err}"
);
}
#[tokio::test]
async fn guards_are_per_recipient_not_global() {
let (state, _t) = test_state().await;
attach(&state, "a").await;
attach(&state, "b").await;
let mut guards = state.peer_guards.lock().await;
let dup = PeerMessage::new("agent:s", "a", "same body");
assert!(guards
.entry("a".into())
.or_insert_with(DeliveryGuard::new)
.admit(&dup, 1_000)
.is_accept());
let to_b = PeerMessage::new("agent:s", "b", "same body");
assert!(guards
.entry("b".into())
.or_insert_with(DeliveryGuard::new)
.admit(&to_b, 1_000)
.is_accept());
}
#[tokio::test]
async fn a_refused_message_still_leaves_an_audit_record() {
let temp = tempfile::tempdir().unwrap();
let journal = temp.path().join("peer-messages.jsonl");
let target = PeerDescriptor {
name: "milo".into(),
reference: None,
kind: PeerKind::CarAgent,
source: PeerSource::Attached,
address: PeerAddress::AttachedAgent {
agent_id: "milo".into(),
},
display_name: None,
capability: None,
last_seen_ms: None,
};
let msg = PeerMessage::new("agent:sender", "milo", "hello");
append_peer_audit_at(
&journal,
&msg,
&target,
&DeliveryOutcome::Refused {
reason: "over the rate budget".into(),
},
);
let body = std::fs::read_to_string(&journal).expect("journal written");
let rec: Value = serde_json::from_str(body.trim()).expect("one json line");
assert_eq!(rec["from"], "agent:sender");
assert_eq!(rec["to"], "milo");
assert_eq!(rec["outcome"]["outcome"], "refused");
assert_eq!(rec["outcome"]["reason"], "over the rate budget");
}
}