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 reasoning_segments: Vec<oxios_ouroboros::ReasoningSegment>,
181 reasoning_bytes: usize,
182 trajectory_steps: Vec<oxios_memory::memory::sona::TrajectoryStep>,
185 pending_tools: std::collections::HashMap<String, (std::time::Instant, usize)>,
189 tool_call_ids: Vec<String>,
192 tool_args_map: std::collections::HashMap<String, String>,
194 tool_error_map: std::collections::HashMap<String, bool>,
196 tool_timestamps: std::collections::HashMap<String, chrono::DateTime<chrono::Utc>>,
198 total_input_tokens: u64,
200 total_output_tokens: u64,
202}
203
204pub struct AgentRuntime {
213 engine_handle: Arc<crate::engine::EngineHandle>,
214 config: AgentRuntimeConfig,
215 kernel_handle: Arc<KernelHandle>,
217 persona_manager: Option<Arc<PersonaManager>>,
219 tool_retriever: Option<Arc<crate::tools::retrieval::ToolRetriever>>,
221 routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
223 persistence_hook: Option<Arc<crate::persistence_hook::PersistenceHook>>,
225 session_msg_counter: Arc<Mutex<HashMap<String, usize>>>,
227}
228
229impl AgentRuntime {
230 pub fn new(
236 engine_handle: Arc<crate::engine::EngineHandle>,
237 kernel_handle: Arc<KernelHandle>,
238 routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
239 ) -> Self {
240 Self {
241 engine_handle,
242 config: AgentRuntimeConfig::default(),
243 kernel_handle,
244 persona_manager: None,
245 tool_retriever: None,
246 routing_stats,
247 persistence_hook: None,
248 session_msg_counter: Arc::new(Mutex::new(HashMap::new())),
249 }
250 }
251
252 pub fn with_persona_manager(mut self, pm: Arc<PersonaManager>) -> Self {
254 self.persona_manager = Some(pm);
255 self
256 }
257
258 pub fn with_config(mut self, config: AgentRuntimeConfig) -> Self {
260 self.config = config;
261 self
262 }
263
264 pub fn with_tool_retriever(
266 mut self,
267 retriever: Arc<crate::tools::retrieval::ToolRetriever>,
268 ) -> Self {
269 self.tool_retriever = Some(retriever);
270 self
271 }
272
273 pub fn with_persistence_hook(
275 mut self,
276 hook: Arc<crate::persistence_hook::PersistenceHook>,
277 ) -> Self {
278 self.persistence_hook = Some(hook);
279 self
280 }
281
282 pub async fn execute_directive(
288 &self,
289 agent_id: AgentId,
290 directive: &Directive,
291 env: &ExecEnv,
292 session_ctx: &mut SessionContext,
293 ) -> Result<ExecutionResult> {
294 let session_id: Option<String> = env
300 .session_id
301 .clone()
302 .or_else(|| Some(agent_id.to_string()));
303 self.execute_directive_with_session(agent_id, directive, env, session_ctx, session_id)
304 .await
305 }
306 pub async fn execute_directive_with_session(
309 &self,
310 agent_id: AgentId,
311 directive: &Directive,
312 env: &ExecEnv,
313 session_ctx: &mut SessionContext,
314 session_id: Option<String>,
315 ) -> Result<ExecutionResult> {
316 self.execute_inner(
317 agent_id,
318 &directive.goal,
319 &directive.original_request,
320 &directive.constraints,
321 &directive.acceptance_criteria,
322 env.cspace_hint.as_deref(),
323 &env.mount_paths,
324 env.workspace_context.as_deref(),
325 session_ctx,
326 session_id,
327 Some(directive),
328 env.model_override.as_deref(),
329 env.role.as_deref(),
330 env.restore_state.as_ref(),
331 )
332 .await
333 }
334
335 #[allow(clippy::too_many_arguments)]
342 async fn execute_inner(
343 &self,
344 agent_id: AgentId,
345 goal: &str,
346 original_request: &str,
347 constraints: &[String],
348 acceptance_criteria: &[String],
349 cspace_hint: Option<&str>,
350 mount_paths: &[std::path::PathBuf],
351 workspace_context: Option<&str>,
352 session_ctx: &mut SessionContext,
353 session_id: Option<String>,
354 persistence_directive: Option<&Directive>,
355 model_override: Option<&str>,
356 role: Option<&str>,
357 restore_state: Option<&serde_json::Value>,
358 ) -> Result<ExecutionResult> {
359 let prompt = build_user_prompt_inner(goal, acceptance_criteria);
360
361 let persona_prompt = self
363 .persona_manager
364 .as_ref()
365 .map(|pm| pm.active_system_prompt())
366 .filter(|s| !s.trim().is_empty());
367
368 let persona_role = self
370 .persona_manager
371 .as_ref()
372 .and_then(|pm| pm.get_active_persona().map(|p| p.role.clone()));
373
374 let cspace = resolve_cspace(
376 cspace_hint,
377 persona_role.as_deref(),
378 Some("worker"),
379 agent_id,
380 );
381
382 let mut system_prompt = build_system_prompt_inner(
385 goal,
386 original_request,
387 constraints,
388 acceptance_criteria,
389 workspace_context,
390 persona_prompt.as_deref(),
391 None,
392 None,
393 );
394
395 let capabilities_xml = if let Some(ref retriever) = self.tool_retriever {
397 match retriever.embedder().embed(goal).await {
398 Ok(query_vec) => {
399 let results = retriever.retrieve(&query_vec, 8);
400 if results.is_empty() {
401 None
402 } else {
403 let xml = crate::tools::retrieval::format_capability_index(&results);
404 tracing::info!(count = results.len(), "Retrieved relevant capabilities");
405 Some(xml)
406 }
407 }
408 Err(e) => {
409 tracing::warn!(error = %e, "Failed to embed goal for retrieval");
410 None
411 }
412 }
413 } else {
414 None
415 };
416
417 let kernel_manifest = {
419 let domains = cspace.active_domains();
420 if domains.is_empty() {
421 None
422 } else {
423 Some(crate::tools::retrieval::build_kernel_manifest(&domains))
424 }
425 };
426
427 if capabilities_xml.is_some() || kernel_manifest.is_some() {
429 system_prompt = build_system_prompt_inner(
430 goal,
431 original_request,
432 constraints,
433 acceptance_criteria,
434 workspace_context,
435 persona_prompt.as_deref(),
436 capabilities_xml.as_deref(),
437 kernel_manifest.as_deref(),
438 );
439 }
440
441 let memory_manager = self.kernel_handle.agents.memory_manager();
443 match memory_manager
444 .recall_with_proactive(goal, &mut session_ctx.recall_timing)
445 .await
446 {
447 Ok(memories) if !memories.is_empty() => {
448 tracing::info!(count = memories.len(), "Recalled memories for task");
449 system_prompt = memory_manager.blend_into_prompt(&memories, &system_prompt);
450 }
451 Ok(_) => tracing::debug!("No memories recalled"),
452 Err(e) => tracing::warn!(error = %e, "Failed to recall memories"),
453 }
454
455 if let Some(sona) = memory_manager.sona_engine() {
457 match sona.adapt(goal).await {
458 Ok(Some(pattern)) if pattern.confidence > 0.5 => {
459 tracing::info!(
460 domain = %pattern.domain,
461 confidence = pattern.confidence,
462 "SONA learned pattern injected"
463 );
464 system_prompt.push_str(&format!(
465 "\n\n## Learned Strategy (confidence: {:.0}%)\n{}\n",
466 pattern.confidence * 100.0,
467 pattern.strategy,
468 ));
469 }
470 Ok(_) => tracing::debug!("No high-confidence SONA pattern found"),
471 Err(e) => tracing::debug!(error = %e, "SONA adapt failed (non-fatal)"),
472 }
473 }
474
475 match self
477 .kernel_handle
478 .knowledge_lens
479 .recall_for_context(goal, 5)
480 .await
481 {
482 Ok(ctx) if !ctx.notes.is_empty() => {
483 tracing::info!(
484 notes = ctx.notes.len(),
485 memories = ctx.memories.len(),
486 "Recalled knowledge context for task"
487 );
488 let knowledge_blend = ctx
489 .notes
490 .iter()
491 .take(3)
492 .map(|n| format!("## {}\n\n{}", n.name, n.content))
493 .collect::<Vec<_>>()
494 .join("\n\n");
495 system_prompt.push_str("\n\n## Relevant Knowledge\n\n");
496 system_prompt.push_str(&knowledge_blend);
497 }
498 Ok(_) => tracing::debug!("No knowledge recalled"),
499 Err(e) => tracing::warn!(error = %e, "Failed to recall knowledge context"),
500 }
501
502 let effective_role = role.or(persona_role.as_deref());
515 let engine = self.engine_handle.get();
516 let model_id = model_override
517 .map(|s| s.to_string())
518 .or_else(|| effective_role.and_then(|r| self.kernel_handle.engine.model_for_role(r)))
519 .unwrap_or_else(|| engine.default_model_id().to_string());
520 engine.resolve_model(&model_id)?;
522 let exec_id = uuid::Uuid::new_v4();
524
525 let mut config = self.config.clone();
530 config.model_id = model_id;
531 let kernel_handle = Arc::clone(&self.kernel_handle);
532
533 let audit_trail: Option<Arc<AuditTrail>> =
535 Some(Arc::clone(&self.kernel_handle.security.audit_trail));
536
537 let (
538 mut final_content,
539 steps_completed,
540 success,
541 trajectory_steps,
542 agent,
543 tool_call_ids,
544 tool_args_map,
545 tool_error_map,
546 tool_timestamps,
547 total_input_tokens,
548 total_output_tokens,
549 reasoning_text,
550 reasoning_segments,
551 ) = {
552 run_agent(
553 &config,
554 &engine,
555 kernel_handle,
556 system_prompt,
557 prompt,
558 exec_id,
559 goal.to_string(),
560 agent_id,
561 cspace,
562 audit_trail,
563 self.routing_stats.clone(),
564 session_id.clone(),
565 mount_paths,
566 restore_state,
567 )
568 .await?
569 };
570
571 if final_content.is_empty() && !trajectory_steps.is_empty() {
578 let tool_summary: Vec<String> = trajectory_steps
579 .iter()
580 .enumerate()
581 .map(|(i, step)| {
582 let truncated = if step.output.len() > 800 {
583 let mut end = 800;
587 while end > 0 && !step.output.is_char_boundary(end) {
588 end -= 1;
589 }
590 format!("{}...", &step.output[..end])
591 } else {
592 step.output.clone()
593 };
594 format!("{}. [{}] {}", i + 1, step.input, truncated)
595 })
596 .collect();
597 let summary_prompt = format!(
598 "도구 실행 결과:\n\n{}\n\n\
599 위 결과를 바탕으로 사용자의 요청에 대해 자연스럽게 한국어로 답변해주세요. \
600 도구의 원시 출력을 그대로 복사하지 말고, 의미 있는 내용만 정리해서 전달하세요.",
601 tool_summary.join("\n")
602 );
603 match agent.run(summary_prompt).await {
604 Ok((response, _events)) => {
605 if !response.content.is_empty() {
606 tracing::info!(exec_id = %exec_id, "Post-execution summary generated");
607 final_content = response.content;
608 }
609 }
610 Err(e) => {
611 tracing::warn!(error = %e, "Post-execution summary failed");
612 }
613 }
614 }
615
616 let tool_calls: Vec<oxios_ouroboros::ToolCallRecord> = trajectory_steps
619 .iter()
620 .enumerate()
621 .map(|(i, step)| {
622 let tc_id = tool_call_ids.get(i).cloned().unwrap_or_default();
623 let args_str = tool_call_ids
624 .get(i)
625 .and_then(|id| tool_args_map.get(id))
626 .cloned()
627 .unwrap_or_default();
628 let is_error = tool_call_ids
629 .get(i)
630 .and_then(|id| tool_error_map.get(id))
631 .copied()
632 .unwrap_or(false);
633 let timestamp = tool_call_ids
634 .get(i)
635 .and_then(|id| tool_timestamps.get(id))
636 .copied();
637 let input_str = truncate_json_str(&args_str, 500);
638 oxios_ouroboros::ToolCallRecord {
639 tool: step.input.clone(),
640 input: input_str,
641 output: step.output.clone(),
642 duration_ms: step.duration_ms,
643 is_error,
644 tool_call_id: tc_id,
645 timestamp,
646 }
647 })
648 .collect();
649
650 tracing::info!(
651 exec_id = %exec_id,
652 steps = steps_completed,
653 success,
654 tool_calls = tool_calls.len(),
655 "AgentRuntime finished"
656 );
657
658 let result = ExecutionResult {
659 output: final_content.clone(),
660 steps_completed,
661 success,
662 tool_calls,
663 failure_class: None,
664 restore_state: None,
665 tokens_input: total_input_tokens,
666 tokens_output: total_output_tokens,
667 model_id: self.engine_handle.get().default_model_id().to_string(),
668 reasoning_text,
669 reasoning_segments,
670 };
671
672 if let Some(directive) = persistence_directive
675 && success
676 && let Some(hook) = &self.persistence_hook
677 {
678 let already_saved_knowledge = trajectory_steps
679 .iter()
680 .any(|s| s.input == "knowledge" && s.output.contains("written successfully"));
681 let hook = hook.clone();
682 let directive_clone = directive.clone();
683 let traj_clone = trajectory_steps.clone();
684 let output_clone = final_content.clone();
685 let sid = session_id.clone();
686 let msg_index = {
689 let mut counter = self.session_msg_counter.lock();
690 let idx = counter.entry(sid.clone().unwrap_or_default()).or_insert(0);
691 let current = *idx;
692 *idx += 1;
693 current
694 };
695 tokio::spawn(async move {
696 match hook
697 .evaluate(
698 &directive_clone,
699 &traj_clone,
700 &output_clone,
701 already_saved_knowledge,
702 )
703 .await
704 {
705 Ok(plan) => {
706 if !plan.memory.is_empty() || !plan.knowledge.is_empty() {
707 tracing::info!(
708 memory = plan.memory.len(),
709 knowledge = plan.knowledge.len(),
710 message_index = msg_index,
711 "PersistenceHook executing plan"
712 );
713 let session_id = sid.unwrap_or_default();
714 hook.execute_plan(plan, &session_id, msg_index).await;
715 }
716 }
717 Err(e) => tracing::warn!(error = %e, "PersistenceHook evaluate failed"),
718 }
719 });
720 }
721
722 Ok(result)
723 }
724}
725
726#[allow(clippy::too_many_arguments)]
731async fn run_agent(
732 config: &AgentRuntimeConfig,
733 engine: &OxiosEngine,
734 kernel_handle: Arc<KernelHandle>,
735 system_prompt: String,
736 prompt: String,
737 exec_id: uuid::Uuid,
738 goal: String,
739 agent_id: AgentId,
740 cspace: crate::capability::CSpace,
741 audit_trail: Option<Arc<AuditTrail>>,
742 routing_stats: Option<Arc<crate::kernel_handle::RoutingStats>>,
743 session_id: Option<String>,
744 mount_paths: &[std::path::PathBuf],
745 restore_state: Option<&serde_json::Value>,
746) -> Result<(
747 String,
748 usize,
749 bool,
750 Vec<oxios_memory::memory::sona::TrajectoryStep>,
751 Arc<Agent>,
752 Vec<String>,
753 std::collections::HashMap<String, String>,
754 std::collections::HashMap<String, bool>,
755 std::collections::HashMap<String, chrono::DateTime<chrono::Utc>>,
756 u64,
757 u64,
758 String,
759 Vec<oxios_ouroboros::ReasoningSegment>,
760)> {
761 let workspace = if !mount_paths.is_empty() {
767 mount_paths[0].clone()
768 } else if let Some(ws) = &config.workspace_dir {
769 ws.clone()
770 } else {
771 std::env::temp_dir()
772 .join("oxios-agent-workspace")
773 .join(agent_id.to_string())
774 };
775
776 let _ = std::fs::create_dir_all(&workspace);
778
779 tracing::debug!(workspace = %workspace.display(), "Agent workspace scoped");
780
781 {
798 use crate::access_manager::{Role, Subject};
799 let agent_name = format!("agent-{agent_id}");
800 let mut am = kernel_handle.exec.access_manager().lock();
801 let perms = am.get_or_create_permissions(&agent_name);
802
803 if let Ok(cwd) = std::env::current_dir() {
805 let cwd_pattern = format!("{}/**", cwd.to_string_lossy().trim_end_matches('/'));
806 if !perms.allowed_paths.iter().any(|p| p == &cwd_pattern) {
807 perms.allow_path(&cwd_pattern);
808 tracing::debug!(
809 agent = %agent_name,
810 path = %cwd_pattern,
811 "Added CWD to agent allowed paths"
812 );
813 }
814 }
815
816 let ws_pattern = format!("{}/**", workspace.to_string_lossy().trim_end_matches('/'));
818 if !perms.allowed_paths.iter().any(|p| p == &ws_pattern) {
819 perms.allow_path(&ws_pattern);
820 }
821
822 for mount_path in mount_paths {
827 let pattern = format!("{}/**", mount_path.to_string_lossy().trim_end_matches('/'));
828 if !perms.allowed_paths.iter().any(|p| p == &pattern) {
829 perms.allow_path(&pattern);
830 tracing::debug!(
831 agent = %agent_name,
832 path = %pattern,
833 "Added Mount path to agent allowed paths (RFC-025)"
834 );
835 }
836 }
837
838 let kernel_ws = kernel_handle
840 .state
841 .workspace_path()
842 .to_string_lossy()
843 .to_string();
844 let kernel_ws_pattern = format!("{}/**", kernel_ws.trim_end_matches('/'));
845 if kernel_ws_pattern != ws_pattern
846 && !perms.allowed_paths.iter().any(|p| p == &kernel_ws_pattern)
847 {
848 perms.allow_path(&kernel_ws_pattern);
849 }
850
851 if !perms.allowed_paths.iter().any(|p| p == "/tmp/**") {
853 perms.allow_path("/tmp/**");
854 }
855
856 let rbac_subject = Subject::Agent(agent_id);
858 am.rbac_manager_mut()
859 .assign_role(rbac_subject, Role::Superuser);
860 }
861
862 let _trace_guard = crate::observability::tracer().start(
864 format!("exec-{}", &exec_id.to_string()[..8]).as_str(),
865 oxi_sdk::SpanKind::Agent,
866 );
867
868 let registry = ToolRegistry::new();
870 let search_cache = Arc::new(SearchCache::new());
871
872 let agent_context = AgentContext {
874 agent_id,
875 agent_name: format!("agent-{agent_id}"),
876 cspace: Arc::new(cspace.clone()),
877 };
878
879 let audit_sink: Arc<dyn crate::access_manager::AuditSink> = if let Some(trail) = audit_trail {
882 let audit_path = kernel_handle
883 .state
884 .workspace_path()
885 .join("audit")
886 .join("access.jsonl");
887 Arc::new(TrailAuditSink::new(trail, audit_path))
888 } else {
889 Arc::new(TracingAuditSink)
890 };
891 let access_gate = Arc::new(AccessGate::new(
893 kernel_handle.exec.access_manager().clone(),
894 Arc::new(kernel_handle.exec.config_snapshot()),
895 audit_sink,
896 ));
897
898 let approval_config = kernel_handle.infra.approval_config_handle();
908 let exec_snapshot = kernel_handle.exec.config_snapshot();
909 let exec_resolver = crate::approval::ExecPolicyResolver {
910 allowed_commands: Arc::new(parking_lot::RwLock::new(
911 exec_snapshot.allowed_commands.clone(),
912 )),
913 };
914 let mut dynamic_resolvers = std::collections::HashMap::new();
915 dynamic_resolvers.insert(
916 "exec".to_string(),
917 Box::new(exec_resolver) as Box<dyn crate::approval::ToolPolicyResolver>,
918 );
919 let approval_gate = Arc::new(crate::approval::ApprovalGate::with_dynamic_resolvers(
920 crate::approval::default_tool_policy_map(),
921 approval_config,
922 vec![Box::new(crate::approval::SecurityBlacklist::new(
923 crate::approval::default_blacklist_rules(),
924 ))],
925 dynamic_resolvers,
926 ));
927 let approval_event_bus = kernel_handle.infra.event_bus_clone();
928 let approval_pending = kernel_handle.infra.pending_tool_approvals();
929 let path_access_pending = kernel_handle.infra.pending_path_access();
930
931 register_tools_from_cspace_gated(
932 ®istry,
933 &kernel_handle,
934 &cspace,
935 search_cache,
936 agent_id,
937 access_gate,
938 agent_context,
939 Some(approval_gate),
940 Some(approval_event_bus),
941 Some(approval_pending),
942 Some(path_access_pending),
943 );
944
945 tracing::info!(
946 exec_id = %exec_id,
947 capabilities = cspace.len(),
948 "Tools registered from CSpace"
949 );
950
951 let agent_config = AgentConfig {
959 name: format!("agent-{agent_id}"),
960 description: None,
961 model_id: config.model_id.clone(),
962 system_prompt: Some(system_prompt.clone()),
963 timeout_seconds: 300,
964 temperature: config
965 .model_params
966 .as_ref()
967 .and_then(|p| p.temperature)
968 .or(Some(0.7)),
969 max_tokens: config
970 .model_params
971 .as_ref()
972 .and_then(|p| p.max_tokens)
973 .map(|v| v as usize)
974 .or(Some(8192)),
975 compaction_strategy: CompactionStrategy::Threshold(0.8),
976 compaction_instruction: None,
977 context_window: 128_000,
978 workspace_dir: Some(workspace.clone()),
979 output_mode: None,
980 provider_options: config.provider_options.clone(),
981 session_id: None,
982 max_tool_result_bytes: config.max_tool_result_bytes,
984 subagent_depth: 0,
989 subagent_runner: Some(
992 crate::subagent_runner::OxiosSubagentRunner::new(engine.oxi().clone())
993 .into_trait_object(),
994 ),
995 ..Default::default()
996 };
997
998 let agent = if config.provider_rpm > 0 {
1013 let resolver: Arc<dyn ProviderResolver> = Arc::new(engine.oxi().clone());
1015 let provider_name = engine.resolve_model(&config.model_id)?.provider;
1016 let provider = engine.pooled_provider(&provider_name, config.provider_rpm)?;
1017
1018 let mut pipeline = oxi_sdk::MiddlewarePipeline::new();
1020 if config.rate_limit_per_minute > 0 {
1021 pipeline = pipeline.push(oxi_sdk::middleware::builtins::RateLimitMiddleware::new(
1022 config.rate_limit_per_minute,
1023 ));
1024 }
1025 if config.token_budget > 0 {
1026 pipeline = pipeline.push(oxi_sdk::middleware::builtins::TokenBudgetMiddleware::new(
1027 config.token_budget,
1028 ));
1029 }
1030 if config.audit_tool_calls {
1031 pipeline = pipeline.push(oxi_sdk::middleware::builtins::LoggingMiddleware::new(
1032 tracing::Level::INFO,
1033 ));
1034 }
1035
1036 let agent = Arc::new(Agent::new_with_resolver(
1038 provider,
1039 agent_config,
1040 Arc::new(registry),
1041 resolver,
1042 ));
1043
1044 if !pipeline.is_empty() {
1046 let terminate_flag = Arc::new(std::sync::atomic::AtomicBool::new(false));
1047 let agent_id_for_hooks = agent_id.to_string();
1048 let hooks = oxi_sdk::middleware::build_hooks(
1049 Arc::new(pipeline),
1050 agent_id_for_hooks,
1051 terminate_flag,
1052 );
1053 agent.set_hooks(hooks);
1054 }
1055
1056 agent
1057 } else {
1058 let mut builder = engine
1060 .oxi()
1061 .agent(agent_config)
1062 .workspace(&workspace)
1063 .system_prompt(system_prompt);
1064
1065 let cspace_tool_arcs: Vec<Arc<dyn oxi_sdk::AgentTool>> = registry
1075 .names()
1076 .into_iter()
1077 .filter_map(|name| registry.get(&name))
1078 .collect();
1079
1080 if let Some(auth) = engine.authorizer() {
1082 builder = builder.authorizer(auth.clone());
1083 }
1084 if let Some(tracer) = engine.tracer() {
1085 builder = builder.tracer(tracer.clone());
1086 }
1087 if let Some(ct) = engine.cost_tracker() {
1088 builder = builder.cost_tracker(ct.clone());
1089 }
1090
1091 if config.rate_limit_per_minute > 0 {
1094 builder = builder.with_rate_limit(config.rate_limit_per_minute);
1095 }
1096 if config.token_budget > 0 {
1097 builder = builder.with_token_budget(config.token_budget);
1098 }
1099 if config.audit_tool_calls {
1100 builder = builder.with_logging();
1101 }
1102
1103 let built = builder.build()?;
1104 let agent = Arc::new(built);
1105
1106 let agent_tools = agent.tools();
1111 for tool in cspace_tool_arcs {
1112 agent_tools.register_arc(tool);
1113 }
1114
1115 agent
1116 };
1117
1118 if let Some(state) = restore_state {
1122 agent.import_state(state.clone()).unwrap_or_else(|e| {
1123 tracing::warn!(agent_id = %agent_id, error = %e, "Failed to restore agent state");
1124 });
1125 }
1126
1127 let exec_state = Arc::new(Mutex::new(ExecuteState::default()));
1129 let exec_state_cb = Arc::clone(&exec_state);
1130 let memory_for_callback: Arc<MemoryManager> = (*kernel_handle.agents.memory_manager()).clone();
1131 let session_id_for_callback = exec_id.to_string();
1132 let model_id_for_callback = config.model_id.clone();
1133 let agent_id_for_callback = agent_id.to_string();
1134 let routing_stats_for_cb = routing_stats.clone();
1135 let transparency_session: Option<String> = session_id.clone();
1138 let kernel_handle_for_cb: Arc<KernelHandle> = Arc::clone(&kernel_handle);
1139 let streaming_sinks_for_cb: Arc<crate::streaming_sink::StreamingSinkRegistry> =
1143 Arc::clone(&kernel_handle.streaming_sinks);
1144 let mut sent_model_for_cb: bool = false;
1146 let result = agent
1147 .run_streaming(prompt, move |event| {
1148 if !sent_model_for_cb
1149 && let Some(ref sid) = transparency_session
1150 && !model_id_for_callback.is_empty()
1151 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1152 {
1153 let _ = tx.try_send(StreamDelta::Model(model_id_for_callback.clone()));
1154 sent_model_for_cb = true;
1155 }
1156 let mut s = exec_state_cb.lock();
1157 match event {
1158 AgentEvent::ToolExecutionStart {
1159 tool_name,
1160 tool_call_id,
1161 args,
1162 context,
1163 ..
1164 } => {
1165 let idx = s.trajectory_steps.len();
1167 s.pending_tools
1168 .insert(tool_call_id.clone(), (std::time::Instant::now(), idx));
1169 s.tool_args_map.insert(
1170 tool_call_id.clone(),
1171 serde_json::to_string(&args).unwrap_or_default(),
1172 );
1173 s.tool_timestamps
1174 .insert(tool_call_id.clone(), chrono::Utc::now());
1175 s.tool_call_ids.push(tool_call_id.clone());
1176 s.trajectory_steps
1177 .push(oxios_memory::memory::sona::TrajectoryStep {
1178 input: tool_name.clone(),
1179 output: String::new(),
1180 duration_ms: 0,
1181 confidence: 0.0,
1182 });
1183 if let Some(ref sid) = transparency_session {
1185 let context_json = context
1186 .as_ref()
1187 .map(serde_json::to_value)
1188 .transpose()
1189 .unwrap_or(None);
1190 let _ =
1191 kernel_handle_for_cb
1192 .infra
1193 .publish(KernelEvent::ToolExecutionStarted {
1194 session_id: sid.clone(),
1195 tool_name: tool_name.clone(),
1196 tool_call_id: tool_call_id.clone(),
1197 tool_args: args.clone(),
1198 context: context_json,
1199 });
1200 }
1201 }
1202 AgentEvent::ToolExecutionUpdate {
1203 tool_call_id,
1204 tool_name,
1205 partial_result,
1206 tab_id,
1207 context,
1208 } => {
1209 if let Some(ref sid) = transparency_session {
1219 let context_json = context
1220 .as_ref()
1221 .map(serde_json::to_value)
1222 .transpose()
1223 .unwrap_or(None);
1224 let _ = kernel_handle_for_cb.infra.publish(
1225 KernelEvent::ToolExecutionProgress {
1226 session_id: sid.clone(),
1227 tool_call_id: tool_call_id.clone(),
1228 tool_name: tool_name.clone(),
1229 progress: partial_result,
1230 tab_id,
1231 context: context_json,
1232 },
1233 );
1234 }
1235 }
1236 AgentEvent::ToolExecutionEnd {
1237 tool_name,
1238 tool_call_id,
1239 is_error,
1240 result,
1241 ..
1242 } => {
1243 if !is_error {
1244 s.steps_completed += 1;
1245 }
1246 let mut duration_ms: u64 = 0;
1248 let mut summary = String::new();
1249 if let Some((start, idx)) = s.pending_tools.remove(tool_call_id.as_str()) {
1250 duration_ms = start.elapsed().as_millis() as u64;
1251 if let Some(step) = s.trajectory_steps.get_mut(idx) {
1252 summary = summarize_tool_result(&result.content, 200);
1253 step.output = summary.clone();
1254 step.duration_ms = duration_ms;
1255 step.confidence = if is_error { 0.3 } else { 0.8 };
1256 }
1257 }
1258 s.tool_error_map.insert(tool_call_id.clone(), is_error);
1259 if let Some(ref sid) = transparency_session {
1261 let _ = kernel_handle_for_cb.infra.publish(
1262 KernelEvent::ToolExecutionFinished {
1263 session_id: sid.clone(),
1264 tool_call_id: tool_call_id.clone(),
1265 tool_name: tool_name.clone(),
1266 duration_ms,
1267 is_error,
1268 output_summary: summary,
1269 },
1270 );
1271 }
1272 }
1273 AgentEvent::AgentEnd {
1274 messages,
1275 stop_reason,
1276 ..
1277 } => {
1278 if let Some(oxi_sdk::Message::Assistant(a)) = messages.last() {
1279 s.final_content = a.text_content();
1280 }
1281 s.success = matches!(stop_reason.as_deref(), Some("Stop") | Some("ToolUse"));
1287 }
1288 AgentEvent::Error { message, .. } => {
1289 s.final_content = message.clone();
1290 s.success = false;
1291 }
1292 AgentEvent::Usage {
1293 input_tokens,
1294 output_tokens,
1295 } => {
1296 s.total_input_tokens += input_tokens as u64;
1298 s.total_output_tokens += output_tokens as u64;
1299
1300 let agent_label = format!("agent-{agent_id_for_callback}");
1302 crate::observability::cost_tracker().record(
1303 &agent_label,
1304 &oxi_sdk::Model::new(
1305 &model_id_for_callback,
1306 &model_id_for_callback,
1307 oxi_sdk::Api::OpenAiCompletions,
1308 "unknown",
1309 "https://unknown.com",
1310 ),
1311 oxi_sdk::TokenUsage {
1312 input: input_tokens as u64,
1313 output: output_tokens as u64,
1314 cache_read: 0,
1315 cache_write: 0,
1316 },
1317 );
1318
1319 if let Some(stats) = &routing_stats_for_cb {
1321 let cost = crate::kernel_handle::engine_api::estimate_cost(
1322 &model_id_for_callback,
1323 input_tokens as u64,
1324 output_tokens as u64,
1325 );
1326 stats.record_model_usage(&model_id_for_callback, cost);
1327 }
1328 if let Some(ref sid) = transparency_session {
1330 let _ = kernel_handle_for_cb
1331 .infra
1332 .publish(KernelEvent::TokenUsageUpdate {
1333 session_id: sid.clone(),
1334 input_tokens: input_tokens as u64,
1335 output_tokens: output_tokens as u64,
1336 });
1337 }
1338 }
1339 AgentEvent::Compaction {
1340 event: CompactionEvent::Completed { result, .. },
1341 } => {
1342 handle_compaction(
1343 result.summary.clone(),
1344 session_id_for_callback.clone(),
1345 memory_for_callback.clone(),
1346 );
1347 if let Some(ref sid) = transparency_session {
1349 let _ =
1350 kernel_handle_for_cb
1351 .infra
1352 .publish(KernelEvent::ReasoningFragment {
1353 session_id: sid.clone(),
1354 content: result.summary.clone(),
1355 source: "compaction".to_string(),
1356 });
1357 }
1358 }
1359 AgentEvent::Compaction {
1360 event: CompactionEvent::Triggered { source, .. },
1361 } => {
1362 if let Some(ref sid) = transparency_session {
1368 let _ =
1369 kernel_handle_for_cb
1370 .infra
1371 .publish(KernelEvent::CompactionTriggered {
1372 session_id: Some(sid.clone()),
1373 source,
1374 });
1375 } else {
1376 let _ =
1377 kernel_handle_for_cb
1378 .infra
1379 .publish(KernelEvent::CompactionTriggered {
1380 session_id: None,
1381 source,
1382 });
1383 }
1384 }
1385 AgentEvent::TextChunk { text } => {
1386 if let Some(ref sid) = transparency_session
1400 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1401 {
1402 let _ = tx.try_send(StreamDelta::Text(text.clone()));
1403 }
1404 }
1405 AgentEvent::Thinking => {
1406 if let Some(ref sid) = transparency_session
1411 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1412 {
1413 let _ = tx.try_send(StreamDelta::Thinking);
1414 }
1415 }
1416 AgentEvent::ThinkingDelta { text } => {
1417 const REASONING_CAP: usize = 4096;
1431 if s.reasoning_text.len() < REASONING_CAP {
1432 s.reasoning_text.push_str(&text);
1433 if s.reasoning_text.len() > REASONING_CAP {
1434 s.reasoning_text.truncate(REASONING_CAP);
1435 }
1436 }
1437 if s.reasoning_bytes < REASONING_CAP {
1443 let budget = REASONING_CAP - s.reasoning_bytes;
1444 let mut chunk = String::new();
1445 for ch in text.chars() {
1446 if chunk.len() + ch.len_utf8() > budget {
1447 break;
1448 }
1449 chunk.push(ch);
1450 }
1451 if !chunk.is_empty() {
1452 s.reasoning_bytes += chunk.len();
1453 let pos = s.trajectory_steps.len();
1454 if s.reasoning_segments
1455 .last()
1456 .is_some_and(|l| l.before_step == pos)
1457 {
1458 s.reasoning_segments
1459 .last_mut()
1460 .unwrap()
1461 .text
1462 .push_str(&chunk);
1463 } else {
1464 s.reasoning_segments
1465 .push(oxios_ouroboros::ReasoningSegment {
1466 before_step: pos,
1467 text: chunk,
1468 });
1469 }
1470 }
1471 }
1472 if let Some(ref sid) = transparency_session
1473 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1474 {
1475 let _ = tx.try_send(StreamDelta::ThinkingDelta(text.clone()));
1476 }
1477 }
1478 AgentEvent::ThinkingEnd => {
1479 if let Some(ref sid) = transparency_session
1485 && let Some(tx) = streaming_sinks_for_cb.lookup(sid)
1486 {
1487 let _ = tx.try_send(StreamDelta::ThinkingEnd);
1488 }
1489 }
1490 AgentEvent::ToolCallDelta {
1491 tool_call_id,
1492 args_delta,
1493 } => {
1494 if let Some(ref sid) = transparency_session {
1499 let _ = kernel_handle_for_cb
1500 .infra
1501 .publish(KernelEvent::ToolArgsDelta {
1502 session_id: sid.clone(),
1503 tool_call_id: tool_call_id.clone(),
1504 args_delta: args_delta.clone(),
1505 });
1506 }
1507 }
1508 _ => {}
1509 }
1510 })
1511 .await;
1512
1513 let circuit = get_llm_circuit_breaker();
1515 if result.is_err() {
1516 circuit.record_failure();
1517 crate::metrics::get_metrics()
1518 .llm_circuit_breaker_state
1519 .set(1.0);
1520 } else {
1521 circuit.record_success();
1522 crate::metrics::get_metrics()
1523 .llm_circuit_breaker_state
1524 .set(0.0);
1525 }
1526
1527 if let Err(e) = result {
1528 tracing::error!(exec_id = %exec_id, error = %e, "Agent failed");
1529 let restore_state = agent.export_state().ok();
1534 return Err(crate::resilience::AgentRunError::wrap(e, restore_state).into());
1535 }
1536
1537 let s = exec_state.lock();
1538 tracing::info!(
1539 exec_id = %exec_id,
1540 steps = s.steps_completed,
1541 success = s.success,
1542 "Agent completed"
1543 );
1544
1545 if !s.trajectory_steps.is_empty()
1548 && let Some(sona) = kernel_handle.agents.memory_manager().sona_engine()
1549 {
1550 let steps = s.trajectory_steps.clone();
1551 let success = s.success;
1552 let sona = Arc::clone(sona);
1553 let domain = infer_domain(&goal);
1554 tokio::spawn(async move {
1555 let verdict = if success {
1556 oxios_memory::memory::sona::Verdict::Success
1557 } else {
1558 oxios_memory::memory::sona::Verdict::Failure
1559 };
1560 let trajectory = oxios_memory::memory::sona::Trajectory::new(steps, verdict, &domain);
1561 if let Err(e) = sona.record(trajectory).await {
1562 tracing::debug!(error = %e, "SONA trajectory recording failed (non-fatal)");
1563 }
1564 });
1565 }
1566
1567 Ok((
1568 s.final_content.clone(),
1569 s.steps_completed,
1570 s.success,
1571 s.trajectory_steps.clone(),
1572 agent,
1573 s.tool_call_ids.clone(),
1574 s.tool_args_map.clone(),
1575 s.tool_error_map.clone(),
1576 s.tool_timestamps.clone(),
1577 s.total_input_tokens,
1578 s.total_output_tokens,
1579 s.reasoning_text.clone(),
1580 s.reasoning_segments.clone(),
1581 ))
1582}
1583
1584fn summarize_tool_result(result: &str, max_len: usize) -> String {
1589 let trimmed = result.trim();
1590 if trimmed.chars().count() <= max_len {
1591 return trimmed.to_string();
1592 }
1593 let first_line = trimmed.lines().next().unwrap_or("");
1595 if first_line.chars().count() <= max_len {
1596 first_line.to_string()
1597 } else {
1598 let take = max_len.saturating_sub(3);
1599 let truncated: String = if take == 0 {
1600 first_line.chars().take(max_len).collect()
1601 } else {
1602 first_line.chars().take(take).collect()
1603 };
1604 format!("{truncated}...")
1605 }
1606}
1607fn truncate_json_str(json_str: &str, max_len: usize) -> String {
1608 if json_str.len() <= max_len {
1609 return json_str.to_string();
1610 }
1611 let take = max_len.saturating_sub(3);
1614 if take == 0 {
1615 return json_str.chars().take(max_len).collect();
1616 }
1617 let truncated: String = json_str.chars().take(take).collect();
1618 format!("{truncated}...")
1619}
1620
1621fn infer_domain(goal: &str) -> String {
1626 let lower = goal.to_lowercase();
1627 let keywords: Vec<&str> = lower.split_whitespace().take(8).collect();
1628
1629 if keywords.iter().any(|k| {
1631 [
1632 "test",
1633 "tests",
1634 "spec",
1635 "testing",
1636 "assert",
1637 "unit test",
1638 "integration",
1639 ]
1640 .contains(k)
1641 }) {
1642 return "testing".to_string();
1643 }
1644 if keywords
1645 .iter()
1646 .any(|k| ["deploy", "release", "publish", "ship"].contains(k))
1647 {
1648 return "deployment".to_string();
1649 }
1650 if keywords
1651 .iter()
1652 .any(|k| ["fix", "bug", "patch", "repair", "debug"].contains(k))
1653 {
1654 return "bugfix".to_string();
1655 }
1656 if keywords
1657 .iter()
1658 .any(|k| ["refactor", "restructure", "reorganize", "rewrite"].contains(k))
1659 {
1660 return "refactoring".to_string();
1661 }
1662 if keywords
1663 .iter()
1664 .any(|k| ["doc", "document", "readme", "guide", "explain"].contains(k))
1665 {
1666 return "documentation".to_string();
1667 }
1668 if keywords
1669 .iter()
1670 .any(|k| ["build", "create", "implement", "add", "make", "new"].contains(k))
1671 {
1672 return "development".to_string();
1673 }
1674 if keywords
1675 .iter()
1676 .any(|k| ["analyze", "review", "audit", "inspect", "check"].contains(k))
1677 {
1678 return "analysis".to_string();
1679 }
1680 if keywords
1681 .iter()
1682 .any(|k| ["config", "setup", "install", "configure", "init"].contains(k))
1683 {
1684 return "configuration".to_string();
1685 }
1686
1687 let meaningful: Vec<&str> = lower
1689 .split_whitespace()
1690 .filter(|w| w.len() > 2)
1691 .take(2)
1692 .collect();
1693 if meaningful.len() >= 2 {
1694 meaningful.join("_")
1695 } else {
1696 "general".to_string()
1697 }
1698}
1699
1700fn handle_compaction(summary: String, session_id: String, memory_manager: Arc<MemoryManager>) {
1706 let entry = MemoryEntry {
1707 id: uuid::Uuid::new_v4().to_string(),
1708 memory_type: MemoryType::Conversation,
1709 tier: crate::memory::MemoryTier::Warm,
1710 content: summary,
1711 content_hash: 0,
1712 source: "compaction".to_string(),
1713 session_id: Some(session_id),
1714 tags: vec![],
1715 importance: 0.5,
1716 pinned: false,
1717 protection: crate::memory::ProtectionLevel::None,
1718 auto_classified: false,
1719 session_appearances: 0,
1720 user_corrected: false,
1721 seen_in_sessions: vec![],
1722 created_at: chrono::Utc::now(),
1723 accessed_at: chrono::Utc::now(),
1724 modified_at: chrono::Utc::now(),
1725 access_count: 0,
1726 decay_score: 1.0,
1727 compaction_level: 0,
1728 compacted_from: vec![],
1729 related_ids: vec![],
1730 contradicts: None,
1731 };
1732 tokio::spawn(async move {
1733 if let Err(e) = memory_manager.remember(entry).await {
1734 tracing::warn!(error = %e, "Failed to save compaction summary");
1735 }
1736 });
1737}
1738
1739#[allow(dead_code)]
1744fn build_directive_system_prompt(
1745 directive: &Directive,
1746 env: &ExecEnv,
1747 persona_prompt: Option<&str>,
1748 capabilities_xml: Option<&str>,
1749 kernel_manifest: Option<&str>,
1750) -> String {
1751 build_system_prompt_inner(
1752 &directive.goal,
1753 &directive.original_request,
1754 &directive.constraints,
1755 &directive.acceptance_criteria,
1756 env.workspace_context.as_deref(),
1757 persona_prompt,
1758 capabilities_xml,
1759 kernel_manifest,
1760 )
1761}
1762
1763const ARTIFACT_PROTOCOL: &str = "\n\n\
1770 ## Artifacts\n\
1771 When you produce substantial, self-contained content the user will want to\n\
1772 view or interact with separately — a complete HTML page, an SVG graphic, a\n\
1773 Mermaid diagram, or an interactive React component — wrap it in an artifact\n\
1774 tag so the UI shows a live preview panel:\n\n\
1775 <lobeArtifact type=\"...\" title=\"...\" identifier=\"...\">\n\
1776 ...the full content...\n\
1777 </lobeArtifact>\n\n\
1778 Use exactly one of these `type` values (others are not recognised):\n\
1779 - `text/html` → an HTML document or fragment\n\
1780 - `image/svg+xml` → an SVG graphic\n\
1781 - `application/lobe.artifacts.mermaid` → a Mermaid diagram\n\
1782 - `application/lobe.artifacts.react` → a React component (JSX/TSX)\n\n\
1783 - `title` — a short human title for the panel.\n\
1784 - `identifier` — a unique kebab-case id, e.g. `sales-dashboard`.\n\
1785 Put the full, runnable content INSIDE the tag (not in a separate fence).\n\n\
1786 Do NOT wrap non-visual code. Shell commands, Python, Rust, JSON, config, or\n\
1787 short snippets that are part of an explanation belong in a normal fenced\n\
1788 code block. Use an artifact only for content the user would open to view,\n\
1789 not copy-and-paste. Limit one artifact per self-contained piece.\n";
1790#[allow(clippy::too_many_arguments)]
1796fn build_system_prompt_inner(
1797 goal: &str,
1798 original_request: &str,
1799 constraints: &[String],
1800 acceptance_criteria: &[String],
1801 workspace_context: Option<&str>,
1802 persona_prompt: Option<&str>,
1803 capabilities_xml: Option<&str>,
1804 kernel_manifest: Option<&str>,
1805) -> String {
1806 let mut prompt = String::from(
1807 "You are an autonomous agent in the Oxios operating system.\n\
1808 You execute Seeds — immutable specifications with goals, constraints, and\n\
1809 acceptance criteria.\n\n\
1810 ## Available Tools\n\
1811 You have the following tools:\n\
1812 - **File tools**: read, write, edit files; grep, find, ls for searching\n\
1813 - **Web tools**: web_search for searching the web, get_search_results for retrieving cached results\n\
1814 - **Exec**: run shell commands\n\
1815 - **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.
1816 - **Knowledge**: knowledge — personal markdown vault for documents and notes\n\
1817 - **Kernel tools**: agent, project, persona, cron, security, budget, resource\n\n\
1818 **Important**: When the task involves fetching information from the internet,\n\
1819 websites, or online services, use `web_search` first — do NOT search local files.\n\
1820 When the task asks to \"get\", \"fetch\", \"find online\", or \"look up\" something\n\
1821 from the web, use `web_search`.\n",
1822 );
1823 prompt.push_str(&format!("\n## Goal\n{}\n", goal));
1824
1825 if !original_request.is_empty() && original_request != goal {
1828 prompt.push_str(&format!(
1829 "\n## User's Original Request\n{}\n",
1830 original_request
1831 ));
1832 }
1833
1834 if !constraints.is_empty() {
1835 prompt.push_str("\n## Constraints\n");
1836 for (i, c) in constraints.iter().enumerate() {
1837 prompt.push_str(&format!("{}. {}\n", i + 1, c));
1838 }
1839 }
1840
1841 if !acceptance_criteria.is_empty() {
1842 prompt.push_str("\n## Acceptance Criteria\n");
1843 for (i, c) in acceptance_criteria.iter().enumerate() {
1844 prompt.push_str(&format!("{}. {}\n", i + 1, c));
1845 }
1846 }
1847
1848 if let Some(ctx) = workspace_context.filter(|s| !s.trim().is_empty()) {
1852 prompt.push_str("\n## Workspace Context\n");
1853 prompt.push_str(ctx);
1854 prompt.push('\n');
1855 }
1856
1857 if let Some(pp) = persona_prompt {
1859 prompt.push_str("\n## Persona\n");
1860 prompt.push_str(pp);
1861 prompt.push('\n');
1862 }
1863
1864 if let Some(xml) = capabilities_xml {
1866 prompt.push_str("\n## Available Capabilities\n");
1867 prompt.push_str("The following capabilities are relevant to your goal. ");
1868 prompt.push_str("Use the `read` tool to load SKILL.md for any program.\n\n");
1869 prompt.push_str(xml);
1870 prompt.push('\n');
1871 }
1872
1873 if let Some(manifest) = kernel_manifest {
1875 prompt.push('\n');
1876 prompt.push_str(manifest);
1877 prompt.push('\n');
1878 }
1879
1880 prompt.push_str(
1882 "\n## Execution Protocol\n\
1883 1. UNDERSTAND — Read the user's request carefully. If it is a simple\n\
1884 greeting, small talk, or a question you can answer from knowledge,\n\
1885 respond naturally and conversationally — no tools needed.\n\
1886 2. PLAN — For complex tasks, outline your approach before acting.\n\
1887 3. EXECUTE — Use tools only when the task actually requires them.\n\
1888 Prefer the simplest approach. Simple requests need no tools.\n\
1889 4. VERIFY — After each action, check the result: created a file? read it back.\n\
1890 5. REPORT — Summarize how each acceptance criterion was met, with evidence.\n\n\
1891 If the request is ambiguous, use the `ask_user` tool (free-text question)\n\
1892 or the `pi-questionnaire` tool (structured choices) to clarify before\n\
1893 executing — do not guess when a single question would resolve the intent.\n\n\
1894 ## Hard Boundaries\n\
1895 - NEVER modify files outside the workspace scope\n\
1896 - NEVER execute destructive commands without confirming scope\n\
1897 - NEVER claim completion without evidence — show the output, not your opinion\n\
1898 - NEVER add features or improvements beyond the goal's scope\n\
1899 - If you cannot complete the task, say so and explain WHY\n\n\
1900 ## Scope Guard\n\
1901 The goal defines your universe. Do not:\n\
1902 - Refactor code the goal didn't mention\n\
1903 - Add tests the goal didn't require\n\
1904 - Change configuration the goal didn't specify\n\
1905 - \"Improve\" anything beyond what the acceptance criteria demand\n\n\
1906 ## Error Handling\n\
1907 - If a tool fails, read the error message carefully before retrying\n\
1908 - If a command fails, do NOT immediately retry with --force or sudo\n\
1909 - If stuck after 3 attempts, report the blocker rather than continuing to fail\n\n\
1910 ## Shape Matching\n\
1911 Match your output to the task: simple task → concise response.\n\
1912 Do not write 50 lines when 5 would do.\n\
1913 Use `exec` for all command execution (git, gh, osascript, etc.).",
1914 );
1915 prompt.push_str(ARTIFACT_PROTOCOL);
1916
1917 prompt
1918}
1919#[allow(dead_code)]
1920fn build_directive_user_prompt(directive: &Directive) -> String {
1921 build_user_prompt_inner(&directive.goal, &directive.acceptance_criteria)
1922}
1923
1924fn build_user_prompt_inner(goal: &str, acceptance_criteria: &[String]) -> String {
1926 format!(
1927 "Execute the following goal:\n\n{}\n\nAcceptance criteria:\n{}",
1928 goal,
1929 acceptance_criteria
1930 .iter()
1931 .enumerate()
1932 .map(|(i, c)| format!("{}. {}", i + 1, c))
1933 .collect::<Vec<_>>()
1934 .join("\n")
1935 )
1936}
1937
1938impl std::fmt::Debug for AgentRuntime {
1939 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1940 f.debug_struct("AgentRuntime")
1941 .field("model_id", &self.engine_handle.get().default_model_id())
1942 .finish()
1943 }
1944}
1945
1946#[cfg(test)]
1947mod tests {
1948 use super::*;
1949 use async_trait::async_trait;
1950 use oxi_sdk::{AgentTool, ToolContext, ToolError};
1951 use serde_json::Value;
1952
1953 struct DummyTool {
1955 name: String,
1956 }
1957
1958 #[async_trait]
1959 impl AgentTool for DummyTool {
1960 fn name(&self) -> &str {
1961 &self.name
1962 }
1963 fn label(&self) -> &str {
1964 &self.name
1965 }
1966 fn description(&self) -> &str {
1967 "Test tool"
1968 }
1969 fn parameters_schema(&self) -> Value {
1970 serde_json::json!({"type": "object"})
1971 }
1972
1973 async fn execute(
1974 &self,
1975 _tool_call_id: &str,
1976 _params: Value,
1977 _shutdown: Option<tokio::sync::oneshot::Receiver<()>>,
1978 _ctx: &ToolContext,
1979 ) -> Result<oxi_sdk::AgentToolResult, ToolError> {
1980 Ok(oxi_sdk::AgentToolResult::success("ok"))
1981 }
1982 }
1983
1984 #[test]
1986 fn test_requires_tools_validation_passes() {
1987 let registry = ToolRegistry::new();
1988
1989 registry.register(DummyTool {
1990 name: "read".into(),
1991 });
1992 registry.register(DummyTool {
1993 name: "exec".into(),
1994 });
1995
1996 let missing = registry.missing(&["read", "exec"]);
1997
1998 assert!(
1999 missing.is_empty(),
2000 "Expected no missing tools, got: {:?}",
2001 missing
2002 );
2003 }
2004
2005 #[test]
2007 fn test_requires_tools_validation_fails() {
2008 let registry = ToolRegistry::new();
2009
2010 registry.register(DummyTool {
2011 name: "read".into(),
2012 });
2013
2014 let missing = registry.missing(&["read", "exec", "nonexistent"]);
2015
2016 assert_eq!(missing, vec!["exec", "nonexistent"]);
2017 }
2018
2019 #[test]
2020 fn test_infer_domain_testing() {
2021 assert_eq!(infer_domain("run all unit tests for the kernel"), "testing");
2022 }
2023
2024 #[test]
2025 fn test_infer_domain_deployment() {
2026 assert_eq!(
2027 infer_domain("deploy the web service to production"),
2028 "deployment"
2029 );
2030 }
2031
2032 #[test]
2033 fn test_infer_domain_bugfix() {
2034 assert_eq!(infer_domain("fix the null pointer error in main"), "bugfix");
2035 }
2036
2037 #[test]
2038 fn test_infer_domain_development() {
2039 assert_eq!(
2040 infer_domain("create a new REST API endpoint"),
2041 "development"
2042 );
2043 }
2044
2045 #[test]
2046 fn test_infer_domain_analysis() {
2047 assert_eq!(
2048 infer_domain("review the code for security issues"),
2049 "analysis"
2050 );
2051 }
2052
2053 #[test]
2054 fn test_infer_domain_fallback() {
2055 let domain = infer_domain("optimize performance metrics");
2056 assert!(!domain.is_empty());
2058 }
2059 #[test]
2060 fn test_system_prompt_includes_artifact_protocol() {
2061 let prompt = build_system_prompt_inner(
2062 "build a dashboard",
2063 "build a dashboard",
2064 &[],
2065 &[],
2066 None,
2067 None,
2068 None,
2069 None,
2070 );
2071 assert!(prompt.contains("<lobeArtifact"));
2074 assert!(prompt.contains("text/html"));
2075 assert!(prompt.contains("image/svg+xml"));
2076 assert!(prompt.contains("application/lobe.artifacts.mermaid"));
2077 assert!(prompt.contains("application/lobe.artifacts.react"));
2078 }
2079}