1use anyhow::Result;
31use oxi_sdk::observability::AuditTrail;
32use oxi_sdk::{
33 Agent, AgentConfig, AgentEvent, CompactionEvent, CompactionStrategy, ProviderResolver,
34};
35use oxi_sdk::{SearchCache, ToolExecutionMode, ToolRegistry};
36use parking_lot::Mutex;
37use std::collections::HashMap;
38use std::sync::Arc;
39use crate::access_manager::{AccessGate, AgentContext, TracingAuditSink, TrailAuditSink};
43use crate::capability::resolve::resolve_cspace;
44use crate::engine::OxiosEngine;
45use crate::memory::{MemoryEntry, MemoryManager, MemoryType};
46use crate::persona::PersonaManager;
47use crate::tools::registration::register_tools_from_cspace_gated;
48
49use crate::KernelHandle;
50use crate::event_bus::KernelEvent;
51use crate::session_context::SessionContext;
52use crate::types::AgentId;
53use oxios_ouroboros::{Directive, ExecEnv, ExecutionResult};
54
55static LLM_CIRCUIT_BREAKER: std::sync::OnceLock<oxi_sdk::ProviderCircuitBreaker> =
57 std::sync::OnceLock::new();
58
59fn get_llm_circuit_breaker() -> &'static oxi_sdk::ProviderCircuitBreaker {
61 LLM_CIRCUIT_BREAKER.get_or_init(|| {
62 oxi_sdk::ProviderCircuitBreaker::new(
63 "global".to_string(),
64 oxi_sdk::CircuitBreakerConfig::default(),
65 )
66 })
67}
68
69#[derive(Debug, Clone)]
75#[non_exhaustive]
76pub enum StreamDelta {
77 Model(String),
81 Text(String),
83
84 Thinking,
87 ThinkingDelta(String),
91 ThinkingEnd,
99}
100
101pub type StreamingSinkTx = std::sync::Arc<tokio::sync::mpsc::Sender<StreamDelta>>;
108
109#[derive(Debug, Clone)]
111pub struct AgentRuntimeConfig {
112 pub model_id: String,
114 pub tool_execution: ToolExecutionMode,
116 pub auto_retry_enabled: bool,
118 pub workspace_dir: Option<std::path::PathBuf>,
120 pub api_key: Option<String>,
122 pub provider_options: Option<oxi_sdk::ProviderOptions>,
124 pub rate_limit_per_minute: usize,
126 pub token_budget: usize,
128 pub audit_tool_calls: bool,
130 pub provider_rpm: u32,
133 pub max_tool_result_bytes: Option<usize>,
138 pub model_params: Option<oxios_ouroboros::ModelParams>,
143 }
147
148impl Default for AgentRuntimeConfig {
149 fn default() -> Self {
150 Self {
151 model_id: String::new(),
152 tool_execution: ToolExecutionMode::Parallel,
153 auto_retry_enabled: true,
154 workspace_dir: None,
155 api_key: None,
156 provider_options: None,
157 rate_limit_per_minute: 0,
158 token_budget: 0,
159 audit_tool_calls: false,
160 provider_rpm: 0,
161 max_tool_result_bytes: None,
162 model_params: None,
163 }
164 }
165}
166
167#[derive(Default)]
169struct ExecuteState {
170 final_content: String,
171 steps_completed: usize,
172 success: bool,
173 reasoning_text: String,
178 trajectory_steps: Vec<oxios_memory::memory::sona::TrajectoryStep>,
181 pending_tools: std::collections::HashMap<String, (std::time::Instant, usize)>,
185 tool_call_ids: Vec<String>,
188 tool_args_map: std::collections::HashMap<String, String>,
190 tool_error_map: std::collections::HashMap<String, bool>,
192 tool_timestamps: std::collections::HashMap<String, chrono::DateTime<chrono::Utc>>,
194 total_input_tokens: u64,
196 total_output_tokens: u64,
198}
199
200pub struct AgentRuntime {
209 engine_handle: Arc<crate::engine::EngineHandle>,
210 config: AgentRuntimeConfig,
211 kernel_handle: Arc<KernelHandle>,
213 persona_manager: Option<Arc<PersonaManager>>,
215 tool_retriever: Option<Arc<crate::tools::retrieval::ToolRetriever>>,
217 routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
219 persistence_hook: Option<Arc<crate::persistence_hook::PersistenceHook>>,
221 session_msg_counter: Arc<Mutex<HashMap<String, usize>>>,
223}
224
225impl AgentRuntime {
226 pub fn new(
232 engine_handle: Arc<crate::engine::EngineHandle>,
233 kernel_handle: Arc<KernelHandle>,
234 routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
235 ) -> Self {
236 Self {
237 engine_handle,
238 config: AgentRuntimeConfig::default(),
239 kernel_handle,
240 persona_manager: None,
241 tool_retriever: None,
242 routing_stats,
243 persistence_hook: None,
244 session_msg_counter: Arc::new(Mutex::new(HashMap::new())),
245 }
246 }
247
248 pub fn with_persona_manager(mut self, pm: Arc<PersonaManager>) -> Self {
250 self.persona_manager = Some(pm);
251 self
252 }
253
254 pub fn with_config(mut self, config: AgentRuntimeConfig) -> Self {
256 self.config = config;
257 self
258 }
259
260 pub fn with_tool_retriever(
262 mut self,
263 retriever: Arc<crate::tools::retrieval::ToolRetriever>,
264 ) -> Self {
265 self.tool_retriever = Some(retriever);
266 self
267 }
268
269 pub fn with_persistence_hook(
271 mut self,
272 hook: Arc<crate::persistence_hook::PersistenceHook>,
273 ) -> Self {
274 self.persistence_hook = Some(hook);
275 self
276 }
277
278 pub async fn execute_directive(
284 &self,
285 agent_id: AgentId,
286 directive: &Directive,
287 env: &ExecEnv,
288 session_ctx: &mut SessionContext,
289 ) -> Result<ExecutionResult> {
290 let session_id: Option<String> = env
296 .session_id
297 .clone()
298 .or_else(|| Some(agent_id.to_string()));
299 self.execute_directive_with_session(agent_id, directive, env, session_ctx, session_id)
300 .await
301 }
302 pub async fn execute_directive_with_session(
305 &self,
306 agent_id: AgentId,
307 directive: &Directive,
308 env: &ExecEnv,
309 session_ctx: &mut SessionContext,
310 session_id: Option<String>,
311 ) -> Result<ExecutionResult> {
312 self.execute_inner(
313 agent_id,
314 &directive.goal,
315 &directive.original_request,
316 &directive.constraints,
317 &directive.acceptance_criteria,
318 env.cspace_hint.as_deref(),
319 &env.mount_paths,
320 env.workspace_context.as_deref(),
321 session_ctx,
322 session_id,
323 Some(directive),
324 env.model_override.as_deref(),
325 env.role.as_deref(),
326 env.restore_state.as_ref(),
327 )
328 .await
329 }
330
331 #[allow(clippy::too_many_arguments)]
338 async fn execute_inner(
339 &self,
340 agent_id: AgentId,
341 goal: &str,
342 original_request: &str,
343 constraints: &[String],
344 acceptance_criteria: &[String],
345 cspace_hint: Option<&str>,
346 mount_paths: &[std::path::PathBuf],
347 workspace_context: Option<&str>,
348 session_ctx: &mut SessionContext,
349 session_id: Option<String>,
350 persistence_directive: Option<&Directive>,
351 model_override: Option<&str>,
352 role: Option<&str>,
353 restore_state: Option<&serde_json::Value>,
354 ) -> Result<ExecutionResult> {
355 let prompt = build_user_prompt_inner(goal, acceptance_criteria);
356
357 let persona_prompt = self
359 .persona_manager
360 .as_ref()
361 .map(|pm| pm.active_system_prompt())
362 .filter(|s| !s.trim().is_empty());
363
364 let persona_role = self
366 .persona_manager
367 .as_ref()
368 .and_then(|pm| pm.get_active_persona().map(|p| p.role.clone()));
369
370 let cspace = resolve_cspace(
372 cspace_hint,
373 persona_role.as_deref(),
374 Some("worker"),
375 agent_id,
376 );
377
378 let mut system_prompt = build_system_prompt_inner(
381 goal,
382 original_request,
383 constraints,
384 acceptance_criteria,
385 workspace_context,
386 persona_prompt.as_deref(),
387 None,
388 None,
389 );
390
391 let capabilities_xml = if let Some(ref retriever) = self.tool_retriever {
393 match retriever.embedder().embed(goal).await {
394 Ok(query_vec) => {
395 let results = retriever.retrieve(&query_vec, 8);
396 if results.is_empty() {
397 None
398 } else {
399 let xml = crate::tools::retrieval::format_capability_index(&results);
400 tracing::info!(count = results.len(), "Retrieved relevant capabilities");
401 Some(xml)
402 }
403 }
404 Err(e) => {
405 tracing::warn!(error = %e, "Failed to embed goal for retrieval");
406 None
407 }
408 }
409 } else {
410 None
411 };
412
413 let kernel_manifest = {
415 let domains = cspace.active_domains();
416 if domains.is_empty() {
417 None
418 } else {
419 Some(crate::tools::retrieval::build_kernel_manifest(&domains))
420 }
421 };
422
423 if capabilities_xml.is_some() || kernel_manifest.is_some() {
425 system_prompt = build_system_prompt_inner(
426 goal,
427 original_request,
428 constraints,
429 acceptance_criteria,
430 workspace_context,
431 persona_prompt.as_deref(),
432 capabilities_xml.as_deref(),
433 kernel_manifest.as_deref(),
434 );
435 }
436
437 let memory_manager = self.kernel_handle.agents.memory_manager();
439 match memory_manager
440 .recall_with_proactive(goal, &mut session_ctx.recall_timing)
441 .await
442 {
443 Ok(memories) if !memories.is_empty() => {
444 tracing::info!(count = memories.len(), "Recalled memories for task");
445 system_prompt = memory_manager.blend_into_prompt(&memories, &system_prompt);
446 }
447 Ok(_) => tracing::debug!("No memories recalled"),
448 Err(e) => tracing::warn!(error = %e, "Failed to recall memories"),
449 }
450
451 if let Some(sona) = memory_manager.sona_engine() {
453 match sona.adapt(goal).await {
454 Ok(Some(pattern)) if pattern.confidence > 0.5 => {
455 tracing::info!(
456 domain = %pattern.domain,
457 confidence = pattern.confidence,
458 "SONA learned pattern injected"
459 );
460 system_prompt.push_str(&format!(
461 "\n\n## Learned Strategy (confidence: {:.0}%)\n{}\n",
462 pattern.confidence * 100.0,
463 pattern.strategy,
464 ));
465 }
466 Ok(_) => tracing::debug!("No high-confidence SONA pattern found"),
467 Err(e) => tracing::debug!(error = %e, "SONA adapt failed (non-fatal)"),
468 }
469 }
470
471 match self
473 .kernel_handle
474 .knowledge_lens
475 .recall_for_context(goal, 5)
476 .await
477 {
478 Ok(ctx) if !ctx.notes.is_empty() => {
479 tracing::info!(
480 notes = ctx.notes.len(),
481 memories = ctx.memories.len(),
482 "Recalled knowledge context for task"
483 );
484 let knowledge_blend = ctx
485 .notes
486 .iter()
487 .take(3)
488 .map(|n| format!("## {}\n\n{}", n.name, n.content))
489 .collect::<Vec<_>>()
490 .join("\n\n");
491 system_prompt.push_str("\n\n## Relevant Knowledge\n\n");
492 system_prompt.push_str(&knowledge_blend);
493 }
494 Ok(_) => tracing::debug!("No knowledge recalled"),
495 Err(e) => tracing::warn!(error = %e, "Failed to recall knowledge context"),
496 }
497
498 let effective_role = role.or(persona_role.as_deref());
511 let engine = self.engine_handle.get();
512 let model_id = model_override
513 .map(|s| s.to_string())
514 .or_else(|| effective_role.and_then(|r| self.kernel_handle.engine.model_for_role(r)))
515 .unwrap_or_else(|| engine.default_model_id().to_string());
516 engine.resolve_model(&model_id)?;
518 let exec_id = uuid::Uuid::new_v4();
520
521 let mut config = self.config.clone();
526 config.model_id = model_id;
527 let kernel_handle = Arc::clone(&self.kernel_handle);
528
529 let audit_trail: Option<Arc<AuditTrail>> =
531 Some(Arc::clone(&self.kernel_handle.security.audit_trail));
532
533 let (
534 mut final_content,
535 steps_completed,
536 success,
537 trajectory_steps,
538 agent,
539 tool_call_ids,
540 tool_args_map,
541 tool_error_map,
542 tool_timestamps,
543 total_input_tokens,
544 total_output_tokens,
545 reasoning_text,
546 ) = {
547 run_agent(
548 &config,
549 &engine,
550 kernel_handle,
551 system_prompt,
552 prompt,
553 exec_id,
554 goal.to_string(),
555 agent_id,
556 cspace,
557 audit_trail,
558 self.routing_stats.clone(),
559 session_id.clone(),
560 mount_paths,
561 restore_state,
562 )
563 .await?
564 };
565
566 if final_content.is_empty() && !trajectory_steps.is_empty() {
573 let tool_summary: Vec<String> = trajectory_steps
574 .iter()
575 .enumerate()
576 .map(|(i, step)| {
577 let truncated = if step.output.len() > 800 {
578 let mut end = 800;
582 while end > 0 && !step.output.is_char_boundary(end) {
583 end -= 1;
584 }
585 format!("{}...", &step.output[..end])
586 } else {
587 step.output.clone()
588 };
589 format!("{}. [{}] {}", i + 1, step.input, truncated)
590 })
591 .collect();
592 let summary_prompt = format!(
593 "도구 실행 결과:\n\n{}\n\n\
594 위 결과를 바탕으로 사용자의 요청에 대해 자연스럽게 한국어로 답변해주세요. \
595 도구의 원시 출력을 그대로 복사하지 말고, 의미 있는 내용만 정리해서 전달하세요.",
596 tool_summary.join("\n")
597 );
598 match agent.run(summary_prompt).await {
599 Ok((response, _events)) => {
600 if !response.content.is_empty() {
601 tracing::info!(exec_id = %exec_id, "Post-execution summary generated");
602 final_content = response.content;
603 }
604 }
605 Err(e) => {
606 tracing::warn!(error = %e, "Post-execution summary failed");
607 }
608 }
609 }
610
611 let tool_calls: Vec<oxios_ouroboros::ToolCallRecord> = trajectory_steps
614 .iter()
615 .enumerate()
616 .map(|(i, step)| {
617 let tc_id = tool_call_ids.get(i).cloned().unwrap_or_default();
618 let args_str = tool_call_ids
619 .get(i)
620 .and_then(|id| tool_args_map.get(id))
621 .cloned()
622 .unwrap_or_default();
623 let is_error = tool_call_ids
624 .get(i)
625 .and_then(|id| tool_error_map.get(id))
626 .copied()
627 .unwrap_or(false);
628 let timestamp = tool_call_ids
629 .get(i)
630 .and_then(|id| tool_timestamps.get(id))
631 .copied();
632 let input_str = truncate_json_str(&args_str, 500);
633 oxios_ouroboros::ToolCallRecord {
634 tool: step.input.clone(),
635 input: input_str,
636 output: step.output.clone(),
637 duration_ms: step.duration_ms,
638 is_error,
639 tool_call_id: tc_id,
640 timestamp,
641 }
642 })
643 .collect();
644
645 tracing::info!(
646 exec_id = %exec_id,
647 steps = steps_completed,
648 success,
649 tool_calls = tool_calls.len(),
650 "AgentRuntime finished"
651 );
652
653 let result = ExecutionResult {
654 output: final_content.clone(),
655 steps_completed,
656 success,
657 tool_calls,
658 failure_class: None,
659 restore_state: None,
660 tokens_input: total_input_tokens,
661 tokens_output: total_output_tokens,
662 model_id: self.engine_handle.get().default_model_id().to_string(),
663 reasoning_text,
664 };
665
666 if let Some(directive) = persistence_directive
669 && success
670 && let Some(hook) = &self.persistence_hook
671 {
672 let already_saved_knowledge = trajectory_steps
673 .iter()
674 .any(|s| s.input == "knowledge" && s.output.contains("written successfully"));
675 let hook = hook.clone();
676 let directive_clone = directive.clone();
677 let traj_clone = trajectory_steps.clone();
678 let output_clone = final_content.clone();
679 let sid = session_id.clone();
680 let msg_index = {
683 let mut counter = self.session_msg_counter.lock();
684 let idx = counter.entry(sid.clone().unwrap_or_default()).or_insert(0);
685 let current = *idx;
686 *idx += 1;
687 current
688 };
689 tokio::spawn(async move {
690 match hook
691 .evaluate(
692 &directive_clone,
693 &traj_clone,
694 &output_clone,
695 already_saved_knowledge,
696 )
697 .await
698 {
699 Ok(plan) => {
700 if !plan.memory.is_empty() || !plan.knowledge.is_empty() {
701 tracing::info!(
702 memory = plan.memory.len(),
703 knowledge = plan.knowledge.len(),
704 message_index = msg_index,
705 "PersistenceHook executing plan"
706 );
707 let session_id = sid.unwrap_or_default();
708 hook.execute_plan(plan, &session_id, msg_index).await;
709 }
710 }
711 Err(e) => tracing::warn!(error = %e, "PersistenceHook evaluate failed"),
712 }
713 });
714 }
715
716 Ok(result)
717 }
718}
719
720#[allow(clippy::too_many_arguments)]
725async fn run_agent(
726 config: &AgentRuntimeConfig,
727 engine: &OxiosEngine,
728 kernel_handle: Arc<KernelHandle>,
729 system_prompt: String,
730 prompt: String,
731 exec_id: uuid::Uuid,
732 goal: String,
733 agent_id: AgentId,
734 cspace: crate::capability::CSpace,
735 audit_trail: Option<Arc<AuditTrail>>,
736 routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
737 session_id: Option<String>,
738 mount_paths: &[std::path::PathBuf],
739 restore_state: Option<&serde_json::Value>,
740) -> Result<(
741 String,
742 usize,
743 bool,
744 Vec<oxios_memory::memory::sona::TrajectoryStep>,
745 Arc<Agent>,
746 Vec<String>,
747 std::collections::HashMap<String, String>,
748 std::collections::HashMap<String, bool>,
749 std::collections::HashMap<String, chrono::DateTime<chrono::Utc>>,
750 u64,
751 u64,
752 String,
753)> {
754 let workspace = if !mount_paths.is_empty() {
760 mount_paths[0].clone()
761 } else if let Some(ws) = &config.workspace_dir {
762 ws.clone()
763 } else {
764 std::env::temp_dir()
765 .join("oxios-agent-workspace")
766 .join(agent_id.to_string())
767 };
768
769 let _ = std::fs::create_dir_all(&workspace);
771
772 tracing::debug!(workspace = %workspace.display(), "Agent workspace scoped");
773
774 {
791 use crate::access_manager::{Role, Subject};
792 let agent_name = format!("agent-{agent_id}");
793 let mut am = kernel_handle.exec.access_manager().lock();
794 let perms = am.get_or_create_permissions(&agent_name);
795
796 if let Ok(cwd) = std::env::current_dir() {
798 let cwd_pattern = format!("{}/**", cwd.to_string_lossy().trim_end_matches('/'));
799 if !perms.allowed_paths.iter().any(|p| p == &cwd_pattern) {
800 perms.allow_path(&cwd_pattern);
801 tracing::debug!(
802 agent = %agent_name,
803 path = %cwd_pattern,
804 "Added CWD to agent allowed paths"
805 );
806 }
807 }
808
809 let ws_pattern = format!("{}/**", workspace.to_string_lossy().trim_end_matches('/'));
811 if !perms.allowed_paths.iter().any(|p| p == &ws_pattern) {
812 perms.allow_path(&ws_pattern);
813 }
814
815 for mount_path in mount_paths {
820 let pattern = format!("{}/**", mount_path.to_string_lossy().trim_end_matches('/'));
821 if !perms.allowed_paths.iter().any(|p| p == &pattern) {
822 perms.allow_path(&pattern);
823 tracing::debug!(
824 agent = %agent_name,
825 path = %pattern,
826 "Added Mount path to agent allowed paths (RFC-025)"
827 );
828 }
829 }
830
831 let kernel_ws = kernel_handle
833 .state
834 .workspace_path()
835 .to_string_lossy()
836 .to_string();
837 let kernel_ws_pattern = format!("{}/**", kernel_ws.trim_end_matches('/'));
838 if kernel_ws_pattern != ws_pattern
839 && !perms.allowed_paths.iter().any(|p| p == &kernel_ws_pattern)
840 {
841 perms.allow_path(&kernel_ws_pattern);
842 }
843
844 if !perms.allowed_paths.iter().any(|p| p == "/tmp/**") {
846 perms.allow_path("/tmp/**");
847 }
848
849 let rbac_subject = Subject::Agent(agent_id);
851 am.rbac_manager_mut()
852 .assign_role(rbac_subject, Role::Superuser);
853 }
854
855 let _trace_guard = crate::observability::tracer().start(
857 format!("exec-{}", &exec_id.to_string()[..8]).as_str(),
858 oxi_sdk::SpanKind::Agent,
859 );
860
861 let registry = ToolRegistry::new();
863 let search_cache = Arc::new(SearchCache::new());
864
865 let agent_context = AgentContext {
867 agent_id,
868 agent_name: format!("agent-{agent_id}"),
869 cspace: Arc::new(cspace.clone()),
870 };
871
872 let audit_sink: Arc<dyn crate::access_manager::AuditSink> = if let Some(trail) = audit_trail {
875 let audit_path = kernel_handle
876 .state
877 .workspace_path()
878 .join("audit")
879 .join("access.jsonl");
880 Arc::new(TrailAuditSink::new(trail, audit_path))
881 } else {
882 Arc::new(TracingAuditSink)
883 };
884
885 let access_gate = Arc::new(AccessGate::new(
887 kernel_handle.exec.access_manager().clone(),
888 Arc::new(kernel_handle.exec.config_snapshot()),
889 audit_sink,
890 ));
891
892 register_tools_from_cspace_gated(
893 ®istry,
894 &kernel_handle,
895 &cspace,
896 search_cache,
897 agent_id,
898 access_gate,
899 agent_context,
900 );
901
902 tracing::info!(
903 exec_id = %exec_id,
904 capabilities = cspace.len(),
905 "Tools registered from CSpace"
906 );
907
908 let agent_config = AgentConfig {
916 name: format!("agent-{agent_id}"),
917 description: None,
918 model_id: config.model_id.clone(),
919 system_prompt: Some(system_prompt.clone()),
920 timeout_seconds: 300,
921 temperature: config
922 .model_params
923 .as_ref()
924 .and_then(|p| p.temperature)
925 .or(Some(0.7)),
926 max_tokens: config
927 .model_params
928 .as_ref()
929 .and_then(|p| p.max_tokens)
930 .map(|v| v as usize)
931 .or(Some(8192)),
932 compaction_strategy: CompactionStrategy::Threshold(0.8),
933 compaction_instruction: None,
934 context_window: 128_000,
935 workspace_dir: Some(workspace.clone()),
936 output_mode: None,
937 provider_options: config.provider_options.clone(),
938 session_id: None,
939 max_tool_result_bytes: config.max_tool_result_bytes,
941 subagent_depth: 0,
946 subagent_runner: Some(
949 crate::subagent_runner::OxiosSubagentRunner::new(engine.oxi().clone())
950 .into_trait_object(),
951 ),
952 ..Default::default()
953 };
954
955 let agent = if config.provider_rpm > 0 {
970 let resolver: Arc<dyn ProviderResolver> = Arc::new(engine.oxi().clone());
972 let provider_name = engine.resolve_model(&config.model_id)?.provider;
973 let provider = engine.pooled_provider(&provider_name, config.provider_rpm)?;
974
975 let mut pipeline = oxi_sdk::MiddlewarePipeline::new();
977 if config.rate_limit_per_minute > 0 {
978 pipeline = pipeline.push(oxi_sdk::middleware::builtins::RateLimitMiddleware::new(
979 config.rate_limit_per_minute,
980 ));
981 }
982 if config.token_budget > 0 {
983 pipeline = pipeline.push(oxi_sdk::middleware::builtins::TokenBudgetMiddleware::new(
984 config.token_budget,
985 ));
986 }
987 if config.audit_tool_calls {
988 pipeline = pipeline.push(oxi_sdk::middleware::builtins::LoggingMiddleware::new(
989 tracing::Level::INFO,
990 ));
991 }
992
993 let agent = Arc::new(Agent::new_with_resolver(
995 provider,
996 agent_config,
997 Arc::new(registry),
998 resolver,
999 ));
1000
1001 if !pipeline.is_empty() {
1003 let terminate_flag = Arc::new(std::sync::atomic::AtomicBool::new(false));
1004 let agent_id_for_hooks = agent_id.to_string();
1005 let hooks = oxi_sdk::middleware::build_hooks(
1006 Arc::new(pipeline),
1007 agent_id_for_hooks,
1008 terminate_flag,
1009 );
1010 agent.set_hooks(hooks);
1011 }
1012
1013 agent
1014 } else {
1015 let mut builder = engine
1017 .oxi()
1018 .agent(agent_config)
1019 .workspace(&workspace)
1020 .system_prompt(system_prompt);
1021
1022 let cspace_tool_arcs: Vec<Arc<dyn oxi_sdk::AgentTool>> = registry
1032 .names()
1033 .into_iter()
1034 .filter_map(|name| registry.get(&name))
1035 .collect();
1036
1037 if let Some(auth) = engine.authorizer() {
1039 builder = builder.authorizer(auth.clone());
1040 }
1041 if let Some(tracer) = engine.tracer() {
1042 builder = builder.tracer(tracer.clone());
1043 }
1044 if let Some(ct) = engine.cost_tracker() {
1045 builder = builder.cost_tracker(ct.clone());
1046 }
1047
1048 if config.rate_limit_per_minute > 0 {
1051 builder = builder.with_rate_limit(config.rate_limit_per_minute);
1052 }
1053 if config.token_budget > 0 {
1054 builder = builder.with_token_budget(config.token_budget);
1055 }
1056 if config.audit_tool_calls {
1057 builder = builder.with_logging();
1058 }
1059
1060 let built = builder.build()?;
1061 let agent = Arc::new(built);
1062
1063 let agent_tools = agent.tools();
1068 for tool in cspace_tool_arcs {
1069 agent_tools.register_arc(tool);
1070 }
1071
1072 agent
1073 };
1074
1075 if let Some(state) = restore_state {
1079 agent.import_state(state.clone()).unwrap_or_else(|e| {
1080 tracing::warn!(agent_id = %agent_id, error = %e, "Failed to restore agent state");
1081 });
1082 }
1083
1084 let exec_state = Arc::new(Mutex::new(ExecuteState::default()));
1086 let exec_state_cb = Arc::clone(&exec_state);
1087 let memory_for_callback: Arc<MemoryManager> = (*kernel_handle.agents.memory_manager()).clone();
1088 let session_id_for_callback = exec_id.to_string();
1089 let model_id_for_callback = config.model_id.clone();
1090 let agent_id_for_callback = agent_id.to_string();
1091 let routing_stats_for_cb = routing_stats.clone();
1092 let transparency_session: Option<String> = session_id.clone();
1095 let kernel_handle_for_cb: Arc<KernelHandle> = Arc::clone(&kernel_handle);
1096 let streaming_sinks_for_cb: Arc<crate::streaming_sink::StreamingSinkRegistry> =
1100 Arc::clone(&kernel_handle.streaming_sinks);
1101 let mut sent_model_for_cb: bool = false;
1103 let result = agent
1104 .run_streaming(prompt, move |event| {
1105 if !sent_model_for_cb
1106 && let Some(ref sid) = transparency_session
1107 && !model_id_for_callback.is_empty()
1108 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1109 {
1110 let _ = tx.try_send(StreamDelta::Model(model_id_for_callback.clone()));
1111 sent_model_for_cb = true;
1112 }
1113 let mut s = exec_state_cb.lock();
1114 match event {
1115 AgentEvent::ToolExecutionStart {
1116 tool_name,
1117 tool_call_id,
1118 args,
1119 context,
1120 ..
1121 } => {
1122 let idx = s.trajectory_steps.len();
1124 s.pending_tools
1125 .insert(tool_call_id.clone(), (std::time::Instant::now(), idx));
1126 s.tool_args_map.insert(
1127 tool_call_id.clone(),
1128 serde_json::to_string(&args).unwrap_or_default(),
1129 );
1130 s.tool_timestamps
1131 .insert(tool_call_id.clone(), chrono::Utc::now());
1132 s.tool_call_ids.push(tool_call_id.clone());
1133 s.trajectory_steps
1134 .push(oxios_memory::memory::sona::TrajectoryStep {
1135 input: tool_name.clone(),
1136 output: String::new(),
1137 duration_ms: 0,
1138 confidence: 0.0,
1139 });
1140 if let Some(ref sid) = transparency_session {
1142 let context_json = context
1143 .as_ref()
1144 .map(serde_json::to_value)
1145 .transpose()
1146 .unwrap_or(None);
1147 let _ =
1148 kernel_handle_for_cb
1149 .infra
1150 .publish(KernelEvent::ToolExecutionStarted {
1151 session_id: sid.clone(),
1152 tool_name: tool_name.clone(),
1153 tool_call_id: tool_call_id.clone(),
1154 tool_args: args.clone(),
1155 context: context_json,
1156 });
1157 }
1158 }
1159 AgentEvent::ToolExecutionUpdate {
1160 tool_call_id,
1161 tool_name,
1162 partial_result,
1163 tab_id,
1164 context,
1165 } => {
1166 if let Some(ref sid) = transparency_session {
1176 let context_json = context
1177 .as_ref()
1178 .map(serde_json::to_value)
1179 .transpose()
1180 .unwrap_or(None);
1181 let _ = kernel_handle_for_cb.infra.publish(
1182 KernelEvent::ToolExecutionProgress {
1183 session_id: sid.clone(),
1184 tool_call_id: tool_call_id.clone(),
1185 tool_name: tool_name.clone(),
1186 progress: partial_result,
1187 tab_id,
1188 context: context_json,
1189 },
1190 );
1191 }
1192 }
1193 AgentEvent::ToolExecutionEnd {
1194 tool_name,
1195 tool_call_id,
1196 is_error,
1197 result,
1198 ..
1199 } => {
1200 if !is_error {
1201 s.steps_completed += 1;
1202 }
1203 let mut duration_ms: u64 = 0;
1205 let mut summary = String::new();
1206 if let Some((start, idx)) = s.pending_tools.remove(tool_call_id.as_str()) {
1207 duration_ms = start.elapsed().as_millis() as u64;
1208 if let Some(step) = s.trajectory_steps.get_mut(idx) {
1209 summary = summarize_tool_result(&result.content, 200);
1210 step.output = summary.clone();
1211 step.duration_ms = duration_ms;
1212 step.confidence = if is_error { 0.3 } else { 0.8 };
1213 }
1214 }
1215 s.tool_error_map.insert(tool_call_id.clone(), is_error);
1216 if let Some(ref sid) = transparency_session {
1218 let _ = kernel_handle_for_cb.infra.publish(
1219 KernelEvent::ToolExecutionFinished {
1220 session_id: sid.clone(),
1221 tool_call_id: tool_call_id.clone(),
1222 tool_name: tool_name.clone(),
1223 duration_ms,
1224 is_error,
1225 output_summary: summary,
1226 },
1227 );
1228 }
1229 }
1230 AgentEvent::AgentEnd {
1231 messages,
1232 stop_reason,
1233 ..
1234 } => {
1235 if let Some(oxi_sdk::Message::Assistant(a)) = messages.last() {
1236 s.final_content = a.text_content();
1237 }
1238 s.success = matches!(stop_reason.as_deref(), Some("Stop") | Some("ToolUse"));
1244 }
1245 AgentEvent::Error { message, .. } => {
1246 s.final_content = message.clone();
1247 s.success = false;
1248 }
1249 AgentEvent::Usage {
1250 input_tokens,
1251 output_tokens,
1252 } => {
1253 s.total_input_tokens += input_tokens as u64;
1255 s.total_output_tokens += output_tokens as u64;
1256
1257 let agent_label = format!("agent-{agent_id_for_callback}");
1259 crate::observability::cost_tracker().record(
1260 &agent_label,
1261 &oxi_sdk::Model::new(
1262 &model_id_for_callback,
1263 &model_id_for_callback,
1264 oxi_sdk::Api::OpenAiCompletions,
1265 "unknown",
1266 "https://unknown.com",
1267 ),
1268 oxi_sdk::TokenUsage {
1269 input: input_tokens as u64,
1270 output: output_tokens as u64,
1271 cache_read: 0,
1272 cache_write: 0,
1273 },
1274 );
1275
1276 if let Some(stats) = &routing_stats_for_cb {
1278 let cost = crate::kernel_handle::engine_api::estimate_cost(
1279 &model_id_for_callback,
1280 input_tokens as u64,
1281 output_tokens as u64,
1282 );
1283 stats.record_model_usage(&model_id_for_callback, cost);
1284 }
1285 if let Some(ref sid) = transparency_session {
1287 let _ = kernel_handle_for_cb
1288 .infra
1289 .publish(KernelEvent::TokenUsageUpdate {
1290 session_id: sid.clone(),
1291 input_tokens: input_tokens as u64,
1292 output_tokens: output_tokens as u64,
1293 });
1294 }
1295 }
1296 AgentEvent::Compaction {
1297 event: CompactionEvent::Completed { result, .. },
1298 } => {
1299 handle_compaction(
1300 result.summary.clone(),
1301 session_id_for_callback.clone(),
1302 memory_for_callback.clone(),
1303 );
1304 if let Some(ref sid) = transparency_session {
1306 let _ =
1307 kernel_handle_for_cb
1308 .infra
1309 .publish(KernelEvent::ReasoningFragment {
1310 session_id: sid.clone(),
1311 content: result.summary.clone(),
1312 source: "compaction".to_string(),
1313 });
1314 }
1315 }
1316 AgentEvent::Compaction {
1317 event: CompactionEvent::Triggered { source, .. },
1318 } => {
1319 if let Some(ref sid) = transparency_session {
1325 let _ =
1326 kernel_handle_for_cb
1327 .infra
1328 .publish(KernelEvent::CompactionTriggered {
1329 session_id: Some(sid.clone()),
1330 source,
1331 });
1332 } else {
1333 let _ =
1334 kernel_handle_for_cb
1335 .infra
1336 .publish(KernelEvent::CompactionTriggered {
1337 session_id: None,
1338 source,
1339 });
1340 }
1341 }
1342 AgentEvent::TextChunk { text } => {
1343 if let Some(ref sid) = transparency_session
1357 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1358 {
1359 let _ = tx.try_send(StreamDelta::Text(text.clone()));
1360 }
1361 }
1362 AgentEvent::Thinking => {
1363 if let Some(ref sid) = transparency_session
1368 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1369 {
1370 let _ = tx.try_send(StreamDelta::Thinking);
1371 }
1372 }
1373 AgentEvent::ThinkingDelta { text } => {
1374 const REASONING_CAP: usize = 4096;
1388 if s.reasoning_text.len() < REASONING_CAP {
1389 s.reasoning_text.push_str(&text);
1390 if s.reasoning_text.len() > REASONING_CAP {
1391 s.reasoning_text.truncate(REASONING_CAP);
1392 }
1393 }
1394 if let Some(ref sid) = transparency_session
1395 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1396 {
1397 let _ = tx.try_send(StreamDelta::ThinkingDelta(text.clone()));
1398 }
1399 }
1400 AgentEvent::ThinkingEnd => {
1401 if let Some(ref sid) = transparency_session
1407 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1408 {
1409 let _ = tx.try_send(StreamDelta::ThinkingEnd);
1410 }
1411 }
1412 AgentEvent::ToolCallDelta {
1413 tool_call_id,
1414 args_delta,
1415 } => {
1416 if let Some(ref sid) = transparency_session {
1421 let _ = kernel_handle_for_cb
1422 .infra
1423 .publish(KernelEvent::ToolArgsDelta {
1424 session_id: sid.clone(),
1425 tool_call_id: tool_call_id.clone(),
1426 args_delta: args_delta.clone(),
1427 });
1428 }
1429 }
1430 _ => {}
1431 }
1432 })
1433 .await;
1434
1435 let circuit = get_llm_circuit_breaker();
1437 if result.is_err() {
1438 circuit.record_failure();
1439 crate::metrics::get_metrics()
1440 .llm_circuit_breaker_state
1441 .set(1.0);
1442 } else {
1443 circuit.record_success();
1444 crate::metrics::get_metrics()
1445 .llm_circuit_breaker_state
1446 .set(0.0);
1447 }
1448
1449 if let Err(e) = result {
1450 tracing::error!(exec_id = %exec_id, error = %e, "Agent failed");
1451 let restore_state = agent.export_state().ok();
1456 return Err(crate::resilience::AgentRunError::wrap(e, restore_state).into());
1457 }
1458
1459 let s = exec_state.lock();
1460 tracing::info!(
1461 exec_id = %exec_id,
1462 steps = s.steps_completed,
1463 success = s.success,
1464 "Agent completed"
1465 );
1466
1467 if !s.trajectory_steps.is_empty()
1470 && let Some(sona) = kernel_handle.agents.memory_manager().sona_engine()
1471 {
1472 let steps = s.trajectory_steps.clone();
1473 let success = s.success;
1474 let sona = Arc::clone(sona);
1475 let domain = infer_domain(&goal);
1476 tokio::spawn(async move {
1477 let verdict = if success {
1478 oxios_memory::memory::sona::Verdict::Success
1479 } else {
1480 oxios_memory::memory::sona::Verdict::Failure
1481 };
1482 let trajectory = oxios_memory::memory::sona::Trajectory::new(steps, verdict, &domain);
1483 if let Err(e) = sona.record(trajectory).await {
1484 tracing::debug!(error = %e, "SONA trajectory recording failed (non-fatal)");
1485 }
1486 });
1487 }
1488
1489 Ok((
1490 s.final_content.clone(),
1491 s.steps_completed,
1492 s.success,
1493 s.trajectory_steps.clone(),
1494 agent,
1495 s.tool_call_ids.clone(),
1496 s.tool_args_map.clone(),
1497 s.tool_error_map.clone(),
1498 s.tool_timestamps.clone(),
1499 s.total_input_tokens,
1500 s.total_output_tokens,
1501 s.reasoning_text.clone(),
1502 ))
1503}
1504
1505fn summarize_tool_result(result: &str, max_len: usize) -> String {
1510 let trimmed = result.trim();
1511 if trimmed.chars().count() <= max_len {
1512 return trimmed.to_string();
1513 }
1514 let first_line = trimmed.lines().next().unwrap_or("");
1516 if first_line.chars().count() <= max_len {
1517 first_line.to_string()
1518 } else {
1519 let take = max_len.saturating_sub(3);
1520 let truncated: String = if take == 0 {
1521 first_line.chars().take(max_len).collect()
1522 } else {
1523 first_line.chars().take(take).collect()
1524 };
1525 format!("{truncated}...")
1526 }
1527}
1528fn truncate_json_str(json_str: &str, max_len: usize) -> String {
1529 if json_str.len() <= max_len {
1530 return json_str.to_string();
1531 }
1532 let take = max_len.saturating_sub(3);
1535 if take == 0 {
1536 return json_str.chars().take(max_len).collect();
1537 }
1538 let truncated: String = json_str.chars().take(take).collect();
1539 format!("{truncated}...")
1540}
1541
1542fn infer_domain(goal: &str) -> String {
1547 let lower = goal.to_lowercase();
1548 let keywords: Vec<&str> = lower.split_whitespace().take(8).collect();
1549
1550 if keywords.iter().any(|k| {
1552 [
1553 "test",
1554 "tests",
1555 "spec",
1556 "testing",
1557 "assert",
1558 "unit test",
1559 "integration",
1560 ]
1561 .contains(k)
1562 }) {
1563 return "testing".to_string();
1564 }
1565 if keywords
1566 .iter()
1567 .any(|k| ["deploy", "release", "publish", "ship"].contains(k))
1568 {
1569 return "deployment".to_string();
1570 }
1571 if keywords
1572 .iter()
1573 .any(|k| ["fix", "bug", "patch", "repair", "debug"].contains(k))
1574 {
1575 return "bugfix".to_string();
1576 }
1577 if keywords
1578 .iter()
1579 .any(|k| ["refactor", "restructure", "reorganize", "rewrite"].contains(k))
1580 {
1581 return "refactoring".to_string();
1582 }
1583 if keywords
1584 .iter()
1585 .any(|k| ["doc", "document", "readme", "guide", "explain"].contains(k))
1586 {
1587 return "documentation".to_string();
1588 }
1589 if keywords
1590 .iter()
1591 .any(|k| ["build", "create", "implement", "add", "make", "new"].contains(k))
1592 {
1593 return "development".to_string();
1594 }
1595 if keywords
1596 .iter()
1597 .any(|k| ["analyze", "review", "audit", "inspect", "check"].contains(k))
1598 {
1599 return "analysis".to_string();
1600 }
1601 if keywords
1602 .iter()
1603 .any(|k| ["config", "setup", "install", "configure", "init"].contains(k))
1604 {
1605 return "configuration".to_string();
1606 }
1607
1608 let meaningful: Vec<&str> = lower
1610 .split_whitespace()
1611 .filter(|w| w.len() > 2)
1612 .take(2)
1613 .collect();
1614 if meaningful.len() >= 2 {
1615 meaningful.join("_")
1616 } else {
1617 "general".to_string()
1618 }
1619}
1620
1621fn handle_compaction(summary: String, session_id: String, memory_manager: Arc<MemoryManager>) {
1627 let entry = MemoryEntry {
1628 id: uuid::Uuid::new_v4().to_string(),
1629 memory_type: MemoryType::Conversation,
1630 tier: crate::memory::MemoryTier::Warm,
1631 content: summary,
1632 content_hash: 0,
1633 source: "compaction".to_string(),
1634 session_id: Some(session_id),
1635 tags: vec![],
1636 importance: 0.5,
1637 pinned: false,
1638 protection: crate::memory::ProtectionLevel::None,
1639 auto_classified: false,
1640 session_appearances: 0,
1641 user_corrected: false,
1642 seen_in_sessions: vec![],
1643 created_at: chrono::Utc::now(),
1644 accessed_at: chrono::Utc::now(),
1645 modified_at: chrono::Utc::now(),
1646 access_count: 0,
1647 decay_score: 1.0,
1648 compaction_level: 0,
1649 compacted_from: vec![],
1650 related_ids: vec![],
1651 contradicts: None,
1652 };
1653 tokio::spawn(async move {
1654 if let Err(e) = memory_manager.remember(entry).await {
1655 tracing::warn!(error = %e, "Failed to save compaction summary");
1656 }
1657 });
1658}
1659
1660#[allow(dead_code)]
1665fn build_directive_system_prompt(
1666 directive: &Directive,
1667 env: &ExecEnv,
1668 persona_prompt: Option<&str>,
1669 capabilities_xml: Option<&str>,
1670 kernel_manifest: Option<&str>,
1671) -> String {
1672 build_system_prompt_inner(
1673 &directive.goal,
1674 &directive.original_request,
1675 &directive.constraints,
1676 &directive.acceptance_criteria,
1677 env.workspace_context.as_deref(),
1678 persona_prompt,
1679 capabilities_xml,
1680 kernel_manifest,
1681 )
1682}
1683
1684const ARTIFACT_PROTOCOL: &str = "\n\n\
1691 ## Artifacts\n\
1692 When you produce substantial, self-contained content the user will want to\n\
1693 view or interact with separately — a complete HTML page, an SVG graphic, a\n\
1694 Mermaid diagram, or an interactive React component — wrap it in an artifact\n\
1695 tag so the UI shows a live preview panel:\n\n\
1696 <lobeArtifact type=\"...\" title=\"...\" identifier=\"...\">\n\
1697 ...the full content...\n\
1698 </lobeArtifact>\n\n\
1699 Use exactly one of these `type` values (others are not recognised):\n\
1700 - `text/html` → an HTML document or fragment\n\
1701 - `image/svg+xml` → an SVG graphic\n\
1702 - `application/lobe.artifacts.mermaid` → a Mermaid diagram\n\
1703 - `application/lobe.artifacts.react` → a React component (JSX/TSX)\n\n\
1704 - `title` — a short human title for the panel.\n\
1705 - `identifier` — a unique kebab-case id, e.g. `sales-dashboard`.\n\
1706 Put the full, runnable content INSIDE the tag (not in a separate fence).\n\n\
1707 Do NOT wrap non-visual code. Shell commands, Python, Rust, JSON, config, or\n\
1708 short snippets that are part of an explanation belong in a normal fenced\n\
1709 code block. Use an artifact only for content the user would open to view,\n\
1710 not copy-and-paste. Limit one artifact per self-contained piece.\n";
1711#[allow(clippy::too_many_arguments)]
1717fn build_system_prompt_inner(
1718 goal: &str,
1719 original_request: &str,
1720 constraints: &[String],
1721 acceptance_criteria: &[String],
1722 workspace_context: Option<&str>,
1723 persona_prompt: Option<&str>,
1724 capabilities_xml: Option<&str>,
1725 kernel_manifest: Option<&str>,
1726) -> String {
1727 let mut prompt = String::from(
1728 "You are an autonomous agent in the Oxios operating system.\n\
1729 You execute Seeds — immutable specifications with goals, constraints, and\n\
1730 acceptance criteria.\n\n\
1731 ## Available Tools\n\
1732 You have the following tools:\n\
1733 - **File tools**: read, write, edit files; grep, find, ls for searching\n\
1734 - **Web tools**: web_search for searching the web, get_search_results for retrieving cached results\n\
1735 - **Exec**: run shell commands\n\
1736 - **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.
1737 - **Knowledge**: knowledge — personal markdown vault for documents and notes\n\
1738 - **Kernel tools**: agent, project, persona, cron, security, budget, resource\n\n\
1739 **Important**: When the task involves fetching information from the internet,\n\
1740 websites, or online services, use `web_search` first — do NOT search local files.\n\
1741 When the task asks to \"get\", \"fetch\", \"find online\", or \"look up\" something\n\
1742 from the web, use `web_search`.\n",
1743 );
1744 prompt.push_str(&format!("\n## Goal\n{}\n", goal));
1745
1746 if !original_request.is_empty() && original_request != goal {
1749 prompt.push_str(&format!(
1750 "\n## User's Original Request\n{}\n",
1751 original_request
1752 ));
1753 }
1754
1755 if !constraints.is_empty() {
1756 prompt.push_str("\n## Constraints\n");
1757 for (i, c) in constraints.iter().enumerate() {
1758 prompt.push_str(&format!("{}. {}\n", i + 1, c));
1759 }
1760 }
1761
1762 if !acceptance_criteria.is_empty() {
1763 prompt.push_str("\n## Acceptance Criteria\n");
1764 for (i, c) in acceptance_criteria.iter().enumerate() {
1765 prompt.push_str(&format!("{}. {}\n", i + 1, c));
1766 }
1767 }
1768
1769 if let Some(ctx) = workspace_context.filter(|s| !s.trim().is_empty()) {
1773 prompt.push_str("\n## Workspace Context\n");
1774 prompt.push_str(ctx);
1775 prompt.push('\n');
1776 }
1777
1778 if let Some(pp) = persona_prompt {
1780 prompt.push_str("\n## Persona\n");
1781 prompt.push_str(pp);
1782 prompt.push('\n');
1783 }
1784
1785 if let Some(xml) = capabilities_xml {
1787 prompt.push_str("\n## Available Capabilities\n");
1788 prompt.push_str("The following capabilities are relevant to your goal. ");
1789 prompt.push_str("Use the `read` tool to load SKILL.md for any program.\n\n");
1790 prompt.push_str(xml);
1791 prompt.push('\n');
1792 }
1793
1794 if let Some(manifest) = kernel_manifest {
1796 prompt.push('\n');
1797 prompt.push_str(manifest);
1798 prompt.push('\n');
1799 }
1800
1801 prompt.push_str(
1803 "\n## Execution Protocol\n\
1804 1. UNDERSTAND — Read the user's request carefully. If it is a simple\n\
1805 greeting, small talk, or a question you can answer from knowledge,\n\
1806 respond naturally and conversationally — no tools needed.\n\
1807 2. PLAN — For complex tasks, outline your approach before acting.\n\
1808 3. EXECUTE — Use tools only when the task actually requires them.\n\
1809 Prefer the simplest approach. Simple requests need no tools.\n\
1810 4. VERIFY — After each action, check the result: created a file? read it back.\n\
1811 5. REPORT — Summarize how each acceptance criterion was met, with evidence.\n\n\
1812 If the request is ambiguous, use the `ask_user` tool (free-text question)\n\
1813 or the `pi-questionnaire` tool (structured choices) to clarify before\n\
1814 executing — do not guess when a single question would resolve the intent.\n\n\
1815 ## Hard Boundaries\n\
1816 - NEVER modify files outside the workspace scope\n\
1817 - NEVER execute destructive commands without confirming scope\n\
1818 - NEVER claim completion without evidence — show the output, not your opinion\n\
1819 - NEVER add features or improvements beyond the goal's scope\n\
1820 - If you cannot complete the task, say so and explain WHY\n\n\
1821 ## Scope Guard\n\
1822 The goal defines your universe. Do not:\n\
1823 - Refactor code the goal didn't mention\n\
1824 - Add tests the goal didn't require\n\
1825 - Change configuration the goal didn't specify\n\
1826 - \"Improve\" anything beyond what the acceptance criteria demand\n\n\
1827 ## Error Handling\n\
1828 - If a tool fails, read the error message carefully before retrying\n\
1829 - If a command fails, do NOT immediately retry with --force or sudo\n\
1830 - If stuck after 3 attempts, report the blocker rather than continuing to fail\n\n\
1831 ## Shape Matching\n\
1832 Match your output to the task: simple task → concise response.\n\
1833 Do not write 50 lines when 5 would do.\n\
1834 Use `exec` for all command execution (git, gh, osascript, etc.).",
1835 );
1836 prompt.push_str(ARTIFACT_PROTOCOL);
1837
1838 prompt
1839}
1840#[allow(dead_code)]
1841fn build_directive_user_prompt(directive: &Directive) -> String {
1842 build_user_prompt_inner(&directive.goal, &directive.acceptance_criteria)
1843}
1844
1845fn build_user_prompt_inner(goal: &str, acceptance_criteria: &[String]) -> String {
1847 format!(
1848 "Execute the following goal:\n\n{}\n\nAcceptance criteria:\n{}",
1849 goal,
1850 acceptance_criteria
1851 .iter()
1852 .enumerate()
1853 .map(|(i, c)| format!("{}. {}", i + 1, c))
1854 .collect::<Vec<_>>()
1855 .join("\n")
1856 )
1857}
1858
1859impl std::fmt::Debug for AgentRuntime {
1860 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1861 f.debug_struct("AgentRuntime")
1862 .field("model_id", &self.engine_handle.get().default_model_id())
1863 .finish()
1864 }
1865}
1866
1867#[cfg(test)]
1868mod tests {
1869 use super::*;
1870 use async_trait::async_trait;
1871 use oxi_sdk::{AgentTool, ToolContext, ToolError};
1872 use serde_json::Value;
1873
1874 struct DummyTool {
1876 name: String,
1877 }
1878
1879 #[async_trait]
1880 impl AgentTool for DummyTool {
1881 fn name(&self) -> &str {
1882 &self.name
1883 }
1884 fn label(&self) -> &str {
1885 &self.name
1886 }
1887 fn description(&self) -> &str {
1888 "Test tool"
1889 }
1890 fn parameters_schema(&self) -> Value {
1891 serde_json::json!({"type": "object"})
1892 }
1893
1894 async fn execute(
1895 &self,
1896 _tool_call_id: &str,
1897 _params: Value,
1898 _shutdown: Option<tokio::sync::oneshot::Receiver<()>>,
1899 _ctx: &ToolContext,
1900 ) -> Result<oxi_sdk::AgentToolResult, ToolError> {
1901 Ok(oxi_sdk::AgentToolResult::success("ok"))
1902 }
1903 }
1904
1905 #[test]
1907 fn test_requires_tools_validation_passes() {
1908 let registry = ToolRegistry::new();
1909
1910 registry.register(DummyTool {
1911 name: "read".into(),
1912 });
1913 registry.register(DummyTool {
1914 name: "exec".into(),
1915 });
1916
1917 let missing = registry.missing(&["read", "exec"]);
1918
1919 assert!(
1920 missing.is_empty(),
1921 "Expected no missing tools, got: {:?}",
1922 missing
1923 );
1924 }
1925
1926 #[test]
1928 fn test_requires_tools_validation_fails() {
1929 let registry = ToolRegistry::new();
1930
1931 registry.register(DummyTool {
1932 name: "read".into(),
1933 });
1934
1935 let missing = registry.missing(&["read", "exec", "nonexistent"]);
1936
1937 assert_eq!(missing, vec!["exec", "nonexistent"]);
1938 }
1939
1940 #[test]
1941 fn test_infer_domain_testing() {
1942 assert_eq!(infer_domain("run all unit tests for the kernel"), "testing");
1943 }
1944
1945 #[test]
1946 fn test_infer_domain_deployment() {
1947 assert_eq!(
1948 infer_domain("deploy the web service to production"),
1949 "deployment"
1950 );
1951 }
1952
1953 #[test]
1954 fn test_infer_domain_bugfix() {
1955 assert_eq!(infer_domain("fix the null pointer error in main"), "bugfix");
1956 }
1957
1958 #[test]
1959 fn test_infer_domain_development() {
1960 assert_eq!(
1961 infer_domain("create a new REST API endpoint"),
1962 "development"
1963 );
1964 }
1965
1966 #[test]
1967 fn test_infer_domain_analysis() {
1968 assert_eq!(
1969 infer_domain("review the code for security issues"),
1970 "analysis"
1971 );
1972 }
1973
1974 #[test]
1975 fn test_infer_domain_fallback() {
1976 let domain = infer_domain("optimize performance metrics");
1977 assert!(!domain.is_empty());
1979 }
1980 #[test]
1981 fn test_system_prompt_includes_artifact_protocol() {
1982 let prompt = build_system_prompt_inner(
1983 "build a dashboard",
1984 "build a dashboard",
1985 &[],
1986 &[],
1987 None,
1988 None,
1989 None,
1990 None,
1991 );
1992 assert!(prompt.contains("<lobeArtifact"));
1995 assert!(prompt.contains("text/html"));
1996 assert!(prompt.contains("image/svg+xml"));
1997 assert!(prompt.contains("application/lobe.artifacts.mermaid"));
1998 assert!(prompt.contains("application/lobe.artifacts.react"));
1999 }
2000}