From dadd12311084c66bb1f9f5830e481eb7537a1e65 Mon Sep 17 00:00:00 2001 From: Steven Enamakel's Droid Date: Thu, 25 Jun 2026 21:37:59 -0700 Subject: [PATCH] feat(agent): expose async subagent control refs (#4138) --- src/openhuman/agent_orchestration/tools.rs | 3 + .../tools/spawn_async_subagent.rs | 159 ++++++-- .../agent_orchestration/tools/wait.rs | 365 ++++++++++++++++++ .../tools/wait_subagent.rs | 144 ++++++- src/openhuman/agent_registry/agents/loader.rs | 8 + .../agents/orchestrator/agent.toml | 2 + .../agents/orchestrator/prompt.md | 22 ++ src/openhuman/threads/schemas_tests.rs | 1 + src/openhuman/tokenjuice/savings.rs | 4 +- src/openhuman/tools/ops.rs | 2 + src/openhuman/tools/ops_tests.rs | 2 + tests/inference_agent_raw_coverage_e2e.rs | 5 +- tests/memory_threads_raw_coverage_e2e.rs | 3 +- 13 files changed, 678 insertions(+), 42 deletions(-) create mode 100644 src/openhuman/agent_orchestration/tools/wait.rs diff --git a/src/openhuman/agent_orchestration/tools.rs b/src/openhuman/agent_orchestration/tools.rs index 5f8feb356..fb008557a 100644 --- a/src/openhuman/agent_orchestration/tools.rs +++ b/src/openhuman/agent_orchestration/tools.rs @@ -25,6 +25,8 @@ mod steer_subagent; #[cfg(test)] #[path = "tools/tools_e2e_tests.rs"] mod tools_e2e_tests; +#[path = "tools/wait.rs"] +mod wait; #[path = "tools/wait_subagent.rs"] mod wait_subagent; #[path = "tools/worker_thread.rs"] @@ -45,4 +47,5 @@ pub use spawn_parallel_agents::SpawnParallelAgentsTool; pub use spawn_subagent::SpawnSubagentTool; pub use spawn_worker_thread::SpawnWorkerThreadTool; pub use steer_subagent::SteerSubagentTool; +pub use wait::{WaitLoopTool, WaitTool}; pub use wait_subagent::WaitSubagentTool; diff --git a/src/openhuman/agent_orchestration/tools/spawn_async_subagent.rs b/src/openhuman/agent_orchestration/tools/spawn_async_subagent.rs index 703be0be0..45a670b94 100644 --- a/src/openhuman/agent_orchestration/tools/spawn_async_subagent.rs +++ b/src/openhuman/agent_orchestration/tools/spawn_async_subagent.rs @@ -295,21 +295,19 @@ impl Tool for SpawnAsyncSubagentTool { definition.id, reuse_decision.as_str() ); - let payload = json!({ - "task_id": running_task_id, - "subagent_session_id": session.subagent_session_id, - "agent_id": definition.id, - "mode": "async", - "worker_thread_id": session.worker_thread_id, - "reused": true, - "reuse_decision": reuse_decision.as_str(), - }); - return Ok(ToolResult::success(format!( - "Continued reusable sub-agent `{}` (subagent_session_id `{}`, task_id `{}`). \ - It is already running and will pick up the new instruction at its next step.\n\n[async_subagent_ref]\n{}\n[/async_subagent_ref]", - payload["agent_id"].as_str().unwrap_or("subagent"), - payload["subagent_session_id"].as_str().unwrap_or(""), + let payload = async_subagent_ref_payload( running_task_id, + &session.subagent_session_id, + &definition.id, + session.worker_thread_id.as_deref(), + true, + reuse_decision.as_str(), + "running", + ); + return Ok(ToolResult::success(format!( + "Continued reusable async sub-agent `{}`. It is already running and will pick up the new instruction at its next step. \ + Use the structured reference below to send more input, wait, or perform a short timeout tick.\n\n[async_subagent_ref]\n{}\n[/async_subagent_ref]", + payload["agent_id"].as_str().unwrap_or("subagent"), serde_json::to_string(&payload).unwrap_or_else(|_| "{}".to_string()) ))); } @@ -641,15 +639,15 @@ impl Tool for SpawnAsyncSubagentTool { status_rx, ); - let payload = json!({ - "task_id": task_id, - "subagent_session_id": durable_session.subagent_session_id, - "agent_id": definition.id, - "mode": "async", - "worker_thread_id": worker_thread_id, - "reused": reusable.is_some(), - "reuse_decision": reuse_decision.as_str(), - }); + let payload = async_subagent_ref_payload( + &task_id, + &durable_session.subagent_session_id, + &definition.id, + worker_thread_id.as_deref(), + reusable.is_some(), + reuse_decision.as_str(), + "running", + ); let payload_json = match serde_json::to_string(&payload) { Ok(serialized) => { log::debug!( @@ -674,14 +672,94 @@ impl Tool for SpawnAsyncSubagentTool { } } +/// Format the user-facing acceptance text around a structured async sub-agent reference. fn format_async_subagent_accepted(agent_id: &str, payload_json: &str) -> String { format!( - "Accepted background sub-agent `{agent_id}`. Do not block on it before answering the user. \ - You may redirect it mid-run with `steer_subagent` and collect its result with `wait_subagent` \ - using the structured reference below.\n\n[async_subagent_ref]\n{payload_json}\n[/async_subagent_ref]" + "Accepted async sub-agent `{agent_id}`. Use the structured reference below to send more input, \ + wait for completion, or perform a short timeout tick to check status. If the user does not need \ + the result now, continue without blocking.\n\n[async_subagent_ref]\n{payload_json}\n[/async_subagent_ref]" ) } +/// Build the machine-readable reference the orchestrator uses to steer, wait, or poll a worker. +fn async_subagent_ref_payload( + task_id: &str, + subagent_session_id: &str, + agent_id: &str, + worker_thread_id: Option<&str>, + reused: bool, + reuse_decision: &str, + status: &str, +) -> serde_json::Value { + json!({ + "task_id": task_id, + "taskId": task_id, + "subagent_session_id": subagent_session_id, + "subagentSessionId": subagent_session_id, + "agent_id": agent_id, + "agentId": agent_id, + "mode": "async", + "status": status, + "worker_thread_id": worker_thread_id, + "workerThreadId": worker_thread_id, + "reused": reused, + "reuse_decision": reuse_decision, + "reuseDecision": reuse_decision, + "instructions": { + "send_message": { + "tool": "steer_subagent", + "description": "Send additional instructions or context to this running async sub-agent.", + "arguments": { + "subagent_session_id": subagent_session_id, + "message": "", + "mode": "steer" + } + }, + "wait": { + "tool": "wait_subagent", + "description": "Block until the async sub-agent finishes, up to the timeout.", + "arguments": { + "subagent_session_id": subagent_session_id, + "timeout_secs": 120 + } + }, + "timeout_tick": { + "tool": "wait_subagent", + "description": "Perform a short status tick without committing the parent to a long wait.", + "arguments": { + "subagent_session_id": subagent_session_id, + "timeout_secs": 1 + } + }, + "delayed_tick": { + "tool": "wait", + "description": "Trigger a delayed callback before checking this async sub-agent again.", + "arguments": { + "duration_secs": 30, + "message": format!("Check async sub-agent {agent_id} status with wait_subagent using subagent_session_id {subagent_session_id}.") + } + }, + "delayed_loop": { + "tool": "wait_loop", + "description": "Trigger repeatable delayed callbacks while this async sub-agent is still relevant.", + "arguments": { + "duration_secs": 30, + "message": format!("Check async sub-agent {agent_id} status with wait_subagent using subagent_session_id {subagent_session_id}."), + "loop_key": subagent_session_id, + "iteration": 1 + } + } + }, + "next_actions": [ + "call steer_subagent to send more input", + "call wait_subagent with timeout_secs to collect the result", + "call wait_subagent with timeout_secs=1 as a timeout tick/status check", + "call wait or wait_loop with the returned message to trigger a delayed status check", + "continue without waiting when the current user reply does not depend on the result" + ] + }) +} + fn add_background_contract(prompt: &str) -> String { format!( "[Background Contract]\n\ @@ -765,12 +843,39 @@ mod tests { .next() .expect("prose before structured reference"); - assert!(prose.contains("Accepted background sub-agent `archivist`")); + assert!(prose.contains("Accepted async sub-agent `archivist`")); assert!(!prose.contains("sub-internal-123")); assert!(message.contains("[async_subagent_ref]")); assert!(message.contains("sub-internal-123")); } + #[test] + fn async_reference_payload_includes_agent_id_and_control_instructions() { + let payload = async_subagent_ref_payload( + "sub-123", + "subsess-456", + "researcher", + Some("thread-worker"), + false, + "created", + "running", + ); + + assert_eq!(payload["agent_id"], "researcher"); + assert_eq!(payload["agentId"], "researcher"); + assert_eq!(payload["instructions"]["wait"]["tool"], "wait_subagent"); + assert_eq!( + payload["instructions"]["timeout_tick"]["arguments"]["timeout_secs"], + 1 + ); + assert_eq!(payload["instructions"]["delayed_tick"]["tool"], "wait"); + assert_eq!(payload["instructions"]["delayed_loop"]["tool"], "wait_loop"); + assert_eq!( + payload["instructions"]["send_message"]["tool"], + "steer_subagent" + ); + } + #[test] fn durable_task_key_defaults_to_prompt_not_display_title() { let args = json!({ diff --git a/src/openhuman/agent_orchestration/tools/wait.rs b/src/openhuman/agent_orchestration/tools/wait.rs new file mode 100644 index 000000000..ad1d5dd79 --- /dev/null +++ b/src/openhuman/agent_orchestration/tools/wait.rs @@ -0,0 +1,365 @@ +//! Tool: `wait` / `wait_loop` - delayed callback ticks for the orchestrator. +//! +//! These tools intentionally do not own scheduling state. They sleep for a +//! bounded duration, then return the caller-provided message as a tool result so +//! the orchestrator can decide whether to act, poll async sub-agents, or call the +//! loop variant again. + +use std::time::Duration; + +use crate::openhuman::tools::traits::{PermissionLevel, Tool, ToolResult, ToolTimeout}; +use async_trait::async_trait; +use serde_json::json; + +const DEFAULT_DURATION_SECS: u64 = 5; +const MAX_DURATION_SECS: u64 = 600; +const MILLIS_PER_SEC: u64 = 1_000; + +/// One-shot delayed callback tool for the orchestrator. +pub struct WaitTool; +/// Repeatable delayed callback tool for orchestrator-controlled polling loops. +pub struct WaitLoopTool; + +impl WaitTool { + pub fn new() -> Self { + Self + } +} + +impl Default for WaitTool { + fn default() -> Self { + Self::new() + } +} + +impl WaitLoopTool { + pub fn new() -> Self { + Self + } +} + +impl Default for WaitLoopTool { + fn default() -> Self { + Self::new() + } +} + +#[async_trait] +impl Tool for WaitTool { + fn name(&self) -> &str { + "wait" + } + + fn description(&self) -> &str { + "Wait for a bounded duration, then return the provided callback message \ + to the orchestrator. Use this to create a delayed tick before checking \ + async work or retrying a condition." + } + + fn parameters_schema(&self) -> serde_json::Value { + wait_schema(false) + } + + fn permission_level(&self) -> PermissionLevel { + PermissionLevel::Execute + } + + fn timeout_policy(&self, args: &serde_json::Value) -> ToolTimeout { + timeout_policy_for_wait(args) + } + + async fn execute(&self, args: serde_json::Value) -> anyhow::Result { + execute_wait(args, false).await + } +} + +#[async_trait] +impl Tool for WaitLoopTool { + fn name(&self) -> &str { + "wait_loop" + } + + fn description(&self) -> &str { + "Wait for a bounded duration, then return the same callback message plus \ + a ready-to-call wait_loop instruction. Use this for deliberate polling \ + loops where the orchestrator decides after each tick whether to repeat." + } + + fn parameters_schema(&self) -> serde_json::Value { + wait_schema(true) + } + + fn permission_level(&self) -> PermissionLevel { + PermissionLevel::Execute + } + + fn timeout_policy(&self, args: &serde_json::Value) -> ToolTimeout { + timeout_policy_for_wait(args) + } + + async fn execute(&self, args: serde_json::Value) -> anyhow::Result { + execute_wait(args, true).await + } +} + +/// Sleep for the requested duration, then return the callback tick payload. +async fn execute_wait(args: serde_json::Value, loop_mode: bool) -> anyhow::Result { + let request = match parse_wait_request(&args) { + Ok(request) => request, + Err(message) => return Ok(ToolResult::error(message)), + }; + let tool_name = if loop_mode { "wait_loop" } else { "wait" }; + log::info!( + "[wait] tool={} duration_ms={} iteration={} loop_key={} message_chars={}", + tool_name, + request.duration_ms, + request.iteration, + request.loop_key.as_deref().unwrap_or("none"), + request.message.chars().count() + ); + + tokio::time::sleep(Duration::from_millis(request.duration_ms)).await; + + log::debug!( + "[wait] elapsed tool={} duration_ms={} iteration={} loop_key={}", + tool_name, + request.duration_ms, + request.iteration, + request.loop_key.as_deref().unwrap_or("none") + ); + + Ok(ToolResult::success(format_wait_tick(&request, loop_mode))) +} + +/// Build the JSON schema shared by `wait` and `wait_loop`. +fn wait_schema(loop_mode: bool) -> serde_json::Value { + let mut properties = serde_json::Map::new(); + properties.insert( + "message".to_string(), + json!({ + "type": "string", + "description": "Callback message to return to the orchestrator after the wait elapses." + }), + ); + properties.insert( + "duration_secs".to_string(), + json!({ + "type": "integer", + "minimum": 1, + "maximum": MAX_DURATION_SECS, + "description": "Seconds to wait before returning the callback. Default 5. Ignored when duration_ms is supplied." + }), + ); + properties.insert( + "duration_ms".to_string(), + json!({ + "type": "integer", + "minimum": 1, + "maximum": MAX_DURATION_SECS * MILLIS_PER_SEC, + "description": "Milliseconds to wait before returning the callback. Use only for short test-sized waits." + }), + ); + + if loop_mode { + properties.insert( + "loop_key".to_string(), + json!({ + "type": "string", + "description": "Optional caller-defined key for the polling loop." + }), + ); + properties.insert( + "iteration".to_string(), + json!({ + "type": "integer", + "minimum": 1, + "description": "Current loop iteration. Defaults to 1; the returned instruction increments it." + }), + ); + } + + json!({ + "type": "object", + "required": ["message"], + "properties": serde_json::Value::Object(properties) + }) +} + +/// Give the harness a deadline that covers the requested wait plus grace. +fn timeout_policy_for_wait(args: &serde_json::Value) -> ToolTimeout { + let duration_ms = duration_ms_from_args(args).unwrap_or(DEFAULT_DURATION_SECS * MILLIS_PER_SEC); + let rounded_secs = duration_ms.saturating_add(MILLIS_PER_SEC - 1) / MILLIS_PER_SEC; + ToolTimeout::Secs(rounded_secs.saturating_add(1)) +} + +/// Extract and clamp the caller's requested wait duration in milliseconds. +fn duration_ms_from_args(args: &serde_json::Value) -> Option { + if let Some(duration_ms) = args.get("duration_ms").and_then(|v| v.as_u64()) { + return Some(duration_ms.clamp(1, MAX_DURATION_SECS * MILLIS_PER_SEC)); + } + args.get("duration_secs") + .and_then(|v| v.as_u64()) + .map(|secs| secs.clamp(1, MAX_DURATION_SECS) * MILLIS_PER_SEC) +} + +/// Validate tool arguments and normalize them into an internal request. +fn parse_wait_request(args: &serde_json::Value) -> Result { + let message = args + .get("message") + .and_then(|v| v.as_str()) + .unwrap_or("") + .trim() + .to_string(); + if message.is_empty() { + return Err("wait: `message` is required".to_string()); + } + + let duration_ms = duration_ms_from_args(args).unwrap_or(DEFAULT_DURATION_SECS * MILLIS_PER_SEC); + let iteration = args + .get("iteration") + .and_then(|v| v.as_u64()) + .unwrap_or(1) + .max(1); + let loop_key = args + .get("loop_key") + .and_then(|v| v.as_str()) + .map(str::trim) + .filter(|s| !s.is_empty()) + .map(str::to_string); + + Ok(WaitRequest { + message, + duration_ms, + iteration, + loop_key, + }) +} + +/// Render the delayed callback as prose plus a machine-readable `[wait_tick]`. +fn format_wait_tick(request: &WaitRequest, loop_mode: bool) -> String { + let seconds = request.duration_ms as f64 / MILLIS_PER_SEC as f64; + let loop_instruction = loop_mode.then(|| { + json!({ + "tool": "wait_loop", + "arguments": { + "message": request.message, + "duration_ms": request.duration_ms, + "loop_key": request.loop_key, + "iteration": request.iteration + 1 + } + }) + }); + let payload = json!({ + "status": "elapsed", + "message": request.message, + "duration_ms": request.duration_ms, + "durationMs": request.duration_ms, + "duration_secs": seconds, + "durationSecs": seconds, + "loop": loop_mode, + "loop_key": request.loop_key, + "loopKey": request.loop_key, + "iteration": request.iteration, + "instructions": { + "callback_message": request.message, + "repeat": loop_instruction + } + }); + + let prefix = if loop_mode { + format!( + "Loop tick {} elapsed after {}ms.", + request.iteration, request.duration_ms + ) + } else { + format!("Wait elapsed after {}ms.", request.duration_ms) + }; + + format!( + "{prefix}\n\n[wait_tick]\n{}\n[/wait_tick]\n\n{}", + serde_json::to_string(&payload).unwrap_or_else(|_| "{}".to_string()), + request.message + ) +} + +#[derive(Debug)] +/// Normalized wait arguments used by both wait tools. +struct WaitRequest { + message: String, + duration_ms: u64, + iteration: u64, + loop_key: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn wait_schema_requires_message() { + let schema = WaitTool::new().parameters_schema(); + let required = schema + .get("required") + .and_then(|v| v.as_array()) + .expect("required list"); + assert!(required.iter().any(|v| v.as_str() == Some("message"))); + assert!(schema["properties"].get("duration_secs").is_some()); + assert!(schema["properties"].get("duration_ms").is_some()); + } + + #[test] + fn wait_loop_schema_includes_loop_controls() { + let schema = WaitLoopTool::new().parameters_schema(); + assert!(schema["properties"].get("loop_key").is_some()); + assert!(schema["properties"].get("iteration").is_some()); + } + + #[test] + fn parse_wait_request_clamps_duration_and_iteration() { + let request = parse_wait_request(&json!({ + "message": "check subagents", + "duration_secs": 9999, + "iteration": 0 + })) + .unwrap(); + assert_eq!(request.duration_ms, MAX_DURATION_SECS * MILLIS_PER_SEC); + assert_eq!(request.iteration, 1); + } + + #[test] + fn missing_message_is_rejected() { + let err = parse_wait_request(&json!({ "duration_secs": 1 })).unwrap_err(); + assert!(err.contains("message")); + } + + #[test] + fn wait_loop_tick_repeats_same_message() { + let request = parse_wait_request(&json!({ + "message": "poll async workers", + "duration_ms": 10, + "loop_key": "workers", + "iteration": 2 + })) + .unwrap(); + let output = format_wait_tick(&request, true); + + assert!(output.contains("Loop tick 2 elapsed")); + assert!(output.contains("\"tool\":\"wait_loop\"")); + assert!(output.contains("\"message\":\"poll async workers\"")); + assert!(output.contains("\"iteration\":3")); + } + + #[tokio::test] + async fn wait_execute_returns_callback_message() { + let res = WaitTool::new() + .execute(json!({ + "message": "time to check", + "duration_ms": 1 + })) + .await + .unwrap(); + assert!(!res.is_error); + assert!(res.output().contains("[wait_tick]")); + assert!(res.output().contains("time to check")); + } +} diff --git a/src/openhuman/agent_orchestration/tools/wait_subagent.rs b/src/openhuman/agent_orchestration/tools/wait_subagent.rs index 4054e1499..9bcd7c15f 100644 --- a/src/openhuman/agent_orchestration/tools/wait_subagent.rs +++ b/src/openhuman/agent_orchestration/tools/wait_subagent.rs @@ -151,8 +151,18 @@ impl Tool for WaitSubagentTool { resolved_task_id, iterations ); + let status = wait_status_payload( + resume_ref.as_ref(), + &resolved_task_id, + "completed", + Some(iterations), + None, + "synthesize the sub-agent output into the parent response", + ); Ok(ToolResult::success(format!( - "Sub-agent completed in {iterations} iteration(s):\n\n{output}" + "Sub-agent `{}` completed in {iterations} iteration(s).\n\n[subagent_wait_result]\n{}\n[/subagent_wait_result]\n\n{output}", + status["agent_id"].as_str().unwrap_or("subagent"), + serde_json::to_string(&status).unwrap_or_else(|_| "{}".to_string()) ))) } Ok(WaitOutcome::Terminal(SubagentStatus::AwaitingUser { question })) => { @@ -161,9 +171,19 @@ impl Tool for WaitSubagentTool { resolved_task_id, question.chars().count() ); + let status = wait_status_payload( + resume_ref.as_ref(), + &resolved_task_id, + "awaiting_user", + None, + Some(&question), + "ask the user for the missing information, then call continue_subagent", + ); let mut message = format!( - "Sub-agent paused for clarification and did not finish: {question}\n\n\ - It cannot proceed unattended. Resume it with continue_subagent once you have an answer." + "Sub-agent `{}` paused for clarification and did not finish: {question}\n\n\ + It cannot proceed unattended. Resume it with continue_subagent once you have an answer.\n\n[subagent_wait_result]\n{}\n[/subagent_wait_result]", + status["agent_id"].as_str().unwrap_or("subagent"), + serde_json::to_string(&status).unwrap_or_else(|_| "{}".to_string()) ); if let Some(reference) = resume_ref { message.push_str("\n\n[subagent_resume_ref]\n"); @@ -191,7 +211,19 @@ impl Tool for WaitSubagentTool { resolved_task_id, error ); - Ok(ToolResult::error(format!("Sub-agent failed: {error}"))) + let status = wait_status_payload( + resume_ref.as_ref(), + &resolved_task_id, + "failed", + None, + Some(&error), + "report the failure or retry with a corrected instruction", + ); + Ok(ToolResult::error(format!( + "Sub-agent `{}` failed: {error}\n\n[subagent_wait_result]\n{}\n[/subagent_wait_result]", + status["agent_id"].as_str().unwrap_or("subagent"), + serde_json::to_string(&status).unwrap_or_else(|_| "{}".to_string()) + ))) } // `Running` is never terminal; treat defensively as a timeout-style result. Ok(WaitOutcome::Terminal(SubagentStatus::Running)) => { @@ -200,9 +232,10 @@ impl Tool for WaitSubagentTool { resolved_task_id, timeout_secs ); - Ok(ToolResult::success(format!( - "Sub-agent is still running after {timeout_secs}s. \ - Continue with other work and call wait_subagent again later, or steer_subagent to redirect it." + Ok(ToolResult::success(format_running_wait_message( + resume_ref.as_ref(), + &resolved_task_id, + timeout_secs, ))) } Ok(WaitOutcome::TimedOut(_)) => { @@ -211,9 +244,10 @@ impl Tool for WaitSubagentTool { resolved_task_id, timeout_secs ); - Ok(ToolResult::success(format!( - "Sub-agent is still running after {timeout_secs}s. \ - Continue with other work and call wait_subagent again later, or steer_subagent to redirect it." + Ok(ToolResult::success(format_running_wait_message( + resume_ref.as_ref(), + &resolved_task_id, + timeout_secs, ))) } Err(WaitError::Unknown) => { @@ -239,6 +273,80 @@ impl Tool for WaitSubagentTool { } } +/// Render a timeout/running wait response with a structured status payload. +fn format_running_wait_message( + reference: Option<&running_subagents::SubagentResumeRef>, + task_id: &str, + timeout_secs: u64, +) -> String { + let status = wait_status_payload( + reference, + task_id, + "running", + None, + None, + "continue other work, call wait_subagent again, or call steer_subagent to send more input", + ); + format!( + "Sub-agent `{}` is still running after {timeout_secs}s.\n\n[subagent_wait_result]\n{}\n[/subagent_wait_result]\n\nContinue with other work and call wait_subagent again later, or steer_subagent to redirect it.", + status["agent_id"].as_str().unwrap_or("subagent"), + serde_json::to_string(&status).unwrap_or_else(|_| "{}".to_string()) + ) +} + +/// Build the machine-readable wait status block returned to the orchestrator. +fn wait_status_payload( + reference: Option<&running_subagents::SubagentResumeRef>, + task_id: &str, + status: &str, + iterations: Option, + detail: Option<&str>, + next_action: &str, +) -> serde_json::Value { + let agent_id = reference.map(|r| r.agent_id.as_str()).unwrap_or("unknown"); + let subagent_session_id = reference.and_then(|r| r.subagent_session_id.as_deref()); + json!({ + "task_id": task_id, + "taskId": task_id, + "subagent_session_id": subagent_session_id, + "subagentSessionId": subagent_session_id, + "agent_id": agent_id, + "agentId": agent_id, + "status": status, + "iterations": iterations, + "detail": detail, + "next_action": next_action, + "nextAction": next_action, + "instructions": { + "send_message": { + "tool": "steer_subagent", + "arguments": { + "subagent_session_id": subagent_session_id, + "task_id": task_id, + "message": "", + "mode": "steer" + } + }, + "wait": { + "tool": "wait_subagent", + "arguments": { + "subagent_session_id": subagent_session_id, + "task_id": task_id, + "timeout_secs": DEFAULT_TIMEOUT_SECS + } + }, + "timeout_tick": { + "tool": "wait_subagent", + "arguments": { + "subagent_session_id": subagent_session_id, + "task_id": task_id, + "timeout_secs": 1 + } + } + } + }) +} + #[cfg(test)] mod tests { use super::*; @@ -269,4 +377,20 @@ mod tests { assert!(res.is_error); assert!(res.output().contains("outside of an agent turn")); } + + #[test] + fn running_wait_message_includes_agent_id_and_tick_instruction() { + let reference = running_subagents::SubagentResumeRef { + task_id: "sub-1".into(), + agent_id: "researcher".into(), + subagent_session_id: Some("subsess-1".into()), + }; + let message = format_running_wait_message(Some(&reference), "sub-1", 1); + + assert!(message.contains("Sub-agent `researcher` is still running")); + assert!(message.contains("[subagent_wait_result]")); + assert!(message.contains("\"agentId\":\"researcher\"")); + assert!(message.contains("\"timeout_tick\"")); + assert!(message.contains("\"timeout_secs\":1")); + } } diff --git a/src/openhuman/agent_registry/agents/loader.rs b/src/openhuman/agent_registry/agents/loader.rs index 8a00afab6..6ffd12862 100644 --- a/src/openhuman/agent_registry/agents/loader.rs +++ b/src/openhuman/agent_registry/agents/loader.rs @@ -573,6 +573,14 @@ mod tests { tools.iter().any(|t| t == "spawn_async_subagent"), "orchestrator must have spawn_async_subagent for sparse background work" ); + assert!( + tools.iter().any(|t| t == "wait"), + "orchestrator must have wait for delayed callback ticks" + ); + assert!( + tools.iter().any(|t| t == "wait_loop"), + "orchestrator must have wait_loop for deliberate polling loops" + ); assert!( !tools.iter().any(|t| t == "spawn_subagent"), "spawn_subagent must not appear — removed in #1141" diff --git a/src/openhuman/agent_registry/agents/orchestrator/agent.toml b/src/openhuman/agent_registry/agents/orchestrator/agent.toml index 0304956bc..0260b3241 100644 --- a/src/openhuman/agent_registry/agents/orchestrator/agent.toml +++ b/src/openhuman/agent_registry/agents/orchestrator/agent.toml @@ -191,6 +191,8 @@ named = [ "spawn_worker_thread", "spawn_async_subagent", "spawn_parallel_agents", + "wait", + "wait_loop", "composio_list_connections", # Inline OAuth connect card (#3993). Free-form `toolkit` slug — the # orchestrator calls this directly to connect a service the user names, diff --git a/src/openhuman/agent_registry/agents/orchestrator/prompt.md b/src/openhuman/agent_registry/agents/orchestrator/prompt.md index b6f6fde78..45dcd8e2d 100644 --- a/src/openhuman/agent_registry/agents/orchestrator/prompt.md +++ b/src/openhuman/agent_registry/agents/orchestrator/prompt.md @@ -98,6 +98,28 @@ external-service writes, financial/market actions, scheduling, desktop control, task that may need clarification. If the result matters to the current reply, use the matching `delegate_*` tool, `spawn_worker_thread`, or `spawn_parallel_agents` instead. +`spawn_async_subagent` returns an `[async_subagent_ref]` block with both `agent_id` +and `agentId`, plus concrete control instructions: + +- To send more input, call `steer_subagent` using the returned + `subagent_session_id` (preferred) or `task_id`. +- To collect the result, call `wait_subagent` using that reference. Use a longer + `timeout_secs` only when the current response depends on the result. +- To perform a non-blocking status tick, call `wait_subagent` with + `timeout_secs: 1`. If it returns `status: "running"`, continue other work or + answer without waiting unless the user specifically needs that result now. +- To delay a status check, call `wait` with a short `duration_secs` and a + concrete `message` such as "check with wait_subagent". + When it returns, treat the message as your callback prompt. +- To keep polling, call `wait_loop` with the same message. Each tick returns a + ready-to-call `wait_loop` instruction with the same message and incremented + iteration; repeat only while the task still needs polling. + +When you spawn multiple async sub-agents, treat them as parallel workers: keep +their refs separate by `subagent_session_id` or `task_id` (`agentId` is only the +worker type), tick or wait on each independently, and synthesise only completed +outputs. Never fabricate a result for a worker that is still running or failed. + ## Connecting external services When the user asks to connect a service (Gmail, Notion, WhatsApp, Calendar, Drive, etc.) or a sub-agent reports `Connection error, try to authenticate`: diff --git a/src/openhuman/threads/schemas_tests.rs b/src/openhuman/threads/schemas_tests.rs index e06355c92..2aeb8c07d 100644 --- a/src/openhuman/threads/schemas_tests.rs +++ b/src/openhuman/threads/schemas_tests.rs @@ -19,6 +19,7 @@ const ALL_FUNCTIONS: &[&str] = &[ "turn_state_clear", "task_board_get", "task_board_put", + "token_usage", ]; #[test] diff --git a/src/openhuman/tokenjuice/savings.rs b/src/openhuman/tokenjuice/savings.rs index db891e513..add57c9e7 100644 --- a/src/openhuman/tokenjuice/savings.rs +++ b/src/openhuman/tokenjuice/savings.rs @@ -199,9 +199,9 @@ mod tests { #[test] fn cost_uses_input_price() { - // agentic-v1 input is $3/Mtok → 1M saved tokens ≈ $3. + // agentic-v1 input pricing is used for saved-token cost estimates. let c = cost_saved_usd("agentic-v1", 1_000_000); - assert!((c - 3.0).abs() < 1e-6, "got {c}"); + assert!((c - 0.435).abs() < 1e-6, "got {c}"); } #[test] diff --git a/src/openhuman/tools/ops.rs b/src/openhuman/tools/ops.rs index 7efa190fc..f46b5d3ba 100644 --- a/src/openhuman/tools/ops.rs +++ b/src/openhuman/tools/ops.rs @@ -176,6 +176,8 @@ pub fn all_tools_with_runtime( // durable `subagent_session_id` (preferred) or transient `task_id`. Box::new(ListSubagentsTool::new()), Box::new(SteerSubagentTool::new()), + Box::new(WaitTool::new()), + Box::new(WaitLoopTool::new()), Box::new(WaitSubagentTool::new()), Box::new(CloseSubagentTool::new()), Box::new(ContinueSubagentTool::new()), diff --git a/src/openhuman/tools/ops_tests.rs b/src/openhuman/tools/ops_tests.rs index 5358608d6..c1c6a8b65 100644 --- a/src/openhuman/tools/ops_tests.rs +++ b/src/openhuman/tools/ops_tests.rs @@ -437,6 +437,8 @@ fn all_tools_default_registry_contains_expected_baseline_surface() { "spawn_subagent", "spawn_async_subagent", "spawn_parallel_agents", + "wait", + "wait_loop", "todo", "plan_exit", "current_time", diff --git a/tests/inference_agent_raw_coverage_e2e.rs b/tests/inference_agent_raw_coverage_e2e.rs index ac667dbe4..20128744f 100644 --- a/tests/inference_agent_raw_coverage_e2e.rs +++ b/tests/inference_agent_raw_coverage_e2e.rs @@ -2795,7 +2795,7 @@ async fn agent_runtime_policy_cost_and_triage_helpers_cover_public_edges() { ); let estimated = openhuman_core::openhuman::agent::cost::estimate_call_cost_usd("agentic-v1", &usage); - assert!((estimated - 18.3).abs() < 1e-6, "got {estimated}"); + assert!((estimated - 1.308625).abs() < 1e-6, "got {estimated}"); let charged = UsageInfo { charged_amount_usd: 0.42, ..usage.clone() @@ -2812,7 +2812,7 @@ async fn agent_runtime_policy_cost_and_triage_helpers_cover_public_edges() { assert_eq!(turn_cost.cached_input_tokens, 2_000_000); assert_eq!(turn_cost.charged_usd, 0.42); assert_eq!(turn_cost.call_count, 2); - assert!(turn_cost.total_usd() > 18.7); + assert!((turn_cost.total_usd() - 1.728625).abs() < 1e-6); let composio = TriggerEnvelope::from_composio( "gmail", @@ -4701,6 +4701,7 @@ async fn agent_subagent_public_types_cover_task_local_and_error_display_paths() mode: SubagentMode::Typed, status: SubagentRunStatus::Completed, final_history: Vec::new(), + usage: Default::default(), }; assert_eq!(outcome.mode.as_str(), "typed"); assert_eq!(outcome.elapsed.as_millis(), 12); diff --git a/tests/memory_threads_raw_coverage_e2e.rs b/tests/memory_threads_raw_coverage_e2e.rs index bfa6cdbeb..dc71def4c 100644 --- a/tests/memory_threads_raw_coverage_e2e.rs +++ b/tests/memory_threads_raw_coverage_e2e.rs @@ -771,7 +771,7 @@ async fn memory_thread_tree_and_sync_controller_schemas_execute_public_handlers( let thread_schemas = all_threads_controller_schemas(); let thread_controllers = all_threads_registered_controllers(); - assert_eq!(thread_schemas.len(), 16); + assert_eq!(thread_schemas.len(), 17); assert_eq!(thread_schemas.len(), thread_controllers.len()); assert_eq!( openhuman_core::openhuman::threads::schemas::schemas("missing").function, @@ -794,6 +794,7 @@ async fn memory_thread_tree_and_sync_controller_schemas_execute_public_handlers( "turn_state_clear", "task_board_get", "task_board_put", + "token_usage", ] { assert!(thread_schemas .iter()