1#![warn(missing_docs)]
2#![warn(clippy::unwrap_used)]
9#![cfg_attr(test, allow(clippy::unwrap_used, clippy::field_reassign_with_default))]
10#![allow(unknown_lints)]
11
12pub mod bootstrap;
19pub mod cli;
20pub mod internal_urls;
21pub mod lsp;
22pub mod main_dispatch;
23pub mod mcp_credentials;
24pub mod print_mode;
25pub mod services;
26pub mod setup_wizard;
27pub mod store;
28
29pub(crate) mod app;
31pub(crate) mod context;
32pub mod discovery;
33pub mod extensions; pub(crate) mod infra;
35pub(crate) mod media;
36pub(crate) mod prompt;
37pub mod rpc_mode;
38pub(crate) mod skills;
39pub mod storage; pub use storage::packages::PackageManager;
42pub use storage::packages::ResourceKind;
43pub mod tools;
44pub mod tui; pub(crate) mod ui;
46pub(crate) mod util;
47
48pub async fn build_oxi_engine() -> anyhow::Result<oxi_sdk::Oxi> {
65 let paths = services::OxiPaths::default_paths()?;
66 services::build_oxi(&paths).await
67}
68
69pub async fn run_port_check() -> anyhow::Result<()> {
76 let oxi = build_oxi_engine().await?;
77 let ports = oxi.ports();
78
79 let entries = ports.state.list("").await?;
81 println!("[state] entries: {}", entries.len());
82
83 let providers = ports.auth.list_providers().await?;
85 println!("[auth] providers with credentials: {:?}", providers);
86
87 let keys = ports.config.list()?;
89 println!("[config] keys: {}", keys.len());
90
91 let skills = ports.skills.list().await?;
93 println!("[skills] {} skill(s) discovered", skills.len());
94 for s in &skills {
95 println!(" - {}: {}", s.name, s.description);
96 }
97
98 let _ = ports
100 .event_bus
101 .publish(&"port-check".to_string(), serde_json::json!({"ok": true}))
102 .await;
103 println!("[event-bus] publish ok (noop bus if not registered)");
104
105 println!("\nport check: ok");
106 Ok(())
107}
108
109#[derive(Debug, Clone)]
111pub struct CompactionContext {
112 pub messages_count: usize,
114 pub tokens_before: usize,
116 pub target_tokens: usize,
118 pub strategy: String,
120}
121
122impl CompactionContext {
123 pub fn new(
125 messages_count: usize,
126 tokens_before: usize,
127 target_tokens: usize,
128 strategy: impl Into<String>,
129 ) -> Self {
130 Self {
131 messages_count,
132 tokens_before,
133 target_tokens,
134 strategy: strategy.into(),
135 }
136 }
137
138 pub fn compression_ratio(&self) -> f32 {
140 if self.tokens_before == 0 {
141 return 1.0;
142 }
143 self.target_tokens as f32 / self.tokens_before as f32
144 }
145}
146
147use crate::store::settings::Settings;
149use anyhow::{Error, Result};
150use oxi_agent::{Agent, AgentConfig, AgentEvent};
151use parking_lot::RwLock;
152use skills::SkillManager;
153use std::sync::Arc;
154
155pub struct App {
164 oxi: oxi_sdk::Oxi,
165 agent: Arc<Agent>,
166 settings: Settings,
167 skills: RwLock<SkillManager>,
168 active_skills: RwLock<Vec<String>>,
169 wasm_ext: Option<std::sync::Arc<crate::extensions::WasmExtensionManager>>,
170 ask_bridge: Option<std::sync::Arc<oxi_agent::tools::ask::AskBridge>>,
171 issue_store: Option<crate::store::issues::FileIssueStore>,
175 ownership_session_id: String,
180 #[allow(dead_code)]
185 liveness_guard: Option<crate::store::issues::liveness::AliveGuard>,
186 persona_body: RwLock<Option<String>>,
189}
190
191fn build_system_prompt(
194 thinking_level: crate::store::settings::ThinkingLevel,
195 skill_contents: &[String],
196 persona_body: Option<&str>,
197) -> String {
198 let skills: Vec<prompt::system_prompt::Skill> = skill_contents
199 .iter()
200 .enumerate()
201 .map(|(i, content)| prompt::system_prompt::Skill {
202 name: format!("skill-{}", i),
203 content: content.clone(),
204 })
205 .collect();
206
207 let options = prompt::system_prompt::BuildSystemPromptOptions {
208 custom_prompt: prompt::system_prompt::thinking_level_prompt(thinking_level),
209 skills,
210 cwd: std::env::current_dir()
211 .map(|p| p.to_string_lossy().to_string())
212 .unwrap_or_default(),
213 persona_prompt: persona_body.map(|s| s.to_string()),
214 ..Default::default()
215 };
216
217 prompt::system_prompt::build_system_prompt(&options)
218}
219
220impl App {
223 pub async fn from_oxi(
237 oxi: oxi_sdk::Oxi,
238 settings: Settings,
239 ownership_session_id: String,
240 ) -> Result<Self> {
241 let persona = match oxi.ports().personas.get("default").await {
246 Ok(Some(p)) if !p.system_prompt.trim().is_empty() => Some(p),
247 Ok(_) => None,
248 Err(e) => {
249 tracing::warn!(error = %e, "default persona lookup failed");
250 None
251 }
252 };
253
254 let model_id = persona
255 .as_ref()
256 .and_then(|p| p.preferred_model.clone())
257 .or_else(|| settings.effective_model(None))
258 .unwrap_or_default();
259 let skills_dir = SkillManager::skills_dir().unwrap_or_else(|_| {
263 dirs::home_dir()
264 .unwrap_or_default()
265 .join(".oxi")
266 .join("skills")
267 });
268 let skills = SkillManager::load_from_dir(&skills_dir).unwrap_or_else(|e| {
269 tracing::debug!("Skills not loaded: {}", e);
270 SkillManager::new()
271 });
272
273 let body_str = persona.as_ref().map(|p| p.system_prompt.clone());
274 let system_prompt = build_system_prompt(settings.thinking_level, &[], body_str.as_deref());
275 let compaction_strategy = if settings.auto_compaction {
276 oxi_sdk::CompactionStrategy::Threshold(0.8)
277 } else {
278 oxi_sdk::CompactionStrategy::Disabled
279 };
280
281 let config = AgentConfig {
282 name: "oxi".to_string(),
283 description: Some("oxi CLI agent".to_string()),
284 model_id: model_id.clone(),
285 system_prompt: Some(system_prompt),
286 timeout_seconds: settings.tool_timeout_seconds,
287 temperature: settings.effective_temperature(),
288 max_tokens: settings.effective_max_tokens(),
289 compaction_strategy,
290 compaction_instruction: None,
291 context_window: 128_000,
292 workspace_dir: Some(
293 std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
294 ),
295 output_mode: None,
296 provider_options: None,
297 session_id: Some(ownership_session_id.clone()),
298 ttsr_engine: None,
299 memory: None,
300 todo: None,
301 agent_pool: None,
302 url_resolver: Some(Arc::new(oxi_sdk::SdkUrlResolver::new(
303 oxi.ports().url_router.clone(),
304 ))),
305 lsp: if crate::lsp::manager::default_servers().is_empty() {
311 None
312 } else {
313 Some(Arc::new(crate::lsp::CliLspProvider::with_defaults(
314 std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")),
315 )))
316 },
317 ..Default::default()
318 };
319
320 let cwd = std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."));
322 let agent = oxi
323 .agent(config)
324 .workspace(cwd)
325 .build()
326 .map_err(|e| Error::msg(format!("agent build failed: {e}")))?;
327 let agent = Arc::new(agent);
328
329 let ask_timeout = if settings.ask_timeout_secs > 0 {
330 Some(std::time::Duration::from_secs(settings.ask_timeout_secs))
331 } else {
332 None
333 };
334 let bridge =
335 std::sync::Arc::new(oxi_agent::tools::ask::AskBridge::with_timeout(ask_timeout));
336 let ask_tool = oxi_agent::tools::ask::AskTool::new(bridge.clone());
337 agent.tools().register_arc(std::sync::Arc::new(ask_tool));
338 let issue_store = std::env::current_dir()
343 .ok()
344 .map(|cwd| crate::store::issues::FileIssueStore::open_from_cwd(&cwd))
345 .and_then(|r| {
346 r.map_err(|e| tracing::warn!("issue store unavailable: {e}"))
347 .ok()
348 });
349
350 if let Some(store) = issue_store.clone() {
352 let tool = std::sync::Arc::new(crate::tools::IssueTool::new(store));
353 agent.tools().register_arc(tool);
354 }
355
356 Ok(Self {
357 oxi,
358 agent,
359 settings,
360 skills: RwLock::new(skills),
361 active_skills: RwLock::new(Vec::new()),
362 wasm_ext: None,
363 ask_bridge: Some(bridge),
364 issue_store,
365 ownership_session_id,
366 liveness_guard: None, persona_body: RwLock::new(persona.as_ref().map(|p| p.system_prompt.clone())),
368 })
369 .map(|mut app| {
370 app.liveness_guard =
374 acquire_ownership_guard(app.issue_store.as_ref(), &app.ownership_session_id);
375 app
376 })
377 }
378
379 pub fn ownership_session_id(&self) -> &str {
382 &self.ownership_session_id
383 }
384
385 pub fn has_liveness_lock(&self) -> bool {
390 self.liveness_guard.is_some()
391 }
392
393 pub fn settings(&self) -> &Settings {
395 &self.settings
396 }
397
398 pub fn set_wasm_ext(
400 &mut self,
401 ext: Option<std::sync::Arc<crate::extensions::WasmExtensionManager>>,
402 ) {
403 self.wasm_ext = ext;
404 }
405
406 pub fn wasm_ext(&self) -> Option<&std::sync::Arc<crate::extensions::WasmExtensionManager>> {
408 self.wasm_ext.as_ref()
409 }
410
411 pub fn issue_store(&self) -> Option<crate::store::issues::FileIssueStore> {
413 self.issue_store.clone()
414 }
415
416 pub fn oxi(&self) -> &oxi_sdk::Oxi {
419 &self.oxi
420 }
421
422 pub fn agent(&self) -> Arc<Agent> {
424 Arc::clone(&self.agent)
425 }
426
427 pub fn agent_tools(&self) -> Arc<oxi_agent::ToolRegistry> {
429 self.agent.tools()
430 }
431
432 pub fn ask_bridge(&self) -> Option<&std::sync::Arc<oxi_agent::tools::ask::AskBridge>> {
434 self.ask_bridge.as_ref()
435 }
436
437 pub fn skills(&self) -> parking_lot::RwLockReadGuard<'_, SkillManager> {
439 self.skills.read()
440 }
441
442 pub fn activate_skill(&self, name: &str) -> Result<(), String> {
444 {
445 let skills = self.skills.read();
446 if skills.get(name).is_none() {
447 return Err(format!("Skill '{}' not found", name));
448 }
449 }
450 let name_lower = name.to_lowercase();
451 {
452 let mut active = self.active_skills.write();
453 if !active.contains(&name_lower) {
454 active.push(name_lower);
455 }
456 }
457 self.rebuild_system_prompt();
458 Ok(())
459 }
460
461 pub fn deactivate_skill(&self, name: &str) {
463 let name_lower = name.to_lowercase();
464 {
465 let mut active = self.active_skills.write();
466 active.retain(|n| n != &name_lower);
467 }
468 self.rebuild_system_prompt();
469 }
470
471 pub fn active_skills(&self) -> Vec<String> {
473 self.active_skills.read().clone()
474 }
475
476 fn rebuild_system_prompt(&self) {
478 let active = self.active_skills.read();
479 let skills = self.skills.read();
480 let contents: Vec<String> = active
481 .iter()
482 .filter_map(|name| skills.get(name).map(|s| s.content.clone()))
483 .collect();
484 let persona = self.persona_body.read().clone();
488 let prompt =
489 build_system_prompt(self.settings.thinking_level, &contents, persona.as_deref());
490 self.agent.set_system_prompt(prompt);
491 }
492
493 pub fn agent_state(&self) -> oxi_agent::AgentState {
495 self.agent.state()
496 }
497
498 pub async fn run_prompt(&self, prompt: String) -> Result<String> {
500 let (response, _events) = self.agent.run(prompt).await?;
501 Ok(response.content)
502 }
503
504 pub async fn run_prompt_with_events<F>(&self, prompt: String, on_event: F) -> Result<String>
506 where
507 F: FnMut(AgentEvent) + Send + 'static,
508 {
509 self.agent.run_streaming(prompt, on_event).await?;
510 let state = self.agent_state();
511 for msg in state.messages.iter().rev() {
512 if let oxi_sdk::Message::Assistant(a) = msg {
513 return Ok(a.text_content());
514 }
515 }
516 Ok(String::new())
517 }
518
519 pub fn reset(&self) {
521 self.agent.reset();
522 }
523
524 pub async fn switch_model(&self, model_id: &str) -> anyhow::Result<()> {
530 let _ = self.agent.switch_model(model_id);
531 Ok(())
532 }
533
534 pub fn model_id(&self) -> String {
536 self.agent.model_id()
537 }
538}
539
540pub(crate) fn acquire_ownership_guard(
551 issue_store: Option<&crate::store::issues::FileIssueStore>,
552 ownership_id: &str,
553) -> Option<crate::store::issues::liveness::AliveGuard> {
554 let store = issue_store?;
555 if ownership_id.is_empty() {
556 return None;
559 }
560 crate::store::issues::liveness::acquire(&store.issues_dir(), ownership_id).ok()
561}
562
563#[cfg(test)]
564mod tests {
565 use super::*;
570 use crate::store::issues::FileIssueStore;
571 use crate::store::issues::liveness;
572
573 fn tmp_store() -> (tempfile::TempDir, FileIssueStore) {
574 let tmp = tempfile::tempdir().unwrap();
575 let dir = tmp.path().join(".oxi").join("issues");
576 std::fs::create_dir_all(&dir).unwrap();
577 (tmp, FileIssueStore::open(dir).unwrap())
578 }
579
580 #[test]
581 fn app_holds_single_liveness_lock() {
582 let (_tmp, store) = tmp_store();
586 let dir = store.issues_dir();
587 let id = "proc-test-app";
588
589 let guard = acquire_ownership_guard(Some(&store), id);
590 assert!(
591 guard.is_some(),
592 "App must acquire the liveness lock for its ownership id"
593 );
594 assert!(
595 liveness::is_session_alive(&dir, id),
596 "after acquire, the session must be live"
597 );
598
599 let second = liveness::acquire(&dir, id);
601 assert!(second.is_err(), "second acquire under same id must fail");
602
603 drop(guard);
604 assert!(
605 !liveness::is_session_alive(&dir, id),
606 "dropping App's guard releases the lock"
607 );
608 }
609
610 #[test]
611 fn acquire_returns_none_without_store() {
612 let dir = tempfile::tempdir().unwrap();
614 let id = "proc-x";
615 assert!(acquire_ownership_guard(None, id).is_none());
616 let _ = dir; }
618
619 #[test]
620 fn acquire_rejects_empty_ownership_id() {
621 let (_tmp, store) = tmp_store();
624 assert!(
625 acquire_ownership_guard(Some(&store), "").is_none(),
626 "empty ownership id must never acquire a lock (#13 guard)"
627 );
628 }
629}