use anyhow::Result;
use oxi_sdk::observability::AuditTrail;
use oxi_sdk::{
Agent, AgentConfig, AgentEvent, CompactionEvent, CompactionStrategy, ProviderResolver,
};
use oxi_sdk::{SearchCache, ToolExecutionMode, ToolRegistry};
use parking_lot::Mutex;
use std::collections::HashMap;
use std::sync::Arc;
use crate::access_manager::{AccessGate, AgentContext, TracingAuditSink, TrailAuditSink};
use crate::capability::resolve::resolve_cspace;
use crate::engine::OxiosEngine;
use crate::memory::{MemoryEntry, MemoryManager, MemoryType};
use crate::persona::PersonaManager;
use crate::tools::registration::register_tools_from_cspace_gated;
use crate::KernelHandle;
use crate::event_bus::KernelEvent;
use crate::session_context::SessionContext;
use crate::types::AgentId;
use oxios_ouroboros::{Directive, ExecEnv, ExecutionResult};
static LLM_CIRCUIT_BREAKER: std::sync::OnceLock<oxi_sdk::ProviderCircuitBreaker> =
std::sync::OnceLock::new();
fn get_llm_circuit_breaker() -> &'static oxi_sdk::ProviderCircuitBreaker {
LLM_CIRCUIT_BREAKER.get_or_init(|| {
oxi_sdk::ProviderCircuitBreaker::new(
"global".to_string(),
oxi_sdk::CircuitBreakerConfig::default(),
)
})
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum StreamDelta {
Model(String),
Text(String),
Thinking,
ThinkingDelta(String),
ThinkingEnd,
}
pub type StreamingSinkTx = std::sync::Arc<tokio::sync::mpsc::Sender<StreamDelta>>;
#[derive(Debug, Clone)]
pub struct AgentRuntimeConfig {
pub model_id: String,
pub tool_execution: ToolExecutionMode,
pub auto_retry_enabled: bool,
pub workspace_dir: Option<std::path::PathBuf>,
pub api_key: Option<String>,
pub provider_options: Option<oxi_sdk::ProviderOptions>,
pub rate_limit_per_minute: usize,
pub token_budget: usize,
pub audit_tool_calls: bool,
pub provider_rpm: u32,
pub max_tool_result_bytes: Option<usize>,
pub model_params: Option<oxios_ouroboros::ModelParams>,
}
impl Default for AgentRuntimeConfig {
fn default() -> Self {
Self {
model_id: String::new(),
tool_execution: ToolExecutionMode::Parallel,
auto_retry_enabled: true,
workspace_dir: None,
api_key: None,
provider_options: None,
rate_limit_per_minute: 0,
token_budget: 0,
audit_tool_calls: false,
provider_rpm: 0,
max_tool_result_bytes: None,
model_params: None,
}
}
}
#[derive(Default)]
struct ExecuteState {
final_content: String,
steps_completed: usize,
success: bool,
reasoning_text: String,
trajectory_steps: Vec<oxios_memory::memory::sona::TrajectoryStep>,
pending_tools: std::collections::HashMap<String, (std::time::Instant, usize)>,
tool_call_ids: Vec<String>,
tool_args_map: std::collections::HashMap<String, String>,
tool_error_map: std::collections::HashMap<String, bool>,
tool_timestamps: std::collections::HashMap<String, chrono::DateTime<chrono::Utc>>,
total_input_tokens: u64,
total_output_tokens: u64,
}
pub struct AgentRuntime {
engine_handle: Arc<crate::engine::EngineHandle>,
config: AgentRuntimeConfig,
kernel_handle: Arc<KernelHandle>,
persona_manager: Option<Arc<PersonaManager>>,
tool_retriever: Option<Arc<crate::tools::retrieval::ToolRetriever>>,
routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
persistence_hook: Option<Arc<crate::persistence_hook::PersistenceHook>>,
session_msg_counter: Arc<Mutex<HashMap<String, usize>>>,
}
impl AgentRuntime {
pub fn new(
engine_handle: Arc<crate::engine::EngineHandle>,
kernel_handle: Arc<KernelHandle>,
routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
) -> Self {
Self {
engine_handle,
config: AgentRuntimeConfig::default(),
kernel_handle,
persona_manager: None,
tool_retriever: None,
routing_stats,
persistence_hook: None,
session_msg_counter: Arc::new(Mutex::new(HashMap::new())),
}
}
pub fn with_persona_manager(mut self, pm: Arc<PersonaManager>) -> Self {
self.persona_manager = Some(pm);
self
}
pub fn with_config(mut self, config: AgentRuntimeConfig) -> Self {
self.config = config;
self
}
pub fn with_tool_retriever(
mut self,
retriever: Arc<crate::tools::retrieval::ToolRetriever>,
) -> Self {
self.tool_retriever = Some(retriever);
self
}
pub fn with_persistence_hook(
mut self,
hook: Arc<crate::persistence_hook::PersistenceHook>,
) -> Self {
self.persistence_hook = Some(hook);
self
}
pub async fn execute_directive(
&self,
agent_id: AgentId,
directive: &Directive,
env: &ExecEnv,
session_ctx: &mut SessionContext,
) -> Result<ExecutionResult> {
let session_id: Option<String> = env
.session_id
.clone()
.or_else(|| Some(agent_id.to_string()));
self.execute_directive_with_session(agent_id, directive, env, session_ctx, session_id)
.await
}
pub async fn execute_directive_with_session(
&self,
agent_id: AgentId,
directive: &Directive,
env: &ExecEnv,
session_ctx: &mut SessionContext,
session_id: Option<String>,
) -> Result<ExecutionResult> {
self.execute_inner(
agent_id,
&directive.goal,
&directive.original_request,
&directive.constraints,
&directive.acceptance_criteria,
env.cspace_hint.as_deref(),
&env.mount_paths,
env.workspace_context.as_deref(),
session_ctx,
session_id,
Some(directive),
env.model_override.as_deref(),
env.role.as_deref(),
env.restore_state.as_ref(),
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn execute_inner(
&self,
agent_id: AgentId,
goal: &str,
original_request: &str,
constraints: &[String],
acceptance_criteria: &[String],
cspace_hint: Option<&str>,
mount_paths: &[std::path::PathBuf],
workspace_context: Option<&str>,
session_ctx: &mut SessionContext,
session_id: Option<String>,
persistence_directive: Option<&Directive>,
model_override: Option<&str>,
role: Option<&str>,
restore_state: Option<&serde_json::Value>,
) -> Result<ExecutionResult> {
let prompt = build_user_prompt_inner(goal, acceptance_criteria);
let persona_prompt = self
.persona_manager
.as_ref()
.map(|pm| pm.active_system_prompt())
.filter(|s| !s.trim().is_empty());
let persona_role = self
.persona_manager
.as_ref()
.and_then(|pm| pm.get_active_persona().map(|p| p.role.clone()));
let cspace = resolve_cspace(
cspace_hint,
persona_role.as_deref(),
Some("worker"),
agent_id,
);
let mut system_prompt = build_system_prompt_inner(
goal,
original_request,
constraints,
acceptance_criteria,
workspace_context,
persona_prompt.as_deref(),
None,
None,
);
let capabilities_xml = if let Some(ref retriever) = self.tool_retriever {
match retriever.embedder().embed(goal).await {
Ok(query_vec) => {
let results = retriever.retrieve(&query_vec, 8);
if results.is_empty() {
None
} else {
let xml = crate::tools::retrieval::format_capability_index(&results);
tracing::info!(count = results.len(), "Retrieved relevant capabilities");
Some(xml)
}
}
Err(e) => {
tracing::warn!(error = %e, "Failed to embed goal for retrieval");
None
}
}
} else {
None
};
let kernel_manifest = {
let domains = cspace.active_domains();
if domains.is_empty() {
None
} else {
Some(crate::tools::retrieval::build_kernel_manifest(&domains))
}
};
if capabilities_xml.is_some() || kernel_manifest.is_some() {
system_prompt = build_system_prompt_inner(
goal,
original_request,
constraints,
acceptance_criteria,
workspace_context,
persona_prompt.as_deref(),
capabilities_xml.as_deref(),
kernel_manifest.as_deref(),
);
}
let memory_manager = self.kernel_handle.agents.memory_manager();
match memory_manager
.recall_with_proactive(goal, &mut session_ctx.recall_timing)
.await
{
Ok(memories) if !memories.is_empty() => {
tracing::info!(count = memories.len(), "Recalled memories for task");
system_prompt = memory_manager.blend_into_prompt(&memories, &system_prompt);
}
Ok(_) => tracing::debug!("No memories recalled"),
Err(e) => tracing::warn!(error = %e, "Failed to recall memories"),
}
if let Some(sona) = memory_manager.sona_engine() {
match sona.adapt(goal).await {
Ok(Some(pattern)) if pattern.confidence > 0.5 => {
tracing::info!(
domain = %pattern.domain,
confidence = pattern.confidence,
"SONA learned pattern injected"
);
system_prompt.push_str(&format!(
"\n\n## Learned Strategy (confidence: {:.0}%)\n{}\n",
pattern.confidence * 100.0,
pattern.strategy,
));
}
Ok(_) => tracing::debug!("No high-confidence SONA pattern found"),
Err(e) => tracing::debug!(error = %e, "SONA adapt failed (non-fatal)"),
}
}
match self
.kernel_handle
.knowledge_lens
.recall_for_context(goal, 5)
.await
{
Ok(ctx) if !ctx.notes.is_empty() => {
tracing::info!(
notes = ctx.notes.len(),
memories = ctx.memories.len(),
"Recalled knowledge context for task"
);
let knowledge_blend = ctx
.notes
.iter()
.take(3)
.map(|n| format!("## {}\n\n{}", n.name, n.content))
.collect::<Vec<_>>()
.join("\n\n");
system_prompt.push_str("\n\n## Relevant Knowledge\n\n");
system_prompt.push_str(&knowledge_blend);
}
Ok(_) => tracing::debug!("No knowledge recalled"),
Err(e) => tracing::warn!(error = %e, "Failed to recall knowledge context"),
}
let effective_role = role.or(persona_role.as_deref());
let engine = self.engine_handle.get();
let model_id = model_override
.map(|s| s.to_string())
.or_else(|| effective_role.and_then(|r| self.kernel_handle.engine.model_for_role(r)))
.unwrap_or_else(|| engine.default_model_id().to_string());
engine.resolve_model(&model_id)?;
let exec_id = uuid::Uuid::new_v4();
let mut config = self.config.clone();
config.model_id = model_id;
let kernel_handle = Arc::clone(&self.kernel_handle);
let audit_trail: Option<Arc<AuditTrail>> =
Some(Arc::clone(&self.kernel_handle.security.audit_trail));
let (
mut final_content,
steps_completed,
success,
trajectory_steps,
agent,
tool_call_ids,
tool_args_map,
tool_error_map,
tool_timestamps,
total_input_tokens,
total_output_tokens,
reasoning_text,
) = {
run_agent(
&config,
&engine,
kernel_handle,
system_prompt,
prompt,
exec_id,
goal.to_string(),
agent_id,
cspace,
audit_trail,
self.routing_stats.clone(),
session_id.clone(),
mount_paths,
restore_state,
)
.await?
};
if final_content.is_empty() && !trajectory_steps.is_empty() {
let tool_summary: Vec<String> = trajectory_steps
.iter()
.enumerate()
.map(|(i, step)| {
let truncated = if step.output.len() > 800 {
let mut end = 800;
while end > 0 && !step.output.is_char_boundary(end) {
end -= 1;
}
format!("{}...", &step.output[..end])
} else {
step.output.clone()
};
format!("{}. [{}] {}", i + 1, step.input, truncated)
})
.collect();
let summary_prompt = format!(
"도구 실행 결과:\n\n{}\n\n\
위 결과를 바탕으로 사용자의 요청에 대해 자연스럽게 한국어로 답변해주세요. \
도구의 원시 출력을 그대로 복사하지 말고, 의미 있는 내용만 정리해서 전달하세요.",
tool_summary.join("\n")
);
match agent.run(summary_prompt).await {
Ok((response, _events)) => {
if !response.content.is_empty() {
tracing::info!(exec_id = %exec_id, "Post-execution summary generated");
final_content = response.content;
}
}
Err(e) => {
tracing::warn!(error = %e, "Post-execution summary failed");
}
}
}
let tool_calls: Vec<oxios_ouroboros::ToolCallRecord> = trajectory_steps
.iter()
.enumerate()
.map(|(i, step)| {
let tc_id = tool_call_ids.get(i).cloned().unwrap_or_default();
let args_str = tool_call_ids
.get(i)
.and_then(|id| tool_args_map.get(id))
.cloned()
.unwrap_or_default();
let is_error = tool_call_ids
.get(i)
.and_then(|id| tool_error_map.get(id))
.copied()
.unwrap_or(false);
let timestamp = tool_call_ids
.get(i)
.and_then(|id| tool_timestamps.get(id))
.copied();
let input_str = truncate_json_str(&args_str, 500);
oxios_ouroboros::ToolCallRecord {
tool: step.input.clone(),
input: input_str,
output: step.output.clone(),
duration_ms: step.duration_ms,
is_error,
tool_call_id: tc_id,
timestamp,
}
})
.collect();
tracing::info!(
exec_id = %exec_id,
steps = steps_completed,
success,
tool_calls = tool_calls.len(),
"AgentRuntime finished"
);
let result = ExecutionResult {
output: final_content.clone(),
steps_completed,
success,
tool_calls,
failure_class: None,
restore_state: None,
tokens_input: total_input_tokens,
tokens_output: total_output_tokens,
model_id: self.engine_handle.get().default_model_id().to_string(),
reasoning_text,
};
if let Some(directive) = persistence_directive
&& success
&& let Some(hook) = &self.persistence_hook
{
let already_saved_knowledge = trajectory_steps
.iter()
.any(|s| s.input == "knowledge" && s.output.contains("written successfully"));
let hook = hook.clone();
let directive_clone = directive.clone();
let traj_clone = trajectory_steps.clone();
let output_clone = final_content.clone();
let sid = session_id.clone();
let msg_index = {
let mut counter = self.session_msg_counter.lock();
let idx = counter.entry(sid.clone().unwrap_or_default()).or_insert(0);
let current = *idx;
*idx += 1;
current
};
tokio::spawn(async move {
match hook
.evaluate(
&directive_clone,
&traj_clone,
&output_clone,
already_saved_knowledge,
)
.await
{
Ok(plan) => {
if !plan.memory.is_empty() || !plan.knowledge.is_empty() {
tracing::info!(
memory = plan.memory.len(),
knowledge = plan.knowledge.len(),
message_index = msg_index,
"PersistenceHook executing plan"
);
let session_id = sid.unwrap_or_default();
hook.execute_plan(plan, &session_id, msg_index).await;
}
}
Err(e) => tracing::warn!(error = %e, "PersistenceHook evaluate failed"),
}
});
}
Ok(result)
}
}
#[allow(clippy::too_many_arguments)]
async fn run_agent(
config: &AgentRuntimeConfig,
engine: &OxiosEngine,
kernel_handle: Arc<KernelHandle>,
system_prompt: String,
prompt: String,
exec_id: uuid::Uuid,
goal: String,
agent_id: AgentId,
cspace: crate::capability::CSpace,
audit_trail: Option<Arc<AuditTrail>>,
routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
session_id: Option<String>,
mount_paths: &[std::path::PathBuf],
restore_state: Option<&serde_json::Value>,
) -> Result<(
String,
usize,
bool,
Vec<oxios_memory::memory::sona::TrajectoryStep>,
Arc<Agent>,
Vec<String>,
std::collections::HashMap<String, String>,
std::collections::HashMap<String, bool>,
std::collections::HashMap<String, chrono::DateTime<chrono::Utc>>,
u64,
u64,
String,
)> {
let workspace = if !mount_paths.is_empty() {
mount_paths[0].clone()
} else if let Some(ws) = &config.workspace_dir {
ws.clone()
} else {
std::env::temp_dir()
.join("oxios-agent-workspace")
.join(agent_id.to_string())
};
let _ = std::fs::create_dir_all(&workspace);
tracing::debug!(workspace = %workspace.display(), "Agent workspace scoped");
{
use crate::access_manager::{Role, Subject};
let agent_name = format!("agent-{agent_id}");
let mut am = kernel_handle.exec.access_manager().lock();
let perms = am.get_or_create_permissions(&agent_name);
if let Ok(cwd) = std::env::current_dir() {
let cwd_pattern = format!("{}/**", cwd.to_string_lossy().trim_end_matches('/'));
if !perms.allowed_paths.iter().any(|p| p == &cwd_pattern) {
perms.allow_path(&cwd_pattern);
tracing::debug!(
agent = %agent_name,
path = %cwd_pattern,
"Added CWD to agent allowed paths"
);
}
}
let ws_pattern = format!("{}/**", workspace.to_string_lossy().trim_end_matches('/'));
if !perms.allowed_paths.iter().any(|p| p == &ws_pattern) {
perms.allow_path(&ws_pattern);
}
for mount_path in mount_paths {
let pattern = format!("{}/**", mount_path.to_string_lossy().trim_end_matches('/'));
if !perms.allowed_paths.iter().any(|p| p == &pattern) {
perms.allow_path(&pattern);
tracing::debug!(
agent = %agent_name,
path = %pattern,
"Added Mount path to agent allowed paths (RFC-025)"
);
}
}
let kernel_ws = kernel_handle
.state
.workspace_path()
.to_string_lossy()
.to_string();
let kernel_ws_pattern = format!("{}/**", kernel_ws.trim_end_matches('/'));
if kernel_ws_pattern != ws_pattern
&& !perms.allowed_paths.iter().any(|p| p == &kernel_ws_pattern)
{
perms.allow_path(&kernel_ws_pattern);
}
if !perms.allowed_paths.iter().any(|p| p == "/tmp/**") {
perms.allow_path("/tmp/**");
}
let rbac_subject = Subject::Agent(agent_id);
am.rbac_manager_mut()
.assign_role(rbac_subject, Role::Superuser);
}
let _trace_guard = crate::observability::tracer().start(
format!("exec-{}", &exec_id.to_string()[..8]).as_str(),
oxi_sdk::SpanKind::Agent,
);
let registry = ToolRegistry::new();
let search_cache = Arc::new(SearchCache::new());
let agent_context = AgentContext {
agent_id,
agent_name: format!("agent-{agent_id}"),
cspace: Arc::new(cspace.clone()),
};
let audit_sink: Arc<dyn crate::access_manager::AuditSink> = if let Some(trail) = audit_trail {
let audit_path = kernel_handle
.state
.workspace_path()
.join("audit")
.join("access.jsonl");
Arc::new(TrailAuditSink::new(trail, audit_path))
} else {
Arc::new(TracingAuditSink)
};
let access_gate = Arc::new(AccessGate::new(
kernel_handle.exec.access_manager().clone(),
Arc::new(kernel_handle.exec.config_snapshot()),
audit_sink,
));
register_tools_from_cspace_gated(
®istry,
&kernel_handle,
&cspace,
search_cache,
agent_id,
access_gate,
agent_context,
);
tracing::info!(
exec_id = %exec_id,
capabilities = cspace.len(),
"Tools registered from CSpace"
);
let agent_config = AgentConfig {
name: format!("agent-{agent_id}"),
description: None,
model_id: config.model_id.clone(),
system_prompt: Some(system_prompt.clone()),
timeout_seconds: 300,
temperature: config
.model_params
.as_ref()
.and_then(|p| p.temperature)
.or(Some(0.7)),
max_tokens: config
.model_params
.as_ref()
.and_then(|p| p.max_tokens)
.map(|v| v as usize)
.or(Some(8192)),
compaction_strategy: CompactionStrategy::Threshold(0.8),
compaction_instruction: None,
context_window: 128_000,
workspace_dir: Some(workspace.clone()),
output_mode: None,
provider_options: config.provider_options.clone(),
session_id: None,
max_tool_result_bytes: config.max_tool_result_bytes,
subagent_depth: 0,
subagent_runner: Some(
crate::subagent_runner::OxiosSubagentRunner::new(engine.oxi().clone())
.into_trait_object(),
),
..Default::default()
};
let agent = if config.provider_rpm > 0 {
let resolver: Arc<dyn ProviderResolver> = Arc::new(engine.oxi().clone());
let provider_name = engine.resolve_model(&config.model_id)?.provider;
let provider = engine.pooled_provider(&provider_name, config.provider_rpm)?;
let mut pipeline = oxi_sdk::MiddlewarePipeline::new();
if config.rate_limit_per_minute > 0 {
pipeline = pipeline.push(oxi_sdk::middleware::builtins::RateLimitMiddleware::new(
config.rate_limit_per_minute,
));
}
if config.token_budget > 0 {
pipeline = pipeline.push(oxi_sdk::middleware::builtins::TokenBudgetMiddleware::new(
config.token_budget,
));
}
if config.audit_tool_calls {
pipeline = pipeline.push(oxi_sdk::middleware::builtins::LoggingMiddleware::new(
tracing::Level::INFO,
));
}
let agent = Arc::new(Agent::new_with_resolver(
provider,
agent_config,
Arc::new(registry),
resolver,
));
if !pipeline.is_empty() {
let terminate_flag = Arc::new(std::sync::atomic::AtomicBool::new(false));
let agent_id_for_hooks = agent_id.to_string();
let hooks = oxi_sdk::middleware::build_hooks(
Arc::new(pipeline),
agent_id_for_hooks,
terminate_flag,
);
agent.set_hooks(hooks);
}
agent
} else {
let mut builder = engine
.oxi()
.agent(agent_config)
.workspace(&workspace)
.system_prompt(system_prompt);
let cspace_tool_arcs: Vec<Arc<dyn oxi_sdk::AgentTool>> = registry
.names()
.into_iter()
.filter_map(|name| registry.get(&name))
.collect();
if let Some(auth) = engine.authorizer() {
builder = builder.authorizer(auth.clone());
}
if let Some(tracer) = engine.tracer() {
builder = builder.tracer(tracer.clone());
}
if let Some(ct) = engine.cost_tracker() {
builder = builder.cost_tracker(ct.clone());
}
if config.rate_limit_per_minute > 0 {
builder = builder.with_rate_limit(config.rate_limit_per_minute);
}
if config.token_budget > 0 {
builder = builder.with_token_budget(config.token_budget);
}
if config.audit_tool_calls {
builder = builder.with_logging();
}
let built = builder.build()?;
let agent = Arc::new(built);
let agent_tools = agent.tools();
for tool in cspace_tool_arcs {
agent_tools.register_arc(tool);
}
agent
};
if let Some(state) = restore_state {
agent.import_state(state.clone()).unwrap_or_else(|e| {
tracing::warn!(agent_id = %agent_id, error = %e, "Failed to restore agent state");
});
}
let exec_state = Arc::new(Mutex::new(ExecuteState::default()));
let exec_state_cb = Arc::clone(&exec_state);
let memory_for_callback: Arc<MemoryManager> = (*kernel_handle.agents.memory_manager()).clone();
let session_id_for_callback = exec_id.to_string();
let model_id_for_callback = config.model_id.clone();
let agent_id_for_callback = agent_id.to_string();
let routing_stats_for_cb = routing_stats.clone();
let transparency_session: Option<String> = session_id.clone();
let kernel_handle_for_cb: Arc<KernelHandle> = Arc::clone(&kernel_handle);
let streaming_sinks_for_cb: Arc<crate::streaming_sink::StreamingSinkRegistry> =
Arc::clone(&kernel_handle.streaming_sinks);
let mut sent_model_for_cb: bool = false;
let result = agent
.run_streaming(prompt, move |event| {
if !sent_model_for_cb
&& let Some(ref sid) = transparency_session
&& !model_id_for_callback.is_empty()
&& let Some(tx) = streaming_sinks_for_cb.lookup(sid)
{
let _ = tx.try_send(StreamDelta::Model(model_id_for_callback.clone()));
sent_model_for_cb = true;
}
let mut s = exec_state_cb.lock();
match event {
AgentEvent::ToolExecutionStart {
tool_name,
tool_call_id,
args,
context,
..
} => {
let idx = s.trajectory_steps.len();
s.pending_tools
.insert(tool_call_id.clone(), (std::time::Instant::now(), idx));
s.tool_args_map.insert(
tool_call_id.clone(),
serde_json::to_string(&args).unwrap_or_default(),
);
s.tool_timestamps
.insert(tool_call_id.clone(), chrono::Utc::now());
s.tool_call_ids.push(tool_call_id.clone());
s.trajectory_steps
.push(oxios_memory::memory::sona::TrajectoryStep {
input: tool_name.clone(),
output: String::new(),
duration_ms: 0,
confidence: 0.0,
});
if let Some(ref sid) = transparency_session {
let context_json = context
.as_ref()
.map(serde_json::to_value)
.transpose()
.unwrap_or(None);
let _ =
kernel_handle_for_cb
.infra
.publish(KernelEvent::ToolExecutionStarted {
session_id: sid.clone(),
tool_name: tool_name.clone(),
tool_call_id: tool_call_id.clone(),
tool_args: args.clone(),
context: context_json,
});
}
}
AgentEvent::ToolExecutionUpdate {
tool_call_id,
tool_name,
partial_result,
tab_id,
context,
} => {
if let Some(ref sid) = transparency_session {
let context_json = context
.as_ref()
.map(serde_json::to_value)
.transpose()
.unwrap_or(None);
let _ = kernel_handle_for_cb.infra.publish(
KernelEvent::ToolExecutionProgress {
session_id: sid.clone(),
tool_call_id: tool_call_id.clone(),
tool_name: tool_name.clone(),
progress: partial_result,
tab_id,
context: context_json,
},
);
}
}
AgentEvent::ToolExecutionEnd {
tool_name,
tool_call_id,
is_error,
result,
..
} => {
if !is_error {
s.steps_completed += 1;
}
let mut duration_ms: u64 = 0;
let mut summary = String::new();
if let Some((start, idx)) = s.pending_tools.remove(tool_call_id.as_str()) {
duration_ms = start.elapsed().as_millis() as u64;
if let Some(step) = s.trajectory_steps.get_mut(idx) {
summary = summarize_tool_result(&result.content, 200);
step.output = summary.clone();
step.duration_ms = duration_ms;
step.confidence = if is_error { 0.3 } else { 0.8 };
}
}
s.tool_error_map.insert(tool_call_id.clone(), is_error);
if let Some(ref sid) = transparency_session {
let _ = kernel_handle_for_cb.infra.publish(
KernelEvent::ToolExecutionFinished {
session_id: sid.clone(),
tool_call_id: tool_call_id.clone(),
tool_name: tool_name.clone(),
duration_ms,
is_error,
output_summary: summary,
},
);
}
}
AgentEvent::AgentEnd {
messages,
stop_reason,
..
} => {
if let Some(oxi_sdk::Message::Assistant(a)) = messages.last() {
s.final_content = a.text_content();
}
s.success = matches!(stop_reason.as_deref(), Some("Stop") | Some("ToolUse"));
}
AgentEvent::Error { message, .. } => {
s.final_content = message.clone();
s.success = false;
}
AgentEvent::Usage {
input_tokens,
output_tokens,
} => {
s.total_input_tokens += input_tokens as u64;
s.total_output_tokens += output_tokens as u64;
let agent_label = format!("agent-{agent_id_for_callback}");
crate::observability::cost_tracker().record(
&agent_label,
&oxi_sdk::Model::new(
&model_id_for_callback,
&model_id_for_callback,
oxi_sdk::Api::OpenAiCompletions,
"unknown",
"https://unknown.com",
),
oxi_sdk::TokenUsage {
input: input_tokens as u64,
output: output_tokens as u64,
cache_read: 0,
cache_write: 0,
},
);
if let Some(stats) = &routing_stats_for_cb {
let cost = crate::kernel_handle::engine_api::estimate_cost(
&model_id_for_callback,
input_tokens as u64,
output_tokens as u64,
);
stats.record_model_usage(&model_id_for_callback, cost);
}
if let Some(ref sid) = transparency_session {
let _ = kernel_handle_for_cb
.infra
.publish(KernelEvent::TokenUsageUpdate {
session_id: sid.clone(),
input_tokens: input_tokens as u64,
output_tokens: output_tokens as u64,
});
}
}
AgentEvent::Compaction {
event: CompactionEvent::Completed { result, .. },
} => {
handle_compaction(
result.summary.clone(),
session_id_for_callback.clone(),
memory_for_callback.clone(),
);
if let Some(ref sid) = transparency_session {
let _ =
kernel_handle_for_cb
.infra
.publish(KernelEvent::ReasoningFragment {
session_id: sid.clone(),
content: result.summary.clone(),
source: "compaction".to_string(),
});
}
}
AgentEvent::Compaction {
event: CompactionEvent::Triggered { source, .. },
} => {
if let Some(ref sid) = transparency_session {
let _ =
kernel_handle_for_cb
.infra
.publish(KernelEvent::CompactionTriggered {
session_id: Some(sid.clone()),
source,
});
} else {
let _ =
kernel_handle_for_cb
.infra
.publish(KernelEvent::CompactionTriggered {
session_id: None,
source,
});
}
}
AgentEvent::TextChunk { text } => {
if let Some(ref sid) = transparency_session
&& let Some(tx) = streaming_sinks_for_cb.lookup(sid)
{
let _ = tx.try_send(StreamDelta::Text(text.clone()));
}
}
AgentEvent::Thinking => {
if let Some(ref sid) = transparency_session
&& let Some(tx) = streaming_sinks_for_cb.lookup(sid)
{
let _ = tx.try_send(StreamDelta::Thinking);
}
}
AgentEvent::ThinkingDelta { text } => {
const REASONING_CAP: usize = 4096;
if s.reasoning_text.len() < REASONING_CAP {
s.reasoning_text.push_str(&text);
if s.reasoning_text.len() > REASONING_CAP {
s.reasoning_text.truncate(REASONING_CAP);
}
}
if let Some(ref sid) = transparency_session
&& let Some(tx) = streaming_sinks_for_cb.lookup(sid)
{
let _ = tx.try_send(StreamDelta::ThinkingDelta(text.clone()));
}
}
AgentEvent::ThinkingEnd => {
if let Some(ref sid) = transparency_session
&& let Some(tx) = streaming_sinks_for_cb.lookup(sid)
{
let _ = tx.try_send(StreamDelta::ThinkingEnd);
}
}
AgentEvent::ToolCallDelta {
tool_call_id,
args_delta,
} => {
if let Some(ref sid) = transparency_session {
let _ = kernel_handle_for_cb
.infra
.publish(KernelEvent::ToolArgsDelta {
session_id: sid.clone(),
tool_call_id: tool_call_id.clone(),
args_delta: args_delta.clone(),
});
}
}
_ => {}
}
})
.await;
let circuit = get_llm_circuit_breaker();
if result.is_err() {
circuit.record_failure();
crate::metrics::get_metrics()
.llm_circuit_breaker_state
.set(1.0);
} else {
circuit.record_success();
crate::metrics::get_metrics()
.llm_circuit_breaker_state
.set(0.0);
}
if let Err(e) = result {
tracing::error!(exec_id = %exec_id, error = %e, "Agent failed");
let restore_state = agent.export_state().ok();
return Err(crate::resilience::AgentRunError::wrap(e, restore_state).into());
}
let s = exec_state.lock();
tracing::info!(
exec_id = %exec_id,
steps = s.steps_completed,
success = s.success,
"Agent completed"
);
if !s.trajectory_steps.is_empty()
&& let Some(sona) = kernel_handle.agents.memory_manager().sona_engine()
{
let steps = s.trajectory_steps.clone();
let success = s.success;
let sona = Arc::clone(sona);
let domain = infer_domain(&goal);
tokio::spawn(async move {
let verdict = if success {
oxios_memory::memory::sona::Verdict::Success
} else {
oxios_memory::memory::sona::Verdict::Failure
};
let trajectory = oxios_memory::memory::sona::Trajectory::new(steps, verdict, &domain);
if let Err(e) = sona.record(trajectory).await {
tracing::debug!(error = %e, "SONA trajectory recording failed (non-fatal)");
}
});
}
Ok((
s.final_content.clone(),
s.steps_completed,
s.success,
s.trajectory_steps.clone(),
agent,
s.tool_call_ids.clone(),
s.tool_args_map.clone(),
s.tool_error_map.clone(),
s.tool_timestamps.clone(),
s.total_input_tokens,
s.total_output_tokens,
s.reasoning_text.clone(),
))
}
fn summarize_tool_result(result: &str, max_len: usize) -> String {
let trimmed = result.trim();
if trimmed.chars().count() <= max_len {
return trimmed.to_string();
}
let first_line = trimmed.lines().next().unwrap_or("");
if first_line.chars().count() <= max_len {
first_line.to_string()
} else {
let take = max_len.saturating_sub(3);
let truncated: String = if take == 0 {
first_line.chars().take(max_len).collect()
} else {
first_line.chars().take(take).collect()
};
format!("{truncated}...")
}
}
fn truncate_json_str(json_str: &str, max_len: usize) -> String {
if json_str.len() <= max_len {
return json_str.to_string();
}
let take = max_len.saturating_sub(3);
if take == 0 {
return json_str.chars().take(max_len).collect();
}
let truncated: String = json_str.chars().take(take).collect();
format!("{truncated}...")
}
fn infer_domain(goal: &str) -> String {
let lower = goal.to_lowercase();
let keywords: Vec<&str> = lower.split_whitespace().take(8).collect();
if keywords.iter().any(|k| {
[
"test",
"tests",
"spec",
"testing",
"assert",
"unit test",
"integration",
]
.contains(k)
}) {
return "testing".to_string();
}
if keywords
.iter()
.any(|k| ["deploy", "release", "publish", "ship"].contains(k))
{
return "deployment".to_string();
}
if keywords
.iter()
.any(|k| ["fix", "bug", "patch", "repair", "debug"].contains(k))
{
return "bugfix".to_string();
}
if keywords
.iter()
.any(|k| ["refactor", "restructure", "reorganize", "rewrite"].contains(k))
{
return "refactoring".to_string();
}
if keywords
.iter()
.any(|k| ["doc", "document", "readme", "guide", "explain"].contains(k))
{
return "documentation".to_string();
}
if keywords
.iter()
.any(|k| ["build", "create", "implement", "add", "make", "new"].contains(k))
{
return "development".to_string();
}
if keywords
.iter()
.any(|k| ["analyze", "review", "audit", "inspect", "check"].contains(k))
{
return "analysis".to_string();
}
if keywords
.iter()
.any(|k| ["config", "setup", "install", "configure", "init"].contains(k))
{
return "configuration".to_string();
}
let meaningful: Vec<&str> = lower
.split_whitespace()
.filter(|w| w.len() > 2)
.take(2)
.collect();
if meaningful.len() >= 2 {
meaningful.join("_")
} else {
"general".to_string()
}
}
fn handle_compaction(summary: String, session_id: String, memory_manager: Arc<MemoryManager>) {
let entry = MemoryEntry {
id: uuid::Uuid::new_v4().to_string(),
memory_type: MemoryType::Conversation,
tier: crate::memory::MemoryTier::Warm,
content: summary,
content_hash: 0,
source: "compaction".to_string(),
session_id: Some(session_id),
tags: vec![],
importance: 0.5,
pinned: false,
protection: crate::memory::ProtectionLevel::None,
auto_classified: false,
session_appearances: 0,
user_corrected: false,
seen_in_sessions: vec![],
created_at: chrono::Utc::now(),
accessed_at: chrono::Utc::now(),
modified_at: chrono::Utc::now(),
access_count: 0,
decay_score: 1.0,
compaction_level: 0,
compacted_from: vec![],
related_ids: vec![],
contradicts: None,
};
tokio::spawn(async move {
if let Err(e) = memory_manager.remember(entry).await {
tracing::warn!(error = %e, "Failed to save compaction summary");
}
});
}
#[allow(dead_code)]
fn build_directive_system_prompt(
directive: &Directive,
env: &ExecEnv,
persona_prompt: Option<&str>,
capabilities_xml: Option<&str>,
kernel_manifest: Option<&str>,
) -> String {
build_system_prompt_inner(
&directive.goal,
&directive.original_request,
&directive.constraints,
&directive.acceptance_criteria,
env.workspace_context.as_deref(),
persona_prompt,
capabilities_xml,
kernel_manifest,
)
}
const ARTIFACT_PROTOCOL: &str = "\n\n\
## Artifacts\n\
When you produce substantial, self-contained content the user will want to\n\
view or interact with separately — a complete HTML page, an SVG graphic, a\n\
Mermaid diagram, or an interactive React component — wrap it in an artifact\n\
tag so the UI shows a live preview panel:\n\n\
<lobeArtifact type=\"...\" title=\"...\" identifier=\"...\">\n\
...the full content...\n\
</lobeArtifact>\n\n\
Use exactly one of these `type` values (others are not recognised):\n\
- `text/html` → an HTML document or fragment\n\
- `image/svg+xml` → an SVG graphic\n\
- `application/lobe.artifacts.mermaid` → a Mermaid diagram\n\
- `application/lobe.artifacts.react` → a React component (JSX/TSX)\n\n\
- `title` — a short human title for the panel.\n\
- `identifier` — a unique kebab-case id, e.g. `sales-dashboard`.\n\
Put the full, runnable content INSIDE the tag (not in a separate fence).\n\n\
Do NOT wrap non-visual code. Shell commands, Python, Rust, JSON, config, or\n\
short snippets that are part of an explanation belong in a normal fenced\n\
code block. Use an artifact only for content the user would open to view,\n\
not copy-and-paste. Limit one artifact per self-contained piece.\n";
#[allow(clippy::too_many_arguments)]
fn build_system_prompt_inner(
goal: &str,
original_request: &str,
constraints: &[String],
acceptance_criteria: &[String],
workspace_context: Option<&str>,
persona_prompt: Option<&str>,
capabilities_xml: Option<&str>,
kernel_manifest: Option<&str>,
) -> String {
let mut prompt = String::from(
"You are an autonomous agent in the Oxios operating system.\n\
You execute Seeds — immutable specifications with goals, constraints, and\n\
acceptance criteria.\n\n\
## Available Tools\n\
You have the following tools:\n\
- **File tools**: read, write, edit files; grep, find, ls for searching\n\
- **Web tools**: web_search for searching the web, get_search_results for retrieving cached results\n\
- **Exec**: run shell commands\n\
- **Memory tools**: memory_write (store facts/preferences), memory_read (list entries), memory_search (find relevant memories) — your cross-session recall. Use memory_write proactively when the user shares preferences, facts, or corrections worth remembering.
- **Knowledge**: knowledge — personal markdown vault for documents and notes\n\
- **Kernel tools**: agent, project, persona, cron, security, budget, resource\n\n\
**Important**: When the task involves fetching information from the internet,\n\
websites, or online services, use `web_search` first — do NOT search local files.\n\
When the task asks to \"get\", \"fetch\", \"find online\", or \"look up\" something\n\
from the web, use `web_search`.\n",
);
prompt.push_str(&format!("\n## Goal\n{}\n", goal));
if !original_request.is_empty() && original_request != goal {
prompt.push_str(&format!(
"\n## User's Original Request\n{}\n",
original_request
));
}
if !constraints.is_empty() {
prompt.push_str("\n## Constraints\n");
for (i, c) in constraints.iter().enumerate() {
prompt.push_str(&format!("{}. {}\n", i + 1, c));
}
}
if !acceptance_criteria.is_empty() {
prompt.push_str("\n## Acceptance Criteria\n");
for (i, c) in acceptance_criteria.iter().enumerate() {
prompt.push_str(&format!("{}. {}\n", i + 1, c));
}
}
if let Some(ctx) = workspace_context.filter(|s| !s.trim().is_empty()) {
prompt.push_str("\n## Workspace Context\n");
prompt.push_str(ctx);
prompt.push('\n');
}
if let Some(pp) = persona_prompt {
prompt.push_str("\n## Persona\n");
prompt.push_str(pp);
prompt.push('\n');
}
if let Some(xml) = capabilities_xml {
prompt.push_str("\n## Available Capabilities\n");
prompt.push_str("The following capabilities are relevant to your goal. ");
prompt.push_str("Use the `read` tool to load SKILL.md for any program.\n\n");
prompt.push_str(xml);
prompt.push('\n');
}
if let Some(manifest) = kernel_manifest {
prompt.push('\n');
prompt.push_str(manifest);
prompt.push('\n');
}
prompt.push_str(
"\n## Execution Protocol\n\
1. UNDERSTAND — Read the user's request carefully. If it is a simple\n\
greeting, small talk, or a question you can answer from knowledge,\n\
respond naturally and conversationally — no tools needed.\n\
2. PLAN — For complex tasks, outline your approach before acting.\n\
3. EXECUTE — Use tools only when the task actually requires them.\n\
Prefer the simplest approach. Simple requests need no tools.\n\
4. VERIFY — After each action, check the result: created a file? read it back.\n\
5. REPORT — Summarize how each acceptance criterion was met, with evidence.\n\n\
If the request is ambiguous, use the `ask_user` tool (free-text question)\n\
or the `pi-questionnaire` tool (structured choices) to clarify before\n\
executing — do not guess when a single question would resolve the intent.\n\n\
## Hard Boundaries\n\
- NEVER modify files outside the workspace scope\n\
- NEVER execute destructive commands without confirming scope\n\
- NEVER claim completion without evidence — show the output, not your opinion\n\
- NEVER add features or improvements beyond the goal's scope\n\
- If you cannot complete the task, say so and explain WHY\n\n\
## Scope Guard\n\
The goal defines your universe. Do not:\n\
- Refactor code the goal didn't mention\n\
- Add tests the goal didn't require\n\
- Change configuration the goal didn't specify\n\
- \"Improve\" anything beyond what the acceptance criteria demand\n\n\
## Error Handling\n\
- If a tool fails, read the error message carefully before retrying\n\
- If a command fails, do NOT immediately retry with --force or sudo\n\
- If stuck after 3 attempts, report the blocker rather than continuing to fail\n\n\
## Shape Matching\n\
Match your output to the task: simple task → concise response.\n\
Do not write 50 lines when 5 would do.\n\
Use `exec` for all command execution (git, gh, osascript, etc.).",
);
prompt.push_str(ARTIFACT_PROTOCOL);
prompt
}
#[allow(dead_code)]
fn build_directive_user_prompt(directive: &Directive) -> String {
build_user_prompt_inner(&directive.goal, &directive.acceptance_criteria)
}
fn build_user_prompt_inner(goal: &str, acceptance_criteria: &[String]) -> String {
format!(
"Execute the following goal:\n\n{}\n\nAcceptance criteria:\n{}",
goal,
acceptance_criteria
.iter()
.enumerate()
.map(|(i, c)| format!("{}. {}", i + 1, c))
.collect::<Vec<_>>()
.join("\n")
)
}
impl std::fmt::Debug for AgentRuntime {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AgentRuntime")
.field("model_id", &self.engine_handle.get().default_model_id())
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
use async_trait::async_trait;
use oxi_sdk::{AgentTool, ToolContext, ToolError};
use serde_json::Value;
struct DummyTool {
name: String,
}
#[async_trait]
impl AgentTool for DummyTool {
fn name(&self) -> &str {
&self.name
}
fn label(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
"Test tool"
}
fn parameters_schema(&self) -> Value {
serde_json::json!({"type": "object"})
}
async fn execute(
&self,
_tool_call_id: &str,
_params: Value,
_shutdown: Option<tokio::sync::oneshot::Receiver<()>>,
_ctx: &ToolContext,
) -> Result<oxi_sdk::AgentToolResult, ToolError> {
Ok(oxi_sdk::AgentToolResult::success("ok"))
}
}
#[test]
fn test_requires_tools_validation_passes() {
let registry = ToolRegistry::new();
registry.register(DummyTool {
name: "read".into(),
});
registry.register(DummyTool {
name: "exec".into(),
});
let missing = registry.missing(&["read", "exec"]);
assert!(
missing.is_empty(),
"Expected no missing tools, got: {:?}",
missing
);
}
#[test]
fn test_requires_tools_validation_fails() {
let registry = ToolRegistry::new();
registry.register(DummyTool {
name: "read".into(),
});
let missing = registry.missing(&["read", "exec", "nonexistent"]);
assert_eq!(missing, vec!["exec", "nonexistent"]);
}
#[test]
fn test_infer_domain_testing() {
assert_eq!(infer_domain("run all unit tests for the kernel"), "testing");
}
#[test]
fn test_infer_domain_deployment() {
assert_eq!(
infer_domain("deploy the web service to production"),
"deployment"
);
}
#[test]
fn test_infer_domain_bugfix() {
assert_eq!(infer_domain("fix the null pointer error in main"), "bugfix");
}
#[test]
fn test_infer_domain_development() {
assert_eq!(
infer_domain("create a new REST API endpoint"),
"development"
);
}
#[test]
fn test_infer_domain_analysis() {
assert_eq!(
infer_domain("review the code for security issues"),
"analysis"
);
}
#[test]
fn test_infer_domain_fallback() {
let domain = infer_domain("optimize performance metrics");
assert!(!domain.is_empty());
}
#[test]
fn test_system_prompt_includes_artifact_protocol() {
let prompt = build_system_prompt_inner(
"build a dashboard",
"build a dashboard",
&[],
&[],
None,
None,
None,
None,
);
assert!(prompt.contains("<lobeArtifact"));
assert!(prompt.contains("text/html"));
assert!(prompt.contains("image/svg+xml"));
assert!(prompt.contains("application/lobe.artifacts.mermaid"));
assert!(prompt.contains("application/lobe.artifacts.react"));
}
}