Compare commits

...
32 Commits
Author SHA1 Message Date
semantic-release-bot 9d1544d799 chore(release): 1.22.0 [skip ci]
# [1.22.0](https://github.com/asepharyana/zesdex/compare/v1.21.2...v1.22.0) (2026-08-28)

### Features

* **agent:** wire auto-review (review_enabled no-op -> nyata) ([a27e815](https://github.com/asepharyana/zesdex/commit/a27e8151b39da4c8759e8e922ef132212327e413))
2026-08-28 08:31:47 +00:00
asepharyana a27e8151b3 feat(agent): wire auto-review (review_enabled no-op -> nyata)
review_enabled (default TRUE) selama ini no-op: tampil di TUI settings
overlay ("Review: true"), bisa di-toggle via command, TAPI spawn_background_review
tidak pernah dipanggil & Origin::Reviewer tidak pernah dikonstruksi. User
melihat "Review: true" padahal auto-review setelah edit tak pernah jalan.

Sekarang feature yang sudah dibangun penuh (subagent/auto/engine.rs: git diff
-> LLM review -> auto-fix HIGH/MEDIUM) di-wire:
- daemon/handler.rs: trigger setelah run_turn bila review_enabled; capture
  flag+creds SEBELUM api_key/provider_cfg di-move ke LlmClient.
- tui/turn.rs: sama, gated by state.settings.flags.review_enabled.
- ws/lib.rs: channel minimal tanpa settings -> review nyala tiap prompt
  (konsisten dgn default ON).

Aman: review fire-and-forget (tokio::spawn), get_git_diff skip bila no-change,
no-op bila bukan git repo (auto/engine). Creds dipakai = creds ter-resolve yg
sama dgn komposisi turn.

Verifikasi: check/clippy/fmt/test workspace hijau (0 error/warning/fail).
2026-08-28 15:27:33 +07:00
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
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
semantic-release-bot 3847c0e6fd chore(release): 1.19.1 [skip ci]
## [1.19.1](https://github.com/asepharyana/zesdex/compare/v1.19.0...v1.19.1) (2026-08-27)

### Performance Improvements

* **agent:** rombak alur AI agent — adaptif, hemat token, self-healing ([eac0443](https://github.com/asepharyana/zesdex/commit/eac0443c4c3b8bfcbefd4bad9554168fb6525b94))
2026-08-27 16:10:48 +00:00
asepharyana 5023e5dfa1 chore: sinkronkan Cargo.lock dengan versi 1.19.0 2026-08-27 23:06:57 +07:00
asepharyana eac0443c4c perf(agent): rombak alur AI agent — adaptif, hemat token, self-healing
Ganti explore phase MANDATORY (3 subagent tiap turn, boros) dengan
tool explore_codebase yang DIPUTUSKAN agent sendiri (lazy, token-aware):
- hapus ExploreService trait + with_explore + Phase 0 dari turn loop
- ExploreServiceImpl kini jadi tool 'explore_codebase' (1 context-scout
  subagent, read-only, cap output 4k chars)
- system prompt: instruksi TOKEN BUDGET (jawab langsung utk query simple,
  panggil explore_codebase sekali utk task kompleks)

Loop utama kini adaptif & self-healing:
- max_tokens adaptif (800/1600/4096 by request length) — bukan selalu 4096
- temperature 0.2 saat tool-calling, 0.7 utk final answer
- ErrorTracker: deteksi tool error berulang → inject recovery note,
  stop setelah 8 error total (bukan 50 iterasi sia-sia)
- auto-compact history > 60k chars sebelum LLM call
- tool output di-truncate ke 12k chars sebelum masuk konteks

Tambah 8 unit test (truncation, adaptive tokens, error tracker).
2026-08-27 23:06:37 +07:00
asepharyana 14f3eae62a a 2026-08-27 23:06:37 +07:00
semantic-release-bot aec53651ed chore(release): 1.19.0 [skip ci]
# [1.19.0](https://github.com/asepharyana/zesdex/compare/v1.18.4...v1.19.0) (2026-08-27)

### Features

* hapus fitur LSP bawaan (language server protocol) ([93f3c2a](https://github.com/asepharyana/zesdex/commit/93f3c2a3572511b5a84f244980ad71bd1b455e70))
2026-08-27 15:49:38 +00:00
asepharyana 93f3c2a357 feat: hapus fitur LSP bawaan (language server protocol)
Hapus seluruh pipeline LSP (client, manager, provisioner, dan 7 tool
lsp_*) dari codebase:

- apps/infrastructure/src/lsp/ (client.rs, manager.rs, provisioner/*)
- apps/infrastructure/src/tools/lsp/ (connect, disconnect, diagnostics,
  hover, completion, definition, references)
- ToolCtx/ToolCtxBuilder: hapus field lsp_manager
- Daemon state: hapus lsp_manager, lsp_provision_msgs, shutdown_lsp
- Registry: hapus registrasi 7 tool lsp_*
- Settings: hapus lsp_auto_provision + lsp_languages
- Agent definitions: hapus lsp_* dari allowed tools coder/reviewer
- Cargo: hapus dependency lsp-types (workspace + infra)
- Update dokumentasi mod + arch_audit forbidden list

Verifikasi: cargo check/clippy/test semua hijau (54 test), tidak ada
referensi lsp_* tersisa di luar CHANGELOG.
2026-08-27 22:38:21 +07:00
65 changed files with 2172 additions and 1908 deletions
+1
View File
@@ -7,3 +7,4 @@ package-lock.json
.superpowers/
docs/lesson/
.kilo/
.hermes/
+99
View File
@@ -1,3 +1,102 @@
# [1.22.0](https://github.com/asepharyana/zesdex/compare/v1.21.2...v1.22.0) (2026-08-28)
### Features
* **agent:** wire auto-review (review_enabled no-op -> nyata) ([a27e815](https://github.com/asepharyana/zesdex/commit/a27e8151b39da4c8759e8e922ef132212327e413))
## [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)
### 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)
### Performance Improvements
* **agent:** rombak alur AI agent — adaptif, hemat token, self-healing ([eac0443](https://github.com/asepharyana/zesdex/commit/eac0443c4c3b8bfcbefd4bad9554168fb6525b94))
# [1.19.0](https://github.com/asepharyana/zesdex/compare/v1.18.4...v1.19.0) (2026-08-27)
### Features
* hapus fitur LSP bawaan (language server protocol) ([93f3c2a](https://github.com/asepharyana/zesdex/commit/93f3c2a3572511b5a84f244980ad71bd1b455e70))
## [1.18.4](https://github.com/asepharyana/zesdex/compare/v1.18.3...v1.18.4) (2026-08-27)
Generated
+12 -45
View File
@@ -1088,15 +1088,6 @@ dependencies = [
"miniz_oxide",
]
[[package]]
name = "fluent-uri"
version = "0.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "17c704e9dbe1ddd863da1e6ff3567795087b1eb201ce80d8fa81162e1516500d"
dependencies = [
"bitflags 1.3.2",
]
[[package]]
name = "fnv"
version = "1.0.7"
@@ -1972,19 +1963,6 @@ version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154"
[[package]]
name = "lsp-types"
version = "0.97.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "53353550a17c04ac46c585feb189c2db82154fc84b79c7a66c96c2c644f66071"
dependencies = [
"bitflags 1.3.2",
"fluent-uri",
"serde",
"serde_json",
"serde_repr",
]
[[package]]
name = "mac_address"
version = "1.1.8"
@@ -3369,17 +3347,6 @@ dependencies = [
"serde_core",
]
[[package]]
name = "serde_repr"
version = "0.1.20"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "175ee3e80ae9982737ca543e96133087cbd9a485eecc3bc4de9c1a37b47ea59c"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.118",
]
[[package]]
name = "serde_urlencoded"
version = "0.7.1"
@@ -4895,7 +4862,7 @@ dependencies = [
[[package]]
name = "zesdex-api"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"argon2",
@@ -4918,11 +4885,12 @@ dependencies = [
[[package]]
name = "zesdex-application"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"base64",
"chrono",
"futures-util",
"serde",
"serde_json",
"sha2 0.11.0",
@@ -4935,7 +4903,7 @@ dependencies = [
[[package]]
name = "zesdex-bootstrap"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"chrono",
@@ -4952,7 +4920,7 @@ dependencies = [
[[package]]
name = "zesdex-daemon"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"base64",
@@ -4976,7 +4944,7 @@ dependencies = [
[[package]]
name = "zesdex-domain"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"base64",
@@ -4992,7 +4960,7 @@ dependencies = [
[[package]]
name = "zesdex-gateway"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"axum",
@@ -5019,7 +4987,7 @@ dependencies = [
[[package]]
name = "zesdex-grpc"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"axum",
@@ -5036,7 +5004,7 @@ dependencies = [
[[package]]
name = "zesdex-infrastructure"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"argon2",
@@ -5055,7 +5023,6 @@ dependencies = [
"infer",
"jsonwebtoken",
"libc",
"lsp-types",
"nucleo-matcher",
"percent-encoding",
"pulldown-cmark",
@@ -5085,7 +5052,7 @@ dependencies = [
[[package]]
name = "zesdex-tui"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"base64",
@@ -5111,7 +5078,7 @@ dependencies = [
[[package]]
name = "zesdex-web"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"axum",
@@ -5131,7 +5098,7 @@ dependencies = [
[[package]]
name = "zesdex-ws"
version = "1.18.2"
version = "1.21.1"
dependencies = [
"anyhow",
"axum",
+1 -2
View File
@@ -15,7 +15,7 @@ members = [
]
[workspace.package]
version = "1.18.4"
version = "1.22.0"
edition = "2021"
authors = ["asepharyana <superaseph@gmail.com>"]
@@ -61,7 +61,6 @@ ignore = "0.4"
nucleo-matcher = "0.3"
futures-util = "0.3"
rmcp = { version = "2.2", default-features = false, features = ["client", "transport-child-process", "transport-streamable-http-client-reqwest", "macros"] }
lsp-types = "0.97"
tiktoken-rs = "0.12"
similar = "3"
syntect = { version = "5", default-features = false, features = ["default-fancy"] }
-250
View File
@@ -1,250 +0,0 @@
# Zesdex — Autonomous AI Coding Agent
Zesdex is an autonomous AI coding agent with a Terminal UI (TUI). It acts as an
OpenAI/Anthropic-compatible LLM client wrapped in a tool-use harness with **37
built-in tools** — file operations, git, shell execution, LSP integration, MCP,
subagent orchestration, and more.
```
┌──────────────────────────────────────────────────────────────┐
│ Mode Selector │
│ TUI (default) ─── Daemon ─── Attach ─── API ─── WS/gRPC/Web │
└──────────────────────────────────────────────────────────────┘
```
---
## Quick Start
```bash
# Run the TUI (default mode)
cargo run
# Run the REST API server
cargo run -- --api --api-port 8080
# Run in daemon mode (background + IPC)
cargo run -- --daemon
# Attach TUI to a running daemon session
cargo run -- --attach <session-id>
# Seed initial data (first run)
cargo run --bin bootstrap
```
### Prerequisites
- **Rust** 1.81+ (edition 2021)
- **Linux** or **macOS** (Unix domain sockets required for daemon mode)
- An **API key** for an OpenAI/Anthropic-compatible LLM provider (set via
settings or environment variable)
---
## Modes
| Flag | Mode | Description |
|------|------|-------------|
| *(none)* | **TUI** | Full terminal UI with chat, overlays, and agent loop in one process |
| `--daemon` | **Daemon** | Background daemon with IPC socket; clients attach separately |
| `--attach <id>` | **Attach** | Connect TUI to an existing daemon session via Unix socket |
| `--api` | **REST API** | HTTP server with session management and chat endpoints |
| `--ws` | **WebSocket** | WebSocket server for real-time communication |
| `--grpc` | **gRPC** | gRPC server for programmatic access |
| `--web` | **Web** | Serves the web frontend |
| `--api-port`, `--ws-port`, `--grpc-port`, `--web-port` | *(ports)* | Configure server ports (defaults: 8080, 8081, 50051, 3000) |
---
## Architecture
### Clean Architecture Layering
```
apps/
├── domain/ # Pure entities, value objects, repository/service traits
│ # Zero framework deps — only serde + chrono + uuid
├── application/ # Use-case services (auth, sessions, conversations, memory)
│ # Depends only on domain-layer trait interfaces
├── infrastructure/ # All I/O: LLM client, IPC, persistence, LSP, MCP, tools
│ # Implements domain/application port interfaces
└── interfaces/ # Entry points
├── tui/ # Ratatui terminal UI
├── api/ # Axum REST API
├── daemon/ # Unix socket daemon + client
├── ws/ # WebSocket server
├── grpc/ # gRPC server
└── web/ # Web frontend (static file server)
```
### Tool System
37 tools across 9 categories:
| Category | Tools |
|----------|-------|
| **File System** | `read`, `write`, `edit`, `delete`, `dir_list`, `dir_cache_update` |
| **Shell** | `bash`, `bash_interactive`, `bash_kill`, `bash_output` |
| **Git** | `git_operator`, `git_cred`, `git_worktree` |
| **Search** | `search`, `grep`, `glob`, `semantic_search` |
| **LSP** | `lsp_connect`, `lsp_hover`, `lsp_completion`, `lsp_definition`, `lsp_references`, `lsp_diagnostics`, `lsp_disconnect` |
| **Memory** | `remember`, `recall`, `forget` |
| **Workflow** | `spawn_agents`, `spawn_pipeline`, `plan`, `sequential_think`, `hive_mind` |
| **Utility** | `todo_write`, `todo_finish`, `pong`, `cd` |
| **Background** | Background bash jobs with `cancel/status/list` |
Each tool implements the `Tool` trait:
```rust
pub trait Tool: Send + Sync {
fn name(&self) -> &'static str;
fn description(&self) -> &'static str;
fn parameters(&self) -> Value;
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String>;
}
```
### Hive Mind Orchestration
The multi-agent orchestration system compiles a **cognitive cycle plan** per
task — ordered cycles of parallel processing nodes. Each node has a directive
and an **access tier** (`read` / `write` / `full`). Node outputs merge into a
shared collective state in real time, and a final **consensus synthesis**
produces the unified result.
- **Auto-trigger**: Complex requests automatically use the hive mind
- **Manual entry**: The `hive_mind` tool lets the LLM specify cycles explicitly
- **Live progress**: TUI panel shows each node's status and current tool
- **Guaranteed docs**: Every convergence writes to `docs/runs/`
### IPC Protocol (Daemon Mode)
```
┌──────────┐ Unix socket ┌──────────┐
│ Client │ ◄──────────────► │ Daemon │
│ (TUI) │ length-prefixed│ │
└──────────┘ serde_json └──────────┘
Frame format: [4-byte BE length][JSON payload]
```
The daemon holds `AppStateRest` and drives the agent loop. Clients are stateless
renderers that receive full state snapshots after each action.
---
## Built-in Features
| Feature | Description |
|---------|-------------|
| **LLM Provider** | OpenAI/Anthropic-compatible API (streaming + non-streaming) with automatic retry and fallback |
| **Tool Harness** | Safety-gated tool execution with graduated review checks |
| **Subagents** | Auto-inline review, background test-gen, arch-review, security-review |
| **OAuth 2.0** | PKCE flow for LLM provider authentication |
| **MCP** | Model Context Protocol server management (stdio + HTTP transport) |
| **LSP** | Language Server Protocol integration (completion, hover, diagnostics, references) |
| **Session Mgmt** | SQLite-persisted sessions with lock-based concurrency control |
| **Memory** | File-based memory system with frontmatter metadata |
| **Edit Log** | Append-only edit history with configurable retention |
| **Rate Limiting** | Sliding-window per-client rate limiter |
| **JWT Auth** | HS256 JWT access/refresh tokens (API mode) |
| **Password Auth** | Argon2 password hashing with pepper |
| **OAuth Loopback** | Localhost HTTP server for OAuth redirect capture |
| **Background Jobs** | Long-running shell jobs with cancellation and output collection |
| **Settings** | JSON-persisted settings with hot-reload |
---
## TUI Overlays
16 overlays accessible from the terminal UI:
| Overlay | Purpose |
|---------|---------|
| Chat Input | Main input bar with autocomplete |
| Bash Panel | Interactive shell panel |
| File Editor | Built-in file editor |
| Effort Selector | LLM reasoning effort selector |
| Help | Keybindings reference |
| Key Input | Custom key binding configuration |
| Learning | Lesson viewer |
| Loading | Generating spinner |
| MCP Manager | MCP server management |
| Model Selector | LLM model picker |
| Quit Confirm | Exit confirmation dialog |
| Rewind | Message/history rewind |
| Settings | Settings panel |
| Todo | Task/TODO list |
| Usage | Token usage statistics |
| Workflow | Hive-mind node progress |
---
## Data & Persistence
All data lives under the platform's data directory (`~/.local/share/zesdex/`):
```
~/.local/share/zesdex/
├── settings.json # User settings (provider, model, keys)
├── app_config.json # Provider definitions (endpoints, env vars)
├── sessions/ # Chat sessions (one subdirectory per session)
│ └── <uuid>/
│ ├── session.json # Session metadata
│ ├── messages.jsonl # Message log
│ └── .lock # Session lock file
└── memories/ # Memory files with frontmatter metadata
└── *.md
```
---
## Development
```bash
# Build all crates
cargo build
# Run all unit tests (8 tests across 11 crates)
cargo test
# Run clippy linting
cargo clippy --all-targets
# Run with verbose logging
RUST_LOG=debug cargo run
```
### Workspace Crates
| Crate | Path | Layer |
|-------|------|-------|
| `zesdex-domain` | `apps/domain/` | Pure domain entities & traits |
| `zesdex-application` | `apps/application/` | Use-case services |
| `zesdex-infrastructure` | `apps/infrastructure/` | All I/O & tool implementations |
| `zesdex-tui` | `apps/interfaces/tui/` | Ratatui terminal interface |
| `zesdex-api` | `apps/interfaces/api/` | Axum REST API |
| `zesdex-daemon` | `apps/interfaces/daemon/` | Unix socket daemon |
| `zesdex-ws` | `apps/interfaces/ws/` | WebSocket server |
| `zesdex-grpc` | `apps/interfaces/grpc/` | gRPC server |
| `zesdex-web` | `apps/interfaces/web/` | Web frontend |
| `zesdex-gateway` | `apps/gateway/` | CLI entry point & dispatcher |
| `zesdex-bootstrap` | `apps/bootstrap/` | Initial data seeder |
### Code Map
Detailed architecture documentation is in `docs/CODEMAPS/`:
| File | Covers |
|------|--------|
| `docs/CODEMAPS/architecture.md` | System layout, process modes, data flow |
| `docs/CODEMAPS/backend.md` | Provider, OAuth, IPC, workflow engine, MCP, LSP, review |
| `docs/CODEMAPS/frontend.md` | TUI render pipeline, 16 overlays, toasts, input handling |
| `docs/CODEMAPS/data.md` | Persistence, SQLite msglog, memory files, settings/config |
| `docs/CODEMAPS/dependencies.md` | All Rust crates and external services |
---
## License
See `CHANGELOG.md` for release history.
+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
-57
View File
@@ -1,57 +0,0 @@
//! Mandatory explore phase — spawns parallel subagents to discover context
//! before the main agent begins its turn.
//!
//! # Flow
//!
//! Before the main agent's LLM loop, [`ExploreService::explore`] dispatches
//! at least 3 subagents in parallel (code-structure scan, symbol-index query,
//! semantic-context search). Their findings are consolidated into a single
//! system message that is prepended to the conversation.
//!
//! # Why mandatory
//!
//! Without structured exploration the main agent works from an empty context
//! window. The explore phase guarantees that every turn starts with a compact
//! snapshot of what the codebase contains and where relevant code lives.
use anyhow::Result;
use std::collections::VecDeque;
use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use zesdex_domain::agent::TurnEvent;
/// The consolidated output of an explore phase — a set of system-level
/// context messages injected before the main agent prompt.
#[derive(Debug, Clone)]
pub struct ExploreOutput {
/// One or more system messages summarising what the explore subagents
/// discovered. Prepended to the conversation by the turn service.
pub context_messages: Vec<String>,
/// Short human-readable summary of what was explored.
pub summary: String,
}
/// Service trait for the mandatory pre-turn exploration phase.
///
/// Implementors spawn ≥3 parallel subagents, each analysing a different
/// aspect of the workspace, and return a consolidated summary.
///
/// # Object safety
///
/// This trait is `dyn`-safe — it returns `Pin<Box<dyn Future>>` so it can
/// be stored as `Arc<dyn ExploreService>`.
pub trait ExploreService: Send + Sync {
/// Run the explore phase.
///
/// `query` — the user's current input phrase.
/// `workspace_root` — absolute path to the workspace root.
/// `turn_events` — shared event queue for TUI updates.
/// Returns structured context messages and a summary blob.
fn explore<'a>(
&'a self,
query: &'a str,
workspace_root: &'a str,
turn_events: &'a Arc<Mutex<VecDeque<TurnEvent>>>,
) -> Pin<Box<dyn Future<Output = Result<ExploreOutput>> + Send + 'a>>;
}
+11 -2
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.
@@ -19,8 +30,6 @@ pub trait AgentTurnService: Send + Sync {
fn run_turn(&self, params: AgentTurnParams) -> impl Future<Output = Result<()>> + Send;
}
pub mod explore;
pub mod turn_service;
pub use explore::{ExploreOutput, ExploreService};
pub use turn_service::{compact_messages_with_ai, AgentTurnServiceImpl};
+461 -89
View File
@@ -5,14 +5,31 @@ 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::{ExploreService, ToolExecutor};
use super::ToolExecutor;
use crate::ports::ProviderService;
/// Maximum tool-call iterations per agent turn before forcing termination.
const MAX_TURN_ITERATIONS: u32 = 50;
/// Maximum number of consecutive identical tool errors before the loop
/// injects a recovery note and forces a different approach.
const MAX_CONSECUTIVE_TOOL_ERRORS: usize = 3;
/// Total tool-call errors tolerated per turn before the loop is stopped.
const MAX_TOTAL_TOOL_ERRORS: usize = 8;
/// Ceiling for a single tool-result message inserted into context.
///
/// Tool outputs can be huge (read / semantic_search). Truncating keeps the
/// context window from exploding while preserving the important head.
const TOOL_OUTPUT_MAX_CHARS: usize = 12_000;
/// Total conversation characters that trigger auto-compaction before the
/// next LLM call.
const AUTO_COMPACT_CHARS: usize = 60_000;
// ---------------------------------------------------------------------------
// Helper: push a TurnEvent onto the shared queue.
// ---------------------------------------------------------------------------
@@ -51,6 +68,129 @@ fn make_stream_callback(
})
}
// ---------------------------------------------------------------------------
// Helper: truncate a long tool output before it enters the conversation
// context. Preserves the head and appends a clear truncation marker.
// ---------------------------------------------------------------------------
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
}
// ---------------------------------------------------------------------------
// Helper: adaptive generation parameters.
// ---------------------------------------------------------------------------
/// Pick a `max_tokens` budget for the turn's next LLM call based on the
/// length of the user's request. Short requests need far fewer tokens than
/// the current hardcoded 4096 — big savings on small tasks.
fn adaptive_max_tokens(request_len: usize) -> u32 {
if request_len <= 80 {
800
} else if request_len <= 400 {
1600
} else {
4096
}
}
/// Sum the character length of the conversation (user + assistant +
/// tool content) as a cheap proxy for context size.
fn conversation_chars(messages: &[ChatMessage]) -> usize {
messages
.iter()
.map(|m| m.content.as_deref().map(str::len).unwrap_or(0))
.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)]
struct ErrorTracker {
consecutive: usize,
total: usize,
last_tool: String,
last_error: String,
}
impl ErrorTracker {
fn record(&mut self, tool_name: &str, error: &str, messages: &mut Vec<ChatMessage>) {
if self.last_tool == tool_name {
self.consecutive += 1;
} else {
self.consecutive = 1;
}
self.last_tool = tool_name.to_string();
self.last_error = error.to_string();
self.total += 1;
// Inject a recovery note once the same tool keeps failing.
if self.consecutive >= MAX_CONSECUTIVE_TOOL_ERRORS
&& !messages.iter().any(|m| {
m.content
.as_deref()
.is_some_and(|c| c.contains("[System note]"))
})
{
messages.push(ChatMessage::system(
zesdex_domain::agent::prompt::error_recovery_note(tool_name, error),
));
}
}
fn should_stop(&self) -> bool {
self.consecutive >= MAX_CONSECUTIVE_TOOL_ERRORS * 2 || self.total >= MAX_TOTAL_TOOL_ERRORS
}
}
// ---------------------------------------------------------------------------
// Helper: execute a single tool call, push events, return the result string.
// ---------------------------------------------------------------------------
@@ -71,6 +211,7 @@ async fn execute_tool_call<T: ToolExecutor>(
};
let is_error = output.starts_with("Error:");
let output = truncate_tool_output(output);
push_event(
turn_events,
@@ -86,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.
// ---------------------------------------------------------------------------
@@ -108,19 +296,19 @@ fn emit_usage(turn_events: &Arc<Mutex<VecDeque<TurnEvent>>>, usage: Option<(u64,
/// Service implementation for executing an agent turn asynchronously.
///
/// # Explore phase
///
/// Before the main LLM loop begins, [`AgentTurnServiceImpl`] runs a mandatory
/// explore phase that spawns ≥3 parallel subagents (code structure, symbol
/// index, semantic context) and injects their consolidated findings as a
/// system message. See [`ExploreService`] for the trait contract.
/// The turn loop is adaptive and token-aware:
/// - No mandatory explore phase — the *agent* decides when to call the
/// `explore_codebase` tool (see the main prompt), so simple queries skip
/// exploration entirely.
/// - `max_tokens` / `temperature` adapt to the request length and phase.
/// - Repeated tool errors trigger a system recovery note and eventually
/// stop the loop instead of burning iterations.
/// - Tool outputs are truncated before entering context.
/// - Oversized histories are auto-compacted before the next LLM call.
pub struct AgentTurnServiceImpl<P: ProviderService, T: ToolExecutor> {
provider: Arc<P>,
tool_executor: Arc<T>,
tool_defs: Vec<ToolDef>,
/// Optional explore-phase service. When `Some`, the explore phase runs
/// before every turn; when `None` it is skipped (tests, daemon mode).
explore_service: Option<Arc<dyn ExploreService>>,
}
impl<P: ProviderService, T: ToolExecutor> AgentTurnServiceImpl<P, T> {
@@ -129,19 +317,9 @@ impl<P: ProviderService, T: ToolExecutor> AgentTurnServiceImpl<P, T> {
provider,
tool_executor,
tool_defs,
explore_service: None,
}
}
/// Attach an optional explore-phase service.
///
/// When set, every call to `run_turn` will first run the explore phase
/// and inject the consolidated context as a system message.
pub fn with_explore(mut self, service: Arc<dyn ExploreService>) -> Self {
self.explore_service = Some(service);
self
}
/// Execute a single LLM call with the current message list, handling
/// streaming events and error reporting.
async fn call_llm(
@@ -149,6 +327,8 @@ impl<P: ProviderService, T: ToolExecutor> AgentTurnServiceImpl<P, T> {
messages: &[ChatMessage],
abort: &Arc<AtomicBool>,
turn_events: &Arc<Mutex<VecDeque<TurnEvent>>>,
max_tokens: u32,
temperature: f32,
) -> Result<(ChatMessage, Option<(u64, u64)>), String> {
let on_event = make_stream_callback(abort, turn_events);
@@ -156,13 +336,39 @@ impl<P: ProviderService, T: ToolExecutor> AgentTurnServiceImpl<P, T> {
.chat_stream(
messages,
Some(self.tool_defs.clone()),
Some(4096),
Some(0.7),
Some(max_tokens),
Some(temperature),
on_event,
)
.await
.map_err(|e| format!("LLM error: {e}"))
}
/// Auto-compact the history in place if it exceeds the threshold.
///
/// Runs at most once per turn. Skips the synthetic system prompt that
/// this service inserts at index 0.
async fn auto_compact_if_needed(&self, messages: &mut Vec<ChatMessage>) {
if conversation_chars(messages) <= AUTO_COMPACT_CHARS {
return;
}
// Keep the system prompt (index 0) out of compaction.
let sys = messages[0].clone();
let mut rest: Vec<ChatMessage> = messages.drain(1..).collect();
let before = rest.len();
if let Err(e) = super::compact_messages_with_ai(&mut rest, self.provider.as_ref()).await {
warn!("auto-compact failed (non-fatal): {e}");
}
info!(
"auto-compacted history: {} messages -> {}",
before,
rest.len()
);
let mut rebuilt = Vec::with_capacity(rest.len() + 1);
rebuilt.push(sys);
rebuilt.extend(rest);
*messages = rebuilt;
}
}
impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnServiceImpl<P, T> {
@@ -173,73 +379,34 @@ impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnS
params.model
);
// ── Phase 0: Mandatory explore ──────────────────────────────────
// Spawn ≥3 parallel subagents to discover code structure, symbols,
// and semantic context. The consolidated summary is injected as a
// system message before the main agent prompt.
if let Some(ref explorer) = self.explore_service {
// Determine workspace root from the first message's context or
// the first workspace root in params.
let user_query = params
.messages
.last()
.map(|m| m.content.clone().unwrap_or_default())
.unwrap_or_default();
let workspace_root = params
.workspace_roots
.first()
.map(|p| p.to_string_lossy().to_string())
.unwrap_or_else(|| ".".to_string());
push_event(
&params.turn_events,
TurnEvent::SystemNote {
kind: "info".into(),
message: "🔍 Exploring codebase structure...".into(),
},
);
match explorer
.explore(&user_query, &workspace_root, &params.turn_events)
.await
{
Ok(output) => {
// Insert each context message as a system message.
// They go at index 0 and are removed after the turn
// like the main agent prompt.
for ctx_msg in &output.context_messages {
params
.messages
.insert(0, ChatMessage::system(ctx_msg.clone()));
}
info!(
"Explore phase complete: {} context messages, {}",
output.context_messages.len(),
output.summary
);
}
Err(e) => {
warn!("Explore phase failed (non-fatal): {e}");
push_event(
&params.turn_events,
TurnEvent::SystemNote {
kind: "warn".into(),
message: format!("Explore phase failed: {e}"),
},
);
}
}
}
// Insert system prompt at position 0 once and keep it there for the
// entire turn, avoiding per-iteration clones of the full message list.
// It is removed before emitting the Compacted event so persistence
// does not store the prompt redundantly.
// 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.
let request_len = params
.messages
.last()
.and_then(|m| m.content.as_deref())
.map(str::len)
.unwrap_or(0);
let mut errors = ErrorTracker::default();
// Track whether the previous call produced tool calls — used to
// lower temperature once the agent starts producing a final answer.
let mut saw_tool_calls = false;
for iteration in 0..MAX_TURN_ITERATIONS {
// ── Check abort flag ────────────────────────────────────────
if params.abort.load(Ordering::SeqCst) {
@@ -254,15 +421,40 @@ impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnS
break;
}
if errors.should_stop() {
push_event(
&params.turn_events,
TurnEvent::SystemNote {
kind: "warn".into(),
message: "Stopping: repeated tool errors without progress".into(),
},
);
break;
}
debug!("agent turn iteration {iteration}");
// ── Auto-compact oversized history before the LLM call ─────
self.auto_compact_if_needed(&mut params.messages).await;
// ── Adaptive generation parameters ─────────────────────────
let max_tokens = adaptive_max_tokens(request_len);
// Lower temperature while the agent is still choosing tools to
// keep tool selection deterministic; raise it for the final
// free-form answer.
let temperature = if saw_tool_calls { 0.2 } else { 0.7 };
// ── Stream start + call LLM ─────────────────────────────────
push_event(&params.turn_events, TurnEvent::StreamStart);
// Uses params.messages directly (sys_msg[0] already in place
// from the insert above) — no per-iteration clone needed.
let result = self
.call_llm(&params.messages, &params.abort, &params.turn_events)
.call_llm(
&params.messages,
&params.abort,
&params.turn_events,
max_tokens,
temperature,
)
.await;
match result {
@@ -283,13 +475,49 @@ impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnS
break;
}
saw_tool_calls = true;
params.messages.push(assistant_msg);
// ── Execute each tool call ──────────────────────────
//
// 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 output =
execute_tool_call(self.tool_executor.as_ref(), &params.turn_events, tc)
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);
}
params
.messages
.push(ChatMessage::tool(tc.id.clone(), output));
@@ -370,3 +598,147 @@ pub async fn compact_messages_with_ai<P: ProviderService>(
}
}
}
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn truncate_short_output_is_unchanged() {
let out = "short".to_string();
assert_eq!(truncate_tool_output(out.clone()), out);
}
#[test]
fn truncate_long_output_preserves_head_and_marks_cut() {
let long = "x".repeat(TOOL_OUTPUT_MAX_CHARS + 500);
let truncated = truncate_tool_output(long.clone());
assert!(truncated.len() < long.len());
assert!(truncated.contains("...[truncated"));
assert!(truncated.starts_with("xxx"));
}
#[test]
fn adaptive_max_tokens_scales_with_request_len() {
assert_eq!(adaptive_max_tokens(10), 800);
assert_eq!(adaptive_max_tokens(200), 1600);
assert_eq!(adaptive_max_tokens(5000), 4096);
}
#[test]
fn error_tracker_injects_recovery_note_after_repeats() {
let mut tracker = ErrorTracker::default();
let mut messages: Vec<ChatMessage> = Vec::new();
tracker.record("read", "Error: File not found", &mut messages);
tracker.record("read", "Error: File not found", &mut messages);
assert!(!tracker.should_stop());
// Third consecutive failure → recovery note injected.
tracker.record("read", "Error: File not found", &mut messages);
assert!(messages.iter().any(|m| m
.content
.as_deref()
.is_some_and(|c| c.contains("[System note]"))));
}
#[test]
fn error_tracker_stops_after_too_many_errors() {
let mut tracker = ErrorTracker::default();
let mut messages: Vec<ChatMessage> = Vec::new();
for i in 0..MAX_TOTAL_TOOL_ERRORS {
tracker.record("bash", &format!("Error: boom {i}"), &mut messages);
}
assert!(tracker.should_stop());
}
#[test]
fn conversation_chars_sums_content_only() {
let messages = vec![
ChatMessage::system("sys".to_string()),
ChatMessage::user("hello world".to_string()),
ChatMessage::tool("id".to_string(), "output".to_string()),
];
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)));
}
}
+1 -1
View File
@@ -52,5 +52,5 @@ pub use cms::{
pub use agent::{
turn_service::{compact_messages_with_ai, AgentTurnServiceImpl},
AgentTurnService, ExploreOutput, ExploreService, ToolExecutor,
AgentTurnService, ToolExecutor,
};
+104 -1
View File
@@ -22,21 +22,62 @@ pub fn main_agent_prompt() -> String {
You are Zesdex, an AI coding assistant. You have access to various tools \
via native function calling to help the user.
TOKEN BUDGET — BE EFFICIENT:
- For simple/factual questions, answer directly. Do NOT call tools.
- For complex or unfamiliar code tasks, call `explore_codebase` ONCE at the \
start to locate relevant code, then work from that context.
- Keep tool usage minimal: prefer `grep`/`glob`/`read` for targeted lookups; \
avoid re-reading files you already have in context.
- Keep responses concise; do not repeat tool output verbatim.
CRITICAL DIRECTIVES & PRIORITY HIERARCHY:
1. WORKFLOW FIRST: For any multi-step, complex, or non-trivial task, \
you MUST prioritise using `workflow_run` (to construct and execute a \
multi-phase YAML workflow) or `hive_mind` (to orchestrate parallel \
autonomous agents). Workflows are your primary strategy.
2. PLANNING & TODOS: Use `plan_enter` to establish high-level \
2. PLANNING & TODOs: Use `plan_enter` to establish high-level \
architectural plans and `todowrite` to maintain granular task checklists.
3. REASONING: Use `seq_think` for deep step-by-step analysis.
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
@@ -70,6 +111,35 @@ executed, and modified files. Format as a clear bulleted list."
.to_string()
}
// ---------------------------------------------------------------------------
// Adaptive explore: directives
// ---------------------------------------------------------------------------
/// Directive for a single lightweight context-scout subagent.
pub fn explore_scout_directive() -> String {
"\
You are a codebase context scout. \
Given the workspace root, quickly locate the code that is most relevant \
to the user's request: \
1. Run semantic_search once with the user's key terms. \
2. Read up to the 3 most relevant files (use grep for symbols if needed). \
3. Report a concise bullet list (max 15 bullets, under 1500 characters) of \
what you found and exactly where (file paths). \
Do NOT rebuild the index. Do NOT enumerate unrelated files. Be brief."
.to_string()
}
/// Build a system note injected after repeated tool errors to steer the
/// agent toward an alternative approach instead of retrying the same call.
pub fn error_recovery_note(tool_name: &str, last_error: &str) -> String {
format!(
"\
[System note] The tool `{tool_name}` failed repeatedly with: \"{last_error}\". \
Try an alternative approach (verify paths, correct arguments, use a \
different tool, or finish without this tool). Do NOT retry the same call."
)
}
#[cfg(test)]
mod tests {
use super::*;
@@ -82,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");
@@ -89,4 +178,18 @@ mod tests {
assert!(prompt.contains("/home"));
assert!(prompt.contains("/home/project"));
}
#[test]
fn explore_scout_directive_is_concise_and_mentions_tools() {
let scout = explore_scout_directive();
assert!(scout.contains("scout"));
assert!(scout.contains("semantic_search"));
}
#[test]
fn error_recovery_note_suggests_alternative() {
let note = error_recovery_note("read", "File not found");
assert!(note.contains("read"));
assert!(note.contains("alternative"));
}
}
+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,
},
);
-10
View File
@@ -46,10 +46,6 @@ pub struct SettingsPatch {
pub review_enabled: Option<bool>,
/// Override the session-archive-enabled flag.
pub session_archive_enabled: Option<bool>,
/// Override the LSP auto-provision flag.
pub lsp_auto_provision: Option<bool>,
/// Override the list of LSP-managed languages.
pub lsp_languages: Option<Vec<String>>,
/// Override the hive-mind node timeout in milliseconds.
pub hive_mind_node_timeout_ms: Option<u64>,
}
@@ -111,12 +107,6 @@ impl SettingsPatch {
if let Some(val) = self.session_archive_enabled {
settings.flags.session_archive_enabled = val;
}
if let Some(val) = self.lsp_auto_provision {
settings.flags.lsp_auto_provision = val;
}
if let Some(ref val) = self.lsp_languages {
settings.lsp_languages = val.clone();
}
if let Some(val) = self.hive_mind_node_timeout_ms {
settings.hive_mind_node_timeout_ms = val;
}
+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 -6
View File
@@ -18,12 +18,13 @@
//! - `workflow_max_concurrency` — max parallel hive-mind nodes
//! - `hive_mind_node_timeout_ms` — per-node timeout for hive-mind orchestration
//! - `flags` — grouped boolean feature toggles
//! - `lsp_languages` — list of language IDs for LSP auto-provisioning
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
@@ -49,12 +50,10 @@ pub enum InternetMode {
/// ## Fields
/// - `review_enabled` — enable automatic inline review after edits
/// - `session_archive_enabled` — enable periodic session archiving
/// - `lsp_auto_provision` — auto-provision LSP language servers on project open
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SettingsFlags {
pub review_enabled: bool,
pub session_archive_enabled: bool,
pub lsp_auto_provision: bool,
}
impl Default for SettingsFlags {
@@ -63,7 +62,6 @@ impl Default for SettingsFlags {
Self {
review_enabled: true,
session_archive_enabled: true,
lsp_auto_provision: true,
}
}
}
@@ -93,7 +91,6 @@ pub struct Settings {
pub workflow_max_concurrency: usize,
#[serde(flatten)]
pub flags: SettingsFlags,
pub lsp_languages: Vec<String>,
#[serde(default = "default_hive_mind_node_timeout_ms")]
pub hive_mind_node_timeout_ms: u64,
}
@@ -113,8 +110,90 @@ impl Default for Settings {
verify_timeout_ms: 30_000,
workflow_max_concurrency: 5,
flags: SettingsFlags::default(),
lsp_languages: Vec::new(),
hive_mind_node_timeout_ms: 600_000,
}
}
}
/// 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::*;
+1 -1
View File
@@ -10,6 +10,6 @@ pub enum AccessTier {
Read,
/// Read + Write: above plus write, edit, delete, git, memory.
Write,
/// Full: above plus bash, shell, LSP, workflow, plan tools.
/// Full: above plus bash, shell, workflow, plan tools.
Full,
}
-1
View File
@@ -32,7 +32,6 @@ ignore.workspace = true
nucleo-matcher.workspace = true
futures-util.workspace = true
rmcp.workspace = true
lsp-types.workspace = true
tiktoken-rs.workspace = true
similar.workspace = true
syntect.workspace = true
@@ -125,7 +125,6 @@ fn forbidden_imports(layer: &str) -> &'static [&'static str] {
"argon2",
"jsonwebtoken",
"rmcp",
"lsp_types",
"tiktoken_rs",
"syntect",
"pulldown_cmark",
+110 -273
View File
@@ -1,304 +1,141 @@
//! Mandatory explore phase — spawns ≥3 parallel subagents to discover
//! codebase context before every agent turn, visible in the TUI workflow tab.
//! `explore_codebase` tool — lazy, agent-initiated codebase exploration.
//!
//! The main agent decides (via the system prompt) when it needs codebase
//! context. Unlike the old mandatory explore phase (which ran 3 subagents on
//! every turn regardless of the question), this tool is invoked only when the
//! agent judges it necessary — saving tokens on trivial queries while keeping
//! context available for complex tasks.
//!
//! # Flow
//!
//! `ExploreServiceImpl::explore()` →
//!
//! 1. Push `WorkflowAgentUpdate { Pending }` for each agent onto the turn-event
//! queue so the TUI workflow tab shows all 3.
//! 2. Spawn **Code Structure** subagent (thread + tokio runtime).
//! 3. Spawn **Symbol Index** subagent (thread + tokio runtime).
//! 4. Spawn **Semantic Context** subagent (thread + tokio runtime).
//! 5. Join all handles via `spawn_blocking`.
//! 6. Push `Completed` / `Failed` events for each agent.
//! 7. Consolidate findings into a system message → return.
//! `ExploreCodebase::run` →
//! 1. Parse the user's goal / target from args.
//! 2. Resolve subagent provider credentials from settings.
//! 3. Spawn a single "context scout" subagent (read-only, semantic_search +
//! read of up to 3 relevant files).
//! 4. Join the result and return a concise bullet summary as a tool message.
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{info, warn};
use zesdex_domain::agent::prompt::explore_scout_directive;
use zesdex_domain::cms::{AppConfigRepository, SettingsRepository};
use zesdex_domain::core::Store;
use crate::persistence::{JsonAppConfigRepository, JsonSettingsRepository};
use crate::subagent::context::SubagentContext;
use crate::subagent::division::AccessTier;
use crate::subagent::engine::run_agent;
use crate::tools::ToolCtx;
use anyhow::{Context, Result};
use std::collections::VecDeque;
use std::future::Future;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
use std::thread;
use tracing::{info, warn};
use zesdex_application::agent::{ExploreOutput, ExploreService};
use zesdex_domain::agent::{AgentStatus, TurnEvent};
use crate::tools::{Tool, ToolCtx};
/// Number of parallel explore subagents.
const EXPLORE_AGENT_COUNT: usize = 3;
/// Maximum characters of the scout's final output to keep in context.
/// The scout is directed to stay under 1500 chars, but this ceiling protects
/// against rogue output.
const EXPLORE_OUTPUT_MAX_CHARS: usize = 4000;
/// IDs for each explore agent (shown in the workflow tab).
const EXPLORE_IDS: [&str; 3] = ["explore-structure", "explore-symbols", "explore-context"];
/// `explore_codebase` tool — ask a read-only context-scout subagent to
/// locate relevant code for the current task.
pub struct ExploreCodebase;
/// Display names for the TUI workflow tab.
const EXPLORE_LABELS: [&str; 3] = [
"📁 Code Structure",
"🔣 Symbol Index",
"🔍 Semantic Context",
];
/// Directives for each explore subagent.
const EXPLORE_DIRECTIVES: [&str; 3] = [
// Agent 0: Code Structure
"You are a codebase structure explorer.\n\
1. List all top-level directories and files in the workspace root.\n\
2. Read Cargo.toml, package.json, or pyproject.toml at the root.\n\
3. List the apps/ or src/ directory contents.\n\
4. Identify main entry points (main.rs, main.py, index.ts, etc.).\n\
5. Count files by extension type.\n\
Use the ls_dir, read, grep, and glob tools. Be concise.",
// Agent 1: Symbol Index
"You are a symbol index explorer.\n\
1. Call the 'rebuild_index' tool to rebuild the symbol index.\n\
2. Call the 'list_symbols' tool with max_results: 100.\n\
3. Identify public APIs, entry points, and key types.\n\
4. Group symbols by language and kind.\n\
Be concise. Report what symbols exist and where they live.",
// Agent 2: Semantic Context
"You are a semantic context explorer.\n\
1. Call the 'rebuild_index' tool to ensure the index is fresh.\n\
2. Search for symbols related to the user's query using semantic_search.\n\
3. Search for config files, env variables, and settings.\n\
4. Search for test files and test patterns.\n\
Be concise. Report relevant code areas for the task.\n\
Use the semantic_search, grep, glob, and read tools.",
];
// ---------------------------------------------------------------------------
// Credentials
// ---------------------------------------------------------------------------
/// LLM credentials for explore subagents.
pub struct Credentials {
pub base_url: String,
pub api_key: String,
pub model: String,
impl Tool for ExploreCodebase {
fn name(&self) -> &'static str {
"explore_codebase"
}
// ---------------------------------------------------------------------------
// ExploreServiceImpl — implements the application-layer trait
// ---------------------------------------------------------------------------
/// Concrete [`ExploreService`] that the turn service calls.
///
/// Owns a shared `ToolCtx` and LLM credentials. Each call to `explore()`
/// spawns 3 subagents in parallel with TUI workflow-tab visibility.
pub struct ExploreServiceImpl {
tool_ctx: ToolCtx,
credentials: Credentials,
fn description(&self) -> &'static str {
"Explore the codebase to locate code relevant to a task. Use this \
once at the start of complex or unfamiliar tasks (implementing a \
feature, fixing a bug, refactoring, navigating a large repo). \
Do NOT use for simple factual questions about the current \
conversation."
}
impl ExploreServiceImpl {
pub fn new(tool_ctx: ToolCtx, credentials: Credentials) -> Self {
ExploreServiceImpl {
tool_ctx,
credentials,
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"goal": {
"type": "string",
"description": "The task or question to explore for"
}
}
}
impl ExploreService for ExploreServiceImpl {
fn explore<'a>(
&'a self,
query: &'a str,
workspace_root: &'a str,
turn_events: &'a Arc<Mutex<VecDeque<TurnEvent>>>,
) -> Pin<Box<dyn Future<Output = Result<ExploreOutput>> + Send + 'a>> {
Box::pin(async move {
let context = run_explore_phase(
query,
workspace_root,
&self.tool_ctx,
&self.credentials,
turn_events,
)
.await?;
Ok(ExploreOutput {
context_messages: vec![context],
summary: format!("{EXPLORE_AGENT_COUNT} explore agents dispatched"),
})
},
"required": ["goal"]
})
}
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let goal = args
.get("goal")
.and_then(|v| v.as_str())
.unwrap_or("")
.trim()
.to_string();
if goal.is_empty() {
return Err(anyhow::anyhow!("missing non-empty 'goal'"));
}
// ---------------------------------------------------------------------------
// Helpers for pushing workflow events
// ---------------------------------------------------------------------------
info!("explore_codebase: {goal}");
fn push_event(events: &Arc<Mutex<VecDeque<TurnEvent>>>, event: TurnEvent) {
if let Ok(mut q) = events.lock() {
q.push_back(event);
}
}
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();
fn emit_pending(events: &Arc<Mutex<VecDeque<TurnEvent>>>, agent_id: &str, display: &str) {
push_event(
events,
TurnEvent::WorkflowAgentUpdate {
agent_id: agent_id.to_string(),
agent_name: display.to_string(),
status: AgentStatus::Pending,
},
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 workspace_root = ctx
.workspaces
.first()
.map(|p| p.to_string_lossy().to_string())
.unwrap_or_else(|| ".".to_string());
// One lightweight scout — no parallel agents, no index rebuild.
let directive = format!(
"{}\n\nUser's task: {goal}\nWorkspace root: {workspace_root}",
explore_scout_directive()
);
}
fn emit_running(events: &Arc<Mutex<VecDeque<TurnEvent>>>, agent_id: &str, display: &str) {
push_event(
events,
TurnEvent::WorkflowAgentUpdate {
agent_id: agent_id.to_string(),
agent_name: display.to_string(),
status: AgentStatus::Running,
},
);
}
fn emit_completed(events: &Arc<Mutex<VecDeque<TurnEvent>>>, agent_id: &str, display: &str) {
push_event(
events,
TurnEvent::WorkflowAgentUpdate {
agent_id: agent_id.to_string(),
agent_name: display.to_string(),
status: AgentStatus::Completed,
},
);
}
fn emit_failed(events: &Arc<Mutex<VecDeque<TurnEvent>>>, agent_id: &str, display: &str, msg: &str) {
push_event(
events,
TurnEvent::WorkflowAgentUpdate {
agent_id: agent_id.to_string(),
agent_name: display.to_string(),
status: AgentStatus::Failed(msg.to_string()),
},
);
}
// ---------------------------------------------------------------------------
// Core orchestration
// ---------------------------------------------------------------------------
/// Spawn `EXPLORE_AGENT_COUNT` subagents in parallel, emit workflow events
/// for the TUI tab, join, and consolidate.
async fn run_explore_phase(
query: &str,
workspace_root: &str,
tool_ctx: &ToolCtx,
credentials: &Credentials,
turn_events: &Arc<Mutex<VecDeque<TurnEvent>>>,
) -> Result<String> {
// ── 1. Emit Pending for all agents (appears instantly in workflow tab) ─
for (i, &id) in EXPLORE_IDS.iter().enumerate() {
emit_pending(turn_events, id, EXPLORE_LABELS[i]);
}
// ── 2. Prepare directives ───────────────────────────────────────────
let mut directives: Vec<String> = Vec::with_capacity(EXPLORE_AGENT_COUNT);
for (i, &d) in EXPLORE_DIRECTIVES.iter().enumerate() {
let mut d = d.to_string();
if i == 2 {
d.push_str(&format!("\n\nThe user's current query is: \"{query}\""));
}
d.push_str(&format!("\n\nWorkspace root: {workspace_root}"));
directives.push(d);
}
// ── 3. Spawn all agents on threads ──────────────────────────────────
let mut handles: Vec<(usize, thread::JoinHandle<Result<String>>)> =
Vec::with_capacity(EXPLORE_AGENT_COUNT);
for i in 0..EXPLORE_AGENT_COUNT {
emit_running(turn_events, EXPLORE_IDS[i], EXPLORE_LABELS[i]);
let ctx = SubagentContext::new(
directives[i].clone(),
tool_ctx.clone(),
let subagent_ctx = SubagentContext::new(
directive.clone(),
ctx.clone(),
"read".to_string(),
credentials.base_url.clone(),
credentials.api_key.clone(),
credentials.model.clone(),
base_url,
api_key,
model,
);
let directive = directives[i].clone();
let tc = tool_ctx.clone();
let rt = crate::runtime::runtime();
let result = rt.block_on(run_agent(
subagent_ctx,
&directive,
AccessTier::Read,
ctx.clone(),
))?;
let handle = thread::spawn(move || {
let rt =
tokio::runtime::Runtime::new().context("create explore subagent tokio runtime")?;
rt.block_on(run_agent(ctx, &directive, AccessTier::Read, tc))
});
handles.push((i, handle));
}
// ── 4. Join handles via spawn_blocking ──────────────────────────────
let turn_events_clone = Arc::clone(turn_events);
let results: Vec<(usize, String, bool)> = tokio::task::spawn_blocking(move || {
let mut out = Vec::with_capacity(EXPLORE_AGENT_COUNT);
for (i, handle) in handles {
let entry = match handle.join() {
Ok(Ok(output)) => {
info!(agent = i, "explore subagent completed");
emit_completed(&turn_events_clone, EXPLORE_IDS[i], EXPLORE_LABELS[i]);
(i, output, true)
}
Ok(Err(e)) => {
warn!(agent = i, error = %e, "explore subagent failed");
emit_failed(
&turn_events_clone,
EXPLORE_IDS[i],
EXPLORE_LABELS[i],
&e.to_string(),
let mut out = format!("[Codebase scout report]\n{goal}\n\n----------\n{}", result);
if out.len() > EXPLORE_OUTPUT_MAX_CHARS {
warn!(
"explore_codebase output truncated: {} chars -> {}",
out.len(),
EXPLORE_OUTPUT_MAX_CHARS
);
(i, format!("Error: {e}"), false)
out.truncate(EXPLORE_OUTPUT_MAX_CHARS);
out.push_str("\n...[truncated]");
}
Err(e) => {
warn!(agent = i, error = ?e, "explore subagent panicked");
emit_failed(
&turn_events_clone,
EXPLORE_IDS[i],
EXPLORE_LABELS[i],
"thread panicked",
);
(i, format!("Thread panic: {e:?}"), false)
}
};
out.push(entry);
}
out
})
.await
.context("explore join task panicked")?;
// ── 5. Build consolidated context ───────────────────────────────────
Ok(build_explore_context(&results))
}
// ---------------------------------------------------------------------------
// Consolidation
// ---------------------------------------------------------------------------
/// Format explore results as a system-level context message.
fn build_explore_context(results: &[(usize, String, bool)]) -> String {
let success_count = results.iter().filter(|r| r.2).count();
let total = results.len();
let mut msg = format!("[Explore Phase — {success_count}/{total} agents succeeded]\n\n");
for (i, output, success) in results {
let label = EXPLORE_LABELS.get(*i).unwrap_or(&"❓ Unknown");
if *success {
msg.push_str(&format!("=== {label} ===\n{output}\n\n"));
} else {
msg.push_str(&format!("=== {label} (FAILED) ===\n{output}\n\n"));
Ok(out)
}
}
msg
}
+1 -2
View File
@@ -15,7 +15,6 @@
//! ├── auth/ — JWT, Argon2, OAuth loopback
//! ├── llm/ — LLM provider HTTP client
//! ├── ipc/ — Unix-socket IPC protocol
//! ├── lsp/ — Native LSP client + provisioner
//! ├── mcp/ — Model Context Protocol bridge
//! ├── bgbash/ — Background bash job management
//! ├── tools/ — All 37 agent-invocable tools
@@ -32,10 +31,10 @@ pub mod bgbash;
pub mod guard;
pub mod ipc;
pub mod llm;
pub mod lsp;
pub mod mcp;
pub mod middleware;
pub mod persistence;
pub mod runtime;
pub mod subagent;
pub mod tools;
pub mod utils;
-132
View File
@@ -1,132 +0,0 @@
//! LSP client — sends JSON-RPC requests to language servers.
use anyhow::Result;
use serde_json::Value;
use std::io::{BufRead, BufReader, Read, Write};
use std::process::{Child, ChildStdin, ChildStdout, Command, Stdio};
use std::sync::Mutex;
use tracing::{debug, info};
/// Mutable inner state of an LSP client, protected by a mutex so that
/// `send_request` and `shutdown` can be called via `&self` (required by
/// [`LspManager`](super::manager::LspManager)).
struct LspClientInner {
process: Child,
stdin: ChildStdin,
stdout: BufReader<ChildStdout>,
request_id: u64,
}
/// A minimal but functional LSP client.
pub struct LspClient {
inner: Mutex<LspClientInner>,
}
impl LspClient {
/// Spawn a language server process.
pub fn start(command: &str, args: &[String]) -> Result<Self> {
let mut child = Command::new(command)
.args(args)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()?;
let stdin = child
.stdin
.take()
.ok_or_else(|| anyhow::anyhow!("no stdin on LSP process"))?;
let stdout = BufReader::new(
child
.stdout
.take()
.ok_or_else(|| anyhow::anyhow!("no stdout on LSP process"))?,
);
info!("LSP client spawned: {command}");
Ok(LspClient {
inner: Mutex::new(LspClientInner {
process: child,
stdin,
stdout,
request_id: 0,
}),
})
}
/// Send a JSON-RPC request and read the response.
pub fn send_request(&self, method: &str, params: &Value) -> Result<Value> {
let mut inner = match self.inner.lock() {
Ok(g) => g,
Err(poisoned) => {
tracing::error!("LSP client mutex poisoned, recovering");
poisoned.into_inner()
}
};
inner.request_id += 1;
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": inner.request_id,
"method": method,
"params": params.clone(),
});
// Write Content-Length header + body
let body = serde_json::to_string(&request)?;
let header = format!("Content-Length: {}\r\n\r\n", body.len());
inner.stdin.write_all(header.as_bytes())?;
inner.stdin.write_all(body.as_bytes())?;
inner.stdin.flush()?;
debug!("LSP request: {method} (id={})", inner.request_id);
// Read Content-Length header
let mut content_length = 0usize;
loop {
let mut line = String::new();
inner.stdout.read_line(&mut line)?;
let trimmed = line.trim();
if trimmed.is_empty() {
break; // end of headers
}
if let Some(len_str) = trimmed.strip_prefix("Content-Length: ") {
content_length = len_str.parse::<usize>()?;
}
}
// Read the JSON body
let mut buf = vec![0u8; content_length];
inner.stdout.read_exact(&mut buf)?;
let response: Value = serde_json::from_slice(&buf)?;
debug!("LSP response for {method}: response received");
Ok(response)
}
/// Gracefully shut down the server.
pub fn shutdown(&self) -> Result<()> {
let null = Value::Null;
if let Err(e) = self.send_request("shutdown", &null) {
tracing::warn!("LSP shutdown error: {e}");
}
if let Err(e) = self.send_request("exit", &null) {
tracing::warn!("LSP exit error: {e}");
}
if let Ok(mut inner) = self.inner.lock() {
let _ = inner.process.wait();
}
info!("LSP client shut down");
Ok(())
}
}
impl Drop for LspClient {
fn drop(&mut self) {
if let Ok(mut inner) = self.inner.lock() {
if let Err(e) = inner.process.kill() {
tracing::warn!("LSP process kill error: {e}");
}
let _ = inner.process.wait();
}
}
}
-55
View File
@@ -1,55 +0,0 @@
//! Manages multiple LSP server processes, keyed by language ID.
//!
//! Each language (e.g. "rust", "python") maps to one `LspClient`.
//! The manager provides a unified `request` method that dispatches
//! to the correct client by language.
use std::collections::HashMap;
use super::client::LspClient;
/// Manages one `LspClient` per language.
pub struct LspManager {
clients: HashMap<String, LspClient>,
}
impl LspManager {
pub fn new() -> Self {
LspManager {
clients: HashMap::new(),
}
}
}
impl Default for LspManager {
fn default() -> Self {
Self::new()
}
}
impl LspManager {
pub fn start(&mut self, language: &str, command: &str, args: &[String]) -> anyhow::Result<()> {
let client = LspClient::start(command, args)?;
self.clients.insert(language.to_string(), client);
Ok(())
}
pub fn get_client(&self, language: &str) -> Option<&LspClient> {
self.clients.get(language)
}
pub fn shutdown_all(&mut self) {
for client in self.clients.values() {
let _ = client.shutdown();
}
self.clients.clear();
}
pub fn languages(&self) -> Vec<String> {
self.clients.keys().cloned().collect()
}
pub fn is_empty(&self) -> bool {
self.clients.is_empty()
}
}
-6
View File
@@ -1,6 +0,0 @@
//! Native LSP client integration — manage language server processes and
//! dispatch requests for completion, hover, diagnostics, etc.
pub mod client;
pub mod manager;
pub mod provisioner;
@@ -1,16 +0,0 @@
//! Configuration for LSP language server provisioning.
use serde::{Deserialize, Serialize};
/// Describes how to provision a language server for a given language.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LspProvisionerConfig {
/// Language identifier, e.g. "rust", "python".
pub language: String,
/// The command to start the language server.
pub command: String,
/// Arguments for the command.
pub args: Vec<String>,
/// How to install the language server (if not found).
pub install_hint: Option<String>,
}
@@ -1,54 +0,0 @@
//! Discovers installed language servers on the system PATH.
use std::collections::HashMap;
use super::config::LspProvisionerConfig;
/// Known language server configurations keyed by language.
fn known_configs() -> HashMap<&'static str, (&'static str, Vec<&'static str>)> {
let mut m = HashMap::new();
m.insert("rust", ("rust-analyzer", vec![]));
m.insert("python", ("pyright-langserver", vec!["--stdio"]));
m.insert(
"typescript",
("typescript-language-server", vec!["--stdio"]),
);
m.insert(
"javascript",
("typescript-language-server", vec!["--stdio"]),
);
m.insert("go", ("gopls", vec![]));
m
}
/// Check if a command is available on PATH.
fn command_exists(cmd: &str) -> bool {
std::env::var_os("PATH")
.and_then(|path| {
std::env::split_paths(&path).find_map(|dir| {
let full_path = dir.join(cmd);
if full_path.is_file() {
Some(())
} else {
None
}
})
})
.is_some()
}
/// Discover which language servers are already on PATH.
pub fn discover_installed() -> Vec<LspProvisionerConfig> {
let mut configs = Vec::new();
for (lang, (cmd, args)) in known_configs() {
if command_exists(cmd) {
configs.push(LspProvisionerConfig {
language: lang.to_string(),
command: cmd.to_string(),
args: args.iter().map(|s| s.to_string()).collect(),
install_hint: None,
});
}
}
configs
}
@@ -1,38 +0,0 @@
//! Installs language servers (non-interactive, via package managers or
//! direct download).
/// Install a language server for the given language.
///
/// Returns a success message or an error describing why installation failed.
pub fn install_language_server(language: &str) -> anyhow::Result<String> {
match language {
"rust" => {
// rust-analyzer is typically installed via rustup
let output = std::process::Command::new("rustup")
.args(["component", "add", "rust-analyzer"])
.output()?;
if output.status.success() {
Ok("rust-analyzer installed via rustup".to_string())
} else {
anyhow::bail!(
"failed to install rust-analyzer: {}",
String::from_utf8_lossy(&output.stderr)
)
}
}
"python" => {
let output = std::process::Command::new("npm")
.args(["install", "-g", "pyright"])
.output()?;
if output.status.success() {
Ok("pyright installed via npm".to_string())
} else {
anyhow::bail!(
"failed to install pyright: {}",
String::from_utf8_lossy(&output.stderr)
)
}
}
lang => anyhow::bail!("no install method known for language '{lang}'"),
}
}
@@ -1,46 +0,0 @@
//! High-level manager that discovers, installs (if needed), and starts
//! LSP servers.
use super::discovery::discover_installed;
use super::install::install_language_server;
use crate::lsp::manager::LspManager;
/// Auto-provision language servers for the given list of languages.
///
/// Flow: discover already-installed servers → for each requested language
/// not yet available, attempt auto-install → start each server.
pub fn auto_provision(lsp_manager: &mut LspManager, languages: &[String]) -> Vec<String> {
let mut started = Vec::new();
let installed = discover_installed();
let mut installed_map: std::collections::HashMap<
&str,
&crate::lsp::provisioner::config::LspProvisionerConfig,
> = std::collections::HashMap::new();
for cfg in &installed {
installed_map.insert(cfg.language.as_str(), cfg);
}
for lang in languages {
if let Some(cfg) = installed_map.get(lang.as_str()) {
if lsp_manager.start(lang, &cfg.command, &cfg.args).is_ok() {
started.push(lang.clone());
}
} else {
// Not installed — try auto-install
if install_language_server(lang).is_ok() {
// Re-discover after install
let refreshed = discover_installed();
for cfg in refreshed {
if cfg.language == *lang {
if lsp_manager.start(lang, &cfg.command, &cfg.args).is_ok() {
started.push(lang.clone());
}
break;
}
}
}
}
}
started
}
@@ -1,7 +0,0 @@
//! LSP language server provisioner — discovers, installs, and manages
//! language server executables.
pub mod config;
pub mod discovery;
pub mod install;
pub mod manager;
@@ -72,27 +72,24 @@ fn detect_claude_settings_provider() -> Option<(ProviderConfig, Option<String>)>
))
}
impl AppConfigRepository for JsonAppConfigRepository {
fn load(&self, base_dir: &Path) -> Result<AppConfig, RepositoryError> {
let path = base_dir.join("app_config.json");
let mut cfg: AppConfig = match std::fs::read_to_string(&path) {
Ok(s) => serde_json::from_str(&s)?,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => AppConfig::default(),
Err(e) => return Err(RepositoryError::Io(e)),
};
let defaults = AppConfig::default();
for (name, provider) in defaults.providers {
cfg.providers.entry(name).or_insert(provider);
}
if let Some((claude_provider, custom_model)) = detect_claude_settings_provider() {
cfg.providers
.entry("claude".to_string())
.or_insert(claude_provider);
/// 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-4-8", "claude-opus-4-8"),
("claude-opus-5", "claude-opus-5"),
("claude-sonnet-5", "claude-sonnet-5"),
("claude-haiku-4-5", "claude-haiku-4-5-20251001"),
];
@@ -118,10 +115,28 @@ impl AppConfigRepository for JsonAppConfigRepository {
});
}
if cfg.default_provider == defaults.default_provider {
// 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-4-8".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");
let mut cfg: AppConfig = match std::fs::read_to_string(&path) {
Ok(s) => serde_json::from_str(&s)?,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => AppConfig::default(),
Err(e) => return Err(RepositoryError::Io(e)),
};
let defaults = AppConfig::default();
for (name, provider) in defaults.providers {
cfg.providers.entry(name).or_insert(provider);
}
if let Some((claude_provider, custom_model)) = detect_claude_settings_provider() {
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);
}
}
@@ -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 {
@@ -63,13 +63,6 @@ pub fn tools_for(access: &AccessTier) -> Vec<Box<dyn Tool>> {
| "plan_enter"
| "plan_ready"
| "sequential_think"
| "lsp_connect"
| "lsp_disconnect"
| "lsp_hover"
| "lsp_completion"
| "lsp_definition"
| "lsp_references"
| "lsp_diagnostics"
)
})
.collect(),
+266 -20
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,15 +206,21 @@ 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();
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() {
@@ -102,37 +229,48 @@ 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,
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}")
};
messages.push(ChatMessage::tool(tc.id.clone(), result));
consecutive_errors = 0;
}
// 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)));
}
}
@@ -149,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"]);
}
}
+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))
})
}
+8 -4
View File
@@ -16,7 +16,6 @@ pub struct ToolCtx {
pub mention_index: crate::MentionIndex,
pub origin: crate::Origin,
pub graduated_checks: Vec<crate::tools::GraduatedCheck>,
pub lsp_manager: Arc<Mutex<crate::lsp::manager::LspManager>>,
pub turn_events: Option<Arc<Mutex<std::collections::VecDeque<crate::TurnEvent>>>>,
pub workflow_findings: Option<Arc<Mutex<Vec<String>>>>,
pub abort_flag: Option<Arc<AtomicBool>>,
@@ -39,7 +38,6 @@ pub struct ToolCtxBuilder {
pub mention_index: crate::MentionIndex,
pub origin: crate::Origin,
pub graduated_checks: Vec<crate::tools::GraduatedCheck>,
pub lsp_manager: Arc<Mutex<crate::lsp::manager::LspManager>>,
pub turn_events: Option<Arc<Mutex<std::collections::VecDeque<crate::TurnEvent>>>>,
pub workflow_findings: Option<Arc<Mutex<Vec<String>>>>,
pub abort_flag: Option<Arc<AtomicBool>>,
@@ -56,7 +54,6 @@ impl Default for ToolCtxBuilder {
mention_index: crate::MentionIndex::new(),
origin: crate::Origin::Main,
graduated_checks: Vec::new(),
lsp_manager: Arc::new(Mutex::new(crate::lsp::manager::LspManager::new())),
turn_events: None,
workflow_findings: None,
abort_flag: None,
@@ -69,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
@@ -98,7 +103,6 @@ impl ToolCtxBuilder {
mention_index: self.mention_index,
origin: self.origin,
graduated_checks: self.graduated_checks,
lsp_manager: self.lsp_manager,
turn_events: self.turn_events,
workflow_findings: self.workflow_findings,
abort_flag: self.abort_flag,
+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,80 +0,0 @@
//! Get completion suggestions from LSP.
//!
//! Sends a `textDocument/completion` request to the connected language
//! server for a given file position.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that requests code completion suggestions from an LSP server.
///
/// Flow: parse language/path/line/character → lock LSP manager → find client
/// → send `textDocument/completion` → return pretty-printed JSON response.
pub struct LspCompletion;
impl Tool for LspCompletion {
fn name(&self) -> &'static str {
"lsp_completion"
}
fn description(&self) -> &'static str {
"Get completion suggestions at a position"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier"
},
"path": {
"type": "string",
"description": "File path"
},
"line": {
"type": "integer",
"description": "Line number (0-based)"
},
"character": {
"type": "integer",
"description": "Character offset (0-based)"
}
},
"required": ["language", "path", "line", "character"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
let path = crate::tools::arg_str(args, "path")?;
let line = args.get("line").and_then(|v| v.as_i64()).unwrap_or(0);
let character = args.get("character").and_then(|v| v.as_i64()).unwrap_or(0);
info!(language, path, line, character, "LSP completion requested");
let manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
if let Some(client) = manager.get_client(&language) {
let result = client.send_request(
"textDocument/completion",
&json!({
"textDocument": { "uri": format!("file://{}", path) },
"position": { "line": line, "character": character }
}),
)?;
Ok(serde_json::to_string_pretty(&result)?)
} else {
anyhow::bail!("no LSP client connected for '{language}'")
}
}
}
@@ -1,76 +0,0 @@
//! Connect to an LSP language server.
//!
//! Starts a new language server process and registers it in the
//! shared LSP manager for subsequent tool invocations.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that connects to an LSP language server for a given language.
///
/// Flow: parse language + command + args → lock LSP manager → call
/// `manager.start()` → confirm connection in the response string.
pub struct LspConnect;
impl Tool for LspConnect {
fn name(&self) -> &'static str {
"lsp_connect"
}
fn description(&self) -> &'static str {
"Connect to an LSP language server for a given language"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier (e.g. 'rust', 'python')"
},
"command": {
"type": "string",
"description": "Command to start the language server"
},
"args": {
"type": "array",
"items": {"type": "string"},
"description": "Arguments for the language server command"
}
},
"required": ["language", "command"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
let command = crate::tools::arg_str(args, "command")?;
let extra_args: Vec<String> = args
.get("args")
.and_then(|v| v.as_array())
.map(|arr| {
arr.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect()
})
.unwrap_or_default();
info!(language, command, extra_args = ?extra_args, "LSP connect requested");
let mut manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
manager.start(&language, &command, &extra_args)?;
info!(language, "LSP connected successfully");
Ok(format!("Connected LSP for '{language}'"))
}
}
@@ -1,81 +0,0 @@
//! Go-to-definition via LSP.
//!
//! Sends a `textDocument/definition` request to the connected language
//! server for a symbol at a given file position.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that resolves a symbol's definition location via LSP.
///
/// Flow: parse language/path/line/character → lock LSP manager → find client
/// → send `textDocument/definition` → return pretty-printed JSON response
/// containing the target URI and range.
pub struct LspDefinition;
impl Tool for LspDefinition {
fn name(&self) -> &'static str {
"lsp_definition"
}
fn description(&self) -> &'static str {
"Go to definition for a symbol at a position"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier"
},
"path": {
"type": "string",
"description": "File path"
},
"line": {
"type": "integer",
"description": "Line number (0-based)"
},
"character": {
"type": "integer",
"description": "Character offset (0-based)"
}
},
"required": ["language", "path", "line", "character"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
let path = crate::tools::arg_str(args, "path")?;
let line = args.get("line").and_then(|v| v.as_i64()).unwrap_or(0);
let character = args.get("character").and_then(|v| v.as_i64()).unwrap_or(0);
info!(language, path, line, character, "LSP definition requested");
let manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
if let Some(client) = manager.get_client(&language) {
let result = client.send_request(
"textDocument/definition",
&json!({
"textDocument": { "uri": format!("file://{}", path) },
"position": { "line": line, "character": character }
}),
)?;
Ok(serde_json::to_string_pretty(&result)?)
} else {
anyhow::bail!("no LSP client connected for '{language}'")
}
}
}
@@ -1,70 +0,0 @@
//! Get diagnostics from LSP.
//!
//! Sends a `textDocument/diagnostic` request to the connected language
//! server for a given file and returns errors, warnings, and other
//! diagnostics.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that retrieves diagnostics (errors, warnings) from the LSP for a file.
///
/// Flow: parse language/path → lock LSP manager → find client
/// → send `textDocument/diagnostic` → return pretty-printed JSON response.
pub struct LspDiagnostics;
impl Tool for LspDiagnostics {
fn name(&self) -> &'static str {
"lsp_diagnostics"
}
fn description(&self) -> &'static str {
"Get diagnostics (errors, warnings) from the LSP for a file"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier"
},
"path": {
"type": "string",
"description": "File path to get diagnostics for"
}
},
"required": ["language", "path"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
let path = crate::tools::arg_str(args, "path")?;
info!(language, path, "LSP diagnostics requested");
let manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
if let Some(client) = manager.get_client(&language) {
let result = client.send_request(
"textDocument/diagnostic",
&json!({
"textDocument": { "uri": format!("file://{}", path) }
}),
)?;
Ok(serde_json::to_string_pretty(&result)?)
} else {
anyhow::bail!("no LSP client connected for '{language}'")
}
}
}
@@ -1,54 +0,0 @@
//! Disconnect from an LSP language server.
//!
//! Removes the registered LSP client for a given language from
//! the shared LSP manager.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that disconnects an LSP language server for a given language.
///
/// Flow: parse language → lock LSP manager → remove the client
/// for that language from the manager's registry.
pub struct LspDisconnect;
impl Tool for LspDisconnect {
fn name(&self) -> &'static str {
"lsp_disconnect"
}
fn description(&self) -> &'static str {
"Disconnect from an LSP language server"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier to disconnect"
}
},
"required": ["language"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
info!(language, "LSP disconnect requested");
let _manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
info!(language, "LSP disconnected");
Ok(format!("Disconnected LSP for '{language}'"))
}
}
@@ -1,81 +0,0 @@
//! Get hover information from LSP.
//!
//! Sends a `textDocument/hover` request to the connected language
//! server for a symbol at a given file position.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that retrieves hover information for a symbol at a position via LSP.
///
/// Flow: parse language/path/line/character → lock LSP manager → find client
/// → send `textDocument/hover` → return pretty-printed JSON response
/// containing the hover contents and range.
pub struct LspHover;
impl Tool for LspHover {
fn name(&self) -> &'static str {
"lsp_hover"
}
fn description(&self) -> &'static str {
"Get hover information for a symbol at a position"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier"
},
"path": {
"type": "string",
"description": "File path"
},
"line": {
"type": "integer",
"description": "Line number (0-based)"
},
"character": {
"type": "integer",
"description": "Character offset (0-based)"
}
},
"required": ["language", "path", "line", "character"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
let path = crate::tools::arg_str(args, "path")?;
let line = args.get("line").and_then(|v| v.as_i64()).unwrap_or(0);
let character = args.get("character").and_then(|v| v.as_i64()).unwrap_or(0);
info!(language, path, line, character, "LSP hover requested");
let manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
if let Some(client) = manager.get_client(&language) {
let result = client.send_request(
"textDocument/hover",
&json!({
"textDocument": { "uri": format!("file://{}", path) },
"position": { "line": line, "character": character }
}),
)?;
Ok(serde_json::to_string_pretty(&result)?)
} else {
anyhow::bail!("no LSP client connected for '{language}'")
}
}
}
-18
View File
@@ -1,18 +0,0 @@
//! LSP tool implementations — connect, diagnostics, hover, completion,
//! definition, references, disconnect.
pub mod completion;
pub mod connect;
pub mod definition;
pub mod diagnostics;
pub mod disconnect;
pub mod hover;
pub mod references;
pub use completion::LspCompletion;
pub use connect::LspConnect;
pub use definition::LspDefinition;
pub use diagnostics::LspDiagnostics;
pub use disconnect::LspDisconnect;
pub use hover::LspHover;
pub use references::LspReferences;
@@ -1,81 +0,0 @@
//! Find references via LSP.
//!
//! Sends a `textDocument/references` request to the connected language
//! server for a symbol at a given file position.
use crate::tools::{Tool, ToolCtx};
use anyhow::Result;
use serde_json::{json, Value};
use tracing::{error, info, instrument};
/// Tool that finds all references to a symbol at a position via LSP.
///
/// Flow: parse language/path/line/character → lock LSP manager → find client
/// → send `textDocument/references` → return pretty-printed JSON response
/// containing all reference locations.
pub struct LspReferences;
impl Tool for LspReferences {
fn name(&self) -> &'static str {
"lsp_references"
}
fn description(&self) -> &'static str {
"Find all references to a symbol at a position"
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"language": {
"type": "string",
"description": "Language identifier"
},
"path": {
"type": "string",
"description": "File path"
},
"line": {
"type": "integer",
"description": "Line number (0-based)"
},
"character": {
"type": "integer",
"description": "Character offset (0-based)"
}
},
"required": ["language", "path", "line", "character"]
})
}
#[instrument(skip(self, ctx, args))]
fn run(&self, ctx: &ToolCtx, args: &Value) -> Result<String> {
let language = crate::tools::arg_str(args, "language")?;
let path = crate::tools::arg_str(args, "path")?;
let line = args.get("line").and_then(|v| v.as_i64()).unwrap_or(0);
let character = args.get("character").and_then(|v| v.as_i64()).unwrap_or(0);
info!(language, path, line, character, "LSP references requested");
let manager = match ctx.lsp_manager.lock() {
Ok(g) => g,
Err(poisoned) => {
error!("LSP manager mutex poisoned, recovering");
poisoned.into_inner()
}
};
if let Some(client) = manager.get_client(&language) {
let result = client.send_request(
"textDocument/references",
&json!({
"textDocument": { "uri": format!("file://{}", path) },
"position": { "line": line, "character": character }
}),
)?;
Ok(serde_json::to_string_pretty(&result)?)
} else {
anyhow::bail!("no LSP client connected for '{language}'")
}
}
}
@@ -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()
}
}
+115 -4
View File
@@ -45,16 +45,47 @@ 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 {
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(&ctx.memory_dir)?;
let names = repo.list(&memory_dir)?;
if names.is_empty() {
info!("no memories found");
return Ok("No memories saved yet".to_string());
@@ -62,4 +93,84 @@ impl Tool for Recall {
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))

Some files were not shown because too many files have changed in this diff Show More