Compare commits

..
12 Commits
Author SHA1 Message Date
semantic-release-bot 88c2a7d000 chore(release): 1.21.2 [skip ci]
## [1.21.2](https://github.com/asepharyana/zesdex/compare/v1.21.1...v1.21.2) (2026-08-28)

### Performance Improvements

* **agent:** symbol index tak pegang mutex global saat rebuild I/O ([3b25e38](https://github.com/asepharyana/zesdex/commit/3b25e3898f9ba2fc1c58b991ad95a8c0bbe98601))
2026-08-28 07:58:48 +00:00
asepharyana 3b25e3898f perf(agent): symbol index tak pegang mutex global saat rebuild I/O
Sebelumnya SemanticSearch/ListSymbols/RebuildIndex menahan SYMBOL_INDEX
Mutex selama full `rebuild` (walk seluruh workspace, bisa detikan) +
selama search. Di main loop yang menjalankan read-only tools paralel,
semantic_search/list_symbols lain jadi BLOCK selama rebuild.

Refactor:
- Global berubah Mutex<Option<SymbolIndex>> -> OnceLock<Mutex<HashMap<
  workspace, SymbolIndex>>> — index per-workspace, jadi pencarian workspace B
  tidak mungkin bocor simbol stale dari A (workspace-awareness kini struktural,
  bukan hanya via needs_rebuild).
- ensure_symbol_index(workspace, force): rebuild dijalankan DI LUAR lock
  (mutex hanya dicek/insert/lookup singkat), lalu hasilnya di-swap-in di bawah
  short lock. Search/list/rebuild-report tak lagi memblock thread lain selama
  walk I/O. Per-workspace key menghilangkan race lintas-workspace dari skema
  swap tunggal.
- Test +1 (test_ensure_symbol_index_per_workspace_isolation): verifikasi dua
  workspace punya index independen, rebuild A tidak menimpa B.

Verifikasi: check/clippy -D warnings/fmt clean; test infra 64 (0 gagal).
2026-08-28 14:54:54 +07:00
semantic-release-bot d23d3855ec chore(release): 1.21.1 [skip ci]
## [1.21.1](https://github.com/asepharyana/zesdex/compare/v1.21.0...v1.21.1) (2026-08-28)

### Bug Fixes

* **agent:** semantic_search symbol index workspace-aware ([104af3a](https://github.com/asepharyana/zesdex/commit/104af3abb9e626c5d00ea87d523248d606d523ee))
2026-08-28 07:23:23 +00:00
asepharyana 104af3abb9 fix(agent): semantic_search symbol index workspace-aware
SymbolIndex global sudah melacak workspace_path tapi SemanticSearch dan
ListSymbols Cuma rebuild saat index kosong (is_empty). Akibat: setelah
mengindeks workspace A, mencari di workspace B diam-diam mengembalikan
simbol STALE dari A — menyesatkan coding agent (referensikan simbol yang
tidak ada di repo aktif).

Fix:
- Tambah SymbolIndex::needs_rebuild(workspace) — true bila index kosong
  ATAU workspace diminta beda dari yang ter-cache.
- Pakai di 2 call site (SemanticSearch::run, ListSymbols::run) menggantikan
  is_empty(), jadi pindah workspace otomatis trigger rebuild.
- test: +1 (test_needs_rebuild_workspace_aware — verifikasi flip workspace
  memicu rebuild bolak-balik A -> B -> A).

Catatan (bukan bug, dilaporkan): mutex SYMBOL_INDEX masih dipegang selama
full rebuild di run() — bottleneck saat semantic_search dipanggil paralel;
perbaikan butuh restrukturisasi double-checked rebuild, tak diubah di sini.

Verifikasi: check/clippy -D warnings/fmt clean; test infra 63 (0 gagal).
2026-08-28 14:19:22 +07:00
semantic-release-bot 62fa85867c chore(release): 1.21.0 [skip ci]
# [1.21.0](https://github.com/asepharyana/zesdex/compare/v1.20.2...v1.21.0) (2026-08-28)

### Features

* **agent:** hive-mind consensus synthesis pakai LLM nyata ([f75ff74](https://github.com/asepharyana/zesdex/commit/f75ff740ac2fa63340656e1e8215440f5e48063d))
2026-08-28 07:18:51 +00:00
asepharyana f75ff740ac feat(agent): hive-mind consensus synthesis pakai LLM nyata
synth_consensus selama ini Cuma concatenate output node lalu dilabeli
"Consensus" — tidak ada sintesis. Kini:

- Resolve kredensial LLM (provider/model/base_url/api_key) dari Store,
  sumber yang sama dgn execute_cycle.
- Kirim prompt sintesis ke model: minta distilasi node outputs jadi satu
  laporan konsensus berisi AGREEMENTS / CONFLICTS / KEY FINDINGS /
  RECOMMENDATION.
- Graceful fallback ke summary concatenation bila panggilan LLM gagal /
  output kosong, supaya sintesis konsensus tidak pernah merusak siklus
  hive-mind (konsisten dgn filosofi isolated-errors utk node).
- Batasi output per-node (MAX_NODE_OUTPUT_CHARS=4000, char-safe via
  truncate_chars) agar prompt tetap bounded.
- test: +2 (truncation char-safe pada output besar multi-byte; concat
  summary memuat semua node id).

Verifikasi: cargo check/clippy -D warnings/fmt clean; test infra 62 (0
gagal). Disk root sudah di-cargo clean (free 43.7GB, turun 98% -> 63%).
2026-08-28 14:15:03 +07:00
semantic-release-bot e349a35716 chore(release): 1.20.2 [skip ci]
## [1.20.2](https://github.com/asepharyana/zesdex/compare/v1.20.1...v1.20.2) (2026-08-28)

### Bug Fixes

* **agent:** recall search + memory_dir fallback + bersihkan dead llm_client ([f1f58b9](https://github.com/asepharyana/zesdex/commit/f1f58b9996f8eb58871d44fdd41c2d863874f253))
2026-08-28 05:57:17 +00:00
asepharyana f1f58b9996 fix(agent): recall search + memory_dir fallback + bersihkan dead llm_client
Hasil audit round 4 (workflow/hive_mind + memory + semantic_search).

- fix(memory): recall.search selama ini TIDAK pernah dipakai — tool
  mengiklankan keyword search di skema tapi run() cuma list semua nama.
  Kini search benar-benar memfilter (cocok di name/description/content,
  case-insensitive), + output 'No memories match' bila kosong.
- fix(memory): ToolCtxBuilder tidak punya setter memory_dir dan tak ada
  call-site yang mengisinya — remember/recall/forget memakai PathBuf kosong
  dan menulis memory ke CWD (bukan lokasi persisten). Tambah setter
  memory_dir + worktrees_dir, dan helper resolve_memory_dir() yang fallback
  ke Store::new().memory_dir bila ctx.memory_dir kosong; dipakai di ketiga
  tool memory.
- refactor(workflow): hapus LlmClient dummy di WorkflowRun (dibuat dengan
  API key kosong + model default + base_url default lalu tak pernah dipakai
  — execute_workflow menerimanya sebagai _llm_client). Kini execute_workflow
  tak ambil parameter tak terpakai; LLM asli tetap lewat execute_primitive
  yang resolve kredensial dengan benar.
- test: +2 (recall search memfilter; resolve_memory_dir fallback/eksplisit).

Catatan audit yang dilaporkan (belum difix): synth_consensus hanya
menggabungkan output (label Consensus menyesatkan, bukan sintesis LLM), dan
semantic_search memegang Mutex index global saat full rebuild (bottleneck
saat paralel) + index tidak workspace-aware.

PENTING (infra): disk root 100% saat kerja. Saya bebaskan ~4.6G dari /tmp +
cache aman (sekai*, verify-z, bun/npm cache). target/debug di repo = 38G —
rampah, perlu cargo clean + rebuild (jangan dibiarkan).
2026-08-28 12:53:32 +07:00
semantic-release-bot f4c02fd64e chore(release): 1.20.1 [skip ci]
## [1.20.1](https://github.com/asepharyana/zesdex/compare/v1.20.0...v1.20.1) (2026-08-28)

### Bug Fixes

* **agent:** subagent patuhi tool-calling contract + truncation char-safe ([e982cbe](https://github.com/asepharyana/zesdex/commit/e982cbeb041baea9cd500e2a29862a7eca9e6e17))
2026-08-28 04:47:39 +00:00
asepharyana e982cbeb04 fix(agent): subagent patuhi tool-calling contract + truncation char-safe
Hasil audit alur AI agent round 3 (fokus correctness & latent crash).

- fix(subagent): engine.rs sebelumnya mengeksekusi tool lalu push
  ChatMessage::tool hasil TANPA mendahuluinya dengan pesan assistant yang
  mendeklarasikan tool_calls → history malformed ([..., tool, tool,
  assistant(text)]). Kontrak OpenAI/Anthropic mensyaratkan pesan assistant
  (berisi tool_calls) sebelum hasil tool. Kini push response_msg
  (assistant + tool_calls + content) sebelum eksekusi, dan hapus push
  assistant content-only di akhir (agar tidak duplikat). Loop utama sudah
  benar; subagent kini selaras.
- fix(utils): &content[..1500] / &content[..1000] di build_rich_context
  dan &diff[..5000] di auto/engine.rs bisa panic saat indeks byte jatuh di
  tengah karakter multi-byte UTF-8 (emoji/CJK/panah). Tambah helper
  truncate_chars() yang memotong per karakter (char-safe) dan pakai di
  3 titik tersebut.
- test: +4 unit test truncate_chars (ASCII, potong, multibyte no-panic,
  emoji).

Catatan audit: subagent/auto (auto-review) & build_rich_context adalah dead
code (spawn_background_review & build_rich_context tidak pernah dipanggil).
Auto-review jangan diaktifkan asal (parser format teks rapuh + tanpa
verifikasi pasca-fix) — dilaporkan, bukan dicolokkan.
2026-08-28 11:43:42 +07:00
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
18 changed files with 905 additions and 108 deletions
+42
View File
@@ -1,3 +1,45 @@
## [1.21.2](https://github.com/asepharyana/zesdex/compare/v1.21.1...v1.21.2) (2026-08-28)
### Performance Improvements
* **agent:** symbol index tak pegang mutex global saat rebuild I/O ([3b25e38](https://github.com/asepharyana/zesdex/commit/3b25e3898f9ba2fc1c58b991ad95a8c0bbe98601))
## [1.21.1](https://github.com/asepharyana/zesdex/compare/v1.21.0...v1.21.1) (2026-08-28)
### Bug Fixes
* **agent:** semantic_search symbol index workspace-aware ([104af3a](https://github.com/asepharyana/zesdex/commit/104af3abb9e626c5d00ea87d523248d606d523ee))
# [1.21.0](https://github.com/asepharyana/zesdex/compare/v1.20.2...v1.21.0) (2026-08-28)
### Features
* **agent:** hive-mind consensus synthesis pakai LLM nyata ([f75ff74](https://github.com/asepharyana/zesdex/commit/f75ff740ac2fa63340656e1e8215440f5e48063d))
## [1.20.2](https://github.com/asepharyana/zesdex/compare/v1.20.1...v1.20.2) (2026-08-28)
### Bug Fixes
* **agent:** recall search + memory_dir fallback + bersihkan dead llm_client ([f1f58b9](https://github.com/asepharyana/zesdex/commit/f1f58b9996f8eb58871d44fdd41c2d863874f253))
## [1.20.1](https://github.com/asepharyana/zesdex/compare/v1.20.0...v1.20.1) (2026-08-28)
### Bug Fixes
* **agent:** subagent patuhi tool-calling contract + truncation char-safe ([e982cbe](https://github.com/asepharyana/zesdex/commit/e982cbeb041baea9cd500e2a29862a7eca9e6e17))
# [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)
Generated
+11 -11
View File
@@ -4862,7 +4862,7 @@ dependencies = [
[[package]]
name = "zesdex-api"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"argon2",
@@ -4885,7 +4885,7 @@ dependencies = [
[[package]]
name = "zesdex-application"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"base64",
@@ -4903,7 +4903,7 @@ dependencies = [
[[package]]
name = "zesdex-bootstrap"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"chrono",
@@ -4920,7 +4920,7 @@ dependencies = [
[[package]]
name = "zesdex-daemon"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"base64",
@@ -4944,7 +4944,7 @@ dependencies = [
[[package]]
name = "zesdex-domain"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"base64",
@@ -4960,7 +4960,7 @@ dependencies = [
[[package]]
name = "zesdex-gateway"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"axum",
@@ -4987,7 +4987,7 @@ dependencies = [
[[package]]
name = "zesdex-grpc"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"axum",
@@ -5004,7 +5004,7 @@ dependencies = [
[[package]]
name = "zesdex-infrastructure"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"argon2",
@@ -5052,7 +5052,7 @@ dependencies = [
[[package]]
name = "zesdex-tui"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"base64",
@@ -5078,7 +5078,7 @@ dependencies = [
[[package]]
name = "zesdex-web"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"axum",
@@ -5098,7 +5098,7 @@ dependencies = [
[[package]]
name = "zesdex-ws"
version = "1.19.4"
version = "1.21.0"
dependencies = [
"anyhow",
"axum",
+1 -1
View File
@@ -15,7 +15,7 @@ members = [
]
[workspace.package]
version = "1.19.6"
version = "1.21.2"
edition = "2021"
authors = ["asepharyana <superaseph@gmail.com>"]
+50 -2
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)]
@@ -341,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.
+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");
+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::*;
@@ -130,7 +130,7 @@ pub fn spawn_background_review(
);
format!(
"{}...\n[diff truncated at {} characters]",
&diff[..MAX_DIFF_CHARS],
crate::utils::truncate_chars(&diff, MAX_DIFF_CHARS),
MAX_DIFF_CHARS
)
} else {
+214 -27
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;
@@ -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 {
@@ -130,7 +220,7 @@ pub async fn run_agent(
.await?;
let content = response_msg.content.clone().unwrap_or_default();
let tool_calls = response_msg.tool_calls.unwrap_or_default();
let tool_calls = response_msg.tool_calls.clone().unwrap_or_default();
// If no tool calls, we're done — return content
if tool_calls.is_empty() {
@@ -139,12 +229,18 @@ 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);
// Push the assistant message (with its tool_calls) BEFORE executing
// so the tool-calling contract is honoured: tool results reference
// the calls declared in the preceding assistant message. Without
// this, the history is malformed (`[...tool, tool, assistant]`).
messages.push(response_msg);
debug!("Subagent executing tool: {tool_name}");
// 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);
for (id, tool_name, result) in results {
debug!("Subagent tool {tool_name} finished");
report_progress(
&tool_ctx,
@@ -155,15 +251,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 +258,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,15 +270,7 @@ pub async fn run_agent(
consecutive_errors = 0;
}
messages.push(ChatMessage::tool(
tc.id.clone(),
truncate_tool_output(result),
));
}
// Add assistant response if there was text content
if !content.is_empty() {
messages.push(ChatMessage::assistant(Some(content)));
messages.push(ChatMessage::tool(id, truncate_tool_output(result)));
}
}
@@ -208,3 +287,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"]);
}
}
+8
View File
@@ -66,6 +66,14 @@ impl ToolCtxBuilder {
self.session_dir = v;
self
}
pub fn memory_dir(mut self, v: PathBuf) -> Self {
self.memory_dir = v;
self
}
pub fn worktrees_dir(mut self, v: PathBuf) -> Self {
self.worktrees_dir = v;
self
}
pub fn workspaces(mut self, v: Vec<PathBuf>) -> Self {
self.workspaces = v;
self
@@ -43,7 +43,8 @@ impl Tool for Forget {
let name = crate::tools::arg_str(args, "name")?;
info!(name, "forget invoked");
let repo = crate::persistence::cms::memory_repo::MarkdownMemoryRepository::new();
repo.delete(&ctx.memory_dir, &name)?;
let memory_dir = crate::tools::memory::resolve_memory_dir(&ctx.memory_dir);
repo.delete(&memory_dir, &name)?;
info!(name, "memory deleted");
Ok(format!("Memory '{}' deleted", name))
}
@@ -1,5 +1,21 @@
//! Memory management tools — remember, recall, forget.
use std::path::PathBuf;
pub mod forget;
pub mod recall;
pub mod remember;
/// Resolve the directory the memory tools should read/write.
///
/// Prefer an explicitly-configured `ToolCtx.memory_dir`. If that is empty
/// (a `ToolCtx` is often built without setting `memory_dir`), fall back to
/// the canonical persistent memory location from `Store` so memories are not
/// silently written into the current working directory.
pub fn resolve_memory_dir(ctx_memory_dir: &std::path::Path) -> PathBuf {
if ctx_memory_dir.as_os_str().is_empty() {
zesdex_domain::core::Store::new().memory_dir
} else {
ctx_memory_dir.to_path_buf()
}
}
+121 -10
View File
@@ -45,21 +45,132 @@ impl Tool for Recall {
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let repo = crate::persistence::cms::memory_repo::MarkdownMemoryRepository::new();
let memory_dir = crate::tools::memory::resolve_memory_dir(&ctx.memory_dir);
let specific_name = args.get("name").and_then(|v| v.as_str());
let search = args.get("search").and_then(|v| v.as_str());
if let Some(name) = specific_name {
info!(name, "recall loading specific memory");
let memory = repo.load(&ctx.memory_dir, name)?;
Ok(serde_json::to_string_pretty(&memory)?)
} else {
info!("recall listing all memories");
let names = repo.list(&ctx.memory_dir)?;
if names.is_empty() {
info!("no memories found");
return Ok("No memories saved yet".to_string());
}
Ok(format!("Available memories:\n{}", names.join("\n")))
let memory = repo.load(&memory_dir, name)?;
return Ok(serde_json::to_string_pretty(&memory)?);
}
if let Some(query) = search {
let query = query.trim().to_lowercase();
info!(search = %query, "recall searching memories");
if query.is_empty() {
return Ok("Search query is empty".to_string());
}
let names = repo.list(&memory_dir)?;
let mut matches: Vec<String> = Vec::new();
for name in &names {
// Load each memory and match against name/description/content.
if let Ok(m) = repo.load(&memory_dir, name) {
let haystack =
format!("{} {} {}", m.name, m.description, m.content).to_lowercase();
if haystack.contains(&query) {
matches.push(m.name);
}
}
}
if matches.is_empty() {
return Ok(format!("No memories match '{query}'"));
}
return Ok(format!(
"Memories matching '{query}' ({}):\n{}",
matches.len(),
matches.join("\n")
));
}
info!("recall listing all memories");
let names = repo.list(&memory_dir)?;
if names.is_empty() {
info!("no memories found");
return Ok("No memories saved yet".to_string());
}
Ok(format!("Available memories:\n{}", names.join("\n")))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::tools::ToolCtxBuilder;
use zesdex_domain::cms::MemoryRepository;
fn save_mem(name: &str, description: &str, content: &str, dir: &std::path::Path) {
let repo = crate::persistence::cms::memory_repo::MarkdownMemoryRepository::new();
let memory = zesdex_domain::cms::Memory {
name: name.to_string(),
description: description.to_string(),
content: content.to_string(),
kind: "reference".to_string(),
created_at: 0,
updated_at: 0,
outcome: None,
lifecycle: "active".to_string(),
scope: None,
before_snippet: None,
after_snippet: None,
provenances: Vec::new(),
};
repo.save(dir, &memory).unwrap();
}
#[test]
fn search_filters_memories_by_keyword() {
let dir = std::env::temp_dir().join(format!("zdx-mem-test-{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&dir).unwrap();
save_mem(
"rust-concurrency",
"tokio spawn patterns",
"how to use async tasks",
&dir,
);
save_mem(
"docker-deploy",
"deploy via compose",
"container orchestration",
&dir,
);
let tool_ctx = ToolCtxBuilder::default().memory_dir(dir.clone()).build();
let args = serde_json::json!({ "search": "tokio" });
let out = Recall.run(&tool_ctx, &args).unwrap();
assert!(
out.contains("rust-concurrency"),
"should match rust-concurrency, got: {out}"
);
assert!(
!out.contains("docker-deploy"),
"docker-deploy should not match tokio"
);
// A query with no match reports so.
let no_match = Recall
.run(&tool_ctx, &serde_json::json!({ "search": "zzzznope" }))
.unwrap();
assert!(no_match.contains("No memories match"), "{no_match}");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn resolve_memory_dir_falls_back_to_store_when_empty() {
// An empty ToolCtx.memory_dir is resolved to the canonical Store path.
let resolved = crate::tools::memory::resolve_memory_dir(std::path::Path::new(""));
assert!(!resolved.as_os_str().is_empty());
assert!(
resolved.ends_with("memory"),
"expected memory dir, got {resolved:?}"
);
// An explicit memory_dir is preserved.
let explicit =
crate::tools::memory::resolve_memory_dir(std::path::Path::new("/tmp/custom-memory"));
assert_eq!(explicit, std::path::Path::new("/tmp/custom-memory"));
}
}
@@ -81,7 +81,8 @@ impl Tool for Remember {
};
let repo = crate::persistence::cms::memory_repo::MarkdownMemoryRepository::new();
repo.save(&ctx.memory_dir, &memory)?;
let memory_dir = crate::tools::memory::resolve_memory_dir(&ctx.memory_dir);
repo.save(&memory_dir, &memory)?;
info!(name, "memory saved");
Ok(format!("Memory '{}' saved", name))
+166 -25
View File
@@ -121,7 +121,66 @@ pub struct CodeSymbol {
}
/// The in-memory symbol index, shared via a global static.
static SYMBOL_INDEX: Mutex<Option<SymbolIndex>> = Mutex::new(None);
///
/// Keyed by workspace path: each workspace owns its own `SymbolIndex`, so
/// searching one repo never leaks stale symbols from another, and a rebuild
/// triggered for workspace A cannot clobber the index of B. The lock is only
/// ever held briefly (to check/insert/look up) — never across the I/O-heavy
/// walk in `rebuild`, which runs locally and is swapped in under a short lock.
static SYMBOL_INDEX: std::sync::OnceLock<Mutex<HashMap<String, SymbolIndex>>> =
std::sync::OnceLock::new();
/// Access the (lazily initialised) global per-workspace symbol index map.
fn symbol_index_map() -> &'static Mutex<HashMap<String, SymbolIndex>> {
SYMBOL_INDEX.get_or_init(|| Mutex::new(HashMap::new()))
}
/// Ensure the per-workspace symbol index is built, returning the symbol count.
///
/// If `force` is true, or the workspace has no cached (non-empty) index yet,
/// the index is rebuilt. The rebuild itself runs OUTSIDE the global lock
/// (the walk can take seconds on a large repo), then the result is stored
/// under a short lock so concurrent searches never block on the I/O. Returns
/// the number of symbols now cached for the workspace.
fn ensure_symbol_index(workspace: &str, force: bool) -> Result<usize> {
let ready = {
let map = symbol_index_map()
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
!force && map.get(workspace).is_some_and(|i| !i.is_empty())
};
if !ready {
// Rebuild locally, off the global lock (I/O heavy).
let mut fresh = SymbolIndex::new();
let count = fresh.rebuild(workspace)?;
// Swap in under a short lock; keep an existing non-empty index if a
// concurrent rebuild already populated this workspace.
let mut map = symbol_index_map()
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
if map.get(workspace).is_none_or(|i| i.is_empty()) {
map.insert(workspace.to_string(), fresh);
}
return Ok(count);
}
let map = symbol_index_map()
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
Ok(map.get(workspace).map_or(0, |i| i.len()))
}
/// Lock and return a borrow to the global per-workspace symbol index map.
///
/// The caller must have called [`ensure_symbol_index`] first, then looks up
/// its workspace key; the lookup is short and in-memory, so holding the guard
/// for the search is fine.
fn symbol_index() -> Result<std::sync::MutexGuard<'static, HashMap<String, SymbolIndex>>> {
symbol_index_map()
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))
}
// ---------------------------------------------------------------------------
// Language-specific regexes (lazily compiled)
@@ -270,6 +329,17 @@ impl SymbolIndex {
self.symbols.is_empty()
}
/// Returns true when the cached index must be rebuilt for the given
/// workspace — either because nothing has been indexed yet, or because the
/// requested workspace differs from the one the index was built for.
///
/// Without this, searching a *different* workspace after the first one
/// silently returns stale symbols from the previously indexed repo
/// (a misleading result for a coding agent).
pub fn needs_rebuild(&self, workspace: &str) -> bool {
self.is_empty() || self.workspace_path.as_deref() != Some(workspace)
}
pub fn len(&self) -> usize {
self.symbols.len()
}
@@ -1250,15 +1320,11 @@ impl Tool for SemanticSearch {
"semantic search"
);
let mut guard = SYMBOL_INDEX
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
let index = guard.get_or_insert_with(SymbolIndex::new);
if rebuild || index.is_empty() {
let count = index.rebuild(&workspace)?;
debug!(symbol_count = count, "symbol index rebuilt");
}
// Build (or load) the per-workspace index without holding the global
// lock across the I/O-heavy walk.
let _count = ensure_symbol_index(&workspace, rebuild)?;
let map = symbol_index()?;
let index = map.get(&workspace).expect("index should be ensured");
// Map kind filter to enum
let target_kind = match kind_filter {
@@ -1398,12 +1464,15 @@ impl Tool for RebuildIndex {
info!("rebuilding multi-language symbol index");
let mut guard = SYMBOL_INDEX
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
let index = guard.get_or_insert_with(SymbolIndex::new);
let count = index.rebuild(&workspace)?;
let by_lang = index.count_by_language();
// Force a rebuild of this workspace's index. The walk runs off the
// global lock (via ensure_symbol_index) so it cannot stall concurrent
// searches.
let count = ensure_symbol_index(&workspace, true)?;
let map = symbol_index()?;
let by_lang = match map.get(&workspace) {
Some(i) => i.count_by_language(),
None => Vec::new(),
};
let mut out = format!(
"Symbol index rebuilt successfully. {} symbols indexed.\n\n",
@@ -1497,15 +1566,11 @@ impl Tool for ListSymbols {
.map(|p| p.to_string_lossy().to_string())
.unwrap_or_else(|| ".".to_string());
let mut guard = SYMBOL_INDEX
.lock()
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
let index = guard.get_or_insert_with(SymbolIndex::new);
if rebuild || index.is_empty() {
let count = index.rebuild(&workspace)?;
info!(symbol_count = count, "symbol index rebuilt for list");
}
// Build (or load) the per-workspace index without holding the global
// lock across the I/O-heavy walk.
let _count = ensure_symbol_index(&workspace, rebuild)?;
let map = symbol_index()?;
let index = map.get(&workspace).expect("index should be ensured");
let target_lang = match lang_filter {
"rust" => Some(Language::Rust),
@@ -1746,4 +1811,80 @@ mod tests {
let index = SymbolIndex::new();
assert!(index.search("anything", 10).is_empty());
}
#[test]
fn test_needs_rebuild_workspace_aware() {
let ws_a = std::env::temp_dir().join(format!("ws_a_{}", uuid::Uuid::new_v4()));
let ws_b = std::env::temp_dir().join(format!("ws_b_{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&ws_a).unwrap();
std::fs::create_dir_all(&ws_b).unwrap();
std::fs::write(ws_a.join("a.rs"), "pub fn fn_in_a() {}\n").unwrap();
std::fs::write(ws_b.join("b.rs"), "pub fn fn_in_b() {}\n").unwrap();
let mut index = SymbolIndex::new();
let a = ws_a.to_string_lossy().to_string();
let b = ws_b.to_string_lossy().to_string();
// Fresh index: needs rebuild for any workspace.
assert!(index.needs_rebuild(&a));
// After rebuilding A, searching A needs no rebuild...
index.rebuild(&a).unwrap();
assert!(!index.needs_rebuild(&a));
// ...but searching B DOES (stale index otherwise).
assert!(
index.needs_rebuild(&b),
"workspace switch must trigger rebuild"
);
// Rebuilding B flips the cached workspace.
index.rebuild(&b).unwrap();
assert!(!index.needs_rebuild(&b));
assert!(index.needs_rebuild(&a));
std::fs::remove_dir_all(&ws_a).ok();
std::fs::remove_dir_all(&ws_b).ok();
}
#[test]
fn test_ensure_symbol_index_per_workspace_isolation() {
let ws_a = std::env::temp_dir().join(format!("iso_a_{}", uuid::Uuid::new_v4()));
let ws_b = std::env::temp_dir().join(format!("iso_b_{}", uuid::Uuid::new_v4()));
std::fs::create_dir_all(&ws_a).unwrap();
std::fs::create_dir_all(&ws_b).unwrap();
std::fs::write(ws_a.join("a.rs"), "pub fn only_in_a() {}\n").unwrap();
std::fs::write(ws_b.join("b.rs"), "pub fn only_in_b() {}\n").unwrap();
let a = ws_a.to_string_lossy().to_string();
let b = ws_b.to_string_lossy().to_string();
// Build A and B independently through the shared global helper.
let count_a = ensure_symbol_index(&a, false).unwrap();
assert!(
count_a >= 1,
"workspace A should index its fn, got {count_a}"
);
let count_b = ensure_symbol_index(&b, false).unwrap();
assert!(
count_b >= 1,
"workspace B should index its fn, got {count_b}"
);
// Rebuilding A must not have clobbered B and vice-versa.
let count_a_again = ensure_symbol_index(&a, true).unwrap();
assert!(count_a_again >= 1);
// Each workspace's cached index is independently correct.
{
let map = symbol_index().unwrap();
let idx_a = map.get(&a).unwrap();
assert!(!idx_a.search("only_in_a", 5).is_empty());
assert!(idx_a.search("only_in_b", 5).is_empty());
let idx_b = map.get(&b).unwrap();
assert!(!idx_b.search("only_in_b", 5).is_empty());
}
std::fs::remove_dir_all(&ws_a).ok();
std::fs::remove_dir_all(&ws_b).ok();
}
}
+1 -8
View File
@@ -9,7 +9,6 @@ use anyhow::Result;
use serde_json::{json, Value};
use tracing::{debug, info, instrument, warn};
use crate::llm::provider::LlmClient;
use crate::tools::{arg_str, Tool, ToolCtx};
use crate::workflow::engine::execution::execute_workflow;
use crate::workflow::hive_mind::cycle::execute_cycle;
@@ -60,14 +59,8 @@ impl Tool for WorkflowRun {
phase_names.join(", ")
);
let llm_client = LlmClient::new(
crate::llm::provider::DEFAULT_API_KEY.to_string(),
zesdex_domain::agent::defaults::DEFAULT_MODEL.to_string(),
None,
);
let rt = crate::runtime::runtime();
let result: Vec<String> =
rt.block_on(async { execute_workflow(&script, ctx, &llm_client).await })?;
let result: Vec<String> = rt.block_on(async { execute_workflow(&script, ctx).await })?;
info!(phase_count = result.len(), "Workflow completed");
Ok(format!(
+60 -3
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) {
@@ -268,7 +275,7 @@ pub fn build_rich_context(root: &Path) -> String {
let p = root.join(file);
if let Ok(content) = std::fs::read_to_string(&p) {
let snippet = if content.len() > 1500 {
format!("{}\n... (truncated)", &content[..1500])
format!("{}\n... (truncated)", truncate_chars(&content, 1500))
} else {
content
};
@@ -294,7 +301,7 @@ pub fn build_rich_context(root: &Path) -> String {
let readme_path = root.join("README.md");
if let Ok(content) = std::fs::read_to_string(&readme_path) {
let snippet = if content.len() > 1000 {
format!("{}\n... (truncated)", &content[..1000])
format!("{}\n... (truncated)", truncate_chars(&content, 1000))
} else {
content
};
@@ -341,3 +348,53 @@ pub fn build_rich_context(root: &Path) -> String {
ctx.trim_end().to_string()
}
/// Truncate a string to at most `max_chars` **characters**, never cutting a
/// multi-byte UTF-8 code point in half.
///
/// `&s[..n]` with `n` a raw byte index panics when `n` lands inside a
/// multi-byte character (e.g. an emoji, `→`, or CJK in a README/diff). This
/// helper slices on character boundaries so content is safely capped at a
/// byte budget while remaining valid UTF-8.
pub fn truncate_chars(s: &str, max_chars: usize) -> String {
if s.chars().count() <= max_chars {
return s.to_string();
}
s.chars().take(max_chars).collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn truncate_chars_leaves_short_strings_unchanged() {
assert_eq!(truncate_chars("short", 100), "short");
assert_eq!(truncate_chars("", 5), "");
}
#[test]
fn truncate_chars_cuts_to_max_chars() {
assert_eq!(truncate_chars("hello world", 5), "hello");
}
#[test]
fn truncate_chars_never_splits_multibyte_utf8() {
// 4 chars each: '→' is 3 bytes. Byte-slicing at 5 would panic; char
// slicing must not.
let s = "a→b→c→d";
let t = truncate_chars(s, 5);
assert_eq!(t, "a→b→c");
assert!(t.chars().count() <= 5);
// No replacement char must appear (valid UTF-8 preserved).
assert!(!t.contains('\u{FFFD}'));
}
#[test]
fn truncate_chars_handles_emoji() {
let s = "🚀🚀🚀🚀";
let t = truncate_chars(s, 2);
assert_eq!(t, "🚀🚀");
assert!(t.chars().count() == 2);
}
}
@@ -3,7 +3,6 @@
use anyhow::Result;
use tracing::{info, instrument};
use crate::llm::provider::LlmClient;
use crate::tools::ToolCtx;
use crate::workflow::engine::primitives::execute_primitive;
use zesdex_domain::workflow::WorkflowScript;
@@ -11,12 +10,8 @@ use zesdex_domain::workflow::WorkflowScript;
/// Execute each phase of a workflow script sequentially.
///
/// Flow: for each phase → execute_primitive → collect result.
#[instrument(skip(tool_ctx, _llm_client))]
pub async fn execute_workflow(
script: &WorkflowScript,
tool_ctx: &ToolCtx,
_llm_client: &LlmClient,
) -> Result<Vec<String>> {
#[instrument(skip(tool_ctx))]
pub async fn execute_workflow(script: &WorkflowScript, tool_ctx: &ToolCtx) -> Result<Vec<String>> {
info!(
"Executing workflow: {} ({} phases)",
script.name,
@@ -1,28 +1,128 @@
//! Consensus synthesis — reconciles multiple node outputs into one assessment.
//!
//! Flow: load settings → resolve LLM credentials → ask the model to distill the
//! node outputs into a single consensus (conflicts, agreements, key findings)
//! → return the synthesized text. If the LLM call fails for any reason, we
//! degrade gracefully to a concatenation-based summary so consensus synthesis
//! never breaks the surrounding hive-mind cycle (mirrors the isolated-errors
//! philosophy used for the nodes themselves).
use anyhow::Result;
use tracing::info;
use tracing::{info, warn};
use crate::tools::ToolCtx;
use zesdex_domain::cms::{AppConfigRepository, SettingsRepository};
use zesdex_domain::core::message::ChatMessage;
use zesdex_domain::core::Store;
use zesdex_domain::workflow::NodeOutput;
/// Synthesize a consensus from all node outputs.
use crate::persistence::{JsonAppConfigRepository, JsonSettingsRepository};
use crate::tools::ToolCtx;
/// Maximum characters of node output to feed into the synthesis prompt per node.
/// Keeps the prompt bounded so a huge/talkative node cannot blow up the request.
const MAX_NODE_OUTPUT_CHARS: usize = 4000;
/// Synthesize a consensus from all node outputs using the LLM.
///
/// Flow: combine node outputs → return consensus text.
/// Uses simple concatenation-based synthesis (avoids LLM call dependency).
/// Flow: combine node outputs → ask the model to reconcile them into a single
/// consensus → return the synthesized text. Falls back to a plain
/// concatenation summary if the LLM is unreachable or the call fails.
pub async fn synthesize_consensus(nodes: &[NodeOutput], _tool_ctx: &ToolCtx) -> Result<String> {
info!("Synthesizing consensus from {} nodes", nodes.len());
let combined = build_combined_body(nodes);
// 1. Resolve LLM credentials (same source of truth as execute_cycle).
let store = Store::new();
let settings = JsonSettingsRepository::new()
.load(&store.base_dir)
.unwrap_or_default();
let app_config = JsonAppConfigRepository::new()
.load(&store.base_dir)
.unwrap_or_default();
let (provider, model) =
crate::subagent::provider::resolve_subagent_provider(&settings, &app_config);
let base_url = app_config
.providers
.get(&provider)
.map(|p| p.api_base.clone())
.unwrap_or_else(|| zesdex_domain::agent::defaults::DEFAULT_API_BASE.to_string());
let api_key = crate::llm::provider::resolve_api_key(&settings, &app_config);
let client = crate::llm::provider::LlmClient::new(api_key, model, Some(base_url));
// 2. Build the synthesis prompt.
let system_msg = ChatMessage::system(
"You are a consensus synthesizer for a multi-agent hive mind. \
Several independent nodes analysed a problem and produced the outputs \
below. Distill them into ONE coherent consensus report with these \
sections:\n\
- AGREEMENTS: points multiple nodes converge on.\n\
- CONFLICTS: contradictory conclusions, with which node(s) support each side.\n\
- KEY FINDINGS: the most important, actionable takeaways.\n\
- RECOMMENDATION: a single recommended next action, or 'no clear consensus' \
if the outputs are too divergent.\n\
Be concise and factual. If a node errored, note it and ignore its content.\n\
Do not invent facts not present in the node outputs.",
);
let user_msg = ChatMessage::user(format!(
"Consolidate these {} node outputs into a single consensus:\n\n{}",
nodes.len(),
combined
));
// 3. Call the model and gracefully degrade on failure.
match call_consensus(&client, &[system_msg, user_msg]).await {
Ok(text) => {
let trimmed = text.trim();
if trimmed.is_empty() {
warn!("consensus LLM returned empty output; falling back to concat summary");
Ok(concat_summary(nodes))
} else {
Ok(format!(
"# Consensus Synthesis\n\n\
Nodes synthesized: {}\n\n{}",
nodes.len(),
trimmed
))
}
}
Err(e) => {
warn!(error = %e, "consensus LLM call failed; falling back to concat summary");
Ok(concat_summary(nodes))
}
}
}
/// Run the LLM consensus call and return the assistant text.
async fn call_consensus(
client: &crate::llm::provider::LlmClient,
messages: &[ChatMessage],
) -> Result<String> {
use zesdex_application::ports::ProviderService;
let (msg, _) = client.chat(messages, None, Some(1024), Some(0.3)).await?;
Ok(msg.content.unwrap_or_default())
}
/// Build the concatenated node-output body for the prompt.
fn build_combined_body(nodes: &[NodeOutput]) -> String {
let mut combined = String::new();
for node in nodes {
let output = crate::utils::truncate_chars(&node.output, MAX_NODE_OUTPUT_CHARS);
combined.push_str(&format!(
"\n## {}{}\n\n{}\n",
node.id, node.directive, node.output
"\n## {}{}\n{}\n",
node.id, node.directive, output
));
}
combined
}
Ok(format!(
"# Consensus Synthesis\n\
/// Fallback: a plain concatenation summary (the behaviour of the original stub).
fn concat_summary(nodes: &[NodeOutput]) -> String {
let combined = build_combined_body(nodes);
format!(
"# Consensus Synthesis\n\n\
Nodes synthesized: {}\n\n\
## Summary\n\
The following node outputs were collected:\n\
@@ -31,5 +131,47 @@ pub async fn synthesize_consensus(nodes: &[NodeOutput], _tool_ctx: &ToolCtx) ->
Review the individual node outputs above for detailed findings.",
nodes.len(),
combined
))
)
}
#[cfg(test)]
mod tests {
use super::*;
use zesdex_domain::workflow::NodeOutput;
fn node(id: &str, output: &str) -> NodeOutput {
NodeOutput {
id: id.to_string(),
directive: "directive".to_string(),
output: output.to_string(),
}
}
#[test]
fn build_combined_body_truncates_oversized_output() {
let long = "é".repeat(MAX_NODE_OUTPUT_CHARS + 500);
let body = build_combined_body(&[node("n1", &long)]);
// Must contain the header and a char-truncated (<= cap) payload without
// panicking on a multi-byte boundary.
assert!(body.contains("## n1"));
// Header "## n1 — directive\n" ~= 20 chars, so char count stays just above cap.
let char_count = body.chars().count();
assert!(
char_count <= MAX_NODE_OUTPUT_CHARS + 50,
"expected body near {MAX_NODE_OUTPUT_CHARS} chars, got {char_count}"
);
assert!(
char_count > 1000,
"expected a many-node output, got small: {char_count}"
);
}
#[test]
fn concat_summary_includes_all_node_ids() {
let nodes = vec![node("node-a", "a out"), node("node-b", "b out")];
let summary = concat_summary(&nodes);
assert!(summary.contains("node-a"));
assert!(summary.contains("node-b"));
assert!(summary.contains("Nodes synthesized: 2"));
}
}