use std::num::NonZeroUsize;
use pretty_assertions::assert_eq;
use rho_sdk::{
tool::{
tool_progress_channel, Tool, ToolAccessMode, ToolContext, ToolExecutionPolicy,
ToolInvocation, ToolPreparationContext, ToolResourceKind,
},
CancellationToken, ToolCallId, Workspace,
};
use super::*;
use crate::app::subagent_host_input::SubagentHostInputBridge;
use crate::{
app::agent_executor::AgentExecutor, config::Config, diagnostics::test_diagnostics,
tools::agent_output::MODEL_NOTIFICATION_BYTES,
};
fn manager(root: &Path) -> SubagentManager {
SubagentManager::new(AgentExecutor::new(
Config::default(),
root.join("rho.toml"),
root.to_path_buf(),
SubagentHostInputBridge::new(),
))
}
fn invocation(arguments: serde_json::Value) -> ToolInvocation {
ToolInvocation::new(ToolCallId::from_string("call-1").unwrap(), arguments)
}
fn tool_context(root: &Path) -> ToolContext {
let (progress, _receiver) = tool_progress_channel(NonZeroUsize::new(4).unwrap());
ToolContext::new(
Some(Workspace::new(root).unwrap()),
CancellationToken::new(),
progress,
)
}
async fn call_agent(tool: &AgentTool, root: &Path, arguments: serde_json::Value) -> ToolOutput {
tool.call(invocation(arguments), tool_context(root))
.await
.expect("agent tool call")
}
#[test]
fn agent_tool_uses_agent_id_terminology() {
let root = tempfile::tempdir().unwrap();
let tool = AgentTool::new(
manager(root.path()),
root.path(),
BackgroundSubagents::Enabled,
);
let spec = tool.spec();
let properties = &spec.input_schema["properties"];
assert!(properties.get("agent_id").is_some());
assert_eq!(
spec.input_schema["required"],
serde_json::json!(["agent_id", "prompt"])
);
}
#[test]
fn delegated_manager_starts_empty() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
assert!(manager.list().is_empty());
assert!(manager.status("missing").is_none());
assert!(!manager.has_running_for_session("session-1"));
}
#[tokio::test]
async fn stopping_unknown_run_is_actionable() {
let root = tempfile::tempdir().unwrap();
let error = manager(root.path()).stop("abcdef").await.unwrap_err();
assert!(error.to_string().contains("unknown delegated run"));
}
#[tokio::test]
async fn background_start_receipt_is_the_registration() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
let tool = AgentTool::new(manager.clone(), root.path(), BackgroundSubagents::Enabled);
let result = call_agent(
&tool,
root.path(),
serde_json::json!({
"agent_id": "default",
"prompt": "background task",
"background": true,
}),
)
.await;
let runs = manager.list();
assert_eq!(runs.len(), 1);
let run_id = &runs[0].id;
assert_eq!(
result.content(),
format!("agent {run_id} (default) started in background\nattach: rho attach {run_id}")
);
}
#[test]
fn background_guidance_is_gated_by_capability() {
let root = tempfile::tempdir().unwrap();
let enabled = AgentTool::new(
manager(root.path()),
root.path(),
BackgroundSubagents::Enabled,
);
let disabled = AgentTool::new(
manager(root.path()),
root.path(),
BackgroundSubagents::Disabled,
);
assert!(enabled.spec().description.contains("background=true"));
assert!(enabled
.spec()
.description
.contains("Independent agent calls in the same batch run together"));
assert!(enabled
.spec()
.description
.contains("issue them in one turn for parallel work"));
assert!(disabled
.spec()
.description
.contains("Independent agent calls in the same batch run together"));
assert!(!disabled.spec().description.contains("background=true"));
assert_eq!(
enabled.spec().input_schema["properties"]["background"]["description"],
"Starts the run and returns an id immediately instead of waiting. Omit or set false to wait for the final result. Independent agent calls in the same batch run together either way."
);
let disabled_spec = disabled.spec();
assert!(!disabled_spec.description.contains("background=true"));
assert!(disabled_spec.input_schema["properties"]
.get("background")
.is_none());
}
fn notification(id: &str, agent_id: &str, state: RunState) -> SubagentNotification {
SubagentNotification {
snapshot: SubagentSnapshot {
id: id.into(),
agent_id: agent_id.into(),
elapsed: Duration::from_secs(5),
status: crate::subagent::RunStatus {
state,
turns: 1,
input_tokens: Some(10),
output_tokens: Some(2),
result: Some(format!("{id} result")),
..crate::subagent::RunStatus::default()
},
done: true,
},
}
}
#[test]
fn notification_prompts_batch_terminal_runs_into_one_message() {
let first = notification("aaa111", "worker", RunState::Ok);
let mut second = notification("bbb222", "reviewer", RunState::Stopped);
second.snapshot.status.error = Some("review stopped before completion".into());
second.snapshot.status.attachment_error = Some("log unavailable".into());
let notifications = vec![first, second];
let (model, display) = notification_prompts(¬ifications);
assert_eq!(model.matches("[agent notification]").count(), 1);
assert!(model.contains("agent aaa111 (worker): ok"));
assert!(model.contains("aaa111 result"));
assert!(model.contains("agent bbb222 (reviewer): stopped"));
assert!(model.contains("error: review stopped before completion"));
assert!(model.contains("attachment error: log unavailable"));
assert!(model.contains("treat its work as unverified"));
assert_eq!(
display,
"agent aaa111 (worker) finished - ok\nagent bbb222 (reviewer) finished - stopped"
);
}
#[test]
fn notification_prompts_bound_many_large_utf8_results_and_keep_run_statuses() {
let notifications = (0..96)
.map(|index| {
let id = format!("run{index:03}");
let mut notification = notification(&id, "worker", RunState::Ok);
notification.snapshot.status.result = Some("🦀".repeat(12 * 1024));
notification
})
.collect::<Vec<_>>();
let (model, _) = notification_prompts(¬ifications);
assert!(
model.len() <= MODEL_NOTIFICATION_BYTES,
"{}-byte notification exceeded the {}-byte budget",
model.len(),
MODEL_NOTIFICATION_BYTES
);
for index in 0..notifications.len() {
assert!(
model.contains(&format!("agent run{index:03} (worker): ok")),
"missing status for run {index}"
);
}
assert!(model.contains("Any omitted or truncated result details remain available"));
assert!(model.contains("`agents status`"));
assert!(model.contains("`rho attach <run-id>`"));
assert_eq!(model, notification_prompts(¬ifications).0);
let newer = (0..96)
.map(|index| {
let id = format!("new{index:03}");
let mut notification = notification(&id, "reviewer", RunState::Ok);
notification.snapshot.status.result = Some("🦀".repeat(12 * 1024));
notification
})
.collect::<Vec<_>>();
let newer = notification_prompts(&newer).0;
let retried_context = merge_notification_context(Some(&model), &newer);
assert!(retried_context.len() <= NOTIFICATION_CONTEXT_BYTES);
assert!(retried_context.contains("agent new000 (reviewer): ok"));
}
async fn spawn_background_run(manager: &SubagentManager, root: &Path) -> String {
let tool = AgentTool::new(manager.clone(), root, BackgroundSubagents::Enabled);
call_agent(
&tool,
root,
serde_json::json!({
"agent_id": "default",
"prompt": "background task",
"background": true,
}),
)
.await;
manager.list().last().unwrap().id.clone()
}
#[tokio::test]
async fn running_queries_are_scoped_to_the_parent_session() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
manager.set_session("session-1".into());
let id = spawn_background_run(&manager, root.path()).await;
assert!(!manager.has_running_for_session("session-2"));
manager.stop(&id).await.unwrap();
}
#[tokio::test]
async fn observed_terminal_run_is_not_redelivered() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
manager.set_session("session-1".into());
let id = spawn_background_run(&manager, root.path()).await;
let snapshot = manager.wait_done(&id).await.unwrap();
assert!(snapshot.done);
let observed = manager.observe(&id).unwrap();
assert!(observed.done);
assert!(manager.take_notifications("session-1").is_empty());
assert!(!manager.has_active_or_pending_notification("session-1"));
}
#[tokio::test]
async fn unobserved_terminal_runs_drain_as_one_batch() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
manager.set_session("session-1".into());
let first = spawn_background_run(&manager, root.path()).await;
let second = spawn_background_run(&manager, root.path()).await;
manager.wait_done(&first).await.unwrap();
manager.wait_done(&second).await.unwrap();
let batch = manager.take_notifications("session-1");
let ids = batch
.iter()
.map(|notification| notification.snapshot.id.clone())
.collect::<Vec<_>>();
assert_eq!(ids, vec![first, second], "batch drains in launch order");
assert!(
manager.take_notifications("session-1").is_empty(),
"a drained batch is observed and never redelivered"
);
}
#[test]
fn claim_terminal_costs_is_idempotent_and_session_scoped() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
manager.insert_completed_for_test("aaa111", "session-1", Some(0.034271));
manager.insert_completed_for_test("bbb222", "session-1", Some(0.01));
manager.insert_completed_for_test("ccc333", "session-2", Some(0.5));
manager.insert_completed_for_test("ddd444", "session-1", None);
assert_eq!(manager.claim_terminal_costs_usd_micros("session-1"), 44_271);
assert_eq!(manager.claim_terminal_costs_usd_micros("session-1"), 0);
assert_eq!(
manager.claim_terminal_costs_usd_micros("session-2"),
500_000
);
}
#[test]
fn lifecycle_tool_schema_is_stable() {
let root = tempfile::tempdir().unwrap();
let tool = AgentsTool::new(manager(root.path()));
let spec = tool.spec();
assert_eq!(spec.name, "agents");
assert_eq!(
spec.input_schema["properties"]["action"]["enum"],
serde_json::json!(["list", "status", "stop"])
);
let _ = test_diagnostics("test", "test");
}
fn one_access(
prepared: &rho_sdk::tool::PreparedToolInvocation<'_>,
) -> (ToolResourceKind, ToolAccessMode) {
let ToolExecutionPolicy::ResourceAware { accesses } = prepared.execution_policy() else {
panic!("expected a resource-aware invocation");
};
assert_eq!(accesses.len(), 1);
(accesses[0].resource().kind(), accesses[0].mode())
}
fn preparation_context(root: &Path) -> ToolPreparationContext {
ToolPreparationContext::new(
Some(Workspace::new(root).unwrap()),
CancellationToken::new(),
)
}
#[tokio::test]
async fn agent_and_agents_prepare_subagent_manager_resources() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
let agent = AgentTool::new(manager.clone(), root.path(), BackgroundSubagents::Enabled);
let agents = AgentsTool::new(manager);
let launch = agent
.prepare(
invocation(serde_json::json!({
"agent_id": "default",
"prompt": "task",
})),
preparation_context(root.path()),
)
.await
.unwrap();
assert_eq!(
one_access(&launch),
(ToolResourceKind::ManagerState, ToolAccessMode::Shared)
);
let background = agent
.prepare(
invocation(serde_json::json!({
"agent_id": "default",
"prompt": "task",
"background": true,
})),
preparation_context(root.path()),
)
.await
.unwrap();
assert_eq!(
one_access(&background),
(ToolResourceKind::ManagerState, ToolAccessMode::Shared)
);
let list = agents
.prepare(
invocation(serde_json::json!({"action": "list"})),
preparation_context(root.path()),
)
.await
.unwrap();
assert_eq!(
one_access(&list),
(ToolResourceKind::ManagerState, ToolAccessMode::Shared)
);
let status = agents
.prepare(
invocation(serde_json::json!({"action": "status", "id": "run-1"})),
preparation_context(root.path()),
)
.await
.unwrap();
assert_eq!(
one_access(&status),
(ToolResourceKind::ManagerState, ToolAccessMode::Shared)
);
let stop = agents
.prepare(
invocation(serde_json::json!({"action": "stop", "id": "run-1"})),
preparation_context(root.path()),
)
.await
.unwrap();
assert_eq!(
one_access(&stop),
(ToolResourceKind::ManagerState, ToolAccessMode::Shared)
);
}
#[tokio::test]
async fn concurrent_background_launches_register_together() {
let root = tempfile::tempdir().unwrap();
let manager = manager(root.path());
let tool = AgentTool::new(manager.clone(), root.path(), BackgroundSubagents::Enabled);
let first = call_agent(
&tool,
root.path(),
serde_json::json!({
"agent_id": "default",
"prompt": "first background task",
"background": true,
}),
);
let second = call_agent(
&tool,
root.path(),
serde_json::json!({
"agent_id": "default",
"prompt": "second background task",
"background": true,
}),
);
let (first, second) = tokio::join!(first, second);
let runs = manager.list();
assert_eq!(runs.len(), 2, "both background launches should register");
assert!(first.content().contains("started in background"));
assert!(second.content().contains("started in background"));
let ids = runs.iter().map(|run| run.id.as_str()).collect::<Vec<_>>();
assert!(ids.iter().any(|id| first.content().contains(id)));
assert!(ids.iter().any(|id| second.content().contains(id)));
}