Skip to main content

qualia_client_core/
local_job_scheduler.rs

1//! Bounded local job queue for cold-path work (ontology import, daemon reload, …).
2//!
3//! Admin/UI layer may allocate; workers run one job at a time by default.
4
5use serde::{Deserialize, Serialize};
6use std::collections::{HashMap, HashSet};
7use std::fs;
8use std::path::{Path, PathBuf};
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::sync::{Arc, Mutex, OnceLock};
11use std::time::{SystemTime, UNIX_EPOCH};
12use tokio::sync::Notify;
13use uuid::Uuid;
14
15pub const MAX_JOBS: usize = 64;
16pub const MAX_COMPLETED_HISTORY: usize = 48;
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
19#[serde(rename_all = "snake_case")]
20pub enum JobStatus {
21    Queued,
22    Running,
23    Completed,
24    Failed,
25    Cancelled,
26}
27
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
29#[serde(rename_all = "snake_case", tag = "kind")]
30pub enum LocalJobKind {
31    OntologyCatalogImport {
32        ontology_id: String,
33    },
34    OntologyUriImport {
35        uri: String,
36        #[serde(default)]
37        ontology_id: Option<String>,
38        #[serde(default)]
39        domain: Option<String>,
40        #[serde(default)]
41        title: Option<String>,
42    },
43    BundledOntologySeed {
44        #[serde(default)]
45        ontology_id: Option<String>,
46    },
47    WorkbenchDaemonSync,
48    DaemonGraphReload,
49    /// Download a local inference model into the configured model store.
50    ModelDownload {
51        url: String,
52        filename: String,
53        model_id: String,
54    },
55    /// Validate, index, memory-map, and activate a GGUF/P64 model.
56    ModelActivation {
57        model_name: String,
58    },
59    /// Download public reference GLBs and compile them into the local anatomy `.10d` cache.
60    AnatomyAssetAcquire {
61        model: String,
62    },
63    /// Run one agent turn as a background job — the native agent processes it locally, or (when the
64    /// chosen agent's backend is remote-MCP) routes it out over MCP. Curated from chat by the person
65    /// (or, with confirmation, by their agent). Runs off the chat thread; the reply lands in the session.
66    AgentTurn {
67        session_id: String,
68        #[serde(default)]
69        agent_slug: Option<String>,
70        /// Roster revision captured when the job was curated.  Execution
71        /// refuses a changed/missing definition instead of silently adopting
72        /// new model, sharing, or remote-placement policy.
73        #[serde(default)]
74        agent_updated_at_unix: Option<u64>,
75        prompt: String,
76    },
77}
78
79impl LocalJobKind {
80    pub fn is_ontology_work(&self) -> bool {
81        matches!(
82            self,
83            LocalJobKind::OntologyCatalogImport { .. }
84                | LocalJobKind::OntologyUriImport { .. }
85                | LocalJobKind::BundledOntologySeed { .. }
86                | LocalJobKind::WorkbenchDaemonSync
87                | LocalJobKind::DaemonGraphReload
88        )
89    }
90}
91
92#[derive(Debug, Clone, Serialize, Deserialize)]
93pub struct LocalJob {
94    pub id: String,
95    pub kind: LocalJobKind,
96    pub status: JobStatus,
97    pub created_at: u64,
98    pub started_at: Option<u64>,
99    pub finished_at: Option<u64>,
100    pub progress: f64,
101    pub message: String,
102    #[serde(default)]
103    pub result: Option<serde_json::Value>,
104    #[serde(default)]
105    pub error: Option<String>,
106    /// Apparatus that should run the work (`did:q42:device:…`). Empty/None → this install.
107    #[serde(default)]
108    pub target_device_id: Option<String>,
109    /// Apparatus that enqueued the work (local device when known).
110    #[serde(default)]
111    pub originating_device_id: Option<String>,
112    /// Person principal who owns the work (not the OS account).
113    #[serde(default)]
114    pub person_id: Option<String>,
115}
116
117#[derive(Debug, Clone, Serialize, Deserialize)]
118pub struct EnqueueJobRequest {
119    #[serde(flatten)]
120    pub kind: LocalJobKind,
121}
122
123#[derive(Debug, Clone, Serialize, Deserialize)]
124pub struct JobQueueSnapshot {
125    pub jobs: Vec<LocalJob>,
126    pub queued: usize,
127    pub running: usize,
128    pub completed: usize,
129    pub failed: usize,
130}
131
132struct SchedulerInner {
133    jobs: Mutex<Vec<LocalJob>>,
134    cancel: Mutex<HashMap<String, Arc<AtomicBool>>>,
135    notify: Notify,
136    worker_started: AtomicBool,
137    store_path: Mutex<PathBuf>,
138}
139
140#[derive(Clone)]
141pub struct LocalJobScheduler {
142    inner: Arc<SchedulerInner>,
143}
144
145static GLOBAL: OnceLock<Arc<LocalJobScheduler>> = OnceLock::new();
146
147fn now_unix() -> u64 {
148    SystemTime::now()
149        .duration_since(UNIX_EPOCH)
150        .map(|d| d.as_secs())
151        .unwrap_or(0)
152}
153
154fn jobs_store_path() -> PathBuf {
155    crate::state::app_meta_dir().join("local-jobs.json")
156}
157
158impl LocalJobScheduler {
159    pub fn new() -> Self {
160        let path = jobs_store_path();
161        let mut jobs = load_jobs(&path);
162        // Recover interrupted runs after restart.
163        for job in jobs.iter_mut() {
164            if job.status == JobStatus::Running {
165                job.status = JobStatus::Queued;
166                job.message = "Re-queued after restart".to_string();
167                job.started_at = None;
168            }
169        }
170        Self {
171            inner: Arc::new(SchedulerInner {
172                jobs: Mutex::new(jobs),
173                cancel: Mutex::new(HashMap::new()),
174                notify: Notify::new(),
175                worker_started: AtomicBool::new(false),
176                store_path: Mutex::new(path),
177            }),
178        }
179    }
180
181    pub fn global() -> Arc<Self> {
182        GLOBAL.get_or_init(|| Arc::new(Self::new())).clone()
183    }
184
185    /// Spawn the background worker that processes queued jobs. The worker must run on a Tokio runtime;
186    /// pass a `Handle` when calling from outside a runtime context (e.g. from a Tauri `setup` callback,
187    /// which runs synchronously and is NOT inside `tokio::spawn`'s implicit runtime). When called from
188    /// within a runtime context, `None` falls back to `tokio::spawn`.
189    pub fn spawn_global_worker_with_runtime(runtime: tokio::runtime::Handle) {
190        Self::spawn_global_worker(Some(runtime));
191    }
192
193    /// Spawn the background worker. Uses `tokio::spawn` — only call this from within a Tokio runtime
194    /// context. For callers outside a runtime (e.g. Tauri's `setup` hook), use
195    /// [`spawn_global_worker_with_runtime`] with an explicit handle.
196    pub fn spawn_global_worker(runtime: Option<tokio::runtime::Handle>) {
197        let scheduler = Self::global();
198        if scheduler
199            .inner
200            .worker_started
201            .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
202            .is_err()
203        {
204            return;
205        }
206        let worker = scheduler.clone();
207        let task = async move {
208            worker.worker_loop().await;
209        };
210        match runtime {
211            Some(handle) => {
212                handle.spawn(task);
213            }
214            None => {
215                tokio::spawn(task);
216            }
217        }
218        // Kick the loop so jobs queued before first notify are processed.
219        scheduler.inner.notify.notify_one();
220    }
221
222    pub fn enqueue(&self, kind: LocalJobKind) -> Result<LocalJob, String> {
223        self.enqueue_for_device(kind, None)
224    }
225
226    /// Patch provenance fields on an existing job (e.g. after fleet accept).
227    pub fn update_job_meta(&self, updated: &LocalJob) -> Result<(), String> {
228        let mut jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
229        if let Some(slot) = jobs.iter_mut().find(|j| j.id == updated.id) {
230            slot.originating_device_id = updated.originating_device_id.clone();
231            slot.person_id = updated.person_id.clone();
232            slot.message = updated.message.clone();
233            slot.target_device_id = updated.target_device_id.clone();
234            self.persist(&jobs)?;
235        }
236        Ok(())
237    }
238
239    /// Enqueue work for a specific apparatus. `None` / local device_id runs here.
240    /// Remote registered devices are delivered via the fleet job path (HTTP + outbox).
241    pub fn enqueue_for_device(
242        &self,
243        kind: LocalJobKind,
244        target_device_id: Option<String>,
245    ) -> Result<LocalJob, String> {
246        let plane = crate::identity_plane::ensure_local_apparatus(None).ok();
247        let local_id = plane.as_ref().map(|p| p.local_device_id.clone());
248        let person_id = plane.as_ref().map(|p| p.person.person_id.clone());
249        let placement = crate::identity_plane::resolve_job_placement(target_device_id.as_deref())?;
250        match &placement {
251            crate::identity_plane::JobPlacement::Local { .. } => {}
252            crate::identity_plane::JobPlacement::RemoteRegistered { device_id, .. } => {
253                let entry =
254                    crate::identity_plane::deliver_or_queue_remote_job(kind.clone(), device_id)?;
255                // Represent remote work as a completed/queued audit job on the origin.
256                let job = LocalJob {
257                    id: entry.id.clone(),
258                    kind,
259                    status: if entry.delivered {
260                        JobStatus::Completed
261                    } else {
262                        JobStatus::Queued
263                    },
264                    created_at: entry.created_at_unix,
265                    started_at: Some(entry.last_attempt_unix),
266                    finished_at: if entry.delivered {
267                        Some(entry.last_attempt_unix)
268                    } else {
269                        None
270                    },
271                    progress: if entry.delivered { 1.0 } else { 0.0 },
272                    message: if entry.delivered {
273                        format!("Delivered to remote apparatus {device_id}")
274                    } else {
275                        format!(
276                            "Queued for remote delivery to {device_id}: {}",
277                            entry.last_error.clone().unwrap_or_default()
278                        )
279                    },
280                    result: serde_json::to_value(&entry).ok(),
281                    error: entry.last_error.clone(),
282                    target_device_id: Some(device_id.clone()),
283                    originating_device_id: local_id,
284                    person_id,
285                };
286                let mut jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
287                jobs.push(job.clone());
288                self.persist(&jobs)?;
289                return Ok(job);
290            }
291            crate::identity_plane::JobPlacement::Unknown { device_id } => {
292                return Err(format!(
293                    "Unknown target device_id {device_id}. Register it in the device fleet or leave target empty for this install."
294                ));
295            }
296        }
297
298        let job = LocalJob {
299            id: Uuid::new_v4().to_string(),
300            kind,
301            status: JobStatus::Queued,
302            created_at: now_unix(),
303            started_at: None,
304            finished_at: None,
305            progress: 0.0,
306            message: "Queued".to_string(),
307            result: None,
308            error: None,
309            target_device_id: local_id.clone(),
310            originating_device_id: local_id,
311            person_id,
312        };
313        {
314            let mut jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
315            if jobs.iter().any(|existing| {
316                matches!(existing.status, JobStatus::Queued | JobStatus::Running)
317                    && existing.kind == job.kind
318            }) {
319                return Err("A matching job is already queued or running".to_string());
320            }
321            if jobs.len() >= MAX_JOBS {
322                self.prune_completed(&mut jobs);
323            }
324            if jobs.len() >= MAX_JOBS {
325                return Err(format!("Job queue full (max {MAX_JOBS})"));
326            }
327            jobs.push(job.clone());
328            self.inner
329                .cancel
330                .lock()
331                .map_err(|e| e.to_string())?
332                .insert(job.id.clone(), Arc::new(AtomicBool::new(false)));
333            self.persist(&jobs)?;
334        }
335        self.inner.notify.notify_one();
336        Ok(job)
337    }
338
339    pub fn snapshot(&self) -> Result<JobQueueSnapshot, String> {
340        let jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
341        let mut queued = 0usize;
342        let mut running = 0usize;
343        let mut completed = 0usize;
344        let mut failed = 0usize;
345        for job in jobs.iter() {
346            match job.status {
347                JobStatus::Queued => queued += 1,
348                JobStatus::Running => running += 1,
349                JobStatus::Completed => completed += 1,
350                JobStatus::Failed | JobStatus::Cancelled => failed += 1,
351            }
352        }
353        Ok(JobQueueSnapshot {
354            jobs: jobs.clone(),
355            queued,
356            running,
357            completed,
358            failed,
359        })
360    }
361
362    pub fn get(&self, id: &str) -> Result<Option<LocalJob>, String> {
363        let jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
364        Ok(jobs.iter().find(|j| j.id == id).cloned())
365    }
366
367    pub fn cancel(&self, id: &str) -> Result<bool, String> {
368        let mut jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
369        let Some(job) = jobs.iter_mut().find(|j| j.id == id) else {
370            return Ok(false);
371        };
372        let download_id = match &job.kind {
373            LocalJobKind::ModelDownload { model_id, .. } => Some(model_id.clone()),
374            _ => None,
375        };
376        if job.status == JobStatus::Queued {
377            job.status = JobStatus::Cancelled;
378            job.finished_at = Some(now_unix());
379            job.message = "Cancelled before start".to_string();
380            self.persist(&jobs)?;
381            return Ok(true);
382        }
383        if job.status == JobStatus::Running {
384            if matches!(job.kind, LocalJobKind::ModelActivation { .. }) {
385                return Err(
386                    "Model activation cannot be interrupted safely after memory mapping starts"
387                        .to_string(),
388                );
389            }
390            if let Some(flag) = self.inner.cancel.lock().map_err(|e| e.to_string())?.get(id) {
391                flag.store(true, Ordering::Relaxed);
392            }
393            job.message = "Cancel requested".to_string();
394            self.persist(&jobs)?;
395            drop(jobs);
396            if let Some(download_id) = download_id {
397                let _ = crate::api::cancel_download(download_id);
398            }
399            return Ok(true);
400        }
401        Ok(false)
402    }
403
404    /// Re-enqueue a completed, failed, or cancelled job with the same bounded inputs.
405    pub fn retry(&self, id: &str) -> Result<LocalJob, String> {
406        let kind = {
407            let jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
408            let job = jobs
409                .iter()
410                .find(|job| job.id == id)
411                .ok_or_else(|| "job not found".to_string())?;
412            if matches!(job.status, JobStatus::Queued | JobStatus::Running) {
413                return Err("job is still active".to_string());
414            }
415            job.kind.clone()
416        };
417        self.enqueue(kind)
418    }
419
420    /// Remove finished history while preserving queued and running work.
421    pub fn clear_finished(&self) -> Result<usize, String> {
422        let mut jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
423        let before = jobs.len();
424        let removed_ids: Vec<String> = jobs
425            .iter()
426            .filter(|job| {
427                matches!(
428                    job.status,
429                    JobStatus::Completed | JobStatus::Failed | JobStatus::Cancelled
430                )
431            })
432            .map(|job| job.id.clone())
433            .collect();
434        jobs.retain(|job| matches!(job.status, JobStatus::Queued | JobStatus::Running));
435        if let Ok(mut cancel) = self.inner.cancel.lock() {
436            for id in removed_ids {
437                cancel.remove(&id);
438            }
439        }
440        self.persist(&jobs)?;
441        Ok(before.saturating_sub(jobs.len()))
442    }
443
444    fn prune_completed(&self, jobs: &mut Vec<LocalJob>) {
445        let mut done: Vec<(u64, String)> = jobs
446            .iter()
447            .filter_map(|j| {
448                if matches!(
449                    j.status,
450                    JobStatus::Completed | JobStatus::Failed | JobStatus::Cancelled
451                ) {
452                    Some((j.finished_at.unwrap_or(j.created_at), j.id.clone()))
453                } else {
454                    None
455                }
456            })
457            .collect();
458        if done.len() <= MAX_COMPLETED_HISTORY {
459            return;
460        }
461        done.sort_by_key(|(t, _)| *t);
462        let drop_n = done.len() - MAX_COMPLETED_HISTORY;
463        let drop_ids: HashSet<String> = done.into_iter().take(drop_n).map(|(_, id)| id).collect();
464        jobs.retain(|j| !drop_ids.contains(&j.id));
465        if let Ok(mut cancel) = self.inner.cancel.lock() {
466            for id in drop_ids {
467                cancel.remove(&id);
468            }
469        }
470    }
471
472    fn persist(&self, jobs: &[LocalJob]) -> Result<(), String> {
473        let path = self
474            .inner
475            .store_path
476            .lock()
477            .map_err(|e| e.to_string())?
478            .clone();
479        if let Some(parent) = path.parent() {
480            fs::create_dir_all(parent).map_err(|e| e.to_string())?;
481        }
482        let json = serde_json::to_string_pretty(jobs).map_err(|e| e.to_string())?;
483        fs::write(path, json).map_err(|e| e.to_string())
484    }
485
486    fn update_job<F>(&self, id: &str, f: F) -> Result<(), String>
487    where
488        F: FnOnce(&mut LocalJob),
489    {
490        let mut jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
491        let Some(job) = jobs.iter_mut().find(|j| j.id == id) else {
492            return Ok(());
493        };
494        f(job);
495        self.persist(&jobs)
496    }
497
498    /// Update visible progress for a running job. Progress is clamped and persisted so recovery and
499    /// agent diagnostics see the same state as the UI.
500    pub fn report_progress(
501        &self,
502        id: &str,
503        progress: f64,
504        message: impl Into<String>,
505    ) -> Result<(), String> {
506        let progress = progress.clamp(0.0, 0.99);
507        let message = message.into();
508        self.update_job(id, |job| {
509            if job.status == JobStatus::Running {
510                job.progress = progress;
511                job.message = message;
512            }
513        })
514    }
515
516    fn cancel_flag(&self, id: &str) -> Option<Arc<AtomicBool>> {
517        self.inner.cancel.lock().ok()?.get(id).cloned()
518    }
519
520    async fn worker_loop(&self) {
521        loop {
522            self.inner.notify.notified().await;
523            loop {
524                let next_id = {
525                    let jobs = match self.inner.jobs.lock() {
526                        Ok(j) => j,
527                        Err(_) => break,
528                    };
529                    jobs.iter()
530                        .find(|j| j.status == JobStatus::Queued)
531                        .map(|j| j.id.clone())
532                };
533                let Some(job_id) = next_id else { break };
534                if let Err(err) = self.run_one(job_id).await {
535                    eprintln!("[local_job_scheduler] worker error: {err}");
536                }
537            }
538        }
539    }
540
541    async fn run_one(&self, job_id: String) -> Result<(), String> {
542        let kind = {
543            let jobs = self.inner.jobs.lock().map_err(|e| e.to_string())?;
544            jobs.iter()
545                .find(|j| j.id == job_id)
546                .map(|j| j.kind.clone())
547                .ok_or_else(|| "job not found".to_string())?
548        };
549
550        self.update_job(&job_id, |job| {
551            job.status = JobStatus::Running;
552            job.started_at = Some(now_unix());
553            job.progress = 0.05;
554            job.message = "Running".to_string();
555        })?;
556
557        if kind.is_ontology_work() {
558            crate::activity_signals::begin_ontology_job();
559        }
560
561        log::info!("JOB|{}|running|{:?}", job_id, kind);
562        let outcome = execute_job(self, &job_id, &kind, self.cancel_flag(&job_id)).await;
563
564        if kind.is_ontology_work() {
565            crate::activity_signals::end_ontology_job();
566        }
567
568        match outcome {
569            Ok(result) => {
570                self.update_job(&job_id, |job| {
571                    job.status = JobStatus::Completed;
572                    job.finished_at = Some(now_unix());
573                    job.progress = 1.0;
574                    job.message = "Completed".to_string();
575                    job.result = Some(result);
576                    job.error = None;
577                })?;
578                log::info!("JOB|{}|completed|Completed", job_id);
579            }
580            Err(err) => {
581                let cancelled = err == "cancelled";
582                if cancelled {
583                    log::warn!("JOB|{}|cancelled|Cancelled", job_id);
584                } else {
585                    log::error!("JOB|{}|failed|{}", job_id, err);
586                }
587                self.update_job(&job_id, |job| {
588                    job.status = if cancelled {
589                        JobStatus::Cancelled
590                    } else {
591                        JobStatus::Failed
592                    };
593                    job.finished_at = Some(now_unix());
594                    job.progress = 1.0;
595                    job.message = if cancelled {
596                        "Cancelled".to_string()
597                    } else {
598                        "Failed".to_string()
599                    };
600                    job.error = Some(err);
601                })?;
602            }
603        }
604        Ok(())
605    }
606}
607
608fn load_jobs(path: &Path) -> Vec<LocalJob> {
609    let Ok(text) = fs::read_to_string(path) else {
610        return Vec::new();
611    };
612    serde_json::from_str(&text).unwrap_or_default()
613}
614
615fn check_cancel(flag: &Option<Arc<AtomicBool>>) -> Result<(), String> {
616    if flag.as_ref().is_some_and(|f| f.load(Ordering::Relaxed)) {
617        Err("cancelled".to_string())
618    } else {
619        Ok(())
620    }
621}
622
623async fn execute_job(
624    scheduler: &LocalJobScheduler,
625    job_id: &str,
626    kind: &LocalJobKind,
627    cancel: Option<Arc<AtomicBool>>,
628) -> Result<serde_json::Value, String> {
629    let state = crate::state::APP_STATE
630        .get()
631        .ok_or("APP_STATE not initialized")?;
632    let storage_path = state
633        .config
634        .lock()
635        .map_err(|e| e.to_string())?
636        .storage_path
637        .clone();
638
639    match kind {
640        LocalJobKind::OntologyCatalogImport { ontology_id } => {
641            check_cancel(&cancel)?;
642            let catalog = crate::api::load_workspace_catalog();
643            let progress = build_import_progress(job_id, state, cancel.clone());
644            let result = crate::resource_import::import_catalog_ontology_with_options(
645                &catalog,
646                ontology_id,
647                Path::new(&storage_path),
648                Some(&progress),
649                true,
650            )
651            .await
652            .map_err(|e| e.to_string())?;
653            check_cancel(&cancel)?;
654            qualia_core_db::daemon_graph::init_daemon_graph(&storage_path);
655            serde_json::to_value(result).map_err(|e| e.to_string())
656        }
657        LocalJobKind::OntologyUriImport {
658            uri,
659            ontology_id,
660            domain,
661            title,
662        } => {
663            check_cancel(&cancel)?;
664            let result = crate::ontology_workbench::import_from_uri(
665                Path::new(&storage_path),
666                uri.clone(),
667                ontology_id.clone(),
668                domain.clone(),
669                title.clone(),
670            )
671            .await?;
672            check_cancel(&cancel)?;
673            serde_json::to_value(result).map_err(|e| e.to_string())
674        }
675        LocalJobKind::BundledOntologySeed { ontology_id } => {
676            check_cancel(&cancel)?;
677            if let Some(id) = ontology_id {
678                let seeded = crate::bundled_ontologies::resolve_bundled_ontology_source(id)
679                    .ok_or_else(|| format!("Bundled source missing for {id}"))?;
680                let catalog = crate::api::load_workspace_catalog();
681                let ont = catalog.find_ontology(id);
682                crate::resource_import::ingest_local_rdf(
683                    &seeded,
684                    id,
685                    Path::new(&storage_path),
686                    ont,
687                )
688                .map_err(|e| e.to_string())?;
689                Ok(serde_json::json!({ "seeded": [id] }))
690            } else {
691                let seeded = crate::bundled_ontologies::seed_bundled_ontologies()?;
692                Ok(serde_json::json!({ "seeded": seeded }))
693            }
694        }
695        LocalJobKind::WorkbenchDaemonSync => {
696            check_cancel(&cancel)?;
697            crate::ontology_workbench::sync_workbench_seeds_to_daemon(Path::new(&storage_path))
698        }
699        LocalJobKind::DaemonGraphReload => {
700            check_cancel(&cancel)?;
701            qualia_core_db::daemon_graph::init_daemon_graph(&storage_path);
702            #[cfg(not(target_arch = "wasm32"))]
703            qualia_core_db::ontology_loader::load_startup_ontologies();
704            Ok(serde_json::json!({
705                "storage_path": storage_path,
706                "daemon_graph": "reloaded"
707            }))
708        }
709        LocalJobKind::ModelDownload {
710            url,
711            filename,
712            model_id,
713        } => {
714            check_cancel(&cancel)?;
715            scheduler.report_progress(
716                job_id,
717                0.02,
718                format!("Connecting to download source for {filename}"),
719            )?;
720            let mut download = Box::pin(crate::api::download_model(
721                url.clone(),
722                filename.clone(),
723                model_id.clone(),
724            ));
725            loop {
726                tokio::select! {
727                    result = &mut download => {
728                        let path = match result {
729                            Ok(path) => path,
730                            Err(_) if cancel
731                                .as_ref()
732                                .is_some_and(|flag| flag.load(Ordering::Relaxed)) =>
733                            {
734                                return Err("cancelled".to_string());
735                            }
736                            Err(error) => return Err(error),
737                        };
738                        check_cancel(&cancel)?;
739                        // Download alone left models as raw files with no install
740                        // manifest, so activation failed with "No install manifest"
741                        // and chat stayed on "no active model". Path-based
742                        // set_active_model finalizes (manifest + index) and activates.
743                        scheduler.report_progress(
744                            job_id,
745                            0.95,
746                            format!("Download complete — indexing and activating {filename}"),
747                        )?;
748                        let path_for_activate = path.clone();
749                        let model_id_for_activate = model_id.clone();
750                        let activate_result = tokio::task::spawn_blocking(move || {
751                            crate::api::set_active_model(path_for_activate.clone()).map_err(|e| {
752                                format!(
753                                    "download ok but activate failed for {model_id_for_activate}: {e}"
754                                )
755                            })?;
756                            Ok::<_, String>(
757                                crate::api::get_active_model().unwrap_or(path_for_activate),
758                            )
759                        })
760                        .await
761                        .map_err(|e| format!("post-download activation task failed: {e}"))?;
762                        let active_path = activate_result?;
763                        check_cancel(&cancel)?;
764                        return Ok(serde_json::json!({
765                            "model_id": model_id,
766                            "filename": filename,
767                            "path": path,
768                            "active": active_path,
769                            "lifecycle": "Active",
770                        }));
771                    }
772                    _ = tokio::time::sleep(std::time::Duration::from_millis(350)) => {
773                        if check_cancel(&cancel).is_err() {
774                            let _ = crate::api::cancel_download(model_id.clone());
775                        }
776                        if let Some(payload) = state
777                            .active_downloads
778                            .lock()
779                            .ok()
780                            .and_then(|downloads| downloads.get(model_id).cloned())
781                        {
782                            let fraction = (payload.progress / 100.0).clamp(0.02, 0.94);
783                            let size = if payload.total_bytes > 0 {
784                                format!(
785                                    "{:.1}/{:.1} MB",
786                                    payload.downloaded_bytes as f64 / 1_000_000.0,
787                                    payload.total_bytes as f64 / 1_000_000.0
788                                )
789                            } else {
790                                format!("{:.1} MB", payload.downloaded_bytes as f64 / 1_000_000.0)
791                            };
792                            scheduler.report_progress(
793                                job_id,
794                                fraction,
795                                format!("Downloading {filename} · {size} · {:.0} KB/s", payload.speed_kbps),
796                            )?;
797                        }
798                    }
799                }
800            }
801        }
802        LocalJobKind::ModelActivation { model_name } => {
803            check_cancel(&cancel)?;
804            scheduler.report_progress(
805                job_id,
806                0.08,
807                format!("Validating and indexing {model_name}"),
808            )?;
809            let selected = model_name.clone();
810            let scheduler_for_task = scheduler.clone();
811            let job_for_task = job_id.to_string();
812            let result = tokio::task::spawn_blocking(move || {
813                scheduler_for_task.report_progress(
814                    &job_for_task,
815                    0.32,
816                    "Mapping model into memory and initialising the inference backend",
817                )?;
818                crate::api::set_active_model(selected)?;
819                scheduler_for_task.report_progress(
820                    &job_for_task,
821                    0.94,
822                    "Model mapped; completing activation record",
823                )?;
824                Ok::<_, String>(())
825            })
826            .await
827            .map_err(|error| format!("model activation task failed: {error}"))?;
828            result?;
829            check_cancel(&cancel)?;
830            Ok(serde_json::json!({
831                "model": model_name,
832                "active": crate::api::get_active_model(),
833            }))
834        }
835        LocalJobKind::AnatomyAssetAcquire { model } => {
836            check_cancel(&cancel)?;
837            let parsed = crate::wellfair::api::parse_anatomy_model(model)?;
838            let scheduler_for_task = scheduler.clone();
839            let job_for_task = job_id.to_string();
840            let storage_for_task = storage_path.clone();
841            let cancel_for_task = cancel.clone();
842            let report = tokio::task::spawn_blocking(move || {
843                qualia_client_core_anatomy_acquire(
844                    &scheduler_for_task,
845                    &job_for_task,
846                    &storage_for_task,
847                    parsed,
848                    cancel_for_task,
849                )
850            })
851            .await
852            .map_err(|error| format!("anatomy acquisition task failed: {error}"))??;
853            check_cancel(&cancel)?;
854            serde_json::to_value(report).map_err(|error| error.to_string())
855        }
856        LocalJobKind::AgentTurn {
857            session_id,
858            agent_slug,
859            agent_updated_at_unix,
860            prompt,
861        } => {
862            check_cancel(&cancel)?;
863            if let (Some(slug), Some(expected_revision)) =
864                (agent_slug.as_deref(), agent_updated_at_unix)
865            {
866                let storage = crate::state::APP_STATE
867                    .get()
868                    .ok_or("Application not initialized")?
869                    .config
870                    .lock()
871                    .map_err(|error| error.to_string())?
872                    .storage_path
873                    .clone();
874                let current =
875                    crate::agent_registry::get_agent(std::path::Path::new(&storage), slug)
876                        .ok_or_else(|| format!("scheduled agent @{slug} no longer exists"))?;
877                if current.updated_at_unix != *expected_revision {
878                    return Err(format!(
879                        "scheduled agent @{slug} changed after this job was queued; review and schedule it again"
880                    ));
881                }
882            }
883            // Route by the chosen agent's backend (local-first). A remote-MCP agent runs its turn out
884            // over MCP (native only); otherwise the native engine runs it and appends the reply.
885            #[cfg(not(target_arch = "wasm32"))]
886            {
887                let backend = crate::api::agent_backend_kind(agent_slug.clone())
888                    .unwrap_or_else(|_| "local".to_string());
889                if backend == "remote" {
890                    let slug = agent_slug.clone().unwrap_or_default();
891                    return crate::api::run_remote_agent_turn(
892                        session_id.clone(),
893                        slug,
894                        prompt.clone(),
895                        false,
896                    );
897                }
898            }
899            let result = crate::chat_inference::run_chat_inference_for_agent(
900                session_id,
901                prompt,
902                agent_slug.as_deref(),
903                None,
904            );
905            if result.committed && !result.text.trim().is_empty() {
906                let _ = crate::api::append_chat_message(
907                    session_id.clone(),
908                    "agent".to_string(),
909                    result.text.clone(),
910                );
911            }
912            serde_json::to_value(&result).map_err(|e| e.to_string())
913        }
914    }
915}
916
917fn qualia_client_core_anatomy_acquire(
918    scheduler: &LocalJobScheduler,
919    job_id: &str,
920    storage_path: &str,
921    model: wellfare_core::anatomy::AnatomyModel,
922    cancel: Option<Arc<AtomicBool>>,
923) -> Result<crate::wellfair::anatomy_assets::AcquireReport, String> {
924    crate::wellfair::anatomy_assets::acquire_body_assets_controlled(
925        storage_path,
926        model,
927        |progress| {
928            let fraction = match progress.stage.as_str() {
929                "discover" => 0.05,
930                "fetch" => 0.10 + 0.58 * (progress.done as f64 / progress.total.max(1) as f64),
931                "compile" => 0.70 + 0.25 * (progress.done as f64 / progress.total.max(1) as f64),
932                "done" => 0.98,
933                _ => 0.05,
934            };
935            let _ = scheduler.report_progress(job_id, fraction, progress.message);
936        },
937        || {
938            cancel
939                .as_ref()
940                .is_some_and(|flag| flag.load(Ordering::Relaxed))
941        },
942    )
943}
944
945fn build_import_progress(
946    job_id: &str,
947    state: &crate::state::AppState,
948    cancel: Option<Arc<AtomicBool>>,
949) -> crate::resource_import::ImportProgressCtx {
950    let flag = cancel.unwrap_or_else(|| Arc::new(AtomicBool::new(false)));
951    if let Ok(mut handles) = state.download_handles.lock() {
952        handles.insert(job_id.to_string(), flag);
953    }
954    crate::resource_import::ImportProgressCtx {
955        id: job_id.to_string(),
956        handles: state.download_handles.clone(),
957        active_downloads: state.active_downloads.clone(),
958        download_events: state.download_events.clone(),
959    }
960}
961
962#[cfg(test)]
963mod tests {
964    use super::*;
965
966    #[test]
967    fn job_kind_deserializes_catalog_import() {
968        let raw = r#"{"kind":"ontology_catalog_import","ontology_id":"shacl"}"#;
969        let req: EnqueueJobRequest = serde_json::from_str(raw).unwrap();
970        assert!(matches!(
971            req.kind,
972            LocalJobKind::OntologyCatalogImport { .. }
973        ));
974    }
975
976    #[test]
977    fn job_kind_deserializes_model_activation() {
978        let raw = r#"{"kind":"model_activation","model_name":"C:\\Models\\local.gguf"}"#;
979        let req: EnqueueJobRequest = serde_json::from_str(raw).unwrap();
980        assert!(matches!(
981            req.kind,
982            LocalJobKind::ModelActivation { model_name }
983                if model_name.ends_with("local.gguf")
984        ));
985    }
986
987    #[test]
988    fn job_kind_deserializes_anatomy_acquisition() {
989        let raw = r#"{"kind":"anatomy_asset_acquire","model":"female"}"#;
990        let req: EnqueueJobRequest = serde_json::from_str(raw).unwrap();
991        assert!(matches!(
992            req.kind,
993            LocalJobKind::AnatomyAssetAcquire { model } if model == "female"
994        ));
995    }
996
997    #[test]
998    fn agent_turn_snapshot_is_backward_compatible() {
999        let raw = r#"{"kind":"agent_turn","session_id":"chat-1","agent_slug":"researcher","prompt":"summarise"}"#;
1000        let req: EnqueueJobRequest = serde_json::from_str(raw).unwrap();
1001        assert!(matches!(
1002            req.kind,
1003            LocalJobKind::AgentTurn {
1004                agent_slug: Some(ref slug),
1005                agent_updated_at_unix: None,
1006                ..
1007            } if slug == "researcher"
1008        ));
1009    }
1010}