Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9d1544d799 | ||
|
|
a27e8151b3 | ||
|
|
88c2a7d000 | ||
|
|
3b25e3898f | ||
|
|
d23d3855ec | ||
|
|
104af3abb9 | ||
|
|
62fa85867c | ||
|
|
f75ff740ac |
@@ -1,3 +1,31 @@
|
|||||||
|
# [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)
|
## [1.20.2](https://github.com/asepharyana/zesdex/compare/v1.20.1...v1.20.2) (2026-08-28)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Generated
+11
-11
@@ -4862,7 +4862,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-api"
|
name = "zesdex-api"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"argon2",
|
"argon2",
|
||||||
@@ -4885,7 +4885,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-application"
|
name = "zesdex-application"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -4903,7 +4903,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-bootstrap"
|
name = "zesdex-bootstrap"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"chrono",
|
"chrono",
|
||||||
@@ -4920,7 +4920,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-daemon"
|
name = "zesdex-daemon"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -4944,7 +4944,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-domain"
|
name = "zesdex-domain"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -4960,7 +4960,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-gateway"
|
name = "zesdex-gateway"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -4987,7 +4987,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-grpc"
|
name = "zesdex-grpc"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -5004,7 +5004,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-infrastructure"
|
name = "zesdex-infrastructure"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"argon2",
|
"argon2",
|
||||||
@@ -5052,7 +5052,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-tui"
|
name = "zesdex-tui"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -5078,7 +5078,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-web"
|
name = "zesdex-web"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -5098,7 +5098,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-ws"
|
name = "zesdex-ws"
|
||||||
version = "1.20.0"
|
version = "1.21.1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
|
|||||||
+1
-1
@@ -15,7 +15,7 @@ members = [
|
|||||||
]
|
]
|
||||||
|
|
||||||
[workspace.package]
|
[workspace.package]
|
||||||
version = "1.20.2"
|
version = "1.22.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
|
|
||||||
|
|||||||
@@ -121,7 +121,66 @@ pub struct CodeSymbol {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// The in-memory symbol index, shared via a global static.
|
/// The in-memory symbol index, shared via a global static.
|
||||||
static SYMBOL_INDEX: Mutex<Option<SymbolIndex>> = Mutex::new(None);
|
///
|
||||||
|
/// Keyed by workspace path: each workspace owns its own `SymbolIndex`, so
|
||||||
|
/// searching one repo never leaks stale symbols from another, and a rebuild
|
||||||
|
/// triggered for workspace A cannot clobber the index of B. The lock is only
|
||||||
|
/// ever held briefly (to check/insert/look up) — never across the I/O-heavy
|
||||||
|
/// walk in `rebuild`, which runs locally and is swapped in under a short lock.
|
||||||
|
static SYMBOL_INDEX: std::sync::OnceLock<Mutex<HashMap<String, SymbolIndex>>> =
|
||||||
|
std::sync::OnceLock::new();
|
||||||
|
|
||||||
|
/// Access the (lazily initialised) global per-workspace symbol index map.
|
||||||
|
fn symbol_index_map() -> &'static Mutex<HashMap<String, SymbolIndex>> {
|
||||||
|
SYMBOL_INDEX.get_or_init(|| Mutex::new(HashMap::new()))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Ensure the per-workspace symbol index is built, returning the symbol count.
|
||||||
|
///
|
||||||
|
/// If `force` is true, or the workspace has no cached (non-empty) index yet,
|
||||||
|
/// the index is rebuilt. The rebuild itself runs OUTSIDE the global lock
|
||||||
|
/// (the walk can take seconds on a large repo), then the result is stored
|
||||||
|
/// under a short lock so concurrent searches never block on the I/O. Returns
|
||||||
|
/// the number of symbols now cached for the workspace.
|
||||||
|
fn ensure_symbol_index(workspace: &str, force: bool) -> Result<usize> {
|
||||||
|
let ready = {
|
||||||
|
let map = symbol_index_map()
|
||||||
|
.lock()
|
||||||
|
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
|
||||||
|
!force && map.get(workspace).is_some_and(|i| !i.is_empty())
|
||||||
|
};
|
||||||
|
|
||||||
|
if !ready {
|
||||||
|
// Rebuild locally, off the global lock (I/O heavy).
|
||||||
|
let mut fresh = SymbolIndex::new();
|
||||||
|
let count = fresh.rebuild(workspace)?;
|
||||||
|
// Swap in under a short lock; keep an existing non-empty index if a
|
||||||
|
// concurrent rebuild already populated this workspace.
|
||||||
|
let mut map = symbol_index_map()
|
||||||
|
.lock()
|
||||||
|
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
|
||||||
|
if map.get(workspace).is_none_or(|i| i.is_empty()) {
|
||||||
|
map.insert(workspace.to_string(), fresh);
|
||||||
|
}
|
||||||
|
return Ok(count);
|
||||||
|
}
|
||||||
|
|
||||||
|
let map = symbol_index_map()
|
||||||
|
.lock()
|
||||||
|
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
|
||||||
|
Ok(map.get(workspace).map_or(0, |i| i.len()))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Lock and return a borrow to the global per-workspace symbol index map.
|
||||||
|
///
|
||||||
|
/// The caller must have called [`ensure_symbol_index`] first, then looks up
|
||||||
|
/// its workspace key; the lookup is short and in-memory, so holding the guard
|
||||||
|
/// for the search is fine.
|
||||||
|
fn symbol_index() -> Result<std::sync::MutexGuard<'static, HashMap<String, SymbolIndex>>> {
|
||||||
|
symbol_index_map()
|
||||||
|
.lock()
|
||||||
|
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))
|
||||||
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
// Language-specific regexes (lazily compiled)
|
// Language-specific regexes (lazily compiled)
|
||||||
@@ -270,6 +329,17 @@ impl SymbolIndex {
|
|||||||
self.symbols.is_empty()
|
self.symbols.is_empty()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Returns true when the cached index must be rebuilt for the given
|
||||||
|
/// workspace — either because nothing has been indexed yet, or because the
|
||||||
|
/// requested workspace differs from the one the index was built for.
|
||||||
|
///
|
||||||
|
/// Without this, searching a *different* workspace after the first one
|
||||||
|
/// silently returns stale symbols from the previously indexed repo
|
||||||
|
/// (a misleading result for a coding agent).
|
||||||
|
pub fn needs_rebuild(&self, workspace: &str) -> bool {
|
||||||
|
self.is_empty() || self.workspace_path.as_deref() != Some(workspace)
|
||||||
|
}
|
||||||
|
|
||||||
pub fn len(&self) -> usize {
|
pub fn len(&self) -> usize {
|
||||||
self.symbols.len()
|
self.symbols.len()
|
||||||
}
|
}
|
||||||
@@ -1250,15 +1320,11 @@ impl Tool for SemanticSearch {
|
|||||||
"semantic search"
|
"semantic search"
|
||||||
);
|
);
|
||||||
|
|
||||||
let mut guard = SYMBOL_INDEX
|
// Build (or load) the per-workspace index without holding the global
|
||||||
.lock()
|
// lock across the I/O-heavy walk.
|
||||||
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
|
let _count = ensure_symbol_index(&workspace, rebuild)?;
|
||||||
let index = guard.get_or_insert_with(SymbolIndex::new);
|
let map = symbol_index()?;
|
||||||
|
let index = map.get(&workspace).expect("index should be ensured");
|
||||||
if rebuild || index.is_empty() {
|
|
||||||
let count = index.rebuild(&workspace)?;
|
|
||||||
debug!(symbol_count = count, "symbol index rebuilt");
|
|
||||||
}
|
|
||||||
|
|
||||||
// Map kind filter to enum
|
// Map kind filter to enum
|
||||||
let target_kind = match kind_filter {
|
let target_kind = match kind_filter {
|
||||||
@@ -1398,12 +1464,15 @@ impl Tool for RebuildIndex {
|
|||||||
|
|
||||||
info!("rebuilding multi-language symbol index");
|
info!("rebuilding multi-language symbol index");
|
||||||
|
|
||||||
let mut guard = SYMBOL_INDEX
|
// Force a rebuild of this workspace's index. The walk runs off the
|
||||||
.lock()
|
// global lock (via ensure_symbol_index) so it cannot stall concurrent
|
||||||
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
|
// searches.
|
||||||
let index = guard.get_or_insert_with(SymbolIndex::new);
|
let count = ensure_symbol_index(&workspace, true)?;
|
||||||
let count = index.rebuild(&workspace)?;
|
let map = symbol_index()?;
|
||||||
let by_lang = index.count_by_language();
|
let by_lang = match map.get(&workspace) {
|
||||||
|
Some(i) => i.count_by_language(),
|
||||||
|
None => Vec::new(),
|
||||||
|
};
|
||||||
|
|
||||||
let mut out = format!(
|
let mut out = format!(
|
||||||
"Symbol index rebuilt successfully. {} symbols indexed.\n\n",
|
"Symbol index rebuilt successfully. {} symbols indexed.\n\n",
|
||||||
@@ -1497,15 +1566,11 @@ impl Tool for ListSymbols {
|
|||||||
.map(|p| p.to_string_lossy().to_string())
|
.map(|p| p.to_string_lossy().to_string())
|
||||||
.unwrap_or_else(|| ".".to_string());
|
.unwrap_or_else(|| ".".to_string());
|
||||||
|
|
||||||
let mut guard = SYMBOL_INDEX
|
// Build (or load) the per-workspace index without holding the global
|
||||||
.lock()
|
// lock across the I/O-heavy walk.
|
||||||
.map_err(|e| anyhow::anyhow!("index lock failed: {e}"))?;
|
let _count = ensure_symbol_index(&workspace, rebuild)?;
|
||||||
let index = guard.get_or_insert_with(SymbolIndex::new);
|
let map = symbol_index()?;
|
||||||
|
let index = map.get(&workspace).expect("index should be ensured");
|
||||||
if rebuild || index.is_empty() {
|
|
||||||
let count = index.rebuild(&workspace)?;
|
|
||||||
info!(symbol_count = count, "symbol index rebuilt for list");
|
|
||||||
}
|
|
||||||
|
|
||||||
let target_lang = match lang_filter {
|
let target_lang = match lang_filter {
|
||||||
"rust" => Some(Language::Rust),
|
"rust" => Some(Language::Rust),
|
||||||
@@ -1746,4 +1811,80 @@ mod tests {
|
|||||||
let index = SymbolIndex::new();
|
let index = SymbolIndex::new();
|
||||||
assert!(index.search("anything", 10).is_empty());
|
assert!(index.search("anything", 10).is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_needs_rebuild_workspace_aware() {
|
||||||
|
let ws_a = std::env::temp_dir().join(format!("ws_a_{}", uuid::Uuid::new_v4()));
|
||||||
|
let ws_b = std::env::temp_dir().join(format!("ws_b_{}", uuid::Uuid::new_v4()));
|
||||||
|
std::fs::create_dir_all(&ws_a).unwrap();
|
||||||
|
std::fs::create_dir_all(&ws_b).unwrap();
|
||||||
|
std::fs::write(ws_a.join("a.rs"), "pub fn fn_in_a() {}\n").unwrap();
|
||||||
|
std::fs::write(ws_b.join("b.rs"), "pub fn fn_in_b() {}\n").unwrap();
|
||||||
|
|
||||||
|
let mut index = SymbolIndex::new();
|
||||||
|
let a = ws_a.to_string_lossy().to_string();
|
||||||
|
let b = ws_b.to_string_lossy().to_string();
|
||||||
|
|
||||||
|
// Fresh index: needs rebuild for any workspace.
|
||||||
|
assert!(index.needs_rebuild(&a));
|
||||||
|
|
||||||
|
// After rebuilding A, searching A needs no rebuild...
|
||||||
|
index.rebuild(&a).unwrap();
|
||||||
|
assert!(!index.needs_rebuild(&a));
|
||||||
|
// ...but searching B DOES (stale index otherwise).
|
||||||
|
assert!(
|
||||||
|
index.needs_rebuild(&b),
|
||||||
|
"workspace switch must trigger rebuild"
|
||||||
|
);
|
||||||
|
|
||||||
|
// Rebuilding B flips the cached workspace.
|
||||||
|
index.rebuild(&b).unwrap();
|
||||||
|
assert!(!index.needs_rebuild(&b));
|
||||||
|
assert!(index.needs_rebuild(&a));
|
||||||
|
|
||||||
|
std::fs::remove_dir_all(&ws_a).ok();
|
||||||
|
std::fs::remove_dir_all(&ws_b).ok();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_ensure_symbol_index_per_workspace_isolation() {
|
||||||
|
let ws_a = std::env::temp_dir().join(format!("iso_a_{}", uuid::Uuid::new_v4()));
|
||||||
|
let ws_b = std::env::temp_dir().join(format!("iso_b_{}", uuid::Uuid::new_v4()));
|
||||||
|
std::fs::create_dir_all(&ws_a).unwrap();
|
||||||
|
std::fs::create_dir_all(&ws_b).unwrap();
|
||||||
|
std::fs::write(ws_a.join("a.rs"), "pub fn only_in_a() {}\n").unwrap();
|
||||||
|
std::fs::write(ws_b.join("b.rs"), "pub fn only_in_b() {}\n").unwrap();
|
||||||
|
|
||||||
|
let a = ws_a.to_string_lossy().to_string();
|
||||||
|
let b = ws_b.to_string_lossy().to_string();
|
||||||
|
|
||||||
|
// Build A and B independently through the shared global helper.
|
||||||
|
let count_a = ensure_symbol_index(&a, false).unwrap();
|
||||||
|
assert!(
|
||||||
|
count_a >= 1,
|
||||||
|
"workspace A should index its fn, got {count_a}"
|
||||||
|
);
|
||||||
|
let count_b = ensure_symbol_index(&b, false).unwrap();
|
||||||
|
assert!(
|
||||||
|
count_b >= 1,
|
||||||
|
"workspace B should index its fn, got {count_b}"
|
||||||
|
);
|
||||||
|
|
||||||
|
// Rebuilding A must not have clobbered B and vice-versa.
|
||||||
|
let count_a_again = ensure_symbol_index(&a, true).unwrap();
|
||||||
|
assert!(count_a_again >= 1);
|
||||||
|
|
||||||
|
// Each workspace's cached index is independently correct.
|
||||||
|
{
|
||||||
|
let map = symbol_index().unwrap();
|
||||||
|
let idx_a = map.get(&a).unwrap();
|
||||||
|
assert!(!idx_a.search("only_in_a", 5).is_empty());
|
||||||
|
assert!(idx_a.search("only_in_b", 5).is_empty());
|
||||||
|
let idx_b = map.get(&b).unwrap();
|
||||||
|
assert!(!idx_b.search("only_in_b", 5).is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
std::fs::remove_dir_all(&ws_a).ok();
|
||||||
|
std::fs::remove_dir_all(&ws_b).ok();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,28 +1,128 @@
|
|||||||
//! Consensus synthesis — reconciles multiple node outputs into one assessment.
|
//! Consensus synthesis — reconciles multiple node outputs into one assessment.
|
||||||
|
//!
|
||||||
|
//! Flow: load settings → resolve LLM credentials → ask the model to distill the
|
||||||
|
//! node outputs into a single consensus (conflicts, agreements, key findings)
|
||||||
|
//! → return the synthesized text. If the LLM call fails for any reason, we
|
||||||
|
//! degrade gracefully to a concatenation-based summary so consensus synthesis
|
||||||
|
//! never breaks the surrounding hive-mind cycle (mirrors the isolated-errors
|
||||||
|
//! philosophy used for the nodes themselves).
|
||||||
|
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use tracing::info;
|
use tracing::{info, warn};
|
||||||
|
|
||||||
use crate::tools::ToolCtx;
|
use zesdex_domain::cms::{AppConfigRepository, SettingsRepository};
|
||||||
|
use zesdex_domain::core::message::ChatMessage;
|
||||||
|
use zesdex_domain::core::Store;
|
||||||
use zesdex_domain::workflow::NodeOutput;
|
use zesdex_domain::workflow::NodeOutput;
|
||||||
|
|
||||||
/// Synthesize a consensus from all node outputs.
|
use crate::persistence::{JsonAppConfigRepository, JsonSettingsRepository};
|
||||||
|
use crate::tools::ToolCtx;
|
||||||
|
|
||||||
|
/// Maximum characters of node output to feed into the synthesis prompt per node.
|
||||||
|
/// Keeps the prompt bounded so a huge/talkative node cannot blow up the request.
|
||||||
|
const MAX_NODE_OUTPUT_CHARS: usize = 4000;
|
||||||
|
|
||||||
|
/// Synthesize a consensus from all node outputs using the LLM.
|
||||||
///
|
///
|
||||||
/// Flow: combine node outputs → return consensus text.
|
/// Flow: combine node outputs → ask the model to reconcile them into a single
|
||||||
/// Uses simple concatenation-based synthesis (avoids LLM call dependency).
|
/// consensus → return the synthesized text. Falls back to a plain
|
||||||
|
/// concatenation summary if the LLM is unreachable or the call fails.
|
||||||
pub async fn synthesize_consensus(nodes: &[NodeOutput], _tool_ctx: &ToolCtx) -> Result<String> {
|
pub async fn synthesize_consensus(nodes: &[NodeOutput], _tool_ctx: &ToolCtx) -> Result<String> {
|
||||||
info!("Synthesizing consensus from {} nodes", nodes.len());
|
info!("Synthesizing consensus from {} nodes", nodes.len());
|
||||||
|
|
||||||
|
let combined = build_combined_body(nodes);
|
||||||
|
|
||||||
|
// 1. Resolve LLM credentials (same source of truth as execute_cycle).
|
||||||
|
let store = Store::new();
|
||||||
|
let settings = JsonSettingsRepository::new()
|
||||||
|
.load(&store.base_dir)
|
||||||
|
.unwrap_or_default();
|
||||||
|
let app_config = JsonAppConfigRepository::new()
|
||||||
|
.load(&store.base_dir)
|
||||||
|
.unwrap_or_default();
|
||||||
|
|
||||||
|
let (provider, model) =
|
||||||
|
crate::subagent::provider::resolve_subagent_provider(&settings, &app_config);
|
||||||
|
let base_url = app_config
|
||||||
|
.providers
|
||||||
|
.get(&provider)
|
||||||
|
.map(|p| p.api_base.clone())
|
||||||
|
.unwrap_or_else(|| zesdex_domain::agent::defaults::DEFAULT_API_BASE.to_string());
|
||||||
|
let api_key = crate::llm::provider::resolve_api_key(&settings, &app_config);
|
||||||
|
|
||||||
|
let client = crate::llm::provider::LlmClient::new(api_key, model, Some(base_url));
|
||||||
|
|
||||||
|
// 2. Build the synthesis prompt.
|
||||||
|
let system_msg = ChatMessage::system(
|
||||||
|
"You are a consensus synthesizer for a multi-agent hive mind. \
|
||||||
|
Several independent nodes analysed a problem and produced the outputs \
|
||||||
|
below. Distill them into ONE coherent consensus report with these \
|
||||||
|
sections:\n\
|
||||||
|
- AGREEMENTS: points multiple nodes converge on.\n\
|
||||||
|
- CONFLICTS: contradictory conclusions, with which node(s) support each side.\n\
|
||||||
|
- KEY FINDINGS: the most important, actionable takeaways.\n\
|
||||||
|
- RECOMMENDATION: a single recommended next action, or 'no clear consensus' \
|
||||||
|
if the outputs are too divergent.\n\
|
||||||
|
Be concise and factual. If a node errored, note it and ignore its content.\n\
|
||||||
|
Do not invent facts not present in the node outputs.",
|
||||||
|
);
|
||||||
|
let user_msg = ChatMessage::user(format!(
|
||||||
|
"Consolidate these {} node outputs into a single consensus:\n\n{}",
|
||||||
|
nodes.len(),
|
||||||
|
combined
|
||||||
|
));
|
||||||
|
|
||||||
|
// 3. Call the model and gracefully degrade on failure.
|
||||||
|
match call_consensus(&client, &[system_msg, user_msg]).await {
|
||||||
|
Ok(text) => {
|
||||||
|
let trimmed = text.trim();
|
||||||
|
if trimmed.is_empty() {
|
||||||
|
warn!("consensus LLM returned empty output; falling back to concat summary");
|
||||||
|
Ok(concat_summary(nodes))
|
||||||
|
} else {
|
||||||
|
Ok(format!(
|
||||||
|
"# Consensus Synthesis\n\n\
|
||||||
|
Nodes synthesized: {}\n\n{}",
|
||||||
|
nodes.len(),
|
||||||
|
trimmed
|
||||||
|
))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
warn!(error = %e, "consensus LLM call failed; falling back to concat summary");
|
||||||
|
Ok(concat_summary(nodes))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Run the LLM consensus call and return the assistant text.
|
||||||
|
async fn call_consensus(
|
||||||
|
client: &crate::llm::provider::LlmClient,
|
||||||
|
messages: &[ChatMessage],
|
||||||
|
) -> Result<String> {
|
||||||
|
use zesdex_application::ports::ProviderService;
|
||||||
|
let (msg, _) = client.chat(messages, None, Some(1024), Some(0.3)).await?;
|
||||||
|
Ok(msg.content.unwrap_or_default())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Build the concatenated node-output body for the prompt.
|
||||||
|
fn build_combined_body(nodes: &[NodeOutput]) -> String {
|
||||||
let mut combined = String::new();
|
let mut combined = String::new();
|
||||||
for node in nodes {
|
for node in nodes {
|
||||||
|
let output = crate::utils::truncate_chars(&node.output, MAX_NODE_OUTPUT_CHARS);
|
||||||
combined.push_str(&format!(
|
combined.push_str(&format!(
|
||||||
"\n## {} — {}\n\n{}\n",
|
"\n## {} — {}\n{}\n",
|
||||||
node.id, node.directive, node.output
|
node.id, node.directive, output
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
combined
|
||||||
|
}
|
||||||
|
|
||||||
Ok(format!(
|
/// Fallback: a plain concatenation summary (the behaviour of the original stub).
|
||||||
"# Consensus Synthesis\n\
|
fn concat_summary(nodes: &[NodeOutput]) -> String {
|
||||||
|
let combined = build_combined_body(nodes);
|
||||||
|
format!(
|
||||||
|
"# Consensus Synthesis\n\n\
|
||||||
Nodes synthesized: {}\n\n\
|
Nodes synthesized: {}\n\n\
|
||||||
## Summary\n\
|
## Summary\n\
|
||||||
The following node outputs were collected:\n\
|
The following node outputs were collected:\n\
|
||||||
@@ -31,5 +131,47 @@ pub async fn synthesize_consensus(nodes: &[NodeOutput], _tool_ctx: &ToolCtx) ->
|
|||||||
Review the individual node outputs above for detailed findings.",
|
Review the individual node outputs above for detailed findings.",
|
||||||
nodes.len(),
|
nodes.len(),
|
||||||
combined
|
combined
|
||||||
))
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use zesdex_domain::workflow::NodeOutput;
|
||||||
|
|
||||||
|
fn node(id: &str, output: &str) -> NodeOutput {
|
||||||
|
NodeOutput {
|
||||||
|
id: id.to_string(),
|
||||||
|
directive: "directive".to_string(),
|
||||||
|
output: output.to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn build_combined_body_truncates_oversized_output() {
|
||||||
|
let long = "é".repeat(MAX_NODE_OUTPUT_CHARS + 500);
|
||||||
|
let body = build_combined_body(&[node("n1", &long)]);
|
||||||
|
// Must contain the header and a char-truncated (<= cap) payload without
|
||||||
|
// panicking on a multi-byte boundary.
|
||||||
|
assert!(body.contains("## n1"));
|
||||||
|
// Header "## n1 — directive\n" ~= 20 chars, so char count stays just above cap.
|
||||||
|
let char_count = body.chars().count();
|
||||||
|
assert!(
|
||||||
|
char_count <= MAX_NODE_OUTPUT_CHARS + 50,
|
||||||
|
"expected body near {MAX_NODE_OUTPUT_CHARS} chars, got {char_count}"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
char_count > 1000,
|
||||||
|
"expected a many-node output, got small: {char_count}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn concat_summary_includes_all_node_ids() {
|
||||||
|
let nodes = vec![node("node-a", "a out"), node("node-b", "b out")];
|
||||||
|
let summary = concat_summary(&nodes);
|
||||||
|
assert!(summary.contains("node-a"));
|
||||||
|
assert!(summary.contains("node-b"));
|
||||||
|
assert!(summary.contains("Nodes synthesized: 2"));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -320,6 +320,18 @@ fn handle_submit_input(state: &mut AppStateRest, text: String) {
|
|||||||
api_base: provider_cfg.as_ref().map(|cfg| cfg.api_base.clone()),
|
api_base: provider_cfg.as_ref().map(|cfg| cfg.api_base.clone()),
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Capture state for the optional background auto-review so it can run
|
||||||
|
// with the same resolved provider. Done here, before `api_key` /
|
||||||
|
// `provider_cfg` are moved into the client below. The reviewer is
|
||||||
|
// fire-and-forget and skips itself when the workspace has no diff.
|
||||||
|
let review_enabled = state.settings.flags.review_enabled;
|
||||||
|
let review_key = api_key.clone();
|
||||||
|
let review_model =
|
||||||
|
zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config);
|
||||||
|
let review_base = provider_cfg.as_ref().map(|c| c.api_base.clone());
|
||||||
|
let review_ws = params.workspace_roots.clone();
|
||||||
|
let review_events = params.turn_events.clone();
|
||||||
|
|
||||||
let client = std::sync::Arc::new(zesdex_infrastructure::llm::provider::LlmClient::new(
|
let client = std::sync::Arc::new(zesdex_infrastructure::llm::provider::LlmClient::new(
|
||||||
api_key,
|
api_key,
|
||||||
zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config),
|
zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config),
|
||||||
@@ -348,6 +360,15 @@ fn handle_submit_input(state: &mut AppStateRest, text: String) {
|
|||||||
use zesdex_application::agent::AgentTurnService;
|
use zesdex_application::agent::AgentTurnService;
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let _ = turn_service.run_turn(params).await;
|
let _ = turn_service.run_turn(params).await;
|
||||||
|
if review_enabled {
|
||||||
|
zesdex_infrastructure::subagent::auto::engine::spawn_background_review(
|
||||||
|
review_ws,
|
||||||
|
review_events,
|
||||||
|
review_key,
|
||||||
|
review_model,
|
||||||
|
review_base,
|
||||||
|
);
|
||||||
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -110,6 +110,16 @@ pub fn spawn_agent_turn(state: &mut AppStateRest, text: String) {
|
|||||||
api_base: api_base.clone(),
|
api_base: api_base.clone(),
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Capture for the optional background auto-review before `api_key` /
|
||||||
|
// `model` / `api_base` move into the client below. The reviewer is
|
||||||
|
// fire-and-forget and skips itself when the workspace has no diff.
|
||||||
|
let review_enabled = state.settings.flags.review_enabled;
|
||||||
|
let review_key = api_key.clone();
|
||||||
|
let review_model = model.clone();
|
||||||
|
let review_base = api_base.clone();
|
||||||
|
let review_ws = workspace_roots.clone();
|
||||||
|
let review_events = turn_events.clone();
|
||||||
|
|
||||||
let client = std::sync::Arc::new(LlmClient::new(api_key, model, api_base));
|
let client = std::sync::Arc::new(LlmClient::new(api_key, model, api_base));
|
||||||
|
|
||||||
let tool_ctx = ToolCtx::builder()
|
let tool_ctx = ToolCtx::builder()
|
||||||
@@ -127,6 +137,15 @@ pub fn spawn_agent_turn(state: &mut AppStateRest, text: String) {
|
|||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let _ = turn_service.run_turn(params).await;
|
let _ = turn_service.run_turn(params).await;
|
||||||
|
if review_enabled {
|
||||||
|
zesdex_infrastructure::subagent::auto::engine::spawn_background_review(
|
||||||
|
review_ws,
|
||||||
|
review_events,
|
||||||
|
review_key,
|
||||||
|
review_model,
|
||||||
|
review_base,
|
||||||
|
);
|
||||||
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -130,6 +130,17 @@ async fn handle_socket(mut socket: WebSocket, state: Arc<WsState>) {
|
|||||||
api_base: None,
|
api_base: None,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Capture for the optional background auto-review
|
||||||
|
// before `api_key`/`model` move into the client.
|
||||||
|
// This minimal WS channel has no settings toggle, so
|
||||||
|
// review fires whenever a prompt runs (consistent
|
||||||
|
// with the default review_enabled=true).
|
||||||
|
let review_key = api_key.clone();
|
||||||
|
let review_model = model.clone();
|
||||||
|
let review_base = None;
|
||||||
|
let review_ws = workspace_roots.clone();
|
||||||
|
let review_events = turn_events.clone();
|
||||||
|
|
||||||
let client = std::sync::Arc::new(
|
let client = std::sync::Arc::new(
|
||||||
zesdex_infrastructure::llm::provider::LlmClient::new(
|
zesdex_infrastructure::llm::provider::LlmClient::new(
|
||||||
api_key, model, None,
|
api_key, model, None,
|
||||||
@@ -159,6 +170,13 @@ async fn handle_socket(mut socket: WebSocket, state: Arc<WsState>) {
|
|||||||
use zesdex_application::agent::AgentTurnService;
|
use zesdex_application::agent::AgentTurnService;
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
let _ = turn_service.run_turn(params).await;
|
let _ = turn_service.run_turn(params).await;
|
||||||
|
zesdex_infrastructure::subagent::auto::engine::spawn_background_review(
|
||||||
|
review_ws,
|
||||||
|
review_events,
|
||||||
|
review_key,
|
||||||
|
review_model,
|
||||||
|
review_base,
|
||||||
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
let tx_clone = tx.clone();
|
let tx_clone = tx.clone();
|
||||||
|
|||||||
Reference in New Issue
Block a user