1use 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
23pub 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 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
263pub 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
322pub 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
331pub 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}