use crate::agentloop::stop::Outcome;
use crate::config::{A2aPeerSpec, McpServerSpec, SwapPolicy};
use crate::wire::intel::Usage;
use serde::{Deserialize, Serialize};
use serde_json::Value;
pub const SUBAGENT_ENV: &str = "AGENT_SUBAGENT";
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum ControlMsg {
Spawn(Box<SpawnPayload>),
Ping { seq: u64 },
Pause,
Resume,
Cancel { reason: String },
Inject { message: String },
SwapIntel(Box<SwapIntel>),
ToolResult {
id: u64,
result: Value,
#[serde(default)]
is_error: bool,
},
BudgetGrant {
id: u64,
ok: bool,
#[serde(default)]
wait_ms: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
model: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
reason: Option<String>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SwapIntel {
pub uri: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub token: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
#[serde(default)]
pub policy: SwapPolicy,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum AgentMsg {
Ready,
Pong { seq: u64 },
Event { event: String, fields: Value },
Usage(Usage),
Turn { outcome: Outcome },
Result { outcome: Outcome },
Failed { error: String },
Gate { node: String, payload: Value },
GateClosed { node: String, via: String },
IntelHealth {
all_down: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
active: Option<IntelActive>,
},
ToolRequest { id: u64, name: String, args: Value },
BudgetRequest { id: u64, estimate: u64 },
TurnDone { turn: Box<TurnResult> },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct IntelActive {
pub index: usize,
pub transport: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SpawnPayload {
pub instruction: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_contract: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub context_seed: Vec<SeedMessage>,
pub intelligence: IntelConfig,
#[serde(default)]
pub mcp_servers: Vec<McpServerSpec>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub a2a_peers: Vec<A2aPeerSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tls_ca: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub aauth: Option<crate::config::AAuthSettings>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub gated_tools: Vec<String>,
pub limits: Limits,
pub telemetry: Telemetry,
pub depth: u32,
#[serde(default)]
pub warm: bool,
#[serde(default)]
pub role: Role,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub turn: Option<Box<TurnSpec>>,
}
pub const ALLOWED_TOOLS_ROLE: &str = "agentd/allowed-tools";
pub fn parse_allowed_tools(content: &str) -> Vec<String> {
serde_json::from_str::<Vec<String>>(content).unwrap_or_default()
}
impl SpawnPayload {
pub fn narrow_tools(&mut self, allow: &[String]) {
self.context_seed.retain(|m| m.role != ALLOWED_TOOLS_ROLE);
self.context_seed.insert(
0,
SeedMessage {
role: ALLOWED_TOOLS_ROLE.to_string(),
content: serde_json::to_string(allow).unwrap_or_else(|_| "[]".to_string()),
},
);
}
pub fn allowed_tools(&self) -> Option<Vec<String>> {
self.context_seed
.iter()
.find(|m| m.role == ALLOWED_TOOLS_ROLE)
.map(|m| parse_allowed_tools(&m.content))
}
}
#[cfg(feature = "workflow")]
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct WorkflowResumeRef {
pub server: String,
pub key: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub seq: Option<u64>,
#[serde(default)]
pub force: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum Role {
#[default]
Agent,
Turn,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum TurnKind {
#[default]
Turn,
Think,
Agent,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct TurnSpec {
#[serde(default)]
pub kind: TurnKind,
pub system: String,
#[serde(default)]
pub messages: Vec<crate::context::Msg>,
#[serde(default)]
pub tools: Vec<crate::wire::intel::ToolDef>,
#[serde(default)]
pub internal: Vec<String>,
#[serde(default)]
pub mcp_routes: std::collections::BTreeMap<String, (String, String)>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_schema: Option<Value>,
#[serde(default)]
pub max_rounds: u32,
#[serde(default)]
pub budget_admission: bool,
#[serde(default)]
pub idempotency_prefix: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tool_meta: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub temperature: Option<f32>,
#[serde(default)]
pub max_tokens_per_call: u32,
#[serde(default)]
pub turn_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct TurnResult {
pub status: String,
#[serde(default)]
pub messages: Vec<crate::context::Msg>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub text: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub value: Option<Value>,
#[serde(default)]
pub usage: Usage,
#[serde(default)]
pub rounds: u32,
#[serde(default)]
pub tool_calls: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finish: Option<Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SeedMessage {
pub role: String,
pub content: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct IntelConfig {
pub uri: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub token: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub model: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub headers: Vec<(String, String)>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub aws_auth: Option<crate::config::AuthSpec>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub dialect: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Limits {
pub max_steps: u32,
pub max_tokens: u64,
pub deadline_ms: u64,
pub max_depth: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub memory_bytes: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cpu_seconds: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub nice: Option<i32>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Telemetry {
pub run_id: String,
pub agent_id: String,
pub agent_path: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub trace_id: Option<String>,
pub log_level: String,
#[serde(default)]
pub log_content: bool,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::agentloop::stop::TerminalStatus;
use crate::json::frame;
use serde_json::json;
use std::io::Cursor;
fn payload() -> SpawnPayload {
SpawnPayload {
instruction: "summarize the file".into(),
output_contract: Some("Return a 3-bullet summary.".into()),
context_seed: vec![SeedMessage {
role: "user".into(),
content: "prior note".into(),
}],
gated_tools: Vec::new(),
intelligence: IntelConfig {
uri: "https://intel.example".into(),
token: Some("secret".into()),
model: Some("m".into()),
headers: Vec::new(),
aws_auth: None,
dialect: None,
},
mcp_servers: vec![McpServerSpec {
name: "fs".into(),
endpoint: "unix:/mcp-fs.sock".into(),
tags: Vec::new(),
..Default::default()
}],
a2a_peers: Vec::new(),
tls_ca: None,
aauth: None,
limits: Limits {
max_steps: 20,
max_tokens: 100_000,
deadline_ms: 600_000,
max_depth: 4,
memory_bytes: None,
cpu_seconds: None,
nice: None,
},
telemetry: Telemetry {
run_id: "r1".into(),
agent_id: "0.1".into(),
agent_path: "0.1".into(),
trace_id: None,
log_level: "info".into(),
log_content: false,
},
depth: 1,
warm: false,
role: crate::subagent::protocol::Role::Agent,
turn: None,
}
}
#[test]
fn control_spawn_frames_roundtrip() {
let mut p = payload();
p.instruction = "line1\nline2".into();
let msg = ControlMsg::Spawn(Box::new(p));
let mut buf = Vec::new();
frame::write_frame(&mut buf, &msg).unwrap();
let mut cur = Cursor::new(buf);
let bytes = frame::read_frame(&mut cur).unwrap().unwrap();
let back: ControlMsg = serde_json::from_slice(&bytes).unwrap();
match back {
ControlMsg::Spawn(p) => assert_eq!(p.instruction, "line1\nline2"),
other => panic!("expected spawn, got {other:?}"),
}
}
#[test]
fn agent_messages_tag_correctly() {
let result = AgentMsg::Result {
outcome: Outcome {
status: TerminalStatus::Completed,
partial: false,
result: json!("done"),
scheduled: Vec::new(),
subscriptions: Vec::new(),
},
};
let s = serde_json::to_string(&result).unwrap();
assert!(s.contains("\"type\":\"result\""));
assert!(s.contains("\"status\":\"completed\""));
let pong = serde_json::to_string(&AgentMsg::Pong { seq: 7 }).unwrap();
assert!(pong.contains("\"type\":\"pong\""));
assert!(pong.contains("\"seq\":7"));
}
#[test]
fn control_ping_cancel_tags() {
assert!(
serde_json::to_string(&ControlMsg::Ping { seq: 1 })
.unwrap()
.contains("\"type\":\"ping\"")
);
assert!(
serde_json::to_string(&ControlMsg::Cancel {
reason: "drain".into()
})
.unwrap()
.contains("\"type\":\"cancel\"")
);
}
#[test]
fn control_swap_intel_roundtrip_and_policy_default() {
let swap = ControlMsg::SwapIntel(Box::new(SwapIntel {
uri: "https://gw-a.example,https://gw-b.example".into(),
token: Some("rotated-secret".into()),
model: Some("claude-haiku-4".into()),
policy: SwapPolicy::RestartTurn,
}));
let s = serde_json::to_string(&swap).unwrap();
assert!(s.contains("\"type\":\"swap_intel\""));
assert!(s.contains("\"policy\":\"restart-turn\""));
let back: ControlMsg = serde_json::from_str(&s).unwrap();
match back {
ControlMsg::SwapIntel(p) => {
assert_eq!(p.uri, "https://gw-a.example,https://gw-b.example");
assert_eq!(p.model.as_deref(), Some("claude-haiku-4"));
assert_eq!(p.policy, SwapPolicy::RestartTurn);
}
other => panic!("expected swap_intel, got {other:?}"),
}
let minimal: SwapIntel = serde_json::from_str(r#"{"uri":"https://a.example"}"#).unwrap();
assert_eq!(minimal.policy, SwapPolicy::FinishOnOld);
assert!(minimal.model.is_none() && minimal.token.is_none());
}
#[test]
fn intel_health_roundtrips_and_carries_no_url_or_secret() {
let down = AgentMsg::IntelHealth {
all_down: true,
active: None,
};
let s = serde_json::to_string(&down).unwrap();
assert!(s.contains("\"type\":\"intel_health\""));
assert!(s.contains("\"all_down\":true"));
assert!(!s.contains("active"));
let back: AgentMsg = serde_json::from_str(&s).unwrap();
assert!(matches!(
back,
AgentMsg::IntelHealth {
all_down: true,
active: None
}
));
let up = AgentMsg::IntelHealth {
all_down: false,
active: Some(IntelActive {
index: 1,
transport: "https".into(),
}),
};
let s = serde_json::to_string(&up).unwrap();
assert!(s.contains("\"all_down\":false"));
assert!(s.contains("\"index\":1"));
assert!(s.contains("\"transport\":\"https\""));
assert!(!s.contains("https://"), "no full URI in the report: {s}");
let back: AgentMsg = serde_json::from_str(&s).unwrap();
match back {
AgentMsg::IntelHealth { all_down, active } => {
assert!(!all_down);
let a = active.unwrap();
assert_eq!(a.index, 1);
assert_eq!(a.transport, "https");
}
other => panic!("expected intel_health, got {other:?}"),
}
}
#[test]
fn narrow_tools_mints_one_supervisor_grant_and_survives_the_wire() {
let mut p = payload();
p.context_seed.insert(
0,
SeedMessage {
role: ALLOWED_TOOLS_ROLE.into(),
content: "[\"*\"]".into(),
},
);
p.narrow_tools(&["knowledge.search".to_string()]);
assert_eq!(
p.context_seed
.iter()
.filter(|m| m.role == ALLOWED_TOOLS_ROLE)
.count(),
1,
"one grant only — the forged `*` is gone"
);
assert_eq!(
p.allowed_tools(),
Some(vec!["knowledge.search".to_string()])
);
assert!(p.context_seed.iter().any(|m| m.content == "prior note"));
let msg = ControlMsg::Spawn(Box::new(p));
let mut buf = Vec::new();
frame::write_frame(&mut buf, &msg).unwrap();
let back: ControlMsg =
serde_json::from_slice(&frame::read_frame(&mut Cursor::new(buf)).unwrap().unwrap())
.unwrap();
match back {
ControlMsg::Spawn(p) => assert_eq!(
p.allowed_tools(),
Some(vec!["knowledge.search".to_string()])
),
other => panic!("expected spawn, got {other:?}"),
}
}
#[test]
fn an_unnarrowed_payload_has_no_grant_and_a_broken_one_grants_nothing() {
assert_eq!(payload().allowed_tools(), None);
assert!(parse_allowed_tools("not json").is_empty());
assert!(parse_allowed_tools("[\"a\",\"b.*\"]").len() == 2);
}
#[test]
fn control_pause_resume_roundtrip() {
let pause = serde_json::to_string(&ControlMsg::Pause).unwrap();
assert_eq!(pause, "{\"type\":\"pause\"}");
let resume = serde_json::to_string(&ControlMsg::Resume).unwrap();
assert_eq!(resume, "{\"type\":\"resume\"}");
assert!(matches!(
serde_json::from_str::<ControlMsg>(&pause).unwrap(),
ControlMsg::Pause
));
assert!(matches!(
serde_json::from_str::<ControlMsg>(&resume).unwrap(),
ControlMsg::Resume
));
}
}