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 foundation;
23pub mod internal_urls;
24pub mod lsp;
25pub mod main_dispatch;
26pub mod mcp_credentials;
27pub mod oauth_listener;
28pub mod oauth_refresh;
29pub mod print_mode;
30pub mod provider_oauth;
31pub mod services;
32pub mod setup_wizard;
33pub mod store;
34pub(crate) mod app;
36pub(crate) mod context;
37pub mod discovery;
38pub mod extensions; pub(crate) mod infra;
40pub(crate) mod media;
41pub(crate) mod prompt;
42pub mod rpc_mode;
43pub(crate) mod skills;
44pub mod storage; pub use storage::packages::PackageManager;
47pub use storage::packages::ResourceKind;
48pub mod tools;
49pub(crate) mod ui;
50pub(crate) mod util;
51
52pub async fn build_oxicode_engine(
67 embedding_provider: Option<std::sync::Arc<dyn oxicode_sdk::ports::EmbeddingProvider>>,
68 hook_runner: Option<std::sync::Arc<dyn oxicode_sdk::ports::HookRunner>>,
69) -> anyhow::Result<oxicode_sdk::Oxicode> {
70 let paths = services::OxicodePaths::default_paths()?;
71 services::build_oxicode(&paths, embedding_provider, hook_runner).await
72}
73
74pub async fn run_port_check() -> anyhow::Result<()> {
81 let oxicode = build_oxicode_engine(None, None).await?;
82 let ports = oxicode.ports();
83
84 let entries = ports.state.list("").await?;
85 println!("[state] entries: {}", entries.len());
86
87 let providers = ports.auth.list_providers().await?;
89 println!("[auth] providers with credentials: {:?}", providers);
90
91 let keys = ports.config.list()?;
93 println!("[config] keys: {}", keys.len());
94
95 let skills = ports.skills.list().await?;
97 println!("[skills] {} skill(s) discovered", skills.len());
98 for s in &skills {
99 println!(" - {}: {}", s.name, s.description);
100 }
101
102 let _ = ports
104 .event_bus
105 .publish(&"port-check".to_string(), serde_json::json!({"ok": true}))
106 .await;
107 println!("[event-bus] publish ok (noop bus if not registered)");
108
109 println!("\nport check: ok");
110 Ok(())
111}
112
113#[derive(Debug, Clone)]
115pub struct CompactionContext {
116 pub messages_count: usize,
118 pub tokens_before: usize,
120 pub target_tokens: usize,
122 pub strategy: String,
124}
125
126impl CompactionContext {
127 pub fn new(
129 messages_count: usize,
130 tokens_before: usize,
131 target_tokens: usize,
132 strategy: impl Into<String>,
133 ) -> Self {
134 Self {
135 messages_count,
136 tokens_before,
137 target_tokens,
138 strategy: strategy.into(),
139 }
140 }
141
142 pub fn compression_ratio(&self) -> f32 {
144 if self.tokens_before == 0 {
145 return 1.0;
146 }
147 self.target_tokens as f32 / self.tokens_before as f32
148 }
149}
150
151use crate::store::settings::Settings;
153use anyhow::{Error, Result};
154use oxicode_agent::{Agent, AgentConfig, AgentEvent};
155use parking_lot::RwLock;
156use skills::SkillManager;
157use std::collections::VecDeque;
158use std::sync::Arc;
159
160#[derive(Clone)]
171pub struct SessionState {
172 pub should_stop: Arc<std::sync::atomic::AtomicBool>,
176 pub steering: Arc<RwLock<VecDeque<oxicode_sdk::Message>>>,
178 pub follow_up: Arc<RwLock<VecDeque<oxicode_sdk::Message>>>,
180}
181
182impl Default for SessionState {
183 fn default() -> Self {
184 Self {
185 should_stop: Arc::new(std::sync::atomic::AtomicBool::new(false)),
186 steering: Arc::new(RwLock::new(VecDeque::new())),
187 follow_up: Arc::new(RwLock::new(VecDeque::new())),
188 }
189 }
190}
191
192pub struct App {
199 oxicode: oxicode_sdk::Oxicode,
200 agent: Arc<Agent>,
201 settings: Settings,
202 skills: RwLock<SkillManager>,
203 active_skills: RwLock<Vec<String>>,
204 wasm_ext: Option<std::sync::Arc<crate::extensions::WasmExtensionManager>>,
205 ask_bridge: Option<std::sync::Arc<oxicode_agent::tools::ask::AskBridge>>,
206 issue_store: Option<oxicode_sdk::FileIssueStore>,
210 ownership_session_id: String,
215 #[allow(dead_code)]
220 liveness_guard: Option<oxicode_sdk::liveness::AliveGuard>,
221 persona_body: RwLock<Option<String>>,
224 session_state: SessionState,
228}
229fn build_system_prompt(
231 thinking_level: crate::store::settings::ThinkingLevel,
232 skill_contents: &[String],
233 persona_body: Option<&str>,
234) -> String {
235 let skills: Vec<prompt::system_prompt::Skill> = skill_contents
236 .iter()
237 .enumerate()
238 .map(|(i, content)| prompt::system_prompt::Skill {
239 name: format!("skill-{}", i),
240 content: content.clone(),
241 })
242 .collect();
243
244 let options = prompt::system_prompt::BuildSystemPromptOptions {
245 custom_prompt: prompt::system_prompt::thinking_level_prompt(thinking_level),
246 skills,
247 cwd: std::env::current_dir()
248 .map(|p| p.to_string_lossy().to_string())
249 .unwrap_or_default(),
250 persona_prompt: persona_body.map(|s| s.to_string()),
251 ..Default::default()
252 };
253
254 prompt::system_prompt::build_system_prompt(&options)
255}
256
257impl App {
260 pub async fn from_oxicode(
280 oxicode: oxicode_sdk::Oxicode,
281 settings: Settings,
282 ownership_session_id: String,
283 session_state: Option<SessionState>,
284 ) -> Result<Self> {
285 let session_state = session_state.unwrap_or_default();
286 let persona = match oxicode.ports().personas.get("default").await {
291 Ok(Some(p)) if !p.system_prompt.trim().is_empty() => Some(p),
292 Ok(_) => None,
293 Err(e) => {
294 tracing::warn!(error = %e, "default persona lookup failed");
295 None
296 }
297 };
298
299 let model_id = persona
300 .as_ref()
301 .and_then(|p| p.preferred_model.clone())
302 .or_else(|| settings.effective_model(None))
303 .unwrap_or_default();
304 let skills_dir = SkillManager::skills_dir().unwrap_or_else(|_| {
308 oxicode_catalog::product_env::home_dir()
309 .unwrap_or_default()
310 .join("skills")
311 });
312 let skills = SkillManager::load_from_dir(&skills_dir).unwrap_or_else(|e| {
313 tracing::debug!("Skills not loaded: {}", e);
314 SkillManager::new()
315 });
316
317 let body_str = persona.as_ref().map(|p| p.system_prompt.clone());
318 let system_prompt = build_system_prompt(settings.thinking_level, &[], body_str.as_deref());
319 let compaction_strategy = if settings.auto_compaction {
320 oxicode_sdk::CompactionStrategy::Threshold(0.8)
321 } else {
322 oxicode_sdk::CompactionStrategy::Disabled
323 };
324
325 let config = AgentConfig {
326 name: "oxicode".to_string(),
327 description: Some("oxicode CLI agent".to_string()),
328 model_id: model_id.clone(),
329 system_prompt: Some(system_prompt),
330 timeout_seconds: settings.tool_timeout_seconds,
331 temperature: settings.effective_temperature(),
332 max_tokens: settings.effective_max_tokens(),
333 compaction_strategy,
334 compaction_instruction: None,
335 context_window: 128_000,
336 workspace_dir: Some(
337 std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
338 ),
339 output_mode: None,
340 provider_options: None,
341 session_id: Some(ownership_session_id.clone()),
342 ttsr_engine: None,
343 memory: None,
344 todo: None,
345 agent_pool: None,
346 url_resolver: Some(Arc::new(oxicode_sdk::SdkUrlResolver::new(
347 oxicode.ports().url_router.clone(),
348 ))),
349 lsp: if crate::lsp::manager::default_servers().is_empty() {
355 None
356 } else {
357 Some(Arc::new(crate::lsp::CliLspProvider::with_defaults(
358 std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
359 )))
360 },
361 ..Default::default()
362 };
363
364 let cwd = std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."));
373
374 let ask_timeout = if settings.ask_timeout_secs > 0 {
383 Some(std::time::Duration::from_secs(settings.ask_timeout_secs))
384 } else {
385 None
386 };
387 let bridge = std::sync::Arc::new(oxicode_agent::tools::ask::AskBridge::with_timeout(
388 ask_timeout,
389 ));
390 let mode_handle = bridge.mode_handle();
391
392 let stop_flag = Arc::clone(&session_state.should_stop);
393 let steering = Arc::clone(&session_state.steering);
394 let follow_up = Arc::clone(&session_state.follow_up);
395 let session_hooks = oxicode_sdk::agent_builder::SessionHookClosures {
396 should_stop_after_turn: Arc::new(move |_| {
397 stop_flag.load(std::sync::atomic::Ordering::SeqCst)
398 }),
399 get_steering_messages: Arc::new(move || {
400 let mut msgs: Vec<oxicode_sdk::Message> = steering.write().drain(..).collect();
401 if oxicode_agent::config::Mode::load(&mode_handle).is_auto() {
404 msgs.push(oxicode_sdk::Message::User(oxicode_sdk::UserMessage::new(
405 "Autonomy mode (auto) is active: proceed autonomously to \
406 completion without asking the user questions. Make \
407 reasonable decisions on your own and keep working.",
408 )));
409 }
410 msgs
411 }),
412 get_follow_up_messages: Arc::new(move || follow_up.write().drain(..).collect()),
413 tool_execution: oxicode_agent::config::ToolExecutionMode::Sequential,
414 };
415
416 let agent = oxicode
417 .agent(config)
418 .workspace(cwd)
419 .with_port_hooks()
420 .with_session_hooks(session_hooks)
421 .build()
422 .map_err(|e| Error::msg(format!("agent build failed: {e}")))?;
423 let agent = Arc::new(agent);
424
425 let ask_tool = oxicode_agent::tools::ask::AskTool::new(bridge.clone());
426 agent.tools().register_arc(std::sync::Arc::new(ask_tool));
427 let issue_store = std::env::current_dir()
432 .ok()
433 .map(|cwd| oxicode_sdk::FileIssueStore::open_from_cwd(&cwd))
434 .and_then(|r| {
435 r.map_err(|e| tracing::warn!("issue store unavailable: {e}"))
436 .ok()
437 });
438
439 if let Some(store) = issue_store.clone() {
441 let tool = std::sync::Arc::new(oxicode_sdk::IssueTool::new(store));
442 agent.tools().register_arc(tool);
443 }
444
445 Ok(Self {
446 oxicode,
447 agent,
448 settings,
449 skills: RwLock::new(skills),
450 active_skills: RwLock::new(Vec::new()),
451 wasm_ext: None,
452 ask_bridge: Some(bridge),
453 issue_store,
454 ownership_session_id,
455 liveness_guard: None, persona_body: RwLock::new(persona.as_ref().map(|p| p.system_prompt.clone())),
457 session_state,
458 })
459 .map(|mut app| {
460 app.liveness_guard =
464 acquire_ownership_guard(app.issue_store.as_ref(), &app.ownership_session_id);
465 app
466 })
467 }
468
469 pub fn ownership_session_id(&self) -> &str {
472 &self.ownership_session_id
473 }
474
475 pub fn has_liveness_lock(&self) -> bool {
480 self.liveness_guard.is_some()
481 }
482
483 pub fn settings(&self) -> &Settings {
485 &self.settings
486 }
487
488 pub fn set_wasm_ext(
490 &mut self,
491 ext: Option<std::sync::Arc<crate::extensions::WasmExtensionManager>>,
492 ) {
493 self.wasm_ext = ext;
494 }
495
496 pub fn wasm_ext(&self) -> Option<&std::sync::Arc<crate::extensions::WasmExtensionManager>> {
498 self.wasm_ext.as_ref()
499 }
500
501 pub fn issue_store(&self) -> Option<oxicode_sdk::FileIssueStore> {
503 self.issue_store.clone()
504 }
505
506 pub fn oxicode(&self) -> &oxicode_sdk::Oxicode {
509 &self.oxicode
510 }
511
512 pub fn catalog(&self) -> std::sync::Arc<dyn oxicode_sdk::ports::catalog::ModelCatalog> {
515 std::sync::Arc::clone(self.oxicode.catalog())
516 }
517
518 pub fn agent(&self) -> Arc<Agent> {
520 Arc::clone(&self.agent)
521 }
522
523 pub fn agent_tools(&self) -> Arc<oxicode_agent::ToolRegistry> {
525 self.agent.tools()
526 }
527
528 pub fn ask_bridge(&self) -> Option<&std::sync::Arc<oxicode_agent::tools::ask::AskBridge>> {
530 self.ask_bridge.as_ref()
531 }
532
533 pub fn skills(&self) -> parking_lot::RwLockReadGuard<'_, SkillManager> {
535 self.skills.read()
536 }
537
538 pub fn activate_skill(&self, name: &str) -> Result<(), String> {
540 {
541 let skills = self.skills.read();
542 if skills.get(name).is_none() {
543 return Err(format!("Skill '{}' not found", name));
544 }
545 }
546 let name_lower = name.to_lowercase();
547 {
548 let mut active = self.active_skills.write();
549 if !active.contains(&name_lower) {
550 active.push(name_lower);
551 }
552 }
553 self.rebuild_system_prompt();
554 Ok(())
555 }
556
557 pub fn deactivate_skill(&self, name: &str) {
559 let name_lower = name.to_lowercase();
560 {
561 let mut active = self.active_skills.write();
562 active.retain(|n| n != &name_lower);
563 }
564 self.rebuild_system_prompt();
565 }
566
567 pub fn active_skills(&self) -> Vec<String> {
569 self.active_skills.read().clone()
570 }
571
572 fn rebuild_system_prompt(&self) {
574 let active = self.active_skills.read();
575 let skills = self.skills.read();
576 let contents: Vec<String> = active
577 .iter()
578 .filter_map(|name| skills.get(name).map(|s| s.content.clone()))
579 .collect();
580 let persona = self.persona_body.read().clone();
584 let prompt =
585 build_system_prompt(self.settings.thinking_level, &contents, persona.as_deref());
586 self.agent.set_system_prompt(prompt);
587 }
588
589 pub fn agent_state(&self) -> oxicode_agent::AgentState {
591 self.agent.state()
592 }
593
594 pub async fn run_prompt(&self, prompt: String) -> Result<String> {
596 let (response, _events) = self.agent.run(prompt).await?;
597 Ok(response.content)
598 }
599
600 pub async fn run_prompt_with_events<F>(&self, prompt: String, on_event: F) -> Result<String>
602 where
603 F: FnMut(AgentEvent) + Send + 'static,
604 {
605 self.agent.run_streaming(prompt, on_event).await?;
606 let state = self.agent_state();
607 for msg in state.messages.iter().rev() {
608 if let oxicode_sdk::Message::Assistant(a) = msg {
609 return Ok(a.text_content());
610 }
611 }
612 Ok(String::new())
613 }
614
615 pub fn reset(&self) {
617 self.agent.reset();
618 }
619
620 pub async fn switch_model(&self, model_id: &str) -> anyhow::Result<()> {
626 let _ = self.agent.switch_model(model_id);
627 Ok(())
628 }
629
630 pub fn model_id(&self) -> String {
632 self.agent.model_id()
633 }
634
635 pub fn session_state(&self) -> &SessionState {
640 &self.session_state
641 }
642
643 pub fn should_stop_flag(&self) -> Arc<std::sync::atomic::AtomicBool> {
645 Arc::clone(&self.session_state.should_stop)
646 }
647
648 pub fn steering_queue(&self) -> Arc<RwLock<VecDeque<oxicode_sdk::Message>>> {
650 Arc::clone(&self.session_state.steering)
651 }
652
653 pub fn follow_up_queue(&self) -> Arc<RwLock<VecDeque<oxicode_sdk::Message>>> {
655 Arc::clone(&self.session_state.follow_up)
656 }
657}
658
659pub(crate) fn acquire_ownership_guard(
670 issue_store: Option<&oxicode_sdk::FileIssueStore>,
671 ownership_id: &str,
672) -> Option<oxicode_sdk::liveness::AliveGuard> {
673 let store = issue_store?;
674 if ownership_id.is_empty() {
675 return None;
678 }
679 match oxicode_sdk::liveness::acquire(&store.issues_dir(), ownership_id) {
680 Ok(guard) => Some(guard),
681 Err(e) => {
682 tracing::warn!(
687 ownership_id,
688 error = %e,
689 "issue liveness flock acquisition failed; this session's \
690 ownership claims are contestable while the lock is missing"
691 );
692 None
693 }
694 }
695}
696
697#[cfg(test)]
698mod tests {
699 use super::*;
704 use oxicode_sdk::FileIssueStore;
705 use oxicode_sdk::liveness;
706
707 fn tmp_store() -> (tempfile::TempDir, FileIssueStore) {
708 let tmp = tempfile::tempdir().unwrap();
709 let dir = tmp.path().join(".oxicode").join("issues");
710 std::fs::create_dir_all(&dir).unwrap();
711 (tmp, FileIssueStore::open(dir).unwrap())
712 }
713
714 #[test]
715 fn app_holds_single_liveness_lock() {
716 let (_tmp, store) = tmp_store();
720 let dir = store.issues_dir();
721 let id = "proc-test-app";
722
723 let guard = acquire_ownership_guard(Some(&store), id);
724 assert!(
725 guard.is_some(),
726 "App must acquire the liveness lock for its ownership id"
727 );
728 assert!(
729 liveness::is_session_alive(&dir, id),
730 "after acquire, the session must be live"
731 );
732
733 let second = liveness::acquire(&dir, id);
735 assert!(second.is_err(), "second acquire under same id must fail");
736
737 drop(guard);
738 assert!(
739 !liveness::is_session_alive(&dir, id),
740 "dropping App's guard releases the lock"
741 );
742 }
743
744 #[test]
745 fn acquire_returns_none_without_store() {
746 let dir = tempfile::tempdir().unwrap();
748 let id = "proc-x";
749 assert!(acquire_ownership_guard(None, id).is_none());
750 let _ = dir; }
752
753 #[test]
754 fn acquire_rejects_empty_ownership_id() {
755 let (_tmp, store) = tmp_store();
758 assert!(
759 acquire_ownership_guard(Some(&store), "").is_none(),
760 "empty ownership id must never acquire a lock (#13 guard)"
761 );
762 }
763}
764pub mod symbols;
765pub mod tui_vt;