Files
zesdex/src/app/subagent/engine.rs
T

125 lines
4.8 KiB
Rust
Raw Normal View History

use tokio::sync::mpsc;
use crate::dto::chat::message::ChatMessage;
use crate::tool::{all_tools, tool_is_risky};
use super::context::SubagentContext;
use super::event::SubagentEvent;
pub const MAX_AGENT_STEPS: usize = 25;
fn tool_call_from_response(response: &str) -> Vec<String> {
let mut calls = Vec::new();
for line in response.lines() {
let trimmed = line.trim();
if let Some(tool_call) = trimmed.strip_prefix("Tool: ") {
calls.push(tool_call.to_string());
}
}
calls
}
pub fn run_subagent(ctx: SubagentContext, tx: mpsc::Sender<SubagentEvent>) -> anyhow::Result<String> {
let mut output = String::new();
let mut messages: Vec<ChatMessage> = Vec::new();
messages.push(ChatMessage::system(ctx.system_prompt.clone()));
let tool_ctx = crate::tool::ToolCtx::builder()
.session_dir(ctx.session_dir.clone())
.origin(crate::app::state::types::Origin::SubAgent)
.build();
let max_steps = ctx.max_steps.min(MAX_AGENT_STEPS);
for step in 0..max_steps {
let api_key = std::env::var("API_KEY").unwrap_or_default();
let model = std::env::var("MODEL").unwrap_or_default();
let client = crate::service::provider::LlmClient::new(api_key, model, None);
let response = match client.chat(&messages) {
Ok(r) => r,
Err(e) => {
let _ = tx.blocking_send(SubagentEvent::StepFailed {
_step: step,
_error: e.to_string(),
});
anyhow::bail!("subagent call failed at step {}: {}", step, e);
}
};
let _ = tx.blocking_send(SubagentEvent::ToolCall {
_tool: "api".to_string(),
_args: serde_json::json!({"response": response}),
});
let tool_calls = tool_call_from_response(&response);
if tool_calls.is_empty() {
output.push_str(&response);
output.push('\n');
let _ = tx.blocking_send(SubagentEvent::StepCompleted {
_step: step,
_output: response.clone(),
});
if !response.contains("Tool:") {
break;
}
} else {
let tools = all_tools();
for tool_name in &tool_calls {
let explicitly_allowed = ctx.allowed_tools.contains(tool_name);
let generally_allowed = ctx.allowed_tools.is_empty() || explicitly_allowed;
if !generally_allowed {
let msg = format!("tool '{}' not allowed for this subagent", tool_name);
messages.push(ChatMessage::tool_result(tool_name.clone(), msg.clone()));
let _ = tx.blocking_send(SubagentEvent::ToolResult {
_tool: tool_name.clone(),
_output: msg,
});
continue;
}
if tool_is_risky(tool_name) && !explicitly_allowed {
let msg = format!("risky tool '{}' requires explicit permission; not allowed for this subagent", tool_name);
messages.push(ChatMessage::tool_result(tool_name.clone(), msg.clone()));
let _ = tx.blocking_send(SubagentEvent::ToolResult {
_tool: tool_name.clone(),
_output: msg,
});
continue;
}
let result = match tools.iter().find(|t| t.name() == *tool_name) {
Some(tool) => tool.run(&tool_ctx, &serde_json::json!({})),
None => Err(anyhow::anyhow!("tool '{}' not found", tool_name)),
};
match result {
Ok(output_text) => {
let _ = tx.blocking_send(SubagentEvent::ToolResult {
_tool: tool_name.clone(),
_output: output_text,
});
}
Err(e) => {
let msg = format!("tool '{}' failed: {}", tool_name, e);
let _ = tx.blocking_send(SubagentEvent::ToolResult {
_tool: tool_name.clone(),
_output: msg,
});
}
}
}
let _ = tx.blocking_send(SubagentEvent::StepCompleted {
_step: step,
_output: response.clone(),
});
}
let assistant_msg = ChatMessage::assistant(Some(response.clone()));
messages.push(assistant_msg);
let user_msg = ChatMessage::user("Continue with the next step based on the tool results above.".to_string());
messages.push(user_msg);
}
let _ = tx.blocking_send(SubagentEvent::Completed { _output: output.clone() });
Ok(output)
}