Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
992e60980c | ||
|
|
74b1ad4302 | ||
|
|
9ad04cf819 | ||
|
|
b46935c606 | ||
|
|
b28a5fe384 | ||
|
|
b66898ea28 | ||
|
|
9aca45cb65 | ||
|
|
4dccf0cee4 | ||
|
|
1f91b44708 | ||
|
|
b25929824a | ||
|
|
3fd9a2b2db | ||
|
|
6a98d52d54 |
@@ -1,3 +1,39 @@
|
|||||||
|
## [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)
|
## [1.19.1](https://github.com/asepharyana/zesdex/compare/v1.19.0...v1.19.1) (2026-08-27)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Generated
+12
-11
@@ -4862,7 +4862,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-api"
|
name = "zesdex-api"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"argon2",
|
"argon2",
|
||||||
@@ -4885,11 +4885,12 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-application"
|
name = "zesdex-application"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
"chrono",
|
"chrono",
|
||||||
|
"futures-util",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
"sha2 0.11.0",
|
"sha2 0.11.0",
|
||||||
@@ -4902,7 +4903,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-bootstrap"
|
name = "zesdex-bootstrap"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"chrono",
|
"chrono",
|
||||||
@@ -4919,7 +4920,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-daemon"
|
name = "zesdex-daemon"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -4943,7 +4944,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-domain"
|
name = "zesdex-domain"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -4959,7 +4960,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-gateway"
|
name = "zesdex-gateway"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -4986,7 +4987,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-grpc"
|
name = "zesdex-grpc"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -5003,7 +5004,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-infrastructure"
|
name = "zesdex-infrastructure"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"argon2",
|
"argon2",
|
||||||
@@ -5051,7 +5052,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-tui"
|
name = "zesdex-tui"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"base64",
|
"base64",
|
||||||
@@ -5077,7 +5078,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-web"
|
name = "zesdex-web"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -5097,7 +5098,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "zesdex-ws"
|
name = "zesdex-ws"
|
||||||
version = "1.18.4"
|
version = "1.19.4"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
|
|||||||
+1
-1
@@ -15,7 +15,7 @@ members = [
|
|||||||
]
|
]
|
||||||
|
|
||||||
[workspace.package]
|
[workspace.package]
|
||||||
version = "1.19.1"
|
version = "1.19.6"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
|
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ uuid.workspace = true
|
|||||||
anyhow.workspace = true
|
anyhow.workspace = true
|
||||||
tracing.workspace = true
|
tracing.workspace = true
|
||||||
tokio.workspace = true
|
tokio.workspace = true
|
||||||
|
futures-util.workspace = true
|
||||||
base64.workspace = true
|
base64.workspace = true
|
||||||
sha2.workspace = true
|
sha2.workspace = true
|
||||||
url.workspace = true
|
url.workspace = true
|
||||||
|
|||||||
@@ -11,6 +11,17 @@ pub trait ToolExecutor: Send + Sync {
|
|||||||
tool_name: &str,
|
tool_name: &str,
|
||||||
args: &serde_json::Value,
|
args: &serde_json::Value,
|
||||||
) -> impl Future<Output = Result<String>> + Send;
|
) -> 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.
|
/// Service for running agent turns asynchronously.
|
||||||
|
|||||||
@@ -187,6 +187,53 @@ async fn execute_tool_call<T: ToolExecutor>(
|
|||||||
output
|
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.
|
// Helper: emit usage event from optional LLM response metadata.
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
@@ -384,10 +431,42 @@ impl<P: ProviderService, T: ToolExecutor> super::AgentTurnService for AgentTurnS
|
|||||||
params.messages.push(assistant_msg);
|
params.messages.push(assistant_msg);
|
||||||
|
|
||||||
// ── Execute each tool call ──────────────────────────
|
// ── 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(),
|
||||||
|
¶ms.turn_events,
|
||||||
|
&tool_calls,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
} else {
|
||||||
|
let mut sequential = Vec::with_capacity(tool_calls.len());
|
||||||
for tc in &tool_calls {
|
for tc in &tool_calls {
|
||||||
let output =
|
let out = execute_tool_call(
|
||||||
execute_tool_call(self.tool_executor.as_ref(), ¶ms.turn_events, tc)
|
self.tool_executor.as_ref(),
|
||||||
|
¶ms.turn_events,
|
||||||
|
tc,
|
||||||
|
)
|
||||||
.await;
|
.await;
|
||||||
|
sequential.push(out);
|
||||||
|
}
|
||||||
|
sequential
|
||||||
|
};
|
||||||
|
|
||||||
|
for (tc, output) in tool_calls.iter().zip(outputs) {
|
||||||
if output.starts_with("Error:") {
|
if output.starts_with("Error:") {
|
||||||
errors.record(&tc.function.name, &output, &mut params.messages);
|
errors.record(&tc.function.name, &output, &mut params.messages);
|
||||||
}
|
}
|
||||||
@@ -479,7 +558,6 @@ pub async fn compact_messages_with_ai<P: ProviderService>(
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn truncate_short_output_is_unchanged() {
|
fn truncate_short_output_is_unchanged() {
|
||||||
let out = "short".to_string();
|
let out = "short".to_string();
|
||||||
@@ -536,4 +614,83 @@ mod tests {
|
|||||||
];
|
];
|
||||||
assert_eq!(conversation_chars(&messages), 3 + 11 + 6);
|
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)));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -80,7 +80,7 @@ impl Default for AppConfig {
|
|||||||
///
|
///
|
||||||
/// ## Defaults
|
/// ## Defaults
|
||||||
/// - Zen provider: `deepseek-v4-flash-free` model
|
/// - 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 role: "default" → zen / deepseek-v4-flash-free, temp 0.7
|
||||||
/// - `default_context_window`: 256,000 tokens
|
/// - `default_context_window`: 256,000 tokens
|
||||||
fn default() -> Self {
|
fn default() -> Self {
|
||||||
@@ -99,7 +99,7 @@ impl Default for AppConfig {
|
|||||||
ProviderConfig {
|
ProviderConfig {
|
||||||
api_base: "https://9router.asepharyana.my.id/v1".to_string(),
|
api_base: "https://9router.asepharyana.my.id/v1".to_string(),
|
||||||
api_key_env: Some("ROUTER_API_KEY".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,
|
default_api_key: None,
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -47,6 +47,7 @@ pub use repository::SettingsRepository;
|
|||||||
pub use service::ConversationService;
|
pub use service::ConversationService;
|
||||||
pub use service::MemoryService;
|
pub use service::MemoryService;
|
||||||
pub use service::SettingsService;
|
pub use service::SettingsService;
|
||||||
|
pub use settings::resolve_effective_model;
|
||||||
pub use settings::InternetMode;
|
pub use settings::InternetMode;
|
||||||
pub use settings::Settings;
|
pub use settings::Settings;
|
||||||
pub use settings::SettingsFlags;
|
pub use settings::SettingsFlags;
|
||||||
|
|||||||
@@ -23,6 +23,8 @@ use std::collections::HashMap;
|
|||||||
|
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
|
use super::app_config::AppConfig;
|
||||||
|
|
||||||
/// Controls how much network access the agent is permitted during a session.
|
/// Controls how much network access the agent is permitted during a session.
|
||||||
///
|
///
|
||||||
/// ## Variants
|
/// ## Variants
|
||||||
@@ -112,3 +114,86 @@ impl Default for Settings {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Pick the effective model name for the main agent.
|
||||||
|
///
|
||||||
|
/// When `settings.provider` is `"claude"` (auto-detected from
|
||||||
|
/// `~/.claude/settings.json`), the provider's `default_model` (or the
|
||||||
|
/// app-level `default_model`) wins over a possibly-stale persisted
|
||||||
|
/// `settings.model`. Otherwise the user's explicit `settings.model` is used.
|
||||||
|
///
|
||||||
|
/// Why: the user's custom Claude endpoint (URL + API key from
|
||||||
|
/// `~/.claude/settings.json`) implies Opus as the model; a stale
|
||||||
|
/// `settings.json` (e.g. "deepseek-v4-flash-free") must not override it.
|
||||||
|
pub fn resolve_effective_model(settings: &Settings, app_config: &AppConfig) -> String {
|
||||||
|
if settings.provider == "claude" {
|
||||||
|
if let Some(m) = app_config
|
||||||
|
.providers
|
||||||
|
.get("claude")
|
||||||
|
.and_then(|p| p.default_model.clone())
|
||||||
|
{
|
||||||
|
return m;
|
||||||
|
}
|
||||||
|
return app_config.default_model.clone();
|
||||||
|
}
|
||||||
|
settings.model.clone()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use crate::cms::app_config::AppConfig;
|
||||||
|
|
||||||
|
fn claude_app_config() -> AppConfig {
|
||||||
|
let mut cfg = AppConfig::default();
|
||||||
|
cfg.providers.insert(
|
||||||
|
"claude".to_string(),
|
||||||
|
crate::cms::ProviderConfig {
|
||||||
|
api_base: "https://9router.example/v1".to_string(),
|
||||||
|
api_key_env: Some("ANTHROPIC_API_KEY".to_string()),
|
||||||
|
default_model: Some("claude-opus-5".to_string()),
|
||||||
|
default_api_key: Some("sk-test".to_string()),
|
||||||
|
},
|
||||||
|
);
|
||||||
|
cfg.default_provider = "claude".to_string();
|
||||||
|
cfg.default_model = "claude-opus-5".to_string();
|
||||||
|
cfg
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn claude_provider_uses_opus_model_over_stale_settings_model() {
|
||||||
|
let settings = Settings {
|
||||||
|
provider: "claude".to_string(),
|
||||||
|
model: "deepseek-v4-flash-free".to_string(), // stale persisted
|
||||||
|
..Settings::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let model = resolve_effective_model(&settings, &claude_app_config());
|
||||||
|
assert_eq!(model, "claude-opus-5");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn non_claude_provider_uses_settings_model() {
|
||||||
|
let settings = Settings {
|
||||||
|
provider: "zen".to_string(),
|
||||||
|
model: "my-model".to_string(),
|
||||||
|
..Settings::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let model = resolve_effective_model(&settings, &AppConfig::default());
|
||||||
|
assert_eq!(model, "my-model");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn claude_falls_back_to_app_default() {
|
||||||
|
let settings = Settings {
|
||||||
|
provider: "claude".to_string(),
|
||||||
|
model: String::new(),
|
||||||
|
..Settings::default()
|
||||||
|
};
|
||||||
|
|
||||||
|
let cfg = AppConfig::default();
|
||||||
|
let model = resolve_effective_model(&settings, &cfg);
|
||||||
|
assert_eq!(model, cfg.default_model);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -15,7 +15,7 @@
|
|||||||
//! read of up to 3 relevant files).
|
//! read of up to 3 relevant files).
|
||||||
//! 4. Join the result and return a concise bullet summary as a tool message.
|
//! 4. Join the result and return a concise bullet summary as a tool message.
|
||||||
|
|
||||||
use anyhow::{Context, Result};
|
use anyhow::Result;
|
||||||
use serde_json::{json, Value};
|
use serde_json::{json, Value};
|
||||||
use tracing::{info, warn};
|
use tracing::{info, warn};
|
||||||
|
|
||||||
@@ -118,7 +118,7 @@ impl Tool for ExploreCodebase {
|
|||||||
model,
|
model,
|
||||||
);
|
);
|
||||||
|
|
||||||
let rt = tokio::runtime::Runtime::new().context("create explore tokio runtime")?;
|
let rt = crate::runtime::runtime();
|
||||||
let result = rt.block_on(run_agent(
|
let result = rt.block_on(run_agent(
|
||||||
subagent_ctx,
|
subagent_ctx,
|
||||||
&directive,
|
&directive,
|
||||||
|
|||||||
@@ -34,6 +34,7 @@ pub mod llm;
|
|||||||
pub mod mcp;
|
pub mod mcp;
|
||||||
pub mod middleware;
|
pub mod middleware;
|
||||||
pub mod persistence;
|
pub mod persistence;
|
||||||
|
pub mod runtime;
|
||||||
pub mod subagent;
|
pub mod subagent;
|
||||||
pub mod tools;
|
pub mod tools;
|
||||||
pub mod utils;
|
pub mod utils;
|
||||||
|
|||||||
@@ -72,27 +72,24 @@ fn detect_claude_settings_provider() -> Option<(ProviderConfig, Option<String>)>
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
impl AppConfigRepository for JsonAppConfigRepository {
|
/// Apply a detected Claude provider + custom model onto an `AppConfig`.
|
||||||
fn load(&self, base_dir: &Path) -> Result<AppConfig, RepositoryError> {
|
///
|
||||||
let path = base_dir.join("app_config.json");
|
/// Pure (no I/O) so it can be unit-tested. Flow:
|
||||||
let mut cfg: AppConfig = match std::fs::read_to_string(&path) {
|
/// 1. Always `insert`s the "claude" provider (refreshing a possibly stale
|
||||||
Ok(s) => serde_json::from_str(&s)?,
|
/// persisted entry with the current base URL + key from settings.json).
|
||||||
Err(e) if e.kind() == std::io::ErrorKind::NotFound => AppConfig::default(),
|
/// 2. Registers known Claude model roles if missing.
|
||||||
Err(e) => return Err(RepositoryError::Io(e)),
|
/// 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.
|
||||||
let defaults = AppConfig::default();
|
fn apply_claude_provider(
|
||||||
for (name, provider) in defaults.providers {
|
cfg: &mut AppConfig,
|
||||||
cfg.providers.entry(name).or_insert(provider);
|
claude_provider: ProviderConfig,
|
||||||
}
|
custom_model: Option<String>,
|
||||||
|
) {
|
||||||
if let Some((claude_provider, custom_model)) = detect_claude_settings_provider() {
|
cfg.providers.insert("claude".to_string(), claude_provider);
|
||||||
cfg.providers
|
|
||||||
.entry("claude".to_string())
|
|
||||||
.or_insert(claude_provider);
|
|
||||||
|
|
||||||
let claude_models: [(&str, &str); 3] = [
|
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-sonnet-5", "claude-sonnet-5"),
|
||||||
("claude-haiku-4-5", "claude-haiku-4-5-20251001"),
|
("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_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)
|
Ok(cfg)
|
||||||
@@ -134,3 +149,106 @@ impl AppConfigRepository for JsonAppConfigRepository {
|
|||||||
Ok(())
|
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");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -23,6 +23,37 @@ use zesdex_domain::subagent_directive;
|
|||||||
/// Maximum number of tool-call iterations before the engine gives up.
|
/// Maximum number of tool-call iterations before the engine gives up.
|
||||||
const MAX_ITERATIONS: u32 = 25;
|
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;
|
||||||
|
|
||||||
|
/// 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
|
/// Emit an `AgentProgress` event onto the turn-event queue, if one is
|
||||||
/// configured in the `ToolCtx`.
|
/// configured in the `ToolCtx`.
|
||||||
fn report_progress(tool_ctx: &ToolCtx, progress: AgentProgress) {
|
fn report_progress(tool_ctx: &ToolCtx, progress: AgentProgress) {
|
||||||
@@ -85,11 +116,17 @@ pub async fn run_agent(
|
|||||||
Some(ctx.base_url.clone()),
|
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
|
// Limited iteration loop so we don't run forever
|
||||||
for iteration in 0..MAX_ITERATIONS {
|
for iteration in 0..MAX_ITERATIONS {
|
||||||
use zesdex_application::ports::ProviderService;
|
use zesdex_application::ports::ProviderService;
|
||||||
let (response_msg, _usage) = client
|
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?;
|
.await?;
|
||||||
|
|
||||||
let content = response_msg.content.clone().unwrap_or_default();
|
let content = response_msg.content.clone().unwrap_or_default();
|
||||||
@@ -113,7 +150,7 @@ pub async fn run_agent(
|
|||||||
&tool_ctx,
|
&tool_ctx,
|
||||||
AgentProgress::running(
|
AgentProgress::running(
|
||||||
"subagent",
|
"subagent",
|
||||||
format!("{}:{}", directive, tool_name),
|
format!("{}:{tool_name}", directive),
|
||||||
Some(tool_name.clone()),
|
Some(tool_name.clone()),
|
||||||
),
|
),
|
||||||
);
|
);
|
||||||
@@ -127,7 +164,29 @@ pub async fn run_agent(
|
|||||||
format!("Unknown tool: {tool_name}")
|
format!("Unknown tool: {tool_name}")
|
||||||
};
|
};
|
||||||
|
|
||||||
messages.push(ChatMessage::tool(tc.id.clone(), result));
|
// 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.to_string();
|
||||||
|
}
|
||||||
|
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 {
|
||||||
|
consecutive_errors = 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
messages.push(ChatMessage::tool(
|
||||||
|
tc.id.clone(),
|
||||||
|
truncate_tool_output(result),
|
||||||
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add assistant response if there was text content
|
// Add assistant response if there was text content
|
||||||
|
|||||||
@@ -63,15 +63,21 @@ impl SubagentProvider {
|
|||||||
|
|
||||||
/// Resolve subagent provider and model from settings.
|
/// 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
|
/// 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]
|
#[instrument]
|
||||||
pub fn resolve_subagent_provider(
|
pub fn resolve_subagent_provider(
|
||||||
settings: &zesdex_domain::cms::Settings,
|
settings: &zesdex_domain::cms::Settings,
|
||||||
app_config: &zesdex_domain::cms::AppConfig,
|
app_config: &zesdex_domain::cms::AppConfig,
|
||||||
) -> (String, String) {
|
) -> (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();
|
let model = settings.model.clone();
|
||||||
|
|
||||||
// Use the default model from the provider config if available
|
// Use the default model from the provider config if available
|
||||||
@@ -80,6 +86,7 @@ pub fn resolve_subagent_provider(
|
|||||||
.providers
|
.providers
|
||||||
.get(&provider)
|
.get(&provider)
|
||||||
.and_then(|p| p.default_model.clone())
|
.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())
|
.unwrap_or_else(|| zesdex_domain::agent::defaults::DEFAULT_MODEL.to_string())
|
||||||
} else {
|
} else {
|
||||||
model
|
model
|
||||||
|
|||||||
@@ -32,7 +32,6 @@ pub fn spawn_subagent(
|
|||||||
) -> thread::JoinHandle<Result<String>> {
|
) -> thread::JoinHandle<Result<String>> {
|
||||||
info!("Spawning subagent: {directive}");
|
info!("Spawning subagent: {directive}");
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
let rt = tokio::runtime::Runtime::new()?;
|
crate::runtime::runtime().block_on(run_agent(ctx, &directive, access, tool_ctx))
|
||||||
rt.block_on(run_agent(ctx, &directive, access, tool_ctx))
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,6 +15,13 @@ impl InfrastructureToolExecutor {
|
|||||||
tools: all_tools(),
|
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 {
|
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)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -59,7 +59,7 @@ pub mod workflow;
|
|||||||
// like `crate::tools::{Tool, ToolCtx}` continue to work.
|
// like `crate::tools::{Tool, ToolCtx}` continue to work.
|
||||||
pub use context::{ToolCtx, ToolCtxBuilder};
|
pub use context::{ToolCtx, ToolCtxBuilder};
|
||||||
pub use graduated::{check_graduated_checks, GraduatedCheck};
|
pub use graduated::{check_graduated_checks, GraduatedCheck};
|
||||||
pub use registry::{all_tools, tool_defs, tool_is_risky};
|
pub use registry::{all_tools, tool_defs, tool_is_parallel_safe, tool_is_risky};
|
||||||
pub use util::{arg_str, execute_cmd, log_write_edit_tool, resolve_path};
|
pub use util::{arg_str, execute_cmd, log_write_edit_tool, resolve_path};
|
||||||
|
|
||||||
/// Common interface every agent-invocable tool implements.
|
/// Common interface every agent-invocable tool implements.
|
||||||
|
|||||||
@@ -123,7 +123,7 @@ impl Tool for ParallelDelegate {
|
|||||||
.collect()
|
.collect()
|
||||||
} else {
|
} else {
|
||||||
// Auto-split using LLM
|
// Auto-split using LLM
|
||||||
let rt = tokio::runtime::Runtime::new()?;
|
let rt = crate::runtime::runtime();
|
||||||
let directives = rt.block_on(auto_split_task(
|
let directives = rt.block_on(auto_split_task(
|
||||||
&task,
|
&task,
|
||||||
max_parallel,
|
max_parallel,
|
||||||
@@ -146,9 +146,13 @@ impl Tool for ParallelDelegate {
|
|||||||
"parallel delegation: starting subagents"
|
"parallel delegation: starting subagents"
|
||||||
);
|
);
|
||||||
|
|
||||||
// Spawn agents in parallel
|
// Spawn agents in parallel — bounded: never more than `max_parallel`
|
||||||
let mut handles = Vec::new();
|
// subagent threads in flight at once (Claude Code-style isolation).
|
||||||
for (i, (directive, access)) in directives.iter().enumerate() {
|
let mut results: Vec<(usize, String, String)> = Vec::new();
|
||||||
|
for batch in directives.chunks(max_parallel) {
|
||||||
|
let mut handles = Vec::with_capacity(batch.len());
|
||||||
|
for (i, (directive, access)) in batch.iter().enumerate() {
|
||||||
|
let global_idx = results.len() + i;
|
||||||
let subagent_ctx = SubagentContext::new(
|
let subagent_ctx = SubagentContext::new(
|
||||||
directive.clone(),
|
directive.clone(),
|
||||||
ctx.clone(),
|
ctx.clone(),
|
||||||
@@ -158,13 +162,12 @@ impl Tool for ParallelDelegate {
|
|||||||
model.clone(),
|
model.clone(),
|
||||||
);
|
);
|
||||||
|
|
||||||
debug!(agent_index = i, access = ?access, "spawning parallel agent");
|
debug!(agent_index = global_idx, access = ?access, "spawning parallel agent");
|
||||||
let handle = spawn_subagent(subagent_ctx, directive.clone(), *access, ctx.clone());
|
let handle = spawn_subagent(subagent_ctx, directive.clone(), *access, ctx.clone());
|
||||||
handles.push((i, handle));
|
handles.push((global_idx, handle));
|
||||||
}
|
}
|
||||||
|
|
||||||
// Join all results
|
// Join this batch before spawning the next.
|
||||||
let mut results: Vec<(usize, String, String)> = Vec::new();
|
|
||||||
for (i, handle) in handles {
|
for (i, handle) in handles {
|
||||||
match handle.join() {
|
match handle.join() {
|
||||||
Ok(Ok(output)) => {
|
Ok(Ok(output)) => {
|
||||||
@@ -185,10 +188,11 @@ impl Tool for ParallelDelegate {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Consolidate results
|
// Consolidate results
|
||||||
if synthesize && results.len() > 1 {
|
if synthesize && results.len() > 1 {
|
||||||
let rt = tokio::runtime::Runtime::new()?;
|
let rt = crate::runtime::runtime();
|
||||||
let consolidated =
|
let consolidated =
|
||||||
rt.block_on(consolidate_results(&results, &base_url, &api_key, &model))?;
|
rt.block_on(consolidate_results(&results, &base_url, &api_key, &model))?;
|
||||||
Ok(format!(
|
Ok(format!(
|
||||||
|
|||||||
@@ -52,6 +52,34 @@ pub fn tool_is_risky(name: &str) -> bool {
|
|||||||
matches!(name, "write" | "delete" | "edit" | "bash" | "git_operator")
|
matches!(name, "write" | "delete" | "edit" | "bash" | "git_operator")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Whether a tool by name is read-only and therefore safe to run *in
|
||||||
|
/// parallel* with other tool calls within the same assistant message.
|
||||||
|
///
|
||||||
|
/// Read-only tools only inspect the workspace (read files, grep, glob,
|
||||||
|
/// semantic search, list symbols, recall memory, web search, directory
|
||||||
|
/// listing). They have no side effects, so concurrent execution cannot
|
||||||
|
/// create data races or conflicting writes.
|
||||||
|
///
|
||||||
|
/// Everything else (edits, writes, deletes, shell, git, planning, memory
|
||||||
|
/// writes, agent/spawn orchestration) stays sequential to preserve
|
||||||
|
/// correctness.
|
||||||
|
pub fn tool_is_parallel_safe(name: &str) -> bool {
|
||||||
|
matches!(
|
||||||
|
name,
|
||||||
|
"read"
|
||||||
|
| "grep"
|
||||||
|
| "glob"
|
||||||
|
| "semantic_search"
|
||||||
|
| "list_symbols"
|
||||||
|
| "web_search"
|
||||||
|
| "recall"
|
||||||
|
| "dir_list"
|
||||||
|
| "pong"
|
||||||
|
| "seq_think"
|
||||||
|
| "dir_cache_update"
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
/// Convert a list of tools into provider-facing `ToolDef` request schema.
|
/// Convert a list of tools into provider-facing `ToolDef` request schema.
|
||||||
pub fn tool_defs(tools: &[Box<dyn super::Tool>]) -> Vec<zesdex_domain::core::ToolDef> {
|
pub fn tool_defs(tools: &[Box<dyn super::Tool>]) -> Vec<zesdex_domain::core::ToolDef> {
|
||||||
tools
|
tools
|
||||||
@@ -66,3 +94,59 @@ pub fn tool_defs(tools: &[Box<dyn super::Tool>]) -> Vec<zesdex_domain::core::Too
|
|||||||
})
|
})
|
||||||
.collect()
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn read_only_tools_are_parallel_safe() {
|
||||||
|
for name in [
|
||||||
|
"read",
|
||||||
|
"grep",
|
||||||
|
"glob",
|
||||||
|
"semantic_search",
|
||||||
|
"list_symbols",
|
||||||
|
"web_search",
|
||||||
|
"recall",
|
||||||
|
"dir_list",
|
||||||
|
"pong",
|
||||||
|
"seq_think",
|
||||||
|
"dir_cache_update",
|
||||||
|
] {
|
||||||
|
assert!(
|
||||||
|
tool_is_parallel_safe(name),
|
||||||
|
"{name} should be parallel-safe"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn mutating_and_shell_tools_are_not_parallel_safe() {
|
||||||
|
for name in [
|
||||||
|
"edit",
|
||||||
|
"write",
|
||||||
|
"delete",
|
||||||
|
"bash",
|
||||||
|
"git_operator",
|
||||||
|
"git_worktree",
|
||||||
|
"remember",
|
||||||
|
"forget",
|
||||||
|
"todowrite",
|
||||||
|
"todofinish",
|
||||||
|
"plan_enter",
|
||||||
|
"plan_ready",
|
||||||
|
"workflow_run",
|
||||||
|
"hive_mind",
|
||||||
|
"spawn_agents",
|
||||||
|
"spawn_pipeline",
|
||||||
|
"parallel_delegate",
|
||||||
|
"explore_codebase",
|
||||||
|
] {
|
||||||
|
assert!(
|
||||||
|
!tool_is_parallel_safe(name),
|
||||||
|
"{name} should NOT be parallel-safe"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -63,7 +63,7 @@ impl crate::tools::Tool for DirCacheUpdate {
|
|||||||
// Persist the resolved paths into the shared DirCache so the TUI
|
// Persist the resolved paths into the shared DirCache so the TUI
|
||||||
// and other tools can read the cached listing without re-scanning.
|
// and other tools can read the cached listing without re-scanning.
|
||||||
let dc = ctx.dir_cache.clone();
|
let dc = ctx.dir_cache.clone();
|
||||||
let rt = tokio::runtime::Runtime::new()?;
|
let rt = crate::runtime::runtime();
|
||||||
rt.block_on(async { dc.write().await.set(resolved).await });
|
rt.block_on(async { dc.write().await.set(resolved).await });
|
||||||
|
|
||||||
info!(count, "directory cache updated");
|
info!(count, "directory cache updated");
|
||||||
|
|||||||
@@ -65,7 +65,7 @@ impl Tool for WorkflowRun {
|
|||||||
zesdex_domain::agent::defaults::DEFAULT_MODEL.to_string(),
|
zesdex_domain::agent::defaults::DEFAULT_MODEL.to_string(),
|
||||||
None,
|
None,
|
||||||
);
|
);
|
||||||
let rt = tokio::runtime::Runtime::new()?;
|
let rt = crate::runtime::runtime();
|
||||||
let result: Vec<String> =
|
let result: Vec<String> =
|
||||||
rt.block_on(async { execute_workflow(&script, ctx, &llm_client).await })?;
|
rt.block_on(async { execute_workflow(&script, ctx, &llm_client).await })?;
|
||||||
|
|
||||||
@@ -229,7 +229,7 @@ impl Tool for HiveMind {
|
|||||||
.ok_or_else(|| anyhow::anyhow!("missing 'cycles' array"))?;
|
.ok_or_else(|| anyhow::anyhow!("missing 'cycles' array"))?;
|
||||||
|
|
||||||
info!("Hive mind starting with {} cycles", cycles_val.len());
|
info!("Hive mind starting with {} cycles", cycles_val.len());
|
||||||
let rt = tokio::runtime::Runtime::new()?;
|
let rt = crate::runtime::runtime();
|
||||||
let mut all_node_outputs = Vec::new();
|
let mut all_node_outputs = Vec::new();
|
||||||
|
|
||||||
for (cycle_idx, cycle_val) in cycles_val.iter().enumerate() {
|
for (cycle_idx, cycle_val) in cycles_val.iter().enumerate() {
|
||||||
|
|||||||
@@ -1,11 +1,15 @@
|
|||||||
//! Hive-mind cycle execution — run one cycle of parallel nodes.
|
//! Hive-mind cycle execution — run one cycle of parallel nodes.
|
||||||
//!
|
//!
|
||||||
//! Flow: load settings → resolve LLM credentials → run all directives in the
|
//! Flow: load settings → resolve LLM credentials → run all directives in the
|
||||||
//! cycle concurrently via try_join_all → collect Vec<NodeOutput>.
|
//! cycle concurrently via a BOUNDED buffer (`buffer_unordered(MAX)`) → collect
|
||||||
|
//! `Vec<NodeOutput>`. Unlike `try_join_all`, a single failing node does NOT
|
||||||
|
//! fail the whole cycle — failed nodes are logged and replaced with an
|
||||||
|
//! `[ERROR]` output so the remaining results are preserved (like Claude
|
||||||
|
//! Code's isolated subagents).
|
||||||
|
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
use futures_util::future::try_join_all;
|
use futures_util::stream::StreamExt;
|
||||||
use tracing::info;
|
use tracing::{info, warn};
|
||||||
|
|
||||||
use zesdex_domain::cms::{AppConfigRepository, SettingsRepository};
|
use zesdex_domain::cms::{AppConfigRepository, SettingsRepository};
|
||||||
use zesdex_domain::core::Store;
|
use zesdex_domain::core::Store;
|
||||||
@@ -15,16 +19,21 @@ use crate::subagent::context::SubagentContext;
|
|||||||
use crate::subagent::division::AccessTier;
|
use crate::subagent::division::AccessTier;
|
||||||
use crate::subagent::engine::run_agent;
|
use crate::subagent::engine::run_agent;
|
||||||
use crate::tools::ToolCtx;
|
use crate::tools::ToolCtx;
|
||||||
use zesdex_domain::workflow::{CognitiveCycle, NodeOutput};
|
use zesdex_domain::workflow::{CognitiveCycle, NodeDirective, NodeOutput};
|
||||||
|
|
||||||
|
/// Maximum number of hive-mind nodes running concurrently per cycle.
|
||||||
|
/// Keeps thread/runtime pressure bounded (Claude Code-style).
|
||||||
|
const MAX_CONCURRENT_NODES: usize = 8;
|
||||||
|
|
||||||
/// Execute one cycle: run each node directive and collect outputs.
|
/// Execute one cycle: run each node directive and collect outputs.
|
||||||
///
|
///
|
||||||
/// Flow:
|
/// Flow:
|
||||||
/// 1. Load `Settings` and `AppConfig` from the store directory.
|
/// 1. Load `Settings` and `AppConfig` from the store directory.
|
||||||
/// 2. Resolve provider, model, base_url, and api_key.
|
/// 2. Resolve provider, model, base_url, and api_key.
|
||||||
/// 3. Spawn all directives concurrently — each builds a `SubagentContext`
|
/// 3. Spawn directives with bounded concurrency — each builds a
|
||||||
/// and calls `run_agent` (Full access).
|
/// `SubagentContext` and calls `run_agent`.
|
||||||
/// 4. `try_join_all` waits for all to complete, then collect `NodeOutput`s.
|
/// 4. Collect `NodeOutput`s; failed nodes are logged and replaced with an
|
||||||
|
/// `[ERROR]` placeholder so the cycle still completes.
|
||||||
pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result<Vec<NodeOutput>> {
|
pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result<Vec<NodeOutput>> {
|
||||||
info!(
|
info!(
|
||||||
"Executing cycle {} with {} directives",
|
"Executing cycle {} with {} directives",
|
||||||
@@ -53,10 +62,9 @@ pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result
|
|||||||
|
|
||||||
let cycle_index = cycle.index;
|
let cycle_index = cycle.index;
|
||||||
|
|
||||||
use zesdex_domain::workflow::NodeDirective;
|
// Run all directives with bounded concurrency. Each node is its own
|
||||||
|
// future; failures are collected, not propagated (isolated errors).
|
||||||
// Run all directives in this cycle concurrently.
|
let tasks: Vec<_> = cycle
|
||||||
let handles: Vec<_> = cycle
|
|
||||||
.directives
|
.directives
|
||||||
.iter()
|
.iter()
|
||||||
.enumerate()
|
.enumerate()
|
||||||
@@ -79,18 +87,33 @@ pub async fn execute_cycle(cycle: &CognitiveCycle, tool_ctx: &ToolCtx) -> Result
|
|||||||
_ => AccessTier::Read,
|
_ => AccessTier::Read,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
let node_id = format!("Node-{}-{}", cycle_index, i);
|
||||||
async move {
|
async move {
|
||||||
let result = run_agent(ctx, &dir, access, tc).await?;
|
match run_agent(ctx, &dir, access, tc).await {
|
||||||
Ok::<NodeOutput, anyhow::Error>(NodeOutput {
|
Ok(output) => Ok::<NodeOutput, anyhow::Error>(NodeOutput {
|
||||||
id: format!("Node-{}-{}", cycle_index, i),
|
id: node_id.clone(),
|
||||||
directive: dir,
|
directive: dir,
|
||||||
output: result,
|
output,
|
||||||
|
}),
|
||||||
|
Err(e) => {
|
||||||
|
warn!(node = %node_id, error = %e, "hive-mind node failed (isolated)");
|
||||||
|
Ok::<NodeOutput, anyhow::Error>(NodeOutput {
|
||||||
|
id: node_id,
|
||||||
|
directive: dir,
|
||||||
|
output: format!("[ERROR] {e}"),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
})
|
})
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
let results = try_join_all(handles).await?;
|
// Bounded concurrency: run at most MAX_CONCURRENT_NODES futures at once.
|
||||||
|
let mut stream = futures_util::stream::iter(tasks).buffer_unordered(MAX_CONCURRENT_NODES);
|
||||||
|
let mut results = Vec::with_capacity(cycle.directives.len());
|
||||||
|
while let Some(node) = stream.next().await {
|
||||||
|
results.push(node?);
|
||||||
|
}
|
||||||
|
|
||||||
Ok(results)
|
Ok(results)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -316,13 +316,13 @@ fn handle_submit_input(state: &mut AppStateRest, text: String) {
|
|||||||
in_flight: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
in_flight: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
|
||||||
abort: state.abort_flag.clone(),
|
abort: state.abort_flag.clone(),
|
||||||
api_key: api_key.clone(),
|
api_key: api_key.clone(),
|
||||||
model: state.settings.model.clone(),
|
model: zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config),
|
||||||
api_base: provider_cfg.as_ref().map(|cfg| cfg.api_base.clone()),
|
api_base: provider_cfg.as_ref().map(|cfg| cfg.api_base.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,
|
||||||
state.settings.model.clone(),
|
zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config),
|
||||||
provider_cfg.map(|cfg| cfg.api_base.clone()),
|
provider_cfg.map(|cfg| cfg.api_base.clone()),
|
||||||
));
|
));
|
||||||
|
|
||||||
@@ -465,14 +465,13 @@ fn handle_compact(state: &mut AppStateRest) {
|
|||||||
.get(provider_name)
|
.get(provider_name)
|
||||||
.cloned()
|
.cloned()
|
||||||
.unwrap_or_default();
|
.unwrap_or_default();
|
||||||
let model = state.settings.model.clone();
|
let model = zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config);
|
||||||
let api_base = provider_cfg.map(|cfg| cfg.api_base.clone());
|
let api_base = provider_cfg.map(|cfg| cfg.api_base.clone());
|
||||||
|
|
||||||
let client = zesdex_infrastructure::llm::provider::LlmClient::new(api_key, model, api_base);
|
let client = zesdex_infrastructure::llm::provider::LlmClient::new(api_key, model, api_base);
|
||||||
|
|
||||||
if let Some(ref mut rt) = state.session_runtime {
|
if let Some(ref mut rt) = state.session_runtime {
|
||||||
let tokio_rt =
|
let tokio_rt = zesdex_infrastructure::runtime::runtime();
|
||||||
tokio::runtime::Runtime::new().expect("create tokio runtime for AI compaction");
|
|
||||||
if let Ok(()) = tokio_rt.block_on(
|
if let Ok(()) = tokio_rt.block_on(
|
||||||
zesdex_application::agent::turn_service::compact_messages_with_ai(
|
zesdex_application::agent::turn_service::compact_messages_with_ai(
|
||||||
&mut rt.messages,
|
&mut rt.messages,
|
||||||
|
|||||||
@@ -80,7 +80,7 @@ pub fn spawn_agent_turn(state: &mut AppStateRest, text: String) {
|
|||||||
// ── Resolve provider configuration ─────────────────────────────────
|
// ── Resolve provider configuration ─────────────────────────────────
|
||||||
let provider_name = &state.settings.provider;
|
let provider_name = &state.settings.provider;
|
||||||
let api_key = resolve_api_key(state, provider_name);
|
let api_key = resolve_api_key(state, provider_name);
|
||||||
let model = state.settings.model.clone();
|
let model = zesdex_domain::cms::resolve_effective_model(&state.settings, &state.app_config);
|
||||||
let api_base = resolve_api_base(state, provider_name);
|
let api_base = resolve_api_base(state, provider_name);
|
||||||
|
|
||||||
// ── Build message list ─────────────────────────────────────────────
|
// ── Build message list ─────────────────────────────────────────────
|
||||||
|
|||||||
Reference in New Issue
Block a user