Compare commits

..
14 Commits
Author SHA1 Message Date
semantic-release-bot 20ce81a6be chore(release): 1.20.0 [skip ci]
# [1.20.0](https://github.com/asepharyana/zesdex/compare/v1.19.6...v1.20.0) (2026-08-28)

### Features

* **agent:** subagent tool paralel + auto-load AGENTS.md + verify cek setelah edit ([21e3ccc](https://github.com/asepharyana/zesdex/commit/21e3ccc891ab886044a5f0a770db710769d1287b))
2026-08-28 02:54:09 +00:00
asepharyana 21e3ccc891 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).
2026-08-28 09:50:21 +07:00
semantic-release-bot 992e60980c chore(release): 1.19.6 [skip ci]
## [1.19.6](https://github.com/asepharyana/zesdex/compare/v1.19.5...v1.19.6) (2026-08-27)

### Performance Improvements

* **agent:** eksekusi tool read-only paralel seperti Claude Code ([74b1ad4](https://github.com/asepharyana/zesdex/commit/74b1ad43020a11ca27e60e3a24de0db5d1ab37b4))
2026-08-27 18:22:01 +00:00
asepharyana 74b1ad4302 perf(agent): eksekusi tool read-only paralel seperti Claude Code
Sebelumnya loop utama mengeksekusi semua tool call satu-per-satu
(sequential for loop). Seperti Claude Code, tool read-only yang
independen dalam satu pesan assistant (read/grep/glob/semantic_search
dsb.) kini dijalankan konkuren dengan bounded parallelism (max 8),
mengurangi latensi per turn secara signifikan untuk beban coding.

- feat(registry): tool_is_parallel_safe() — whitelist tool read-only
  yang aman dijalankan paralel; tool mutating/shell tetap sequential
- feat(executor): is_parallel_safe() delegasi ke registry; ToolExecutor
  trait Default=false (konservatif)
- fix(application): execute_tool_calls_in_parallel() — join_all +
  semaphore bounded 8, hasil dikumpulkan dalam URUTAN panggilan asli
  (kontrak OpenAI/Anthropic tool-result ordering)
- loop utama: batch paralel hanya jika SEMUA tool parallel-safe; jika
  ada satu tool mutating, jatuh balik ke jalur sequential aman
- test: +2 registry test, +2 application test (konkurensi & urutan,
  fallback batch mutating)
2026-08-28 01:17:26 +07:00
semantic-release-bot 9ad04cf819 chore(release): 1.19.5 [skip ci]
## [1.19.5](https://github.com/asepharyana/zesdex/compare/v1.19.4...v1.19.5) (2026-08-27)

### Bug Fixes

* **api:** model Opus default pakai claude-opus-5 (bukan -4-8) ([b28a5fe](https://github.com/asepharyana/zesdex/commit/b28a5fe384fd45255a6b249c7febd3b4ffc0fd2f))
* **api:** update zesdex packages to version 1.19.4 ([b46935c](https://github.com/asepharyana/zesdex/commit/b46935c606d4f68ea227c7db27e4c2aa9b4e373c))
2026-08-27 17:36:59 +00:00
asepharyana b46935c606 fix(api): update zesdex packages to version 1.19.4 2026-08-28 00:32:27 +07:00
asepharyana b28a5fe384 fix(api): model Opus default pakai claude-opus-5 (bukan -4-8)
Model terbaru di 9router adalah claude-opus-5. Update semua jalur
model default Opus:

- fix(app_config_repo): fallback default_model custom_model.unwrap_or
  -> claude-opus-5; model_roles list claude-opus-5
- fix(app_config): router provider default_model -> claude-opus-5
- fix(settings test): assertion claude-opus-5
- fix(data): ~/.local/share/zesdex/settings.json model -> claude-opus-5
2026-08-28 00:32:06 +07:00
semantic-release-bot b66898ea28 chore(release): 1.19.4 [skip ci]
## [1.19.4](https://github.com/asepharyana/zesdex/compare/v1.19.3...v1.19.4) (2026-08-27)

### Bug Fixes

* **api:** model claude selalu pakai Opus dari settings.json, bukan deepseek ([9aca45c](https://github.com/asepharyana/zesdex/commit/9aca45cb65d6913d14fecc10d70c180934f69d74))
2026-08-27 17:05:10 +00:00
asepharyana 9aca45cb65 fix(api): model claude selalu pakai Opus dari settings.json, bukan deepseek
TUI turn.rs & daemon handler.rs ambil model langsung dari
settings.model (tersimpan 'deepseek-v4-flash-free' di
~/.local/share/zesdex/settings.json) padahal provider sudah 'claude'.

- feat(domain): resolve_effective_model() — saat provider claude, model
  diambil dari app_config provider claude (default_model=claude-opus-4-8
  hasil deteksi ~/.claude/settings.json), menang atas settings.model basi.
  Provider non-claude tetap hormati settings.model user.
- fix(tui): turn.rs pakai resolve_effective_model (bukan settings.model)
- fix(daemon): handler.rs run_turn + compaction pakai resolve_effective_model
- fix(data): ~/.local/share/zesdex/settings.json model deepseek -> claude-opus-4-8
- test: 3 unit test resolve_effective_model
2026-08-28 00:00:32 +07:00
semantic-release-bot 4dccf0cee4 chore(release): 1.19.3 [skip ci]
## [1.19.3](https://github.com/asepharyana/zesdex/compare/v1.19.2...v1.19.3) (2026-08-27)

### Bug Fixes

* **api:** model Opus pakai URL + API custom dari ~/.claude/settings.json ([1f91b44](https://github.com/asepharyana/zesdex/commit/1f91b447080e7106201d77585edb7d38b74e34bb))
2026-08-27 16:49:19 +00:00
asepharyana 1f91b44708 fix(api): model Opus pakai URL + API custom dari ~/.claude/settings.json
Perbaiki provider claude agar selalu refresh dari settings.json dan
menjadi default (claude-opus-4-8) setiap startup:

- fix(app_config_repo): ganti or_insert -> insert untuk provider claude —
  base_url/key dari ~/.claude/settings.json selalu di-refresh, tidak
  tertutup snapshot lama app_config.json.
- fix(app_config_repo): hapus kondisi default_provider == default — saat
  settings.json terdeteksi, default_provider='claude' dan
  default_model='claude-opus-4-8' SELALU di-set (sebelumnya skip kalau
  user pernah ganti provider).
- fix(subagent/provider): resolve_subagent_provider fallback ke
  app_config.default_provider/default_model kalau settings.provider/model
  kosong — subagent ikut pakai Opus.
- test: 4 unit test (parse settings.json, refresh stale provider, custom
  model, env fallback). Verified live: settings.json terbaca (9router URL
  + key).
2026-08-27 23:45:30 +07:00
semantic-release-bot b25929824a chore(release): 1.19.2 [skip ci]
## [1.19.2](https://github.com/asepharyana/zesdex/compare/v1.19.1...v1.19.2) (2026-08-27)

### Performance Improvements

* **agent:** stabilkan async & parallel — satu runtime, bounded concurrency, isolasi error ([6a98d52](https://github.com/asepharyana/zesdex/commit/6a98d52d54a69f78710d852a69dda3ac0a4ead31))
2026-08-27 16:30:31 +00:00
asepharyana 3fd9a2b2db chore: sinkronkan Cargo.lock dengan versi 1.19.1 2026-08-27 23:26:37 +07:00
asepharyana 6a98d52d54 perf(agent): stabilkan async & parallel — satu runtime, bounded concurrency, isolasi error
Seperti Claude Code: satu runtime shared, concurrency dibatasi, error
subagent terisolasi (satu node gagal tidak menggagalkan cycle).

- feat(runtime): global tokio runtime via OnceLock — ganti 9+ titik
  Runtime::new() per tool call (spawn, parallel_delegate, workflow,
  explore, dir_cache, daemon handler). Hemat resource, hilangkan panic
  path Runtime::new().expect() di daemon compaction.
- fix(workflow): execute_cycle ganti try_join_all (fail-fast) →
  buffer_unordered(8) + isolasi error per node; node gagal di-log dan
  diganti [ERROR], hasil node lain tetap dipakai (Claude Code-style).
- fix(parallel_delegate): spawn subagent dibatasi per batch max_parallel
  (tidak unbounded threads).
- perf(subagent): run_agent adaptif max_tokens (800/1600/4096), temp 0.2,
  truncate tool output 12k, error-recovery note utk tool error berulang.
- test: runtime singleton + block_on (2 test).
2026-08-27 23:25:44 +07:00
28 changed files with 1107 additions and 145 deletions
+43
View File
@@ -1,3 +1,46 @@
# [1.20.0](https://github.com/asepharyana/zesdex/compare/v1.19.6...v1.20.0) (2026-08-28)
### Features
* **agent:** subagent tool paralel + auto-load AGENTS.md + verify cek setelah edit ([21e3ccc](https://github.com/asepharyana/zesdex/commit/21e3ccc891ab886044a5f0a770db710769d1287b))
## [1.19.6](https://github.com/asepharyana/zesdex/compare/v1.19.5...v1.19.6) (2026-08-27)
### Performance Improvements
* **agent:** eksekusi tool read-only paralel seperti Claude Code ([74b1ad4](https://github.com/asepharyana/zesdex/commit/74b1ad43020a11ca27e60e3a24de0db5d1ab37b4))
## [1.19.5](https://github.com/asepharyana/zesdex/compare/v1.19.4...v1.19.5) (2026-08-27)
### Bug Fixes
* **api:** model Opus default pakai claude-opus-5 (bukan -4-8) ([b28a5fe](https://github.com/asepharyana/zesdex/commit/b28a5fe384fd45255a6b249c7febd3b4ffc0fd2f))
* **api:** update zesdex packages to version 1.19.4 ([b46935c](https://github.com/asepharyana/zesdex/commit/b46935c606d4f68ea227c7db27e4c2aa9b4e373c))
## [1.19.4](https://github.com/asepharyana/zesdex/compare/v1.19.3...v1.19.4) (2026-08-27)
### Bug Fixes
* **api:** model claude selalu pakai Opus dari settings.json, bukan deepseek ([9aca45c](https://github.com/asepharyana/zesdex/commit/9aca45cb65d6913d14fecc10d70c180934f69d74))
## [1.19.3](https://github.com/asepharyana/zesdex/compare/v1.19.2...v1.19.3) (2026-08-27)
### Bug Fixes
* **api:** model Opus pakai URL + API custom dari ~/.claude/settings.json ([1f91b44](https://github.com/asepharyana/zesdex/commit/1f91b447080e7106201d77585edb7d38b74e34bb))
## [1.19.2](https://github.com/asepharyana/zesdex/compare/v1.19.1...v1.19.2) (2026-08-27)
### Performance Improvements
* **agent:** stabilkan async & parallel — satu runtime, bounded concurrency, isolasi error ([6a98d52](https://github.com/asepharyana/zesdex/commit/6a98d52d54a69f78710d852a69dda3ac0a4ead31))
## [1.19.1](https://github.com/asepharyana/zesdex/compare/v1.19.0...v1.19.1) (2026-08-27)
Generated
+12 -11
View File
@@ -4862,7 +4862,7 @@ dependencies = [
[[package]]
name = "zesdex-api"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"argon2",
@@ -4885,11 +4885,12 @@ dependencies = [
[[package]]
name = "zesdex-application"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"base64",
"chrono",
"futures-util",
"serde",
"serde_json",
"sha2 0.11.0",
@@ -4902,7 +4903,7 @@ dependencies = [
[[package]]
name = "zesdex-bootstrap"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"chrono",
@@ -4919,7 +4920,7 @@ dependencies = [
[[package]]
name = "zesdex-daemon"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"base64",
@@ -4943,7 +4944,7 @@ dependencies = [
[[package]]
name = "zesdex-domain"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"base64",
@@ -4959,7 +4960,7 @@ dependencies = [
[[package]]
name = "zesdex-gateway"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"axum",
@@ -4986,7 +4987,7 @@ dependencies = [
[[package]]
name = "zesdex-grpc"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"axum",
@@ -5003,7 +5004,7 @@ dependencies = [
[[package]]
name = "zesdex-infrastructure"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"argon2",
@@ -5051,7 +5052,7 @@ dependencies = [
[[package]]
name = "zesdex-tui"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"base64",
@@ -5077,7 +5078,7 @@ dependencies = [
[[package]]
name = "zesdex-web"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"axum",
@@ -5097,7 +5098,7 @@ dependencies = [
[[package]]
name = "zesdex-ws"
version = "1.18.4"
version = "1.19.5"
dependencies = [
"anyhow",
"axum",
+1 -1
View File
@@ -15,7 +15,7 @@ members = [
]
[workspace.package]
version = "1.19.1"
version = "1.20.0"
edition = "2021"
authors = ["asepharyana <superaseph@gmail.com>"]
+1
View File
@@ -16,6 +16,7 @@ uuid.workspace = true
anyhow.workspace = true
tracing.workspace = true
tokio.workspace = true
futures-util.workspace = true
base64.workspace = true
sha2.workspace = true
url.workspace = true
+11
View File
@@ -11,6 +11,17 @@ pub trait ToolExecutor: Send + Sync {
tool_name: &str,
args: &serde_json::Value,
) -> impl Future<Output = Result<String>> + Send;
/// Whether a tool is *read-only* and therefore safe to run concurrently
/// with other read-only tools in the same assistant message.
///
/// Defaults to `false` (sequential) so a caller that does not know the
/// tool surface stays conservative. Concrete executors that know their
/// tools override this — e.g. return `true` for `read`/`grep`/`glob`.
fn is_parallel_safe(&self, tool_name: &str) -> bool {
let _ = tool_name;
false
}
}
/// Service for running agent turns asynchronously.
+212 -7
View File
@@ -5,7 +5,7 @@ use tracing::{debug, info, warn};
use zesdex_domain::agent::{AgentTurnParams, TurnEvent};
use zesdex_domain::core::{ChatMessage, StreamEvent, ToolDef};
use zesdex_domain::main_agent_prompt;
use zesdex_domain::main_agent_prompt_with_project_context;
use super::ToolExecutor;
use crate::ports::ProviderService;
@@ -111,6 +111,46 @@ fn conversation_chars(messages: &[ChatMessage]) -> usize {
.sum()
}
/// The maximum combined size (characters) of project-rule files injected into
/// the system prompt, so a huge AGENTS.md cannot blow the context window.
const PROJECT_CONTEXT_MAX_CHARS: usize = 12_000;
/// Case-insensitive rule filenames auto-loaded from the workspace root into
/// the system prompt, matching the Claude-Code/AGENTS.md convention.
const RULE_FILENAMES: [&str; 6] = [
"AGENTS.md",
"agent.md",
"CLAUDE.md",
"claude.md",
".cursorrules",
".zesdexrules",
];
/// Build a compact "project context" block from the repo's convention files
/// (AGENTS.md, CLAUDE.md, .cursorrules, …) found at the workspace root.
///
/// Follows the Claude-Code convention of loading AGENTS.md at startup so the
/// model starts each turn with the repo's rules. Reads are best-effort and
/// capped at [`PROJECT_CONTEXT_MAX_CHARS`] total; missing files are skipped.
fn build_project_context(root: &std::path::Path) -> String {
let mut ctx = String::new();
for file in RULE_FILENAMES {
let full = root.join(file);
if let Ok(content) = std::fs::read_to_string(&full) {
ctx.push_str(&format!("\n### {file}\n```\n{}\n```", content.trim()));
}
}
let context = ctx.trim().to_string();
if context.len() <= PROJECT_CONTEXT_MAX_CHARS {
return context;
}
context
.chars()
.take(PROJECT_CONTEXT_MAX_CHARS)
.collect::<String>()
+ "\n...[project context truncated]"
}
/// Track repeated tool-call errors so the loop can recover instead of
/// burning iterations retrying the same failing tool.
#[derive(Default)]
@@ -187,6 +227,53 @@ async fn execute_tool_call<T: ToolExecutor>(
output
}
// ---------------------------------------------------------------------------
// Helper: bounded-parallel execution of read-only tool calls.
// ---------------------------------------------------------------------------
/// Maximum number of read-only tool calls executed concurrently in a single
/// assistant batch. Models rarely emit more than a handful of reads per
/// message; this cap keeps resource usage bounded while still removing the
/// serial round-trip latency of many independent lookups.
const MAX_PARALLEL_TOOLS: usize = 8;
fn tool_executor_ref<T: ToolExecutor>(tool_executor: &T) -> &T {
tool_executor
}
/// Execute a batch of *read-only* tool calls concurrently (bounded by
/// [`MAX_PARALLEL_TOOLS`]) and return their outputs **in the original call
/// order**.
///
/// Order preservation matters: OpenAI/Anthropic tool-calling contracts expect
/// tool-result messages to appear in the same order as the `tool_calls`
/// emitted in the assistant message. Without it, the model sees shuffled
/// results and loses track of which result belongs to which call.
///
/// Each call still pushes its `TurnEvent::ToolResult` (so the TUI shows each
/// tool as it completes) but the returned `Vec` is ordered by the input index.
async fn execute_tool_calls_in_parallel<T: ToolExecutor>(
tool_executor: &T,
turn_events: &Arc<Mutex<VecDeque<TurnEvent>>>,
tool_calls: &[zesdex_domain::core::ToolCall],
) -> Vec<String> {
let semaphore = Arc::new(tokio::sync::Semaphore::new(MAX_PARALLEL_TOOLS));
let executor_ref = tool_executor_ref(tool_executor);
let futures = tool_calls.iter().map(|tc| {
let tc = tc.clone();
let events = turn_events.clone();
let sem = semaphore.clone();
async move {
// Acquire a permit to bound concurrency across the batch.
let _permit = sem.acquire_owned().await;
execute_tool_call(executor_ref, &events, &tc).await
}
});
futures_util::future::join_all(futures).await
}
// ---------------------------------------------------------------------------
// Helper: emit usage event from optional LLM response metadata.
// ---------------------------------------------------------------------------
@@ -294,9 +381,17 @@ impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnS
// Insert system prompt at position 0 once and keep it there for the
// entire turn, avoiding per-iteration clones of the full message list.
// Auto-load repo conventions (AGENTS.md / CLAUDE.md / .cursorrules)
// from the first workspace root, like Claude Code does at startup.
let project_context = params
.workspace_roots
.first()
.map(|root| build_project_context(root))
.unwrap_or_default();
let system_prompt = main_agent_prompt_with_project_context(&project_context);
params
.messages
.insert(0, ChatMessage::system(main_agent_prompt()));
.insert(0, ChatMessage::system(system_prompt));
let original_count = params.messages.len();
// Estimate request complexity from the last user message.
@@ -384,10 +479,42 @@ impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnS
params.messages.push(assistant_msg);
// ── Execute each tool call ──────────────────────────
for tc in &tool_calls {
let output =
execute_tool_call(self.tool_executor.as_ref(), &params.turn_events, tc)
.await;
//
// If the whole batch is made of *read-only* tools
// (read/grep/glob/…), run it concurrently with bounded
// parallelism — a big latency win for coding turns that
// emit several independent lookups in one message. Any
// single mutating tool forces the whole batch back to the
// safe sequential path so writes never race.
//
// Results are always collected in the original call order
// to honour the tool-calling contract.
let parallel = tool_calls.len() > 1
&& tool_calls
.iter()
.all(|tc| self.tool_executor.is_parallel_safe(&tc.function.name));
let outputs: Vec<String> = if parallel {
execute_tool_calls_in_parallel(
self.tool_executor.as_ref(),
&params.turn_events,
&tool_calls,
)
.await
} else {
let mut sequential = Vec::with_capacity(tool_calls.len());
for tc in &tool_calls {
let out = execute_tool_call(
self.tool_executor.as_ref(),
&params.turn_events,
tc,
)
.await;
sequential.push(out);
}
sequential
};
for (tc, output) in tool_calls.iter().zip(outputs) {
if output.starts_with("Error:") {
errors.record(&tc.function.name, &output, &mut params.messages);
}
@@ -479,7 +606,6 @@ pub async fn compact_messages_with_ai<P: ProviderService>(
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn truncate_short_output_is_unchanged() {
let out = "short".to_string();
@@ -536,4 +662,83 @@ mod tests {
];
assert_eq!(conversation_chars(&messages), 3 + 11 + 6);
}
/// A fake executor that reports parallel-safety for read-only tools and
/// whose `execute` sleeps on the first call to prove the batch runs
/// concurrently (a sequential loop would pay the sleep per call).
struct FakeExecutor;
impl ToolExecutor for FakeExecutor {
async fn execute(&self, name: &str, _args: &serde_json::Value) -> anyhow::Result<String> {
if name == "read" {
// 30ms sleep on every read; a parallel batch of 3 would
// finish in ~30ms instead of ~90ms sequentially.
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
}
Ok(format!("out:{name}"))
}
fn is_parallel_safe(&self, name: &str) -> bool {
matches!(name, "read" | "grep")
}
}
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()),
},
}
}
#[test]
fn parallel_batch_runs_concurrently_and_preserves_order() {
let executor = FakeExecutor;
let events = Arc::new(Mutex::new(VecDeque::new()));
let calls = vec![tc("read", 1), tc("grep", 2), tc("read", 3)];
// All three are parallel-safe.
assert!(calls
.iter()
.all(|c| executor.is_parallel_safe(&c.function.name)));
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_time()
.build()
.unwrap();
let started = std::time::Instant::now();
let outputs = rt.block_on(execute_tool_calls_in_parallel(&executor, &events, &calls));
let elapsed = started.elapsed();
// Results are in *original* call order (read, grep, read).
assert_eq!(
outputs,
vec![
"out:read".to_string(),
"out:grep".to_string(),
"out:read".to_string()
]
);
// Two reads sleep 30ms each; sequential would take ~60ms+ for the
// two reads, parallel keeps the whole batch under 60ms.
assert!(
elapsed < std::time::Duration::from_millis(60),
"batch took {elapsed:?}, expected parallel execution"
);
assert!(elapsed >= std::time::Duration::from_millis(25));
}
#[test]
fn mutating_batch_falls_back_to_sequential_path() {
// A batch containing a mutating tool is not eligible for the parallel
// path, so the main loop keeps results ordered and side-effects safe.
let executor = FakeExecutor;
let calls = [tc("read", 1), tc("edit", 2)];
assert!(!calls
.iter()
.all(|c| executor.is_parallel_safe(&c.function.name)));
}
}
+52
View File
@@ -41,10 +41,43 @@ architectural plans and `todowrite` to maintain granular task checklists.
4. TOOL EXECUTION: Execute individual tools (file edits, terminal commands) \
within or guided by your workflows. If an error occurs, analyse and fix it.
VERIFY AFTER EDIT (CLAUDE-CODE STYLE):
- After modifying code (edit/write), run the repo's check command via `bash` \
before ending the turn: `cargo check` / `cargo clippy` / `cargo test` for Rust, \
or the equivalent lint/test (`bun run lint && bun run test`, `npm test`, etc.) \
for other stacks. Pick the project's actual verify command (see PROJECT \
CONTEXT / AGENTS.md when present).
- If the check fails, fix the errors you can see and re-run; only end the turn \
after the check passes or you cannot resolve a failure yourself (then report it \
explicitly).
- Do NOT claim code compiles or works without running a real check.
Respond conversationally, concisely, and helpfully."
.to_string()
}
/// Build the main-agent system prompt including an injected block of project
/// context (AGENTS.md / CLAUDE.md / project rules).
///
/// Like Claude Code, which loads AGENTS.md at startup so the model starts with
/// the repo's conventions, this wraps [`main_agent_prompt`] and appends a
/// clearly-delimited `## PROJECT CONTEXT` section carrying the rules the user
/// keeps next to their code. When `project_context` is empty the returned
/// prompt is identical to [`main_agent_prompt`], so callers can fall back
/// safely.
pub fn main_agent_prompt_with_project_context(project_context: &str) -> String {
let base = main_agent_prompt();
let context = project_context.trim();
if context.is_empty() {
return base;
}
format!(
"{base}\n\n\
## PROJECT CONTEXT (repo rules — follow these conventions)\n\
{context}"
)
}
/// Build a subagent directive prompt.
///
/// The directive is embedded in a system message that also communicates the
@@ -119,6 +152,25 @@ mod tests {
assert!(prompt.contains("WORKFLOW FIRST"));
}
#[test]
fn project_context_prompt_appends_context_and_keeps_base() {
let base = main_agent_prompt();
let with_ctx = main_agent_prompt_with_project_context("## AGENTS.md\nUse cargo clippy.");
assert!(with_ctx.contains("Zesdex"), "base prompt must be preserved");
assert!(with_ctx.contains("PROJECT CONTEXT"));
assert!(with_ctx.contains("Use cargo clippy."));
assert!(with_ctx.contains(&base));
// The base section should appear before the context section.
assert!(with_ctx.find("PROJECT CONTEXT").unwrap() > with_ctx.find("Zesdex").unwrap());
}
#[test]
fn empty_project_context_returns_base_prompt() {
let base = main_agent_prompt();
assert_eq!(main_agent_prompt_with_project_context(""), base);
assert_eq!(main_agent_prompt_with_project_context(" "), base);
}
#[test]
fn subagent_directive_includes_directive_text() {
let prompt = subagent_directive("test directive", "/home", "/home/project");
+2 -2
View File
@@ -80,7 +80,7 @@ impl Default for AppConfig {
///
/// ## Defaults
/// - Zen provider: `deepseek-v4-flash-free` model
/// - Router provider: `claude-opus-4-8` model
/// - Router provider: `claude-opus-5` model
/// - Default role: "default" → zen / deepseek-v4-flash-free, temp 0.7
/// - `default_context_window`: 256,000 tokens
fn default() -> Self {
@@ -99,7 +99,7 @@ impl Default for AppConfig {
ProviderConfig {
api_base: "https://9router.asepharyana.my.id/v1".to_string(),
api_key_env: Some("ROUTER_API_KEY".to_string()),
default_model: Some("claude-opus-4-8".to_string()),
default_model: Some("claude-opus-5".to_string()),
default_api_key: None,
},
);
+1
View File
@@ -47,6 +47,7 @@ pub use repository::SettingsRepository;
pub use service::ConversationService;
pub use service::MemoryService;
pub use service::SettingsService;
pub use settings::resolve_effective_model;
pub use settings::InternetMode;
pub use settings::Settings;
pub use settings::SettingsFlags;
+85
View File
@@ -23,6 +23,8 @@ use std::collections::HashMap;
use serde::{Deserialize, Serialize};
use super::app_config::AppConfig;
/// Controls how much network access the agent is permitted during a session.
///
/// ## Variants
@@ -112,3 +114,86 @@ impl Default for Settings {
}
}
}
/// Pick the effective model name for the main agent.
///
/// When `settings.provider` is `"claude"` (auto-detected from
/// `~/.claude/settings.json`), the provider's `default_model` (or the
/// app-level `default_model`) wins over a possibly-stale persisted
/// `settings.model`. Otherwise the user's explicit `settings.model` is used.
///
/// Why: the user's custom Claude endpoint (URL + API key from
/// `~/.claude/settings.json`) implies Opus as the model; a stale
/// `settings.json` (e.g. "deepseek-v4-flash-free") must not override it.
pub fn resolve_effective_model(settings: &Settings, app_config: &AppConfig) -> String {
if settings.provider == "claude" {
if let Some(m) = app_config
.providers
.get("claude")
.and_then(|p| p.default_model.clone())
{
return m;
}
return app_config.default_model.clone();
}
settings.model.clone()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cms::app_config::AppConfig;
fn claude_app_config() -> AppConfig {
let mut cfg = AppConfig::default();
cfg.providers.insert(
"claude".to_string(),
crate::cms::ProviderConfig {
api_base: "https://9router.example/v1".to_string(),
api_key_env: Some("ANTHROPIC_API_KEY".to_string()),
default_model: Some("claude-opus-5".to_string()),
default_api_key: Some("sk-test".to_string()),
},
);
cfg.default_provider = "claude".to_string();
cfg.default_model = "claude-opus-5".to_string();
cfg
}
#[test]
fn claude_provider_uses_opus_model_over_stale_settings_model() {
let settings = Settings {
provider: "claude".to_string(),
model: "deepseek-v4-flash-free".to_string(), // stale persisted
..Settings::default()
};
let model = resolve_effective_model(&settings, &claude_app_config());
assert_eq!(model, "claude-opus-5");
}
#[test]
fn non_claude_provider_uses_settings_model() {
let settings = Settings {
provider: "zen".to_string(),
model: "my-model".to_string(),
..Settings::default()
};
let model = resolve_effective_model(&settings, &AppConfig::default());
assert_eq!(model, "my-model");
}
#[test]
fn claude_falls_back_to_app_default() {
let settings = Settings {
provider: "claude".to_string(),
model: String::new(),
..Settings::default()
};
let cfg = AppConfig::default();
let model = resolve_effective_model(&settings, &cfg);
assert_eq!(model, cfg.default_model);
}
}
+4 -1
View File
@@ -58,6 +58,9 @@ pub use agent::*;
// Sub-module items need explicit re-exports
pub use agent::defaults::*;
pub use agent::progress::AgentProgress;
pub use agent::prompt::{compaction_prompt, main_agent_prompt, subagent_directive};
pub use agent::prompt::{
compaction_prompt, main_agent_prompt, main_agent_prompt_with_project_context,
subagent_directive,
};
pub use subagent::*;
pub use workflow::*;
@@ -15,7 +15,7 @@
//! read of up to 3 relevant files).
//! 4. Join the result and return a concise bullet summary as a tool message.
use anyhow::{Context, Result};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{info, warn};
@@ -118,7 +118,7 @@ impl Tool for ExploreCodebase {
model,
);
let rt = tokio::runtime::Runtime::new().context("create explore tokio runtime")?;
let rt = crate::runtime::runtime();
let result = rt.block_on(run_agent(
subagent_ctx,
&directive,
+1
View File
@@ -34,6 +34,7 @@ pub mod llm;
pub mod mcp;
pub mod middleware;
pub mod persistence;
pub mod runtime;
pub mod subagent;
pub mod tools;
pub mod utils;
@@ -72,6 +72,55 @@ fn detect_claude_settings_provider() -> Option<(ProviderConfig, Option<String>)>
))
}
/// Apply a detected Claude provider + custom model onto an `AppConfig`.
///
/// Pure (no I/O) so it can be unit-tested. Flow:
/// 1. Always `insert`s the "claude" provider (refreshing a possibly stale
/// persisted entry with the current base URL + key from settings.json).
/// 2. Registers known Claude model roles if missing.
/// 3. Always sets `default_provider = "claude"` and
/// `default_model = custom_model.unwrap_or("claude-opus-5")` so Opus
/// is the default whenever `~/.claude/settings.json` is present.
fn apply_claude_provider(
cfg: &mut AppConfig,
claude_provider: ProviderConfig,
custom_model: Option<String>,
) {
cfg.providers.insert("claude".to_string(), claude_provider);
let claude_models: [(&str, &str); 3] = [
("claude-opus-5", "claude-opus-5"),
("claude-sonnet-5", "claude-sonnet-5"),
("claude-haiku-4-5", "claude-haiku-4-5-20251001"),
];
for (role_name, model_name) in &claude_models {
cfg.model_roles
.entry(role_name.to_string())
.or_insert(ModelRole {
provider: "claude".to_string(),
model: model_name.to_string(),
max_tokens: Some(8192),
context_window: Some(200_000),
temperature: Some(0.7),
});
}
if let Some(custom) = &custom_model {
cfg.model_roles.entry(custom.clone()).or_insert(ModelRole {
provider: "claude".to_string(),
model: custom.clone(),
max_tokens: Some(8192),
context_window: Some(200_000),
temperature: Some(0.7),
});
}
// Always prefer the Claude provider + Opus model when settings.json
// is present — this is the user's explicit custom endpoint choice.
cfg.default_provider = "claude".to_string();
cfg.default_model = custom_model.unwrap_or_else(|| "claude-opus-5".to_string());
}
impl AppConfigRepository for JsonAppConfigRepository {
fn load(&self, base_dir: &Path) -> Result<AppConfig, RepositoryError> {
let path = base_dir.join("app_config.json");
@@ -87,41 +136,7 @@ impl AppConfigRepository for JsonAppConfigRepository {
}
if let Some((claude_provider, custom_model)) = detect_claude_settings_provider() {
cfg.providers
.entry("claude".to_string())
.or_insert(claude_provider);
let claude_models: [(&str, &str); 3] = [
("claude-opus-4-8", "claude-opus-4-8"),
("claude-sonnet-5", "claude-sonnet-5"),
("claude-haiku-4-5", "claude-haiku-4-5-20251001"),
];
for (role_name, model_name) in &claude_models {
cfg.model_roles
.entry(role_name.to_string())
.or_insert(ModelRole {
provider: "claude".to_string(),
model: model_name.to_string(),
max_tokens: Some(8192),
context_window: Some(200_000),
temperature: Some(0.7),
});
}
if let Some(custom) = &custom_model {
cfg.model_roles.entry(custom.clone()).or_insert(ModelRole {
provider: "claude".to_string(),
model: custom.clone(),
max_tokens: Some(8192),
context_window: Some(200_000),
temperature: Some(0.7),
});
}
if cfg.default_provider == defaults.default_provider {
cfg.default_provider = "claude".to_string();
cfg.default_model = custom_model.unwrap_or_else(|| "claude-opus-4-8".to_string());
}
apply_claude_provider(&mut cfg, claude_provider, custom_model);
}
Ok(cfg)
@@ -134,3 +149,106 @@ impl AppConfigRepository for JsonAppConfigRepository {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
fn claude_provider(base: &str, key: Option<&str>) -> ProviderConfig {
ProviderConfig {
api_base: base.to_string(),
api_key_env: Some("ANTHROPIC_API_KEY".to_string()),
default_model: Some("claude-opus-5".to_string()),
default_api_key: key.map(|s| s.to_string()),
}
}
#[test]
fn claude_settings_parse_env() {
let parsed: ClaudeSettings = serde_json::from_str(
r#"{"env":{"ANTHROPIC_BASE_URL":"https://9router.example/v1","ANTHROPIC_API_KEY":"sk-test"}}"#,
)
.unwrap();
let env = parsed.env.unwrap();
assert_eq!(
env.anthropic_base_url.as_deref(),
Some("https://9router.example/v1")
);
assert_eq!(env.anthropic_api_key.as_deref(), Some("sk-test"));
}
#[test]
fn apply_claude_refreshes_stale_provider_and_sets_opus_default() {
// Simulate a previously-persisted app_config.json with a STALE claude
// provider + non-opus default (e.g. user had switched provider).
let mut cfg = AppConfig {
providers: {
let mut m = HashMap::new();
m.insert(
"claude".to_string(),
claude_provider("https://old.example/v1", Some("sk-old")),
);
m
},
model_roles: HashMap::new(),
default_provider: "router".to_string(),
default_model: "other-model".to_string(),
default_context_window: 256_000,
};
// Detect returned a fresh provider from ~/.claude/settings.json.
apply_claude_provider(
&mut cfg,
claude_provider("https://9router.example/v1", Some("sk-new")),
None,
);
let claude = cfg.providers.get("claude").unwrap();
assert_eq!(claude.api_base, "https://9router.example/v1");
assert_eq!(claude.default_api_key.as_deref(), Some("sk-new"));
// Insert (not or_insert) → stale entry refreshed.
assert_eq!(cfg.default_provider, "claude");
assert_eq!(cfg.default_model, "claude-opus-5");
// Claude model roles registered.
assert!(cfg.model_roles.contains_key("claude-opus-5"));
assert!(cfg.model_roles.contains_key("claude-sonnet-5"));
assert!(cfg.model_roles.contains_key("claude-haiku-4-5"));
}
#[test]
fn apply_claude_honors_custom_model_from_settings() {
let mut cfg = AppConfig::default();
apply_claude_provider(
&mut cfg,
claude_provider("https://9router.example/v1", Some("sk-new")),
Some("claude-opus-5".to_string()),
);
assert_eq!(cfg.default_model, "claude-opus-5");
assert!(cfg.model_roles.contains_key("claude-opus-5"));
}
#[test]
fn detect_uses_env_creds_as_fallback() {
// When ~/.claude/settings.json is absent/unreadable, the env-var
// fallback should produce a "claude" provider. Set env vars, call
// detect, and assert the resulting provider uses them.
std::env::set_var("ANTHROPIC_BASE_URL", "https://env.example/v1");
std::env::set_var("ANTHROPIC_API_KEY", "sk-env");
match detect_claude_settings_provider() {
Some((provider, _custom)) => {
// If the real settings.json exists it wins (base could be the
// real 9router URL); otherwise env creds are used. Either way,
// the provider must have api_key_env pointing at ANTHROPIC_API_KEY.
assert_eq!(provider.api_key_env.as_deref(), Some("ANTHROPIC_API_KEY"));
}
None => {
// No file + no env (shouldn't happen since we just set env).
panic!("expected env fallback to produce a provider");
}
}
std::env::remove_var("ANTHROPIC_BASE_URL");
std::env::remove_var("ANTHROPIC_API_KEY");
}
}
+62
View File
@@ -0,0 +1,62 @@
//! Process-wide shared Tokio runtime for sync → async bridging.
//!
//! Many `Tool::run` implementations are synchronous but need to drive async
//! work (LLM calls, subagent execution). Creating a fresh
//! [`tokio::runtime::Runtime`] on every call is expensive (spawns a thread
//! pool + runtime each time) and can fail randomly under thread pressure.
//!
//! # Flow
//!
//! [`runtime()`] returns a lazily-initialised process-wide runtime created
//! exactly once via [`std::sync::OnceLock`]. Callers use
//! `runtime().block_on(...)` exactly like they would with a local runtime —
//! the only difference is the runtime is shared, so the cost is paid once per
//! process instead of once per tool call.
//!
//! # Safety
//!
//! `block_on` panics if called from within a running Tokio runtime. The
//! tools that use this helper are synchronous (`Tool::run`), so this is safe
//! in practice. Async code should never call `runtime().block_on`.
use std::sync::OnceLock;
/// Maximum worker threads for the shared runtime. Kept modest — tools are
/// mostly I/O-bound and rarely need more concurrency than this.
const RUNTIME_WORKER_THREADS: usize = 8;
static SHARED_RUNTIME: OnceLock<tokio::runtime::Runtime> = OnceLock::new();
/// Return the process-wide shared Tokio runtime, initialising it on first use.
///
/// The runtime is configured with `worker_threads = 8` and
/// `enable_all()` (time + IO drivers) so streams, timers, and network calls
/// all work. If initialisation fails (extremely rare — resource exhaustion at
/// startup), the process aborts with a clear message rather than returning
/// an error on every subsequent call.
pub fn runtime() -> &'static tokio::runtime::Runtime {
SHARED_RUNTIME.get_or_init(|| {
tokio::runtime::Builder::new_multi_thread()
.worker_threads(RUNTIME_WORKER_THREADS)
.thread_name("zesdex-shared-rt")
.enable_all()
.build()
.expect("failed to create shared tokio runtime")
})
}
#[cfg(test)]
mod tests {
use super::runtime;
#[test]
fn runtime_is_singleton() {
assert!(std::ptr::eq(runtime(), runtime()));
}
#[test]
fn runtime_blocks_and_resolves() {
let val = runtime().block_on(async { 6 * 7 });
assert_eq!(val, 42);
}
}
+260 -15
View File
@@ -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;
@@ -23,6 +24,126 @@ use zesdex_domain::subagent_directive;
/// Maximum number of tool-call iterations before the engine gives up.
const MAX_ITERATIONS: u32 = 25;
/// A single tool-result message is truncated before entering the subagent's
/// context so it cannot blow the window (matches the main turn service).
const TOOL_OUTPUT_MAX_CHARS: usize = 12_000;
/// Maximum consecutive identical tool errors before the engine injects a
/// 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 {
800
} else if directive_len <= 400 {
1600
} else {
4096
}
}
fn truncate_tool_output(output: String) -> String {
if output.len() <= TOOL_OUTPUT_MAX_CHARS {
return output;
}
let mut result: String = output.chars().take(TOOL_OUTPUT_MAX_CHARS).collect();
result.push_str(&format!(
"\n...[truncated {} chars]",
output.len() - TOOL_OUTPUT_MAX_CHARS
));
result
}
/// Emit an `AgentProgress` event onto the turn-event queue, if one is
/// configured in the `ToolCtx`.
fn report_progress(tool_ctx: &ToolCtx, progress: AgentProgress) {
@@ -85,11 +206,17 @@ pub async fn run_agent(
Some(ctx.base_url.clone()),
);
let max_tokens = adaptive_max_tokens(directive.len());
// Track repeated tool errors so the agent can recover from a dead end.
let mut consecutive_errors = 0usize;
let mut last_tool = String::new();
// Limited iteration loop so we don't run forever
for iteration in 0..MAX_ITERATIONS {
use zesdex_application::ports::ProviderService;
let (response_msg, _usage) = client
.chat(&messages, Some(defs.clone()), Some(4096), None)
.chat(&messages, Some(defs.clone()), Some(max_tokens), Some(0.2))
.await?;
let content = response_msg.content.clone().unwrap_or_default();
@@ -102,32 +229,42 @@ 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,
AgentProgress::running(
"subagent",
format!("{}:{}", directive, tool_name),
format!("{}:{tool_name}", directive),
Some(tool_name.clone()),
),
);
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}"),
// Error-recovery: if the same tool keeps failing, inject a
// system note steering the model to a different approach.
if result.starts_with("Error:") {
if last_tool.as_str() == tool_name.as_str() {
consecutive_errors += 1;
} else {
consecutive_errors = 1;
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),
));
consecutive_errors = 0;
}
} else {
format!("Unknown tool: {tool_name}")
};
consecutive_errors = 0;
}
messages.push(ChatMessage::tool(tc.id.clone(), result));
messages.push(ChatMessage::tool(id, truncate_tool_output(result)));
}
// Add assistant response if there was text content
@@ -149,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"]);
}
}
+10 -3
View File
@@ -63,15 +63,21 @@ impl SubagentProvider {
/// Resolve subagent provider and model from settings.
///
/// Flow: reads `settings.provider` and `settings.model` → if model is empty,
/// Flow: reads `settings.provider` and `settings.model` → if provider is
/// empty, falls back to `app_config.default_provider` → if model is empty,
/// falls back to the provider config's `default_model` → if that is also
/// empty, uses the domain default model constant.
/// empty, uses `app_config.default_model` → finally the domain default model
/// constant.
#[instrument]
pub fn resolve_subagent_provider(
settings: &zesdex_domain::cms::Settings,
app_config: &zesdex_domain::cms::AppConfig,
) -> (String, String) {
let provider = settings.provider.clone();
let provider = if settings.provider.is_empty() {
app_config.default_provider.clone()
} else {
settings.provider.clone()
};
let model = settings.model.clone();
// Use the default model from the provider config if available
@@ -80,6 +86,7 @@ pub fn resolve_subagent_provider(
.providers
.get(&provider)
.and_then(|p| p.default_model.clone())
.or_else(|| Some(app_config.default_model.clone()))
.unwrap_or_else(|| zesdex_domain::agent::defaults::DEFAULT_MODEL.to_string())
} else {
model
+1 -2
View File
@@ -32,7 +32,6 @@ pub fn spawn_subagent(
) -> thread::JoinHandle<Result<String>> {
info!("Spawning subagent: {directive}");
thread::spawn(move || {
let rt = tokio::runtime::Runtime::new()?;
rt.block_on(run_agent(ctx, &directive, access, tool_ctx))
crate::runtime::runtime().block_on(run_agent(ctx, &directive, access, tool_ctx))
})
}
+11
View File
@@ -15,6 +15,13 @@ impl InfrastructureToolExecutor {
tools: all_tools(),
}
}
/// Whether a tool is read-only and safe to execute concurrently with
/// other parallel-safe tools. Delegates to the registry so the main
/// turn loop and subagent engine share one source of truth.
pub fn is_parallel_safe(tool_name: &str) -> bool {
crate::tools::tool_is_parallel_safe(tool_name)
}
}
impl ToolExecutor for InfrastructureToolExecutor {
@@ -37,4 +44,8 @@ impl ToolExecutor for InfrastructureToolExecutor {
}
}
}
fn is_parallel_safe(&self, tool_name: &str) -> bool {
Self::is_parallel_safe(tool_name)
}
}
+1 -1
View File
@@ -59,7 +59,7 @@ pub mod workflow;
// like `crate::tools::{Tool, ToolCtx}` continue to work.
pub use context::{ToolCtx, ToolCtxBuilder};
pub use graduated::{check_graduated_checks, GraduatedCheck};
pub use registry::{all_tools, tool_defs, tool_is_risky};
pub use registry::{all_tools, tool_defs, tool_is_parallel_safe, tool_is_risky};
pub use util::{arg_str, execute_cmd, log_write_edit_tool, resolve_path};
/// Common interface every agent-invocable tool implements.
@@ -123,7 +123,7 @@ impl Tool for ParallelDelegate {
.collect()
} else {
// Auto-split using LLM
let rt = tokio::runtime::Runtime::new()?;
let rt = crate::runtime::runtime();
let directives = rt.block_on(auto_split_task(
&task,
max_parallel,
@@ -146,49 +146,53 @@ impl Tool for ParallelDelegate {
"parallel delegation: starting subagents"
);
// Spawn agents in parallel
let mut handles = Vec::new();
for (i, (directive, access)) in directives.iter().enumerate() {
let subagent_ctx = SubagentContext::new(
directive.clone(),
ctx.clone(),
format!("{access:?}"),
base_url.clone(),
api_key.clone(),
model.clone(),
);
debug!(agent_index = i, access = ?access, "spawning parallel agent");
let handle = spawn_subagent(subagent_ctx, directive.clone(), *access, ctx.clone());
handles.push((i, handle));
}
// Join all results
// Spawn agents in parallel — bounded: never more than `max_parallel`
// subagent threads in flight at once (Claude Code-style isolation).
let mut results: Vec<(usize, String, String)> = Vec::new();
for (i, handle) in handles {
match handle.join() {
Ok(Ok(output)) => {
info!(agent_index = i, "parallel agent completed");
results.push((i, directives[i].0.clone(), output));
}
Ok(Err(e)) => {
warn!(agent_index = i, error = %e, "parallel agent failed");
results.push((i, directives[i].0.clone(), format!("[ERROR] {e}")));
}
Err(e) => {
warn!(agent_index = i, error = ?e, "parallel agent panicked");
results.push((
i,
directives[i].0.clone(),
"[ERROR] Agent panicked".to_string(),
));
for batch in directives.chunks(max_parallel) {
let mut handles = Vec::with_capacity(batch.len());
for (i, (directive, access)) in batch.iter().enumerate() {
let global_idx = results.len() + i;
let subagent_ctx = SubagentContext::new(
directive.clone(),
ctx.clone(),
format!("{access:?}"),
base_url.clone(),
api_key.clone(),
model.clone(),
);
debug!(agent_index = global_idx, access = ?access, "spawning parallel agent");
let handle = spawn_subagent(subagent_ctx, directive.clone(), *access, ctx.clone());
handles.push((global_idx, handle));
}
// Join this batch before spawning the next.
for (i, handle) in handles {
match handle.join() {
Ok(Ok(output)) => {
info!(agent_index = i, "parallel agent completed");
results.push((i, directives[i].0.clone(), output));
}
Ok(Err(e)) => {
warn!(agent_index = i, error = %e, "parallel agent failed");
results.push((i, directives[i].0.clone(), format!("[ERROR] {e}")));
}
Err(e) => {
warn!(agent_index = i, error = ?e, "parallel agent panicked");
results.push((
i,
directives[i].0.clone(),
"[ERROR] Agent panicked".to_string(),
));
}
}
}
}
// Consolidate results
if synthesize && results.len() > 1 {
let rt = tokio::runtime::Runtime::new()?;
let rt = crate::runtime::runtime();
let consolidated =
rt.block_on(consolidate_results(&results, &base_url, &api_key, &model))?;
Ok(format!(
+84
View File
@@ -52,6 +52,34 @@ pub fn tool_is_risky(name: &str) -> bool {
matches!(name, "write" | "delete" | "edit" | "bash" | "git_operator")
}
/// Whether a tool by name is read-only and therefore safe to run *in
/// parallel* with other tool calls within the same assistant message.
///
/// Read-only tools only inspect the workspace (read files, grep, glob,
/// semantic search, list symbols, recall memory, web search, directory
/// listing). They have no side effects, so concurrent execution cannot
/// create data races or conflicting writes.
///
/// Everything else (edits, writes, deletes, shell, git, planning, memory
/// writes, agent/spawn orchestration) stays sequential to preserve
/// correctness.
pub fn tool_is_parallel_safe(name: &str) -> bool {
matches!(
name,
"read"
| "grep"
| "glob"
| "semantic_search"
| "list_symbols"
| "web_search"
| "recall"
| "dir_list"
| "pong"
| "seq_think"
| "dir_cache_update"
)
}
/// Convert a list of tools into provider-facing `ToolDef` request schema.
pub fn tool_defs(tools: &[Box<dyn super::Tool>]) -> Vec<zesdex_domain::core::ToolDef> {
tools
@@ -66,3 +94,59 @@ pub fn tool_defs(tools: &[Box<dyn super::Tool>]) -> Vec<zesdex_domain::core::Too
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn read_only_tools_are_parallel_safe() {
for name in [
"read",
"grep",
"glob",
"semantic_search",
"list_symbols",
"web_search",
"recall",
"dir_list",
"pong",
"seq_think",
"dir_cache_update",
] {
assert!(
tool_is_parallel_safe(name),
"{name} should be parallel-safe"
);
}
}
#[test]
fn mutating_and_shell_tools_are_not_parallel_safe() {
for name in [
"edit",
"write",
"delete",
"bash",
"git_operator",
"git_worktree",
"remember",
"forget",
"todowrite",
"todofinish",
"plan_enter",
"plan_ready",
"workflow_run",
"hive_mind",
"spawn_agents",
"spawn_pipeline",
"parallel_delegate",
"explore_codebase",
] {
assert!(
!tool_is_parallel_safe(name),
"{name} should NOT be parallel-safe"
);
}
}
}
@@ -63,7 +63,7 @@ impl crate::tools::Tool for DirCacheUpdate {
// Persist the resolved paths into the shared DirCache so the TUI
// and other tools can read the cached listing without re-scanning.
let dc = ctx.dir_cache.clone();
let rt = tokio::runtime::Runtime::new()?;
let rt = crate::runtime::runtime();
rt.block_on(async { dc.write().await.set(resolved).await });
info!(count, "directory cache updated");
+2 -2
View File
@@ -65,7 +65,7 @@ impl Tool for WorkflowRun {
zesdex_domain::agent::defaults::DEFAULT_MODEL.to_string(),
None,
);
let rt = tokio::runtime::Runtime::new()?;
let rt = crate::runtime::runtime();
let result: Vec<String> =
rt.block_on(async { execute_workflow(&script, ctx, &llm_client).await })?;
@@ -229,7 +229,7 @@ impl Tool for HiveMind {
.ok_or_else(|| anyhow::anyhow!("missing 'cycles' array"))?;
info!("Hive mind starting with {} cycles", cycles_val.len());
let rt = tokio::runtime::Runtime::new()?;
let rt = crate::runtime::runtime();
let mut all_node_outputs = Vec::new();
for (cycle_idx, cycle_val) in cycles_val.iter().enumerate() {
+8 -1
View File
@@ -221,7 +221,14 @@ pub fn build_rich_context(root: &Path) -> String {
));
// 2. Custom Rules
let rule_files = [".cursorrules", ".zesdexrules", "claude.md", "agent.md"];
let rule_files = [
"AGENTS.md",
"CLAUDE.md",
".cursorrules",
".zesdexrules",
"claude.md",
"agent.md",
];
for file in rule_files {
let p = root.join(file);
if let Ok(content) = std::fs::read_to_string(&p) {
@@ -1,11 +1,15 @@
//! Hive-mind cycle execution — run one cycle of parallel nodes.
//!
//! Flow: load settings → resolve LLM credentials → run all directives in the
//! cycle concurrently via try_join_all → collect Vec<NodeOutput>.
//! cycle concurrently via a BOUNDED buffer (`buffer_unordered(MAX)`) → collect
//! `Vec<NodeOutput>`. Unlike `try_join_all`, a single failing node does NOT
//! fail the whole cycle — failed nodes are logged and replaced with an
//! `[ERROR]` output so the remaining results are preserved (like Claude
//! Code's isolated subagents).
use anyhow::Result;
use futures_util::future::try_join_all;
use tracing::info;
use futures_util::stream::StreamExt;
use tracing::{info, warn};
use zesdex_domain::cms::{AppConfigRepository, SettingsRepository};
use zesdex_domain::core::Store;
@@ -15,16 +19,21 @@ use crate::subagent::context::SubagentContext;
use crate::subagent::division::AccessTier;
use crate::subagent::engine::run_agent;
use crate::tools::ToolCtx;
use zesdex_domain::workflow::{CognitiveCycle, NodeOutput};
use zesdex_domain::workflow::{CognitiveCycle, NodeDirective, NodeOutput};
/// Maximum number of hive-mind nodes running concurrently per cycle.
/// Keeps thread/runtime pressure bounded (Claude Code-style).
const MAX_CONCURRENT_NODES: usize = 8;
/// Execute one cycle: run each node directive and collect outputs.
///
/// Flow:
/// 1. Load `Settings` and `AppConfig` from the store directory.
/// 2. Resolve provider, model, base_url, and api_key.
/// 3. Spawn all directives concurrently — each builds a `SubagentContext`
/// and calls `run_agent` (Full access).
/// 4. `try_join_all` waits for all to complete, then collect `NodeOutput`s.
/// 3. Spawn directives with bounded concurrency — each builds a
/// `SubagentContext` and calls `run_agent`.
/// 4. Collect `NodeOutput`s; failed nodes are logged and replaced with an
/// `[ERROR]` placeholder so the cycle still completes.
pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result<Vec<NodeOutput>> {
info!(
"Executing cycle {} with {} directives",
@@ -53,10 +62,9 @@ pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result
let cycle_index = cycle.index;
use zesdex_domain::workflow::NodeDirective;
// Run all directives in this cycle concurrently.
let handles: Vec<_> = cycle
// Run all directives with bounded concurrency. Each node is its own
// future; failures are collected, not propagated (isolated errors).
let tasks: Vec<_> = cycle
.directives
.iter()
.enumerate()
@@ -79,18 +87,33 @@ pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result
_ => AccessTier::Read,
};
let node_id = format!("Node-{}-{}", cycle_index, i);
async move {
let result = run_agent(ctx, &dir, access, tc).await?;
Ok::<NodeOutput, anyhow::Error>(NodeOutput {
id: format!("Node-{}-{}", cycle_index, i),
directive: dir,
output: result,
})
match run_agent(ctx, &dir, access, tc).await {
Ok(output) => Ok::<NodeOutput, anyhow::Error>(NodeOutput {
id: node_id.clone(),
directive: dir,
output,
}),
Err(e) => {
warn!(node = %node_id, error = %e, "hive-mind node failed (isolated)");
Ok::<NodeOutput, anyhow::Error>(NodeOutput {
id: node_id,
directive: dir,
output: format!("[ERROR] {e}"),
})
}
}
}
})
.collect();
let results = try_join_all(handles).await?;
// Bounded concurrency: run at most MAX_CONCURRENT_NODES futures at once.
let mut stream = futures_util::stream::iter(tasks).buffer_unordered(MAX_CONCURRENT_NODES);
let mut results = Vec::with_capacity(cycle.directives.len());
while let Some(node) = stream.next().await {
results.push(node?);
}
Ok(results)
}
+4 -5
View File
@@ -316,13 +316,13 @@ fn handle_submit_input(state: &mut AppStateRest, text: String) {
in_flight: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
abort: state.abort_flag.clone(),
api_key: api_key.clone(),
model: state.settings.model.clone(),
model: zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config),
api_base: provider_cfg.as_ref().map(|cfg| cfg.api_base.clone()),
};
let client = std::sync::Arc::new(zesdex_infrastructure::llm::provider::LlmClient::new(
api_key,
state.settings.model.clone(),
zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config),
provider_cfg.map(|cfg| cfg.api_base.clone()),
));
@@ -465,14 +465,13 @@ fn handle_compact(state: &mut AppStateRest) {
.get(provider_name)
.cloned()
.unwrap_or_default();
let model = state.settings.model.clone();
let model = zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config);
let api_base = provider_cfg.map(|cfg| cfg.api_base.clone());
let client = zesdex_infrastructure::llm::provider::LlmClient::new(api_key, model, api_base);
if let Some(ref mut rt) = state.session_runtime {
let tokio_rt =
tokio::runtime::Runtime::new().expect("create tokio runtime for AI compaction");
let tokio_rt = zesdex_infrastructure::runtime::runtime();
if let Ok(()) = tokio_rt.block_on(
zesdex_application::agent::turn_service::compact_messages_with_ai(
&mut rt.messages,
+1 -1
View File
@@ -80,7 +80,7 @@ pub fn spawn_agent_turn(state: &mut AppStateRest, text: String) {
// ── Resolve provider configuration ─────────────────────────────────
let provider_name = &state.settings.provider;
let api_key = resolve_api_key(state, provider_name);
let model = state.settings.model.clone();
let model = zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config);
let api_base = resolve_api_base(state, provider_name);
// ── Build message list ─────────────────────────────────────────────