Files
openhuman/tests/raw_coverage/agent_tool_loop_raw_coverage_e2e.rs
T

649 lines
20 KiB
Rust

use async_trait::async_trait;
use openhuman_core::core::event_bus::{init_global, request_native_global, DEFAULT_CAPACITY};
use openhuman_core::openhuman::agent::bus::{
register_agent_handlers, AgentTurnRequest, AgentTurnResponse, AGENT_RUN_TURN_METHOD,
};
use openhuman_core::openhuman::agent::debug::{dump_agent_prompt, DumpPromptOptions};
use openhuman_core::openhuman::agent::dispatcher::XmlToolDispatcher;
use openhuman_core::openhuman::agent::{Agent, AgentBuilder};
use openhuman_core::openhuman::config::{AgentConfig, MultimodalConfig, MultimodalFileConfig};
use openhuman_core::openhuman::context::prompt::LearnedContextData;
use openhuman_core::openhuman::agent::messages::ChatMessage;
use openhuman_core::openhuman::memory::{
Memory, MemoryCategory, MemoryEntry, NamespaceSummary, RecallOpts,
};
use openhuman_core::openhuman::tools::{PermissionLevel, Tool, ToolContent, ToolResult, ToolScope};
use serde_json::json;
use std::collections::{HashSet, VecDeque};
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use tinyagents::harness::message::{AssistantMessage, ContentBlock, Message, MessageDelta};
use tinyagents::harness::model::{
ChatModel, ModelProfile, ModelRequest, ModelResponse, ModelStream, ModelStreamItem,
};
use tinyagents::harness::tool::{ToolCall, ToolDelta};
use tinyagents::harness::usage::Usage;
#[derive(Clone, Debug)]
struct CapturedTurn {
messages: Vec<Message>,
tool_names: Vec<String>,
}
struct ScriptedModel {
responses: Mutex<VecDeque<anyhow::Result<ModelResponse>>>,
turns: Mutex<Vec<CapturedTurn>>,
profile: ModelProfile,
stream_events: Vec<ModelStreamItem>,
always_fail: Option<String>,
}
impl Default for ScriptedModel {
fn default() -> Self {
Self {
responses: Mutex::new(VecDeque::new()),
turns: Mutex::new(Vec::new()),
profile: ModelProfile::default(),
stream_events: Vec::new(),
always_fail: None,
}
}
}
impl ScriptedModel {
fn new(responses: Vec<ModelResponse>) -> Arc<Self> {
Arc::new(Self {
responses: Mutex::new(responses.into_iter().map(Ok).collect()),
..Self::default()
})
}
fn failing(message: &str) -> Arc<Self> {
Arc::new(Self {
always_fail: Some(message.to_string()),
..Self::default()
})
}
fn turns(&self) -> Vec<CapturedTurn> {
self.turns.lock().unwrap().clone()
}
}
#[async_trait]
impl ChatModel<()> for ScriptedModel {
fn profile(&self) -> Option<&ModelProfile> {
Some(&self.profile)
}
async fn invoke(
&self,
_state: &(),
request: ModelRequest,
) -> tinyagents::Result<ModelResponse> {
self.capture(request);
self.pop_response()
}
async fn stream(&self, _state: &(), request: ModelRequest) -> tinyagents::Result<ModelStream> {
self.capture(request);
let response = self.pop_response()?;
let mut items = vec![ModelStreamItem::Started];
items.extend(self.stream_events.iter().cloned());
items.push(ModelStreamItem::Completed(response));
Ok(Box::pin(futures::stream::iter(items)))
}
}
impl ScriptedModel {
fn capture(&self, request: ModelRequest) {
self.turns.lock().unwrap().push(CapturedTurn {
messages: request.messages,
tool_names: request.tools.iter().map(|tool| tool.name.clone()).collect(),
});
}
fn pop_response(&self) -> tinyagents::Result<ModelResponse> {
if let Some(message) = &self.always_fail {
return Err(tinyagents::TinyAgentsError::Model(message.clone()));
}
self.responses
.lock()
.unwrap()
.pop_front()
.unwrap_or_else(|| Ok(ModelResponse::assistant("")))
.map_err(|error| tinyagents::TinyAgentsError::Model(error.to_string()))
}
}
struct StaticTool {
name: &'static str,
output: &'static str,
is_error: bool,
scope: ToolScope,
permission: PermissionLevel,
cap: Option<usize>,
}
impl StaticTool {
fn ok(name: &'static str, output: &'static str) -> Box<dyn Tool> {
Box::new(Self {
name,
output,
is_error: false,
scope: ToolScope::All,
permission: PermissionLevel::ReadOnly,
cap: None,
})
}
fn err(name: &'static str, output: &'static str) -> Box<dyn Tool> {
Box::new(Self {
name,
output,
is_error: true,
scope: ToolScope::All,
permission: PermissionLevel::ReadOnly,
cap: None,
})
}
fn cli_only(name: &'static str) -> Box<dyn Tool> {
Box::new(Self {
name,
output: "cli-only-output",
is_error: false,
scope: ToolScope::CliRpcOnly,
permission: PermissionLevel::ReadOnly,
cap: None,
})
}
fn capped(name: &'static str, output: &'static str, cap: usize) -> Box<dyn Tool> {
Box::new(Self {
name,
output,
is_error: false,
scope: ToolScope::All,
permission: PermissionLevel::ReadOnly,
cap: Some(cap),
})
}
}
#[async_trait]
impl Tool for StaticTool {
fn name(&self) -> &str {
self.name
}
fn description(&self) -> &str {
"round15 deterministic test tool"
}
fn parameters_schema(&self) -> serde_json::Value {
json!({
"type": "object",
"properties": {
"value": { "type": "string" }
}
})
}
async fn execute(&self, args: serde_json::Value) -> anyhow::Result<ToolResult> {
let suffix = args
.get("value")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
let body = if suffix.is_empty() {
self.output.to_string()
} else {
format!("{}:{suffix}", self.output)
};
Ok(ToolResult {
content: vec![ToolContent::Text { text: body }],
is_error: self.is_error,
markdown_formatted: None,
})
}
fn permission_level(&self) -> PermissionLevel {
self.permission
}
fn scope(&self) -> ToolScope {
self.scope
}
fn max_result_size_chars(&self) -> Option<usize> {
self.cap
}
}
#[derive(Default)]
struct NoopMemory {
entries: Mutex<Vec<MemoryEntry>>,
}
#[async_trait]
impl Memory for NoopMemory {
fn name(&self) -> &str {
"round15-noop"
}
async fn store(
&self,
namespace: &str,
key: &str,
content: &str,
category: MemoryCategory,
session_id: Option<&str>,
) -> anyhow::Result<()> {
self.entries.lock().unwrap().push(MemoryEntry {
id: format!("{namespace}:{key}"),
key: key.to_string(),
content: content.to_string(),
namespace: Some(namespace.to_string()),
category,
timestamp: "2026-05-29T00:00:00Z".to_string(),
session_id: session_id.map(str::to_string),
score: Some(1.0),
taint: Default::default(),
});
Ok(())
}
async fn recall(
&self,
_query: &str,
limit: usize,
_opts: RecallOpts<'_>,
) -> anyhow::Result<Vec<MemoryEntry>> {
Ok(self
.entries
.lock()
.unwrap()
.iter()
.take(limit)
.cloned()
.collect())
}
async fn get(&self, namespace: &str, key: &str) -> anyhow::Result<Option<MemoryEntry>> {
Ok(self
.entries
.lock()
.unwrap()
.iter()
.find(|entry| entry.namespace.as_deref() == Some(namespace) && entry.key == key)
.cloned())
}
async fn list(
&self,
namespace: Option<&str>,
category: Option<&MemoryCategory>,
_session_id: Option<&str>,
) -> anyhow::Result<Vec<MemoryEntry>> {
Ok(self
.entries
.lock()
.unwrap()
.iter()
.filter(|entry| namespace.is_none_or(|ns| entry.namespace.as_deref() == Some(ns)))
.filter(|entry| category.is_none_or(|cat| &entry.category == cat))
.cloned()
.collect())
}
async fn forget(&self, namespace: &str, key: &str) -> anyhow::Result<bool> {
let mut entries = self.entries.lock().unwrap();
let before = entries.len();
entries
.retain(|entry| !(entry.namespace.as_deref() == Some(namespace) && entry.key == key));
Ok(entries.len() != before)
}
async fn namespace_summaries(&self) -> anyhow::Result<Vec<NamespaceSummary>> {
Ok(vec![NamespaceSummary {
namespace: "round15".to_string(),
count: self.entries.lock().unwrap().len(),
last_updated: Some("2026-05-29T00:00:00Z".to_string()),
}])
}
async fn count(&self) -> anyhow::Result<usize> {
Ok(self.entries.lock().unwrap().len())
}
async fn health_check(&self) -> bool {
true
}
}
fn text_response(text: &str) -> ModelResponse {
let mut usage = Usage::new(11, 7);
usage.cache_read_tokens = 3;
ModelResponse::assistant(text).with_usage(usage)
}
fn tool_response(name: &str, arguments: serde_json::Value) -> ModelResponse {
let mut usage = Usage::new(13, 5);
usage.cache_read_tokens = 2;
ModelResponse {
message: AssistantMessage {
id: None,
content: vec![
ContentBlock::Text("using native tool".to_string()),
ContentBlock::thinking("private scratchpad"),
],
tool_calls: vec![ToolCall::new(format!("call-{name}"), name, arguments)],
usage: Some(usage),
},
usage: Some(usage),
finish_reason: Some("tool_calls".to_string()),
raw: None,
resolved_model: None,
continue_turn: None,
}
}
async fn run_bus_turn(
model: Arc<dyn ChatModel<()>>,
tools: Vec<Box<dyn Tool>>,
max_tool_iterations: usize,
visible_tool_names: Option<HashSet<String>>,
) -> Result<AgentTurnResponse, String> {
init_global(DEFAULT_CAPACITY);
register_agent_handlers();
request_native_global::<AgentTurnRequest, AgentTurnResponse>(
AGENT_RUN_TURN_METHOD,
AgentTurnRequest {
turn_model_source: openhuman_core::openhuman::tinyagents::TurnModelSource::from_model(
model,
),
history: vec![ChatMessage::system("system"), ChatMessage::user("run")],
tools_registry: Arc::new(tools),
provider_name: "round15".to_string(),
model: "gpt-4o-mini".to_string(),
temperature: 0.0,
silent: true,
channel_name: "round15".to_string(),
multimodal: MultimodalConfig::default(),
multimodal_files: MultimodalFileConfig::default(),
max_tool_iterations,
on_delta: None,
target_agent_id: Some("orchestrator".to_string()),
visible_tool_names,
extra_tools: Vec::new(),
on_progress: None,
origin: openhuman_core::openhuman::agent::turn_origin::AgentTurnOrigin::Cli,
},
)
.await
.map_err(|err| err.to_string())
}
#[tokio::test]
async fn bus_turn_native_tools_dedups_streams_and_records_tool_messages() {
let provider = Arc::new(ScriptedModel {
responses: Mutex::new(
vec![
Ok(tool_response("echo", json!({ "value": "alpha" }))),
Ok(text_response("final native answer")),
]
.into(),
),
turns: Mutex::new(Vec::new()),
profile: ModelProfile {
provider: Some("round15".to_string()),
tool_calling: true,
parallel_tool_calls: true,
streaming: true,
streaming_tool_chunks: true,
..ModelProfile::default()
},
always_fail: None,
stream_events: vec![
ModelStreamItem::MessageDelta(MessageDelta::text("draft ")),
ModelStreamItem::MessageDelta(MessageDelta::reasoning("thinking")),
ModelStreamItem::ToolCallDelta(ToolDelta {
call_id: "call-echo".to_string(),
tool_name: Some("echo".to_string()),
content: String::new(),
}),
ModelStreamItem::ToolCallDelta(ToolDelta {
call_id: "call-echo".to_string(),
tool_name: None,
content: "{\"value\":\"alpha\"}".to_string(),
}),
],
});
let response = run_bus_turn(
provider.clone(),
vec![
StaticTool::ok("echo", "first"),
StaticTool::ok("echo", "duplicate"),
StaticTool::ok("other", "unused"),
],
4,
None,
)
.await
.unwrap();
assert_eq!(response.text, "final native answer");
let turns = provider.turns();
assert_eq!(turns[0].tool_names, vec!["echo", "other"]);
assert!(
turns[1]
.messages
.iter()
.any(|msg| matches!(msg, Message::Tool(_)) && msg.text().contains("first:alpha")),
"second native request should carry a role=tool result message"
);
}
#[tokio::test]
async fn bus_turn_prompt_mode_covers_invisible_cli_only_and_unknown_tools() {
let mut visible = HashSet::new();
visible.insert("allowed".to_string());
let invisible_provider = ScriptedModel::new(vec![
tool_response("hidden", json!({ "value": "x" })),
text_response("after invisible"),
]);
let invisible_response = run_bus_turn(
invisible_provider.clone(),
vec![StaticTool::ok("allowed", "allowed")],
4,
Some(visible),
)
.await
.unwrap();
assert_eq!(invisible_response.text, "after invisible");
let provider = ScriptedModel::new(vec![
tool_response("cli_only", json!({ "value": "x" })),
tool_response("missing", json!({ "value": "x" })),
text_response("recovered"),
]);
let response = run_bus_turn(
provider.clone(),
vec![
StaticTool::ok("allowed", "allowed"),
StaticTool::cli_only("cli_only"),
],
6,
None,
)
.await
.unwrap();
assert_eq!(response.text, "recovered");
let joined = provider
.turns()
.into_iter()
.flat_map(|turn| turn.messages)
.map(|msg| msg.text().to_string())
.collect::<Vec<_>>()
.join("\n");
let invisible_joined = invisible_provider
.turns()
.into_iter()
.flat_map(|turn| turn.messages)
.map(|msg| msg.text().to_string())
.collect::<Vec<_>>()
.join("\n");
// Unknown-tool recovery now flows through the tinyagents
// `UnknownToolPolicy::ReturnToolError` path (issue #4249), which emits
// `unknown tool `<name>` (arguments: …)` instead of the legacy
// `Unknown tool: <name>` wording.
assert!(
invisible_joined.contains("unknown tool") && invisible_joined.contains("hidden"),
"invisible tool call should surface a crate unknown-tool result naming `hidden`"
);
assert!(joined.contains("only available via explicit CLI/RPC invocation"));
assert!(
joined.contains("unknown tool") && joined.contains("missing"),
"unknown tool call should surface a crate unknown-tool result naming `missing`"
);
}
#[tokio::test]
async fn bus_turn_halts_on_repeated_tool_error_and_truncates_capped_result() {
let provider = ScriptedModel::new(vec![
tool_response("capper", json!({ "value": "" })),
tool_response("fail", json!({ "value": "same" })),
tool_response("fail", json!({ "value": "same" })),
tool_response("fail", json!({ "value": "same" })),
]);
let response = run_bus_turn(
provider.clone(),
vec![
StaticTool::err("fail", "boom"),
StaticTool::capped("capper", "abcdefghijklmnopqrstuvwxyz", 5),
],
8,
None,
)
.await
.unwrap();
assert!(response.text.contains("retried 3 times"));
assert!(response.text.contains("boom:same"));
let joined = provider
.turns()
.into_iter()
.flat_map(|turn| turn.messages)
.map(|msg| msg.text().to_string())
.collect::<Vec<_>>()
.join("\n");
assert!(joined.contains("[truncated by tool cap: 21 more chars not shown]"));
}
#[tokio::test]
async fn bus_turn_surfaces_provider_error_and_iteration_cap() {
let provider_error = run_bus_turn(
ScriptedModel::failing("provider unavailable"),
vec![StaticTool::ok("echo", "ok")],
2,
None,
)
.await
.err()
.expect("provider error should surface");
assert!(provider_error.contains("provider unavailable"));
let capped = run_bus_turn(
ScriptedModel::new(vec![tool_response("missing", json!({ "value": "x" }))]),
vec![StaticTool::ok("echo", "ok")],
1,
None,
)
.await
.err()
.expect("iteration cap should surface");
assert!(capped.contains("maximum tool iterations"));
}
#[tokio::test]
async fn agent_builder_prompt_and_debug_dump_cover_public_session_paths() {
let workspace = round15_workspace("session-prompt");
std::fs::create_dir_all(&workspace).unwrap();
std::fs::write(workspace.join("PROFILE.md"), "Round15 profile").unwrap();
std::fs::write(workspace.join("MEMORY.md"), "Round15 memory").unwrap();
let provider = ScriptedModel::new(vec![text_response("unused")]);
let mut config = AgentConfig::default();
config.max_tool_iterations = 2;
config.max_history_messages = 4;
let agent = AgentBuilder::new()
.chat_model(provider)
.tools(vec![StaticTool::ok("echo", "ok")])
.memory(Arc::new(NoopMemory::default()))
.tool_dispatcher(Box::new(XmlToolDispatcher))
.config(config)
.workspace_dir(workspace.clone())
.agent_definition_name("round15/orchestrator")
.event_context("round15-session", "round15-channel")
.omit_profile(false)
.omit_memory_md(false)
.build()
.unwrap();
assert_eq!(agent.agent_definition_name(), "round15/orchestrator");
assert!(agent.session_key().contains("round15_orchestrator"));
let prompt = agent
.build_system_prompt(LearnedContextData::default())
.unwrap();
assert!(prompt.contains("Round15 profile"));
assert!(prompt.contains("Round15 memory"));
assert!(prompt.contains("echo"));
let dump_err = dump_agent_prompt(DumpPromptOptions {
agent_id: "integrations_agent".to_string(),
toolkit: None,
workspace_dir_override: Some(workspace),
model_override: Some("round15-model".to_string()),
})
.await
.unwrap_err()
.to_string();
assert!(dump_err.contains("integrations_agent requires a `toolkit` argument"));
}
#[tokio::test]
async fn agent_turn_blank_final_response_is_typed_error() {
let workspace = round15_workspace("blank-final");
std::fs::create_dir_all(&workspace).unwrap();
let provider = ScriptedModel::new(vec![ModelResponse::assistant("")]);
let mut agent = Agent::builder()
.chat_model(provider)
.tools(vec![])
.memory(Arc::new(NoopMemory::default()))
.tool_dispatcher(Box::new(XmlToolDispatcher))
.config(AgentConfig {
max_tool_iterations: 1,
..AgentConfig::default()
})
.workspace_dir(workspace)
.build()
.unwrap();
let err = agent.turn("blank please").await.unwrap_err().to_string();
assert!(err.contains("empty response"));
}
fn round15_workspace(label: &str) -> PathBuf {
std::env::current_dir()
.unwrap()
.join("target")
.join(format!(
"agent-tool-loop-round15-{label}-{}",
uuid::Uuid::new_v4()
))
}