From 91634d48c146ec3a5dbb1a44ad0708d007a66d59 Mon Sep 17 00:00:00 2001 From: Steven Enamakel <31011319+senamakel@users.noreply.github.com> Date: Sat, 16 May 2026 13:17:42 -0700 Subject: [PATCH] test(rust): deepen parallel subagent harness coverage (#1957) --- .github/workflows/e2e-reusable.yml | 18 +- AGENTS.md | 2 + .../agent/harness/builtin_definitions.rs | 64 +- .../tools/impl/agent/archetype_delegation.rs | 61 ++ .../tools/impl/agent/spawn_parallel_agents.rs | 382 +---------- .../impl/agent/spawn_parallel_agents_test.rs | 620 ++++++++++++++++++ .../tools/impl/agent/spawn_subagent.rs | 79 +++ 7 files changed, 833 insertions(+), 393 deletions(-) create mode 100644 src/openhuman/tools/impl/agent/spawn_parallel_agents_test.rs diff --git a/.github/workflows/e2e-reusable.yml b/.github/workflows/e2e-reusable.yml index 7745a7761..f6b00caad 100644 --- a/.github/workflows/e2e-reusable.yml +++ b/.github/workflows/e2e-reusable.yml @@ -125,9 +125,9 @@ jobs: if ! command -v appium >/dev/null 2>&1; then npm install -g appium@3 fi - if ! appium driver list --installed 2>/dev/null | grep -q chromium; then - appium driver install --source=npm appium-chromium-driver - fi + # `appium driver list --installed` can miss cached installs on some + # Appium builds; install idempotently and ignore "already installed". + appium driver install --source=npm appium-chromium-driver >/dev/null 2>&1 || true - name: Build E2E app run: pnpm --filter openhuman-app test:e2e:build @@ -293,9 +293,9 @@ jobs: if ! command -v appium >/dev/null 2>&1; then npm install -g appium@3 fi - if ! appium driver list --installed 2>/dev/null | grep -q chromium; then - appium driver install --source=npm appium-chromium-driver - fi + # `appium driver list --installed` can miss cached installs on some + # Appium builds; install idempotently and ignore "already installed". + appium driver install --source=npm appium-chromium-driver >/dev/null 2>&1 || true - name: Build E2E app run: pnpm --filter openhuman-app test:e2e:build @@ -387,9 +387,9 @@ jobs: if ! command -v appium >/dev/null 2>&1; then npm install -g appium@3 fi - if ! appium driver list --installed 2>/dev/null | grep -q chromium; then - appium driver install --source=npm appium-chromium-driver - fi + # `appium driver list --installed` can miss cached installs on some + # Appium builds; install idempotently and ignore "already installed". + appium driver install --source=npm appium-chromium-driver >/dev/null 2>&1 || true - name: Build E2E app run: pnpm --filter openhuman-app test:e2e:build diff --git a/AGENTS.md b/AGENTS.md index 71af2296f..0c8e463eb 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -247,6 +247,8 @@ run_case() { } ``` +- **Rust test file naming**: when extracting Rust tests out of an implementation file, prefer a sibling `*_test.rs` file wired in with `#[cfg(test)] #[path = "..._test.rs"] mod tests;`. Do not create ad hoc `_test/` or `_tests/` directories for single-module Rust tests unless a broader multi-file test fixture truly requires a directory. + ### Test authoring checklist - Add/update unit tests for logic changes before stacking additional features. diff --git a/src/openhuman/agent/harness/builtin_definitions.rs b/src/openhuman/agent/harness/builtin_definitions.rs index 35d2a4f26..c86146f7d 100644 --- a/src/openhuman/agent/harness/builtin_definitions.rs +++ b/src/openhuman/agent/harness/builtin_definitions.rs @@ -8,7 +8,9 @@ //! Custom TOML definitions loaded later by //! [`super::definition_loader`] override any built-in with the same id. -use super::definition::{AgentDefinition, DefinitionSource}; +use super::definition::AgentDefinition; +#[cfg(test)] +use super::definition::DefinitionSource; /// All built-in definitions, in stable order. /// @@ -30,7 +32,10 @@ pub fn all() -> Vec { let mut defs = crate::openhuman::agent::agents::load_builtins() .expect("built-in agent TOML must always parse (see agents/*/agent.toml)"); #[cfg(test)] - defs.push(test_inherit_echo_def()); + { + defs.push(test_inherit_echo_def()); + defs.push(test_inherit_parallel_worker_def()); + } defs } @@ -72,6 +77,40 @@ pub(crate) fn test_inherit_echo_def() -> AgentDefinition { } } +/// Test-only sub-agent: inherits the parent's provider and exposes a +/// single named tool so long-running parallel fan-out tests can drive +/// repeated nested tool calls through the real sub-agent loop. +#[cfg(test)] +pub(crate) fn test_inherit_parallel_worker_def() -> AgentDefinition { + use super::definition::{ModelSpec, PromptSource, SandboxMode, ToolScope}; + AgentDefinition { + id: "__test_inherit_parallel_worker".into(), + when_to_use: "test-only parallel sub-agent that inherits the parent provider".into(), + display_name: None, + system_prompt: PromptSource::Inline("You are a test parallel worker.".into()), + omit_identity: true, + omit_memory_context: true, + omit_safety_preamble: true, + omit_skills_catalog: true, + omit_profile: true, + omit_memory_md: true, + model: ModelSpec::Inherit, + temperature: 0.0, + tools: ToolScope::Named(vec!["fixture_step".into()]), + disallowed_tools: vec![], + skill_filter: None, + extra_tools: vec![], + max_iterations: 6, + max_result_chars: None, + timeout_secs: None, + sandbox_mode: SandboxMode::None, + background: false, + subagents: vec![], + delegate_name: None, + source: DefinitionSource::Builtin, + } +} + #[cfg(test)] mod tests { use super::*; @@ -79,10 +118,10 @@ mod tests { #[test] fn all_definitions_present() { let defs = all(); - // +1 for the cfg(test) `__test_inherit_echo` def appended by all(). + // +2 for the cfg(test) inherit-based test defs appended by all(). assert_eq!( defs.len(), - crate::openhuman::agent::agents::BUILTINS.len() + 1 + crate::openhuman::agent::agents::BUILTINS.len() + 2 ); } @@ -99,6 +138,23 @@ mod tests { ); } + #[test] + fn test_inherit_parallel_worker_is_present_and_inherits() { + use super::super::definition::{ModelSpec, ToolScope}; + let def = all() + .into_iter() + .find(|d| d.id == "__test_inherit_parallel_worker") + .expect("test-only parallel worker must be registered in test builds"); + assert!( + matches!(def.model, ModelSpec::Inherit), + "must be Inherit so the sub-agent uses the parent's (mock) provider" + ); + assert!( + matches!(def.tools, ToolScope::Named(ref names) if names == &vec!["fixture_step".to_string()]), + "parallel worker must expose only the fixture_step tool" + ); + } + #[test] fn all_builtin_ids_are_stamped_builtin_source() { for def in all() { diff --git a/src/openhuman/tools/impl/agent/archetype_delegation.rs b/src/openhuman/tools/impl/agent/archetype_delegation.rs index 50a4043fb..55d1601f1 100644 --- a/src/openhuman/tools/impl/agent/archetype_delegation.rs +++ b/src/openhuman/tools/impl/agent/archetype_delegation.rs @@ -58,3 +58,64 @@ impl Tool for ArchetypeDelegationTool { super::dispatch_subagent(&self.agent_id, &self.tool_name, &prompt, None).await } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::openhuman::agent::harness::definition::AgentDefinitionRegistry; + + fn sample_tool() -> ArchetypeDelegationTool { + ArchetypeDelegationTool { + tool_name: "delegate_researcher".to_string(), + agent_id: "researcher".to_string(), + tool_description: "Use for web and docs research.".to_string(), + } + } + + #[test] + fn metadata_methods_expose_name_description_and_system_category() { + let tool = sample_tool(); + assert_eq!(tool.name(), "delegate_researcher"); + assert_eq!(tool.description(), "Use for web and docs research."); + assert_eq!(tool.permission_level(), PermissionLevel::Execute); + assert_eq!(tool.category(), ToolCategory::System); + } + + #[test] + fn parameters_schema_requires_prompt_only() { + let tool = sample_tool(); + let schema = tool.parameters_schema(); + assert_eq!(schema["type"], "object"); + assert_eq!(schema["required"], json!(["prompt"])); + assert_eq!(schema["properties"]["prompt"]["type"], "string"); + } + + #[tokio::test] + async fn execute_rejects_missing_or_blank_prompt() { + let tool = sample_tool(); + + let missing = tool.execute(json!({})).await.unwrap(); + assert!(missing.is_error); + assert!(missing.output().contains("`prompt` is required")); + + let blank = tool.execute(json!({ "prompt": " " })).await.unwrap(); + assert!(blank.is_error); + assert!(blank.output().contains("`prompt` is required")); + } + + #[tokio::test] + async fn execute_accepts_non_empty_prompt_and_reaches_dispatch_path() { + let _ = AgentDefinitionRegistry::init_global_builtins(); + let tool = sample_tool(); + let result = tool + .execute(json!({ "prompt": "find the answer" })) + .await + .unwrap(); + + let out = result.output(); + assert!( + !out.contains("`prompt` is required"), + "non-empty prompt should bypass local validation, got: {out}" + ); + } +} diff --git a/src/openhuman/tools/impl/agent/spawn_parallel_agents.rs b/src/openhuman/tools/impl/agent/spawn_parallel_agents.rs index d409c349a..6b4997410 100644 --- a/src/openhuman/tools/impl/agent/spawn_parallel_agents.rs +++ b/src/openhuman/tools/impl/agent/spawn_parallel_agents.rs @@ -486,383 +486,5 @@ fn with_ownership_boundary(prompt: &str, ownership: Option<&str>) -> String { } #[cfg(test)] -mod tests { - use super::*; - use crate::openhuman::agent::harness::fork_context::{ - with_parent_context, ParentExecutionContext, - }; - use crate::openhuman::context::prompt::ToolCallFormat; - use crate::openhuman::memory::{ - Memory, MemoryCategory, MemoryEntry, NamespaceSummary, RecallOpts, - }; - use crate::openhuman::providers::{ChatRequest, ChatResponse, Provider}; - use std::sync::{ - atomic::{AtomicUsize, Ordering}, - Arc, - }; - use tokio::time::{sleep, Duration}; - - #[test] - fn metadata_methods_expose_execute_permission_and_schema() { - let tool = SpawnParallelAgentsTool::default(); - assert_eq!(tool.name(), "spawn_parallel_agents"); - assert!(tool.description().contains("independent sub-agent tasks")); - assert_eq!(tool.permission_level(), PermissionLevel::Execute); - let schema = tool.parameters_schema(); - assert_eq!(schema["required"][0], "tasks"); - assert_eq!(schema["properties"]["tasks"]["minItems"], 2); - } - - #[test] - fn ownership_boundary_is_prepended_when_present() { - let prompt = with_ownership_boundary("implement tests", Some("files: src/foo.rs")); - assert!(prompt.starts_with("[Ownership Boundary]")); - assert!(prompt.contains("files: src/foo.rs")); - assert!(prompt.contains("[Task]\nimplement tests")); - } - - #[tokio::test] - async fn rejects_single_task() { - let tool = SpawnParallelAgentsTool::new(); - let result = tool - .execute(json!({ - "tasks": [{ "agent_id": "researcher", "prompt": "only one" }] - })) - .await - .unwrap(); - assert!(result.is_error); - assert!(result.output().contains("at least two")); - } - - #[tokio::test] - async fn rejects_missing_or_invalid_tasks_before_parent_lookup() { - let tool = SpawnParallelAgentsTool::new(); - - let missing = tool.execute(json!({})).await.expect_err("missing tasks"); - assert!(missing.to_string().contains("Missing 'tasks'")); - - let invalid = tool - .execute(json!({ "tasks": "not an array" })) - .await - .expect_err("invalid tasks"); - assert!(invalid.to_string().contains("Invalid tasks array")); - } - - #[tokio::test] - async fn rejects_two_tasks_outside_agent_turn() { - let tool = SpawnParallelAgentsTool::new(); - let result = tool - .execute(json!({ - "tasks": [ - { "agent_id": "researcher", "prompt": "one" }, - { "agent_id": "planner", "prompt": "two" } - ] - })) - .await - .expect("tool result"); - assert!(result.is_error); - assert!(result.output().contains("outside of an agent turn")); - } - - struct NoopProvider; - - #[async_trait] - impl Provider for NoopProvider { - async fn chat_with_system( - &self, - _system_prompt: Option<&str>, - _message: &str, - _model: &str, - _temperature: f64, - ) -> anyhow::Result { - Ok("ok".into()) - } - - async fn chat( - &self, - _request: ChatRequest<'_>, - _model: &str, - _temperature: f64, - ) -> anyhow::Result { - Ok(ChatResponse { - text: Some("ok".into()), - tool_calls: Vec::new(), - usage: None, - }) - } - } - - struct ConcurrentProvider { - active: AtomicUsize, - max_active: AtomicUsize, - } - - impl ConcurrentProvider { - fn new() -> Arc { - Arc::new(Self { - active: AtomicUsize::new(0), - max_active: AtomicUsize::new(0), - }) - } - - fn max_active(&self) -> usize { - self.max_active.load(Ordering::SeqCst) - } - - fn observe_active(&self, current: usize) { - let mut observed = self.max_active.load(Ordering::SeqCst); - while current > observed { - match self.max_active.compare_exchange( - observed, - current, - Ordering::SeqCst, - Ordering::SeqCst, - ) { - Ok(_) => break, - Err(next) => observed = next, - } - } - } - } - - #[async_trait] - impl Provider for ConcurrentProvider { - async fn chat_with_system( - &self, - _system_prompt: Option<&str>, - _message: &str, - _model: &str, - _temperature: f64, - ) -> anyhow::Result { - Ok("ok".into()) - } - - async fn chat( - &self, - _request: ChatRequest<'_>, - _model: &str, - _temperature: f64, - ) -> anyhow::Result { - let current = self.active.fetch_add(1, Ordering::SeqCst) + 1; - self.observe_active(current); - sleep(Duration::from_millis(50)).await; - self.active.fetch_sub(1, Ordering::SeqCst); - Ok(ChatResponse { - text: Some("parallel ok".into()), - tool_calls: Vec::new(), - usage: None, - }) - } - } - - struct NoopMemory; - - #[async_trait] - impl Memory for NoopMemory { - async fn store( - &self, - _namespace: &str, - _key: &str, - _content: &str, - _category: MemoryCategory, - _session_id: Option<&str>, - ) -> anyhow::Result<()> { - Ok(()) - } - - async fn recall( - &self, - _query: &str, - _limit: usize, - _opts: RecallOpts<'_>, - ) -> anyhow::Result> { - Ok(Vec::new()) - } - - async fn get(&self, _namespace: &str, _key: &str) -> anyhow::Result> { - Ok(None) - } - - async fn list( - &self, - _namespace: Option<&str>, - _category: Option<&MemoryCategory>, - _session_id: Option<&str>, - ) -> anyhow::Result> { - Ok(Vec::new()) - } - - async fn forget(&self, _namespace: &str, _key: &str) -> anyhow::Result { - Ok(false) - } - - async fn namespace_summaries(&self) -> anyhow::Result> { - Ok(Vec::new()) - } - - async fn count(&self) -> anyhow::Result { - Ok(0) - } - - async fn health_check(&self) -> bool { - true - } - - fn name(&self) -> &str { - "noop" - } - } - - fn parent_context_with_provider( - max_parallel_tools: usize, - provider: Arc, - ) -> ParentExecutionContext { - let agent_config = crate::openhuman::config::AgentConfig { - max_parallel_tools, - ..Default::default() - }; - ParentExecutionContext { - provider, - all_tools: Arc::new(Vec::new()), - all_tool_specs: Arc::new(Vec::new()), - model_name: "test-model".into(), - temperature: 0.2, - workspace_dir: std::env::temp_dir(), - memory: Arc::new(NoopMemory), - agent_config, - skills: Arc::new(Vec::new()), - memory_context: Arc::new(None), - session_id: "session-test".into(), - channel: "test".into(), - connected_integrations: Vec::new(), - tool_call_format: ToolCallFormat::PFormat, - session_key: "0_test".into(), - session_parent_prefix: None, - on_progress: None, - } - } - - fn parent_context(max_parallel_tools: usize) -> ParentExecutionContext { - parent_context_with_provider(max_parallel_tools, Arc::new(NoopProvider)) - } - - #[tokio::test] - async fn rejects_more_tasks_than_parent_parallel_limit() { - let tool = SpawnParallelAgentsTool::new(); - let parent = parent_context(2); - let result = with_parent_context(parent, async { - tool.execute(json!({ - "tasks": [ - { "agent_id": "researcher", "prompt": "one" }, - { "agent_id": "planner", "prompt": "two" }, - { "agent_id": "critic", "prompt": "three" } - ] - })) - .await - }) - .await - .expect("tool result"); - assert!(result.is_error); - assert!(result.output().contains("max_parallel_tools")); - } - - #[tokio::test] - async fn collects_immediate_task_validation_failures() { - let _ = AgentDefinitionRegistry::init_global_builtins(); - let tool = SpawnParallelAgentsTool::new(); - let parent = parent_context(4); - - let result = with_parent_context(parent, async { - tool.execute(json!({ - "tasks": [ - { "agent_id": " ", "prompt": "missing agent", "ownership": "files: none" }, - { "agent_id": "__missing_agent__", "prompt": "unknown agent" }, - { "agent_id": "integrations_agent", "prompt": "needs toolkit" } - ] - })) - .await - }) - .await - .expect("tool result"); - - assert!(!result.is_error, "{}", result.output()); - let body: serde_json::Value = serde_json::from_str(&result.output()).expect("json output"); - assert_eq!(body["parallel_agents"]["total"], 3); - assert_eq!(body["parallel_agents"]["failed"], 3); - let errors = body["parallel_agents"]["results"] - .as_array() - .expect("results") - .iter() - .map(|result| result["error"].as_str().unwrap_or_default()) - .collect::>(); - assert!(errors - .iter() - .any(|error| error.contains("agent_id and prompt"))); - assert!(errors - .iter() - .any(|error| error.contains("unknown agent_id"))); - assert!(errors - .iter() - .any(|error| error.contains("requires toolkit"))); - } - - // After upstream PR #1858 (`feat(ai): unified per-workload provider - // routing + chat-provider factory`), the subagent runner resolves a - // real provider via `resolve_subagent_provider` using the loaded - // `Config`'s workload routing instead of inheriting `parent.provider` - // for `ModelSpec::Hint` agents like `researcher` / `planner`. That - // bypasses this test's mock `ConcurrentProvider`: on CI the resolved - // OpenHuman backend is unreachable so subagents fail with - // `succeeded == 0`; locally a configured cloud provider may answer - // with text that isn't `"parallel ok"`. The right fix is a - // workspace-isolated test fixture that forces the workload factory - // to fall back to the parent provider — tracked as a follow-up. - #[ignore = "subagent_runner provider routing bypasses mock provider after upstream PR #1858"] - #[tokio::test] - async fn runs_valid_tasks_concurrently_and_collects_successes() { - let _ = AgentDefinitionRegistry::init_global_builtins(); - let tool = SpawnParallelAgentsTool::new(); - let provider_impl = ConcurrentProvider::new(); - let provider: Arc = provider_impl.clone(); - let parent = parent_context_with_provider(4, provider); - - let result = with_parent_context(parent, async { - tool.execute(json!({ - "tasks": [ - { - "agent_id": "researcher", - "prompt": "summarize alpha", - "ownership": "files: alpha.rs" - }, - { - "agent_id": "planner", - "prompt": "plan beta", - "ownership": "files: beta.rs" - } - ] - })) - .await - }) - .await - .expect("tool result"); - - assert!(!result.is_error, "{}", result.output()); - let body: serde_json::Value = serde_json::from_str(&result.output()).expect("json output"); - assert_eq!(body["parallel_agents"]["total"], 2); - assert_eq!(body["parallel_agents"]["succeeded"], 2); - assert_eq!(body["parallel_agents"]["failed"], 0); - let results = body["parallel_agents"]["results"] - .as_array() - .expect("results"); - assert_eq!(results.len(), 2); - assert!(results.iter().all(|result| result["success"] == true)); - assert!(results - .iter() - .all(|result| result["output"].as_str() == Some("parallel ok"))); - assert!( - provider_impl.max_active() >= 2, - "expected overlapping provider calls, max_active={}", - provider_impl.max_active() - ); - } -} +#[path = "spawn_parallel_agents_test.rs"] +mod tests; diff --git a/src/openhuman/tools/impl/agent/spawn_parallel_agents_test.rs b/src/openhuman/tools/impl/agent/spawn_parallel_agents_test.rs new file mode 100644 index 000000000..aa364457d --- /dev/null +++ b/src/openhuman/tools/impl/agent/spawn_parallel_agents_test.rs @@ -0,0 +1,620 @@ +use super::*; +use crate::openhuman::agent::dispatcher::NativeToolDispatcher; +use crate::openhuman::agent::harness::definition::AgentDefinitionRegistry; +use crate::openhuman::agent::harness::fork_context::{with_parent_context, ParentExecutionContext}; +use crate::openhuman::agent::Agent; +use crate::openhuman::config::AgentConfig; +use crate::openhuman::context::prompt::ToolCallFormat; +use crate::openhuman::memory::{Memory, MemoryCategory, MemoryEntry, NamespaceSummary, RecallOpts}; +use crate::openhuman::providers::traits::ProviderCapabilities; +use crate::openhuman::providers::{ + ChatRequest, ChatResponse, ConversationMessage, Provider, ToolCall, +}; +use crate::openhuman::tools::{PermissionLevel, Tool, ToolResult}; +use async_trait::async_trait; +use parking_lot::Mutex; +use serde_json::json; +use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, +}; +use tokio::time::{sleep, Duration}; + +const PARENT_PROMPT_CANARY: &str = "parallel-fanout-e2e-canary"; +const RESEARCH_PROMPT_CANARY: &str = "research-branch-canary"; +const PLANNER_PROMPT_CANARY: &str = "planner-branch-canary"; +const RESEARCH_DONE_CANARY: &str = "research-finished-canary"; +const PLANNER_DONE_CANARY: &str = "planner-finished-canary"; +const FINAL_CANARY: &str = "parallel-summary-canary"; + +#[test] +fn metadata_methods_expose_execute_permission_and_schema() { + let tool = SpawnParallelAgentsTool::default(); + assert_eq!(tool.name(), "spawn_parallel_agents"); + assert!(tool.description().contains("independent sub-agent tasks")); + assert_eq!(tool.permission_level(), PermissionLevel::Execute); + let schema = tool.parameters_schema(); + assert_eq!(schema["required"][0], "tasks"); + assert_eq!(schema["properties"]["tasks"]["minItems"], 2); +} + +#[test] +fn ownership_boundary_is_prepended_when_present() { + let prompt = with_ownership_boundary("implement tests", Some("files: src/foo.rs")); + assert!(prompt.starts_with("[Ownership Boundary]")); + assert!(prompt.contains("files: src/foo.rs")); + assert!(prompt.contains("[Task]\nimplement tests")); +} + +#[tokio::test] +async fn rejects_single_task() { + let tool = SpawnParallelAgentsTool::new(); + let result = tool + .execute(json!({ + "tasks": [{ "agent_id": "researcher", "prompt": "only one" }] + })) + .await + .unwrap(); + assert!(result.is_error); + assert!(result.output().contains("at least two")); +} + +#[tokio::test] +async fn rejects_missing_or_invalid_tasks_before_parent_lookup() { + let tool = SpawnParallelAgentsTool::new(); + + let missing = tool.execute(json!({})).await.expect_err("missing tasks"); + assert!(missing.to_string().contains("Missing 'tasks'")); + + let invalid = tool + .execute(json!({ "tasks": "not an array" })) + .await + .expect_err("invalid tasks"); + assert!(invalid.to_string().contains("Invalid tasks array")); +} + +#[tokio::test] +async fn rejects_two_tasks_outside_agent_turn() { + let tool = SpawnParallelAgentsTool::new(); + let result = tool + .execute(json!({ + "tasks": [ + { "agent_id": "researcher", "prompt": "one" }, + { "agent_id": "planner", "prompt": "two" } + ] + })) + .await + .expect("tool result"); + assert!(result.is_error); + assert!(result.output().contains("outside of an agent turn")); +} + +struct NoopProvider; + +#[async_trait] +impl Provider for NoopProvider { + async fn chat_with_system( + &self, + _system_prompt: Option<&str>, + _message: &str, + _model: &str, + _temperature: f64, + ) -> anyhow::Result { + Ok("ok".into()) + } + + async fn chat( + &self, + _request: ChatRequest<'_>, + _model: &str, + _temperature: f64, + ) -> anyhow::Result { + Ok(ChatResponse { + text: Some("ok".into()), + tool_calls: Vec::new(), + usage: None, + }) + } +} + +struct NoopMemory; + +#[async_trait] +impl Memory for NoopMemory { + async fn store( + &self, + _namespace: &str, + _key: &str, + _content: &str, + _category: MemoryCategory, + _session_id: Option<&str>, + ) -> anyhow::Result<()> { + Ok(()) + } + + async fn recall( + &self, + _query: &str, + _limit: usize, + _opts: RecallOpts<'_>, + ) -> anyhow::Result> { + Ok(Vec::new()) + } + + async fn get(&self, _namespace: &str, _key: &str) -> anyhow::Result> { + Ok(None) + } + + async fn list( + &self, + _namespace: Option<&str>, + _category: Option<&MemoryCategory>, + _session_id: Option<&str>, + ) -> anyhow::Result> { + Ok(Vec::new()) + } + + async fn forget(&self, _namespace: &str, _key: &str) -> anyhow::Result { + Ok(false) + } + + async fn namespace_summaries(&self) -> anyhow::Result> { + Ok(Vec::new()) + } + + async fn count(&self) -> anyhow::Result { + Ok(0) + } + + async fn health_check(&self) -> bool { + true + } + + fn name(&self) -> &str { + "noop" + } +} + +fn parent_context_with_provider( + max_parallel_tools: usize, + provider: Arc, +) -> ParentExecutionContext { + let agent_config = AgentConfig { + max_parallel_tools, + ..Default::default() + }; + ParentExecutionContext { + provider, + all_tools: Arc::new(Vec::new()), + all_tool_specs: Arc::new(Vec::new()), + model_name: "test-model".into(), + temperature: 0.2, + workspace_dir: std::env::temp_dir(), + memory: Arc::new(NoopMemory), + agent_config, + skills: Arc::new(Vec::new()), + memory_context: Arc::new(None), + session_id: "session-test".into(), + channel: "test".into(), + connected_integrations: Vec::new(), + tool_call_format: ToolCallFormat::PFormat, + session_key: "0_test".into(), + session_parent_prefix: None, + on_progress: None, + } +} + +fn parent_context(max_parallel_tools: usize) -> ParentExecutionContext { + parent_context_with_provider(max_parallel_tools, Arc::new(NoopProvider)) +} + +#[tokio::test] +async fn rejects_more_tasks_than_parent_parallel_limit() { + let tool = SpawnParallelAgentsTool::new(); + let parent = parent_context(2); + let result = with_parent_context(parent, async { + tool.execute(json!({ + "tasks": [ + { "agent_id": "researcher", "prompt": "one" }, + { "agent_id": "planner", "prompt": "two" }, + { "agent_id": "critic", "prompt": "three" } + ] + })) + .await + }) + .await + .expect("tool result"); + assert!(result.is_error); + assert!(result.output().contains("max_parallel_tools")); +} + +#[tokio::test] +async fn collects_immediate_task_validation_failures() { + let _ = AgentDefinitionRegistry::init_global_builtins(); + let tool = SpawnParallelAgentsTool::new(); + let parent = parent_context(4); + + let result = with_parent_context(parent, async { + tool.execute(json!({ + "tasks": [ + { "agent_id": " ", "prompt": "missing agent", "ownership": "files: none" }, + { "agent_id": "__missing_agent__", "prompt": "unknown agent" }, + { "agent_id": "integrations_agent", "prompt": "needs toolkit" } + ] + })) + .await + }) + .await + .expect("tool result"); + + assert!(!result.is_error, "{}", result.output()); + let body: serde_json::Value = serde_json::from_str(&result.output()).expect("json output"); + assert_eq!(body["parallel_agents"]["total"], 3); + assert_eq!(body["parallel_agents"]["failed"], 3); + let errors = body["parallel_agents"]["results"] + .as_array() + .expect("results") + .iter() + .map(|result| result["error"].as_str().unwrap_or_default()) + .collect::>(); + assert!(errors + .iter() + .any(|error| error.contains("agent_id and prompt"))); + assert!(errors + .iter() + .any(|error| error.contains("unknown agent_id"))); + assert!(errors + .iter() + .any(|error| error.contains("requires toolkit"))); +} + +#[derive(Default)] +struct FixtureStepState { + calls: AtomicUsize, +} + +struct FixtureStepTool { + state: Arc, +} + +#[async_trait] +impl Tool for FixtureStepTool { + fn name(&self) -> &str { + "fixture_step" + } + + fn description(&self) -> &str { + "Fixture tool used by parallel subagent tests." + } + + fn parameters_schema(&self) -> serde_json::Value { + json!({ + "type": "object", + "required": ["branch", "step"], + "properties": { + "branch": { "type": "string" }, + "step": { "type": "integer" } + } + }) + } + + async fn execute(&self, args: serde_json::Value) -> anyhow::Result { + let branch = args + .get("branch") + .and_then(|v| v.as_str()) + .unwrap_or("unknown"); + let step = args.get("step").and_then(|v| v.as_u64()).unwrap_or(0); + self.state.calls.fetch_add(1, Ordering::SeqCst); + Ok(ToolResult::success(format!("{branch}-step-{step}-ok"))) + } + + fn permission_level(&self) -> PermissionLevel { + PermissionLevel::None + } +} + +#[derive(Default)] +struct ParallelHarnessState { + total_calls: AtomicUsize, + active_subagent_calls: AtomicUsize, + max_active_subagent_calls: AtomicUsize, + seen_payloads: Mutex>, +} + +#[derive(Clone, Default)] +struct ParallelHarnessProvider { + state: Arc, +} + +impl ParallelHarnessProvider { + fn total_calls(&self) -> usize { + self.state.total_calls.load(Ordering::SeqCst) + } + + fn max_active_subagent_calls(&self) -> usize { + self.state.max_active_subagent_calls.load(Ordering::SeqCst) + } + + fn record_active_peak(&self, current: usize) { + let mut observed = self.state.max_active_subagent_calls.load(Ordering::SeqCst); + while current > observed { + match self.state.max_active_subagent_calls.compare_exchange( + observed, + current, + Ordering::SeqCst, + Ordering::SeqCst, + ) { + Ok(_) => break, + Err(next) => observed = next, + } + } + } + + async fn respond_for_subagent(&self, flattened: &str) -> anyhow::Result { + let current = self + .state + .active_subagent_calls + .fetch_add(1, Ordering::SeqCst) + + 1; + self.record_active_peak(current); + sleep(Duration::from_millis(25)).await; + + let response = (|| -> anyhow::Result { + if flattened.contains(RESEARCH_PROMPT_CANARY) { + if flattened.contains("research-step-3-ok") { + Ok(text_response(RESEARCH_DONE_CANARY)) + } else if flattened.contains("research-step-2-ok") { + Ok(tool_response( + "fixture_step", + json!({ "branch": "research", "step": 3 }), + )) + } else if flattened.contains("research-step-1-ok") { + Ok(tool_response( + "fixture_step", + json!({ "branch": "research", "step": 2 }), + )) + } else { + Ok(tool_response( + "fixture_step", + json!({ "branch": "research", "step": 1 }), + )) + } + } else if flattened.contains(PLANNER_PROMPT_CANARY) { + if flattened.contains("planner-step-3-ok") { + Ok(text_response(PLANNER_DONE_CANARY)) + } else if flattened.contains("planner-step-2-ok") { + Ok(tool_response( + "fixture_step", + json!({ "branch": "planner", "step": 3 }), + )) + } else if flattened.contains("planner-step-1-ok") { + Ok(tool_response( + "fixture_step", + json!({ "branch": "planner", "step": 2 }), + )) + } else { + Ok(tool_response( + "fixture_step", + json!({ "branch": "planner", "step": 1 }), + )) + } + } else { + anyhow::bail!("unexpected subagent payload: {flattened}"); + } + })(); + + self.state + .active_subagent_calls + .fetch_sub(1, Ordering::SeqCst); + response + } +} + +#[async_trait] +impl Provider for ParallelHarnessProvider { + fn capabilities(&self) -> ProviderCapabilities { + ProviderCapabilities { + native_tool_calling: true, + vision: false, + } + } + + async fn chat_with_system( + &self, + _system_prompt: Option<&str>, + _message: &str, + _model: &str, + _temperature: f64, + ) -> anyhow::Result { + Ok("ok".into()) + } + + async fn chat( + &self, + request: ChatRequest<'_>, + _model: &str, + _temperature: f64, + ) -> anyhow::Result { + self.state.total_calls.fetch_add(1, Ordering::SeqCst); + let flattened = request + .messages + .iter() + .map(|m| format!("{}:{}", m.role, m.content)) + .collect::>() + .join("\n"); + self.state.seen_payloads.lock().push(flattened.clone()); + + if flattened.contains(PARENT_PROMPT_CANARY) { + if flattened.contains(RESEARCH_DONE_CANARY) && flattened.contains(PLANNER_DONE_CANARY) { + return Ok(text_response(format!( + "{FINAL_CANARY}: merged {RESEARCH_DONE_CANARY} and {PLANNER_DONE_CANARY}" + ))); + } + + return Ok(tool_response( + "spawn_parallel_agents", + json!({ + "tasks": [ + { + "agent_id": "__test_inherit_parallel_worker", + "prompt": format!("Work the research branch: {RESEARCH_PROMPT_CANARY}"), + "ownership": "scope: research" + }, + { + "agent_id": "__test_inherit_parallel_worker", + "prompt": format!("Work the planning branch: {PLANNER_PROMPT_CANARY}"), + "ownership": "scope: planning" + } + ] + }), + )); + } + + self.respond_for_subagent(&flattened).await + } +} + +fn text_response(text: impl Into) -> ChatResponse { + ChatResponse { + text: Some(text.into()), + tool_calls: Vec::new(), + usage: None, + } +} + +fn tool_response(name: &str, arguments: serde_json::Value) -> ChatResponse { + ChatResponse { + text: Some(String::new()), + tool_calls: vec![ToolCall { + id: format!("call-{name}"), + name: name.to_string(), + arguments: arguments.to_string(), + }], + usage: None, + } +} + +#[tokio::test] +async fn agent_turn_runs_long_parallel_subagent_flow_with_many_nested_tool_calls() { + AgentDefinitionRegistry::init_global_builtins().unwrap(); + + let workspace = tempfile::TempDir::new().expect("temp workspace"); + let workspace_path = workspace.path().to_path_buf(); + let provider = ParallelHarnessProvider::default(); + let fixture_state = Arc::new(FixtureStepState::default()); + + let memory_cfg = crate::openhuman::config::MemoryConfig { + backend: "none".into(), + ..crate::openhuman::config::MemoryConfig::default() + }; + let mem: Arc = + Arc::from(crate::openhuman::memory::create_memory(&memory_cfg, &workspace_path).unwrap()); + + let tools: Vec> = vec![ + Box::new(SpawnParallelAgentsTool::new()), + Box::new(FixtureStepTool { + state: Arc::clone(&fixture_state), + }), + ]; + + let mut agent = Agent::builder() + .provider(Box::new(provider.clone())) + .tools(tools) + .memory(mem) + .tool_dispatcher(Box::new(NativeToolDispatcher)) + .workspace_dir(workspace_path) + .build() + .unwrap(); + + let response = agent + .turn("Run a long parallel delegation pass. parallel-fanout-e2e-canary") + .await + .unwrap_or_else(|err| { + panic!( + "agent turn failed: {err}\nseen payloads:\n{}", + provider.state.seen_payloads.lock().join("\n---\n") + ) + }); + + assert!( + response.contains(FINAL_CANARY), + "final orchestrator response should contain the synthesis canary: {response}" + ); + assert!( + response.contains(RESEARCH_DONE_CANARY) && response.contains(PLANNER_DONE_CANARY), + "final response should include both subagent completions: {response}" + ); + assert_eq!( + fixture_state.calls.load(Ordering::SeqCst), + 6, + "expected three nested tool calls per parallel subagent" + ); + assert!( + provider.max_active_subagent_calls() >= 2, + "expected overlapping subagent provider calls, max_active={}", + provider.max_active_subagent_calls() + ); + assert!( + provider.total_calls() >= 10, + "expected parent + subagent loop to hit the provider many times, total_calls={}", + provider.total_calls() + ); + + let history = agent.history(); + let mut saw_parallel_call = false; + let mut saw_parallel_result = false; + let mut iterations = Vec::new(); + + for message in history { + match message { + ConversationMessage::AssistantToolCalls { tool_calls, .. } => { + if tool_calls + .iter() + .any(|call| call.name == "spawn_parallel_agents") + { + saw_parallel_call = true; + } + } + ConversationMessage::ToolResults(results) => { + for result in results { + if !result.content.contains("\"parallel_agents\"") { + continue; + } + saw_parallel_result = true; + let payload: serde_json::Value = + serde_json::from_str(&result.content).expect("parallel tool result json"); + assert_eq!(payload["parallel_agents"]["succeeded"], 2); + assert_eq!(payload["parallel_agents"]["failed"], 0); + + let results = payload["parallel_agents"]["results"] + .as_array() + .expect("parallel results array"); + assert_eq!(results.len(), 2); + for item in results { + assert_eq!(item["success"], true); + iterations.push( + item["iterations"] + .as_u64() + .expect("parallel result iterations"), + ); + } + } + } + _ => {} + } + } + + assert!( + saw_parallel_call, + "parent history should record spawn_parallel_agents" + ); + assert!( + saw_parallel_result, + "parent history should record the parallel tool result" + ); + assert_eq!( + iterations, + vec![4, 4], + "each subagent should run three tool calls plus a final completion iteration" + ); +} diff --git a/src/openhuman/tools/impl/agent/spawn_subagent.rs b/src/openhuman/tools/impl/agent/spawn_subagent.rs index ae5aaa284..47a23691c 100644 --- a/src/openhuman/tools/impl/agent/spawn_subagent.rs +++ b/src/openhuman/tools/impl/agent/spawn_subagent.rs @@ -783,4 +783,83 @@ mod tests { // Should list at least one valid built-in. assert!(out.contains("code_executor") || out.contains("researcher")); } + + #[test] + fn classify_subagent_failure_reframes_upstream_provider_outages() { + let msg = SpawnSubagentTool::classify_subagent_failure( + "provider call failed: all providers/models failed: upstream unavailable", + ); + assert!(msg.contains("upstream inference unavailable")); + assert!(msg.contains("NOT a Composio/integration auth issue")); + } + + #[tokio::test] + async fn dedicated_thread_flag_is_rejected_explicitly() { + let tool = SpawnSubagentTool; + let result = tool + .execute(json!({ + "agent_id": "researcher", + "prompt": "find x", + "dedicated_thread": true, + })) + .await + .unwrap(); + assert!(result.is_error); + assert!(result.output().contains("temporarily disabled")); + } + + #[tokio::test] + async fn legacy_archetype_alias_is_accepted_for_lookup() { + let _ = AgentDefinitionRegistry::init_global_builtins(); + let tool = SpawnSubagentTool; + let result = tool + .execute(json!({ + "archetype": "totally_made_up", + "prompt": "x", + })) + .await + .unwrap(); + assert!(result.is_error); + assert!(result + .output() + .contains("unknown agent_id 'totally_made_up'")); + } + + #[tokio::test] + async fn integrations_agent_requires_toolkit_argument() { + let _ = AgentDefinitionRegistry::init_global_builtins(); + let tool = SpawnSubagentTool; + let result = tool + .execute(json!({ + "agent_id": "integrations_agent", + "prompt": "check gmail", + })) + .await + .unwrap(); + assert!(result.is_error); + let out = result.output(); + assert!(out.contains("`toolkit` argument is required")); + assert!(out.contains("currently-connected toolkits")); + } + + #[tokio::test] + async fn integrations_agent_rejects_toolkit_outside_allowlist() { + let _ = AgentDefinitionRegistry::init_global_builtins(); + let tool = SpawnSubagentTool; + let toolkit = "totally_not_a_real_toolkit_slug"; + let result = tool + .execute(json!({ + "agent_id": "integrations_agent", + "prompt": "check gmail", + "toolkit": toolkit, + })) + .await + .unwrap(); + assert!(result.is_error); + let out = result.output(); + assert!(out.contains(&format!( + "toolkit '{toolkit}' is not in the backend allowlist" + ))); + assert!(out.contains("Valid toolkits")); + } }