mod tool_types;
mod turn;
pub use tool_types::*;
use crate::compaction::{
apply_compaction, apply_compaction_result, compact_messages_semantically, estimate_messages,
CompactionConfig,
};
use crate::cost::{PricingRegistry, SessionCost, TokenUsage};
use crate::guardrails::{
plan_tool_effect_batches, GuardrailConfig, GuardrailDecision, SelfHealingRetry, ToolGuardrails,
};
use crate::hooks::HookRegistry;
use crate::mode::{self, Profile, Scope};
use crate::models::ModelRegistry;
use crate::permissions::{
Approver, AsyncApprover, Authorizer, Decision, PlanApprover, PlanDecision, PlanProposal,
Policy, PolicyAuthorizer,
};
use crate::provider::{Message, Provider, Role};
use crate::todo::{TodoConfig, TodoState};
use moka::future::Cache;
use parking_lot::RwLock;
use serde::Serialize;
use sha2::{Digest, Sha256};
use std::sync::Arc;
use std::time::Instant;
#[cfg(feature = "providers")]
use tracing::error;
use tracing::{debug, info, warn};
#[derive(Debug, Clone, Serialize)]
#[serde(tag = "type")]
pub enum Event {
AgentStart,
ContextUsage {
used_tokens: usize,
context_window: usize,
auto_compact_at: usize,
},
Usage {
model: String,
usage: TokenUsage,
estimated: bool,
},
CompactionStart {
reason: String,
before_tokens: usize,
},
CompactionEnd {
reason: String,
result: crate::compaction::CompactionResult,
},
SkillActivated {
id: String,
name: String,
},
ToolSource {
tool: String,
source: ToolSource,
},
TurnStart {
turn: usize,
},
MessageStart {
role: Role,
},
MessageDelta {
delta: String,
},
MessageEnd {
role: Role,
content: String,
},
ToolCall(ToolCall),
ApprovalRequired(crate::permissions::ApprovalRequest),
PlanProposed(crate::permissions::PlanProposal),
PlanDecided {
decision: crate::permissions::PlanDecision,
},
GuardrailWarning {
tool: String,
reason: String,
},
GuardrailStop {
tool: String,
reason: String,
},
SelfHealing {
attempt: u8,
max_attempts: u8,
errors: Vec<String>,
},
ToolExecutionStart(ToolCall),
ToolExecutionEnd(ToolResult),
TodoUpdated {
todos: TodoState,
},
TurnEnd {
turn: usize,
},
TurnEnded {
turn: usize,
metadata: TurnEndMetadata,
},
CacheAudit(CacheAudit),
GateResult(GateResult),
MemoryRecalled {
recalls: Vec<MemoryRecall>,
},
AgentEnd,
Error(String),
BudgetExceeded {
reason: String,
},
}
#[derive(Debug, Clone, Serialize)]
pub struct TurnEndMetadata {
pub open_todos_remain: bool,
pub finish_reason: Option<String>,
pub trailing_tool_intent: bool,
pub final_message_mid_thought: bool,
pub turn_complete: bool,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum CacheDivergence {
SystemPrompt,
Tools,
Message { index: usize },
}
#[derive(Debug, Clone, Serialize)]
pub struct CacheAudit {
pub stable_prefix_bytes: usize,
pub total_prompt_bytes: usize,
pub first_divergence: Option<CacheDivergence>,
}
#[derive(Debug, Clone, Serialize)]
pub struct QualityGateConfig {
pub command: String,
pub max_output_bytes: usize,
}
impl QualityGateConfig {
pub fn new(command: impl Into<String>) -> Self {
Self {
command: command.into(),
max_output_bytes: 10 * 1024,
}
}
}
#[derive(Debug, Clone, Serialize)]
pub struct GateResult {
pub command: String,
pub success: bool,
pub exit_code: Option<i32>,
pub output: String,
pub skipped_unchanged: bool,
}
pub trait SemanticEmbedder: Send + Sync {
fn embed(&self, text: &str) -> Option<Vec<f32>>;
}
#[derive(Clone)]
pub struct SemanticRecallConfig {
pub embedder: Arc<dyn SemanticEmbedder>,
pub top_k: usize,
pub threshold: f32,
}
#[derive(Debug, Clone, Serialize)]
pub struct MemoryRecall {
pub id: String,
pub summary: String,
pub similarity: f32,
}
#[derive(Clone)]
struct PromptFingerprint {
system: Option<String>,
tools: Vec<serde_json::Value>,
messages: Vec<Message>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ToolSource {
Builtin,
Mcp { server: String },
ComputerUse,
}
pub type Subscriber = Arc<dyn Fn(&Event) + Send + Sync>;
#[derive(Debug, Clone, Default, Serialize)]
pub struct AgentBudget {
pub max_cost: Option<f64>,
pub max_duration_seconds: Option<u64>,
pub reserve_budget: Option<f64>,
pub reserve_budget_fraction: Option<f64>,
}
impl AgentBudget {
pub fn effective_max_cost(&self) -> Option<f64> {
let max = self.max_cost?;
let reserve = match (self.reserve_budget, self.reserve_budget_fraction) {
(Some(usd), Some(frac)) => usd + (max * frac),
(Some(usd), None) => usd,
(None, Some(frac)) => max * frac,
(None, None) => 0.0,
};
Some((max - reserve).max(0.0))
}
pub fn exceeded(&self, start: Option<Instant>, total_cost: f64) -> Option<String> {
if let Some(max_dur) = self.max_duration_seconds {
if let Some(start) = start {
let elapsed = start.elapsed().as_secs();
if elapsed >= max_dur {
return Some(format!("time budget exceeded: {elapsed}s >= {max_dur}s"));
}
}
}
if let Some(max) = self.effective_max_cost() {
if total_cost >= max {
return Some(format!(
"cost budget exceeded: ${total_cost:.4} >= ${max:.4}"
));
}
}
None
}
}
pub struct Agent {
pub model: String,
pub model_registry: ModelRegistry,
pub reasoning_effort: Option<String>,
pub system_prompt: Option<String>,
base_system_prompt: Option<String>,
pub tools: Arc<ToolRegistry>,
pub policy: Policy,
pub scope: Scope,
scope_profile: Option<Profile>,
pub hooks: Option<HookRegistry>,
pub approver: Option<Arc<dyn Approver>>,
pub async_approver: Option<Arc<dyn AsyncApprover>>,
pub authorizer: Option<Arc<dyn Authorizer>>,
pub plan_approver: Option<Arc<dyn PlanApprover>>,
pub guardrails: Option<GuardrailConfig>,
pub self_healing: Option<SelfHealingRetry>,
pub provider: Option<Arc<dyn Provider>>,
pub max_tool_iterations: usize,
pub auto_compact_after: usize,
pub todo_config: Option<TodoConfig>,
pub todo_state: Arc<RwLock<TodoState>>,
pub quality_gate: Option<QualityGateConfig>,
#[cfg(feature = "graph-memory")]
pub semantic_recall: Option<SemanticRecallConfig>,
pub workspace_root: std::path::PathBuf,
#[cfg(feature = "autoresearch")]
pub autoresearch_controller:
Option<crate::autoresearch_controller::AutoresearchControllerHandle>,
extra_allowed_tools: Vec<String>,
pub sandbox: Option<Arc<crate::sandbox::SandboxManager>>,
pub os_sandbox: Option<Arc<crate::sandbox::OsSandboxRunner>>,
#[cfg(feature = "ipc")]
lsp: Arc<crate::lsp::LspManager>,
os_sandbox_failed: bool,
#[cfg(feature = "skills")]
pub skill_registry: Option<crate::skill_engine::SkillRegistry>,
#[cfg(feature = "skills")]
pub skill_engine: Option<crate::skill_engine::SkillEngine>,
#[cfg(feature = "graph-memory")]
pub graph_memory: Option<crate::graph_memory::GraphMemory>,
#[cfg(feature = "graph-memory")]
pub auto_dream: bool,
#[cfg(feature = "zkr-memory")]
pub self_improve: Option<crate::self_improve::SelfImprove>,
#[cfg(feature = "personality")]
pub personality: Option<crate::personality::Personality>,
turn_cancellation: CancellationHandle,
subscribers: Vec<Subscriber>,
pub messages: Arc<RwLock<Vec<Message>>>,
tool_cache: Cache<String, ToolResult>,
pub budget: Option<AgentBudget>,
pub pricing_registry: PricingRegistry,
pub cache_stats: crate::prompt_cache::CacheStatsTracker,
session_cost: SessionCost,
budget_start: Option<Instant>,
cache_audit_enabled: bool,
previous_prompt_fingerprint: Option<PromptFingerprint>,
last_gate_workspace_hash: Option<String>,
}
impl Agent {
pub fn new() -> Self {
let mut agent = Self {
model: "gpt-4o".into(),
model_registry: ModelRegistry::new(),
reasoning_effort: None,
system_prompt: None,
base_system_prompt: None,
tools: Arc::new(ToolRegistry::new()),
policy: Policy::workspace_write(),
scope: Scope::Coding,
scope_profile: None,
hooks: None,
approver: None,
async_approver: None,
plan_approver: None,
guardrails: None,
self_healing: None,
authorizer: None,
provider: None,
max_tool_iterations: 50,
auto_compact_after: 0,
todo_config: None,
todo_state: Arc::new(RwLock::new(TodoState::default())),
quality_gate: None,
#[cfg(feature = "graph-memory")]
semantic_recall: None,
workspace_root: std::env::current_dir().unwrap_or_else(|_| ".".into()),
#[cfg(feature = "autoresearch")]
autoresearch_controller: None,
extra_allowed_tools: Vec::new(),
sandbox: None,
os_sandbox: None,
#[cfg(feature = "ipc")]
lsp: Arc::new(crate::lsp::LspManager::new()),
os_sandbox_failed: false,
#[cfg(feature = "skills")]
skill_registry: None,
#[cfg(feature = "skills")]
skill_engine: None,
#[cfg(feature = "graph-memory")]
graph_memory: None,
#[cfg(feature = "graph-memory")]
auto_dream: false,
#[cfg(feature = "zkr-memory")]
self_improve: None,
#[cfg(feature = "personality")]
personality: None,
turn_cancellation: CancellationHandle::new(),
subscribers: Vec::new(),
messages: Arc::new(RwLock::new(Vec::new())),
tool_cache: Cache::builder()
.max_capacity(10_000)
.time_to_live(std::time::Duration::from_secs(3600))
.time_to_idle(std::time::Duration::from_secs(900))
.build(),
budget: None,
pricing_registry: PricingRegistry::new(),
cache_stats: crate::prompt_cache::CacheStatsTracker::new(),
session_cost: SessionCost::new(),
budget_start: None,
cache_audit_enabled: false,
previous_prompt_fingerprint: None,
last_gate_workspace_hash: None,
};
agent.ensure_userspace_sandbox();
if agent.policy.enable_os_sandbox {
if let Err(e) = agent.enable_os_sandbox() {
agent.os_sandbox_failed = true;
tracing::warn!("OS sandbox unavailable — shell tools will be blocked: {e}");
}
}
agent
}
pub fn set_model(&mut self, model: impl Into<String>) {
self.model = model.into();
}
pub fn set_model_registry(&mut self, registry: ModelRegistry) {
self.model_registry = registry;
}
pub fn model_registry(&self) -> &ModelRegistry {
&self.model_registry
}
pub fn set_reasoning_effort(&mut self, effort: Option<String>) {
self.reasoning_effort = effort;
}
pub fn set_system_prompt(&mut self, prompt: impl Into<String>) {
self.base_system_prompt = Some(prompt.into());
self.refresh_system_prompt();
}
pub fn set_tools(&mut self, tools: ToolRegistry) {
self.tools = Arc::new(tools);
}
#[cfg(feature = "ipc")]
pub fn set_lsp_manager(&mut self, lsp: Arc<crate::lsp::LspManager>) {
self.lsp = lsp;
}
pub fn set_policy(&mut self, policy: Policy) {
self.policy = policy;
self.authorizer = None;
self.ensure_userspace_sandbox();
if self.policy.enable_os_sandbox && self.os_sandbox.is_none() && !self.os_sandbox_failed {
if let Err(e) = self.enable_os_sandbox() {
self.os_sandbox_failed = true;
tracing::warn!("OS sandbox unavailable — shell tools will be blocked: {e}");
}
}
}
pub fn set_scope(&mut self, scope: Scope) {
self.scope = scope;
let profile = mode::profile(scope);
self.policy.apply_scope(&profile.policy);
self.authorizer = None;
self.ensure_userspace_sandbox();
if self.policy.enable_os_sandbox && self.os_sandbox.is_none() && !self.os_sandbox_failed {
if let Err(e) = self.enable_os_sandbox() {
self.os_sandbox_failed = true;
tracing::warn!("OS sandbox unavailable — shell tools will be blocked: {e}");
}
}
self.scope_profile = Some(profile);
self.refresh_system_prompt();
}
pub fn set_hooks(&mut self, hooks: HookRegistry) {
self.hooks = Some(hooks);
}
pub fn set_approver(&mut self, approver: Arc<dyn Approver>) {
self.approver = Some(approver);
}
pub fn set_async_approver(&mut self, approver: Arc<dyn AsyncApprover>) {
self.async_approver = Some(approver);
}
pub fn clear_async_approver(&mut self) {
self.async_approver = None;
}
pub fn set_authorizer(&mut self, authorizer: Arc<dyn Authorizer>) {
self.authorizer = Some(authorizer);
}
pub fn set_plan_approver(&mut self, approver: Arc<dyn PlanApprover>) {
self.plan_approver = Some(approver);
}
pub fn clear_plan_approver(&mut self) {
self.plan_approver = None;
}
pub fn set_guardrails(&mut self, config: GuardrailConfig) {
self.guardrails = Some(config);
}
pub fn clear_guardrails(&mut self) {
self.guardrails = None;
}
pub fn set_self_healing(&mut self, max_attempts: u8) {
self.self_healing = Some(SelfHealingRetry::new(max_attempts));
}
pub fn clear_self_healing(&mut self) {
self.self_healing = None;
}
pub fn clear_authorizer(&mut self) {
self.authorizer = None;
}
pub fn set_provider(&mut self, provider: Arc<dyn Provider>) {
self.provider = Some(provider);
}
pub fn set_todo_config(&mut self, config: TodoConfig) {
self.todo_config = Some(config);
}
pub fn clear_todo_config(&mut self) {
self.todo_config = None;
}
pub fn set_todo_state(&mut self, state: TodoState) {
*self.todo_state.write() = state;
}
pub fn todos(&self) -> TodoState {
self.todo_state.read().clone()
}
pub fn enable_cache_audit(&mut self, enabled: bool) {
self.cache_audit_enabled = enabled;
if !enabled {
self.previous_prompt_fingerprint = None;
}
}
pub fn set_quality_gate(&mut self, config: QualityGateConfig) {
self.quality_gate = Some(config);
self.last_gate_workspace_hash = None;
}
pub fn clear_quality_gate(&mut self) {
self.quality_gate = None;
self.last_gate_workspace_hash = None;
}
#[cfg(feature = "graph-memory")]
pub fn set_semantic_recall(&mut self, config: SemanticRecallConfig) {
self.semantic_recall = Some(config);
}
pub fn set_workspace_root(&mut self, path: impl Into<std::path::PathBuf>) {
let new_root = path.into();
let mut sandbox_config = self.sandbox.as_ref().map(|sb| sb.config());
if let Some(config) = sandbox_config.as_mut() {
config.workspace_root = new_root.clone();
}
let mut os_config = self.os_sandbox.as_ref().map(|os| os.config().clone());
if let Some(config) = os_config.as_mut() {
config.workspace = new_root.clone();
}
self.workspace_root = new_root;
self.authorizer = None;
self.sandbox = Some(Arc::new(match sandbox_config {
Some(config) => crate::sandbox::SandboxManager::from_config(config),
None => {
let mut sb = crate::sandbox::SandboxManager::new(
crate::sandbox::SandboxProfile::Workspace,
self.workspace_root.clone(),
);
sb.set_allow_network(true);
sb
}
}));
self.tool_cache.invalidate_all();
self.os_sandbox = None;
self.os_sandbox_failed = false;
if self.policy.enable_os_sandbox {
let result = match os_config {
Some(config) => crate::sandbox::OsSandboxRunner::new(config)
.map(Arc::new)
.map(|runner| self.os_sandbox = Some(runner)),
None => self.enable_os_sandbox().map(|_| ()),
};
if let Err(e) = result {
self.os_sandbox_failed = true;
tracing::warn!("OS sandbox unavailable after workspace change — shell tools will be blocked: {e}");
}
}
}
#[cfg(feature = "autoresearch")]
pub fn set_autoresearch_controller(
&mut self,
controller: crate::autoresearch_controller::AutoresearchControllerHandle,
) {
self.autoresearch_controller = Some(controller);
}
#[cfg(feature = "autoresearch")]
pub fn autoresearch_controller(
&self,
) -> Option<crate::autoresearch_controller::AutoresearchControllerHandle> {
self.autoresearch_controller.clone()
}
#[cfg(feature = "autoresearch")]
pub fn clear_autoresearch_controller(&mut self) {
self.autoresearch_controller = None;
}
pub fn allow_extra_tools<I, S>(&mut self, names: I)
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
self.extra_allowed_tools
.extend(names.into_iter().map(Into::into));
}
pub fn set_extra_allowed_tools<I, S>(&mut self, names: I)
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
self.extra_allowed_tools = names.into_iter().map(Into::into).collect();
}
pub fn extra_allowed_tools(&self) -> &[String] {
&self.extra_allowed_tools
}
#[cfg(feature = "zkr-memory")]
pub fn set_self_improve(&mut self, improve: crate::self_improve::SelfImprove) {
self.self_improve = Some(improve);
}
#[cfg(feature = "personality")]
pub fn set_personality(&mut self, personality: crate::personality::Personality) {
self.personality = Some(personality);
}
pub fn cancel(&self) {
self.turn_cancellation.cancel();
}
pub fn cancellation_handle(&self) -> CancellationHandle {
self.turn_cancellation.clone()
}
pub fn load_project_context(&mut self) {
if let Some(instr) = crate::context::load_project_instructions(&self.workspace_root) {
self.base_system_prompt = crate::context::compose_system_prompt(
self.base_system_prompt.as_deref(),
&instr.content,
);
self.refresh_system_prompt();
}
}
fn refresh_system_prompt(&mut self) {
self.system_prompt = self.scope_profile.as_ref().map_or_else(
|| self.base_system_prompt.clone(),
|profile| {
Some(mode::compose_prompt(
self.base_system_prompt.as_deref(),
profile,
))
},
);
}
pub fn set_sandbox(&mut self, sb: Arc<crate::sandbox::SandboxManager>) {
self.sandbox = Some(sb);
}
pub fn set_os_sandbox(&mut self, os: Arc<crate::sandbox::OsSandboxRunner>) {
self.os_sandbox = Some(os);
}
pub fn set_budget(&mut self, budget: AgentBudget) {
self.budget = Some(budget);
}
pub fn set_pricing_registry(&mut self, registry: PricingRegistry) {
self.pricing_registry = registry;
}
pub fn total_cost(&self) -> f64 {
self.session_cost.total_cost()
}
pub fn session_cost(&self) -> &SessionCost {
&self.session_cost
}
pub fn cache_stats(&self) -> crate::prompt_cache::CacheStats {
self.cache_stats.stats()
}
fn check_budget(&self) -> Option<String> {
self.budget
.as_ref()
.and_then(|b| b.exceeded(self.budget_start, self.session_cost.total_cost()))
}
pub fn ensure_userspace_sandbox(&mut self) {
if self.sandbox.is_none() {
let mut sb = crate::sandbox::SandboxManager::new(
crate::sandbox::SandboxProfile::Workspace,
self.workspace_root.clone(),
);
sb.set_allow_network(true);
self.sandbox = Some(Arc::new(sb));
}
}
pub fn enable_os_sandbox(&mut self) -> Result<(), crate::sandbox::SandboxError> {
self.ensure_userspace_sandbox();
let mode = crate::sandbox::detect_sandbox();
if matches!(mode, crate::sandbox::OsSandbox::UserspaceOnly) {
return Err(crate::sandbox::SandboxError::PathDenied(
"no seatbelt/bwrap on this host".into(),
));
}
let config = crate::sandbox::OsSandboxConfig::new(mode, self.workspace_root.clone());
let runner = crate::sandbox::OsSandboxRunner::new(config)?;
self.os_sandbox = Some(Arc::new(runner));
self.policy.enable_os_sandbox = true;
Ok(())
}
#[cfg(feature = "skills")]
pub fn set_skill_registry(&mut self, registry: crate::skill_engine::SkillRegistry) {
self.skill_registry = Some(registry);
}
#[cfg(feature = "skills")]
pub fn set_skill_engine(&mut self, engine: crate::skill_engine::SkillEngine) {
self.skill_engine = Some(engine);
}
#[cfg(feature = "graph-memory")]
pub fn set_graph_memory(&mut self, graph: crate::graph_memory::GraphMemory) {
self.graph_memory = Some(graph);
}
#[cfg(feature = "graph-memory")]
pub fn enable_auto_dream(&mut self, enabled: bool) {
self.auto_dream = enabled;
}
pub fn subscribe(&mut self, callback: impl Fn(&Event) + Send + Sync + 'static) {
self.subscribers.push(Arc::new(callback));
}
fn emit(&self, event: Event) {
if self.subscribers.is_empty() {
return;
}
for sub in &self.subscribers {
sub(&event);
}
}
pub fn clear_messages(&self) {
self.messages.write().clear();
}
pub fn messages_handle(&self) -> Arc<RwLock<Vec<Message>>> {
Arc::clone(&self.messages)
}
pub fn message_count(&self) -> usize {
self.messages.read().len()
}
pub fn context_window(&self) -> usize {
let model = self
.provider
.as_ref()
.and_then(|provider| {
self.model_registry
.get_for_provider(provider.id(), &self.model)
})
.or_else(|| self.model_registry.get(&self.model));
model
.map(|model| model.context_window)
.filter(|window| *window > 0)
.unwrap_or(CompactionConfig::DEFAULT_CONTEXT_WINDOW)
}
pub fn auto_compact_threshold(&self) -> usize {
if self.auto_compact_after == 0 {
let context_window = self.context_window();
context_window.saturating_sub(context_window / 10)
} else {
self.auto_compact_after
}
}
pub fn context_tokens(&self) -> usize {
estimate_messages(&self.messages.read())
+ self
.system_prompt
.as_deref()
.map(crate::compaction::estimate_tokens)
.unwrap_or(0)
}
fn compaction_config(&self) -> CompactionConfig {
if self.auto_compact_after == 0 {
let context_window = self.context_window();
let reserve = context_window / 10;
CompactionConfig::new(context_window, reserve, reserve)
} else {
let reserve = (self.auto_compact_after / 4).max(32);
CompactionConfig::new(self.auto_compact_after + reserve, reserve, reserve)
}
}
pub async fn prompt(&mut self, text: &str) -> Result<(), AgentError> {
let provider = self.provider.clone().ok_or(AgentError::NoProvider)?;
let tokens = self.context_tokens();
let context_window = self.context_window();
let auto_compact_at = self.auto_compact_threshold();
self.emit(Event::ContextUsage {
used_tokens: tokens,
context_window,
auto_compact_at,
});
if tokens >= auto_compact_at {
if let Err(error) = self
.compact_semantically("auto-compact before prompt", provider.as_ref())
.await
{
warn!("automatic context compaction failed: {error}");
}
}
let redactor = crate::secrets::Redactor::new();
let safe_text = redactor.redact(text);
let active_skills = self.activate_skills_for_prompt(&safe_text);
self.messages.write().push(Message::user(safe_text.clone()));
self.emit(Event::AgentStart);
self.budget_start = Some(Instant::now());
self.before_prompt_hooks(&safe_text).await;
let mut tool_ctx = self.tool_context();
tool_ctx.provider = Some(provider.clone());
tool_ctx.tools = Some(Arc::clone(&self.tools));
let pending_scope = Arc::new(parking_lot::Mutex::new(None));
tool_ctx.pending_scope = Some(Arc::clone(&pending_scope));
tool_ctx.todo_state = Some(Arc::clone(&self.todo_state));
tool_ctx.todo_config = self.todo_config.clone();
tool_ctx.todo_updates = self
.todo_config
.as_ref()
.map(|_| Arc::new(parking_lot::Mutex::new(Vec::new())));
let ctx = Arc::new(tool_ctx);
#[cfg(feature = "zkr-memory")]
let mut tool_error_seen = false;
let mut guardrails = self.guardrails.clone().map(ToolGuardrails::new);
let mut self_healing = self.self_healing.clone();
let mut plan_approved = false;
let last_finish_reason: Option<String> = None;
for iteration in 0..self.max_tool_iterations {
if let Some(reason) = self.check_budget() {
self.emit(Event::BudgetExceeded {
reason: reason.clone(),
});
return Err(AgentError::BudgetExceeded(reason));
}
self.emit(Event::TurnStart { turn: iteration });
let messages: Vec<Message> = self.messages.read().clone();
let base_system =
turn::append_active_skills(self.system_prompt.clone(), active_skills.as_deref());
#[cfg(feature = "graph-memory")]
let base_system = self.append_semantic_recalls(base_system, &safe_text);
#[cfg(feature = "zkr-memory")]
let system = if let Some(improve) = &self.self_improve {
let base = base_system.as_deref().unwrap_or("");
match improve.augment(&safe_text, base).await {
Ok(augmented) => Some(augmented),
Err(error) => {
warn!("self-improve augmentation failed: {error}");
base_system.clone()
}
}
} else {
base_system
};
#[cfg(not(feature = "zkr-memory"))]
let system = base_system;
#[cfg(feature = "personality")]
let system = if let Some(pers) = &self.personality {
let base = system.as_deref().unwrap_or("");
match pers.augment(&safe_text, base).await {
Ok(augmented) => Some(augmented),
Err(error) => {
warn!("personality augmentation failed: {error}");
system
}
}
} else {
system
};
let tool_definitions = self.tools.definitions();
self.audit_prompt(&messages, &system, &tool_definitions);
#[cfg_attr(not(feature = "providers"), allow(unused_mut))]
let mut tool_calls: Vec<ToolCall> = Vec::new();
let mut assistant_content;
#[cfg_attr(not(feature = "providers"), allow(unused_mut))]
let mut provider_usage: Option<TokenUsage> = None;
self.emit(Event::MessageStart {
role: Role::Assistant,
});
#[cfg(feature = "providers")]
{
assistant_content = String::new();
use crate::provider::StreamEvent;
use futures::StreamExt;
let mut attempts = 0;
let reasoning_effort = self.reasoning_effort.as_deref().filter(|_| {
self.model_registry
.supports_reasoning_effort_for(provider.id(), &self.model)
});
let stream = loop {
let result = ctx
.cancellation
.run(provider.stream(
&messages,
&system,
&self.model,
&tool_definitions,
reasoning_effort,
))
.await
.map_err(|_| AgentError::Cancelled)?;
match result {
Ok(stream) => break stream,
Err(e) if e.is_transient() && attempts < 2 => {
attempts += 1;
ctx.cancellation
.run(tokio::time::sleep(std::time::Duration::from_millis(
250 * (1 << attempts),
)))
.await
.map_err(|_| AgentError::Cancelled)?;
}
Err(e) => {
error!("provider stream error: {e}");
self.emit(Event::Error(e.to_string()));
return Err(AgentError::Provider(e.to_string()));
}
}
};
let mut stream = stream;
loop {
let next = ctx
.cancellation
.run(stream.next())
.await
.map_err(|_| AgentError::Cancelled)?;
let Some(event_result) = next else {
break;
};
match event_result {
Ok(StreamEvent::Delta(delta)) => {
assistant_content.push_str(&delta);
}
Ok(StreamEvent::ToolCall(call)) => {
tool_calls.push(call.clone());
self.emit(Event::ToolSource {
tool: call.name.clone(),
source: tool_source(&call.name),
});
self.emit(Event::ToolCall(redact_tool_call(&call)));
}
Ok(StreamEvent::Usage(usage)) => {
provider_usage = Some(usage);
}
Ok(StreamEvent::Done) => break,
Err(e) => {
error!("stream error: {e}");
self.emit(Event::Error(e.to_string()));
return Err(AgentError::Provider(e.to_string()));
}
}
}
}
#[cfg(not(feature = "providers"))]
{
let _ = (&provider, &messages, &system);
assistant_content =
"[providers feature not enabled — enable with --features providers]"
.to_string();
}
let redacted_assistant = redactor.redact(&assistant_content);
if !redacted_assistant.is_empty() {
self.emit(Event::MessageDelta {
delta: redacted_assistant.clone(),
});
}
assistant_content = redacted_assistant;
self.emit(Event::MessageEnd {
role: Role::Assistant,
content: assistant_content.clone(),
});
self.messages.write().push(Message::assistant_with_tools(
assistant_content.clone(),
tool_calls.clone(),
));
let input_tokens = estimate_messages(&messages)
+ system
.as_deref()
.map(crate::compaction::estimate_tokens)
.unwrap_or(0);
let output_tokens = crate::compaction::estimate_tokens(&assistant_content);
let estimated = provider_usage.is_none();
let usage = provider_usage.unwrap_or(TokenUsage {
input_tokens,
output_tokens,
cache_read_tokens: 0,
cache_write_tokens: 0,
});
if !estimated {
self.cache_stats.record_tokens(usage);
}
self.session_cost
.record(&self.model, usage, &self.pricing_registry);
self.emit(Event::Usage {
model: self.model.clone(),
usage,
estimated,
});
self.emit(Event::ContextUsage {
used_tokens: self.context_tokens(),
context_window,
auto_compact_at,
});
if let Some(reason) = self.check_budget() {
self.emit(Event::BudgetExceeded {
reason: reason.clone(),
});
return Err(AgentError::BudgetExceeded(reason));
}
if tool_calls.is_empty() {
if let Some(result) = self.run_quality_gate().await {
let failed = !result.success && !result.skipped_unchanged;
let output = result.output.clone();
self.emit(Event::GateResult(result));
if failed {
self.messages.write().push(Message::user(format!(
"Quality gate failed. Fix the failure, then continue. Bounded gate output:\n{output}"
)));
continue;
}
}
self.emit_turn_end(
iteration,
last_finish_reason.as_deref(),
false,
&assistant_content,
);
#[cfg(feature = "zkr-memory")]
if let Some(improve) = &self.self_improve {
let outcome = if tool_error_seen { "error" } else { "success" };
let lesson = if tool_error_seen {
"avoid repeating the failing tool"
} else {
"continue the current strategy"
};
if let Err(error) = improve
.record(&safe_text, &assistant_content, outcome, lesson)
.await
{
warn!("self-improve reflection failed: {error}");
}
}
#[cfg(feature = "personality")]
if let Some(pers) = &self.personality {
let epoch = (iteration + 1) as u64;
let assistant_event = crate::personality::ConversationEvent {
epoch,
participant: "agent".to_string(),
event_kind: if tool_error_seen { "error" } else { "message" }.to_string(),
content: assistant_content.chars().take(500).collect(),
};
if let Err(error) = pers.record_event(&assistant_event).await {
warn!("personality assistant event recording failed: {error}");
}
let risk = pers.assess_risk("user", &assistant_content).await;
if risk.recommendation == crate::personality::RiskRecommendation::Abort {
warn!(
"personality risk assessment: ABORT (overall {}bps) — {:?}",
risk.overall_risk_basis_points, risk
);
} else if risk.recommendation == crate::personality::RiskRecommendation::Refine
{
debug!(
"personality risk assessment: REFINE (overall {}bps)",
risk.overall_risk_basis_points
);
}
let hyp = crate::personality::MindHypothesis {
participant: "user".to_string(),
belief: format!(
"user sent: {}",
safe_text.chars().take(100).collect::<String>()
),
emotion: if tool_error_seen {
Some("frustrated".into())
} else {
None
},
goal: None,
predicted_reaction: Some(
if tool_error_seen {
"likely frustrated by errors"
} else {
"likely satisfied with response"
}
.into(),
),
confidence_basis_points: if tool_error_seen { 4000 } else { 7000 },
valid_until: None,
};
if let Err(error) = pers.record_hypothesis(&hyp).await {
warn!("personality ToM recording failed: {error}");
}
}
break;
}
if !plan_approved {
if let Some(gate) = self.plan_approver.clone() {
let proposal = PlanProposal {
prompt: safe_text.clone(),
plan: assistant_content.clone(),
calls: tool_calls.iter().map(redact_tool_call).collect(),
turn: iteration,
};
self.emit(Event::PlanProposed(proposal.clone()));
let decision = gate.approve_plan(&proposal).await;
self.emit(Event::PlanDecided {
decision: decision.clone(),
});
match decision {
PlanDecision::Approve => plan_approved = true,
PlanDecision::Reject(reason) => {
info!("plan rejected: {reason}");
self.messages.write().push(Message::user(format!(
"The plan was rejected: {reason}. Do not run it."
)));
self.emit_turn_end(
iteration,
last_finish_reason.as_deref(),
true,
&assistant_content,
);
break;
}
PlanDecision::Revise(guidance) => {
info!("plan revision requested: {guidance}");
self.messages.write().push(Message::user(format!(
"Do not run that plan. Revise it: {guidance}"
)));
self.emit_turn_end(
iteration,
last_finish_reason.as_deref(),
true,
&assistant_content,
);
continue;
}
}
} else {
plan_approved = true;
}
}
let results = self.execute_tools_parallel(&tool_calls, &ctx).await;
if let Some(updates) = &ctx.todo_updates {
for update in std::mem::take(&mut *updates.lock()) {
self.emit(Event::TodoUpdated { todos: update });
}
}
for result in &results {
#[cfg(feature = "zkr-memory")]
{
tool_error_seen |= result.is_error;
}
self.messages
.write()
.push(Message::tool(&result.id, &result.content));
}
if let Some(scope) = pending_scope.lock().take() {
self.set_scope(scope);
}
let mut stopped: Option<(String, String)> = None;
if let Some(rails) = guardrails.as_mut() {
for (call, result) in tool_calls.iter().zip(results.iter()) {
match rails.observe(&call.name, &call.arguments, result.is_error) {
GuardrailDecision::Proceed => {}
GuardrailDecision::Warn(reason) => {
warn!("guardrail warning on '{}': {reason}", call.name);
self.emit(Event::GuardrailWarning {
tool: call.name.clone(),
reason: reason.clone(),
});
self.messages
.write()
.push(Message::user(format!("Guardrail warning: {reason}")));
}
GuardrailDecision::Stop(reason) => {
warn!("guardrail stop on '{}': {reason}", call.name);
stopped = Some((call.name.clone(), reason));
break;
}
}
}
}
if let Some((tool, reason)) = stopped {
self.emit(Event::GuardrailStop {
tool,
reason: reason.clone(),
});
self.messages
.write()
.push(Message::user(format!("Stopped by guardrail: {reason}")));
self.emit_turn_end(
iteration,
last_finish_reason.as_deref(),
true,
&assistant_content,
);
break;
}
let errors: Vec<String> = results
.iter()
.filter(|r| r.is_error)
.map(|r| r.content.clone())
.collect();
if !errors.is_empty() {
if let Some(healer) = self_healing.as_mut() {
if healer.should_retry() {
let message = healer.build_healing_message(&errors);
self.emit(Event::SelfHealing {
attempt: healer.attempts_used,
max_attempts: healer.max_attempts,
errors,
});
self.messages.write().push(Message::user(message));
}
}
}
self.emit_turn_end(
iteration,
last_finish_reason.as_deref(),
true,
&assistant_content,
);
}
self.after_prompt_hooks().await;
self.emit(Event::AgentEnd);
Ok(())
}
async fn execute_tools_parallel(
&self,
calls: &[ToolCall],
ctx: &Arc<ToolContext>,
) -> Vec<ToolResult> {
let effects: Vec<ToolEffect> = calls
.iter()
.map(|c| {
let name = normalize_tool_name(&c.name);
self.tools.effect_of(name)
})
.collect();
let batches = plan_tool_effect_batches(&effects);
let mut results: Vec<Option<ToolResult>> = vec![None; calls.len()];
let mut join_failures: Vec<Option<String>> = vec![None; calls.len()];
for batch in batches {
if batch.len() == 1 {
let idx = batch[0];
let original = &calls[idx];
self.emit(Event::ToolExecutionStart(redact_tool_call(original)));
let (call, result) = self.execute_single_tool(original, ctx).await;
if result.requires_approval() {
self.emit(Event::ApprovalRequired(
crate::permissions::ApprovalRequest::from_call(
&redact_tool_call(&call),
&self.policy,
),
));
}
self.emit(Event::ToolExecutionEnd(result.clone()));
results[idx] = Some(result);
continue;
}
let tools = Arc::clone(&self.tools);
let policy = self.policy.clone();
let scope_profile = self.scope_profile.clone();
let approver = self.approver.clone();
let async_approver = self.async_approver.clone();
let authorizer = self.authorizer.clone();
let tool_cache = self.tool_cache.clone();
let extra_allowed_tools = self.extra_allowed_tools.clone();
let mut join_set = tokio::task::JoinSet::new();
for idx in batch {
let original = &calls[idx];
let call = match self.apply_before_tool_hooks(original) {
Ok(c) => c,
Err(reason) => {
self.emit(Event::ToolExecutionStart(redact_tool_call(original)));
let result = ToolResult::err(&original.id, reason);
self.emit(Event::ToolExecutionEnd(result.clone()));
results[idx] = Some(result);
continue;
}
};
self.emit(Event::ToolExecutionStart(redact_tool_call(&call)));
let ctx = Arc::clone(ctx);
let tools = Arc::clone(&tools);
let policy = policy.clone();
let scope_profile = scope_profile.clone();
let approver = approver.clone();
let async_approver = async_approver.clone();
let authorizer = authorizer.clone();
let tool_cache = tool_cache.clone();
let extra_allowed_tools = extra_allowed_tools.clone();
join_set.spawn(async move {
let result = Agent::run_tool_call(
&tools,
&policy,
authorizer.as_deref(),
scope_profile.as_ref(),
approver.clone(),
async_approver.as_deref(),
&tool_cache,
&call,
&ctx,
&extra_allowed_tools,
)
.await;
(idx, call, result)
});
}
while let Some(joined) = join_set.join_next().await {
match joined {
Ok((idx, call, result)) => {
if result.requires_approval() {
self.emit(Event::ApprovalRequired(
crate::permissions::ApprovalRequest::from_call(
&redact_tool_call(&call),
&self.policy,
),
));
}
self.emit(Event::ToolExecutionEnd(result.clone()));
results[idx] = Some(result);
}
Err(e) => {
warn!("parallel tool task join error: {e}");
if let Some((idx, _)) = results
.iter()
.enumerate()
.find(|(_, result)| result.is_none())
{
join_failures[idx] = Some(format!("parallel tool task failed: {e}"));
}
}
}
}
}
results
.into_iter()
.enumerate()
.map(|(i, r)| {
r.unwrap_or_else(|| {
ToolResult::err(
calls.get(i).map(|c| c.id.as_str()).unwrap_or(""),
join_failures[i]
.as_deref()
.unwrap_or("tool execution failed"),
)
})
})
.collect()
}
fn apply_before_tool_hooks(&self, call: &ToolCall) -> Result<ToolCall, String> {
match &self.hooks {
Some(hooks) => hooks.run_before_tool(call),
None => Ok(call.clone()),
}
}
async fn execute_single_tool(
&self,
call: &ToolCall,
ctx: &Arc<ToolContext>,
) -> (ToolCall, ToolResult) {
let call = match self.apply_before_tool_hooks(call) {
Ok(c) => c,
Err(reason) => {
let id = call.id.clone();
return (call.clone(), ToolResult::err(&id, reason));
}
};
let result = Self::run_tool_call(
self.tools.as_ref(),
&self.policy,
self.authorizer.as_deref(),
self.scope_profile.as_ref(),
self.approver.clone(),
self.async_approver.as_deref(),
&self.tool_cache,
&call,
ctx,
&self.extra_allowed_tools,
)
.await;
(call, result)
}
#[allow(clippy::too_many_arguments)]
async fn run_tool_call(
tools: &ToolRegistry,
policy: &Policy,
authorizer: Option<&dyn Authorizer>,
scope_profile: Option<&Profile>,
approver: Option<Arc<dyn Approver>>,
async_approver: Option<&dyn AsyncApprover>,
tool_cache: &Cache<String, ToolResult>,
call: &ToolCall,
ctx: &Arc<ToolContext>,
extra_allowed_tools: &[String],
) -> ToolResult {
let resolved_name = normalize_tool_name(&call.name).to_string();
if let Some(profile) = scope_profile {
if !mode::tool_allowed_with_extra(profile, extra_allowed_tools, &call.name)
&& !mode::tool_allowed_with_extra(profile, extra_allowed_tools, &resolved_name)
{
let msg = format!("tool not in scope {}: {}", profile.scope.name(), call.name);
return ToolResult::err(&call.id, msg);
}
}
let mut decision = match authorizer {
Some(auth) => auth.authorize(
policy,
&resolved_name,
&call.arguments,
None,
Some(ctx.workspace_root.as_path()),
),
None => PolicyAuthorizer::new().authorize(
policy,
&resolved_name,
&call.arguments,
None,
Some(ctx.workspace_root.as_path()),
),
};
if decision == Decision::Ask {
let ask_call = ToolCall {
id: call.id.clone(),
name: resolved_name.clone(),
arguments: call.arguments.clone(),
};
if let Some(app) = async_approver {
decision = app.approve(&ask_call).await;
} else if let Some(app) = approver {
decision = tokio::task::spawn_blocking(move || app.approve(&ask_call))
.await
.unwrap_or(Decision::Deny);
}
}
match decision {
Decision::Deny => ToolResult::err(&call.id, "denied by policy"),
Decision::Ask => {
ToolResult::approval_required(&call.id)
}
Decision::Allow => {
let effect = tools.effect_of(&resolved_name);
let cache_key = format!(
"{}:{}:{}",
ctx.workspace_root.display(),
resolved_name,
call.arguments
);
if effect == ToolEffect::Read {
if let Some(cached) = tool_cache.get(&cache_key).await {
debug!("tool cache hit: {}", resolved_name);
return ToolResult::ok(&call.id, cached.content);
}
}
let mut result = match tools.execute(&resolved_name, ctx, &call.arguments).await {
Some(r) => r,
None => ToolResult::err(&call.id, format!("unknown tool: {}", call.name)),
};
result.id = call.id.clone();
result.content = crate::secrets::Redactor::new().redact(&result.content);
if !result.is_error {
match effect {
ToolEffect::Read => {
tool_cache.insert(cache_key, result.clone()).await;
}
ToolEffect::Write | ToolEffect::Process => {
tool_cache.invalidate_all();
}
ToolEffect::Network => {}
}
}
result
}
}
}
pub fn compact(&self, reason: &str) {
info!("compacting context: {reason}");
if self.message_count() <= 2 {
return;
}
let before_tokens = self.context_tokens();
self.emit(Event::CompactionStart {
reason: reason.to_string(),
before_tokens,
});
let result = {
let mut msgs = self.messages.write();
let result = apply_compaction(&mut msgs, &self.compaction_config());
if !result.summary.is_empty() {
msgs.push(Message::system(format!("[compact reason: {reason}]")));
}
result
};
self.emit(Event::CompactionEnd {
reason: reason.to_string(),
result,
});
}
fn tool_context(&self) -> ToolContext {
let mut tool_ctx = ToolContext::new(self.workspace_root.clone());
tool_ctx.os_sandbox_required = self.policy.enable_os_sandbox && self.os_sandbox.is_none();
tool_ctx.cancellation = self.turn_cancellation.reset();
#[cfg(feature = "ipc")]
{
tool_ctx.lsp = Some(Arc::clone(&self.lsp));
}
if let Some(sandbox) = self.sandbox.clone() {
tool_ctx = tool_ctx.with_sandbox(sandbox);
}
if let Some(os_sandbox) = self.os_sandbox.clone() {
tool_ctx = tool_ctx.with_os_sandbox(os_sandbox);
}
tool_ctx
}
fn emit_turn_end(
&self,
turn: usize,
finish_reason: Option<&str>,
trailing_tool_intent: bool,
content: &str,
) {
let trimmed = content.trim_end();
let final_message_mid_thought = !trimmed.is_empty()
&& !trailing_tool_intent
&& !trimmed.ends_with(['.', '!', '?', ':', ';', '`', ')', ']', '}'])
&& !trimmed.ends_with("…");
let open_todos_remain = self
.todo_state
.read()
.items
.iter()
.any(|todo| todo.status != crate::todo::TodoStatus::Completed);
self.emit(Event::TurnEnded {
turn,
metadata: TurnEndMetadata {
open_todos_remain,
finish_reason: finish_reason.map(str::to_owned),
trailing_tool_intent,
final_message_mid_thought,
turn_complete: !trailing_tool_intent && !final_message_mid_thought,
},
});
}
fn audit_prompt(
&mut self,
messages: &[Message],
system: &Option<String>,
tools: &[serde_json::Value],
) {
if !self.cache_audit_enabled {
return;
}
let previous = self.previous_prompt_fingerprint.as_ref();
let first_divergence = previous.and_then(|previous| {
if previous.system != *system {
Some(CacheDivergence::SystemPrompt)
} else if previous.tools != tools {
Some(CacheDivergence::Tools)
} else {
let first = previous
.messages
.iter()
.zip(messages)
.position(|(before, now)| before != now);
first
.or_else(|| {
(previous.messages.len() != messages.len()).then_some(messages.len())
})
.map(|index| CacheDivergence::Message { index })
}
});
let serialized = serde_json::to_vec(&(system, tools, messages)).unwrap_or_default();
let previous_serialized = previous
.and_then(|previous| {
serde_json::to_vec(&(&previous.system, &previous.tools, &previous.messages)).ok()
})
.unwrap_or_default();
let stable_prefix_bytes = serialized
.iter()
.zip(&previous_serialized)
.take_while(|(left, right)| left == right)
.count();
self.emit(Event::CacheAudit(CacheAudit {
stable_prefix_bytes,
total_prompt_bytes: serialized.len(),
first_divergence,
}));
self.previous_prompt_fingerprint = Some(PromptFingerprint {
system: system.clone(),
tools: tools.to_vec(),
messages: messages.to_vec(),
});
}
async fn run_quality_gate(&mut self) -> Option<GateResult> {
let config = self.quality_gate.clone()?;
let hash = workspace_hash(&self.workspace_root);
if self.last_gate_workspace_hash.as_ref() == Some(&hash) {
return Some(GateResult {
command: config.command,
success: true,
exit_code: Some(0),
output: String::new(),
skipped_unchanged: true,
});
}
let root = self.workspace_root.clone();
let command = config.command.clone();
let command_for_error = command.clone();
let cap = config.max_output_bytes;
let result = tokio::task::spawn_blocking(move || {
std::process::Command::new("sh")
.arg("-lc")
.arg(&command)
.current_dir(root)
.output()
.map(|output| {
let mut combined = output.stdout;
combined.extend_from_slice(&output.stderr);
GateResult {
command,
success: output.status.success(),
exit_code: output.status.code(),
output: tail_bytes(&combined, cap),
skipped_unchanged: false,
}
})
.unwrap_or_else(|error| GateResult {
command: command_for_error,
success: false,
exit_code: None,
output: error.to_string(),
skipped_unchanged: false,
})
})
.await
.unwrap_or_else(|error| GateResult {
command: config.command,
success: false,
exit_code: None,
output: error.to_string(),
skipped_unchanged: false,
});
self.last_gate_workspace_hash = Some(hash);
Some(result)
}
#[cfg(feature = "graph-memory")]
fn append_semantic_recalls(&self, base: Option<String>, query: &str) -> Option<String> {
let (Some(config), Some(graph), Some(vector)) = (
self.semantic_recall.as_ref(),
self.graph_memory.as_ref(),
self.semantic_recall
.as_ref()
.and_then(|config| config.embedder.embed(query)),
) else {
return base;
};
let recalls = graph.recall_by_embedding(&vector, config.top_k, config.threshold);
if recalls.is_empty() {
return base;
}
let events = recalls
.iter()
.map(|recall| MemoryRecall {
id: recall.id.clone(),
summary: recall.summary.clone(),
similarity: recall.similarity,
})
.collect();
self.emit(Event::MemoryRecalled { recalls: events });
let suffix = recalls
.iter()
.map(|recall| format!("- {}", recall.summary))
.collect::<Vec<_>>()
.join("\n");
Some(match base {
Some(base) => format!("{base}\n\n# Recalled Memory (cache-safe suffix)\n{suffix}"),
None => format!("# Recalled Memory (cache-safe suffix)\n{suffix}"),
})
}
async fn compact_semantically(
&self,
reason: &str,
provider: &dyn Provider,
) -> Result<(), crate::provider::ProviderError> {
info!("compacting context: {reason}");
if self.message_count() <= 2 {
return Ok(());
}
let before_tokens = self.context_tokens();
let snapshot = self.messages.read().clone();
let result = compact_messages_semantically(
&snapshot,
&self.compaction_config(),
provider,
&self.model,
)
.await?;
if result.removed_count == 0 {
return Ok(());
}
{
let mut messages = self.messages.write();
if !apply_compaction_result(&mut messages, &snapshot, &result) {
return Ok(());
}
messages.push(Message::system(format!("[compact reason: {reason}]")));
}
self.emit(Event::CompactionStart {
reason: reason.to_string(),
before_tokens,
});
self.emit(Event::CompactionEnd {
reason: reason.to_string(),
result,
});
Ok(())
}
}
#[cfg(any(feature = "providers", test))]
fn tool_source(name: &str) -> ToolSource {
if let Some(rest) = name.strip_prefix("mcp__") {
return ToolSource::Mcp {
server: rest.split("__").next().unwrap_or(rest).to_string(),
};
}
if name.starts_with("cu_") {
return ToolSource::ComputerUse;
}
ToolSource::Builtin
}
fn redact_tool_call(call: &ToolCall) -> ToolCall {
ToolCall {
id: call.id.clone(),
name: call.name.clone(),
arguments: crate::secrets::Redactor::new().redact(&call.arguments),
}
}
fn tail_bytes(bytes: &[u8], cap: usize) -> String {
let start = bytes.len().saturating_sub(cap);
let mut text = String::from_utf8_lossy(&bytes[start..]).into_owned();
if start > 0 {
text.insert_str(0, "[output truncated; tail retained]\n");
}
text
}
fn workspace_hash(root: &std::path::Path) -> String {
let mut hasher = Sha256::new();
let output = std::process::Command::new("git")
.args(["diff", "--binary", "HEAD"])
.current_dir(root)
.output();
match output {
Ok(output) => {
hasher.update(&output.stdout);
hasher.update(&output.stderr);
if let Ok(untracked) = std::process::Command::new("git")
.args(["ls-files", "--others", "--exclude-standard"])
.current_dir(root)
.output()
{
for path in String::from_utf8_lossy(&untracked.stdout).lines() {
hasher.update(path.as_bytes());
if let Ok(bytes) = std::fs::read(root.join(path)) {
hasher.update(&bytes);
}
}
}
}
Err(error) => hasher.update(error.to_string().as_bytes()),
}
format!("{:x}", hasher.finalize())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::models::ModelInfo;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
static PARALLEL_DELAY_CALLS: AtomicUsize = AtomicUsize::new(0);
fn delay_read_tool(name: &str) -> ToolDefinition {
ToolDefinition::new_boxed(
name,
"delay read",
"{}",
Box::new(|_ctx, _args| {
Box::pin(async {
PARALLEL_DELAY_CALLS.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(40)).await;
ToolResult::ok("id", "ok")
})
}),
)
.with_effect(ToolEffect::Read)
}
#[tokio::test]
async fn parallel_read_tools_run_concurrently() {
PARALLEL_DELAY_CALLS.store(0, Ordering::SeqCst);
let mut registry = ToolRegistry::new();
registry.register(delay_read_tool("a"));
registry.register(delay_read_tool("b"));
let mut agent = Agent::new();
agent.set_tools(registry);
agent.set_policy(Policy::full_access());
let ctx = Arc::new(ToolContext::new("."));
let calls = vec![
ToolCall {
id: "1".into(),
name: "a".into(),
arguments: "{}".into(),
},
ToolCall {
id: "2".into(),
name: "b".into(),
arguments: "{}".into(),
},
];
let start = std::time::Instant::now();
let results = agent.execute_tools_parallel(&calls, &ctx).await;
assert_eq!(results.len(), 2);
assert!(results.iter().all(|r| !r.is_error));
assert_eq!(PARALLEL_DELAY_CALLS.load(Ordering::SeqCst), 2);
assert!(start.elapsed() < Duration::from_millis(70));
}
#[test]
fn cache_audit_reports_structured_divergence() {
let mut agent = Agent::new();
agent.enable_cache_audit(true);
let seen = Arc::new(parking_lot::Mutex::new(Vec::new()));
let capture = Arc::clone(&seen);
agent.subscribe(move |event| {
if let Event::CacheAudit(audit) = event {
capture.lock().push(audit.clone());
}
});
agent.audit_prompt(&[Message::user("one")], &Some("system".into()), &[]);
agent.audit_prompt(&[Message::user("two")], &Some("system".into()), &[]);
let events = seen.lock();
assert_eq!(events.len(), 2);
assert!(matches!(
events[1].first_divergence,
Some(CacheDivergence::Message { index: 0 })
));
}
#[test]
fn tail_biased_gate_output_is_capped() {
let output = tail_bytes(b"0123456789", 4);
assert!(output.contains("6789"));
assert!(output.contains("truncated"));
}
static CACHE_READ_CALLS: AtomicUsize = AtomicUsize::new(0);
static CACHE_WRITE_CALLS: AtomicUsize = AtomicUsize::new(0);
#[tokio::test]
async fn cache_not_used_for_write_effect() {
CACHE_READ_CALLS.store(0, Ordering::SeqCst);
CACHE_WRITE_CALLS.store(0, Ordering::SeqCst);
let mut registry = ToolRegistry::new();
registry.register(
ToolDefinition::new_boxed(
"r",
"read",
"{}",
Box::new(|_ctx, _args| {
Box::pin(async {
CACHE_READ_CALLS.fetch_add(1, Ordering::SeqCst);
ToolResult::ok("id", "data")
})
}),
)
.with_effect(ToolEffect::Read),
);
registry.register(
ToolDefinition::new_boxed(
"w",
"write",
"{}",
Box::new(|_ctx, _args| {
Box::pin(async {
CACHE_WRITE_CALLS.fetch_add(1, Ordering::SeqCst);
ToolResult::ok("id", "wrote")
})
}),
)
.with_effect(ToolEffect::Write),
);
let mut agent = Agent::new();
agent.set_tools(registry);
agent.set_policy(Policy::full_access());
let ctx = Arc::new(ToolContext::new("."));
let read_call = ToolCall {
id: "1".into(),
name: "r".into(),
arguments: "{}".into(),
};
let write_call = ToolCall {
id: "2".into(),
name: "w".into(),
arguments: "{}".into(),
};
agent.execute_single_tool(&read_call, &ctx).await;
agent.execute_single_tool(&read_call, &ctx).await;
assert_eq!(CACHE_READ_CALLS.load(Ordering::SeqCst), 1);
agent.execute_single_tool(&write_call, &ctx).await;
agent.execute_single_tool(&write_call, &ctx).await;
assert_eq!(CACHE_WRITE_CALLS.load(Ordering::SeqCst), 2);
agent.execute_single_tool(&read_call, &ctx).await;
assert_eq!(CACHE_READ_CALLS.load(Ordering::SeqCst), 2);
}
#[test]
fn compact_uses_token_aware_compaction() {
let mut agent = Agent::new();
agent.auto_compact_after = 50;
{
let mut msgs = agent.messages.write();
msgs.push(Message::system("sys"));
for i in 0..20 {
msgs.push(Message::user(
format!("old message {i} ",) + &"x".repeat(80),
));
msgs.push(Message::assistant("reply".repeat(40)));
}
msgs.push(Message::user("recent tail"));
}
agent.compact("test");
let msgs = agent.messages.read();
assert!(msgs.len() < 42);
assert!(msgs.iter().any(|m| m.content.contains("context compacted")));
assert!(msgs.iter().any(|m| m.content.contains("recent tail")));
}
#[test]
fn default_compaction_threshold_tracks_model_window() {
let mut agent = Agent::new();
agent.set_model("gemini-2.0-flash");
agent.set_model_registry(ModelRegistry::from_models([ModelInfo::new(
"google",
"gemini-2.0-flash",
1_048_576,
8_192,
)]));
assert_eq!(agent.context_window(), 1_048_576);
assert_eq!(agent.auto_compact_threshold(), 943_719);
}
#[test]
fn explicit_compaction_threshold_is_preserved() {
let mut agent = Agent::new();
agent.auto_compact_after = 50;
assert_eq!(agent.auto_compact_threshold(), 50);
}
#[test]
fn compaction_emits_lifecycle() {
let mut agent = Agent::new();
agent.auto_compact_after = 50;
agent
.messages
.write()
.extend((0..10).map(|_| Message::user("x".repeat(100))));
let events = Arc::new(parking_lot::Mutex::new(Vec::new()));
let received = Arc::clone(&events);
agent.subscribe(move |event| {
received.lock().push(event.clone());
});
agent.compact("test");
let events = events.lock();
assert!(events
.iter()
.any(|event| matches!(event, Event::CompactionStart { .. })));
assert!(events
.iter()
.any(|event| matches!(event, Event::CompactionEnd { .. })));
}
#[test]
fn tool_sources_are_classified() {
assert_eq!(tool_source("read"), ToolSource::Builtin);
assert_eq!(tool_source("cu_click"), ToolSource::ComputerUse);
assert_eq!(
tool_source("mcp__supabase__query"),
ToolSource::Mcp {
server: "supabase".into()
}
);
}
#[test]
fn cancellation_handle_cancels_reset_turn() {
let handle = CancellationHandle::new();
let external = handle.clone();
let token = handle.reset();
external.cancel();
assert!(token.is_canceled());
}
#[test]
fn set_scope_preserves_host_shell_policy() {
let mut agent = Agent::new();
agent.set_policy(
Policy::workspace_write()
.with_shell_allow(["git *", "cargo test*"])
.with_shell_deny(["sudo *"])
.with_enforce_dangerous_shell(false),
);
agent.set_scope(Scope::Research);
assert_eq!(
agent.policy.mode,
crate::permissions::PermissionMode::ReadOnly
);
assert_eq!(
agent.policy.shell_allow,
vec!["git *".to_string(), "cargo test*".to_string()]
);
assert_eq!(agent.policy.shell_deny, vec!["sudo *".to_string()]);
assert!(!agent.policy.enforce_dangerous_shell);
assert!(!agent.policy.enable_os_sandbox);
agent.set_scope(Scope::Coding);
assert_eq!(
agent.policy.mode,
crate::permissions::PermissionMode::WorkspaceWrite
);
assert_eq!(
agent.policy.shell_allow,
vec!["git *".to_string(), "cargo test*".to_string()]
);
}
#[test]
fn set_scope_replaces_the_previous_scope_prompt() {
let mut agent = Agent::new();
agent.set_system_prompt("host instructions");
agent.set_scope(Scope::Plan);
assert!(agent
.system_prompt
.as_deref()
.is_some_and(|prompt| prompt.contains("multi-step plan")));
agent.set_scope(Scope::Research);
let prompt = agent.system_prompt.as_deref().expect("system prompt");
assert!(prompt.contains("Explore and explain"));
assert!(!prompt.contains("multi-step plan"));
assert_eq!(prompt.matches("host instructions").count(), 1);
}
#[test]
fn changing_workspace_refreshes_custom_sandbox_and_cache_boundary() {
let first = tempfile::tempdir().unwrap();
let second = tempfile::tempdir().unwrap();
let mut agent = Agent::new();
let mut sandbox = crate::sandbox::SandboxManager::new(
crate::sandbox::SandboxProfile::Custom,
first.path().to_path_buf(),
);
sandbox.set_allow_network(false);
agent.set_sandbox(Arc::new(sandbox));
agent.set_authorizer(Arc::new(crate::permissions::PolicyAuthorizer::new()));
agent.set_workspace_root(second.path());
let current = agent.sandbox.as_ref().expect("sandbox attached");
assert_eq!(current.workspace_root(), second.path());
assert!(current.validate_network().is_err());
assert_eq!(agent.tool_cache.entry_count(), 0);
assert!(agent.authorizer.is_none());
}
#[test]
fn active_skills_are_added_to_a_turn_prompt_without_mutating_the_base() {
let base = Some("host instructions".to_string());
let prompt = turn::append_active_skills(base.clone(), Some("skill instructions"));
assert!(prompt
.as_deref()
.is_some_and(|value| value.contains("# Active Skills")));
assert_eq!(base.as_deref(), Some("host instructions"));
}
#[cfg(feature = "ipc")]
#[test]
fn agent_tool_context_gets_lsp_manager() {
let agent = Agent::new();
let tool_ctx = agent.tool_context();
assert!(tool_ctx.lsp.is_some());
}
#[tokio::test]
async fn tool_result_id_matches_call_id() {
let mut tools = ToolRegistry::new();
tools.register(
ToolDefinition::new_boxed(
"echo_id",
"echo",
"{}",
Box::new(|_ctx, _args| Box::pin(async { ToolResult::ok("wrong-id", "ok") })),
)
.with_effect(ToolEffect::Read),
);
let mut agent = Agent::new();
agent.set_policy(Policy::full_access());
agent.tools = std::sync::Arc::new(tools);
let ctx = std::sync::Arc::new(ToolContext::new(agent.workspace_root.clone()));
let call = ToolCall {
id: "call_xyz".into(),
name: "echo_id".into(),
arguments: "{}".into(),
};
let (_c, result) = agent.execute_single_tool(&call, &ctx).await;
assert_eq!(result.id, "call_xyz");
assert_eq!(result.content, "ok");
}
#[test]
fn approval_required_results_are_typed() {
let result = ToolResult::approval_required("call_approval");
assert!(result.requires_approval());
assert_eq!(result.error_kind, Some(ToolErrorKind::ApprovalRequired));
}
#[tokio::test]
async fn h1_os_sandbox_required_flag_blocks_bash() {
let mut agent = Agent::new();
agent.set_policy(Policy::workspace_write()); agent.os_sandbox_failed = true;
let ctx = std::sync::Arc::new({
let mut tc = ToolContext::new(agent.workspace_root.clone());
tc.os_sandbox_required = true;
tc
});
let result = crate::tools::fs::exec_bash(ctx, r#"{"command":"echo hi"}"#.to_string());
let result = result.await;
assert!(result.is_error);
assert!(result.content.contains("OS sandbox required"));
}
#[test]
fn messages_handle_shares_the_agent_history() {
let agent = Agent::new();
let handle = agent.messages_handle();
handle.write().push(Message::user("from host"));
assert_eq!(agent.message_count(), 1);
agent
.messages
.write()
.push(Message::assistant("from agent"));
assert_eq!(handle.read().len(), 2);
agent.clear_messages();
assert!(handle.read().is_empty());
}
#[cfg(feature = "providers")]
struct SteeringProvider {
handle: Arc<RwLock<Vec<Message>>>,
calls: Arc<parking_lot::Mutex<Vec<Vec<String>>>>,
}
#[cfg(feature = "providers")]
#[async_trait::async_trait]
impl crate::provider::Provider for SteeringProvider {
fn id(&self) -> &str {
"steering"
}
fn name(&self) -> &str {
"steering"
}
async fn stream(
&self,
messages: &[Message],
_system: &Option<String>,
_model: &str,
_tools: &[serde_json::Value],
_reasoning_effort: Option<&str>,
) -> Result<crate::provider::StreamResult, crate::provider::ProviderError> {
let seen: Vec<String> = messages.iter().map(|m| m.content.clone()).collect();
let first = {
let mut calls = self.calls.lock();
calls.push(seen);
calls.len() == 1
};
if first {
self.handle.write().push(Message::user("steer"));
Ok(Box::new(futures::stream::iter([
Ok(crate::provider::StreamEvent::ToolCall(ToolCall {
id: "call_1".into(),
name: "noop".into(),
arguments: "{}".into(),
})),
Ok(crate::provider::StreamEvent::Done),
])))
} else {
Ok(Box::new(futures::stream::iter([Ok(
crate::provider::StreamEvent::Done,
)])))
}
}
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn handle_append_is_seen_mid_turn() {
let mut registry = ToolRegistry::new();
registry.register(
ToolDefinition::new_boxed(
"noop",
"noop",
"{}",
Box::new(|_ctx, _args| Box::pin(async { ToolResult::ok("call_1", "ok") })),
)
.with_effect(ToolEffect::Read),
);
let mut agent = Agent::new();
agent.set_tools(registry);
agent.set_policy(Policy::full_access());
let calls = Arc::new(parking_lot::Mutex::new(Vec::new()));
agent.set_provider(Arc::new(SteeringProvider {
handle: agent.messages_handle(),
calls: Arc::clone(&calls),
}));
agent.prompt("hello").await.unwrap();
let calls = calls.lock();
assert!(calls.len() >= 2, "expected a second tool iteration");
assert!(!calls[0].iter().any(|c| c == "steer"));
assert!(
calls[1].iter().any(|c| c == "steer"),
"mid-turn append not observed on the next iteration: {:?}",
calls[1]
);
}
#[cfg(feature = "providers")]
struct RepeatingProvider {
calls: Arc<std::sync::atomic::AtomicUsize>,
limit: usize,
}
#[cfg(feature = "providers")]
#[async_trait::async_trait]
impl crate::provider::Provider for RepeatingProvider {
fn id(&self) -> &str {
"repeating"
}
fn name(&self) -> &str {
"repeating"
}
async fn stream(
&self,
_messages: &[Message],
_system: &Option<String>,
_model: &str,
_tools: &[serde_json::Value],
_reasoning_effort: Option<&str>,
) -> Result<crate::provider::StreamResult, crate::provider::ProviderError> {
let n = self.calls.fetch_add(1, Ordering::SeqCst);
if n >= self.limit {
return Ok(Box::new(futures::stream::iter([Ok(
crate::provider::StreamEvent::Done,
)])));
}
Ok(Box::new(futures::stream::iter([
Ok(crate::provider::StreamEvent::Delta("working on it".into())),
Ok(crate::provider::StreamEvent::ToolCall(ToolCall {
id: "call_1".into(),
name: "flaky".into(),
arguments: "{}".into(),
})),
Ok(crate::provider::StreamEvent::Done),
])))
}
}
#[cfg(feature = "providers")]
fn looping_agent(
limit: usize,
) -> (
Agent,
Arc<std::sync::atomic::AtomicUsize>,
Arc<std::sync::atomic::AtomicUsize>,
) {
let tool_runs = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let runs = Arc::clone(&tool_runs);
let mut registry = ToolRegistry::new();
registry.register(
ToolDefinition::new_boxed(
"flaky",
"always fails",
"{}",
Box::new(move |_ctx, _args| {
let runs = Arc::clone(&runs);
Box::pin(async move {
runs.fetch_add(1, Ordering::SeqCst);
ToolResult::err("call_1", "disk on fire")
})
}),
)
.with_effect(ToolEffect::Read),
);
let mut agent = Agent::new();
agent.set_tools(registry);
agent.set_policy(Policy::full_access());
agent.max_tool_iterations = 12;
let provider_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
agent.set_provider(Arc::new(RepeatingProvider {
calls: Arc::clone(&provider_calls),
limit,
}));
(agent, tool_runs, provider_calls)
}
#[cfg(feature = "providers")]
fn event_sink(agent: &mut Agent) -> Arc<parking_lot::Mutex<Vec<String>>> {
let seen = Arc::new(parking_lot::Mutex::new(Vec::new()));
let sink = Arc::clone(&seen);
agent.subscribe(move |event| {
let label = match event {
Event::GuardrailWarning { tool, .. } => format!("warn:{tool}"),
Event::GuardrailStop { tool, .. } => format!("stop:{tool}"),
Event::SelfHealing { attempt, .. } => format!("heal:{attempt}"),
Event::PlanProposed(p) => format!("plan_proposed:{}", p.calls.len()),
Event::PlanDecided { decision } => format!("plan_decided:{decision:?}"),
_ => return,
};
sink.lock().push(label);
});
seen
}
#[cfg(feature = "providers")]
fn message_texts(agent: &Agent) -> Vec<String> {
agent
.messages
.read()
.iter()
.map(|m| m.content.clone())
.collect()
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn defaults_leave_the_loop_untouched() {
let (mut agent, tool_runs, _) = looping_agent(3);
let seen = event_sink(&mut agent);
agent.prompt("go").await.unwrap();
assert_eq!(tool_runs.load(Ordering::SeqCst), 3);
assert!(
seen.lock().is_empty(),
"unconfigured agent emitted guardrail events: {:?}",
seen.lock()
);
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn guardrails_warn_on_a_repeated_call_and_tell_the_model() {
let (mut agent, _, _) = looping_agent(4);
agent.set_guardrails(GuardrailConfig {
warnings_enabled: true,
hard_stop_enabled: false,
same_tool_failure_warn_after: 1,
..GuardrailConfig::default()
});
let seen = event_sink(&mut agent);
agent.prompt("go").await.unwrap();
assert!(
seen.lock().iter().any(|l| l == "warn:flaky"),
"no guardrail warning: {:?}",
seen.lock()
);
assert!(
message_texts(&agent)
.iter()
.any(|m| m.starts_with("Guardrail warning:")),
"the warning never reached the model"
);
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn guardrails_stop_a_runaway_turn() {
let (mut agent, tool_runs, _) = looping_agent(usize::MAX);
agent.set_guardrails(GuardrailConfig {
warnings_enabled: false,
hard_stop_enabled: true,
same_tool_failure_halt_after: 2,
..GuardrailConfig::default()
});
let seen = event_sink(&mut agent);
agent.prompt("go").await.unwrap();
assert!(
seen.lock().iter().any(|l| l == "stop:flaky"),
"no guardrail stop: {:?}",
seen.lock()
);
let runs = tool_runs.load(Ordering::SeqCst);
assert!(
runs < 12,
"guardrail did not end the turn: {runs} tool runs against a 12 iteration cap"
);
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn self_healing_reprompts_within_its_budget() {
let (mut agent, _, _) = looping_agent(6);
agent.set_self_healing(2);
let seen = event_sink(&mut agent);
agent.prompt("go").await.unwrap();
let heals: Vec<String> = seen
.lock()
.iter()
.filter(|l| l.starts_with("heal:"))
.cloned()
.collect();
assert_eq!(
heals,
vec!["heal:1", "heal:2"],
"healing budget not honoured"
);
assert!(
message_texts(&agent)
.iter()
.any(|m| m.contains("The following tool call(s) failed")),
"no healing message reached the model"
);
}
#[cfg(feature = "providers")]
struct FixedPlanApprover(PlanDecision);
#[cfg(feature = "providers")]
#[async_trait::async_trait]
impl crate::permissions::PlanApprover for FixedPlanApprover {
async fn approve_plan(&self, _proposal: &PlanProposal) -> PlanDecision {
self.0.clone()
}
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn a_rejected_plan_executes_no_tools() {
let (mut agent, tool_runs, _) = looping_agent(usize::MAX);
agent.set_plan_approver(Arc::new(FixedPlanApprover(PlanDecision::Reject(
"wrong approach".into(),
))));
let seen = event_sink(&mut agent);
agent.prompt("go").await.unwrap();
assert_eq!(
tool_runs.load(Ordering::SeqCst),
0,
"a rejected plan still ran its tools"
);
assert!(seen.lock().iter().any(|l| l == "plan_proposed:1"));
assert!(
message_texts(&agent)
.iter()
.any(|m| m.contains("The plan was rejected")),
"the rejection reason never reached the model"
);
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn an_approved_plan_runs_and_is_gated_once() {
let (mut agent, tool_runs, _) = looping_agent(3);
agent.set_plan_approver(Arc::new(crate::permissions::AlwaysApprovePlan));
let seen = event_sink(&mut agent);
agent.prompt("go").await.unwrap();
assert_eq!(tool_runs.load(Ordering::SeqCst), 3);
let proposals = seen
.lock()
.iter()
.filter(|l| l.starts_with("plan_proposed"))
.count();
assert_eq!(proposals, 1, "the plan gate fired more than once per turn");
}
#[cfg(feature = "providers")]
#[tokio::test]
async fn a_revised_plan_goes_back_to_the_model_unrun() {
let (mut agent, tool_runs, provider_calls) = looping_agent(usize::MAX);
agent.max_tool_iterations = 3;
agent.set_plan_approver(Arc::new(FixedPlanApprover(PlanDecision::Revise(
"use the other tool".into(),
))));
agent.prompt("go").await.unwrap();
assert_eq!(
tool_runs.load(Ordering::SeqCst),
0,
"a plan awaiting revision still ran"
);
assert!(
provider_calls.load(Ordering::SeqCst) > 1,
"revision did not loop back to the model"
);
assert!(
message_texts(&agent)
.iter()
.any(|m| m.contains("use the other tool")),
"the revision guidance never reached the model"
);
}
}
impl Default for Agent {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, thiserror::Error)]
pub enum AgentError {
#[error("provider error: {0}")]
Provider(String),
#[error("tool error: {0}")]
Tool(String),
#[error("no provider configured")]
NoProvider,
#[error("agent cancelled")]
Cancelled,
#[error("budget exceeded: {0}")]
BudgetExceeded(String),
}