use std::{ffi::OsString, num::NonZeroUsize, sync::MutexGuard};
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,
};
struct IsolatedRhoHome {
_dir: tempfile::TempDir,
_guard: MutexGuard<'static, ()>,
previous: Option<OsString>,
}
impl IsolatedRhoHome {
fn new() -> Self {
let guard = crate::paths::process_env_lock();
let dir = tempfile::tempdir().expect("rho home tempdir");
let previous = std::env::var_os("RHO_HOME");
std::env::set_var("RHO_HOME", dir.path());
Self {
_dir: dir,
_guard: guard,
previous,
}
}
}
impl Drop for IsolatedRhoHome {
fn drop(&mut self) {
match &self.previous {
Some(value) => std::env::set_var("RHO_HOME", value),
None => std::env::remove_var("RHO_HOME"),
}
}
}
struct ManagerFixture {
manager: SubagentManager,
_rho_home: IsolatedRhoHome,
}
impl ManagerFixture {
fn new(root: &Path) -> Self {
let rho_home = IsolatedRhoHome::new();
Self {
manager: SubagentManager::new(AgentExecutor::new(
Config::default(),
root.join("rho.toml"),
root.to_path_buf(),
SubagentHostInputBridge::new(),
)),
_rho_home: rho_home,
}
}
fn manager(&self) -> SubagentManager {
self.manager.clone()
}
}
fn manager(root: &Path) -> ManagerFixture {
ManagerFixture::new(root)
}
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_fixture = manager(root.path());
let tool = AgentTool::new(
_tool_fixture.manager(),
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 fixture = manager(root.path());
let manager = fixture.manager();
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 fixture = manager(root.path());
let error = fixture.manager().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 fixture = manager(root.path());
let manager = fixture.manager();
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 fixture = manager(root.path());
let enabled = AgentTool::new(fixture.manager(), root.path(), BackgroundSubagents::Enabled);
let disabled = AgentTool::new(
fixture.manager(),
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!(enabled
.spec()
.description
.contains("Issuing a foreground agent beside other tools does not background it"));
assert!(enabled
.spec()
.description
.contains("can delay the rest of that batch until the run finishes"));
assert!(disabled
.spec()
.description
.contains("Independent agent calls in the same batch run together"));
assert!(disabled
.spec()
.description
.contains("can delay the rest of that batch until the run finishes"));
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. Only background=true backgrounds a run; parallel batching does not. 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 fixture = manager(root.path());
let manager = fixture.manager();
manager.bind_parent_session(crate::subagent::RunPlacement::for_parent_session(
"session-1",
None,
));
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 fixture = manager(root.path());
let manager = fixture.manager();
manager.bind_parent_session(crate::subagent::RunPlacement::for_parent_session(
"session-1",
None,
));
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 fixture = manager(root.path());
let manager = fixture.manager();
manager.bind_parent_session(crate::subagent::RunPlacement::for_parent_session(
"session-1",
None,
));
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 fixture = manager(root.path());
let manager = fixture.manager();
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 fixture = manager(root.path());
let tool = AgentsTool::new(fixture.manager());
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 fixture = manager(root.path());
let manager = fixture.manager();
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 fixture = manager(root.path());
let manager = fixture.manager();
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)));
}