use edgecrab_plugins::{
build_plugin_skill_prompt, discover_plugins, extract_pre_llm_context, hermes_supports_hook,
invoke_hermes_hook,
};
use edgecrab_tools::config_ref::AppConfigRef;
use edgecrab_tools::read_materialized_set;
use edgecrab_tools::registry::{
ApprovalRequest, ApprovalResponse, DelegationEvent, ToolContext, ToolRegistry,
build_wire_llm_definitions, to_llm_definitions_with_mode,
};
use edgecrab_tools::{ToolSchemaMode, deferred_tool_error_response, is_deferred_not_on_wire};
use edgecrab_types::trajectory::{
TrajectoryMetadata, convert_scratchpad_to_think, save_trajectory,
};
use edgecrab_types::{
AgentError, Content, Cost, Message, Role, ToolError, ToolErrorResponse, Trajectory, Usage,
};
use edgequake_llm::traits::CacheControl;
use edgequake_llm::{CachePromptConfig, LLMProvider};
use std::collections::HashMap;
use std::sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
};
use tokio_util::sync::CancellationToken;
use tracing::Instrument;
use crate::agent::{Agent, ConversationResult, SessionState, resolve_tool_policy};
use crate::compression::{
CompressionParams, CompressionStatus, check_compression_status_for_estimate,
};
use crate::config::edgecrab_home;
use crate::context_references::expand_context_refs_with_policy;
use crate::model_router::{RoutingThresholds, SmartRoutingConfig, resolve_turn_route};
use crate::pricing::{CanonicalUsage, estimate_cost};
use crate::progress_sink::emit_activity;
use crate::prompt_builder::{
PromptBlocks, PromptBuilder, load_global_soul, load_memory_sections, load_preloaded_skills,
load_skill_prompt_parts,
};
#[cfg(test)]
use crate::provider_call::api_call_streaming;
use crate::provider_call::{
ApiCallContext, FINISH_REASON_STREAM_INTERRUPTED, api_call_with_retry, assistant_display_text,
estimate_request_prompt_tokens, response_has_visible_output,
};
use crate::sub_agent_runner::CoreSubAgentRunner;
use crate::task_class::{
apply_visual_ux_session_preview, classify_from_messages, task_class_advisory,
};
use crate::turn_dispatch::{
ToolTurnFinalizeParams, TurnDispatchTrackers, TurnDispatchTrackersView, apply_guardrail_result,
finalize_tool_turn, forward_process_watch_event,
};
use crate::turn_dispatch_policy::pre_dispatch_decision;
const MAX_RETRIES: u32 = 3;
const SKILL_REFLECTION_THRESHOLD: u32 = 5;
#[derive(Debug, Clone, Copy, Default)]
struct TodoStateSnapshot {
active: usize,
blocked: usize,
}
#[derive(Debug, Default, Clone)]
struct RunProgressState {
pending_approvals: Arc<AtomicUsize>,
pending_clarifications: Arc<AtomicUsize>,
child_runs_in_flight: Arc<AtomicUsize>,
}
fn snapshot_todo_state(todo_store: &edgecrab_tools::TodoStore) -> TodoStateSnapshot {
let items = todo_store.read();
TodoStateSnapshot {
active: items
.iter()
.filter(|item| item.status == "not-started" || item.status == "in-progress")
.count(),
blocked: items.iter().filter(|item| item.status == "blocked").count(),
}
}
fn saturating_dec(counter: &AtomicUsize) {
let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |value| {
Some(value.saturating_sub(1))
});
}
fn completion_options_for(config: &crate::agent::AgentConfig) -> edgequake_llm::CompletionOptions {
edgequake_llm::CompletionOptions {
max_tokens: config.model_config.max_tokens.map(|tokens| tokens as usize),
temperature: config.temperature.or(config.model_config.temperature),
reasoning_effort: config.reasoning_effort.clone(),
..Default::default()
}
}
fn provider_prefers_nonstreaming_tool_turns(provider: &dyn LLMProvider) -> bool {
crate::local_provider_policy::prefers_nonstreaming_tool_turns(provider)
}
pub(crate) fn should_use_native_streaming(
provider: &dyn LLMProvider,
tool_defs: &[edgequake_llm::ToolDefinition],
streaming_enabled: bool,
event_tx_present: bool,
) -> bool {
if !streaming_enabled || !event_tx_present || !provider.supports_tool_streaming() {
return false;
}
if tool_defs.is_empty() {
return true;
}
!provider_prefers_nonstreaming_tool_turns(provider)
}
#[derive(Clone, Debug, Default)]
pub struct ExecuteLoopDelegateCtx {
pub depth: u32,
pub agent_id: String,
pub parent_id: Option<String>,
}
#[allow(clippy::too_many_arguments)]
fn build_tool_context(
cwd: &std::path::Path,
app_config_ref: AppConfigRef,
cancel: &CancellationToken,
state_db: &Option<Arc<edgecrab_state::SessionDb>>,
platform: edgecrab_types::Platform,
capability_grants: Option<edgecrab_tools::CapabilityGrants>,
process_table: &Arc<edgecrab_tools::ProcessTable>,
provider: Option<Arc<dyn edgequake_llm::LLMProvider>>,
tool_registry: Option<Arc<ToolRegistry>>,
sub_agent_runner: Option<Arc<dyn edgecrab_tools::SubAgentRunner>>,
delegation_event_tx: Option<tokio::sync::mpsc::UnboundedSender<DelegationEvent>>,
clarify_tx: Option<
tokio::sync::mpsc::UnboundedSender<edgecrab_tools::registry::ClarifyRequest>,
>,
approval_tx: Option<tokio::sync::mpsc::UnboundedSender<ApprovalRequest>>,
tool_progress_tx: Option<
tokio::sync::mpsc::UnboundedSender<edgecrab_tools::ToolProgressUpdate>,
>,
watch_notification_tx: Option<
tokio::sync::mpsc::UnboundedSender<edgecrab_tools::process_table::WatchEvent>,
>,
gateway_sender: Option<Arc<dyn edgecrab_tools::registry::GatewaySender>>,
origin_chat: Option<edgecrab_types::OriginChat>,
current_tool_call_id: Option<String>,
current_tool_name: Option<String>,
conversation_session_id: &str,
todo_store: Option<Arc<edgecrab_tools::TodoStore>>,
injected_messages: Option<Arc<tokio::sync::Mutex<Vec<Message>>>>,
mutation_turn: Option<Arc<edgecrab_tools::MutationTurnState>>,
lsp_gate: Option<Arc<dyn edgecrab_tools::LspGate>>,
delegate_ctx: Option<ExecuteLoopDelegateCtx>,
kanban_task_id: Option<String>,
materialized_tools: Option<Arc<std::sync::RwLock<edgecrab_tools::MaterializedToolSet>>>,
skills_zone_dirty: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> ToolContext {
let (delegate_depth, delegate_agent_id, delegate_parent_id) = match &delegate_ctx {
Some(d) => (d.depth, Some(d.agent_id.clone()), d.parent_id.clone()),
None => (0, None, None),
};
let on_skills_changed: Option<Arc<dyn Fn() + Send + Sync>> = skills_zone_dirty.map(|flag| {
Arc::new(move || {
crate::prompt_builder::invalidate_skills_cache();
flag.store(true, std::sync::atomic::Ordering::Relaxed);
}) as Arc<dyn Fn() + Send + Sync>
});
ToolContext {
task_id: uuid::Uuid::new_v4().to_string(),
cwd: cwd.to_path_buf(),
session_id: conversation_session_id.to_string(),
user_task: None,
cancel: cancel.clone(),
config: app_config_ref,
state_db: state_db.clone(),
platform,
capability_grants,
process_table: Some(process_table.clone()),
provider,
tool_registry,
delegate_depth,
delegate_agent_id,
delegate_parent_id,
sub_agent_runner,
delegation_event_tx,
clarify_tx,
approval_tx,
on_skills_changed,
gateway_sender,
origin_chat: origin_chat.clone(),
session_key: Some(
origin_chat
.as_ref()
.map(edgecrab_types::OriginChat::session_key)
.unwrap_or_else(|| conversation_session_id.to_string()),
),
todo_store,
current_tool_call_id,
current_tool_name,
injected_messages,
tool_progress_tx,
watch_notification_tx,
mutation_turn,
lsp_gate,
kanban_task_id,
materialized_tools,
}
}
enum LoopAction {
Continue,
Done(String),
PartialAbort { reason: String },
GuardrailHalt,
}
const MAX_DELEGATE_TASK_CALLS_PER_TURN: usize = 3;
#[derive(Clone)]
struct DispatchContext {
cwd: std::path::PathBuf,
registry: Option<Arc<ToolRegistry>>,
cancel: CancellationToken,
state_db: Option<Arc<edgecrab_state::SessionDb>>,
platform: edgecrab_types::Platform,
capability_grants: Option<edgecrab_tools::CapabilityGrants>,
process_table: Arc<edgecrab_tools::ProcessTable>,
provider: Option<Arc<dyn edgequake_llm::LLMProvider>>,
gateway_sender: Option<Arc<dyn edgecrab_tools::registry::GatewaySender>>,
sub_agent_runner: Option<Arc<dyn edgecrab_tools::SubAgentRunner>>,
event_tx: Option<tokio::sync::mpsc::UnboundedSender<crate::StreamEvent>>,
delegation_event_tx: Option<tokio::sync::mpsc::UnboundedSender<DelegationEvent>>,
clarify_tx:
Option<tokio::sync::mpsc::UnboundedSender<edgecrab_tools::registry::ClarifyRequest>>,
approval_tx: Option<tokio::sync::mpsc::UnboundedSender<ApprovalRequest>>,
origin_chat: Option<edgecrab_types::OriginChat>,
app_config_ref: AppConfigRef,
conversation_session_id: String,
todo_store: Option<Arc<edgecrab_tools::TodoStore>>,
capability_suppressions: Arc<Mutex<HashMap<String, ToolErrorResponse>>>,
discovered_plugins: Option<Arc<edgecrab_plugins::PluginDiscovery>>,
spill_seq: Arc<crate::tool_result_spill::SpillSequence>,
context_engine: Option<Arc<dyn crate::context_engine::ContextEngine>>,
engine_tool_names: Arc<std::collections::HashSet<String>>,
mutation_turn: Arc<edgecrab_tools::MutationTurnState>,
lsp_gate: Option<Arc<dyn edgecrab_tools::LspGate>>,
tool_progress_tx:
Option<tokio::sync::mpsc::UnboundedSender<edgecrab_tools::ToolProgressUpdate>>,
watch_notification_tx:
Option<tokio::sync::mpsc::UnboundedSender<edgecrab_tools::process_table::WatchEvent>>,
delegate_ctx: Option<ExecuteLoopDelegateCtx>,
kanban_task_id: Option<String>,
materialized_tools: Option<Arc<std::sync::RwLock<edgecrab_tools::MaterializedToolSet>>>,
subdirectory_hints:
Option<Arc<tokio::sync::Mutex<crate::subdirectory_hints::SubdirectoryHintTracker>>>,
skills_zone_dirty: Option<Arc<std::sync::atomic::AtomicBool>>,
}
fn post_write_lsp_gate(
app_config_ref: &edgecrab_tools::AppConfigRef,
) -> Option<Arc<dyn edgecrab_tools::LspGate>> {
if app_config_ref.lsp_enabled {
Some(Arc::new(edgecrab_lsp::EdgecrabLspGate))
} else {
None
}
}
impl Agent {
pub(crate) async fn execute_loop(
&self,
user_message: &str,
system_message: Option<&str>,
history: Option<Vec<Message>>,
event_tx: Option<&tokio::sync::mpsc::UnboundedSender<crate::StreamEvent>>,
cwd_override: Option<&std::path::Path>,
delegate_ctx: Option<ExecuteLoopDelegateCtx>,
) -> Result<ConversationResult, AgentError> {
tracing::info!(
msg_len = user_message.len(),
has_event_tx = event_tx.is_some(),
"execute_loop: entered"
);
let _conversation_guard = self.conversation_lock.lock().await;
tracing::info!("execute_loop: acquired conversation_lock");
self.budget.reset();
let cancel = {
let mut guard = self.cancel.lock().expect("cancel mutex not poisoned");
if guard.is_cancelled() {
*guard = CancellationToken::new();
}
guard.clone()
};
let mut steer_rx = {
let mut guard = self.steer_rx.lock().expect("steer_rx mutex not poisoned");
guard.take()
};
if let Some(tx) = event_tx.cloned() {
*self
.steer_event_tx
.lock()
.expect("steer_event_tx mutex not poisoned") = Some(tx);
}
self.steer_pending
.store(0, std::sync::atomic::Ordering::Relaxed);
let mutation_turn = Arc::new(edgecrab_tools::MutationTurnState::new());
mutation_turn.clear();
edgecrab_tools::tools::checkpoint::checkpoint_new_turn();
let config = self.config.read().await.clone();
let provider = self.provider.read().await.clone();
let tool_registry = self.tool_registry.read().await.clone();
let mut session = {
let mut shared = self.session.write().await;
if let Some(hist) = history {
shared.messages = hist;
}
shared.clone()
};
session.last_run_outcome = None;
let conversation_session_id = session
.session_id
.clone()
.or_else(|| config.session_id.clone())
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
if session.session_id.is_none() {
session.session_id = Some(conversation_session_id.clone());
}
let cwd = cwd_override
.map(std::path::Path::to_path_buf)
.unwrap_or_else(|| {
std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."))
});
let materialized_tools = Arc::new(std::sync::RwLock::new(
edgecrab_tools::MaterializedToolSet::new(),
));
let subdirectory_hints = if config.skip_context_files {
None
} else {
Some(Arc::new(tokio::sync::Mutex::new(
crate::subdirectory_hints::SubdirectoryHintTracker::new(cwd.clone()),
)))
};
let gateway_running = self.gateway_sender.read().await.is_some();
let tool_policy = resolve_tool_policy(&config);
let expanded_enabled = tool_policy.expanded_enabled.clone();
let expanded_disabled = tool_policy.expanded_disabled.clone();
let mut app_config_ref = config.to_app_config_ref(gateway_running, &tool_policy);
let lsp_gate = post_write_lsp_gate(&app_config_ref);
if !config.terminal_env_passthrough.is_empty() {
edgecrab_tools::tools::backends::local::register_env_passthrough(
&config.terminal_env_passthrough,
);
}
let (mut active_tool_defs, tool_names_for_prompt, session_tool_schemas) =
if let Some(ref registry) = tool_registry {
let ctx = build_tool_context(
&cwd,
app_config_ref.clone(),
&cancel,
&self.state_db,
config.platform,
config.capability_grants,
&self.process_table,
Some(provider.clone()),
tool_registry.clone(),
None,
None,
None, None, None, None, self.gateway_sender.read().await.clone(),
config.origin_chat.clone(),
None, None, "schema-resolution", Some(self.todo_store.clone()),
None, None, None, delegate_ctx.clone(),
config.kanban_task_id.clone(),
Some(materialized_tools.clone()),
Some(self.skills_zone_dirty.clone()),
);
let enabled_filter = if config.enabled_toolsets.is_empty()
|| edgecrab_tools::toolsets::contains_all_sentinel(&config.enabled_toolsets)
|| expanded_enabled.is_empty()
{
None
} else {
Some(expanded_enabled.as_slice())
};
let disabled_filter = if expanded_disabled.is_empty() {
None
} else {
Some(expanded_disabled.as_slice())
};
let schemas = registry.get_definitions(enabled_filter, disabled_filter, &ctx);
let names: Vec<String> = schemas.iter().map(|s| s.name.clone()).collect();
let materialized = read_materialized_set(Some(&materialized_tools));
let hot_full =
crate::local_provider_policy::is_local_inference_provider(provider.name());
let wire = build_wire_llm_definitions(
registry,
&ctx,
enabled_filter,
disabled_filter,
config.tool_schema_mode,
&materialized,
hot_full,
);
(wire, names, Some(schemas))
} else {
(Vec::new(), Vec::new(), None)
};
let effective_schema_mode = edgecrab_tools::resolve_effective_schema_mode(
config.tool_schema_mode,
tool_names_for_prompt.len(),
);
app_config_ref.tool_schema_mode = effective_schema_mode;
let hot_schema_full =
crate::local_provider_policy::is_local_inference_provider(provider.name());
let engine_for_dispatch = self.context_engine.clone();
let engine_tool_names: std::collections::HashSet<String> =
if let Some(ref engine) = self.context_engine {
let engine_schemas = engine.get_tool_schemas();
if !engine_schemas.is_empty() {
let capped = &engine_schemas[..engine_schemas
.len()
.min(crate::context_engine::MAX_ENGINE_TOOLS)];
active_tool_defs
.extend(to_llm_definitions_with_mode(capped, effective_schema_mode));
capped.iter().map(|s| s.name.clone()).collect()
} else {
std::collections::HashSet::new()
}
} else {
std::collections::HashSet::new()
};
let engine_tool_names = Arc::new(engine_tool_names);
let discovered_plugins = discover_plugins(&config.plugins_config, config.platform).ok();
if session.cached_system_prompt.is_none() {
if let Some(explicit) = system_message {
session.cached_system_prompt = Some(explicit.to_string());
} else {
let home = edgecrab_home();
let memory_sections = if config.skip_memory {
Vec::new()
} else {
load_memory_sections(&home)
};
let platform_str = config.platform.to_string();
let mut disabled_skills = config.skills_config.disabled.clone();
if let Some(platform_disabled) =
config.skills_config.platform_disabled.get(&platform_str)
{
disabled_skills.extend(platform_disabled.iter().cloned());
}
let toolsets_for_prompt = if let Some(registry) = tool_registry.as_ref() {
available_toolsets_for_prompt(registry, &tool_names_for_prompt)
} else {
Vec::new()
};
let skill_parts = load_skill_prompt_parts(
&home,
&disabled_skills,
Some(&tool_names_for_prompt),
Some(&toolsets_for_prompt),
);
let preloaded_content = load_preloaded_skills(
&home,
&config.skills_config.external_dirs,
&config.skills_config.preloaded,
Some(&conversation_session_id),
);
let plugin_skill_prompt = discovered_plugins
.as_ref()
.and_then(build_plugin_skill_prompt);
let extra_skill_dynamic: Option<String> =
match (preloaded_content.is_empty(), plugin_skill_prompt) {
(false, Some(plugin_summary)) => {
Some(format!("{preloaded_content}\n\n{plugin_summary}"))
}
(false, None) => Some(preloaded_content),
(true, Some(plugin_summary)) => Some(plugin_summary),
(true, None) => None,
};
let global_soul = load_global_soul(&home);
let has_filesystem_sensitive_tools = tool_names_for_prompt.iter().any(|name| {
matches!(
name.as_str(),
"read_file"
| "write_file"
| "patch"
| "search_files"
| "terminal"
| "execute_code"
)
});
let execution_guidance = has_filesystem_sensitive_tools.then(|| {
edgecrab_tools::describe_execution_filesystem(&app_config_ref, &cwd)
.render_prompt_block()
});
let wire_for_prompt =
if effective_schema_mode == edgecrab_tools::ToolSchemaMode::Indexed {
let mut wire: Vec<String> = edgecrab_tools::INDEXED_HOT_TOOLS
.iter()
.map(|s| (*s).to_string())
.collect();
wire.push(edgecrab_tools::TOOL_SEARCH_NAME.to_string());
wire.extend(read_materialized_set(Some(&materialized_tools)));
Some(wire)
} else {
None
};
let blocks = PromptBuilder::new(config.platform)
.skip_context_files(config.skip_context_files)
.execution_environment_guidance(execution_guidance)
.available_tools(tool_names_for_prompt.clone())
.wire_tools(wire_for_prompt)
.tool_schema_mode(effective_schema_mode)
.model_name(Some(config.model.clone()))
.session_id(Some(conversation_session_id.clone()))
.build_blocks(
global_soul.as_deref(), Some(&cwd),
&memory_sections,
skill_parts.as_ref(),
extra_skill_dynamic.as_deref(),
);
let mut dynamic = blocks.dynamic;
if let Some(ref custom_prompt) = config.custom_system_prompt {
if !dynamic.is_empty() {
dynamic.push_str("\n\n");
}
dynamic.push_str(custom_prompt);
}
if let Some(ref addon) = config.personality_addon {
if !dynamic.is_empty() {
dynamic.push_str("\n\n");
}
dynamic.push_str(&format!("## Personality\n\n{addon}"));
}
if effective_schema_mode == edgecrab_tools::ToolSchemaMode::Indexed
&& let Some(ref schemas) = session_tool_schemas
{
let (_, deferred) = edgecrab_tools::partition_schemas(
schemas,
&std::collections::HashSet::new(),
);
let mut categories: Vec<String> = Vec::new();
if let Some(ref registry) = tool_registry {
let mut seen = std::collections::BTreeSet::new();
for schema in &deferred {
if let Some(toolset) = registry.toolset_for_tool(&schema.name)
&& seen.insert(toolset.clone())
{
categories.push(toolset);
}
}
}
let index = edgecrab_tools::format_deferred_index(deferred.len(), &categories);
if !index.is_empty() {
if !dynamic.is_empty() {
dynamic.push_str("\n\n");
}
dynamic.push_str(&index);
}
}
let stable = blocks.stable;
let semi_stable = blocks.semi_stable;
let combined = PromptBlocks {
stable: stable.clone(),
semi_stable: semi_stable.clone(),
dynamic,
}
.combined();
session.cached_stable_prompt = Some(stable);
session.cached_semi_stable_prompt = if semi_stable.is_empty() {
None
} else {
Some(semi_stable)
};
session.cached_system_prompt = Some(combined);
}
}
let is_first_turn = session.messages.is_empty();
let discovered_plugins = discovered_plugins.map(Arc::new);
if let Some(discovery) = discovered_plugins.as_ref()
&& is_first_turn
{
for plugin in discovery
.plugins
.iter()
.filter(|plugin| hermes_supports_hook(plugin, "on_session_start"))
{
if let Err(error) = invoke_hermes_hook(
plugin,
"on_session_start",
serde_json::json!({
"session_id": &conversation_session_id,
"model": &config.model,
"platform": config.platform.to_string(),
}),
)
.await
{
tracing::warn!(plugin = %plugin.name, ?error, "Hermes on_session_start hook failed");
}
}
}
if is_first_turn {
crate::lifecycle_hooks::emit_global(
crate::lifecycle_hooks::LifecycleEvent::SessionStart,
serde_json::json!({
"session_id": &conversation_session_id,
"model": &config.model,
"platform": config.platform.to_string(),
}),
);
}
let context_path_policy = app_config_ref.file_path_policy(&cwd);
let mut expansion =
expand_context_refs_with_policy(user_message, &cwd, &context_path_policy);
if !expansion.refs_found.is_empty() {
tracing::debug!(
refs = expansion.refs_found.len(),
errors = expansion.errors.len(),
"expanded @context references"
);
}
for err in &expansion.errors {
tracing::warn!(error = %err, "context reference expansion error");
}
if !expansion.refs_found.is_empty() {
let context_window =
CompressionParams::from_model_config(&config.model, &config.compression)
.context_window;
let injected_chars = expansion.expanded.len().saturating_sub(user_message.len());
let injected_tokens = injected_chars / 4;
let hard_limit = context_window / 2; let soft_limit = context_window / 4;
if injected_tokens > hard_limit {
tracing::warn!(
injected_tokens,
hard_limit,
"@context injection exceeds 50% of context window — stripping injected content"
);
let notice = format!(
"{user_message}\n\n[Warning: @context injection (~{injected_tokens} tokens) \
exceeds the 50% context-window limit ({hard_limit} tokens). \
Injected content was removed to protect the context budget.]"
);
expansion.expanded = notice;
expansion.budget_blocked = true;
} else if injected_tokens > soft_limit {
tracing::warn!(
injected_tokens,
soft_limit,
"@context injection exceeds 25% of context window — approaching budget limit"
);
expansion.budget_warning = true;
}
}
if let Some(discovery) = discovered_plugins.as_ref() {
let history_json =
serde_json::to_value(&session.messages).unwrap_or_else(|_| serde_json::json!([]));
let mut injected_context = Vec::new();
for plugin in discovery
.plugins
.iter()
.filter(|plugin| hermes_supports_hook(plugin, "pre_llm_call"))
{
match invoke_hermes_hook(
plugin,
"pre_llm_call",
serde_json::json!({
"session_id": &conversation_session_id,
"user_message": &expansion.expanded,
"conversation_history": history_json,
"is_first_turn": is_first_turn,
"model": &config.model,
"platform": config.platform.to_string(),
}),
)
.await
{
Ok(results) => injected_context.extend(extract_pre_llm_context(&results)),
Err(error) => {
tracing::warn!(plugin = %plugin.name, ?error, "Hermes pre_llm_call hook failed");
}
}
}
if !injected_context.is_empty() {
expansion.expanded = format!(
"{}\n\n{}",
expansion.expanded,
injected_context.join("\n\n")
);
}
}
let smart_routing = SmartRoutingConfig {
enabled: config.model_config.smart_routing.enabled,
cheap_model: config.model_config.smart_routing.cheap_model.clone(),
cheap_base_url: config.model_config.smart_routing.cheap_base_url.clone(),
cheap_api_key_env: config.model_config.smart_routing.cheap_api_key_env.clone(),
thresholds: RoutingThresholds::default(),
};
let route = resolve_turn_route(&expansion.expanded, &config.model_config, &smart_routing);
if route.is_primary {
session.smart_routing_strong_turns =
session.smart_routing_strong_turns.saturating_add(1);
} else {
session.smart_routing_cheap_turns = session.smart_routing_cheap_turns.saturating_add(1);
}
if let Some(ref label) = route.label {
tracing::info!(route = %label, "model routing decision");
}
let (mut effective_provider, smart_routed_provider_active) = if !route.is_primary {
if let Some((prov_name, model_name)) = route.model.split_once('/') {
let canonical = edgecrab_tools::vision_models::normalize_provider_name(prov_name);
let cheap_opt: Option<Arc<dyn LLMProvider>> = if canonical == "vscode-copilot" {
match edgecrab_tools::create_provider_for_model(&canonical, model_name) {
Ok(p) => Some(p),
Err(e) => {
tracing::warn!(error = %e, "failed to create copilot provider, using primary");
None
}
}
} else {
let is_gemini_canonical = matches!(
canonical.as_str(),
"google" | "gemini" | "vertex" | "vertexai"
);
let primary_is_vertex = provider.name() == "vertex-ai";
let (effective_canonical, effective_model) =
if is_gemini_canonical && primary_is_vertex {
tracing::info!(
cheap_model = %route.model,
"smart routing: using Vertex AI endpoint for cheap Gemini model \
(primary is vertex-ai)"
);
let bare = model_name.strip_prefix("vertexai:").unwrap_or(model_name);
("vertexai", bare)
} else {
(canonical.as_str(), model_name)
};
match edgecrab_tools::create_provider_for_model(
effective_canonical,
effective_model,
) {
Ok(p) => Some(p),
Err(e) => {
tracing::warn!(error = %e, "failed to create cheap model provider, using primary");
None
}
}
};
match cheap_opt {
Some(cheap) => {
tracing::info!(model = %route.model, "using smart-routed cheap model");
(cheap, true)
}
None => (provider.clone(), false),
}
} else {
(provider.clone(), false)
}
} else {
(provider.clone(), false)
};
app_config_ref.local_write_create_dirs =
crate::local_provider_policy::effective_local_write_create_dirs(
app_config_ref.local_write_create_dirs,
effective_provider.name(),
);
crate::local_provider_policy::log_local_harness_activated(
effective_provider.name(),
!active_tool_defs.is_empty(),
app_config_ref.local_write_create_dirs,
);
if crate::local_provider_policy::local_tool_harness_active(
effective_provider.name(),
!active_tool_defs.is_empty(),
) {
edgecrab_tools::tool_call_pipeline::log_pipeline_activated();
}
let injection_threats = crate::prompt_builder::scan_for_injection(&expansion.expanded);
if !injection_threats.is_empty() {
tracing::warn!(
threats = injection_threats.len(),
"prompt injection patterns detected in user input"
);
for threat in &injection_threats {
tracing::warn!(
pattern = %threat.pattern_name,
severity = ?threat.severity,
"injection threat"
);
}
}
session.messages.push(Message::user(&expansion.expanded));
session.user_turn_count += 1;
let mut tool_defs_dirty = false;
if effective_schema_mode == ToolSchemaMode::Indexed
&& tool_names_for_prompt.iter().any(|n| n == "tool_search")
&& let Some(ref schemas) = session_tool_schemas
{
let mat = read_materialized_set(Some(&materialized_tools));
let hits = edgecrab_tools::prefetch_tools_for_user_message(
&expansion.expanded,
schemas,
&mat,
edgecrab_tools::DEFAULT_PREFETCH_LIMIT,
);
if !hits.is_empty() {
let _ = edgecrab_tools::materialize_tool_names(
&hits,
schemas,
&materialized_tools,
app_config_ref.max_materialized_tools,
edgecrab_tools::MaterializeSchemaStyle::Compact,
);
tool_defs_dirty = true;
tracing::debug!(
count = hits.len(),
tools = ?hits,
"indexed schema: turn-start tool prefetch"
);
}
}
if !session.task_class_advisory_sent {
let class = classify_from_messages(&session.messages);
if matches!(class, crate::task_class::TaskClass::VisualUx)
&& effective_schema_mode == ToolSchemaMode::Indexed
&& let Some(ref schemas) = session_tool_schemas
{
let visual_tools = vec![
"browser_navigate".to_string(),
"browser_snapshot".to_string(),
];
let outcome = edgecrab_tools::materialize_tool_names(
&visual_tools,
schemas,
&materialized_tools,
app_config_ref.max_materialized_tools,
edgecrab_tools::MaterializeSchemaStyle::Compact,
);
if !outcome.activated.is_empty() {
tool_defs_dirty = true;
tracing::debug!(
tools = ?outcome.activated,
"indexed schema: VisualUx browser verify prefetch"
);
}
}
if apply_visual_ux_session_preview(class, &config.security_preview) {
tracing::info!(
?class,
"harness: session-scoped security.preview enabled for visual_ux task"
);
}
if let Some(advisory) = task_class_advisory(class, Some(cwd.as_path())) {
session.messages.push(Message::user(&advisory));
session.task_class_advisory_sent = true;
}
}
let initial_turn_tool_call_count = session.session_tool_call_count;
let sub_agent_runner: Option<Arc<dyn edgecrab_tools::SubAgentRunner>> =
if let Some(ref registry) = tool_registry {
Some(Arc::new(CoreSubAgentRunner::new(
provider.clone(),
registry.clone(),
config.platform,
config.model.clone(),
)))
} else {
None
};
let (clarify_req_tx, mut clarify_req_rx) =
tokio::sync::mpsc::unbounded_channel::<edgecrab_tools::registry::ClarifyRequest>();
let (approval_req_tx, mut approval_req_rx) =
tokio::sync::mpsc::unbounded_channel::<ApprovalRequest>();
let (delegation_req_tx, mut delegation_req_rx) =
tokio::sync::mpsc::unbounded_channel::<DelegationEvent>();
let run_progress = RunProgressState::default();
if let Some(ev_tx) = event_tx {
let clarify_ev_tx = ev_tx.clone();
let pending_clarifications = run_progress.pending_clarifications.clone();
tokio::spawn(async move {
while let Some(req) = clarify_req_rx.recv().await {
pending_clarifications.fetch_add(1, Ordering::Relaxed);
let (answer_tx, answer_rx) = tokio::sync::oneshot::channel::<String>();
if clarify_ev_tx
.send(crate::StreamEvent::Clarify {
question: req.question,
choices: req.choices,
response_tx: answer_tx,
})
.is_ok()
{
let answer = answer_rx.await.unwrap_or_default();
let _ = req.response_tx.send(answer);
} else {
let _ = req.response_tx.send(String::new());
}
saturating_dec(&pending_clarifications);
}
});
let approval_ev_tx = ev_tx.clone();
let pending_approvals = run_progress.pending_approvals.clone();
tokio::spawn(async move {
while let Some(req) = approval_req_rx.recv().await {
pending_approvals.fetch_add(1, Ordering::Relaxed);
emit_activity(
Some(&approval_ev_tx),
edgecrab_tools::tool_progress_tail::format_approval_waiting(&req.command),
);
let (decision_tx, decision_rx) =
tokio::sync::oneshot::channel::<crate::ApprovalChoice>();
if approval_ev_tx
.send(crate::StreamEvent::Approval {
command: req.command,
full_command: req.full_command,
reasons: req.reasons,
kind: req.kind,
response_tx: decision_tx,
})
.is_ok()
{
let mapped = match decision_rx.await {
Ok(crate::ApprovalChoice::Once) => ApprovalResponse::Once,
Ok(crate::ApprovalChoice::Session) => ApprovalResponse::Session,
Ok(crate::ApprovalChoice::Always) => ApprovalResponse::Always,
Ok(crate::ApprovalChoice::Deny) | Err(_) => ApprovalResponse::Deny,
};
let _ = req.response_tx.send(mapped);
} else {
let _ = req.response_tx.send(ApprovalResponse::Deny);
}
saturating_dec(&pending_approvals);
}
});
let delegation_ev_tx = ev_tx.clone();
let child_runs_in_flight = run_progress.child_runs_in_flight.clone();
tokio::spawn(async move {
while let Some(req) = delegation_req_rx.recv().await {
match req {
DelegationEvent::TaskStarted {
task_index,
task_count,
goal,
depth,
agent_id,
parent_id,
} => {
child_runs_in_flight.fetch_add(1, Ordering::Relaxed);
let _ = delegation_ev_tx.send(crate::StreamEvent::SubAgentStart {
task_index,
task_count,
goal,
depth,
agent_id,
parent_id,
});
}
DelegationEvent::Thinking {
task_index,
task_count,
text,
} => {
let _ = delegation_ev_tx.send(crate::StreamEvent::SubAgentReasoning {
task_index,
task_count,
text,
});
}
DelegationEvent::ToolCalled {
task_index,
task_count,
tool_name,
args_json,
} => {
let _ = delegation_ev_tx.send(crate::StreamEvent::SubAgentToolExec {
task_index,
task_count,
name: tool_name,
args_json,
});
}
DelegationEvent::TaskFinished {
task_index,
task_count,
status,
duration_ms,
summary,
api_calls,
model,
} => {
saturating_dec(&child_runs_in_flight);
let _ = delegation_ev_tx.send(crate::StreamEvent::SubAgentFinish {
task_index,
task_count,
status,
duration_ms,
summary,
api_calls,
model,
});
}
}
}
});
}
let clarify_tx_for_dispatch = if event_tx.is_some() {
Some(clarify_req_tx)
} else {
None
};
let approval_tx_for_dispatch = if event_tx.is_some() {
Some(approval_req_tx)
} else {
None
};
let delegation_tx_for_dispatch = if event_tx.is_some() {
Some(delegation_req_tx)
} else {
None
};
let turn_started_at = std::time::Instant::now();
self.publish_session_state(&session).await;
tracing::info!(
session_id = %conversation_session_id,
messages = session.messages.len(),
has_system_prompt = session.cached_system_prompt.is_some(),
"execute_loop: entering main conversation_loop"
);
let mut final_response = String::new();
let mut interrupted = false;
let mut budget_exhausted = false;
let mut invalid_tool_budget_exhausted = false;
let mut tool_errors_acc: Vec<edgecrab_types::ToolErrorRecord> = Vec::new();
let mut turn = crate::turn_prologue::TurnPrologueState::begin(&config.harness);
let mut phase = crate::turn_phase::TurnPhaseTracker::begin();
let mut completion_reopen_gate = crate::completion_reopen::CompletionReopenGate::new(
config.harness.max_completion_reopens,
);
let capability_suppressions: Arc<Mutex<HashMap<String, ToolErrorResponse>>> =
Arc::new(Mutex::new(HashMap::new()));
let spill_seq = Arc::new(crate::tool_result_spill::SpillSequence::new());
const MAX_COMPRESSION_LLM_FAILURES: u32 = 3;
let shadow_judge_cfg = config.shadow_judge.clone();
let mut shadow_judge_invocations: u32 = 0;
let (shadow_judge_provider, shadow_judge_model) = if shadow_judge_cfg.enabled {
crate::shadow_judge::resolve_shadow_provider_and_model(
&shadow_judge_cfg,
config.auxiliary.model.as_deref(),
effective_provider.clone(),
&config.model,
)
} else {
(effective_provider.clone(), config.model.clone())
};
let turn_tool_progress_tx = make_tool_progress_tx(event_tx);
let watch_notification_tx = if let Some(ev_tx) = event_tx.cloned() {
let (tx, mut rx) =
tokio::sync::mpsc::unbounded_channel::<edgecrab_tools::process_table::WatchEvent>();
tokio::spawn(async move {
while let Some(event) = rx.recv().await {
forward_process_watch_event(event, &ev_tx);
}
});
Some(tx)
} else {
None
};
'conversation_loop: loop {
phase.transition(crate::turn_phase::TurnPhase::Preflight);
if tool_defs_dirty {
active_tool_defs = if let Some(ref registry) = tool_registry {
let schema_ctx = build_tool_context(
&cwd,
app_config_ref.clone(),
&cancel,
&self.state_db,
config.platform,
config.capability_grants,
&self.process_table,
Some(effective_provider.clone()),
tool_registry.clone(),
None,
None,
None,
None,
None,
None,
self.gateway_sender.read().await.clone(),
config.origin_chat.clone(),
None,
None,
&conversation_session_id,
Some(self.todo_store.clone()),
None,
None,
None,
delegate_ctx.clone(),
config.kanban_task_id.clone(),
Some(materialized_tools.clone()),
Some(self.skills_zone_dirty.clone()),
);
let enabled_filter = if config.enabled_toolsets.is_empty()
|| edgecrab_tools::toolsets::contains_all_sentinel(&config.enabled_toolsets)
|| expanded_enabled.is_empty()
{
None
} else {
Some(expanded_enabled.as_slice())
};
let disabled_filter = if expanded_disabled.is_empty() {
None
} else {
Some(expanded_disabled.as_slice())
};
let materialized = read_materialized_set(Some(&materialized_tools));
build_wire_llm_definitions(
registry,
&schema_ctx,
enabled_filter,
disabled_filter,
effective_schema_mode,
&materialized,
hot_schema_full,
)
} else {
Vec::new()
};
tool_defs_dirty = false;
}
if !self.budget.try_consume() {
tracing::warn!(
used = self.budget.used(),
max = self.budget.max(),
"iteration budget exhausted"
);
budget_exhausted = true;
break;
}
if cancel.is_cancelled() {
interrupted = true;
break;
}
sanitize_orphaned_tool_results(&mut session.messages);
if let Some(ref reg) = tool_registry {
edgecrab_tools::tool_call_pipeline::sanitize_assistant_tool_calls_for_api(
&mut session.messages,
reg,
);
}
strip_budget_warnings_from_history(&mut session.messages);
if app_config_ref.computer_use_keep_last_n_screenshots > 0 {
session.messages = crate::compression::prune_computer_use_screenshots(
&session.messages,
app_config_ref.computer_use_keep_last_n_screenshots,
);
}
if crate::local_provider_policy::is_local_inference_provider(effective_provider.name())
{
let _ = effective_provider.refresh_model_metadata().await;
}
let mut compression_params =
CompressionParams::from_model_config(&config.model, &config.compression);
if crate::local_provider_policy::is_local_inference_provider(effective_provider.name())
{
let live_context = effective_provider.max_context_length();
if live_context > 0 && live_context < compression_params.context_window {
compression_params.context_window = live_context;
compression_params.reapply_threshold_floor(config.compression.threshold);
}
}
let estimated_prompt_tokens = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
);
let threshold_tokens = compression_params.threshold_tokens();
turn.observe_opening_tokens(estimated_prompt_tokens);
let protect_first = crate::compression::effective_protect_first_n(
session.compression_runtime.compression_count,
session.messages.iter().any(|m| {
m.text_content()
.contains(crate::compression::SUMMARY_PREFIX)
}),
);
let protect_last = compression_params.protect_last_n;
let preflight_gate = crate::turn_prologue::should_run_preflight_estimate(
session.messages.len(),
protect_first,
protect_last,
estimated_prompt_tokens,
threshold_tokens,
);
let defer_preflight = crate::compression::should_defer_preflight_to_real_usage(
&session.compression_runtime,
estimated_prompt_tokens,
threshold_tokens,
);
if defer_preflight {
tracing::debug!(
estimated_prompt_tokens,
threshold_tokens,
"deferring rough-estimate preflight compression to real usage"
);
}
let blocked =
crate::compression::automatic_compression_blocked(&session.compression_runtime);
match check_compression_status_for_estimate(
estimated_prompt_tokens,
&compression_params,
) {
CompressionStatus::NeedsCompression if blocked => {
turn.note_preflight(false);
turn.trace_preflight(preflight_gate, "blocked");
tracing::warn!(
ineffective = session.compression_runtime.ineffective_compression_count,
fallback = session.compression_runtime.fallback_compression_streak,
"compression skipped — anti-thrashing cooldown active"
);
}
CompressionStatus::NeedsCompression if defer_preflight => {
turn.note_preflight(false);
turn.trace_preflight(preflight_gate, "deferred");
}
CompressionStatus::NeedsCompression => {
let _compression_guard = self.compression_lock.lock().await;
if !preflight_gate {
tracing::debug!(
messages = session.messages.len(),
estimated_prompt_tokens,
"preflight compress despite count gate (token threshold hit)"
);
}
turn.note_preflight(true);
turn.trace_preflight(preflight_gate, "running");
emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_compression_started(),
);
tracing::info!(
messages = session.messages.len(),
estimated_prompt_tokens,
"compressing context before API call"
);
let spill_ctx = crate::compression::PruneSpillContext {
session_id: &conversation_session_id,
cwd: &cwd,
config: &crate::tool_result_spill::SpillConfig {
enabled: app_config_ref.result_spill,
threshold: app_config_ref.result_spill_threshold,
preview_lines: app_config_ref.result_spill_preview_lines,
},
seq: &spill_seq,
};
let compression_count = session.compression_runtime.compression_count;
let mut used_fallback = false;
if turn.compression_llm_failures >= MAX_COMPRESSION_LLM_FAILURES {
emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_compression_circuit_breaker(
turn.compression_llm_failures,
),
);
tracing::warn!(
failures = turn.compression_llm_failures,
"compression circuit breaker active — using structural fallback only"
);
session.messages = crate::compression::compress_structural_only_counted(
&session.messages,
&compression_params,
Some(&spill_ctx),
compression_count,
);
used_fallback = true;
} else {
let (compressed, llm_succeeded) =
crate::compression::compress_with_llm_counted(
&session.messages,
&compression_params,
&provider,
Some(&spill_ctx),
compression_count,
None,
)
.await;
session.messages = compressed;
if llm_succeeded {
turn.compression_llm_failures = 0;
} else {
turn.compression_llm_failures += 1;
used_fallback = true;
tracing::warn!(
failures = turn.compression_llm_failures,
"LLM compression fell back to structural, tracking for circuit breaker (FP29)"
);
}
}
crate::progress_sink::emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_compression_done(
session.messages.len(),
),
);
let rough_after = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
);
session.finish_compression(
self.todo_store.as_ref(),
rough_after,
used_fallback,
);
crate::lifecycle_hooks::emit_global(
crate::lifecycle_hooks::LifecycleEvent::CompressAfter,
serde_json::json!({
"session_id": &conversation_session_id,
"message_count": session.messages.len(),
"used_fallback": used_fallback,
}),
);
let recomputed_prompt_tokens = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
);
if check_compression_status_for_estimate(
recomputed_prompt_tokens,
&compression_params,
) == CompressionStatus::Ok
{
turn.pressure_warned = false;
}
self.publish_session_state(&session).await;
}
CompressionStatus::PressureWarning if !turn.pressure_warned => {
tracing::warn!(
estimated_tokens = estimated_prompt_tokens,
threshold_tokens,
"context approaching compression threshold"
);
if let Some(tx) = event_tx {
let _ = tx.send(crate::StreamEvent::ContextPressure {
estimated_tokens: estimated_prompt_tokens,
threshold_tokens,
});
}
turn.pressure_warned = true;
}
_ => {}
}
if crate::local_provider_policy::local_tool_harness_active(
effective_provider.name(),
!active_tool_defs.is_empty(),
) {
let spill_config = local_prune_spill_config(&app_config_ref);
let spill_ctx = build_prune_spill_context(
&conversation_session_id,
&cwd,
&spill_config,
&spill_seq,
);
let prompt_for_compress = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
);
if let Some((compressed, tokens_before, tokens_after)) =
crate::local_provider_policy::try_local_midband_structural_compress(
&session.messages,
&compression_params,
effective_provider.max_context_length(),
prompt_for_compress,
Some(&spill_ctx),
)
{
crate::local_provider_policy::log_local_structural_compress(
effective_provider.as_ref(),
tokens_before,
tokens_after,
);
session.messages = compressed;
session.finish_compression(
self.todo_store.as_ref(),
tokens_after,
true, );
crate::progress_sink::emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_local_structural_compress_notice(
effective_provider.name(),
tokens_before,
tokens_after,
),
);
self.publish_session_state(&session).await;
}
}
if crate::local_provider_policy::local_tool_harness_active(
effective_provider.name(),
!active_tool_defs.is_empty(),
) {
let spill_config = local_prune_spill_config(&app_config_ref);
let spill_ctx = build_prune_spill_context(
&conversation_session_id,
&cwd,
&spill_config,
&spill_seq,
);
if let Some((pruned_messages, tools_pruned, prompt_before, prompt_after)) =
try_local_structural_prune_request(
crate::local_provider_policy::LocalStructuralPrunePhase::Preflight,
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
&spill_ctx,
effective_provider.max_context_length(),
)
{
crate::local_provider_policy::log_local_prefill_prune(
effective_provider.as_ref(),
prompt_before,
prompt_after,
tools_pruned,
"preflight",
);
session.messages = pruned_messages;
edgecrab_tools::read_tracker::reset_read_dedup(&conversation_session_id);
emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_local_prefill_prune_notice(
effective_provider.name(),
prompt_before,
prompt_after,
tools_pruned,
"preflight",
),
);
self.publish_session_state(&session).await;
}
}
let resolved_cache = crate::prompt_cache_policy::resolve_prompt_cache(
effective_provider.name(),
effective_provider.model(),
config.model_config.base_url.as_deref(),
config.model_config.prompt_caching,
&config.cache.prompt_prefix,
);
let goal_block = crate::goals::render_goal_block(
&self
.goal_store
.active(&conversation_session_id)
.unwrap_or_default(),
);
let last_user_text = session
.messages
.iter()
.rev()
.find(|m| m.role == Role::User)
.map(|m| m.text_content())
.unwrap_or_default();
let mcp_block =
edgecrab_tools::tools::mcp_client::render_mcp_api_context(&last_user_text);
let base_completion_options = completion_options_for(&config);
let max_mutation_payload_bytes = app_config_ref.max_write_payload_bytes();
let local_abs_max = app_config_ref.local_max_tool_turn_tokens;
let completion_options = crate::local_provider_policy::effective_completion_options(
&base_completion_options,
effective_provider.as_ref(),
!active_tool_defs.is_empty(),
max_mutation_payload_bytes,
local_abs_max,
);
let prompt_tokens_for_plan = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
);
let local_tool_turn_plan = if active_tool_defs.is_empty() {
None
} else {
crate::local_provider_policy::local_tool_turn_plan(
effective_provider.as_ref(),
&completion_options,
prompt_tokens_for_plan,
max_mutation_payload_bytes,
base_completion_options.reasoning_effort.as_deref(),
local_abs_max,
)
};
if let Some(ref plan) = local_tool_turn_plan {
crate::local_provider_policy::log_local_tool_turn_plan(plan);
emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_local_tool_turn_preflight(
&plan.log_line(),
),
);
}
let messages_for_api = {
let mut extended = session.messages.clone();
if !goal_block.is_empty() {
extended.push(Message::user(&goal_block));
}
if !mcp_block.is_empty() {
extended.push(Message::user(&mcp_block));
}
if extended.len() == session.messages.len() {
session.messages.clone()
} else {
extended
}
};
let mut chat_messages = build_api_chat_messages(
&session,
&messages_for_api,
resolved_cache.as_ref(),
effective_provider.as_ref(),
&app_config_ref,
);
crate::observability::log_harness_api_iteration(
&conversation_session_id,
session.api_call_count,
effective_provider.name(),
effective_provider.model(),
active_tool_defs.len(),
chat_messages.len(),
);
crate::lifecycle_hooks::emit_global(
crate::lifecycle_hooks::LifecycleEvent::TurnBefore,
serde_json::json!({
"session_id": &conversation_session_id,
"iteration": session.api_call_count,
"model": effective_provider.model(),
}),
);
let native_streaming_active = should_use_native_streaming(
effective_provider.as_ref(),
&active_tool_defs,
config.streaming,
event_tx.is_some(),
) && (!session.native_tool_streaming_disabled
|| active_tool_defs.is_empty());
let local_turn_budget = local_tool_turn_plan
.as_ref()
.map(|plan| (plan.max_tool_argument_bytes, plan.max_tokens));
let api_tool_defs = edgecrab_tools::registry::annotate_llm_definitions_for_local_turn(
active_tool_defs.clone(),
effective_provider.name(),
local_turn_budget,
);
let classify_session = crate::provider_call::classify_session_for_call(
effective_provider.as_ref(),
&chat_messages,
prompt_tokens_for_plan as u32,
);
phase.transition(crate::turn_phase::TurnPhase::ModelCall);
let api_outcome = match api_call_with_retry(
&effective_provider,
&chat_messages,
&api_tool_defs,
MAX_RETRIES,
ApiCallContext {
options: Some(&completion_options),
cancel: &cancel,
event_tx,
use_native_streaming: native_streaming_active,
discovered_plugins: discovered_plugins.as_deref(),
conversation_session_id: &conversation_session_id,
platform: config.platform,
api_call_count: session.api_call_count,
session: classify_session,
},
)
.await
{
Ok(outcome) => {
tracing::info!(
elapsed_ms = turn_started_at.elapsed().as_millis() as u64,
"execute_loop: api_call_with_retry succeeded"
);
Ok(outcome)
}
Err(AgentError::Interrupted) => {
interrupted = true;
break 'conversation_loop;
}
Err(primary_err) => {
let err_text = primary_err.to_string();
if crate::multimodal_tool_content::is_tool_content_rejection_error(&err_text) {
let key = crate::multimodal_tool_content::provider_model_key(
effective_provider.name(),
effective_provider.model(),
);
if !session.tool_result_image_downgrades.contains(&key) {
session.tool_result_image_downgrades.insert(key);
tracing::warn!(
provider = effective_provider.name(),
model = effective_provider.model(),
"provider rejected tool-result images; retrying with text-only tool messages"
);
chat_messages = build_api_chat_messages(
&session,
&messages_for_api,
resolved_cache.as_ref(),
effective_provider.as_ref(),
&app_config_ref,
);
crate::multimodal_tool_content::downgrade_tool_images_in_chat_messages(
&mut chat_messages,
Some((
effective_provider.name(),
effective_provider.model(),
&mut session.tool_result_image_downgrades,
)),
);
match api_call_with_retry(
&effective_provider,
&chat_messages,
&active_tool_defs,
MAX_RETRIES,
ApiCallContext {
options: Some(&completion_options),
cancel: &cancel,
event_tx,
use_native_streaming: native_streaming_active,
discovered_plugins: discovered_plugins.as_deref(),
conversation_session_id: &conversation_session_id,
platform: config.platform,
api_call_count: session.api_call_count,
session: classify_session,
},
)
.await
{
Ok(outcome) => Ok(outcome),
Err(AgentError::Interrupted) => {
interrupted = true;
break 'conversation_loop;
}
Err(retry_err) => Err(retry_err),
}
} else {
Err(primary_err)
}
} else {
Err(primary_err)
}
}
};
let mut api_outcome = match api_outcome {
Ok(outcome) => {
crate::lifecycle_hooks::emit_global(
crate::lifecycle_hooks::LifecycleEvent::TurnAfter,
serde_json::json!({
"session_id": &conversation_session_id,
"iteration": session.api_call_count,
"has_tool_calls": outcome.response.has_tool_calls(),
}),
);
outcome
}
Err(primary_err) => 'recover: {
if native_streaming_active {
self.publish_session_state(&session).await;
return Err(primary_err);
}
let mut primary_err = primary_err;
if smart_routed_provider_active {
tracing::warn!(
routed_model = %route.model,
primary_model = %config.model_config.default_model,
routed_error = %primary_err,
"smart-routed model failed before visible output, retrying primary model"
);
let primary_native_streaming = should_use_native_streaming(
provider.as_ref(),
&active_tool_defs,
config.streaming,
event_tx.is_some(),
);
match api_call_with_retry(
&provider,
&chat_messages,
&active_tool_defs,
MAX_RETRIES,
ApiCallContext {
options: Some(&completion_options),
cancel: &cancel,
event_tx,
use_native_streaming: primary_native_streaming,
discovered_plugins: discovered_plugins.as_deref(),
conversation_session_id: &conversation_session_id,
platform: config.platform,
api_call_count: session.api_call_count,
session: classify_session,
},
)
.await
{
Ok(outcome) => break 'recover outcome,
Err(AgentError::Interrupted) => {
interrupted = true;
break 'conversation_loop;
}
Err(primary_retry_err) => {
tracing::error!(
routed_model = %route.model,
routed_error = %primary_err,
primary_retry_error = %primary_retry_err,
"smart-routed model and primary retry both failed"
);
primary_err = primary_retry_err;
}
}
}
if let Some(ref fb) = config.model_config.fallback {
let fb_route = crate::model_router::fallback_route(fb);
tracing::warn!(
primary_error = %primary_err,
fallback = %fb_route.model,
"primary API failed, trying fallback"
);
if let Some((fb_prov_name, fb_model_name)) = fb_route.model.split_once('/')
{
let fb_canonical =
edgecrab_tools::vision_models::normalize_provider_name(
fb_prov_name,
);
let fb_prov_opt: Option<Arc<dyn LLMProvider>> =
edgecrab_tools::create_provider_for_model(
&fb_canonical,
fb_model_name,
)
.ok();
if let Some(fb_prov) = fb_prov_opt {
let fallback_native_streaming = should_use_native_streaming(
fb_prov.as_ref(),
&active_tool_defs,
config.streaming,
event_tx.is_some(),
);
let fb_completion_options =
crate::local_provider_policy::effective_completion_options(
&completion_options_for(&config),
fb_prov.as_ref(),
!active_tool_defs.is_empty(),
max_mutation_payload_bytes,
local_abs_max,
);
let fb_tool_defs =
edgecrab_tools::registry::annotate_llm_definitions_for_local_turn(
active_tool_defs.clone(),
fb_prov.name(),
local_turn_budget,
);
match api_call_with_retry(
&fb_prov,
&chat_messages,
&fb_tool_defs,
1,
ApiCallContext {
options: Some(&fb_completion_options),
cancel: &cancel,
event_tx,
use_native_streaming: fallback_native_streaming,
discovered_plugins: discovered_plugins.as_deref(),
conversation_session_id: &conversation_session_id,
platform: config.platform,
api_call_count: session.api_call_count,
session: crate::provider_call::classify_session_for_call(
fb_prov.as_ref(),
&chat_messages,
prompt_tokens_for_plan as u32,
),
},
)
.await
{
Ok(outcome) => outcome,
Err(AgentError::Interrupted) => {
interrupted = true;
break 'conversation_loop;
}
Err(fb_err) => {
tracing::error!(fallback_error = %fb_err, "fallback also failed");
self.publish_session_state(&session).await;
return Err(primary_err);
}
}
} else {
self.publish_session_state(&session).await;
return Err(primary_err);
}
} else {
self.publish_session_state(&session).await;
return Err(primary_err);
}
} else {
self.publish_session_state(&session).await;
return Err(primary_err);
}
}
};
if let Some(replacement) = api_outcome.provider_replacement.take() {
let replaces_primary = Arc::ptr_eq(&effective_provider, &provider);
effective_provider = replacement.clone();
if replaces_primary {
*self.provider.write().await = replacement;
}
tracing::info!(
provider = effective_provider.name(),
model = effective_provider.model(),
"installed credential-rotated provider at turn boundary"
);
}
if api_outcome.disabled_native_tool_streaming {
session.native_tool_streaming_disabled = true;
}
let response = api_outcome.response;
phase.transition(crate::turn_phase::TurnPhase::Response);
if cancel.is_cancelled() {
interrupted = true;
break;
}
session.api_call_count += 1;
session.session_input_tokens += response.prompt_tokens as u64;
session.session_output_tokens += response.completion_tokens as u64;
if let Some(cache_tokens) = response.cache_hit_tokens {
session.session_cache_read_tokens += cache_tokens as u64;
}
if let Some(cache_write) = response.cache_write_tokens {
session.session_cache_write_tokens += cache_write as u64;
}
if let Some(reasoning_tokens) = response.thinking_tokens {
session.session_reasoning_tokens += reasoning_tokens as u64;
}
session.last_prompt_tokens =
response.prompt_tokens as u64 + response.cache_hit_tokens.unwrap_or(0) as u64;
if session.last_prompt_tokens == 0 {
let estimated = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
) as u64;
if estimated > 0 {
session.last_prompt_tokens = estimated;
session.session_input_tokens += estimated;
}
}
{
let rough_now = estimate_request_prompt_tokens(
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
);
let threshold_now =
CompressionParams::from_model_config(&config.model, &config.compression)
.threshold_tokens();
crate::compression::update_compression_from_response(
&mut session.compression_runtime,
session.last_prompt_tokens,
threshold_now,
rough_now,
);
}
if !active_tool_defs.is_empty()
&& response.finish_reason.as_deref() == Some("length")
&& !response.has_tool_calls()
{
let max_tokens = completion_options.max_tokens.unwrap_or(0);
let thinking_tokens = response.thinking_tokens.unwrap_or(0);
let length_metrics = crate::local_provider_policy::LocalLlmResponseMetrics {
elapsed_ms: 0,
finish_reason: Some("length".to_string()),
prompt_tokens: response.prompt_tokens,
completion_tokens: response.completion_tokens,
thinking_tokens: Some(thinking_tokens),
tool_call_count: 0,
content_len: response.content.len(),
has_reasoning_content: response.thinking_content.is_some(),
max_tokens: Some(max_tokens),
tool_choice_required: true,
};
crate::local_provider_policy::log_local_tool_length_failure(
provider.as_ref(),
&length_metrics,
);
if crate::local_provider_policy::is_local_inference_provider(provider.name()) {
emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_local_length_without_tools_notice(
provider.name(),
response.completion_tokens,
thinking_tokens,
max_tokens,
),
);
}
if crate::local_provider_policy::is_local_inference_provider(provider.name()) {
let spill_config = local_prune_spill_config(&app_config_ref);
let spill_ctx = build_prune_spill_context(
&conversation_session_id,
&cwd,
&spill_config,
&spill_seq,
);
if let Some((pruned_messages, tools_pruned, prompt_before, prompt_after)) =
try_local_structural_prune_request(
crate::local_provider_policy::LocalStructuralPrunePhase::LengthRecovery,
session.cached_system_prompt.as_deref(),
&session.messages,
&active_tool_defs,
&spill_ctx,
provider.max_context_length(),
)
{
crate::local_provider_policy::log_local_prefill_prune(
provider.as_ref(),
prompt_before,
prompt_after,
tools_pruned,
"length_recovery",
);
session.messages = pruned_messages;
edgecrab_tools::read_tracker::reset_read_dedup(&conversation_session_id);
emit_activity(
event_tx,
edgecrab_tools::tool_progress_tail::format_local_prefill_prune_notice(
provider.name(),
prompt_before,
prompt_after,
tools_pruned,
"length_recovery",
),
);
}
}
let recovery =
edgecrab_tools::mutation_turn_policy::continuation_user_message(
edgecrab_tools::mutation_turn_policy::ContinuationFailureClass::LengthWithoutTools,
&[],
app_config_ref.max_write_payload_bytes(),
Some(provider.as_ref()),
);
session.messages.push(Message::user(&recovery));
self.publish_session_state(&session).await;
turn.trackers.failure.record_success();
continue;
}
if response.finish_reason.as_deref() == Some(FINISH_REASON_STREAM_INTERRUPTED)
&& !response.has_tool_calls()
{
tracing::warn!(
"tool-call draft interrupted before delivery — injecting incremental-edit recovery"
);
emit_activity(
event_tx,
"↻ Tool draft interrupted — use scaffold + patch steps (see recovery message)",
);
let recovery =
edgecrab_tools::mutation_turn_policy::continuation_user_message(
edgecrab_tools::mutation_turn_policy::ContinuationFailureClass::StreamInterruptedNoTools,
&[],
app_config_ref.max_write_payload_bytes(),
Some(provider.as_ref()),
);
session.messages.push(Message::user(&recovery));
continue;
}
if !response_has_visible_output(&response)
&& response.finish_reason.as_deref() != Some("length")
{
tracing::info!("empty response from LLM, nudging to continue");
session.messages.push(Message::user(
"[system: your response was empty — please provide a response]",
));
continue;
}
let dctx = DispatchContext {
cwd: cwd.clone(),
registry: tool_registry.clone(),
cancel: cancel.clone(),
state_db: self.state_db.clone(),
platform: config.platform,
capability_grants: config.capability_grants,
process_table: self.process_table.clone(),
provider: Some(provider.clone()),
gateway_sender: self.gateway_sender.read().await.clone(),
sub_agent_runner: sub_agent_runner.clone(),
event_tx: event_tx.cloned(),
delegation_event_tx: delegation_tx_for_dispatch.clone(),
clarify_tx: clarify_tx_for_dispatch.clone(),
approval_tx: approval_tx_for_dispatch.clone(),
origin_chat: config.origin_chat.clone(),
app_config_ref: app_config_ref.clone(),
conversation_session_id: conversation_session_id.clone(),
todo_store: Some(self.todo_store.clone()),
capability_suppressions: capability_suppressions.clone(),
discovered_plugins: discovered_plugins.clone(),
spill_seq: spill_seq.clone(),
context_engine: engine_for_dispatch.clone(),
engine_tool_names: engine_tool_names.clone(),
mutation_turn: Arc::clone(&mutation_turn),
lsp_gate: lsp_gate.clone(),
tool_progress_tx: turn_tool_progress_tx.clone(),
watch_notification_tx: watch_notification_tx.clone(),
delegate_ctx: delegate_ctx.clone(),
kanban_task_id: config.kanban_task_id.clone(),
materialized_tools: Some(materialized_tools.clone()),
subdirectory_hints: subdirectory_hints.clone(),
skills_zone_dirty: Some(self.skills_zone_dirty.clone()),
};
if response.has_tool_calls() {
phase.transition(crate::turn_phase::TurnPhase::ToolDispatch);
}
let action = match process_response(
&response,
&mut session,
&dctx,
&mut tool_errors_acc,
&mut turn.trackers,
)
.await
{
Ok(action) => action,
Err(err) => {
self.publish_session_state(&session).await;
return Err(err);
}
};
self.publish_session_state(&session).await;
match action {
LoopAction::Done(text) => {
if response.finish_reason.as_deref() == Some(FINISH_REASON_STREAM_INTERRUPTED) {
tracing::warn!(
partial_len = text.len(),
"streamed tool call was interrupted after visible output; continuing via a safe non-streaming recovery turn"
);
emit_activity(
event_tx,
"↻ Retrying tool call via non-streaming completion — split large payloads into patch steps.",
);
let recovery =
edgecrab_tools::mutation_turn_policy::continuation_user_message(
edgecrab_tools::mutation_turn_policy::ContinuationFailureClass::StreamInterruptedAfterPartial,
&[],
app_config_ref.max_write_payload_bytes(),
Some(provider.as_ref()),
);
session.messages.push(Message::user(&recovery));
self.publish_session_state(&session).await;
continue;
}
if response.finish_reason.as_deref() == Some("length") {
if !active_tool_defs.is_empty() && !response.has_tool_calls() {
tracing::warn!(
target: "edgecrab::local_llm",
"finish_reason=length without tool_calls on tool turn — recovery already handled"
);
self.publish_session_state(&session).await;
continue;
}
tracing::info!(
partial_len = text.len(),
"response truncated (finish_reason=length), auto-continuing"
);
session.messages.push(Message::user(
"[system: your response was truncated due to length — please continue exactly where you left off]",
));
continue;
}
let todo = snapshot_todo_state(&self.todo_store);
let task_class = classify_from_messages(&session.messages);
let provisional_harness = crate::turn_epilogue::build_turn_harness_snapshot(
crate::turn_epilogue::TurnHarnessBuildParams {
messages: &session.messages,
mutation_turn: &mutation_turn,
cwd: &cwd,
post_mutation_oracles: false,
harness_advisory: &turn.trackers.harness_advisory,
guardrail_halt: turn.trackers.guardrail_halt,
task_class,
},
);
let goal_state = self
.goal_store
.active(&conversation_session_id)
.unwrap_or_default();
let goal_contract =
(!goal_state.contract.is_empty()).then_some(&goal_state.contract);
let provisional_outcome = crate::turn_epilogue::assess_turn_outcome(
crate::turn_epilogue::TurnAssessParams {
final_response: &text,
messages: &session.messages,
interrupted: false,
budget_exhausted: false,
invalid_tool_budget_exhausted: false,
pending_approval: run_progress
.pending_approvals
.load(Ordering::Relaxed)
> 0,
pending_clarification: run_progress
.pending_clarifications
.load(Ordering::Relaxed)
> 0,
active_todos: todo.active,
blocked_todos: todo.blocked,
child_runs_in_flight: run_progress
.child_runs_in_flight
.load(Ordering::Relaxed),
harness: provisional_harness,
harness_config: &config.harness,
goal_contract,
harness_contract_verify: false,
cwd: &cwd,
evidence: turn.trackers.evidence.assess_snapshot(),
},
);
match crate::completion_reopen::decide_completion_reopen(
&provisional_outcome,
&session.messages,
turn.trackers.evidence.assess_snapshot(),
&completion_reopen_gate,
) {
crate::completion_reopen::ReopenDecision::Reopen => {
completion_reopen_gate.reopens_used =
completion_reopen_gate.reopens_used.saturating_add(1);
tracing::info!(
state = provisional_outcome.state.as_str(),
reopen = completion_reopen_gate.reopens_used,
max_reopens = completion_reopen_gate.max_reopens,
active_tasks = todo.active,
blocked_tasks = todo.blocked,
pending_approvals =
run_progress.pending_approvals.load(Ordering::Relaxed),
pending_clarifications =
run_progress.pending_clarifications.load(Ordering::Relaxed),
child_runs =
run_progress.child_runs_in_flight.load(Ordering::Relaxed),
"model returned final text before the harness considered the task complete; continuing the loop"
);
session.messages.push(Message::user(
&crate::turn_epilogue::completion_follow_up_message(
&provisional_outcome,
),
));
self.publish_session_state(&session).await;
continue;
}
crate::completion_reopen::ReopenDecision::CapReached => {
tracing::warn!(
state = provisional_outcome.state.as_str(),
max_reopens = completion_reopen_gate.max_reopens,
"completion reopen cap reached — ending turn (025)"
);
final_response = text;
turn.trackers.evidence.escalate("reopen_cap");
break;
}
crate::completion_reopen::ReopenDecision::DoNotReopen => {}
}
if shadow_judge_cfg.enabled
&& shadow_judge_invocations < shadow_judge_cfg.max_per_session
&& session.messages.len() >= shadow_judge_cfg.min_messages_before_enable
{
shadow_judge_invocations += 1;
tracing::debug!(
invocation = shadow_judge_invocations,
max = shadow_judge_cfg.max_per_session,
messages = session.messages.len(),
model = %shadow_judge_model,
"shadow judge: invoking completion oracle"
);
let verdict = crate::shadow_judge::run_shadow_judge(
&shadow_judge_provider,
&shadow_judge_model,
&session.messages,
&shadow_judge_cfg,
)
.await;
if let Some(verdict) = verdict {
session.session_input_tokens += u64::from(verdict.input_tokens);
session.session_output_tokens += u64::from(verdict.output_tokens);
tracing::info!(
is_complete = verdict.is_complete,
confidence = verdict.confidence,
reason = %verdict.reason,
invocation = shadow_judge_invocations,
"shadow judge: verdict"
);
let confidence_above_threshold =
verdict.confidence >= shadow_judge_cfg.confidence_threshold;
if !verdict.is_complete && confidence_above_threshold {
let hint = verdict.steering_hint.as_deref().unwrap_or(
"Continue working until all parts of the request are complete.",
);
let msg = build_shadow_judge_message(hint, &verdict.reason);
tracing::info!(
hint = %hint,
"shadow judge: vetoing completion, injecting continuation nudge"
);
if let Some(tx) = event_tx {
let _ = tx.send(crate::StreamEvent::ToolExec {
name: "shadow_judge".to_string(),
args_json: serde_json::json!({
"verdict": "incomplete",
"confidence": verdict.confidence,
"reason": verdict.reason,
})
.to_string(),
tool_call_id: "sj".to_string(),
});
}
session.messages.push(Message::user(&msg));
self.publish_session_state(&session).await;
continue 'conversation_loop;
}
}
}
final_response = text;
break;
}
LoopAction::GuardrailHalt => {
tracing::warn!("tool-loop guardrail halt — ending turn");
final_response = String::new();
break;
}
LoopAction::PartialAbort { reason } => {
tracing::warn!(%reason, "invalid tool retry budget exhausted — partial abort");
emit_activity(
event_tx,
format!(
"⚠️ {reason} (max {} invalid-tool retries)",
edgecrab_tools::MAX_INVALID_TOOL_RETRIES
),
);
session.invalid_tool_call_retries = 0;
invalid_tool_budget_exhausted = true;
final_response = reason;
break;
}
LoopAction::Continue => {
if self.take_skills_zone_dirty() {
session.invalidate_skills_zone();
}
if crate::task_class::document_done_latch_ready(&session.messages) {
tracing::info!(
"document done latch: artifact evidence present — ending turn"
);
final_response = session
.messages
.iter()
.rev()
.find(|m| m.role == Role::Assistant)
.map(|m| m.text_content())
.filter(|t| !t.trim().is_empty())
.unwrap_or_else(|| {
"Document artifact delivered and verified on disk.".into()
});
emit_activity(
event_tx,
"✓ Document artifact ready — ending turn (Done latch).",
);
self.publish_session_state(&session).await;
break;
}
if turn.trackers.evidence.visual_evidence_complete()
|| turn.trackers.evidence.media_evidence_complete()
{
tracing::info!("visual/media evidence latch complete — ending turn");
final_response = session
.messages
.iter()
.rev()
.find(|m| m.role == Role::Assistant)
.map(|m| m.text_content())
.filter(|t| !t.trim().is_empty())
.unwrap_or_else(|| "Visual evidence latched — task verified.".into());
emit_activity(event_tx, "✓ Evidence latch complete — ending turn (Done).");
self.publish_session_state(&session).await;
break;
}
if turn.trackers.evidence.should_hard_stop() {
tracing::warn!(
phase = ?turn.trackers.evidence.phase,
"evidence hard stop — ending turn (GuardrailHalt)"
);
turn.trackers.guardrail_halt = true;
final_response = session
.messages
.iter()
.rev()
.find(|m| m.role == Role::Assistant)
.map(|m| m.text_content())
.filter(|t| !t.trim().is_empty())
.unwrap_or_else(|| turn.trackers.evidence.allowed_action_message());
emit_activity(event_tx, "⏹ Verification escalated — ending turn.");
self.publish_session_state(&session).await;
break;
}
if let Some(ref mut rx) = steer_rx
&& let Some((steer_msg, steer_kind)) =
crate::steering::drain_pending_steers(rx)
{
self.steer_pending
.store(0, std::sync::atomic::Ordering::Relaxed);
tracing::info!(
kind = %steer_kind,
len = steer_msg.len(),
"steering message injected at tool boundary"
);
session.messages.push(Message::user(&steer_msg));
if let Some(tx) = event_tx {
let _ =
tx.send(crate::StreamEvent::SteerApplied { message: steer_msg });
}
if matches!(steer_kind, crate::steering::SteeringKind::Stop) {
tracing::info!("steering STOP: signalling cancel token");
cancel.cancel();
}
}
if config.harness.inject_budget_pressure
&& let Some(warning) =
get_budget_warning(session.api_call_count, config.max_iterations)
{
inject_budget_warning(&mut session.messages, &warning);
}
tool_defs_dirty = true;
self.publish_session_state(&session).await;
continue;
}
}
}
if budget_exhausted && final_response.is_empty() {
let chat_messages = build_chat_messages(
session.cached_system_prompt.as_deref(),
&session.messages,
None,
false,
);
let msg = crate::turn_epilogue::synthesize_budget_exhausted_message(
&effective_provider,
chat_messages,
self.budget.used(),
self.budget.max(),
)
.await;
tracing::warn!(
max = self.budget.max(),
"emitting budget-exhausted summary response"
);
session.messages.push(Message::assistant(&msg));
if let Some(tx) = event_tx {
let _ = tx.send(crate::StreamEvent::Token(msg.clone()));
}
final_response = msg;
}
phase.transition(crate::turn_phase::TurnPhase::Epilogue);
if !interrupted
&& !final_response.is_empty()
&& crate::config::file_mutation_verifier_enabled(config.file_mutation_verifier)
{
let footer = mutation_turn.render_turn_footer();
if !footer.is_empty() {
if let Some(tx) = event_tx {
let _ = tx.send(crate::StreamEvent::Footer(footer.clone()));
}
session.messages.push(Message::user(&format!(
"[file-mutation-verifier]\n{footer}"
)));
final_response = format!("{}\n\n{}", final_response.trim_end(), footer);
self.publish_session_state(&session).await;
}
}
let task_class = classify_from_messages(&session.messages);
let harness_snapshot = crate::turn_epilogue::build_turn_harness_snapshot(
crate::turn_epilogue::TurnHarnessBuildParams {
messages: &session.messages,
mutation_turn: &mutation_turn,
cwd: &cwd,
post_mutation_oracles: crate::config::harness_post_mutation_oracles_enabled(
config.harness_post_mutation_oracles,
),
harness_advisory: &turn.trackers.harness_advisory,
guardrail_halt: turn.trackers.guardrail_halt,
task_class,
},
);
if let Some(gate_footer) = harness_snapshot.render_gate_footer() {
if let Some(tx) = event_tx {
let _ = tx.send(crate::StreamEvent::Footer(gate_footer.clone()));
}
session
.messages
.push(Message::user(&format!("[harness-gates]\n{gate_footer}")));
final_response = format!("{}\n\n{}", final_response.trim_end(), gate_footer);
self.publish_session_state(&session).await;
}
let turn_tool_calls = session
.session_tool_call_count
.saturating_sub(initial_turn_tool_call_count);
if !interrupted
&& !config.skip_memory
&& tool_registry.is_some()
&& turn_tool_calls >= SKILL_REFLECTION_THRESHOLD
{
let bg_ctx = BackgroundReflectionCtx {
messages: session.messages.clone(),
system_prompt: session.cached_system_prompt.clone(),
tool_defs: active_tool_defs.clone(),
cwd: cwd.clone(),
registry: tool_registry.as_ref().map(Arc::clone),
cancel: cancel.clone(),
state_db: self.state_db.clone(),
platform: config.platform,
capability_grants: config.capability_grants,
process_table: Arc::clone(&self.process_table),
provider: Arc::clone(&effective_provider),
gateway_sender: self.gateway_sender.read().await.clone(),
sub_agent_runner: sub_agent_runner.clone(),
app_config_ref: app_config_ref.clone(),
conversation_session_id: conversation_session_id.clone(),
origin_chat: config.origin_chat.clone(),
todo_store: Some(self.todo_store.clone()),
};
tokio::spawn(run_learning_reflection_bg(bg_ctx));
}
self.publish_session_state(&session).await;
let todo = snapshot_todo_state(&self.todo_store);
let final_goal_state = self
.goal_store
.active(&conversation_session_id)
.unwrap_or_default();
let final_goal_contract =
(!final_goal_state.contract.is_empty()).then_some(&final_goal_state.contract);
let mut run_outcome =
crate::turn_epilogue::assess_turn_outcome(crate::turn_epilogue::TurnAssessParams {
final_response: &final_response,
messages: &session.messages,
interrupted,
budget_exhausted,
invalid_tool_budget_exhausted,
pending_approval: run_progress.pending_approvals.load(Ordering::Relaxed) > 0,
pending_clarification: run_progress.pending_clarifications.load(Ordering::Relaxed)
> 0,
active_todos: todo.active,
blocked_todos: todo.blocked,
child_runs_in_flight: run_progress.child_runs_in_flight.load(Ordering::Relaxed),
harness: harness_snapshot.clone(),
harness_config: &config.harness,
goal_contract: final_goal_contract,
harness_contract_verify: true,
cwd: &cwd,
evidence: turn.trackers.evidence.assess_snapshot(),
});
run_outcome = crate::turn_epilogue::enrich_turn_outcome(
run_outcome,
&session.messages,
&harness_snapshot,
effective_provider.name() == "vscode-copilot" && session.native_tool_streaming_disabled,
);
session.last_run_outcome = Some(run_outcome.clone());
crate::observability::log_harness_completion(
&conversation_session_id,
&config.platform.to_string(),
run_outcome.state.as_str(),
run_outcome.exit_reason.as_str(),
harness_snapshot.blocks_completion(),
harness_snapshot.oracle_failures.len(),
harness_snapshot.unresolved_mutation_failures.len(),
);
self.publish_session_state(&session).await;
crate::progress_sink::emit_run_finished(event_tx, run_outcome.clone());
let session_id = session
.session_id
.clone()
.or_else(|| config.session_id.clone())
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
session.session_id = Some(session_id.clone());
self.publish_session_state(&session).await;
let usage = Usage {
input_tokens: session.session_input_tokens,
output_tokens: session.session_output_tokens,
cache_read_tokens: session.session_cache_read_tokens,
cache_write_tokens: session.session_cache_write_tokens,
reasoning_tokens: session.session_reasoning_tokens,
..Default::default()
};
let canonical_usage = CanonicalUsage {
input_tokens: session.session_input_tokens,
output_tokens: session.session_output_tokens,
cache_read_tokens: session.session_cache_read_tokens,
cache_write_tokens: session.session_cache_write_tokens,
reasoning_tokens: session.session_reasoning_tokens,
};
let cost_result = estimate_cost(&canonical_usage, &config.model);
let cost = Cost {
input_cost: canonical_usage.input_tokens as f64 * cost_result.amount_usd.unwrap_or(0.0)
/ canonical_usage.total_tokens().max(1) as f64,
output_cost: canonical_usage.output_tokens as f64
* cost_result.amount_usd.unwrap_or(0.0)
/ canonical_usage.total_tokens().max(1) as f64,
total_cost: cost_result.amount_usd.unwrap_or(0.0),
..Default::default()
};
if let Some(ref db) = self.state_db {
let title = session
.messages
.iter()
.find(|m| m.role == Role::User)
.map(|m| {
let t = m.text_content();
if t.len() > 80 {
format!("{}…", crate::safe_truncate(&t, 80))
} else {
t
}
})
.unwrap_or_else(|| "Untitled session".to_string());
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs_f64();
let existing = db.get_session(&session_id).ok().flatten();
let (source, routing_key) = match &config.origin_chat {
Some(origin) => (origin.platform.clone(), Some(origin.chat_id.clone())),
None => (
existing
.as_ref()
.map(|s| s.source.clone())
.unwrap_or_else(|| "cli".to_string()),
existing.as_ref().and_then(|s| s.user_id.clone()),
),
};
let started_at = existing.as_ref().map(|s| s.started_at).unwrap_or(now);
let title = existing
.as_ref()
.and_then(|s| s.title.clone())
.filter(|t| !t.trim().is_empty())
.unwrap_or(title);
let record = edgecrab_state::SessionRecord {
id: session_id.clone(),
source,
user_id: routing_key,
model: Some(config.model.clone()),
system_prompt: session.cached_system_prompt.clone(),
stable_system_prompt: session.cached_stable_prompt.clone(),
semi_stable_system_prompt: session.cached_semi_stable_prompt.clone(),
parent_session_id: existing.as_ref().and_then(|s| s.parent_session_id.clone()),
started_at,
ended_at: Some(now),
end_reason: Some(run_outcome.exit_reason.as_str().to_string()),
message_count: session.messages.len() as i64,
tool_call_count: derive_tool_call_count(&session) as i64,
input_tokens: session.session_input_tokens as i64,
output_tokens: session.session_output_tokens as i64,
cache_read_tokens: session.session_cache_read_tokens as i64,
cache_write_tokens: session.session_cache_write_tokens as i64,
reasoning_tokens: session.session_reasoning_tokens as i64,
last_prompt_tokens: session.last_prompt_tokens as i64,
estimated_cost_usd: cost_result.amount_usd,
title: Some(title),
};
if let Err(e) = db.save_session_with_messages(&record, &session.messages, now) {
tracing::warn!(error = %e, "failed to atomically save session to state DB");
}
if session.user_turn_count == 1 && !final_response.is_empty() {
let user_snippet = user_message.chars().take(500).collect::<String>();
let asst_snippet = final_response.chars().take(500).collect::<String>();
let db_clone = db.clone();
let sid_clone = session_id.clone();
let prov_clone = effective_provider.clone();
tokio::spawn(async move {
auto_title_session(db_clone, sid_clone, user_snippet, asst_snippet, prov_clone)
.await;
});
}
}
let completed = run_outcome.is_success();
if let Some(discovery) = discovered_plugins.as_ref() {
for plugin in discovery
.plugins
.iter()
.filter(|plugin| hermes_supports_hook(plugin, "on_session_end"))
{
if let Err(error) = invoke_hermes_hook(
plugin,
"on_session_end",
serde_json::json!({
"session_id": &session_id,
"completed": completed,
"interrupted": interrupted,
"model": &config.model,
"platform": config.platform.to_string(),
"completion_state": run_outcome.state.as_str(),
"exit_reason": run_outcome.exit_reason.as_str(),
"active_tasks": run_outcome.active_tasks,
"blocked_tasks": run_outcome.blocked_tasks,
}),
)
.await
{
tracing::warn!(plugin = %plugin.name, ?error, "Hermes on_session_end hook failed");
}
}
}
if config.save_trajectories {
let trajectory_dir = edgecrab_home().join("trajectories");
let trajectory_path = trajectory_dir.join(if completed {
"trajectory_samples.jsonl"
} else {
"failed_trajectories.jsonl"
});
if let Err(e) = std::fs::create_dir_all(&trajectory_dir) {
tracing::warn!(error = %e, path = %trajectory_dir.display(), "failed to create trajectory directory");
} else {
let trajectory = build_trajectory(
&session_id,
&config.model,
&session.messages,
session.api_call_count,
cost.total_cost,
completed,
turn_started_at.elapsed().as_secs_f64(),
);
if let Err(e) = save_trajectory(&trajectory_path, &trajectory) {
tracing::warn!(error = %e, path = %trajectory_path.display(), "failed to save trajectory");
}
}
}
let messages = session.messages.clone();
let api_calls = session.api_call_count;
let model = config.model.clone();
if let Some(rx) = steer_rx {
let mut guard = self.steer_rx.lock().expect("steer_rx mutex not poisoned");
*guard = Some(rx);
}
*self
.steer_event_tx
.lock()
.expect("steer_event_tx mutex not poisoned") = None;
phase.transition(crate::turn_phase::TurnPhase::Complete);
Ok(ConversationResult {
final_response,
messages,
session_id,
api_calls,
interrupted,
budget_exhausted,
run_outcome,
model,
usage,
cost,
tool_errors: tool_errors_acc,
})
}
}
fn build_api_chat_messages(
session: &SessionState,
messages_for_api: &[Message],
cache: Option<&crate::prompt_cache_policy::ResolvedPromptCache>,
provider: &dyn LLMProvider,
app_cfg: &edgecrab_tools::config_ref::AppConfigRef,
) -> Vec<edgequake_llm::ChatMessage> {
let messages_for_api =
crate::compression::ensure_api_safe_tool_pairs(messages_for_api.to_vec());
let messages_for_api = messages_for_api.as_slice();
let attach = crate::multimodal_tool_content::should_attach_computer_use_screenshot(
provider.name(),
provider.model(),
app_cfg,
&session.tool_result_image_downgrades,
);
let native_inner_layout = cache.is_some_and(|resolved| resolved.decision.native_inner_layout);
let mut out = match (
session.cached_stable_prompt.as_deref(),
cache,
native_inner_layout,
) {
(Some(stable), Some(resolved), true) => {
let combined = session.cached_system_prompt.as_deref().unwrap_or("");
let semi = session.cached_semi_stable_prompt.as_deref().unwrap_or("");
let dynamic = split_dynamic_after_cache_prefixes(combined, stable, semi);
build_chat_messages_blocks(
stable,
semi,
dynamic,
messages_for_api,
Some(resolved),
attach,
)
}
_ => build_chat_messages(
session.cached_system_prompt.as_deref(),
messages_for_api,
cache,
attach,
),
};
if crate::local_provider_policy::should_normalize_api_messages_for_kv(provider.name()) {
crate::local_provider_policy::normalize_api_messages_for_kv(&mut out);
}
out
}
fn append_conversation_messages(
out: &mut Vec<edgequake_llm::ChatMessage>,
messages: &[Message],
attach_computer_use_images: bool,
) {
let last_disk_image_idx =
crate::multimodal_tool_content::last_computer_use_disk_capture_index(messages);
for (idx, m) in messages.iter().enumerate() {
let text = m.text_content();
match m.role {
Role::System => out.push(edgequake_llm::ChatMessage::system(&text)),
Role::User => out.push(edgequake_llm::ChatMessage::user(&text)),
Role::Assistant => {
if let Some(ref tool_calls) = m.tool_calls
&& !tool_calls.is_empty()
{
let llm_calls: Vec<edgequake_llm::ToolCall> =
tool_calls.iter().map(|tc| tc.to_llm()).collect();
out.push(edgequake_llm::ChatMessage::assistant_with_tools(
&text, llm_calls,
));
continue;
}
out.push(edgequake_llm::ChatMessage::assistant(&text));
}
Role::Tool => {
let tool_call_id = m.tool_call_id.as_deref().unwrap_or("unknown");
let mut chat_msg = edgequake_llm::ChatMessage::tool_result(tool_call_id, &text);
chat_msg.name = m.name.clone();
crate::multimodal_tool_content::enrich_tool_chat_message(
&mut chat_msg,
m,
attach_computer_use_images,
last_disk_image_idx == Some(idx),
);
out.push(chat_msg);
}
}
}
}
pub fn split_dynamic_after_cache_prefixes<'a>(
combined: &'a str,
stable: &'a str,
semi_stable: &'a str,
) -> &'a str {
let mut rest = combined;
if !stable.is_empty() {
if !rest.starts_with(stable) {
return combined;
}
rest = rest[stable.len()..].trim_start_matches('\n');
}
if !semi_stable.is_empty() {
if !rest.starts_with(semi_stable) {
return combined;
}
rest = rest[semi_stable.len()..].trim_start_matches('\n');
}
rest
}
fn stable_cache_control(cache_config: Option<&CachePromptConfig>) -> Option<CacheControl> {
let ttl = cache_config.and_then(|c| c.cache_ttl.as_deref());
Some(match ttl {
Some("1h") => CacheControl::ephemeral_ttl("1h"),
Some("5m") => CacheControl::ephemeral_ttl("5m"),
_ => CacheControl::ephemeral(),
})
}
pub fn build_chat_messages_blocks(
stable: &str,
semi_stable: &str,
dynamic: &str,
messages: &[Message],
cache: Option<&crate::prompt_cache_policy::ResolvedPromptCache>,
attach_computer_use_images: bool,
) -> Vec<edgequake_llm::ChatMessage> {
let mut out = Vec::with_capacity(messages.len() + 3);
let cache_config = cache.map(|r| &r.config);
if !stable.is_empty() {
let mut sys = edgequake_llm::ChatMessage::system(stable);
sys.cache_control = stable_cache_control(cache_config);
out.push(sys);
}
if !semi_stable.is_empty() {
let mut sys = edgequake_llm::ChatMessage::system(semi_stable);
sys.cache_control = Some(CacheControl::ephemeral_ttl("5m"));
out.push(sys);
}
if !dynamic.is_empty() {
out.push(edgequake_llm::ChatMessage::system(dynamic));
}
append_conversation_messages(&mut out, messages, attach_computer_use_images);
if let Some(resolved) = cache {
crate::prompt_cache_policy::apply_prompt_cache_breakpoints(
&mut out,
resolved.decision,
&resolved.config,
true, );
}
out
}
pub fn build_chat_messages(
system_prompt: Option<&str>,
messages: &[Message],
cache: Option<&crate::prompt_cache_policy::ResolvedPromptCache>,
attach_computer_use_images: bool,
) -> Vec<edgequake_llm::ChatMessage> {
let mut out = Vec::with_capacity(messages.len() + 1);
if let Some(sys) = system_prompt {
out.push(edgequake_llm::ChatMessage::system(sys));
}
append_conversation_messages(&mut out, messages, attach_computer_use_images);
if let Some(resolved) = cache {
crate::prompt_cache_policy::apply_prompt_cache_breakpoints(
&mut out,
resolved.decision,
&resolved.config,
false,
);
}
out
}
#[inline]
fn parse_tool_error_response(result: &str) -> Option<ToolErrorResponse> {
edgecrab_types::parse_tool_error_payload(result)
}
fn tool_attempt_fingerprint(name: &str, args_json: &str) -> String {
let normalized_args = serde_json::from_str::<serde_json::Value>(args_json)
.ok()
.and_then(|value| serde_json::to_string(&value).ok())
.unwrap_or_else(|| args_json.trim().to_string());
format!("{name}:{normalized_args}")
}
fn invalid_args_missing_fields_suppression_key(
name: &str,
args_json: &str,
required_fields: &[String],
) -> Option<String> {
let args = serde_json::from_str::<serde_json::Value>(args_json).ok()?;
let obj = args.as_object()?;
let mut missing: Vec<String> = required_fields
.iter()
.filter(|field| !obj.contains_key(field.as_str()))
.cloned()
.collect();
if missing.is_empty() {
return None;
}
missing.sort();
Some(format!("invalid_args:{name}:missing:{}", missing.join(",")))
}
fn invalid_args_semantic_key(
registry: &ToolRegistry,
name: &str,
args_json: &str,
) -> Option<String> {
let required = registry.required_fields_for_tool(name)?;
invalid_args_missing_fields_suppression_key(name, args_json, &required)
}
fn format_preflight_tool_error(
reg: &ToolRegistry,
name: &str,
args_json: &str,
err: &ToolError,
) -> String {
if matches!(err.core_error(), ToolError::InvalidArgs { .. })
&& let Some(mut enriched) = reg.enrich_invalid_args_error(name, err, Some(args_json))
{
if enriched.suppression_key.is_none() {
enriched.suppression_key = enriched
.required_fields
.as_ref()
.and_then(|fields| {
invalid_args_missing_fields_suppression_key(name, args_json, fields)
})
.or_else(|| invalid_args_semantic_key(reg, name, args_json));
}
return serde_json::to_string(&enriched).expect("preflight enriched error serializes");
}
err.to_llm_response()
}
#[inline]
fn is_suppressed_argument_retry(payload: &ToolErrorResponse) -> bool {
payload.category == "arguments" && payload.code == "suppressed_repeated_tool_error"
}
fn suppressed_retry_response(
name: &str,
args_json: &str,
prior: &ToolErrorResponse,
) -> ToolErrorResponse {
let suggested_action = prior.suggested_action.clone().or_else(|| {
if prior.category == "arguments" {
Some(format!(
"Correct the JSON arguments for `{name}` before retrying. Include all required fields and valid values in the next tool call."
))
} else {
Some(
"Change the approach or complete the required prerequisite before retrying."
.to_string(),
)
}
});
let mut error_msg = format!(
"EdgeCrab already saw the same `{name}` call fail earlier in this conversation. \
Repeating identical arguments would be flaky, so that retry was suppressed.\n\
Original error [{code}]: {original_error}",
code = prior.code,
original_error = prior.error,
);
if let Some(ref hint) = prior.usage_hint {
error_msg.push_str(&format!("\nHint: {hint}"));
}
if let Some(ref action) = suggested_action {
error_msg.push_str(&format!("\nSuggested fix: {action}"));
}
if let Some(ref alt_tool) = prior.suggested_tool {
error_msg.push_str(&format!("\nAlternative tool: {alt_tool}"));
}
ToolErrorResponse {
response_type: "tool_error".into(),
category: prior.category.clone(),
code: "suppressed_repeated_tool_error".into(),
code_num: 1099,
error: error_msg,
retryable: false,
suppress_retry: true,
suppression_key: Some(tool_attempt_fingerprint(name, args_json)),
tool: Some(name.to_string()),
suggested_tool: prior.suggested_tool.clone(),
suggested_action,
required_fields: prior.required_fields.clone(),
usage_hint: prior.usage_hint.clone(),
recovery_feedback: prior.recovery_feedback.clone(),
}
}
#[inline]
fn is_tool_error(result: &str) -> bool {
parse_tool_error_response(result).is_some() || result.starts_with("Tool error:")
}
fn emit_tool_done(
tx: Option<&tokio::sync::mpsc::UnboundedSender<crate::StreamEvent>>,
tool_call_id: &str,
name: &str,
args_json: &str,
tool_result: &str,
duration_ms: u64,
is_error: bool,
) {
crate::stream_observability::log_tool_done(tool_call_id, name, duration_ms, is_error);
crate::progress_sink::emit_optional(
tx,
crate::StreamEvent::ToolDone {
tool_call_id: tool_call_id.to_string(),
name: name.to_string(),
args_json: args_json.to_string(),
result_preview: summarize_tool_result_preview(name, tool_result, is_error),
duration_ms,
is_error,
},
);
if is_error
&& let Some(server_name) =
edgecrab_tools::tools::mcp_client::parse_mcp_oauth_required_server(tool_result)
{
crate::progress_sink::emit_optional(
tx,
crate::StreamEvent::McpOAuthRequired { server_name },
);
}
}
fn make_tool_progress_tx(
event_tx: Option<&tokio::sync::mpsc::UnboundedSender<crate::StreamEvent>>,
) -> Option<tokio::sync::mpsc::UnboundedSender<edgecrab_tools::ToolProgressUpdate>> {
let event_tx = event_tx.cloned()?;
let (tool_progress_tx, mut tool_progress_rx) =
tokio::sync::mpsc::unbounded_channel::<edgecrab_tools::ToolProgressUpdate>();
tokio::spawn(async move {
while let Some(update) = tool_progress_rx.recv().await {
let _ = event_tx.send(crate::StreamEvent::ToolProgress {
tool_call_id: update.tool_call_id,
name: update.tool_name,
message: update.message,
});
}
});
Some(tool_progress_tx)
}
fn build_shadow_judge_message(steering_hint: &str, reason: &str) -> String {
format!(
"[system: verification check indicates the task is not yet complete — {reason}. \
{steering_hint} \
Continue working and only stop once all parts of the original request are done with concrete evidence.]"
)
}
fn build_prune_spill_context<'a>(
conversation_session_id: &'a str,
cwd: &'a std::path::Path,
spill_config: &'a crate::tool_result_spill::SpillConfig,
spill_seq: &'a crate::tool_result_spill::SpillSequence,
) -> crate::compression::PruneSpillContext<'a> {
crate::compression::PruneSpillContext::new(
conversation_session_id,
cwd,
spill_config,
spill_seq,
)
}
fn local_prune_spill_config(
app_config_ref: &edgecrab_tools::config_ref::AppConfigRef,
) -> crate::tool_result_spill::SpillConfig {
crate::tool_result_spill::SpillConfig {
enabled: app_config_ref.result_spill,
threshold: app_config_ref.result_spill_threshold,
preview_lines: app_config_ref.result_spill_preview_lines,
}
}
fn try_local_structural_prune_request(
phase: crate::local_provider_policy::LocalStructuralPrunePhase,
system_prompt: Option<&str>,
messages: &[Message],
tool_defs: &[edgequake_llm::ToolDefinition],
spill_ctx: &crate::compression::PruneSpillContext<'_>,
active_context_length: usize,
) -> Option<(Vec<Message>, usize, usize, usize)> {
let prompt_before = estimate_request_prompt_tokens(system_prompt, messages, tool_defs);
let (pruned_messages, outcome) =
crate::local_provider_policy::try_apply_structural_tool_output_prune(
phase,
prompt_before,
active_context_length,
messages,
Some(spill_ctx),
)?;
let prompt_after = estimate_request_prompt_tokens(system_prompt, &pruned_messages, tool_defs);
Some((
pruned_messages,
outcome.tools_pruned,
prompt_before,
prompt_after,
))
}
fn available_toolsets_for_prompt(
registry: &edgecrab_tools::registry::ToolRegistry,
tool_names: &[String],
) -> Vec<String> {
let mut toolsets: Vec<String> = tool_names
.iter()
.filter_map(|name| registry.toolset_for_tool(name))
.collect();
toolsets.sort();
toolsets.dedup();
toolsets
}
fn summarize_tool_result_preview(name: &str, tool_result: &str, is_error: bool) -> Option<String> {
crate::tool_result_summary::summarize_tool_result_preview(name, tool_result, is_error)
}
fn get_budget_warning(api_call_count: u32, max_iterations: u32) -> Option<String> {
if max_iterations == 0 {
return None;
}
let progress = api_call_count as f64 / max_iterations as f64;
if progress >= 0.9 {
Some(format!(
"[URGENT: {}% of iteration budget used ({}/{}). You MUST provide a final response NOW — do not make further tool calls.]",
(progress * 100.0) as u32,
api_call_count,
max_iterations
))
} else if progress >= 0.7 {
Some(format!(
"[BUDGET: {}% of iteration budget used ({}/{}). Start wrapping up — avoid multi-step tool chains.]",
(progress * 100.0) as u32,
api_call_count,
max_iterations
))
} else {
None
}
}
fn inject_budget_warning(messages: &mut Vec<Message>, warning: &str) {
if let Some(msg) = messages.iter_mut().rev().find(|m| m.role == Role::Tool) {
let current = msg.text_content();
let new_content = if let Ok(mut v) = serde_json::from_str::<serde_json::Value>(¤t) {
if let Some(obj) = v.as_object_mut() {
obj.insert(
"_budget_warning".to_string(),
serde_json::Value::String(warning.to_string()),
);
serde_json::to_string(&v).unwrap_or_else(|_| format!("{}\n\n{}", current, warning))
} else {
format!("{}\n\n{}", current, warning)
}
} else {
format!("{}\n\n{}", current, warning)
};
msg.content = Some(Content::Text(new_content));
} else {
tracing::debug!("no tool messages found, injecting budget warning as user message");
messages.push(Message::user(warning));
}
}
fn strip_budget_warnings_from_history(messages: &mut Vec<Message>) {
for msg in messages.iter_mut().filter(|m| m.role == Role::Tool) {
let current = match &msg.content {
Some(Content::Text(t)) => t.clone(),
_ => continue,
};
if let Ok(mut v) = serde_json::from_str::<serde_json::Value>(¤t) {
if let Some(obj) = v.as_object_mut()
&& obj.remove("_budget_warning").is_some()
&& let Ok(cleaned) = serde_json::to_string(obj)
{
msg.content = Some(Content::Text(cleaned));
}
continue;
}
let cleaned = strip_budget_text_suffix(¤t);
if cleaned.len() < current.len() {
msg.content = Some(Content::Text(cleaned));
}
}
messages.retain(|m| {
if m.role != Role::User {
return true;
}
let text = m.text_content();
!(text.starts_with("[BUDGET:") || text.starts_with("[URGENT:"))
});
}
fn strip_budget_text_suffix(text: &str) -> String {
if !text.contains("\n\n[BUDGET:") && !text.contains("\n\n[URGENT:") {
return text.to_string();
}
let mut result = text.to_string();
loop {
let before = result.len();
for marker in &["\n\n[BUDGET:", "\n\n[URGENT:"] {
if let Some(pos) = result.rfind(marker) {
result.truncate(pos);
}
}
if result.len() == before {
break;
}
}
result
}
fn build_trajectory(
session_id: &str,
model: &str,
messages: &[Message],
api_calls: u32,
total_cost: f64,
completed: bool,
duration_seconds: f64,
) -> Trajectory {
let normalized_messages = normalize_messages_for_trajectory(messages);
let total_tokens = normalized_messages
.iter()
.map(|message| message.text_content().len() as u64 / 4)
.sum();
Trajectory {
session_id: session_id.to_string(),
model: model.to_string(),
timestamp: chrono::Utc::now().to_rfc3339(),
messages: normalized_messages,
metadata: TrajectoryMetadata {
task_id: None,
total_tokens,
total_cost,
api_calls,
tools_used: collect_used_tools(messages),
completed,
duration_seconds,
},
}
}
fn normalize_messages_for_trajectory(messages: &[Message]) -> Vec<Message> {
messages
.iter()
.cloned()
.map(|mut message| {
if let Some(Content::Text(text)) = &message.content {
message.content = Some(Content::Text(convert_scratchpad_to_think(text)));
}
if let Some(reasoning) = &message.reasoning {
message.reasoning = Some(convert_scratchpad_to_think(reasoning));
}
message
})
.collect()
}
fn collect_used_tools(messages: &[Message]) -> Vec<String> {
let mut tools = Vec::new();
for message in messages.iter().filter(|message| message.role == Role::Tool) {
if let Some(name) = &message.name
&& !tools.iter().any(|existing| existing == name)
{
tools.push(name.clone());
}
}
tools
}
#[inline]
fn extract_tool_error_text(result: &str) -> String {
if let Some(payload) = parse_tool_error_response(result) {
return payload.error;
}
result.to_string()
}
fn remember_tool_suppression(
suppressions: &Arc<Mutex<HashMap<String, ToolErrorResponse>>>,
name: &str,
args_json: &str,
result: &str,
) {
let Some(payload) = parse_tool_error_response(result) else {
return;
};
if !payload.suppress_retry {
return;
}
let mut guard = suppressions
.lock()
.expect("capability suppression cache lock poisoned");
guard.insert(tool_attempt_fingerprint(name, args_json), payload.clone());
if let Some(extra_key) = payload.suppression_key.clone() {
guard.insert(extra_key, payload);
}
}
fn append_tool_result_to_session(
session: &mut SessionState,
dctx: &DispatchContext,
tool_call_id: &str,
tool_name: &str,
tool_result: &str,
) {
let (provider, model) = dctx
.provider
.as_ref()
.map(|p| (p.name().to_string(), p.model().to_string()))
.unwrap_or_else(|| {
if let Some((p, m)) = edgecrab_tools::vision_models::parse_provider_model_spec(
&dctx.app_config_ref.active_model,
) {
(p, m)
} else {
("unknown".into(), dctx.app_config_ref.active_model.clone())
}
});
let store_images = crate::multimodal_tool_content::should_store_computer_use_images_in_session(
tool_name,
&provider,
&model,
&dctx.app_config_ref,
&session.tool_result_image_downgrades,
);
let delimit = edgecrab_security::tool_output_delimiters_enabled();
let (body, scan) =
edgecrab_security::prepare_tool_result_body(tool_call_id, tool_result, delimit);
if !matches!(scan.verdict, edgecrab_security::ThreatVerdict::Allow) {
let kinds: Vec<&str> = scan.findings.iter().map(|f| f.pattern_id).collect();
tracing::warn!(
tool = tool_name,
tool_call_id,
threats = ?kinds,
"tool output matched promptware patterns — delimited but not suppressed"
);
}
session
.messages
.push(Message::tool_result_for_session_policy(
tool_call_id,
tool_name,
&body,
store_images,
));
}
async fn process_response(
response: &edgequake_llm::LLMResponse,
session: &mut SessionState,
dctx: &DispatchContext,
tool_errors: &mut Vec<edgecrab_types::ToolErrorRecord>,
trackers: &mut TurnDispatchTrackers,
) -> Result<LoopAction, AgentError> {
if response.has_tool_calls() {
let max_delegate_calls = match dctx.app_config_ref.delegation_max_subagents {
0 => MAX_DELEGATE_TASK_CALLS_PER_TURN,
configured => configured
.min(MAX_DELEGATE_TASK_CALLS_PER_TURN as u32)
.max(1) as usize,
};
let effective_tool_calls =
cap_delegate_task_calls(&response.tool_calls, max_delegate_calls);
let our_tool_calls: Vec<edgecrab_types::ToolCall> = effective_tool_calls
.iter()
.map(|tc| {
dctx.registry.as_ref().map_or_else(
|| edgecrab_types::ToolCall::from_llm(tc),
|reg| edgecrab_tools::tool_call_pipeline::normalize_incoming_tool_call(reg, tc),
)
})
.collect();
let assistant_text = assistant_display_text(response);
let mut assistant_msg =
Message::assistant_with_tool_calls(&assistant_text, our_tool_calls.clone());
if let Some(ref thinking) = response.thinking_content {
assistant_msg.reasoning = Some(thinking.clone());
}
session.messages.push(assistant_msg);
session.session_tool_call_count += effective_tool_calls.len() as u32;
if let Some(reg) = dctx.registry.as_ref() {
let batch = edgecrab_tools::tool_call_pipeline::classify_unknown_tool_batch(
reg,
&dctx.engine_tool_names,
&our_tool_calls,
session.invalid_tool_call_retries,
);
if !batch.unknown_names.is_empty() {
session.invalid_tool_call_retries = batch.retry_count;
let invalid_preview = batch.unknown_names[0].chars().take(80).collect::<String>();
tracing::warn!(
invalid = %invalid_preview,
retry = batch.retry_count,
max = edgecrab_tools::MAX_INVALID_TOOL_RETRIES,
"unknown tool name after repair — sending structured error to model"
);
let unknown_set: std::collections::HashSet<String> =
batch.unknown_names.into_iter().collect();
for tc in &our_tool_calls {
let tool_result = if unknown_set.contains(&tc.function.name) {
edgecrab_tools::tool_call_pipeline::unknown_tool_error_response(
reg,
&tc.function.name,
)
} else {
"Skipped: another tool call in this turn used an invalid name. \
Please retry with a registered tool name."
.to_string()
};
append_tool_result_to_session(
session,
dctx,
&tc.id,
&tc.function.name,
&tool_result,
);
}
trackers.dedup.end_turn();
if batch.should_abort {
return Ok(LoopAction::PartialAbort {
reason: format!("Model generated invalid tool call: {invalid_preview}"),
});
}
return Ok(LoopAction::Continue);
}
session.invalid_tool_call_retries = 0;
}
let tool_turn_start = session.messages.len();
crate::turn_prologue::TurnPrologueState::reset_tool_dispatch(trackers);
let dispatch_calls: Vec<edgequake_llm::ToolCall> = effective_tool_calls
.iter()
.zip(our_tool_calls.iter())
.map(|(raw, normalized)| edgequake_llm::ToolCall {
id: raw.id.clone(),
call_type: raw.call_type.clone(),
function: edgequake_llm::FunctionCall {
name: normalized.function.name.clone(),
arguments: normalized.function.arguments.clone(),
},
thought_signature: raw.thought_signature.clone(),
})
.collect();
let mut sequential_calls = Vec::new();
let mut argument_loop_blocked = false;
let mut parallel_submitted: Vec<(String, String)> = Vec::new();
let mut eligible: Vec<edgecrab_types::ToolCall> = Vec::new();
for tc in &dispatch_calls {
if let Err(violation) = edgecrab_tools::mutation_turn_policy::check_tool_argument_budget(
&tc.function.name,
&tc.function.arguments,
dctx.app_config_ref.max_write_payload_bytes(),
dctx.provider.as_deref(),
) {
tracing::warn!(
tool = %violation.tool_name,
argument_bytes = violation.argument_bytes,
max_bytes = violation.max_bytes,
"rejecting tool call before dispatch — argument exceeds one-completion budget"
);
let tool_err = edgecrab_tools::recovery_catalog::tool_argument_budget_exceeded(
&violation.tool_name,
violation.argument_bytes,
violation.max_bytes,
violation.estimated_tokens,
);
let tool_result = tool_err.to_llm_response();
if should_count_failure_for_escalation(&tool_result) {
trackers
.failure
.record_failure(&extract_tool_error_text(&tool_result));
} else {
trackers.failure.record_success();
}
trackers
.dedup
.record(&tc.function.name, &tc.function.arguments, &tool_result);
tool_errors.push(edgecrab_types::ToolErrorRecord {
turn: session.api_call_count,
tool_name: tc.function.name.clone(),
arguments: tc.function.arguments.clone(),
error: extract_tool_error_text(&tool_result),
tool_result: tool_result.clone(),
});
append_tool_result_to_session(
session,
dctx,
&tc.id,
&tc.function.name,
&tool_result,
);
continue;
}
eligible.push(edgecrab_types::ToolCall {
id: tc.id.clone(),
r#type: tc.call_type.clone(),
function: edgecrab_types::FunctionCall {
name: tc.function.name.clone(),
arguments: tc.function.arguments.clone(),
},
thought_signature: tc.thought_signature.clone(),
});
}
let batch_plan = crate::tool_batch::plan_tool_batch(
dctx.registry.as_ref().map(|r| r.as_ref()),
&eligible,
);
let parallel_cap =
crate::tool_batch::parallel_max_workers(None, batch_plan.parallel.len().max(1));
let mut parallel_tasks =
crate::tool_batch::JoinSetToolBatch::new(Some(parallel_cap), batch_plan.parallel.len());
if batch_plan.has_parallel() {
tracing::debug!(
parallel = batch_plan.parallel.len(),
sequential = batch_plan.sequential.len(),
parallel_cap,
"tool batch plan (014 WS-B)"
);
}
for planned in &batch_plan.parallel {
if let Some(tx) = dctx.event_tx.as_ref() {
let _ = tx.send(crate::StreamEvent::ToolExec {
tool_call_id: planned.id.clone(),
name: planned.name.clone(),
args_json: planned.arguments.clone(),
});
}
let trackers_view = TurnDispatchTrackersView {
harness_advisory: &trackers.harness_advisory,
tool_guardrail: &trackers.tool_guardrail,
evidence: &trackers.evidence,
};
if let Some(blocked) = pre_dispatch_decision(
&trackers_view,
&session.messages,
&planned.name,
&planned.arguments,
&dctx.conversation_session_id,
) {
emit_tool_done(
dctx.event_tx.as_ref(),
&planned.id,
&planned.name,
&planned.arguments,
&blocked,
0,
true,
);
append_tool_result_to_session(session, dctx, &planned.id, &planned.name, &blocked);
trackers.record_tool_outcome(&planned.name, &planned.arguments, &blocked, true);
trackers
.dedup
.record(&planned.name, &planned.arguments, &blocked);
continue;
}
parallel_submitted.push((planned.id.clone(), planned.name.clone()));
let call = planned.clone();
let task_call = call.clone();
let mut inner = dctx.clone();
inner.event_tx = None; inner.delegation_event_tx = None;
crate::tool_batch::ToolBatchDispatch::dispatch(
&mut parallel_tasks,
call,
Box::pin(async move {
dispatch_single_tool(
&task_call.id,
&task_call.name,
&task_call.arguments,
&inner,
)
.await
}),
);
}
for planned in &batch_plan.sequential {
if let Some(tx) = dctx.event_tx.as_ref() {
let _ = tx.send(crate::StreamEvent::ToolExec {
tool_call_id: planned.id.clone(),
name: planned.name.clone(),
args_json: planned.arguments.clone(),
});
}
sequential_calls.push(edgequake_llm::ToolCall {
id: planned.id.clone(),
call_type: "function".into(),
function: edgequake_llm::FunctionCall {
name: planned.name.clone(),
arguments: planned.arguments.clone(),
},
thought_signature: None,
});
}
let mut received_parallel_ids: std::collections::HashSet<String> =
std::collections::HashSet::new();
while let Some(join_result) = parallel_tasks.join_next().await {
let task = match join_result {
Ok(crate::tool_batch::ToolBatchTaskOutcome::Completed(task)) => task,
Ok(crate::tool_batch::ToolBatchTaskOutcome::ConcurrencyGateClosed(call)) => {
crate::tool_batch::ToolBatchTask {
call,
output: (
"Tool batch concurrency semaphore closed".to_string(),
Vec::new(),
),
duration_ms: 0,
}
}
Err(e) => {
tracing::error!(error = %e, "parallel tool task panicked");
continue;
}
};
let tc_id = task.call.id;
let tc_name = task.call.name;
let args_json = task.call.arguments;
let (tool_result, injected_messages) = task.output;
let duration_ms = task.duration_ms;
let tool_result = apply_guardrail_result(
&mut trackers.tool_guardrail,
&tc_name,
&args_json,
&tool_result,
is_tool_error(&tool_result),
);
let is_error = is_tool_error(&tool_result);
trackers.record_tool_outcome(&tc_name, &args_json, &tool_result, is_error);
emit_tool_done(
dctx.event_tx.as_ref(),
&tc_id,
&tc_name,
&args_json,
&tool_result,
duration_ms,
is_error,
);
if is_error {
remember_tool_suppression(
&dctx.capability_suppressions,
&tc_name,
&args_json,
&tool_result,
);
tool_errors.push(edgecrab_types::ToolErrorRecord {
turn: session.api_call_count,
tool_name: tc_name.clone(),
arguments: args_json.clone(),
error: extract_tool_error_text(&tool_result),
tool_result: tool_result.clone(),
});
if let Some(payload) = parse_tool_error_response(&tool_result)
&& is_suppressed_argument_retry(&payload)
{
argument_loop_blocked = true;
}
if should_count_failure_for_escalation(&tool_result) {
trackers
.failure
.record_failure(&extract_tool_error_text(&tool_result));
} else {
trackers.failure.record_success();
}
} else {
trackers.failure.record_success();
}
received_parallel_ids.insert(tc_id.clone());
trackers.dedup.record(&tc_name, &args_json, &tool_result);
append_tool_result_to_session(session, dctx, &tc_id, &tc_name, &tool_result);
crate::compression::maybe_prune_computer_use_screenshots(
&mut session.messages,
dctx.app_config_ref.computer_use_keep_last_n_screenshots,
);
session.messages.extend(injected_messages);
}
for (tc_id, tc_name) in ¶llel_submitted {
if !received_parallel_ids.contains(tc_id) {
tracing::warn!(
tool_call_id = %tc_id,
tool_name = %tc_name,
"injecting error result for panicked parallel tool task"
);
session.messages.push(Message::tool_result(
tc_id,
tc_name,
&format!("Tool error: task panicked — internal error executing '{tc_name}'"),
));
}
}
for tc in sequential_calls {
if let Some(cached) = trackers
.dedup
.check_duplicate(&tc.function.name, &tc.function.arguments)
.map(|s| s.to_owned())
{
tracing::info!(
tool = %tc.function.name,
"duplicate tool call detected — returning cached result (FP11)"
);
let dedup_result = format!(
"{cached}\n\n[Note: This is a cached result — you already called `{}` with identical arguments in the previous turn. Try a different approach or different arguments.]",
tc.function.name
);
emit_tool_done(
dctx.event_tx.as_ref(),
&tc.id,
&tc.function.name,
&tc.function.arguments,
&dedup_result,
0,
false,
);
session.messages.push(Message::tool_result(
&tc.id,
&tc.function.name,
&dedup_result,
));
trackers
.dedup
.record(&tc.function.name, &tc.function.arguments, &cached);
continue;
}
let started = std::time::Instant::now();
let trackers_view = TurnDispatchTrackersView {
harness_advisory: &trackers.harness_advisory,
tool_guardrail: &trackers.tool_guardrail,
evidence: &trackers.evidence,
};
let (mut tool_result, injected_messages) = if let Some(blocked) = pre_dispatch_decision(
&trackers_view,
&session.messages,
&tc.function.name,
&tc.function.arguments,
&dctx.conversation_session_id,
) {
(blocked, Vec::new())
} else {
dispatch_single_tool(&tc.id, &tc.function.name, &tc.function.arguments, dctx).await
};
tool_result = apply_guardrail_result(
&mut trackers.tool_guardrail,
&tc.function.name,
&tc.function.arguments,
&tool_result,
is_tool_error(&tool_result),
);
let duration_ms = started.elapsed().as_millis() as u64;
let is_error = is_tool_error(&tool_result);
trackers.record_tool_outcome(
&tc.function.name,
&tc.function.arguments,
&tool_result,
is_error,
);
emit_tool_done(
dctx.event_tx.as_ref(),
&tc.id,
&tc.function.name,
&tc.function.arguments,
&tool_result,
duration_ms,
is_error,
);
if is_error {
remember_tool_suppression(
&dctx.capability_suppressions,
&tc.function.name,
&tc.function.arguments,
&tool_result,
);
tool_errors.push(edgecrab_types::ToolErrorRecord {
turn: session.api_call_count,
tool_name: tc.function.name.clone(),
arguments: tc.function.arguments.clone(),
error: extract_tool_error_text(&tool_result),
tool_result: tool_result.clone(),
});
if let Some(payload) = parse_tool_error_response(&tool_result)
&& is_suppressed_argument_retry(&payload)
{
argument_loop_blocked = true;
}
if should_count_failure_for_escalation(&tool_result) {
trackers
.failure
.record_failure(&extract_tool_error_text(&tool_result));
} else {
trackers.failure.record_success();
}
} else {
trackers.failure.record_success();
}
trackers
.dedup
.record(&tc.function.name, &tc.function.arguments, &tool_result);
append_tool_result_to_session(session, dctx, &tc.id, &tc.function.name, &tool_result);
crate::compression::maybe_prune_computer_use_screenshots(
&mut session.messages,
dctx.app_config_ref.computer_use_keep_last_n_screenshots,
);
session.messages.extend(injected_messages);
}
let tool_names: Vec<&str> = effective_tool_calls
.iter()
.map(|tc| tc.function.name.as_str())
.collect();
let browser_results: Vec<String> = session.messages[tool_turn_start..]
.iter()
.filter(|m| m.role == Role::Tool && m.name.as_deref() == Some("browser_navigate"))
.map(|m| m.text_content())
.collect();
let browser_refs: Vec<&str> = browser_results.iter().map(String::as_str).collect();
let process_ports = dctx.process_table.list_running_http_server_ports().await;
let known_dev_ports = edgecrab_tools::dev_server::merge_dev_server_ports(
&dctx.conversation_session_id,
&process_ports,
);
let blocked_tools: Vec<String> = if argument_loop_blocked {
effective_tool_calls
.iter()
.map(|tc| tc.function.name.clone())
.collect()
} else {
Vec::new()
};
finalize_tool_turn(
trackers,
ToolTurnFinalizeParams {
messages: &mut session.messages,
tool_turn_start,
tool_names: &tool_names,
browser_navigate_results: &browser_refs,
known_dev_ports: &known_dev_ports,
result_turn_budget_chars: dctx.app_config_ref.result_turn_budget_chars,
spill_config: crate::tool_result_spill::SpillConfig {
enabled: dctx.app_config_ref.result_spill,
threshold: dctx.app_config_ref.result_spill_threshold,
preview_lines: dctx.app_config_ref.result_spill_preview_lines,
},
session_id: &dctx.conversation_session_id,
cwd: &dctx.cwd,
spill_seq: &dctx.spill_seq,
max_write_payload_bytes: dctx.app_config_ref.max_write_payload_bytes(),
provider: dctx.provider.as_deref(),
argument_loop_blocked,
blocked_tool_names: blocked_tools,
},
)
.await;
if trackers.failure.should_escalate() {
let escalation = trackers.failure.escalation_message();
tracing::warn!(
consecutive_failures = trackers.failure.count,
"consecutive failure escalation triggered"
);
session.messages.push(Message::user(&escalation));
trackers.failure.record_success();
}
trackers.dedup.end_turn();
checkpoint_session_messages_mid_turn(dctx, session);
if trackers.guardrail_halt {
return Ok(LoopAction::GuardrailHalt);
}
return Ok(LoopAction::Continue);
}
let text = assistant_display_text(response);
let mut msg = Message::assistant(&text);
if let Some(ref thinking) = response.thinking_content {
msg.reasoning = Some(thinking.clone());
}
session.messages.push(msg);
if let Some(discovery) = dctx.discovered_plugins.as_ref() {
let history_json =
serde_json::to_value(&session.messages).unwrap_or_else(|_| serde_json::json!([]));
for plugin in discovery
.plugins
.iter()
.filter(|plugin| hermes_supports_hook(plugin, "post_llm_call"))
{
if let Err(error) = invoke_hermes_hook(
plugin,
"post_llm_call",
serde_json::json!({
"session_id": &dctx.conversation_session_id,
"user_message": "",
"assistant_response": &text,
"conversation_history": history_json,
"model": "",
"platform": dctx.platform.to_string(),
}),
)
.await
{
tracing::warn!(plugin = %plugin.name, ?error, "Hermes post_llm_call hook failed");
}
}
}
Ok(LoopAction::Done(text))
}
fn cap_delegate_task_calls(
tool_calls: &[edgequake_llm::ToolCall],
max_delegate_calls: usize,
) -> Vec<edgequake_llm::ToolCall> {
let delegate_count = tool_calls
.iter()
.filter(|tc| tc.function.name == "delegate_task")
.count();
if delegate_count <= max_delegate_calls {
return tool_calls.to_vec();
}
let mut kept_delegates = 0usize;
let mut truncated = Vec::with_capacity(tool_calls.len());
for tc in tool_calls {
if tc.function.name == "delegate_task" {
if kept_delegates < max_delegate_calls {
truncated.push(tc.clone());
kept_delegates += 1;
}
} else {
truncated.push(tc.clone());
}
}
tracing::warn!(
delegate_count,
max_delegate_calls,
"truncated excess delegate_task tool calls in a single turn"
);
truncated
}
#[inline]
fn should_count_failure_for_escalation(tool_result: &str) -> bool {
let Some(payload) = parse_tool_error_response(tool_result) else {
return true;
};
if is_suppressed_argument_retry(&payload) {
return false;
}
payload.category == "execution" || payload.category == "permission"
}
async fn dispatch_single_tool(
tool_call_id: &str,
name: &str,
args_json: &str,
dctx: &DispatchContext,
) -> (String, Vec<Message>) {
let span = crate::stream_observability::tool_execution_span(
tool_call_id,
name,
&dctx.conversation_session_id,
dctx.platform,
dctx.capability_grants,
);
async { dispatch_single_tool_impl(tool_call_id, name, args_json, dctx).await }
.instrument(span)
.await
}
async fn dispatch_single_tool_impl(
tool_call_id: &str,
name: &str,
args_json: &str,
dctx: &DispatchContext,
) -> (String, Vec<Message>) {
let Some(reg) = dctx.registry.as_ref() else {
return (
format!(
"Tool '{}' execution is not yet wired (no ToolRegistry provided).",
name
),
Vec::new(),
);
};
let resolved = reg.resolve_tool_call_name(name);
let lookup_name = if resolved.canonical.is_empty() {
name.to_string()
} else {
resolved.canonical
};
let attempt_key = tool_attempt_fingerprint(&lookup_name, args_json);
let semantic_key = invalid_args_semantic_key(reg, &lookup_name, args_json);
let prior = {
let guard = dctx
.capability_suppressions
.lock()
.expect("capability suppression cache lock poisoned");
guard.get(&attempt_key).cloned().or_else(|| {
semantic_key
.as_ref()
.and_then(|key| guard.get(key).cloned())
})
};
if let Some(prior) = prior {
if prior.code != "tool_unavailable" {
return (
serde_json::to_string(&suppressed_retry_response(&lookup_name, args_json, &prior))
.expect("suppressed retry payload serializes"),
Vec::new(),
);
}
}
if dctx.app_config_ref.tool_schema_mode == ToolSchemaMode::Indexed
&& is_deferred_not_on_wire(
&lookup_name,
&read_materialized_set(dctx.materialized_tools.as_ref()),
)
{
return (deferred_tool_error_response(&lookup_name), Vec::new());
}
let prepared = match edgecrab_tools::tool_call_pipeline::prepare_tool_call(reg, name, args_json)
{
Ok(prepared) => {
if prepared.name_repaired || prepared.original_name.trim() != prepared.name {
tracing::info!(
original = %prepared.original_name,
canonical = %prepared.name,
repaired = prepared.name_repaired,
"tool call normalized via tool_call_pipeline"
);
}
prepared
}
Err(e) => {
return (
format_preflight_tool_error(reg, &lookup_name, args_json, &e),
Vec::new(),
);
}
};
let name = prepared.name.as_str();
let normalized_args_json = prepared.args_json;
let prepared_args = prepared.args;
crate::lifecycle_hooks::emit_global(
crate::lifecycle_hooks::LifecycleEvent::ToolBefore,
serde_json::json!({
"session_id": &dctx.conversation_session_id,
"tool_name": name,
"args": &prepared_args,
}),
);
if let Some(tx) = dctx.event_tx.as_ref() {
let ctx_json = serde_json::json!({
"event": "tool:pre",
"tool_name": name,
"args_json": normalized_args_json,
"session_id": &dctx.conversation_session_id,
})
.to_string();
let _ = tx.send(crate::StreamEvent::HookEvent {
event: "tool:pre".to_string(),
context_json: ctx_json,
});
}
if let Some(discovery) = dctx.discovered_plugins.as_ref() {
for plugin in discovery
.plugins
.iter()
.filter(|plugin| hermes_supports_hook(plugin, "pre_tool_call"))
{
if let Err(error) = invoke_hermes_hook(
plugin,
"pre_tool_call",
serde_json::json!({
"tool_name": name,
"args": prepared_args.clone(),
"task_id": &dctx.conversation_session_id,
}),
)
.await
{
tracing::warn!(plugin = %plugin.name, ?error, "Hermes pre_tool_call hook failed");
}
}
}
if dctx.engine_tool_names.contains(name)
&& let Some(ref engine) = dctx.context_engine
{
let args = prepared_args.clone();
match engine.handle_tool_call(name, args).await {
Some(Ok(output)) => return (output, Vec::new()),
Some(Err(e)) => {
return (
edgecrab_types::ToolError::ExecutionFailed {
tool: name.to_string(),
message: e.to_string(),
}
.to_llm_response(),
Vec::new(),
);
}
None => {} }
}
let injected_messages = Arc::new(tokio::sync::Mutex::new(Vec::new()));
let ctx = build_tool_context(
&dctx.cwd,
dctx.app_config_ref.clone(),
&dctx.cancel,
&dctx.state_db,
dctx.platform,
dctx.capability_grants,
&dctx.process_table,
dctx.provider.clone(),
dctx.registry.clone(), dctx.sub_agent_runner.clone(),
dctx.delegation_event_tx.clone(),
dctx.clarify_tx.clone(),
dctx.approval_tx.clone(),
dctx.tool_progress_tx.clone(),
dctx.watch_notification_tx.clone(),
dctx.gateway_sender.clone(),
dctx.origin_chat.clone(),
Some(tool_call_id.to_string()),
Some(name.to_string()),
&dctx.conversation_session_id,
dctx.todo_store.clone(),
Some(injected_messages.clone()),
Some(Arc::clone(&dctx.mutation_turn)),
dctx.lsp_gate.clone(),
dctx.delegate_ctx.clone(),
dctx.kanban_task_id.clone(),
dctx.materialized_tools.clone(),
dctx.skills_zone_dirty.clone(),
);
let args_for_mutation = prepared_args.clone();
let result = match reg.dispatch(name, prepared_args, &ctx).await {
Ok(output) => output,
Err(ref e) if matches!(e.core_error(), ToolError::InvalidArgs { .. }) => {
if let Some(mut enriched) =
reg.enrich_invalid_args_error(name, e, Some(&normalized_args_json))
{
if enriched.suppression_key.is_none()
&& let Some(ref required_fields) = enriched.required_fields
{
enriched.suppression_key = invalid_args_missing_fields_suppression_key(
name,
&normalized_args_json,
required_fields,
);
}
serde_json::to_string(&enriched).expect("enriched error serializes")
} else {
e.to_llm_response()
}
}
Err(e) => e.to_llm_response(),
};
let mut result = result;
let tool_failed = is_tool_error(&result);
if !tool_failed
&& let Some(set) = dctx.materialized_tools.as_ref()
&& let Ok(mut guard) = set.write()
{
guard.touch(name);
}
if !tool_failed && let Some(tracker) = dctx.subdirectory_hints.as_ref() {
let mut guard = tracker.lock().await;
if let Some(hint) = guard.check_tool_call(name, &args_for_mutation) {
result.push_str(&hint);
}
}
dctx.mutation_turn
.record_tool_outcome(name, &args_for_mutation, &result, tool_failed);
crate::lifecycle_hooks::emit_global(
crate::lifecycle_hooks::LifecycleEvent::ToolAfter,
serde_json::json!({
"session_id": &dctx.conversation_session_id,
"tool_name": name,
"is_error": is_tool_error(&result),
}),
);
if let Some(tx) = dctx.event_tx.as_ref() {
let is_error = is_tool_error(&result);
let ctx_json = serde_json::json!({
"event": "tool:post",
"tool_name": name,
"session_id": &dctx.conversation_session_id,
"is_error": is_error,
})
.to_string();
let _ = tx.send(crate::StreamEvent::HookEvent {
event: "tool:post".to_string(),
context_json: ctx_json,
});
}
if let Some(discovery) = dctx.discovered_plugins.as_ref() {
for plugin in discovery
.plugins
.iter()
.filter(|plugin| hermes_supports_hook(plugin, "post_tool_call"))
{
if let Err(error) = invoke_hermes_hook(
plugin,
"post_tool_call",
serde_json::json!({
"tool_name": name,
"args": &args_for_mutation,
"result": &result,
"task_id": &dctx.conversation_session_id,
}),
)
.await
{
tracing::warn!(plugin = %plugin.name, ?error, "Hermes post_tool_call hook failed");
}
}
}
let queued_messages = {
let mut guard = injected_messages.lock().await;
std::mem::take(&mut *guard)
};
let result = if !is_tool_error(&result) {
let spill_config = crate::tool_result_spill::SpillConfig {
enabled: dctx.app_config_ref.result_spill,
threshold: dctx.app_config_ref.result_spill_threshold,
preview_lines: dctx.app_config_ref.result_spill_preview_lines,
};
let spill_ctx =
edgecrab_tools::artifact_spill::spill_context_from_args(name, &args_for_mutation);
match crate::tool_result_spill::maybe_spill(
name,
tool_call_id,
result,
&dctx.conversation_session_id,
&dctx.cwd,
&spill_config,
&dctx.spill_seq,
Some(&spill_ctx),
) {
crate::tool_result_spill::SpillOutcome::Inline(s) => s,
crate::tool_result_spill::SpillOutcome::Spilled { stub, .. } => stub,
}
} else {
result
};
if tool_failed
&& dctx.app_config_ref.tool_schema_mode == ToolSchemaMode::Indexed
&& let Some(mat) = dctx.materialized_tools.as_ref()
{
let targets =
edgecrab_tools::recovery_catalog::tools_to_materialize_from_error_json(&result);
if !targets.is_empty() {
let schemas = reg.get_definitions(None, None, &ctx);
let outcome = edgecrab_tools::materialize_tool_names(
&targets,
&schemas,
mat,
dctx.app_config_ref.max_materialized_tools,
edgecrab_tools::MaterializeSchemaStyle::Compact,
);
if !outcome.activated.is_empty() {
tracing::debug!(
tools = ?outcome.activated,
"indexed schema: recovery auto-materialize"
);
}
}
}
(result, queued_messages)
}
async fn auto_title_session(
db: Arc<edgecrab_state::SessionDb>,
session_id: String,
user_snippet: String,
assistant_snippet: String,
provider: Arc<dyn LLMProvider>,
) {
match db.get_session(&session_id) {
Ok(Some(rec)) => {
if let Some(ref existing) = rec.title {
if !existing.is_empty() && existing.len() < 80 && !existing.ends_with('…') {
tracing::debug!("session already has a title, skipping auto-title");
return;
}
}
}
_ => return,
}
let prompt = format!(
"Generate a short, descriptive title (3-7 words) for a conversation that starts with:\n\
User: {user_snippet}\n\nAssistant: {assistant_snippet}\n\n\
Return ONLY the title. No quotes, no punctuation at the end, no prefixes."
);
let messages = vec![
edgequake_llm::ChatMessage::system(
"You generate ultra-short session titles. Respond with ONLY the title, nothing else.",
),
edgequake_llm::ChatMessage::user(&prompt),
];
match provider.chat(&messages, None).await {
Ok(resp) => {
let mut title = resp.content.trim().to_string();
title = title.trim_matches(|c| c == '"' || c == '\'').to_string();
if title.to_lowercase().starts_with("title:") {
title = title[6..].trim().to_string();
}
if title.len() > 80 {
title = format!("{}…", crate::safe_truncate(&title, 77));
}
if !title.is_empty() {
if let Err(e) = db.update_session_title(&session_id, &title) {
tracing::debug!(error = %e, "auto-title DB update failed");
} else {
tracing::debug!(title, "auto-generated session title");
}
}
}
Err(e) => tracing::debug!(error = %e, "auto-title generation failed"),
}
}
struct BackgroundReflectionCtx {
messages: Vec<Message>,
system_prompt: Option<String>,
tool_defs: Vec<edgequake_llm::ToolDefinition>,
cwd: std::path::PathBuf,
registry: Option<Arc<ToolRegistry>>,
cancel: CancellationToken,
state_db: Option<Arc<edgecrab_state::SessionDb>>,
platform: edgecrab_types::Platform,
capability_grants: Option<edgecrab_tools::CapabilityGrants>,
process_table: Arc<edgecrab_tools::ProcessTable>,
provider: Arc<dyn edgequake_llm::LLMProvider>,
gateway_sender: Option<Arc<dyn edgecrab_tools::registry::GatewaySender>>,
sub_agent_runner: Option<Arc<dyn edgecrab_tools::SubAgentRunner>>,
app_config_ref: AppConfigRef,
conversation_session_id: String,
origin_chat: Option<edgecrab_types::OriginChat>,
todo_store: Option<Arc<edgecrab_tools::TodoStore>>,
}
async fn run_learning_reflection_bg(ctx: BackgroundReflectionCtx) {
let dctx = DispatchContext {
cwd: ctx.cwd.clone(),
registry: ctx.registry.clone(),
cancel: ctx.cancel.clone(),
state_db: ctx.state_db.clone(),
platform: ctx.platform,
capability_grants: ctx.capability_grants,
process_table: ctx.process_table.clone(),
provider: Some(Arc::clone(&ctx.provider)),
gateway_sender: ctx.gateway_sender.clone(),
sub_agent_runner: ctx.sub_agent_runner.clone(),
event_tx: None, delegation_event_tx: None,
clarify_tx: None, approval_tx: None, origin_chat: ctx.origin_chat.clone(),
app_config_ref: ctx.app_config_ref.clone(),
conversation_session_id: ctx.conversation_session_id.clone(),
todo_store: ctx.todo_store.clone(),
capability_suppressions: Arc::new(Mutex::new(HashMap::new())),
discovered_plugins: None,
spill_seq: Arc::new(crate::tool_result_spill::SpillSequence::new()),
context_engine: None,
engine_tool_names: Arc::new(std::collections::HashSet::new()),
mutation_turn: Arc::new(edgecrab_tools::MutationTurnState::new()),
lsp_gate: post_write_lsp_gate(&ctx.app_config_ref),
tool_progress_tx: None,
watch_notification_tx: None,
delegate_ctx: None,
kanban_task_id: None,
materialized_tools: None,
subdirectory_hints: None,
skills_zone_dirty: None,
};
let mut session = SessionState {
messages: ctx.messages,
cached_system_prompt: ctx.system_prompt,
..Default::default()
};
run_learning_reflection(&mut session, &ctx.tool_defs, &ctx.provider, &dctx).await;
}
async fn run_learning_reflection(
session: &mut SessionState,
tool_defs: &[edgequake_llm::ToolDefinition],
provider: &Arc<dyn LLMProvider>,
dctx: &DispatchContext,
) {
const REFLECTION_PROMPT: &str = "\
[system: learning checkpoint] This session used multiple tool calls. \
Please reflect briefly (1-2 sentences of thinking, not shown to the user): \
Did you discover a reusable workflow, debugging technique, or non-trivial \
pattern worth saving? If yes, call skill_manage(action='create', name='...', \
content='---\\nname: ...\\ndescription: ...\\n---\\n# Steps\\n...') to save it. \
Did you learn something important about the user, their project, or environment \
that should persist? If yes, call memory_write to record it. \
If nothing is worth saving, respond with exactly 'reflection: nothing to save' \
and stop — do NOT call any tools.";
session.messages.push(Message::user(REFLECTION_PROMPT));
let chat_messages = build_chat_messages(
session.cached_system_prompt.as_deref(),
&session.messages,
None, false, );
let response = match provider
.chat_with_tools(&chat_messages, tool_defs, None, None)
.await
{
Ok(r) => r,
Err(e) => {
tracing::debug!(error = %e, "learning reflection API call failed (non-fatal)");
session.messages.pop();
return;
}
};
let mut _reflection_tool_errors: Vec<edgecrab_types::ToolErrorRecord> = Vec::new();
let mut _reflection_trackers = TurnDispatchTrackers::new(3);
if let Err(e) = process_response(
&response,
session,
dctx,
&mut _reflection_tool_errors,
&mut _reflection_trackers,
)
.await
{
tracing::debug!(error = %e, "learning reflection tool dispatch failed (non-fatal)");
}
}
fn stamp_unset_message_timestamps(messages: &mut [Message], now: f64) {
let mut seq = 0u32;
for msg in messages.iter_mut() {
if msg.created_at.is_none() {
msg.created_at = Some(now + f64::from(seq) * 1e-6);
seq = seq.saturating_add(1);
}
}
}
fn checkpoint_session_messages_mid_turn(dctx: &DispatchContext, session: &mut SessionState) {
let Some(ref db) = dctx.state_db else {
return;
};
if session.messages.is_empty() {
return;
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs_f64();
stamp_unset_message_timestamps(&mut session.messages, now);
let existing = db.get_session(&dctx.conversation_session_id).ok().flatten();
let (source, user_id) = match &dctx.origin_chat {
Some(origin) => (origin.platform.clone(), Some(origin.chat_id.clone())),
None => (
existing
.as_ref()
.map(|s| s.source.clone())
.unwrap_or_else(|| "cli".to_string()),
existing.as_ref().and_then(|s| s.user_id.clone()),
),
};
let title = existing.as_ref().and_then(|s| s.title.clone()).or_else(|| {
session
.messages
.iter()
.find(|m| m.role == Role::User)
.map(|m| {
let t = m.text_content();
if t.len() > 80 {
format!("{}…", crate::safe_truncate(&t, 80))
} else {
t
}
})
});
let record = edgecrab_state::SessionRecord {
id: dctx.conversation_session_id.clone(),
source,
user_id,
model: Some(dctx.app_config_ref.active_model.clone()),
system_prompt: session.cached_system_prompt.clone(),
stable_system_prompt: session.cached_stable_prompt.clone(),
semi_stable_system_prompt: session.cached_semi_stable_prompt.clone(),
parent_session_id: existing.as_ref().and_then(|s| s.parent_session_id.clone()),
started_at: existing.as_ref().map(|s| s.started_at).unwrap_or(now),
ended_at: None,
end_reason: None,
message_count: session.messages.len() as i64,
tool_call_count: derive_tool_call_count(session) as i64,
input_tokens: session.session_input_tokens as i64,
output_tokens: session.session_output_tokens as i64,
cache_read_tokens: session.session_cache_read_tokens as i64,
cache_write_tokens: session.session_cache_write_tokens as i64,
reasoning_tokens: session.session_reasoning_tokens as i64,
last_prompt_tokens: session.last_prompt_tokens as i64,
estimated_cost_usd: existing
.as_ref()
.and_then(|s| s.estimated_cost_usd)
.or(Some(0.0)),
title,
};
if let Err(e) = db.save_session_with_messages(&record, &session.messages, now) {
tracing::debug!(error = %e, "mid-turn session checkpoint failed (non-fatal)");
}
}
fn derive_tool_call_count(session: &crate::agent::SessionState) -> u32 {
let from_messages = session
.messages
.iter()
.filter(|m| matches!(m.role, Role::Tool))
.count() as u32;
from_messages.max(session.session_tool_call_count)
}
fn sanitize_orphaned_tool_results(messages: &mut Vec<Message>) {
use std::collections::HashSet;
let mut valid_ids: HashSet<String> = HashSet::new();
for msg in messages.iter() {
if msg.role == Role::Assistant
&& let Some(ref calls) = msg.tool_calls
{
for tc in calls {
valid_ids.insert(tc.id.clone());
}
}
}
let before = messages.len();
messages.retain(|msg| {
if msg.role == Role::Tool {
msg.tool_call_id
.as_ref()
.is_some_and(|id| valid_ids.contains(id))
} else {
true
}
});
let removed = before - messages.len();
if removed > 0 {
tracing::info!(removed, "sanitized orphaned tool result messages");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::AgentBuilder;
use crate::goals::GoalStore;
use async_trait::async_trait;
use edgecrab_tools::{ProcessTable, ToolRegistry};
use edgequake_llm::traits::{StreamChunk, StreamUsage};
use edgequake_llm::{ChatMessage, CompletionOptions, FunctionCall, ToolChoice, ToolDefinition};
use serde_json::json;
use tempfile::TempDir;
#[derive(Clone)]
struct StreamingUsageProvider {
chunks: Vec<StreamChunk>,
}
#[derive(Clone)]
struct OrphanRejectingProvider;
#[derive(Clone)]
struct RetryCountingProvider {
provider_name: &'static str,
attempts: Arc<std::sync::atomic::AtomicUsize>,
last_options: Arc<Mutex<Option<CompletionOptions>>>,
}
struct FlakyToolStreamProvider {
attempts: Arc<std::sync::atomic::AtomicUsize>,
}
#[derive(Clone)]
struct FirstChunkTimeoutProvider {
stream_attempts: Arc<std::sync::atomic::AtomicUsize>,
nonstream_attempts: Arc<std::sync::atomic::AtomicUsize>,
}
#[derive(Clone)]
struct ToolStreamingRejectedProvider {
stream_attempts: Arc<std::sync::atomic::AtomicUsize>,
nonstream_attempts: Arc<std::sync::atomic::AtomicUsize>,
}
#[derive(Clone)]
struct StaticResponseProvider;
#[async_trait]
impl LLMProvider for StreamingUsageProvider {
fn name(&self) -> &str {
"streaming-usage-test"
}
fn model(&self) -> &str {
"streaming-usage-test-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
messages: &[ChatMessage],
options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.chat_with_tools(messages, &[], None, options).await
}
async fn chat_with_tools(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new("non-stream", self.model()))
}
async fn chat_with_tools_stream(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<
futures::stream::BoxStream<'static, edgequake_llm::Result<StreamChunk>>,
> {
use futures::StreamExt;
Ok(futures::stream::iter(self.chunks.clone().into_iter().map(Ok)).boxed())
}
fn supports_tool_streaming(&self) -> bool {
true
}
}
#[async_trait]
impl LLMProvider for OrphanRejectingProvider {
fn name(&self) -> &str {
"orphan-rejecting-test"
}
fn model(&self) -> &str {
"orphan-rejecting-test-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
messages: &[ChatMessage],
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.assert_no_orphaned_tool_result(messages)?;
Ok(edgequake_llm::LLMResponse::new(
"clean history",
self.model(),
))
}
}
impl OrphanRejectingProvider {
fn assert_no_orphaned_tool_result(
&self,
messages: &[ChatMessage],
) -> edgequake_llm::Result<()> {
let mut valid_tool_ids = std::collections::HashSet::new();
for message in messages {
if matches!(message.role, edgequake_llm::ChatRole::Assistant)
&& let Some(tool_calls) = &message.tool_calls
{
for tool_call in tool_calls {
valid_tool_ids.insert(tool_call.id.clone());
}
}
let tool_call_id = message.tool_call_id.as_deref();
if matches!(message.role, edgequake_llm::ChatRole::Tool)
&& tool_call_id.is_none_or(|id| !valid_tool_ids.contains(id))
{
return Err(edgequake_llm::LlmError::ApiError(
"orphaned tool result reached provider".into(),
));
}
}
Ok(())
}
}
#[async_trait]
impl LLMProvider for RetryCountingProvider {
fn name(&self) -> &str {
self.provider_name
}
fn model(&self) -> &str {
"retry-counting-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
_messages: &[ChatMessage],
options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.attempts
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
*self.last_options.lock().expect("lock") = options.cloned();
Err(edgequake_llm::LlmError::NetworkError(
"synthetic network failure".into(),
))
}
fn supports_tool_streaming(&self) -> bool {
true
}
}
#[async_trait]
impl LLMProvider for FlakyToolStreamProvider {
fn name(&self) -> &str {
"flaky-tool-stream"
}
fn model(&self) -> &str {
"flaky-tool-stream-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
messages: &[ChatMessage],
options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.chat_with_tools(messages, &[], None, options).await
}
async fn chat_with_tools(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new("non-stream", self.model()))
}
async fn chat_with_tools_stream(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<
futures::stream::BoxStream<'static, edgequake_llm::Result<StreamChunk>>,
> {
use futures::StreamExt;
let attempt = self
.attempts
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let chunks = if attempt == 0 {
vec![
StreamChunk::ToolCallDelta {
index: 0,
id: Some("call_write".into()),
function_name: None,
function_arguments: None,
thought_signature: None,
},
StreamChunk::Finished {
reason: "stop".into(),
ttft_ms: None,
usage: None,
},
]
} else {
vec![
StreamChunk::Content("recovered".into()),
StreamChunk::Finished {
reason: "stop".into(),
ttft_ms: None,
usage: None,
},
]
};
Ok(futures::stream::iter(chunks.into_iter().map(Ok)).boxed())
}
fn supports_tool_streaming(&self) -> bool {
true
}
}
#[async_trait]
impl LLMProvider for FirstChunkTimeoutProvider {
fn name(&self) -> &str {
"first-chunk-timeout"
}
fn model(&self) -> &str {
"first-chunk-timeout-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
messages: &[ChatMessage],
options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.chat_with_tools(messages, &[], None, options).await
}
async fn chat_with_tools(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.nonstream_attempts
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Ok(edgequake_llm::LLMResponse::new(
"fallback after stalled stream",
self.model(),
))
}
async fn chat_with_tools_stream(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<
futures::stream::BoxStream<'static, edgequake_llm::Result<StreamChunk>>,
> {
use futures::StreamExt;
self.stream_attempts
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Ok(futures::stream::pending().boxed())
}
fn supports_tool_streaming(&self) -> bool {
true
}
}
#[async_trait]
impl LLMProvider for ToolStreamingRejectedProvider {
fn name(&self) -> &str {
"tool-streaming-rejected"
}
fn model(&self) -> &str {
"tool-streaming-rejected-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
messages: &[ChatMessage],
options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.chat_with_tools(messages, &[], None, options).await
}
async fn chat_with_tools(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.nonstream_attempts
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Ok(edgequake_llm::LLMResponse::new(
"tool fallback",
self.model(),
))
}
async fn chat_with_tools_stream(
&self,
_messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<
futures::stream::BoxStream<'static, edgequake_llm::Result<StreamChunk>>,
> {
self.stream_attempts
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Err(edgequake_llm::LlmError::InvalidRequest(
"Tool calling is not supported in streaming mode".into(),
))
}
fn supports_tool_streaming(&self) -> bool {
true
}
}
#[derive(Clone, Default)]
struct GoalCapturingProvider {
last_user_tail: Arc<Mutex<Option<String>>>,
}
#[async_trait]
impl LLMProvider for GoalCapturingProvider {
fn name(&self) -> &str {
"goal-capturing-test"
}
fn model(&self) -> &str {
"goal-capturing-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
messages: &[ChatMessage],
options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.chat_with_tools(messages, &[], None, options).await
}
async fn chat_with_tools(
&self,
messages: &[ChatMessage],
_tools: &[ToolDefinition],
_tool_choice: Option<ToolChoice>,
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
let tail = messages.iter().rev().find_map(|m| {
matches!(m.role, edgequake_llm::ChatRole::User).then(|| m.content.clone())
});
*self.last_user_tail.lock().expect("lock") = tail;
Ok(edgequake_llm::LLMResponse::new("ok", self.model()))
}
}
#[async_trait]
impl LLMProvider for StaticResponseProvider {
fn name(&self) -> &str {
"static-response-test"
}
fn model(&self) -> &str {
"static-response-model"
}
fn max_context_length(&self) -> usize {
8192
}
async fn complete(
&self,
prompt: &str,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new(prompt, self.model()))
}
async fn complete_with_options(
&self,
prompt: &str,
_options: &CompletionOptions,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
self.complete(prompt).await
}
async fn chat(
&self,
_messages: &[ChatMessage],
_options: Option<&CompletionOptions>,
) -> edgequake_llm::Result<edgequake_llm::LLMResponse> {
Ok(edgequake_llm::LLMResponse::new("ok", self.model()))
}
}
fn write_api_hook_plugin(dir: &std::path::Path) {
std::fs::write(
dir.join("plugin.yaml"),
r#"
name: api-hooks
version: "1.0.0"
description: API hook recorder
provides_hooks:
- pre_api_request
- post_api_request
"#,
)
.expect("manifest");
std::fs::write(
dir.join("__init__.py"),
r#"
import json
from pathlib import Path
def _append(event_name, payload):
target = Path(__file__).with_name("api-hooks.jsonl")
with target.open("a", encoding="utf-8") as handle:
handle.write(json.dumps({"event": event_name, "payload": payload}) + "\n")
def register(ctx):
ctx.register_hook("pre_api_request", lambda **kwargs: _append("pre_api_request", kwargs))
ctx.register_hook("post_api_request", lambda **kwargs: _append("post_api_request", kwargs))
"#,
)
.expect("plugin");
}
fn api_hook_plugin(dir: &std::path::Path) -> edgecrab_plugins::DiscoveredPlugin {
let manifest = edgecrab_plugins::parse_hermes_manifest(dir).expect("manifest");
edgecrab_plugins::DiscoveredPlugin {
name: manifest.name.clone(),
version: manifest.version.clone(),
description: manifest.description.clone(),
compatibility: None,
kind: edgecrab_plugins::PluginKind::Hermes,
status: edgecrab_plugins::PluginStatus::Available,
path: dir.to_path_buf(),
manifest: Some(edgecrab_plugins::synthesize_hermes_manifest(dir, &manifest)),
skill: None,
tools: Vec::new(),
hooks: vec!["post_api_request".into(), "pre_api_request".into()],
trust_level: edgecrab_plugins::TrustLevel::Unverified,
install_source: None,
enabled: true,
source: edgecrab_plugins::SkillSource::User,
missing_env: Vec::new(),
related_skills: Vec::new(),
cli_commands: Vec::new(),
}
}
#[tokio::test]
async fn api_call_streaming_preserves_authoritative_usage() {
let provider: Arc<dyn LLMProvider> = Arc::new(StreamingUsageProvider {
chunks: vec![
StreamChunk::Content("streamed answer".to_string()),
StreamChunk::Finished {
reason: "stop".to_string(),
ttft_ms: None,
usage: Some(
StreamUsage::new(11, 7)
.with_cache_hit_tokens(2)
.with_thinking_tokens(5),
),
},
],
});
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let tokens_sent = std::sync::atomic::AtomicBool::new(false);
let response = api_call_streaming(
&provider,
&[ChatMessage::user("hello")],
&[],
None,
Some(&CompletionOptions {
max_tokens: Some(256),
..Default::default()
}),
&tx,
&tokens_sent,
Some(1),
)
.await
.expect("stream call");
assert_eq!(response.prompt_tokens, 11);
assert_eq!(response.completion_tokens, 7);
assert_eq!(response.total_tokens, 18);
assert_eq!(response.cache_hit_tokens, Some(2));
assert_eq!(response.thinking_tokens, Some(5));
assert_eq!(response.finish_reason.as_deref(), Some("stop"));
}
#[tokio::test]
async fn api_call_streaming_estimates_usage_when_provider_omits_it() {
let provider: Arc<dyn LLMProvider> = Arc::new(StreamingUsageProvider {
chunks: vec![
StreamChunk::ThinkingContent {
text: "reasoning trace".to_string(),
tokens_used: Some(3),
budget_total: None,
},
StreamChunk::Content("streamed answer".to_string()),
StreamChunk::Finished {
reason: "stop".to_string(),
ttft_ms: None,
usage: None,
},
],
});
let tool_defs = vec![ToolDefinition::function(
"echo",
"Echo input",
json!({
"type": "object",
"properties": {
"text": {"type": "string"}
},
"required": ["text"]
}),
)];
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let tokens_sent = std::sync::atomic::AtomicBool::new(false);
let response = api_call_streaming(
&provider,
&[
ChatMessage::system("system"),
ChatMessage::user("hello world"),
],
&tool_defs,
None,
Some(&CompletionOptions {
max_tokens: Some(512),
..Default::default()
}),
&tx,
&tokens_sent,
Some(1),
)
.await
.expect("stream call");
assert!(
response.prompt_tokens > 0,
"prompt tokens should be estimated"
);
assert!(
response.completion_tokens > 0,
"completion tokens should be estimated"
);
assert_eq!(response.thinking_tokens, Some(3));
assert_eq!(response.finish_reason.as_deref(), Some("stop"));
}
#[tokio::test]
async fn api_call_streaming_rejects_tool_calls_without_arguments() {
unsafe { std::env::set_var("EDGECRAB_TOOL_ARGS_EMPTY_FALLBACK", "0") };
let provider: Arc<dyn LLMProvider> = Arc::new(StreamingUsageProvider {
chunks: vec![
StreamChunk::ToolCallDelta {
index: 0,
id: Some("call_execute".into()),
function_name: Some("execute_code".into()),
function_arguments: None,
thought_signature: None,
},
StreamChunk::Finished {
reason: "stop".to_string(),
ttft_ms: None,
usage: None,
},
],
});
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let tokens_sent = std::sync::atomic::AtomicBool::new(false);
let err = api_call_streaming(
&provider,
&[ChatMessage::user("hello")],
&[],
None,
None,
&tx,
&tokens_sent,
Some(1),
)
.await
.expect_err("missing streamed arguments must be rejected");
assert!(
err.to_string().contains("finished without arguments"),
"unexpected error: {err}"
);
assert!(
!tokens_sent.load(std::sync::atomic::Ordering::Relaxed),
"tool-call deltas alone must not count as visible streamed output"
);
}
#[tokio::test]
async fn api_call_with_retry_recovers_after_visible_streamed_tool_json_breaks() {
unsafe { std::env::set_var("EDGECRAB_TOOL_ARGS_EMPTY_FALLBACK", "0") };
let provider: Arc<dyn LLMProvider> = Arc::new(StreamingUsageProvider {
chunks: vec![
StreamChunk::Content(
"Perfect! Now I have sufficient information. Let me create a comprehensive audit document:".into(),
),
StreamChunk::ToolCallDelta {
index: 0,
id: Some("call_write".into()),
function_name: Some("write_file".into()),
function_arguments: Some("{\"path\":".into()),
thought_signature: None,
},
StreamChunk::Finished {
reason: "tool_use".to_string(),
ttft_ms: None,
usage: None,
},
],
});
let cancel = CancellationToken::new();
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let tool_defs = vec![ToolDefinition::function(
"write_file",
"Write a file",
json!({
"type": "object",
"properties": {
"path": {"type": "string"},
"content": {"type": "string"}
},
"required": ["path", "content"]
}),
)];
let outcome = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&tool_defs,
0,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: Some(&tx),
use_native_streaming: true,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect("visible partial text should be preserved and recovered instead of crashing");
assert!(!outcome.disabled_native_tool_streaming);
assert_eq!(
outcome.response.finish_reason.as_deref(),
Some(FINISH_REASON_STREAM_INTERRUPTED)
);
assert!(outcome.response.tool_calls.is_empty());
assert!(outcome.response.content.contains("sufficient information"));
}
#[tokio::test]
async fn api_call_with_retry_does_not_double_retry_local_inference_transport_errors() {
let attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let provider: Arc<dyn LLMProvider> = Arc::new(RetryCountingProvider {
provider_name: "lmstudio",
attempts: attempts.clone(),
last_options: Arc::new(Mutex::new(None)),
});
let cancel = CancellationToken::new();
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let err = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&[],
3,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: Some(&tx),
use_native_streaming: false,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect_err("local transport failure should fail fast");
assert!(matches!(err, AgentError::Llm(_)));
assert!(
err.to_string()
.contains("did not start a duplicate request"),
"unexpected error: {err}"
);
assert_eq!(
attempts.load(std::sync::atomic::Ordering::SeqCst),
1,
"local inference transport errors must not trigger outer retries"
);
let mut saw_notice = false;
while let Ok(event) = rx.try_recv() {
if let crate::StreamEvent::ActivityNotice(ref msg) = event {
assert!(msg.contains("lmstudio"));
saw_notice = true;
}
}
assert!(saw_notice, "expected local transport stall ActivityNotice");
}
#[tokio::test]
async fn api_call_with_retry_does_not_double_retry_copilot_requests() {
let attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let provider: Arc<dyn LLMProvider> = Arc::new(RetryCountingProvider {
provider_name: "vscode-copilot",
attempts: attempts.clone(),
last_options: Arc::new(Mutex::new(None)),
});
let cancel = CancellationToken::new();
let err = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&[],
3,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: None,
use_native_streaming: false,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect_err("copilot request should fail");
assert!(matches!(err, AgentError::Llm(_)));
assert_eq!(
attempts.load(std::sync::atomic::Ordering::SeqCst),
1,
"Copilot already retries internally; the outer loop must not multiply attempts"
);
}
#[tokio::test]
async fn api_call_with_retry_forwards_completion_options() {
let attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let last_options = Arc::new(Mutex::new(None));
let provider: Arc<dyn LLMProvider> = Arc::new(RetryCountingProvider {
provider_name: "options-test-provider",
attempts,
last_options: last_options.clone(),
});
let cancel = CancellationToken::new();
let options = CompletionOptions {
max_tokens: Some(2048),
temperature: Some(0.1),
reasoning_effort: Some("low".into()),
..Default::default()
};
let _ = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&[],
0,
ApiCallContext {
options: Some(&options),
cancel: &cancel,
event_tx: None,
use_native_streaming: false,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await;
let recorded = last_options.lock().expect("lock").clone().expect("options");
assert_eq!(recorded.max_tokens, Some(2048));
assert_eq!(recorded.temperature, Some(0.1));
assert_eq!(recorded.reasoning_effort.as_deref(), Some("low"));
}
#[tokio::test]
async fn api_call_with_retry_retries_streaming_after_malformed_tool_name_stream() {
let attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let provider: Arc<dyn LLMProvider> = Arc::new(FlakyToolStreamProvider {
attempts: attempts.clone(),
});
let cancel = CancellationToken::new();
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let tool_defs = vec![ToolDefinition::function(
"write_file",
"Write a file",
json!({
"type": "object",
"properties": {
"path": {"type": "string"},
"content": {"type": "string"}
},
"required": ["path", "content"]
}),
)];
let outcome = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&tool_defs,
1,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: Some(&tx),
use_native_streaming: true,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect("malformed tool stream should retry streaming and recover");
let response = outcome.response;
assert_eq!(response.content, "recovered");
assert!(!outcome.disabled_native_tool_streaming);
assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 2);
}
#[tokio::test]
async fn api_call_with_retry_falls_back_when_streamed_tools_are_rejected() {
let stream_attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let nonstream_attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let provider: Arc<dyn LLMProvider> = Arc::new(ToolStreamingRejectedProvider {
stream_attempts: stream_attempts.clone(),
nonstream_attempts: nonstream_attempts.clone(),
});
let cancel = CancellationToken::new();
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let tool_defs = vec![ToolDefinition::function(
"write_file",
"Write a file",
json!({
"type": "object",
"properties": {
"path": {"type": "string"},
"content": {"type": "string"}
},
"required": ["path", "content"]
}),
)];
let outcome = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&tool_defs,
1,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: Some(&tx),
use_native_streaming: true,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect("tool-stream capability miss should downgrade cleanly");
assert_eq!(outcome.response.content, "tool fallback");
assert!(outcome.disabled_native_tool_streaming);
assert_eq!(stream_attempts.load(std::sync::atomic::Ordering::SeqCst), 1);
assert_eq!(
nonstream_attempts.load(std::sync::atomic::Ordering::SeqCst),
1
);
}
#[tokio::test]
async fn api_call_with_retry_falls_back_after_stream_stalls_before_first_chunk() {
let stream_attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let nonstream_attempts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let provider: Arc<dyn LLMProvider> = Arc::new(FirstChunkTimeoutProvider {
stream_attempts: stream_attempts.clone(),
nonstream_attempts: nonstream_attempts.clone(),
});
let cancel = CancellationToken::new();
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
let outcome = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&[],
1,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: Some(&tx),
use_native_streaming: true,
discovered_plugins: None,
conversation_session_id: "test-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect("stalled first chunk should fall back to non-streaming");
assert_eq!(outcome.response.content, "fallback after stalled stream");
assert_eq!(stream_attempts.load(std::sync::atomic::Ordering::SeqCst), 1);
assert_eq!(
nonstream_attempts.load(std::sync::atomic::Ordering::SeqCst),
1
);
}
#[tokio::test]
async fn api_call_with_retry_invokes_hermes_api_hooks() {
let temp = TempDir::new().expect("tempdir");
write_api_hook_plugin(temp.path());
let provider: Arc<dyn LLMProvider> = Arc::new(StaticResponseProvider);
let cancel = CancellationToken::new();
let discovery = edgecrab_plugins::PluginDiscovery {
plugins: vec![api_hook_plugin(temp.path())],
};
let outcome = api_call_with_retry(
&provider,
&[ChatMessage::user("hello")],
&[],
0,
ApiCallContext {
options: None,
cancel: &cancel,
event_tx: None,
use_native_streaming: false,
discovered_plugins: Some(&discovery),
conversation_session_id: "api-hook-session",
platform: edgecrab_types::Platform::Cli,
api_call_count: 0,
session: crate::failover::ClassifyContext::default(),
},
)
.await
.expect("api call");
assert_eq!(outcome.response.content, "ok");
let log = std::fs::read_to_string(temp.path().join("api-hooks.jsonl")).expect("hook log");
assert!(log.contains("pre_api_request"));
assert!(log.contains("post_api_request"));
assert!(log.contains("api-hook-session"));
}
#[test]
fn completion_options_include_model_budget_and_reasoning_policy() {
let config = crate::agent::AgentConfig {
temperature: Some(0.2),
reasoning_effort: Some("medium".into()),
model_config: crate::config::ModelConfig {
max_tokens: Some(3072),
..Default::default()
},
..Default::default()
};
let options = completion_options_for(&config);
assert_eq!(options.max_tokens, Some(3072));
assert_eq!(options.temperature, Some(0.2));
assert_eq!(options.reasoning_effort.as_deref(), Some("medium"));
}
#[test]
fn native_streaming_policy_enables_copilot_tool_streaming() {
let copilot_provider: Arc<dyn LLMProvider> = Arc::new(RetryCountingProvider {
provider_name: "vscode-copilot",
attempts: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
last_options: Arc::new(Mutex::new(None)),
});
let lmstudio_provider: Arc<dyn LLMProvider> = Arc::new(RetryCountingProvider {
provider_name: "lmstudio",
attempts: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
last_options: Arc::new(Mutex::new(None)),
});
let streaming_provider: Arc<dyn LLMProvider> = Arc::new(FlakyToolStreamProvider {
attempts: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
});
let tool_defs = vec![ToolDefinition::function(
"write_file",
"Write a file",
json!({
"type": "object",
"properties": {
"path": {"type": "string"},
"content": {"type": "string"}
},
"required": ["path", "content"]
}),
)];
assert!(
should_use_native_streaming(copilot_provider.as_ref(), &tool_defs, true, true),
"Copilot tool turns should stream via edgequake-llm SSE"
);
assert!(
!should_use_native_streaming(lmstudio_provider.as_ref(), &tool_defs, true, true),
"LM Studio tool turns should use non-streaming to avoid buffered arg stalls"
);
assert!(
should_use_native_streaming(streaming_provider.as_ref(), &[], true, true),
"Plain-text turns can still use native streaming"
);
}
#[test]
fn cap_delegate_task_calls_truncates_excess_and_preserves_other_calls() {
let delegate = |id: &str| edgequake_llm::ToolCall {
id: id.into(),
call_type: "function".into(),
function: FunctionCall {
name: "delegate_task".into(),
arguments: "{}".into(),
},
thought_signature: None,
};
let terminal = edgequake_llm::ToolCall {
id: "tool-terminal".into(),
call_type: "function".into(),
function: FunctionCall {
name: "terminal".into(),
arguments: r#"{"command":"pwd"}"#.into(),
},
thought_signature: None,
};
let tool_calls = vec![
delegate("delegate-1"),
terminal.clone(),
delegate("delegate-2"),
delegate("delegate-3"),
delegate("delegate-4"),
];
let capped = cap_delegate_task_calls(&tool_calls, 3);
assert_eq!(capped.len(), 4);
assert_eq!(capped[0].id, "delegate-1");
assert_eq!(capped[1].id, "tool-terminal");
assert_eq!(capped[2].id, "delegate-2");
assert_eq!(capped[3].id, "delegate-3");
}
#[tokio::test]
async fn execute_loop_basic() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.build()
.expect("build");
let result = agent
.execute_loop("hello", Some("Be helpful."), None, None, None, None)
.await
.expect("loop");
assert!(!result.final_response.is_empty());
assert_eq!(result.api_calls, 1);
assert!(!result.interrupted);
}
#[tokio::test]
async fn execute_loop_with_history() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.build()
.expect("build");
let history = vec![
Message::user("previous question"),
Message::assistant("previous answer"),
];
let result = agent
.execute_loop("follow-up", None, Some(history), None, None, None)
.await
.expect("loop");
assert_eq!(result.messages.len(), 4);
}
#[tokio::test]
async fn execute_loop_sanitizes_history_before_provider_call() {
let provider: Arc<dyn LLMProvider> = Arc::new(OrphanRejectingProvider);
let agent = AgentBuilder::new("mock")
.provider(provider)
.build()
.expect("build");
let history = vec![
Message::user("previous question"),
Message::tool_result("orphan-id", "read_file", "stale output"),
];
let result = agent
.execute_loop("follow-up", None, Some(history), None, None, None)
.await
.expect("loop");
assert_eq!(result.final_response, "clean history");
assert!(
result
.messages
.iter()
.all(|message| message.tool_call_id.as_deref() != Some("orphan-id")),
"orphaned tool result should be removed before persistence"
);
}
#[tokio::test]
async fn execute_loop_injects_goal_block_without_persisting() {
use crate::goals::InMemoryGoalStore;
let goal_store = Arc::new(InMemoryGoalStore::new());
let provider = Arc::new(GoalCapturingProvider::default());
let captures = provider.last_user_tail.clone();
let agent = AgentBuilder::new("mock")
.provider(provider)
.goal_store(goal_store.clone())
.build()
.expect("build");
agent.chat("hello").await.expect("first turn");
let sid = agent
.session
.read()
.await
.session_id
.clone()
.expect("session id");
goal_store
.set_goal(&sid, "Refactor payment service", 20)
.expect("set goal");
agent.chat("continue").await.expect("second turn");
let captured = captures.lock().expect("lock").clone().expect("captured");
assert!(captured.contains("[GOAL CONTEXT"));
assert!(captured.contains("Refactor payment service"));
let persisted = agent.messages().await;
assert!(
!persisted
.iter()
.any(|m| m.text_content().contains("[GOAL CONTEXT")),
"goal block must not be persisted in session messages"
);
}
#[tokio::test]
async fn execute_loop_injects_goal_after_compression() {
use crate::goals::InMemoryGoalStore;
let goal_store = Arc::new(InMemoryGoalStore::new());
let provider = Arc::new(GoalCapturingProvider::default());
let captures = provider.last_user_tail.clone();
let agent = AgentBuilder::new("mock")
.provider(provider)
.goal_store(goal_store.clone())
.build()
.expect("build");
agent.chat("seed").await.expect("seed");
let sid = agent
.session
.read()
.await
.session_id
.clone()
.expect("session id");
goal_store
.set_goal(&sid, "Stay on mission", 20)
.expect("set goal");
agent.force_compress().await;
agent.chat("after compress").await.expect("post compress");
let captured = captures.lock().expect("lock").clone().expect("captured");
assert!(captured.contains("Stay on mission"));
}
#[tokio::test]
async fn execute_loop_uses_cwd_override_for_context_discovery() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.build()
.expect("build");
let workspace = TempDir::new().expect("workspace");
std::fs::write(
workspace.path().join("AGENTS.md"),
"# Workspace Rules\n\nUse the override workspace.",
)
.expect("write AGENTS.md");
agent
.execute_loop("hello", None, None, None, Some(workspace.path()), None)
.await
.expect("loop");
let session = agent.session.read().await;
let prompt = session
.cached_system_prompt
.as_deref()
.expect("cached system prompt");
assert!(prompt.contains("Use the override workspace."));
}
#[test]
fn build_trajectory_normalizes_reasoning_and_collects_tools() {
let messages = vec![
Message::user("hello"),
Message::assistant("<REASONING_SCRATCHPAD>plan</REASONING_SCRATCHPAD>done"),
Message::tool_result("call_1", "read_file", "contents"),
Message::tool_result("call_2", "read_file", "more contents"),
];
let trajectory =
build_trajectory("session-1", "provider/model", &messages, 2, 0.25, true, 1.5);
assert_eq!(trajectory.session_id, "session-1");
assert_eq!(trajectory.metadata.api_calls, 2);
assert_eq!(
trajectory.metadata.tools_used,
vec!["read_file".to_string()]
);
assert!(
trajectory.messages[1]
.text_content()
.contains("<think>plan</think>")
);
}
#[tokio::test]
async fn execute_loop_resets_preexisting_interrupt() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.build()
.expect("build");
agent.interrupt();
let result = agent
.execute_loop("hello", None, None, None, None, None)
.await
.expect("loop");
assert!(
!result.interrupted,
"pre-loop interrupt must be reset, not permanently sticky"
);
assert!(!result.final_response.is_empty());
assert!(!agent.is_cancelled());
}
#[tokio::test]
async fn execute_loop_budget_exhaust() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(1)
.build()
.expect("build");
let result = agent
.execute_loop("hello", None, None, None, None, None)
.await
.expect("loop");
assert_eq!(result.api_calls, 1);
}
#[test]
fn build_chat_messages_prepends_system() {
let messages = vec![Message::user("hi")];
let chat_msgs = build_chat_messages(Some("system prompt"), &messages, None, true);
assert_eq!(chat_msgs.len(), 2);
}
#[test]
fn build_chat_messages_no_system() {
let messages = vec![Message::user("hi")];
let chat_msgs = build_chat_messages(None, &messages, None, true);
assert_eq!(chat_msgs.len(), 1);
}
#[test]
fn build_chat_messages_omits_computer_use_image_when_downgraded() {
use edgecrab_types::multimodal_disk_image_from_content;
let envelope = serde_json::json!({
"_multimodal": true,
"_image_path": "/tmp/test-capture.png",
"_image_mime": "image/png",
"text_summary": "capture summary",
"content": [{"type": "text", "text": "capture summary"}]
});
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("test-capture.png");
std::fs::write(&path, b"\x89PNG\r\n\x1a\n").expect("png");
let body = envelope
.to_string()
.replace("/tmp/test-capture.png", path.display().to_string().as_str());
let messages = vec![Message::tool_result("tc1", "computer_use", &body)];
let mut session = SessionState::default();
session
.tool_result_image_downgrades
.insert(("anthropic".into(), "claude-opus-4.6".into()));
let cfg = edgecrab_tools::config_ref::AppConfigRef::default();
let attach = crate::multimodal_tool_content::should_attach_computer_use_screenshot(
"anthropic",
"claude-opus-4.6",
&cfg,
&session.tool_result_image_downgrades,
);
assert!(!attach);
let chat = build_chat_messages(None, &messages, None, attach);
let tool = chat
.iter()
.find(|m| m.role == edgequake_llm::ChatRole::Tool)
.expect("tool");
assert!(tool.images.is_none());
assert!(multimodal_disk_image_from_content(&tool.content).is_some());
}
#[test]
fn build_chat_messages_attaches_disk_capture_when_policy_allows() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("test-capture.png");
std::fs::write(&path, b"\x89PNG\r\n\x1a\n").expect("png");
let body = format!(
r#"{{"_multimodal":true,"_image_path":"{}","_image_mime":"image/png","text_summary":"capture summary","content":[{{"type":"text","text":"capture summary"}}]}}"#,
path.display()
);
let messages = vec![Message::tool_result("tc1", "computer_use", &body)];
let chat = build_chat_messages(None, &messages, None, true);
let tool = chat
.iter()
.find(|m| m.role == edgequake_llm::ChatRole::Tool)
.expect("tool");
assert!(tool.images.as_ref().is_some_and(|imgs| !imgs.is_empty()));
}
fn test_resolved_cache(native_inner: bool) -> crate::prompt_cache_policy::ResolvedPromptCache {
crate::prompt_cache_policy::ResolvedPromptCache {
decision: crate::prompt_cache_policy::PromptCacheDecision {
should_cache: true,
native_inner_layout: native_inner,
},
config: CachePromptConfig {
enabled: true,
cache_ttl: Some("1h".into()),
..Default::default()
},
}
}
#[test]
fn build_chat_messages_with_cache_config() {
let messages = vec![Message::user(
"a long user message that is at least one thousand chars. "
.repeat(20)
.as_str(),
)];
let resolved = test_resolved_cache(true);
let chat_msgs =
build_chat_messages(Some("system prompt"), &messages, Some(&resolved), true);
assert_eq!(chat_msgs.len(), 2);
assert!(chat_msgs[0].cache_control.is_some());
}
#[test]
fn build_chat_messages_blocks_emits_two_system_messages() {
let msgs = vec![Message::user("hello")];
let out = build_chat_messages_blocks("STABLE", "", "DYNAMIC", &msgs, None, true);
assert_eq!(out.len(), 3);
}
#[test]
fn build_chat_messages_blocks_emits_three_system_messages_with_semi_stable() {
let msgs = vec![Message::user("hello")];
let out = build_chat_messages_blocks("STABLE", "SEMI", "DYNAMIC", &msgs, None, true);
assert_eq!(out.len(), 4);
assert_eq!(out[0].content, "STABLE");
assert_eq!(out[1].content, "SEMI");
assert_eq!(out[2].content, "DYNAMIC");
}
#[test]
fn build_chat_messages_blocks_stable_has_cache_control() {
let msgs: Vec<Message> = vec![];
let out =
build_chat_messages_blocks("STABLE CONTENT", "", "DYNAMIC CONTENT", &msgs, None, true);
assert_eq!(out.len(), 2);
assert!(
out[0].cache_control.is_some(),
"stable system message must carry cache_control"
);
assert!(
out[1].cache_control.is_none(),
"dynamic system message must NOT carry cache_control"
);
}
#[test]
fn build_chat_messages_blocks_semi_stable_has_5m_cache_control() {
let msgs: Vec<Message> = vec![];
let out =
build_chat_messages_blocks("STABLE", "SKILLS INDEX", "DYNAMIC", &msgs, None, true);
assert_eq!(out.len(), 3);
assert_eq!(
out[1].cache_control.as_ref().and_then(|c| c.ttl.as_deref()),
Some("5m"),
"semi-stable block must use 5m cache tier"
);
assert!(out[2].cache_control.is_none());
}
#[test]
fn build_chat_messages_blocks_dynamic_has_no_cache_control() {
let msgs: Vec<Message> = vec![];
let out = build_chat_messages_blocks("", "", "DYNAMIC ONLY", &msgs, None, true);
assert_eq!(out.len(), 1);
assert!(
out[0].cache_control.is_none(),
"dynamic-only system message must not be cached"
);
}
#[test]
fn build_chat_messages_blocks_with_cache_config_does_not_double_mark_system() {
let msgs = vec![Message::user("hi")];
let resolved = test_resolved_cache(true);
let out = build_chat_messages_blocks("STABLE", "", "DYNAMIC", &msgs, Some(&resolved), true);
assert!(out[0].cache_control.is_some());
assert!(out[1].cache_control.is_none());
}
#[test]
fn split_dynamic_after_cache_prefixes_extracts_suffix() {
let stable = "STABLE CONTENT";
let semi = "SKILLS INDEX";
let combined = "STABLE CONTENT\n\nSKILLS INDEX\n\nDYNAMIC CONTENT";
let dynamic = split_dynamic_after_cache_prefixes(combined, stable, semi);
assert_eq!(dynamic, "DYNAMIC CONTENT");
}
#[test]
fn split_dynamic_after_cache_prefixes_stable_only() {
let stable = "STABLE CONTENT";
let combined = "STABLE CONTENT\n\nDYNAMIC CONTENT";
let dynamic = split_dynamic_after_cache_prefixes(combined, stable, "");
assert_eq!(dynamic, "DYNAMIC CONTENT");
}
#[test]
fn split_dynamic_after_cache_prefixes_fallback_when_prefix_mismatch() {
let combined = "SOMETHING ELSE\n\nDYNAMIC";
let dynamic = split_dynamic_after_cache_prefixes(combined, "STABLE", "");
assert_eq!(dynamic, combined);
}
#[test]
fn split_dynamic_after_cache_prefixes_empty_stable_returns_combined() {
let combined = "ALL CONTENT";
let dynamic = split_dynamic_after_cache_prefixes(combined, "", "");
assert_eq!(dynamic, combined);
}
#[test]
fn prompt_cache_config_is_provider_aware() {
let prefix = crate::config::PromptPrefixCacheConfig::default();
assert!(
crate::prompt_cache_policy::resolve_prompt_cache(
"mock",
"gpt-4o",
None,
true,
&prefix,
)
.is_none(),
"non-Anthropic providers should not receive Anthropic cache markers"
);
assert!(
crate::prompt_cache_policy::provider_supports_prompt_caching(
"anthropic",
"claude-sonnet-4",
None,
)
);
}
#[test]
fn envelope_layout_uses_single_system_block() {
let mut session = SessionState::default();
session.cached_stable_prompt = Some("STABLE".into());
session.cached_system_prompt = Some("STABLE\nDYNAMIC".into());
let msgs = vec![Message::user("hi")];
let decision = crate::prompt_cache_policy::PromptCacheDecision {
should_cache: true,
native_inner_layout: false,
};
let resolved = crate::prompt_cache_policy::ResolvedPromptCache {
decision,
config: CachePromptConfig::default(),
};
let out = build_api_chat_messages(
&session,
&msgs,
Some(&resolved),
&edgequake_llm::MockProvider::new(),
&edgecrab_tools::config_ref::AppConfigRef::default(),
);
let system_count = out
.iter()
.filter(|m| matches!(m.role, edgequake_llm::ChatRole::System))
.count();
assert_eq!(
system_count, 1,
"envelope-layout providers must not emit two system blocks"
);
}
#[test]
fn stable_cache_control_uses_1h_ttl_from_config() {
let cfg = CachePromptConfig {
cache_ttl: Some("1h".to_string()),
..Default::default()
};
let cc = stable_cache_control(Some(&cfg)).expect("cache marker");
assert_eq!(cc.ttl.as_deref(), Some("1h"));
}
#[test]
fn available_toolsets_for_prompt_deduplicates_registry_matches() {
let registry = edgecrab_tools::registry::ToolRegistry::new();
let toolsets = available_toolsets_for_prompt(
®istry,
&[
"read_file".to_string(),
"write_file".to_string(),
"read_file".to_string(),
],
);
assert_eq!(toolsets, vec!["file".to_string()]);
}
#[test]
fn sanitize_removes_orphaned_tool_results() {
let mut messages = vec![
Message::user("hi"),
Message::tool_result("orphan-id", "read_file", "file content"),
Message::assistant("hello"),
];
sanitize_orphaned_tool_results(&mut messages);
assert_eq!(messages.len(), 2);
assert_eq!(messages[0].role, Role::User);
assert_eq!(messages[1].role, Role::Assistant);
}
#[test]
fn sanitize_keeps_valid_tool_results() {
let tc = edgecrab_types::ToolCall {
id: "valid-id".into(),
r#type: "function".into(),
function: edgecrab_types::FunctionCall {
name: "read_file".into(),
arguments: "{}".into(),
},
thought_signature: None,
};
let mut messages = vec![
Message::user("hi"),
Message::assistant_with_tool_calls("calling tool", vec![tc]),
Message::tool_result("valid-id", "read_file", "file content"),
];
sanitize_orphaned_tool_results(&mut messages);
assert_eq!(messages.len(), 3, "valid tool result should be kept");
}
#[tokio::test]
async fn budget_exhaustion_at_gate_returns_synthetic_response() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(0)
.build()
.expect("build");
let result = agent
.execute_loop("do something", Some("Be helpful."), None, None, None, None)
.await
.expect("loop should not error on budget exhaustion");
assert!(
!result.final_response.is_empty(),
"budget-exhausted agent must not return empty response"
);
assert!(
result.final_response.contains("iteration limit"),
"synthetic response should mention 'iteration limit'; got: '{}'",
result.final_response
);
assert!(
result.budget_exhausted,
"budget_exhausted must be true when loop exits via budget gate"
);
assert!(!result.interrupted, "interrupted must be false");
assert_eq!(
result.api_calls, 0,
"no API calls should occur with budget=0"
);
}
#[tokio::test]
async fn chat_never_returns_empty_on_budget_exhaustion() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(0)
.build()
.expect("build");
let response = agent
.chat("do a lot of things")
.await
.expect("chat should not error");
assert!(
!response.is_empty(),
"chat() must not return empty string on budget exhaustion"
);
}
#[tokio::test]
async fn normal_completion_resets_budget_exhausted_flag() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(10)
.build()
.expect("build");
let result = agent
.execute_loop("hello", Some("Be helpful."), None, None, None, None)
.await
.expect("loop");
assert!(!result.final_response.is_empty());
assert!(
!result.budget_exhausted,
"normal completion must not set budget_exhausted"
);
assert!(!result.interrupted);
}
#[tokio::test]
async fn budget_exactly_one_produces_response() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(1)
.build()
.expect("build");
let result = agent
.execute_loop("hello", None, None, None, None, None)
.await
.expect("loop");
assert!(!result.final_response.is_empty());
assert!(
!result.budget_exhausted,
"text response was produced, not exhausted"
);
assert_eq!(result.api_calls, 1);
}
#[tokio::test]
async fn budget_exhausted_exactly_on_tool_turn_boundary() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(0)
.build()
.expect("build");
let result = agent
.execute_loop("run", None, None, None, None, None)
.await
.expect("loop");
assert!(result.budget_exhausted, "budget_exhausted must be true");
assert!(
!result.final_response.is_empty(),
"synthetic response must not be empty"
);
assert_eq!(result.api_calls, 0, "no API calls before budget gate");
}
#[tokio::test]
async fn multi_turn_tool_chain_completes_with_sufficient_budget() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(10)
.build()
.expect("build");
let result = agent
.execute_loop("do something and respond", None, None, None, None, None)
.await
.expect("loop");
assert!(
!result.final_response.is_empty(),
"response must be non-empty"
);
assert!(
!result.budget_exhausted,
"should complete normally without budget exhaustion"
);
assert!(!result.interrupted);
assert_eq!(result.api_calls, 1, "one API call for a text-only response");
}
#[test]
fn inject_budget_warning_appends_to_tool_message_json() {
let mut messages = vec![
Message::user("task"),
Message::tool_result("id1", "read_file", r#"{"output": "some content"}"#),
];
inject_budget_warning(&mut messages, "[URGENT: wrap up]");
let last = messages.last().expect("budget warning target exists");
let text = last.text_content();
assert!(
text.contains("_budget_warning"),
"budget warning should be injected into JSON tool message; got: {text}"
);
assert!(
text.contains("wrap up"),
"warning text should be present; got: {text}"
);
}
#[test]
fn inject_budget_warning_appends_to_tool_message_plain() {
let mut messages = vec![
Message::user("task"),
Message::tool_result("id1", "read_file", "plain text output"),
];
inject_budget_warning(&mut messages, "[URGENT: wrap up]");
let last = messages.last().expect("budget warning target exists");
let text = last.text_content();
assert!(
text.contains("wrap up"),
"plain-text warning should be appended; got: {text}"
);
}
#[test]
fn inject_budget_warning_falls_back_to_user_message_when_no_tools() {
let mut messages = vec![
Message::user("hello"),
Message::assistant("how can I help?"),
];
let before = messages.len();
inject_budget_warning(&mut messages, "[BUDGET: 70%]");
assert_eq!(
messages.len(),
before + 1,
"should inject a new user message as fallback"
);
let last = messages.last().expect("fallback warning message exists");
assert_eq!(last.role, Role::User);
assert!(last.text_content().contains("70%"));
}
#[test]
fn get_budget_warning_none_below_70_percent() {
assert!(
get_budget_warning(6, 10).is_none(),
"60% should produce no warning"
);
}
#[test]
fn get_budget_warning_at_70_percent() {
let w = get_budget_warning(7, 10);
assert!(w.is_some(), "70% should produce BUDGET warning");
assert!(
w.expect("70% warning should exist").contains("BUDGET"),
"should say BUDGET"
);
}
#[test]
fn get_budget_warning_at_90_percent() {
let w = get_budget_warning(9, 10);
assert!(w.is_some(), "90% should produce URGENT warning");
assert!(
w.expect("90% warning should exist").contains("URGENT"),
"should say URGENT"
);
}
#[test]
fn get_budget_warning_zero_max_iterations() {
assert!(
get_budget_warning(5, 0).is_none(),
"zero max_iterations should produce no warning (avoid div-by-zero)"
);
}
#[test]
fn strip_budget_warnings_strips_json_key() {
let tool_content = r#"{"result":"ok","_budget_warning":"[BUDGET: 70% ...]"}"#;
let mut messages = vec![
Message::user("task"),
Message {
role: Role::Tool,
content: Some(Content::Text(tool_content.to_string())),
name: Some("my_tool".to_string()),
tool_call_id: Some("call-1".to_string()),
..Default::default()
},
];
strip_budget_warnings_from_history(&mut messages);
let text = messages[1].text_content();
assert!(
!text.contains("_budget_warning"),
"JSON key should be removed"
);
assert!(text.contains("\"result\":\"ok\""), "other fields preserved");
}
#[test]
fn strip_budget_warnings_strips_plain_text_suffix() {
let tool_content = "Here is the file content.\n\n[BUDGET: 70% of iteration budget used (7/10). Start wrapping up.]";
let mut messages = vec![
Message::user("task"),
Message {
role: Role::Tool,
content: Some(Content::Text(tool_content.to_string())),
name: Some("read_file".to_string()),
tool_call_id: Some("call-2".to_string()),
..Default::default()
},
];
strip_budget_warnings_from_history(&mut messages);
let text = messages[1].text_content();
assert_eq!(text, "Here is the file content.");
assert!(!text.contains("BUDGET"), "BUDGET text should be removed");
}
#[test]
fn strip_budget_warnings_strips_urgent_text_suffix() {
let tool_content =
"Result data\n\n[URGENT: 90% of iteration budget used (9/10). You MUST respond now.]";
let mut messages = vec![Message {
role: Role::Tool,
content: Some(Content::Text(tool_content.to_string())),
name: Some("terminal".to_string()),
tool_call_id: Some("call-3".to_string()),
..Default::default()
}];
strip_budget_warnings_from_history(&mut messages);
let text = messages[0].text_content();
assert_eq!(text, "Result data");
}
#[test]
fn strip_budget_warnings_removes_standalone_user_message() {
let mut messages = vec![
Message::user("write me a poem"),
Message::assistant("Here is a poem."),
Message::user("[BUDGET: 70% of iteration budget used. Start wrapping up.]"),
];
strip_budget_warnings_from_history(&mut messages);
assert_eq!(messages.len(), 2, "standalone budget user message removed");
assert_eq!(messages[0].text_content(), "write me a poem");
}
#[test]
fn strip_budget_warnings_removes_standalone_urgent_user_message() {
let mut messages = vec![
Message::user("help"),
Message::user(
"[URGENT: 90% of iteration budget used (9/10). You MUST provide a final response NOW — do not make further tool calls.]",
),
];
strip_budget_warnings_from_history(&mut messages);
assert_eq!(messages.len(), 1);
assert_eq!(messages[0].text_content(), "help");
}
#[test]
fn strip_budget_warnings_noop_on_clean_history() {
let mut messages = vec![
Message::user("task"),
Message {
role: Role::Tool,
content: Some(Content::Text(r#"{"result":"clean"}"#.to_string())),
name: Some("tool".to_string()),
tool_call_id: Some("call-4".to_string()),
..Default::default()
},
Message::assistant("done"),
];
let original_len = messages.len();
strip_budget_warnings_from_history(&mut messages);
assert_eq!(
messages.len(),
original_len,
"no messages removed from clean history"
);
assert_eq!(messages[1].text_content(), r#"{"result":"clean"}"#);
}
#[test]
fn strip_budget_warnings_strips_multiple_stacked_warnings() {
let tool_content = "actual content\n\n[BUDGET: 70% ...]\n\n[URGENT: 90% ...]";
let mut messages = vec![Message {
role: Role::Tool,
content: Some(Content::Text(tool_content.to_string())),
name: Some("tool".to_string()),
tool_call_id: Some("call-5".to_string()),
..Default::default()
}];
strip_budget_warnings_from_history(&mut messages);
assert_eq!(messages[0].text_content(), "actual content");
}
#[test]
fn strip_budget_text_suffix_no_op_when_absent() {
let text = "just normal content";
assert_eq!(strip_budget_text_suffix(text), text);
}
#[test]
fn sanitize_handles_empty_messages() {
let mut messages: Vec<Message> = Vec::new();
sanitize_orphaned_tool_results(&mut messages);
assert_eq!(messages.len(), 0);
}
#[test]
fn sanitize_removes_multiple_orphans() {
let mut messages = vec![
Message::user("hi"),
Message::tool_result("orphan-1", "read_file", "content-a"),
Message::tool_result("orphan-2", "write_file", "content-b"),
Message::assistant("done"),
];
sanitize_orphaned_tool_results(&mut messages);
assert_eq!(messages.len(), 2, "both orphans should be removed");
}
#[test]
fn sanitize_handles_tool_result_without_tool_call_id() {
let mut msg = Message::tool_result("some-id", "read_file", "data");
msg.tool_call_id = None; let mut messages = vec![Message::user("hi"), msg, Message::assistant("done")];
sanitize_orphaned_tool_results(&mut messages);
assert_eq!(messages.len(), 2, "None-id tool result should be removed");
}
#[test]
fn summarize_tool_result_preview_prefers_terminal_body() {
let preview = summarize_tool_result_preview(
"terminal",
"[terminal_result status=success backend=local cwd=/tmp exit_code=0]\nhello world\n",
false,
)
.expect("preview");
assert_eq!(preview, "hello world");
}
#[test]
fn summarize_tool_result_preview_extracts_error_text() {
let preview = summarize_tool_result_preview(
"terminal",
"Tool error: permission denied while executing command",
true,
)
.expect("preview");
assert!(preview.contains("permission denied"));
}
#[test]
fn summarize_tool_result_preview_summarizes_web_search_results() {
let preview = summarize_tool_result_preview(
"web_search",
r#"{"success":true,"backend":"Brave","results":[{"title":"A"},{"title":"B"}]}"#,
false,
)
.expect("preview");
assert_eq!(preview, "2 results via Brave");
}
#[test]
fn summarize_tool_result_preview_shows_web_search_fallback() {
let preview = summarize_tool_result_preview(
"web_search",
r#"{"success":true,"backend":"ddgs","fallback_from":"tavily","results":[{"title":"A"}]}"#,
false,
)
.expect("preview");
assert_eq!(preview, "1 result via ddgs (fallback from tavily)");
}
#[test]
fn summarize_tool_result_preview_shows_skipped_tool_override() {
let preview = summarize_tool_result_preview(
"web_search",
r#"{"success":true,"backend":"ddgs","skipped_tool_override":"parallel","results":[{"title":"A"}]}"#,
false,
)
.expect("preview");
assert_eq!(preview, "(ignored parallel) 1 result via ddgs");
}
#[test]
fn summarize_tool_result_preview_summarizes_todo_state() {
let preview = summarize_tool_result_preview(
"todo",
r#"{"todos":[],"summary":{"total":4,"completed":2,"in_progress":1,"not_started":1,"cancelled":0}}"#,
false,
)
.expect("preview");
assert_eq!(preview, "2/4 done, 1 in progress");
}
#[test]
fn summarize_tool_result_preview_supports_manage_todo_list_alias() {
let preview = summarize_tool_result_preview(
"manage_todo_list",
r#"{"todos":[],"summary":{"total":3,"completed":1,"in_progress":1,"not_started":1,"cancelled":0}}"#,
false,
)
.expect("preview");
assert_eq!(preview, "1/3 done, 1 in progress");
}
#[test]
fn summarize_tool_result_preview_summarizes_delegate_batch() {
let preview = summarize_tool_result_preview(
"delegate_task",
r#"{"results":[{"status":"success"},{"status":"completed"},{"status":"error"}],"total_duration_seconds":1.25}"#,
false,
)
.expect("preview");
assert_eq!(preview, "2/3 task(s) completed in 1.25s");
}
#[test]
fn summarize_tool_result_preview_summarizes_reported_task_status() {
let preview = summarize_tool_result_preview(
"report_task_status",
r#"{"status":"in_progress","summary":"wired the TUI banners","remaining_steps":["run tests"]}"#,
false,
)
.expect("preview");
assert_eq!(preview, "progress: wired the TUI banners · 1 step(s) left");
}
#[test]
fn build_chat_messages_tool_role_uses_tool_call_id() {
let tc = edgecrab_types::ToolCall {
id: "tc-abc".into(),
r#type: "function".into(),
function: edgecrab_types::FunctionCall {
name: "read_file".into(),
arguments: "{}".into(),
},
thought_signature: None,
};
let messages = vec![
Message::user("read something"),
Message::assistant_with_tool_calls("sure", vec![tc]),
Message::tool_result("tc-abc", "read_file", "contents"),
];
let chat_msgs = build_chat_messages(None, &messages, None, true);
assert_eq!(chat_msgs.len(), 3);
}
#[test]
fn build_chat_messages_empty_input() {
let chat_msgs = build_chat_messages(None, &[], None, true);
assert_eq!(
chat_msgs.len(),
0,
"empty messages with no system → 0 chat messages"
);
}
fn make_dispatch_context_for_test(
registry: &Arc<ToolRegistry>,
cancel: &CancellationToken,
state_db: &Option<Arc<edgecrab_state::SessionDb>>,
process_table: &Arc<ProcessTable>,
capability_suppressions: Arc<Mutex<HashMap<String, ToolErrorResponse>>>,
) -> DispatchContext {
DispatchContext {
cwd: std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
registry: Some(Arc::clone(registry)),
cancel: cancel.clone(),
state_db: state_db.clone(),
platform: edgecrab_types::Platform::Cli,
capability_grants: None,
process_table: Arc::clone(process_table),
provider: None,
gateway_sender: None,
sub_agent_runner: None,
event_tx: None,
delegation_event_tx: None,
clarify_tx: None,
approval_tx: None,
origin_chat: None,
app_config_ref: AppConfigRef::default(),
conversation_session_id: "test-conversation".into(),
todo_store: None,
capability_suppressions,
discovered_plugins: None,
spill_seq: Arc::new(crate::tool_result_spill::SpillSequence::new()),
context_engine: None,
engine_tool_names: Arc::new(std::collections::HashSet::new()),
mutation_turn: Arc::new(edgecrab_tools::MutationTurnState::new()),
lsp_gate: None,
tool_progress_tx: None,
watch_notification_tx: None,
delegate_ctx: None,
kanban_task_id: None,
materialized_tools: None,
subdirectory_hints: None,
skills_zone_dirty: None,
}
}
#[tokio::test]
async fn dispatch_single_tool_blocks_deferred_until_materialized() {
use edgecrab_tools::ToolSchemaMode;
use std::sync::{Arc, RwLock};
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let mut dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions,
);
dctx.app_config_ref.tool_schema_mode = ToolSchemaMode::Indexed;
let (blocked, _) = dispatch_single_tool(
"call-browser",
"browser_navigate",
r#"{"url":"https://example.com"}"#,
&dctx,
)
.await;
assert!(
blocked.contains("tool_search"),
"deferred tool should require materialization: {blocked}"
);
let mut set = edgecrab_tools::MaterializedToolSet::new();
set.insert("browser_navigate", 0);
dctx.materialized_tools = Some(Arc::new(RwLock::new(set)));
let (allowed, _) = dispatch_single_tool(
"call-browser-2",
"browser_navigate",
r#"{"url":"https://example.com"}"#,
&dctx,
)
.await;
assert!(
!allowed.contains("not on your wire schema"),
"materialized tool should pass wire gate: {allowed}"
);
}
#[tokio::test]
async fn dispatch_single_tool_uses_dispatch_context_cwd() {
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let mut dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions,
);
let workspace = TempDir::new().expect("workspace");
std::fs::write(workspace.path().join("proof.txt"), "dispatch cwd works").expect("write");
dctx.cwd = workspace.path().to_path_buf();
let (result, injected_messages) = dispatch_single_tool(
"call-read-file",
"read_file",
r#"{"path":"proof.txt","line_numbers":false}"#,
&dctx,
)
.await;
assert!(injected_messages.is_empty());
assert!(result.contains("dispatch cwd works"), "got: {result}");
}
#[tokio::test]
async fn dispatch_single_tool_returns_structured_json_error() {
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions,
);
let (result, injected_messages) =
dispatch_single_tool("call-read-file", "read_file", "{}", &dctx).await;
assert!(injected_messages.is_empty());
let parsed = parse_tool_error_response(&result).expect("structured tool error");
assert_eq!(parsed.response_type, "tool_error");
assert_eq!(parsed.category, "arguments");
assert_eq!(parsed.code, "invalid_arguments");
assert_eq!(parsed.tool.as_deref(), Some("read_file"));
}
#[tokio::test]
async fn dispatch_single_tool_suppresses_repeated_capability_retry() {
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions.clone(),
);
let args_json = r#"{"command":"top"}"#;
let (first, first_injected) =
dispatch_single_tool("call-terminal-1", "terminal", args_json, &dctx).await;
assert!(first_injected.is_empty());
let first_payload = parse_tool_error_response(&first).expect("structured error");
assert_eq!(first_payload.code, "non_interactive_terminal_required");
remember_tool_suppression(&capability_suppressions, "terminal", args_json, &first);
let (second, second_injected) =
dispatch_single_tool("call-terminal-2", "terminal", args_json, &dctx).await;
assert!(second_injected.is_empty());
let second_payload = parse_tool_error_response(&second).expect("structured error");
assert_eq!(second_payload.code, "suppressed_repeated_tool_error");
assert!(second_payload.error.contains("same `terminal` call fail"));
}
#[tokio::test]
async fn dispatch_single_tool_suppresses_repeated_invalid_argument_retry() {
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let mut dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions.clone(),
);
let workspace = TempDir::new().expect("workspace");
std::fs::write(workspace.path().join("audit.md"), "existing content").expect("seed");
dctx.cwd = workspace.path().to_path_buf();
let args_json = r#"{"path":"audit.md"}"#;
let (first, first_injected) =
dispatch_single_tool("call-write-1", "write_file", args_json, &dctx).await;
assert!(first_injected.is_empty());
let first_payload = parse_tool_error_response(&first).expect("structured error");
assert_eq!(first_payload.code, "invalid_arguments");
remember_tool_suppression(&capability_suppressions, "write_file", args_json, &first);
let (second, second_injected) =
dispatch_single_tool("call-write-2", "write_file", args_json, &dctx).await;
assert!(second_injected.is_empty());
let second_payload = parse_tool_error_response(&second).expect("structured error");
assert_eq!(second_payload.code, "suppressed_repeated_tool_error");
assert_eq!(second_payload.category, "arguments");
assert!(second_payload.error.contains("same `write_file` call fail"));
}
#[tokio::test]
async fn dispatch_single_tool_accepts_write_file_file_path_alias() {
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let mut dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions,
);
let workspace = TempDir::new().expect("workspace");
dctx.cwd = workspace.path().to_path_buf();
let args_json = r#"{"file_path":"build_ppt.py","content":"print('ppt')\n"}"#;
let (result, injected) =
dispatch_single_tool("call-write-alias", "write_file", args_json, &dctx).await;
assert!(injected.is_empty());
let payload: serde_json::Value = serde_json::from_str(&result).expect("json");
assert_eq!(payload.get("ok"), Some(&serde_json::json!(true)));
assert!(workspace.path().join("build_ppt.py").exists());
}
#[tokio::test]
async fn dispatch_single_tool_suppresses_semantic_invalid_argument_retry() {
let registry = Arc::new(ToolRegistry::new());
let cancel = CancellationToken::new();
let state_db = None;
let process_table = Arc::new(ProcessTable::new());
let capability_suppressions = Arc::new(Mutex::new(HashMap::new()));
let dctx = make_dispatch_context_for_test(
®istry,
&cancel,
&state_db,
&process_table,
capability_suppressions.clone(),
);
let first_args = r#"{"content":"first"}"#;
let second_args = r#"{"content":"second","if_exists":"overwrite"}"#;
let (first, first_injected) =
dispatch_single_tool("call-write-semantic-1", "write_file", first_args, &dctx).await;
assert!(first_injected.is_empty());
let first_payload = parse_tool_error_response(&first).expect("structured error");
assert_eq!(first_payload.code, "invalid_arguments");
remember_tool_suppression(&capability_suppressions, "write_file", first_args, &first);
let (second, second_injected) =
dispatch_single_tool("call-write-semantic-2", "write_file", second_args, &dctx).await;
assert!(second_injected.is_empty());
let second_payload = parse_tool_error_response(&second).expect("structured error");
assert_eq!(second_payload.code, "suppressed_repeated_tool_error");
assert_eq!(second_payload.category, "arguments");
assert!(
second_payload.error.contains("same `write_file` call fail"),
"semantic suppression should block varied malformed retries"
);
}
#[tokio::test]
async fn cancellation_sets_interrupted_not_budget_exhausted() {
let provider: Arc<dyn LLMProvider> = Arc::new(edgequake_llm::MockProvider::new());
let agent = AgentBuilder::new("mock")
.provider(provider)
.max_iterations(100)
.build()
.expect("build");
let result = agent
.execute_loop("hello", None, None, None, None, None)
.await
.expect("loop");
assert!(
!result.budget_exhausted,
"normal completion must not set budget_exhausted"
);
assert!(
!result.interrupted,
"normal completion must not set interrupted"
);
assert!(!result.final_response.is_empty());
}
#[test]
fn suppressed_retry_includes_original_error_and_hints() {
use edgecrab_types::ToolErrorResponse;
let prior = ToolErrorResponse {
response_type: "tool_error".into(),
category: "arguments".into(),
code: "invalid_args".into(),
code_num: 1002,
error: "missing field `path`".into(),
retryable: true,
suppress_retry: false,
suppression_key: None,
tool: Some("read_file".into()),
suggested_tool: Some("search_files".into()),
suggested_action: None,
required_fields: Some(vec!["path".into()]),
usage_hint: Some("Required: path: string".into()),
recovery_feedback: None,
};
let resp = suppressed_retry_response("read_file", r#"{"wrong":"args"}"#, &prior);
assert!(
resp.error.contains("missing field `path`"),
"should include original error"
);
assert!(
resp.error.contains("Required: path: string"),
"should include usage hint"
);
assert!(
resp.error.contains("search_files"),
"should include alternative tool"
);
assert_eq!(
resp.required_fields.as_deref(),
Some(&["path".to_string()][..])
);
assert!(resp.usage_hint.is_some());
}
#[test]
fn invalid_args_missing_fields_key_detects_missing_required_fields() {
let required = vec!["path".to_string(), "content".to_string()];
let key = invalid_args_missing_fields_suppression_key(
"write_file",
r#"{"content":"hello"}"#,
&required,
)
.expect("missing path should generate a semantic suppression key");
assert_eq!(key, "invalid_args:write_file:missing:path");
}
#[test]
fn dispatch_tool_name_normalization_clean() {
let reg = edgecrab_tools::ToolRegistry::new();
let resolved = reg.resolve_tool_call_name("web_extract");
assert_eq!(resolved.canonical, "web_extract");
assert!(!resolved.repaired);
}
#[test]
fn dispatch_tool_name_normalization_strips_channel_token() {
let reg = edgecrab_tools::ToolRegistry::new();
let resolved = reg.resolve_tool_call_name("web_extract<|channel|>commentary");
assert_eq!(resolved.canonical, "web_extract");
}
#[test]
fn dispatch_tool_name_normalization_strips_im_end_token() {
let reg = edgecrab_tools::ToolRegistry::new();
let resolved = reg.resolve_tool_call_name("read_file<|im_end|>");
assert_eq!(resolved.canonical, "read_file");
}
#[test]
fn dispatch_tool_name_normalization_spaces() {
let reg = edgecrab_tools::ToolRegistry::new();
assert_eq!(
reg.resolve_tool_call_name("read file").canonical,
"read_file"
);
}
#[test]
fn dispatch_tool_name_normalization_hyphens() {
let reg = edgecrab_tools::ToolRegistry::new();
assert_eq!(
reg.resolve_tool_call_name("web-extract").canonical,
"web_extract"
);
}
#[test]
fn dispatch_tool_name_normalization_combined() {
let reg = edgecrab_tools::ToolRegistry::new();
assert_eq!(
reg.resolve_tool_call_name("apply patch<|channel|>action")
.canonical,
"apply_patch"
);
}
#[test]
fn dispatch_tool_name_normalization_trims_whitespace() {
let reg = edgecrab_tools::ToolRegistry::new();
assert_eq!(
reg.resolve_tool_call_name(" write_file ").canonical,
"write_file"
);
}
#[test]
fn dispatch_tool_name_only_token_yields_empty() {
let reg = edgecrab_tools::ToolRegistry::new();
assert_eq!(
reg.resolve_tool_call_name("<|channel|>commentary")
.canonical,
""
);
}
}