Skip to main content

qualia_client_core/
resource_import.rs

1//! Catalog ontology download → `.q42` import pipeline (in-process, no subprocess).
2//!
3//! Mirrors `qualia-cli resources import-ontology` but writes under `{storage}/Index/`.
4//! After a successful compile the raw RDF/OWL source file is removed by default so only
5//! `{ontology_id}.q42` (+ `.q42.meta.json` sidecar) remain under `Index/`.
6
7use std::collections::HashMap;
8use std::io::{Read, Write};
9use std::path::{Path, PathBuf};
10use std::sync::atomic::{AtomicBool, Ordering};
11use std::sync::{Arc, Mutex};
12
13use qualia_core_db::{
14    ingest,
15    resource_catalog::{OntologyResource, ResourceCatalog},
16    wal::WriteAheadLog,
17};
18use serde::Serialize;
19use sha2::{Digest, Sha256};
20
21use crate::state::ProgressPayload;
22
23/// Optional download/import progress sink (shared with LLM Hub via `get_active_downloads`).
24pub struct ImportProgressCtx {
25    pub id: String,
26    pub handles: Arc<Mutex<HashMap<String, Arc<AtomicBool>>>>,
27    pub active_downloads: Arc<Mutex<HashMap<String, ProgressPayload>>>,
28    pub download_events: tokio::sync::broadcast::Sender<ProgressPayload>,
29}
30
31impl ImportProgressCtx {
32    pub fn emit(&self, payload: ProgressPayload) {
33        let _ = self.download_events.send(payload.clone());
34        self.active_downloads
35            .lock()
36            .unwrap()
37            .insert(self.id.clone(), payload);
38    }
39
40    pub fn clear(&self) {
41        self.handles.lock().unwrap().remove(&self.id);
42        self.active_downloads.lock().unwrap().remove(&self.id);
43    }
44}
45
46#[derive(Debug)]
47pub enum ImportError {
48    NotFound(String),
49    NoDownloadUrl(String),
50    Download(String),
51    Ingest(String),
52    Wal(String),
53    Io(std::io::Error),
54    Json(serde_json::Error),
55}
56
57impl std::fmt::Display for ImportError {
58    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
59        match self {
60            ImportError::NotFound(id) => write!(f, "Ontology not found in catalog: {id}"),
61            ImportError::NoDownloadUrl(id) => write!(f, "No download URL for ontology: {id}"),
62            ImportError::Download(e) => write!(f, "Download failed: {e}"),
63            ImportError::Ingest(e) => write!(f, "Ingest failed: {e}"),
64            ImportError::Wal(e) => write!(f, "WAL write failed: {e}"),
65            ImportError::Io(e) => write!(f, "IO error: {e}"),
66            ImportError::Json(e) => write!(f, "JSON error: {e}"),
67        }
68    }
69}
70
71impl From<std::io::Error> for ImportError {
72    fn from(e: std::io::Error) -> Self {
73        ImportError::Io(e)
74    }
75}
76
77impl From<serde_json::Error> for ImportError {
78    fn from(e: serde_json::Error) -> Self {
79        ImportError::Json(e)
80    }
81}
82
83#[derive(Debug, Clone, Serialize)]
84pub struct OntologyImportResult {
85    pub ontology_id: String,
86    pub source_path: String,
87    pub q42_path: String,
88    pub wal_path: String,
89    pub quin_count: u64,
90    pub catalog_quins: usize,
91    pub sha256: String,
92    pub imported_at: u64,
93    /// True when the downloaded RDF/OWL file was deleted after `.q42` compile.
94    pub source_removed: bool,
95}
96
97#[derive(Serialize)]
98struct OntologyMetaSidecar {
99    ontology_id: String,
100    quin_count: u64,
101    sha256: String,
102    imported_at: u64,
103    source_path: String,
104    q42_path: String,
105}
106
107pub fn index_dir(storage_root: &Path) -> PathBuf {
108    storage_root.join("Index")
109}
110
111fn ontology_source_filename(ont: &OntologyResource) -> String {
112    ont.download
113        .local_filename()
114        .unwrap_or_else(|| format!("{}.{}", ont.id, ont.format))
115}
116
117fn unix_now() -> u64 {
118    std::time::SystemTime::now()
119        .duration_since(std::time::UNIX_EPOCH)
120        .unwrap_or_default()
121        .as_secs()
122}
123
124fn sha256_file(path: &Path) -> Result<String, std::io::Error> {
125    let mut file = std::fs::File::open(path)?;
126    let mut hasher = Sha256::new();
127    let mut buf = [0u8; 65_536];
128    loop {
129        let n = file.read(&mut buf)?;
130        if n == 0 {
131            break;
132        }
133        hasher.update(&buf[..n]);
134    }
135    Ok(hex::encode(hasher.finalize()))
136}
137
138fn write_meta_sidecar(
139    ontology_id: &str,
140    q42_path: &Path,
141    source_path: &Path,
142    quin_count: u64,
143    sha256: &str,
144    imported_at: u64,
145) -> Result<(), ImportError> {
146    let meta_path = q42_path.with_extension("q42.meta.json");
147    let meta = OntologyMetaSidecar {
148        ontology_id: ontology_id.to_string(),
149        quin_count,
150        sha256: sha256.to_string(),
151        imported_at,
152        source_path: source_path.to_string_lossy().into_owned(),
153        q42_path: q42_path.to_string_lossy().into_owned(),
154    };
155    let json = serde_json::to_string_pretty(&meta)?;
156    std::fs::write(&meta_path, json)?;
157    Ok(())
158}
159
160pub async fn stream_download(url: &str, dest: &Path) -> Result<(), String> {
161    stream_download_with_progress(url, dest, None)
162        .await
163        .map_err(|e| e.to_string())
164}
165
166pub async fn stream_download_with_progress(
167    url: &str,
168    dest: &Path,
169    progress: Option<&ImportProgressCtx>,
170) -> Result<(), ImportError> {
171    let client = reqwest::Client::new();
172    let response = client
173        .get(url)
174        .header(
175            "User-Agent",
176            concat!("qualiaDB-client/", env!("CARGO_PKG_VERSION")),
177        )
178        .send()
179        .await
180        .map_err(|e| ImportError::Download(e.to_string()))?
181        .error_for_status()
182        .map_err(|e| ImportError::Download(e.to_string()))?;
183
184    let total_bytes = response.content_length().unwrap_or(0);
185    if let Some(ctx) = progress {
186        ctx.emit(ProgressPayload {
187            id: ctx.id.clone(),
188            progress: 0.0,
189            downloaded_bytes: 0,
190            total_bytes,
191            speed_kbps: 0.0,
192            status: "downloading".to_string(),
193        });
194    }
195
196    if let Some(parent) = dest.parent() {
197        std::fs::create_dir_all(parent)?;
198    }
199
200    let mut file = std::fs::File::create(dest).map_err(|e| {
201        ImportError::Io(std::io::Error::new(
202            e.kind(),
203            format!("Cannot create {}: {}", dest.display(), e),
204        ))
205    })?;
206
207    let mut stream = response.bytes_stream();
208    use futures_util::StreamExt;
209    let mut downloaded: u64 = 0;
210    let mut last_report = std::time::Instant::now();
211    let mut last_downloaded: u64 = 0;
212
213    while let Some(chunk) = stream.next().await {
214        if let Some(ctx) = progress {
215            if let Some(flag) = ctx.handles.lock().unwrap().get(&ctx.id) {
216                if flag.load(Ordering::Relaxed) {
217                    let _ = std::fs::remove_file(dest);
218                    ctx.emit(ProgressPayload {
219                        id: ctx.id.clone(),
220                        progress: 0.0,
221                        downloaded_bytes: downloaded,
222                        total_bytes,
223                        speed_kbps: 0.0,
224                        status: "cancelled".to_string(),
225                    });
226                    ctx.clear();
227                    return Err(ImportError::Download("Cancelled".to_string()));
228                }
229            }
230        }
231
232        let chunk = chunk.map_err(|e| ImportError::Download(e.to_string()))?;
233        file.write_all(&chunk)?;
234        downloaded += chunk.len() as u64;
235
236        if let Some(ctx) = progress {
237            let now = std::time::Instant::now();
238            if now.duration_since(last_report).as_millis() >= 200 {
239                let elapsed = now.duration_since(last_report).as_secs_f64().max(0.001);
240                let speed_kbps = ((downloaded - last_downloaded) as f64 / 1024.0) / elapsed;
241                let progress_pct = if total_bytes > 0 {
242                    (downloaded as f64 / total_bytes as f64) * 100.0
243                } else {
244                    0.0
245                };
246                ctx.emit(ProgressPayload {
247                    id: ctx.id.clone(),
248                    progress: progress_pct,
249                    downloaded_bytes: downloaded,
250                    total_bytes,
251                    speed_kbps,
252                    status: "downloading".to_string(),
253                });
254                last_report = now;
255                last_downloaded = downloaded;
256            }
257        }
258    }
259
260    Ok(())
261}
262
263/// Ingest a local RDF file into `{storage}/Index/{ontology_id}.q42`.
264pub fn ingest_local_rdf(
265    source_path: &Path,
266    ontology_id: &str,
267    storage_root: &Path,
268    ont: Option<&OntologyResource>,
269) -> Result<u64, ImportError> {
270    let index = index_dir(storage_root);
271    std::fs::create_dir_all(&index)?;
272
273    let q42_path = index.join(format!("{ontology_id}.q42"));
274    let in_str = source_path.to_string_lossy().into_owned();
275    let out_str = q42_path.to_string_lossy().into_owned();
276
277    let quin_count = ingest::streaming_import_rdf(&in_str, &out_str)
278        .map_err(|e| ImportError::Ingest(e.to_string()))?;
279
280    let imported_at = unix_now();
281    let sha256 = sha256_file(&q42_path)?;
282
283    if let Some(ont) = ont {
284        append_ontology_wal(storage_root, ont, &q42_path, imported_at)?;
285    }
286
287    write_meta_sidecar(
288        ontology_id,
289        &q42_path,
290        source_path,
291        quin_count,
292        &sha256,
293        imported_at,
294    )?;
295
296    Ok(quin_count)
297}
298
299fn append_ontology_wal(
300    storage_root: &Path,
301    ont: &OntologyResource,
302    q42_path: &Path,
303    timestamp: u64,
304) -> Result<usize, ImportError> {
305    let wal_path = index_dir(storage_root).join("ontologies.wal");
306    let mut wal = WriteAheadLog::open(&wal_path)
307        .map_err(|e| ImportError::Wal(format!("Cannot open {}: {}", wal_path.display(), e)))?;
308
309    let prov = ont.provenance_quin(timestamp, &q42_path.to_string_lossy());
310    wal.append_mutation(&prov)
311        .map_err(|e| ImportError::Wal(e.to_string()))?;
312
313    let catalog_quins = ont.to_quins();
314    for q in &catalog_quins {
315        wal.append_mutation(q)
316            .map_err(|e| ImportError::Wal(e.to_string()))?;
317    }
318
319    Ok(catalog_quins.len())
320}
321
322/// Download a catalog ontology and compile it to `.q42` under `{storage}/Index/`.
323pub async fn import_catalog_ontology(
324    catalog: &ResourceCatalog,
325    id: &str,
326    storage_root: &Path,
327) -> Result<OntologyImportResult, ImportError> {
328    import_catalog_ontology_with_options(catalog, id, storage_root, None, true).await
329}
330
331/// Download + compile with optional progress reporting and post-ingest source cleanup.
332pub async fn import_catalog_ontology_with_options(
333    catalog: &ResourceCatalog,
334    id: &str,
335    storage_root: &Path,
336    progress: Option<&ImportProgressCtx>,
337    delete_source_after_ingest: bool,
338) -> Result<OntologyImportResult, ImportError> {
339    let ont = catalog
340        .find_ontology(id)
341        .ok_or_else(|| ImportError::NotFound(id.to_string()))?;
342
343    if let Some(source_path) = crate::bundled_ontologies::resolve_bundled_ontology_source(id) {
344        if let Some(ctx) = progress {
345            ctx.emit(ProgressPayload {
346                id: ctx.id.clone(),
347                progress: 100.0,
348                downloaded_bytes: 0,
349                total_bytes: 0,
350                speed_kbps: 0.0,
351                status: "processing".to_string(),
352            });
353        }
354
355        let index = index_dir(storage_root);
356        std::fs::create_dir_all(&index)?;
357        let quin_count = ingest_local_rdf(&source_path, id, storage_root, Some(ont))?;
358        let q42_path = index.join(format!("{id}.q42"));
359        let wal_path = index.join("ontologies.wal");
360        let catalog_quins = ont.to_quins().len();
361        let sha256 = sha256_file(&q42_path)?;
362        let imported_at = unix_now();
363
364        if let Some(ctx) = progress {
365            ctx.emit(ProgressPayload {
366                id: ctx.id.clone(),
367                progress: 100.0,
368                downloaded_bytes: 0,
369                total_bytes: 0,
370                speed_kbps: 0.0,
371                status: "complete".to_string(),
372            });
373            ctx.clear();
374        }
375
376        return Ok(OntologyImportResult {
377            ontology_id: id.to_string(),
378            source_path: source_path.to_string_lossy().into_owned(),
379            q42_path: q42_path.to_string_lossy().into_owned(),
380            wal_path: wal_path.to_string_lossy().into_owned(),
381            quin_count,
382            catalog_quins,
383            sha256,
384            imported_at,
385            source_removed: false,
386        });
387    }
388
389    let url = ont
390        .download
391        .resolved_url()
392        .ok_or_else(|| ImportError::NoDownloadUrl(id.to_string()))?;
393
394    let index = index_dir(storage_root);
395    std::fs::create_dir_all(&index)?;
396
397    let filename = ontology_source_filename(ont);
398    let source_path = index.join(&filename);
399
400    stream_download_with_progress(&url, &source_path, progress).await?;
401
402    if let Some(ctx) = progress {
403        ctx.emit(ProgressPayload {
404            id: ctx.id.clone(),
405            progress: 100.0,
406            downloaded_bytes: std::fs::metadata(&source_path)
407                .map(|m| m.len())
408                .unwrap_or(0),
409            total_bytes: std::fs::metadata(&source_path)
410                .map(|m| m.len())
411                .unwrap_or(0),
412            speed_kbps: 0.0,
413            status: "processing".to_string(),
414        });
415    }
416
417    let quin_count = ingest_local_rdf(&source_path, id, storage_root, Some(ont))?;
418    let q42_path = index.join(format!("{id}.q42"));
419    let wal_path = index.join("ontologies.wal");
420    let catalog_quins = ont.to_quins().len();
421    let sha256 = sha256_file(&q42_path)?;
422    let imported_at = unix_now();
423    let source_path_str = source_path.to_string_lossy().into_owned();
424
425    let mut source_removed = false;
426    if delete_source_after_ingest && source_path.is_file() {
427        if std::fs::remove_file(&source_path).is_ok() {
428            source_removed = true;
429        }
430    }
431
432    if let Some(ctx) = progress {
433        ctx.emit(ProgressPayload {
434            id: ctx.id.clone(),
435            progress: 100.0,
436            downloaded_bytes: 0,
437            total_bytes: 0,
438            speed_kbps: 0.0,
439            status: "complete".to_string(),
440        });
441        ctx.clear();
442    }
443
444    Ok(OntologyImportResult {
445        ontology_id: id.to_string(),
446        source_path: source_path_str,
447        q42_path: q42_path.to_string_lossy().into_owned(),
448        wal_path: wal_path.to_string_lossy().into_owned(),
449        quin_count,
450        catalog_quins,
451        sha256,
452        imported_at,
453        source_removed,
454    })
455}