1use 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 ModelDownload {
51 url: String,
52 filename: String,
53 model_id: String,
54 },
55 ModelActivation {
57 model_name: String,
58 },
59 AnatomyAssetAcquire {
61 model: String,
62 },
63 AgentTurn {
67 session_id: String,
68 #[serde(default)]
69 agent_slug: Option<String>,
70 #[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 #[serde(default)]
108 pub target_device_id: Option<String>,
109 #[serde(default)]
111 pub originating_device_id: Option<String>,
112 #[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 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 pub fn spawn_global_worker_with_runtime(runtime: tokio::runtime::Handle) {
190 Self::spawn_global_worker(Some(runtime));
191 }
192
193 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 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 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 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 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 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 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 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 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 #[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}