1#![allow(
4 missing_docs,
5 dead_code,
6 clippy::field_reassign_with_default,
7 clippy::unwrap_used,
8 clippy::let_and_return,
9 clippy::borrow_interior_mutable_const,
10 clippy::derivable_impls,
11 clippy::new_without_default,
12 unknown_lints
13)]
14pub mod bootstrap;
21pub mod cli;
22pub mod internal_urls;
23pub mod lsp;
24pub mod main_dispatch;
25pub mod mcp_credentials;
26pub mod print_mode;
27pub mod services;
28pub mod setup_wizard;
29pub mod store;
30
31pub(crate) mod app;
33pub(crate) mod context;
34pub mod discovery;
35pub mod extensions; pub(crate) mod infra;
37pub(crate) mod media;
38pub(crate) mod prompt;
39pub mod rpc_mode;
40pub(crate) mod skills;
41pub mod storage; pub use storage::packages::PackageManager;
44pub use storage::packages::ResourceKind;
45pub mod tools;
46pub(crate) mod ui;
47pub(crate) mod util;
48
49pub async fn build_oxicode_engine(
64 embedding_provider: Option<std::sync::Arc<dyn oxicode_sdk::ports::EmbeddingProvider>>,
65 hook_runner: Option<std::sync::Arc<dyn oxicode_sdk::ports::HookRunner>>,
66) -> anyhow::Result<oxicode_sdk::Oxicode> {
67 let paths = services::OxicodePaths::default_paths()?;
68 services::build_oxicode(&paths, embedding_provider, hook_runner).await
69}
70
71pub async fn run_port_check() -> anyhow::Result<()> {
78 let oxicode = build_oxicode_engine(None, None).await?;
79 let ports = oxicode.ports();
80
81 let entries = ports.state.list("").await?;
82 println!("[state] entries: {}", entries.len());
83
84 let providers = ports.auth.list_providers().await?;
86 println!("[auth] providers with credentials: {:?}", providers);
87
88 let keys = ports.config.list()?;
90 println!("[config] keys: {}", keys.len());
91
92 let skills = ports.skills.list().await?;
94 println!("[skills] {} skill(s) discovered", skills.len());
95 for s in &skills {
96 println!(" - {}: {}", s.name, s.description);
97 }
98
99 let _ = ports
101 .event_bus
102 .publish(&"port-check".to_string(), serde_json::json!({"ok": true}))
103 .await;
104 println!("[event-bus] publish ok (noop bus if not registered)");
105
106 println!("\nport check: ok");
107 Ok(())
108}
109
110#[derive(Debug, Clone)]
112pub struct CompactionContext {
113 pub messages_count: usize,
115 pub tokens_before: usize,
117 pub target_tokens: usize,
119 pub strategy: String,
121}
122
123impl CompactionContext {
124 pub fn new(
126 messages_count: usize,
127 tokens_before: usize,
128 target_tokens: usize,
129 strategy: impl Into<String>,
130 ) -> Self {
131 Self {
132 messages_count,
133 tokens_before,
134 target_tokens,
135 strategy: strategy.into(),
136 }
137 }
138
139 pub fn compression_ratio(&self) -> f32 {
141 if self.tokens_before == 0 {
142 return 1.0;
143 }
144 self.target_tokens as f32 / self.tokens_before as f32
145 }
146}
147
148use crate::store::settings::Settings;
150use anyhow::{Error, Result};
151use oxicode_agent::{Agent, AgentConfig, AgentEvent};
152use parking_lot::RwLock;
153use skills::SkillManager;
154use std::collections::VecDeque;
155use std::sync::Arc;
156
157#[derive(Clone)]
168pub struct SessionState {
169 pub should_stop: Arc<std::sync::atomic::AtomicBool>,
173 pub steering: Arc<RwLock<VecDeque<oxicode_sdk::Message>>>,
175 pub follow_up: Arc<RwLock<VecDeque<oxicode_sdk::Message>>>,
177}
178
179impl Default for SessionState {
180 fn default() -> Self {
181 Self {
182 should_stop: Arc::new(std::sync::atomic::AtomicBool::new(false)),
183 steering: Arc::new(RwLock::new(VecDeque::new())),
184 follow_up: Arc::new(RwLock::new(VecDeque::new())),
185 }
186 }
187}
188
189pub struct App {
196 oxicode: oxicode_sdk::Oxicode,
197 agent: Arc<Agent>,
198 settings: Settings,
199 skills: RwLock<SkillManager>,
200 active_skills: RwLock<Vec<String>>,
201 wasm_ext: Option<std::sync::Arc<crate::extensions::WasmExtensionManager>>,
202 ask_bridge: Option<std::sync::Arc<oxicode_agent::tools::ask::AskBridge>>,
203 issue_store: Option<crate::store::issues::FileIssueStore>,
207 ownership_session_id: String,
212 #[allow(dead_code)]
217 liveness_guard: Option<crate::store::issues::liveness::AliveGuard>,
218 persona_body: RwLock<Option<String>>,
221 session_state: SessionState,
225}
226fn build_system_prompt(
228 thinking_level: crate::store::settings::ThinkingLevel,
229 skill_contents: &[String],
230 persona_body: Option<&str>,
231) -> String {
232 let skills: Vec<prompt::system_prompt::Skill> = skill_contents
233 .iter()
234 .enumerate()
235 .map(|(i, content)| prompt::system_prompt::Skill {
236 name: format!("skill-{}", i),
237 content: content.clone(),
238 })
239 .collect();
240
241 let options = prompt::system_prompt::BuildSystemPromptOptions {
242 custom_prompt: prompt::system_prompt::thinking_level_prompt(thinking_level),
243 skills,
244 cwd: std::env::current_dir()
245 .map(|p| p.to_string_lossy().to_string())
246 .unwrap_or_default(),
247 persona_prompt: persona_body.map(|s| s.to_string()),
248 ..Default::default()
249 };
250
251 prompt::system_prompt::build_system_prompt(&options)
252}
253
254impl App {
257 pub async fn from_oxicode(
277 oxicode: oxicode_sdk::Oxicode,
278 settings: Settings,
279 ownership_session_id: String,
280 session_state: Option<SessionState>,
281 ) -> Result<Self> {
282 let session_state = session_state.unwrap_or_default();
283 let persona = match oxicode.ports().personas.get("default").await {
288 Ok(Some(p)) if !p.system_prompt.trim().is_empty() => Some(p),
289 Ok(_) => None,
290 Err(e) => {
291 tracing::warn!(error = %e, "default persona lookup failed");
292 None
293 }
294 };
295
296 let model_id = persona
297 .as_ref()
298 .and_then(|p| p.preferred_model.clone())
299 .or_else(|| settings.effective_model(None))
300 .unwrap_or_default();
301 let skills_dir = SkillManager::skills_dir().unwrap_or_else(|_| {
305 dirs::home_dir()
306 .unwrap_or_default()
307 .join(".oxicode")
308 .join("skills")
309 });
310 let skills = SkillManager::load_from_dir(&skills_dir).unwrap_or_else(|e| {
311 tracing::debug!("Skills not loaded: {}", e);
312 SkillManager::new()
313 });
314
315 let body_str = persona.as_ref().map(|p| p.system_prompt.clone());
316 let system_prompt = build_system_prompt(settings.thinking_level, &[], body_str.as_deref());
317 let compaction_strategy = if settings.auto_compaction {
318 oxicode_sdk::CompactionStrategy::Threshold(0.8)
319 } else {
320 oxicode_sdk::CompactionStrategy::Disabled
321 };
322
323 let config = AgentConfig {
324 name: "oxicode".to_string(),
325 description: Some("oxicode CLI agent".to_string()),
326 model_id: model_id.clone(),
327 system_prompt: Some(system_prompt),
328 timeout_seconds: settings.tool_timeout_seconds,
329 temperature: settings.effective_temperature(),
330 max_tokens: settings.effective_max_tokens(),
331 compaction_strategy,
332 compaction_instruction: None,
333 context_window: 128_000,
334 workspace_dir: Some(
335 std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
336 ),
337 output_mode: None,
338 provider_options: None,
339 session_id: Some(ownership_session_id.clone()),
340 ttsr_engine: None,
341 memory: None,
342 todo: None,
343 agent_pool: None,
344 url_resolver: Some(Arc::new(oxicode_sdk::SdkUrlResolver::new(
345 oxicode.ports().url_router.clone(),
346 ))),
347 lsp: if crate::lsp::manager::default_servers().is_empty() {
353 None
354 } else {
355 Some(Arc::new(crate::lsp::CliLspProvider::with_defaults(
356 std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
357 )))
358 },
359 ..Default::default()
360 };
361
362 let cwd = std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."));
371
372 let ask_timeout = if settings.ask_timeout_secs > 0 {
381 Some(std::time::Duration::from_secs(settings.ask_timeout_secs))
382 } else {
383 None
384 };
385 let bridge = std::sync::Arc::new(oxicode_agent::tools::ask::AskBridge::with_timeout(
386 ask_timeout,
387 ));
388 let mode_handle = bridge.mode_handle();
389
390 let stop_flag = Arc::clone(&session_state.should_stop);
391 let steering = Arc::clone(&session_state.steering);
392 let follow_up = Arc::clone(&session_state.follow_up);
393 let session_hooks = oxicode_sdk::agent_builder::SessionHookClosures {
394 should_stop_after_turn: Arc::new(move |_| {
395 stop_flag.load(std::sync::atomic::Ordering::SeqCst)
396 }),
397 get_steering_messages: Arc::new(move || {
398 let mut msgs: Vec<oxicode_sdk::Message> = steering.write().drain(..).collect();
399 if oxicode_agent::config::Mode::load(&mode_handle).is_auto() {
402 msgs.push(oxicode_sdk::Message::User(oxicode_sdk::UserMessage::new(
403 "Autonomy mode (auto) is active: proceed autonomously to \
404 completion without asking the user questions. Make \
405 reasonable decisions on your own and keep working.",
406 )));
407 }
408 msgs
409 }),
410 get_follow_up_messages: Arc::new(move || follow_up.write().drain(..).collect()),
411 tool_execution: oxicode_agent::config::ToolExecutionMode::Sequential,
412 };
413
414 let agent = oxicode
415 .agent(config)
416 .workspace(cwd)
417 .with_port_hooks()
418 .with_session_hooks(session_hooks)
419 .build()
420 .map_err(|e| Error::msg(format!("agent build failed: {e}")))?;
421 let agent = Arc::new(agent);
422
423 let ask_tool = oxicode_agent::tools::ask::AskTool::new(bridge.clone());
424 agent.tools().register_arc(std::sync::Arc::new(ask_tool));
425 let issue_store = std::env::current_dir()
430 .ok()
431 .map(|cwd| crate::store::issues::FileIssueStore::open_from_cwd(&cwd))
432 .and_then(|r| {
433 r.map_err(|e| tracing::warn!("issue store unavailable: {e}"))
434 .ok()
435 });
436
437 if let Some(store) = issue_store.clone() {
439 let tool = std::sync::Arc::new(crate::tools::IssueTool::new(store));
440 agent.tools().register_arc(tool);
441 }
442
443 Ok(Self {
444 oxicode,
445 agent,
446 settings,
447 skills: RwLock::new(skills),
448 active_skills: RwLock::new(Vec::new()),
449 wasm_ext: None,
450 ask_bridge: Some(bridge),
451 issue_store,
452 ownership_session_id,
453 liveness_guard: None, persona_body: RwLock::new(persona.as_ref().map(|p| p.system_prompt.clone())),
455 session_state,
456 })
457 .map(|mut app| {
458 app.liveness_guard =
462 acquire_ownership_guard(app.issue_store.as_ref(), &app.ownership_session_id);
463 app
464 })
465 }
466
467 pub fn ownership_session_id(&self) -> &str {
470 &self.ownership_session_id
471 }
472
473 pub fn has_liveness_lock(&self) -> bool {
478 self.liveness_guard.is_some()
479 }
480
481 pub fn settings(&self) -> &Settings {
483 &self.settings
484 }
485
486 pub fn set_wasm_ext(
488 &mut self,
489 ext: Option<std::sync::Arc<crate::extensions::WasmExtensionManager>>,
490 ) {
491 self.wasm_ext = ext;
492 }
493
494 pub fn wasm_ext(&self) -> Option<&std::sync::Arc<crate::extensions::WasmExtensionManager>> {
496 self.wasm_ext.as_ref()
497 }
498
499 pub fn issue_store(&self) -> Option<crate::store::issues::FileIssueStore> {
501 self.issue_store.clone()
502 }
503
504 pub fn oxicode(&self) -> &oxicode_sdk::Oxicode {
507 &self.oxicode
508 }
509
510 pub fn catalog(&self) -> std::sync::Arc<dyn oxicode_sdk::ports::catalog::ModelCatalog> {
513 std::sync::Arc::clone(self.oxicode.catalog())
514 }
515
516 pub fn agent(&self) -> Arc<Agent> {
518 Arc::clone(&self.agent)
519 }
520
521 pub fn agent_tools(&self) -> Arc<oxicode_agent::ToolRegistry> {
523 self.agent.tools()
524 }
525
526 pub fn ask_bridge(&self) -> Option<&std::sync::Arc<oxicode_agent::tools::ask::AskBridge>> {
528 self.ask_bridge.as_ref()
529 }
530
531 pub fn skills(&self) -> parking_lot::RwLockReadGuard<'_, SkillManager> {
533 self.skills.read()
534 }
535
536 pub fn activate_skill(&self, name: &str) -> Result<(), String> {
538 {
539 let skills = self.skills.read();
540 if skills.get(name).is_none() {
541 return Err(format!("Skill '{}' not found", name));
542 }
543 }
544 let name_lower = name.to_lowercase();
545 {
546 let mut active = self.active_skills.write();
547 if !active.contains(&name_lower) {
548 active.push(name_lower);
549 }
550 }
551 self.rebuild_system_prompt();
552 Ok(())
553 }
554
555 pub fn deactivate_skill(&self, name: &str) {
557 let name_lower = name.to_lowercase();
558 {
559 let mut active = self.active_skills.write();
560 active.retain(|n| n != &name_lower);
561 }
562 self.rebuild_system_prompt();
563 }
564
565 pub fn active_skills(&self) -> Vec<String> {
567 self.active_skills.read().clone()
568 }
569
570 fn rebuild_system_prompt(&self) {
572 let active = self.active_skills.read();
573 let skills = self.skills.read();
574 let contents: Vec<String> = active
575 .iter()
576 .filter_map(|name| skills.get(name).map(|s| s.content.clone()))
577 .collect();
578 let persona = self.persona_body.read().clone();
582 let prompt =
583 build_system_prompt(self.settings.thinking_level, &contents, persona.as_deref());
584 self.agent.set_system_prompt(prompt);
585 }
586
587 pub fn agent_state(&self) -> oxicode_agent::AgentState {
589 self.agent.state()
590 }
591
592 pub async fn run_prompt(&self, prompt: String) -> Result<String> {
594 let (response, _events) = self.agent.run(prompt).await?;
595 Ok(response.content)
596 }
597
598 pub async fn run_prompt_with_events<F>(&self, prompt: String, on_event: F) -> Result<String>
600 where
601 F: FnMut(AgentEvent) + Send + 'static,
602 {
603 self.agent.run_streaming(prompt, on_event).await?;
604 let state = self.agent_state();
605 for msg in state.messages.iter().rev() {
606 if let oxicode_sdk::Message::Assistant(a) = msg {
607 return Ok(a.text_content());
608 }
609 }
610 Ok(String::new())
611 }
612
613 pub fn reset(&self) {
615 self.agent.reset();
616 }
617
618 pub async fn switch_model(&self, model_id: &str) -> anyhow::Result<()> {
624 let _ = self.agent.switch_model(model_id);
625 Ok(())
626 }
627
628 pub fn model_id(&self) -> String {
630 self.agent.model_id()
631 }
632
633 pub fn session_state(&self) -> &SessionState {
638 &self.session_state
639 }
640
641 pub fn should_stop_flag(&self) -> Arc<std::sync::atomic::AtomicBool> {
643 Arc::clone(&self.session_state.should_stop)
644 }
645
646 pub fn steering_queue(&self) -> Arc<RwLock<VecDeque<oxicode_sdk::Message>>> {
648 Arc::clone(&self.session_state.steering)
649 }
650
651 pub fn follow_up_queue(&self) -> Arc<RwLock<VecDeque<oxicode_sdk::Message>>> {
653 Arc::clone(&self.session_state.follow_up)
654 }
655}
656
657pub(crate) fn acquire_ownership_guard(
668 issue_store: Option<&crate::store::issues::FileIssueStore>,
669 ownership_id: &str,
670) -> Option<crate::store::issues::liveness::AliveGuard> {
671 let store = issue_store?;
672 if ownership_id.is_empty() {
673 return None;
676 }
677 crate::store::issues::liveness::acquire(&store.issues_dir(), ownership_id).ok()
678}
679
680#[cfg(test)]
681mod tests {
682 use super::*;
687 use crate::store::issues::FileIssueStore;
688 use crate::store::issues::liveness;
689
690 fn tmp_store() -> (tempfile::TempDir, FileIssueStore) {
691 let tmp = tempfile::tempdir().unwrap();
692 let dir = tmp.path().join(".oxicode").join("issues");
693 std::fs::create_dir_all(&dir).unwrap();
694 (tmp, FileIssueStore::open(dir).unwrap())
695 }
696
697 #[test]
698 fn app_holds_single_liveness_lock() {
699 let (_tmp, store) = tmp_store();
703 let dir = store.issues_dir();
704 let id = "proc-test-app";
705
706 let guard = acquire_ownership_guard(Some(&store), id);
707 assert!(
708 guard.is_some(),
709 "App must acquire the liveness lock for its ownership id"
710 );
711 assert!(
712 liveness::is_session_alive(&dir, id),
713 "after acquire, the session must be live"
714 );
715
716 let second = liveness::acquire(&dir, id);
718 assert!(second.is_err(), "second acquire under same id must fail");
719
720 drop(guard);
721 assert!(
722 !liveness::is_session_alive(&dir, id),
723 "dropping App's guard releases the lock"
724 );
725 }
726
727 #[test]
728 fn acquire_returns_none_without_store() {
729 let dir = tempfile::tempdir().unwrap();
731 let id = "proc-x";
732 assert!(acquire_ownership_guard(None, id).is_none());
733 let _ = dir; }
735
736 #[test]
737 fn acquire_rejects_empty_ownership_id() {
738 let (_tmp, store) = tmp_store();
741 assert!(
742 acquire_ownership_guard(Some(&store), "").is_none(),
743 "empty ownership id must never acquire a lock (#13 guard)"
744 );
745 }
746}
747pub mod symbols;
748pub mod tui_vt;