Skip to main content

qualia_client_core/api/
model.rs

1//! Active model + lifecycle
2
3#![allow(non_snake_case)]
4
5use super::*;
6
7use crate::state::*;
8use futures_util::StreamExt;
9use std::io::Write;
10use std::path::{Path, PathBuf};
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::{Arc, Mutex};
13
14pub fn active_model_path() -> PathBuf {
15    app_meta_dir().join("active_model.json")
16}
17
18fn legacy_active_model_path() -> PathBuf {
19    app_meta_dir().join("active_model.txt")
20}
21
22pub fn load_active_model_record_from_disk() -> Option<crate::model_lifecycle::ActiveModelRecord> {
23    let json_path = active_model_path();
24    if let Ok(text) = std::fs::read_to_string(&json_path) {
25        if let Ok(record) = serde_json::from_str(&text) {
26            return Some(record);
27        }
28    }
29
30    // Migrate legacy bare filename.
31    let legacy = std::fs::read_to_string(legacy_active_model_path())
32        .ok()
33        .map(|s| s.trim().to_string())
34        .filter(|s| !s.is_empty())?;
35
36    let state = crate::state::APP_STATE.get()?;
37    let storage = state.config.lock().unwrap().storage_path.clone();
38    let model_id = legacy
39        .trim_end_matches(".gguf")
40        .rsplit(['/', '\\'])
41        .next()
42        .unwrap_or(&legacy)
43        .to_string();
44
45    if let Some(manifest) =
46        crate::model_lifecycle::load_install_manifest(Path::new(&storage), &model_id)
47    {
48        let record = crate::model_lifecycle::ActiveModelRecord {
49            model_id: manifest.model_id,
50            gguf_path: manifest.gguf_path,
51            profile_id: manifest.profile_id,
52            quantization: manifest.quantization,
53            lifecycle_state: crate::model_lifecycle::lifecycle_label(
54                crate::model_lifecycle::get_model_lifecycle_state(),
55            )
56            .to_string(),
57            modality: manifest.modality,
58            architecture: manifest.architecture,
59            mmproj_path: manifest.mmproj_path,
60            context_window: manifest.context_window,
61        };
62        let _ = persist_active_model_record(&record);
63        let _ = std::fs::remove_file(legacy_active_model_path());
64        return Some(record);
65    }
66
67    None
68}
69
70pub fn load_active_model_from_disk() -> Option<String> {
71    load_active_model_record_from_disk().map(|r| r.gguf_path)
72}
73
74fn persist_active_model_record(
75    record: &crate::model_lifecycle::ActiveModelRecord,
76) -> Result<(), String> {
77    let meta = app_meta_dir();
78    std::fs::create_dir_all(&meta).map_err(|e| e.to_string())?;
79    let json = serde_json::to_string_pretty(record).map_err(|e| e.to_string())?;
80    std::fs::write(active_model_path(), json).map_err(|e| e.to_string())
81}
82
83pub fn clear_active_model_record() {
84    let _ = std::fs::remove_file(active_model_path());
85    let _ = std::fs::remove_file(legacy_active_model_path());
86}
87
88pub fn restore_active_model_on_startup() {
89    let state = crate::state::APP_STATE.get().unwrap();
90    let storage = state.config.lock().unwrap().storage_path.clone();
91    let storage_path = Path::new(&storage);
92
93    if let Some(record) = load_active_model_record_from_disk() {
94        if Path::new(&record.gguf_path).is_file() {
95            log::info!(
96                "LLM_LOAD|startup|0.00|Restoring active model {}",
97                record.model_id
98            );
99            match crate::model_lifecycle::activate_model_for_id(&record.model_id, storage_path) {
100                Ok(active) => {
101                    *state.active_model.lock().unwrap() = Some(active.gguf_path.clone());
102                    return;
103                }
104                Err(err) => {
105                    log::error!(
106                        "LLM_LOAD|failed|1.00|Startup restore failed for {}: {}",
107                        record.model_id,
108                        err
109                    );
110                    *state.active_model.lock().unwrap() = None;
111                    clear_active_model_record();
112                }
113            }
114        }
115    }
116
117    let catalog = load_workspace_catalog();
118    let prefs = crate::model_preferences::ensure_preferences(storage_path, &catalog);
119    if prefs.auto_select {
120        let _ = try_apply_model_preference("chat");
121    }
122}
123
124pub fn get_active_model() -> Option<String> {
125    let state = crate::state::APP_STATE.get().unwrap();
126    state.active_model.lock().unwrap().clone()
127}
128
129pub fn get_model_lifecycle_status() -> Result<serde_json::Value, String> {
130    let state = crate::state::APP_STATE.get().unwrap();
131    let path = state.active_model.lock().unwrap().clone();
132    let active = load_active_model_record_from_disk().or_else(|| {
133        path.as_ref()
134            .map(|gguf| crate::model_lifecycle::ActiveModelRecord {
135                model_id: gguf
136                    .rsplit(['/', '\\'])
137                    .next()
138                    .unwrap_or(gguf)
139                    .trim_end_matches(".gguf")
140                    .to_string(),
141                gguf_path: gguf.clone(),
142                profile_id: 0,
143                quantization: String::new(),
144                lifecycle_state: crate::model_lifecycle::lifecycle_label(
145                    crate::model_lifecycle::get_model_lifecycle_state(),
146                )
147                .to_string(),
148                modality: "text".to_string(),
149                architecture: None,
150                mmproj_path: None,
151                context_window: 4096,
152            })
153    });
154    let status = crate::model_lifecycle::get_model_status(active);
155    serde_json::to_value(status).map_err(|e| e.to_string())
156}
157
158pub fn set_active_model(model_name: String) -> Result<(), String> {
159    let state = crate::state::APP_STATE.get().unwrap();
160    let storage = state.config.lock().unwrap().storage_path.clone();
161    let storage_path = Path::new(&storage);
162
163    let result = if Path::new(&model_name).is_file() {
164        crate::model_lifecycle::finalize_local_gguf(Path::new(&model_name), storage_path)
165            .map_err(|e| e.to_string())
166    } else {
167        let model_id = model_name
168            .trim_end_matches(".gguf")
169            .rsplit(['/', '\\'])
170            .next()
171            .unwrap_or(model_name.as_str())
172            .to_string();
173        crate::model_lifecycle::activate_model_for_id(&model_id, storage_path)
174            .map_err(|e| e.to_string())
175    };
176    let record = match result {
177        Ok(record) => record,
178        Err(err) => {
179            *state.active_model.lock().unwrap() = None;
180            clear_active_model_record();
181            return Err(err);
182        }
183    };
184
185    persist_active_model_record(&record)?;
186    *state.active_model.lock().unwrap() = Some(record.gguf_path.clone());
187    Ok(())
188}
189
190/// Evict the resident model from memory without deleting on-disk GGUF files.
191pub fn unload_active_model() -> Result<(), String> {
192    let state = crate::state::APP_STATE.get().unwrap();
193    if let Some(record) = load_active_model_record_from_disk() {
194        crate::model_lifecycle::unload_active_model(Some(record.profile_id));
195    } else {
196        crate::model_lifecycle::unload_active_model(None);
197    }
198    *state.active_model.lock().unwrap() = None;
199    clear_active_model_record();
200    Ok(())
201}
202
203static MODEL_ACTIVATION_IN_PROGRESS: AtomicBool = AtomicBool::new(false);
204static MODEL_ACTIVATION_ERROR: Mutex<Option<String>> = Mutex::new(None);
205
206/// Activate a model on a background thread so the Flutter FRB caller is not blocked.
207pub fn set_active_model_async(model_name: String) -> Result<(), String> {
208    spawn_model_activation(move || set_active_model(model_name))
209}
210
211pub fn try_apply_model_preference_async(task: &str) -> Result<(), String> {
212    let task = task.to_string();
213    spawn_model_activation(move || try_apply_model_preference(&task))
214}
215
216fn spawn_model_activation(
217    work: impl FnOnce() -> Result<(), String> + Send + 'static,
218) -> Result<(), String> {
219    if MODEL_ACTIVATION_IN_PROGRESS.load(Ordering::Acquire) {
220        return Err("Model activation already in progress".to_string());
221    }
222    MODEL_ACTIVATION_IN_PROGRESS.store(true, Ordering::Release);
223    if let Ok(mut slot) = MODEL_ACTIVATION_ERROR.lock() {
224        *slot = None;
225    }
226    crate::system_telemetry::start_activation_telemetry("Loading model");
227    std::thread::Builder::new()
228        .name("qualia-model-activate".into())
229        .spawn(move || {
230            let result = work();
231            if let Err(err) = result {
232                if let Ok(mut slot) = MODEL_ACTIVATION_ERROR.lock() {
233                    *slot = Some(err);
234                }
235            }
236            crate::system_telemetry::stop_activation_telemetry();
237            MODEL_ACTIVATION_IN_PROGRESS.store(false, Ordering::Release);
238        })
239        .map_err(|e| format!("Failed to spawn model activation thread: {e}"))?;
240    Ok(())
241}
242
243pub fn is_model_activation_in_progress() -> bool {
244    MODEL_ACTIVATION_IN_PROGRESS.load(Ordering::Acquire)
245}
246
247pub fn take_model_activation_error() -> Option<String> {
248    MODEL_ACTIVATION_ERROR.lock().ok()?.take()
249}
250
251pub fn get_inference_backend_settings() -> crate::inference_backend::InferenceBackendSettings {
252    crate::inference_backend::load_inference_backend_settings()
253}
254
255pub fn save_inference_backend_settings(
256    settings: crate::inference_backend::InferenceBackendSettings,
257) -> Result<(), String> {
258    crate::inference_backend::save_inference_backend_settings(&settings)
259}
260
261/// Probe the configured Ollama endpoint (tags + reachability).
262pub fn probe_ollama_status() -> crate::ollama_harness::OllamaStatus {
263    crate::ollama_harness::probe_configured_ollama()
264}
265
266pub async fn probe_ollama_status_async() -> crate::ollama_harness::OllamaStatus {
267    crate::ollama_harness::probe_configured_ollama_async().await
268}
269
270/// List models on the configured Ollama host (empty vec if unreachable).
271pub fn list_ollama_models() -> Vec<crate::ollama_harness::OllamaModelInfo> {
272    crate::ollama_harness::probe_configured_ollama().models
273}
274
275/// One-shot Ollama generate using persisted settings (for smoke / ETL hooks).
276pub fn ollama_generate(
277    system: String,
278    prompt: String,
279) -> Result<crate::ollama_harness::OllamaGenerateResult, String> {
280    crate::ollama_harness::OllamaHarness::from_loaded_settings().generate(&system, &prompt)
281}
282
283pub async fn ollama_generate_async(
284    system: String,
285    prompt: String,
286) -> Result<crate::ollama_harness::OllamaGenerateResult, String> {
287    crate::ollama_harness::OllamaHarness::from_loaded_settings()
288        .generate_async(&system, &prompt)
289        .await
290}
291
292pub fn get_model_preferences() -> crate::model_preferences::ModelPreferences {
293    let state = crate::state::APP_STATE.get().unwrap();
294    let storage = state.config.lock().unwrap().storage_path.clone();
295    let catalog = load_workspace_catalog();
296    crate::model_preferences::ensure_preferences(Path::new(&storage), &catalog)
297}
298
299pub fn save_model_preferences(
300    prefs: crate::model_preferences::ModelPreferences,
301) -> Result<(), String> {
302    let state = crate::state::APP_STATE.get().unwrap();
303    let storage = state.config.lock().unwrap().storage_path.clone();
304    crate::model_preferences::save_preferences(Path::new(&storage), &prefs)
305}
306
307pub fn list_installed_llm_ids() -> Vec<String> {
308    let state = crate::state::APP_STATE.get().unwrap();
309    let storage = state.config.lock().unwrap().storage_path.clone();
310    crate::model_preferences::list_installed_model_ids(Path::new(&storage))
311}
312
313pub fn resolve_model_preference(
314    task: &str,
315) -> Option<crate::model_preferences::ResolvedModelPreference> {
316    let state = crate::state::APP_STATE.get().unwrap();
317    let storage = state.config.lock().unwrap().storage_path.clone();
318    let catalog = load_workspace_catalog();
319    let prefs = get_model_preferences();
320    let task = crate::model_preferences::ModelTask::from_str_lossy(task);
321    crate::model_preferences::resolve_preference(Path::new(&storage), &catalog, &prefs, task)
322}
323
324pub fn try_apply_model_preference(task: &str) -> Result<(), String> {
325    let state = crate::state::APP_STATE.get().unwrap();
326    let storage = state.config.lock().unwrap().storage_path.clone();
327    let catalog = load_workspace_catalog();
328    let prefs = get_model_preferences();
329    let task = crate::model_preferences::ModelTask::from_str_lossy(task);
330    let record = match crate::model_preferences::apply_preference(
331        Path::new(&storage),
332        &catalog,
333        &prefs,
334        task,
335    ) {
336        Ok(record) => record,
337        Err(err) => {
338            *state.active_model.lock().unwrap() = None;
339            clear_active_model_record();
340            return Err(err);
341        }
342    };
343    persist_active_model_record(&record)?;
344    *state.active_model.lock().unwrap() = Some(record.gguf_path.clone());
345    Ok(())
346}
347
348pub async fn install_catalog_llm(id: String) -> Result<serde_json::Value, String> {
349    let state = crate::state::APP_STATE.get().unwrap();
350    let storage_path = state.config.lock().unwrap().storage_path.clone();
351    let handles = state.download_handles.clone();
352    let active_dl = state.active_downloads.clone();
353    let catalog = load_workspace_catalog();
354
355    let model = catalog
356        .find_llm(&id)
357        .ok_or_else(|| format!("LLM not found in catalog: {id}"))?;
358    let url = model
359        .download
360        .resolved_url()
361        .ok_or_else(|| format!("No download URL for: {id}"))?;
362    let filename = model
363        .download
364        .local_filename()
365        .unwrap_or_else(|| format!("{id}.gguf"));
366
367    let models_dir = PathBuf::from(&storage_path).join("Models");
368    std::fs::create_dir_all(&models_dir).map_err(|e| e.to_string())?;
369    let dest_path = models_dir.join(&filename);
370
371    let cancelled = Arc::new(AtomicBool::new(false));
372    handles
373        .lock()
374        .unwrap()
375        .insert(id.clone(), cancelled.clone());
376
377    let client = reqwest::Client::new();
378    let response = client.get(&url).send().await.map_err(|e| {
379        handles.lock().unwrap().remove(&id);
380        active_dl.lock().unwrap().remove(&id);
381        e.to_string()
382    })?;
383    let total_bytes = response.content_length().unwrap_or(0);
384    let starting_payload = ProgressPayload {
385        id: id.clone(),
386        progress: 0.0,
387        downloaded_bytes: 0,
388        total_bytes,
389        speed_kbps: 0.0,
390        status: "downloading".to_string(),
391    };
392    let _ = state.download_events.send(starting_payload.clone());
393    active_dl
394        .lock()
395        .unwrap()
396        .insert(id.clone(), starting_payload);
397
398    let mut dest = std::fs::File::create(&dest_path).map_err(|e| e.to_string())?;
399    let mut stream = response.bytes_stream();
400    let mut downloaded: u64 = 0;
401    let mut last_report = std::time::Instant::now();
402    let mut last_downloaded: u64 = 0;
403
404    while let Some(chunk) = stream.next().await {
405        if cancelled.load(Ordering::Relaxed) {
406            let _ = std::fs::remove_file(&dest_path);
407            let payload = ProgressPayload {
408                id: id.clone(),
409                progress: 0.0,
410                downloaded_bytes: downloaded,
411                total_bytes,
412                speed_kbps: 0.0,
413                status: "cancelled".to_string(),
414            };
415            let _ = state.download_events.send(payload.clone());
416            handles.lock().unwrap().remove(&id);
417            active_dl.lock().unwrap().remove(&id);
418            return Err("Cancelled".to_string());
419        }
420        let chunk = chunk.map_err(|e| e.to_string())?;
421        dest.write_all(&chunk).map_err(|e| e.to_string())?;
422        downloaded += chunk.len() as u64;
423
424        let now = std::time::Instant::now();
425        if now.duration_since(last_report).as_millis() >= 200 {
426            let elapsed = now.duration_since(last_report).as_secs_f64().max(0.001);
427            let speed_kbps = ((downloaded - last_downloaded) as f64 / 1024.0) / elapsed;
428            let progress = if total_bytes > 0 {
429                (downloaded as f64 / total_bytes as f64) * 100.0
430            } else {
431                0.0
432            };
433            let payload = ProgressPayload {
434                id: id.clone(),
435                progress,
436                downloaded_bytes: downloaded,
437                total_bytes,
438                speed_kbps,
439                status: "downloading".to_string(),
440            };
441            let _ = state.download_events.send(payload.clone());
442            active_dl.lock().unwrap().insert(id.clone(), payload);
443            last_report = now;
444            last_downloaded = downloaded;
445        }
446    }
447
448    let processing_payload = ProgressPayload {
449        id: id.clone(),
450        progress: 100.0,
451        downloaded_bytes: downloaded,
452        total_bytes,
453        speed_kbps: 0.0,
454        status: "processing".to_string(),
455    };
456    let _ = state.download_events.send(processing_payload.clone());
457    active_dl
458        .lock()
459        .unwrap()
460        .insert(id.clone(), processing_payload);
461
462    let mut mmproj_path: Option<PathBuf> = None;
463    if model.is_multimodal() {
464        let vp = model.vision_projector.as_ref().ok_or_else(|| {
465            "Multimodal catalog entry missing vision_projector download".to_string()
466        })?;
467        let vp_url = vp
468            .resolved_url()
469            .ok_or_else(|| "No download URL for vision projector".to_string())?;
470        let vp_name = vp
471            .local_filename()
472            .unwrap_or_else(|| format!("{id}-mmproj.gguf"));
473        let vp_dest = models_dir.join(&vp_name);
474        crate::resource_import::stream_download(&vp_url, &vp_dest)
475            .await
476            .map_err(|e| e.to_string())?;
477        mmproj_path = Some(vp_dest);
478    }
479
480    let result = crate::model_lifecycle::finalize_llm_install(
481        model,
482        &dest_path,
483        mmproj_path.as_deref(),
484        Path::new(&storage_path),
485    )
486    .map_err(|e| e.to_string())?;
487
488    // New installs should be immediately usable in chat (not left at MappedToDisk).
489    if let Ok(record) = crate::model_lifecycle::activate_model_for_id(&id, Path::new(&storage_path))
490    {
491        let _ = persist_active_model_record(&record);
492        *state.active_model.lock().unwrap() = Some(record.gguf_path.clone());
493    }
494
495    let done_payload = ProgressPayload {
496        id: id.clone(),
497        progress: 100.0,
498        downloaded_bytes: downloaded,
499        total_bytes,
500        speed_kbps: 0.0,
501        status: "complete".to_string(),
502    };
503    let _ = state.download_events.send(done_payload.clone());
504    handles.lock().unwrap().remove(&id);
505    active_dl.lock().unwrap().remove(&id);
506
507    serde_json::to_value(result).map_err(|e| e.to_string())
508}