use serde::{Deserialize, Serialize};
use crate::capability_types::AgentCapabilityConfig;
use crate::mcp_server::ScopedMcpServer;
use crate::session::ExecutionSession;
use crate::session_services::SessionStorageStore;
use crate::typed_id::SessionId;
pub const ARD_ATTACHMENT_KV_PREFIX: &str = "ard_attach:";
pub const ARD_DISCOVERY_KV_PREFIX: &str = "ard_disco:";
pub const ARD_ATTACHMENT_RESOURCE_KIND: &str = "ard_attachment";
pub fn attachment_kv_key(slug: &str) -> String {
format!("{ARD_ATTACHMENT_KV_PREFIX}{slug}")
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ArdAttachmentTarget {
McpServer {
name: String,
server: ScopedMcpServer,
},
ExternalAgent { agent: serde_json::Value },
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct ArdAttachment {
pub urn: String,
pub display_name: String,
pub media_type: String,
pub registry_id: String,
pub target: ArdAttachmentTarget,
}
impl ArdAttachment {
pub fn slug(&self) -> String {
urn_slug(&self.urn)
}
}
pub fn urn_slug(urn: &str) -> String {
urn.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
.collect()
}
pub fn merge_attachment_into_session(session: &mut ExecutionSession, attachment: &ArdAttachment) {
match &attachment.target {
ArdAttachmentTarget::McpServer { name, server } => {
session.mcp_servers.insert(name.clone(), server.clone());
}
ArdAttachmentTarget::ExternalAgent { agent } => {
merge_external_agent(session, agent);
}
}
}
fn merge_external_agent(session: &mut ExecutionSession, agent: &serde_json::Value) {
let Some(agent_id) = agent.get("id").and_then(|v| v.as_str()) else {
return;
};
let cap = session
.capabilities
.iter_mut()
.find(|c| c.capability_id() == crate::capabilities::A2A_AGENT_DELEGATION_CAPABILITY_ID);
let cap = match cap {
Some(cap) => cap,
None => {
session
.capabilities
.push(AgentCapabilityConfig::with_config(
crate::capabilities::A2A_AGENT_DELEGATION_CAPABILITY_ID,
serde_json::json!({ "agents": [] }),
));
session
.capabilities
.last_mut()
.expect("just pushed a2a capability config")
}
};
if !cap.config_value().is_object() {
cap.set_config(serde_json::json!({ "agents": [] }));
}
let config = cap.config_mut();
if !config["agents"].is_array() {
config["agents"] = serde_json::Value::Array(Vec::new());
}
let agents = config["agents"]
.as_array_mut()
.expect("agents array ensured above");
if let Some(existing) = agents
.iter_mut()
.find(|a| a.get("id").and_then(|v| v.as_str()) == Some(agent_id))
{
*existing = agent.clone();
} else {
agents.push(agent.clone());
}
}
pub async fn load_session_attachments(
storage: &dyn SessionStorageStore,
session_id: SessionId,
) -> Vec<ArdAttachment> {
let keys = match storage.list_keys(session_id).await {
Ok(keys) => keys,
Err(e) => {
tracing::warn!("failed to list session keys for ARD attachments: {e}");
return Vec::new();
}
};
let mut attachments = Vec::new();
for info in keys {
if !info.key.starts_with(ARD_ATTACHMENT_KV_PREFIX) {
continue;
}
match storage.get_value(session_id, &info.key).await {
Ok(Some(raw)) => match serde_json::from_str::<ArdAttachment>(&raw) {
Ok(attachment) => attachments.push(attachment),
Err(e) => tracing::warn!(key = %info.key, "skipping malformed ARD attachment: {e}"),
},
Ok(None) => {}
Err(e) => tracing::warn!(key = %info.key, "failed to read ARD attachment: {e}"),
}
}
attachments
}
pub async fn apply_session_attachments(
storage: &dyn SessionStorageStore,
session: &mut ExecutionSession,
) {
let attachments = load_session_attachments(storage, session.id).await;
for attachment in &attachments {
merge_attachment_into_session(session, attachment);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::mcp_server::McpServerTransportType;
fn test_session() -> ExecutionSession {
ExecutionSession::with_own_workspace(SessionId::new(), crate::typed_id::HarnessId::new())
}
fn mcp_attachment(urn: &str, name: &str) -> ArdAttachment {
ArdAttachment {
urn: urn.to_string(),
display_name: "Docs MCP".to_string(),
media_type: "application/mcp-server+json".to_string(),
registry_id: "public".to_string(),
target: ArdAttachmentTarget::McpServer {
name: name.to_string(),
server: ScopedMcpServer {
transport_type: McpServerTransportType::Http,
url: "https://docs.example.com/mcp".to_string(),
..Default::default()
},
},
}
}
fn agent_attachment(urn: &str, id: &str) -> ArdAttachment {
ArdAttachment {
urn: urn.to_string(),
display_name: "Concierge".to_string(),
media_type: "application/a2a-agent-card+json".to_string(),
registry_id: "public".to_string(),
target: ArdAttachmentTarget::ExternalAgent {
agent: serde_json::json!({
"id": id,
"name": "Concierge",
"base_url": "https://agent.example.com",
}),
},
}
}
#[test]
fn slug_is_key_safe() {
assert_eq!(
urn_slug("urn:ai:acme.com:agent:assistant"),
"urn_ai_acme_com_agent_assistant"
);
}
#[test]
fn merges_mcp_server_into_session() {
let mut session = test_session();
merge_attachment_into_session(&mut session, &mcp_attachment("urn:ai:x:y:z", "docs"));
assert_eq!(
session.mcp_servers["docs"].url,
"https://docs.example.com/mcp"
);
}
#[test]
fn mcp_merge_is_idempotent_by_name() {
let mut session = test_session();
merge_attachment_into_session(&mut session, &mcp_attachment("urn:ai:x:y:z", "docs"));
let mut updated = mcp_attachment("urn:ai:x:y:z", "docs");
let ArdAttachmentTarget::McpServer { server, .. } = &mut updated.target else {
unreachable!()
};
server.url = "https://updated.example.com/mcp".into();
merge_attachment_into_session(&mut session, &updated);
assert_eq!(session.mcp_servers.len(), 1);
assert_eq!(
session.mcp_servers["docs"].url,
"https://updated.example.com/mcp"
);
}
#[test]
fn merges_external_agent_into_new_capability() {
let mut session = test_session();
merge_attachment_into_session(&mut session, &agent_attachment("urn:ai:x:y:c", "concierge"));
let cap = session
.capabilities
.iter()
.find(|c| c.capability_id() == "a2a_agent_delegation")
.expect("a2a capability added");
let agents = cap.config_value()["agents"].as_array().unwrap();
assert_eq!(agents.len(), 1);
assert_eq!(
agents[0],
serde_json::json!({"id": "concierge", "name": "Concierge", "base_url": "https://agent.example.com"})
);
}
#[test]
fn external_agent_merge_is_idempotent_by_id() {
let mut session = test_session();
merge_attachment_into_session(&mut session, &agent_attachment("urn:ai:x:y:c", "concierge"));
let mut updated = agent_attachment("urn:ai:x:y:c", "concierge");
let ArdAttachmentTarget::ExternalAgent { agent } = &mut updated.target else {
unreachable!()
};
agent["base_url"] = serde_json::json!("https://updated.example.com");
merge_attachment_into_session(&mut session, &updated);
let cap = session
.capabilities
.iter()
.find(|c| c.capability_id() == "a2a_agent_delegation")
.unwrap();
assert_eq!(cap.config_value()["agents"].as_array().unwrap().len(), 1);
assert_eq!(
cap.config_value()["agents"][0]["base_url"],
"https://updated.example.com"
);
assert_eq!(
session
.capabilities
.iter()
.filter(|c| c.capability_id() == "a2a_agent_delegation")
.count(),
1
);
}
#[test]
fn external_agent_appends_to_existing_capability() {
let mut session = test_session();
session
.capabilities
.push(AgentCapabilityConfig::with_config(
"a2a_agent_delegation",
serde_json::json!({ "agents": [ {"id": "preexisting", "name": "Pre"} ], "max_depth": 3 }),
));
merge_attachment_into_session(&mut session, &agent_attachment("urn:ai:x:y:c", "concierge"));
let cap = session
.capabilities
.iter()
.find(|c| c.capability_id() == "a2a_agent_delegation")
.unwrap();
assert_eq!(
cap.config_value(),
&serde_json::json!({
"agents": [
{"id": "preexisting", "name": "Pre"},
{"id": "concierge", "name": "Concierge", "base_url": "https://agent.example.com"}
],
"max_depth": 3
})
);
}
}