feat(agent): subagent tool paralel + auto-load AGENTS.md + verify cek setelah edit
Lanjutan audit alur AI agent (round 2), mengisi celah yang tersisa dari
perpbaikan paralel tool di loop utama (74b1ad4) agar lebih mirip Claude Code.
- feat(subagent): eksekusi batch tool read-only paralel di subagent engine
(engine.rs). Tool::run sinkron, jadi pakai scoped OS thread (bounded
window 8); hasil dipertahankan dalam urutan panggilan asli. Batch dengan
tool mutating jatuh balik ke jalur sequential aman.
- feat(agent): auto-load AGENTS.md/CLAUDE.md/.cursorrules ke system prompt
tiap turn (seperti Claude Code load AGENTS.md saat startup). Fungsi
main_agent_prompt_with_project_context menempel blok PROJECT CONTEXT;
dibaca dari workspace root pertama & dibatasi 12k char.
- feat(prompt): arahan VERIFY AFTER EDIT — setelah edit/write, agent wajib
jalankan cargo check/clippy/test (atau lint/test sesuai stack) via bash
sebelum mengakhiri turn; perbaiki error yang terlihat, jangan klaim
'compiles/works' tanpa hasil nyata.
- feat(infra): build_rich_context kini membaca AGENTS.md & CLAUDE.md juga
(untuk explore_codebase/scout).
- test: +3 subagent engine (order paralel, kecepatan konkuren, fallback
mutating), +2 domain prompt (konteks proyek & fallback kosong).
This commit is contained in:
@@ -14,7 +14,8 @@ use tracing::{debug, info, instrument};
|
||||
use crate::llm::provider::LlmClient;
|
||||
use crate::subagent::context::SubagentContext;
|
||||
use crate::subagent::division::{tools_for, AccessTier};
|
||||
use crate::tools::{tool_defs, ToolCtx};
|
||||
use crate::tools::{tool_defs, Tool, ToolCtx};
|
||||
use serde_json::Value;
|
||||
use zesdex_domain::agent::progress::AgentProgress;
|
||||
use zesdex_domain::core::tool_call::sanitize_tool_arguments;
|
||||
use zesdex_domain::core::ChatMessage;
|
||||
@@ -31,6 +32,95 @@ const TOOL_OUTPUT_MAX_CHARS: usize = 12_000;
|
||||
/// recovery note steering the model to a different approach.
|
||||
const MAX_CONSECUTIVE_TOOL_ERRORS: usize = 3;
|
||||
|
||||
/// Maximum number of read-only tool calls executed concurrently in a single
|
||||
/// subagent batch. Read-only tools (read/grep/glob/…) block on disk I/O, so
|
||||
/// running them on parallel OS threads removes the serial round-trip latency
|
||||
/// for a batch of independent lookups, mirroring the main turn loop.
|
||||
const MAX_PARALLEL_TOOLS: usize = 8;
|
||||
|
||||
/// Execute a batch of tool calls, running read-only tools concurrently when
|
||||
/// the whole batch is parallel-safe.
|
||||
///
|
||||
/// Returns one `(tool_call_id, tool_name, result)` per call **in the original
|
||||
/// call order** (OpenAI/Anthropic tool-result ordering contract). `Tool::run`
|
||||
/// is synchronous, so real parallelism comes from scoped OS threads; `Tool`
|
||||
/// and `ToolCtx` are `Send + Sync`, so the borrowed references can be shared
|
||||
/// across the short-lived scoped threads.
|
||||
///
|
||||
/// If any single tool in the batch mutates state (edit/write/bash/git/…), the
|
||||
/// whole batch falls back to the safe sequential path so writes never race.
|
||||
fn execute_tool_batch(
|
||||
tools: &[Box<dyn Tool>],
|
||||
tool_ctx: &ToolCtx,
|
||||
tool_calls: &[zesdex_domain::core::ToolCall],
|
||||
) -> Vec<(String, String, String)> {
|
||||
let parallel = tool_calls.len() > 1
|
||||
&& tool_calls
|
||||
.iter()
|
||||
.all(|tc| crate::tools::tool_is_parallel_safe(&tc.function.name));
|
||||
|
||||
if !parallel {
|
||||
// Sequential fallback (kept identical to the historical behavior).
|
||||
return tool_calls
|
||||
.iter()
|
||||
.map(|tc| {
|
||||
let tool_name = tc.function.name.clone();
|
||||
let args = sanitize_tool_arguments(&tc.function.arguments);
|
||||
let result = run_one_tool(tools, tool_ctx, &tool_name, &args);
|
||||
(tc.id.clone(), tool_name, result)
|
||||
})
|
||||
.collect();
|
||||
}
|
||||
|
||||
// Bounded parallel path: process the batch in windows of
|
||||
// `MAX_PARALLEL_TOOLS` so concurrency stays bounded, joining each window
|
||||
// before the next so results stay in original order.
|
||||
let mut ordered = Vec::with_capacity(tool_calls.len());
|
||||
for window in tool_calls.chunks(MAX_PARALLEL_TOOLS) {
|
||||
let window_results = std::thread::scope(|s| {
|
||||
let handles: Vec<_> = window
|
||||
.iter()
|
||||
.map(|tc| {
|
||||
let tool_name = tc.function.name.clone();
|
||||
let args = sanitize_tool_arguments(&tc.function.arguments);
|
||||
s.spawn(move || {
|
||||
debug!("Subagent executing tool: {tool_name}");
|
||||
run_one_tool(tools, tool_ctx, &tool_name, &args)
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
handles
|
||||
.into_iter()
|
||||
.map(|h| {
|
||||
h.join()
|
||||
.unwrap_or_else(|_| "Error: tool panicked".to_string())
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
});
|
||||
for (tc, result) in window.iter().zip(window_results) {
|
||||
ordered.push((tc.id.clone(), tc.function.name.clone(), result));
|
||||
}
|
||||
}
|
||||
ordered
|
||||
}
|
||||
|
||||
/// Run a single synchronous tool call and capture its result string.
|
||||
fn run_one_tool(
|
||||
tools: &[Box<dyn Tool>],
|
||||
tool_ctx: &ToolCtx,
|
||||
tool_name: &str,
|
||||
args: &Value,
|
||||
) -> String {
|
||||
if let Some(tool) = tools.iter().find(|t| t.name() == tool_name) {
|
||||
match tool.run(tool_ctx, args) {
|
||||
Ok(output) => output,
|
||||
Err(e) => format!("Error: {e}"),
|
||||
}
|
||||
} else {
|
||||
format!("Unknown tool: {tool_name}")
|
||||
}
|
||||
}
|
||||
|
||||
/// Pick a `max_tokens` budget proportional to the directive's length.
|
||||
fn adaptive_max_tokens(directive_len: usize) -> u32 {
|
||||
if directive_len <= 80 {
|
||||
@@ -139,12 +229,12 @@ pub async fn run_agent(
|
||||
return Ok(content);
|
||||
}
|
||||
|
||||
// Execute tool calls
|
||||
for tc in &tool_calls {
|
||||
let tool_name = &tc.function.name;
|
||||
let args = sanitize_tool_arguments(&tc.function.arguments);
|
||||
// Execute tool calls — read-only batches run concurrently (bounded,
|
||||
// order preserved); any mutating tool forces the safe sequential path.
|
||||
let results = execute_tool_batch(&tools, &tool_ctx, &tool_calls);
|
||||
|
||||
debug!("Subagent executing tool: {tool_name}");
|
||||
for (id, tool_name, result) in results {
|
||||
debug!("Subagent tool {tool_name} finished");
|
||||
|
||||
report_progress(
|
||||
&tool_ctx,
|
||||
@@ -155,15 +245,6 @@ pub async fn run_agent(
|
||||
),
|
||||
);
|
||||
|
||||
let result = if let Some(tool) = tools.iter().find(|t| t.name() == tool_name) {
|
||||
match tool.run(&tool_ctx, &args) {
|
||||
Ok(output) => output,
|
||||
Err(e) => format!("Error: {e}"),
|
||||
}
|
||||
} else {
|
||||
format!("Unknown tool: {tool_name}")
|
||||
};
|
||||
|
||||
// Error-recovery: if the same tool keeps failing, inject a
|
||||
// system note steering the model to a different approach.
|
||||
if result.starts_with("Error:") {
|
||||
@@ -171,11 +252,11 @@ pub async fn run_agent(
|
||||
consecutive_errors += 1;
|
||||
} else {
|
||||
consecutive_errors = 1;
|
||||
last_tool = tool_name.to_string();
|
||||
last_tool = tool_name.clone();
|
||||
}
|
||||
if consecutive_errors >= MAX_CONSECUTIVE_TOOL_ERRORS {
|
||||
messages.push(ChatMessage::system(
|
||||
zesdex_domain::agent::prompt::error_recovery_note(tool_name, &result),
|
||||
zesdex_domain::agent::prompt::error_recovery_note(&tool_name, &result),
|
||||
));
|
||||
consecutive_errors = 0;
|
||||
}
|
||||
@@ -183,10 +264,7 @@ pub async fn run_agent(
|
||||
consecutive_errors = 0;
|
||||
}
|
||||
|
||||
messages.push(ChatMessage::tool(
|
||||
tc.id.clone(),
|
||||
truncate_tool_output(result),
|
||||
));
|
||||
messages.push(ChatMessage::tool(id, truncate_tool_output(result)));
|
||||
}
|
||||
|
||||
// Add assistant response if there was text content
|
||||
@@ -208,3 +286,111 @@ pub async fn run_agent(
|
||||
"Subagent reached iteration limit ({MAX_ITERATIONS})"
|
||||
))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::tools::ToolCtxBuilder;
|
||||
use serde_json::json;
|
||||
|
||||
/// A deterministic mock tool whose `run` returns its own name (opting into
|
||||
/// an optional sleep to make parallel-vs-sequential observable).
|
||||
struct MockTool {
|
||||
name: &'static str,
|
||||
sleep_ms: u64,
|
||||
}
|
||||
|
||||
impl MockTool {
|
||||
fn new(name: &'static str, sleep_ms: u64) -> Self {
|
||||
Self { name, sleep_ms }
|
||||
}
|
||||
}
|
||||
|
||||
impl Tool for MockTool {
|
||||
fn name(&self) -> &'static str {
|
||||
self.name
|
||||
}
|
||||
fn description(&self) -> &'static str {
|
||||
"mock tool for tests"
|
||||
}
|
||||
fn parameters(&self) -> Value {
|
||||
json!({"type":"object","properties":{}})
|
||||
}
|
||||
fn run(&self, _ctx: &ToolCtx, _args: &Value) -> Result<String> {
|
||||
if self.sleep_ms > 0 {
|
||||
std::thread::sleep(std::time::Duration::from_millis(self.sleep_ms));
|
||||
}
|
||||
Ok(self.name.to_string())
|
||||
}
|
||||
}
|
||||
|
||||
fn tc(name: &str, id: usize) -> zesdex_domain::core::ToolCall {
|
||||
zesdex_domain::core::ToolCall {
|
||||
id: format!("call_{id}"),
|
||||
type_: "function".to_string(),
|
||||
function: zesdex_domain::core::ToolFunction {
|
||||
name: name.to_string(),
|
||||
arguments: serde_json::Value::String(String::new()),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
fn ctx() -> ToolCtx {
|
||||
ToolCtxBuilder::default().build()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_batch_preserves_original_order() {
|
||||
let tools: Vec<Box<dyn Tool>> = vec![
|
||||
Box::new(MockTool::new("read", 0)),
|
||||
Box::new(MockTool::new("grep", 0)),
|
||||
];
|
||||
let calls = vec![tc("read", 1), tc("grep", 2), tc("read", 3)];
|
||||
|
||||
let results = execute_tool_batch(&tools, &ctx(), &calls);
|
||||
|
||||
// Results keep the assistant's original call order.
|
||||
let names: Vec<&str> = results.iter().map(|(_, n, _)| n.as_str()).collect();
|
||||
assert_eq!(names, vec!["read", "grep", "read"]);
|
||||
// IDs follow the same original order (ordering contract).
|
||||
let ids: Vec<&str> = results.iter().map(|(id, _, _)| id.as_str()).collect();
|
||||
assert_eq!(ids, vec!["call_1", "call_2", "call_3"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_read_batch_is_faster_than_sequential() {
|
||||
// Both reads sleep 30ms each. Parallel should finish ~30ms (both run
|
||||
// at once), sequential would take ~60ms.
|
||||
let tools: Vec<Box<dyn Tool>> = vec![Box::new(MockTool::new("read", 30))];
|
||||
let calls = vec![tc("read", 1), tc("read", 2)];
|
||||
|
||||
let started = std::time::Instant::now();
|
||||
let results = execute_tool_batch(&tools, &ctx(), &calls);
|
||||
let elapsed = started.elapsed();
|
||||
|
||||
assert_eq!(results.len(), 2);
|
||||
assert!(
|
||||
elapsed < std::time::Duration::from_millis(55),
|
||||
"parallel read batch took {elapsed:?}, expected concurrent execution"
|
||||
);
|
||||
assert!(elapsed >= std::time::Duration::from_millis(25));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mutating_tool_forces_sequential_batch() {
|
||||
// A batch containing a mutating tool ("write") must NOT run in
|
||||
// parallel — the single 30ms read runs alone, then the write runs.
|
||||
let tools: Vec<Box<dyn Tool>> = vec![
|
||||
Box::new(MockTool::new("read", 30)),
|
||||
Box::new(MockTool::new("write", 0)),
|
||||
];
|
||||
let calls = vec![tc("read", 1), tc("write", 2)];
|
||||
|
||||
let results = execute_tool_batch(&tools, &ctx(), &calls);
|
||||
|
||||
let names: Vec<&str> = results.iter().map(|(_, n, _)| n.as_str()).collect();
|
||||
assert_eq!(names, vec!["read", "write"]);
|
||||
let ids: Vec<&str> = results.iter().map(|(id, _, _)| id.as_str()).collect();
|
||||
assert_eq!(ids, vec!["call_1", "call_2"]);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user