1use std::collections::{HashMap, HashSet};
7
8use chrono::{DateTime, Utc};
9use serde::{Deserialize, Serialize};
10use uuid::Uuid;
11
12use khive_storage::types::{EdgeFilter, LinkId, PageRequest};
13use khive_storage::{EdgeRelation, EntityFilter};
14
15use crate::error::{RuntimeError, RuntimeResult};
16use crate::runtime::{KhiveRuntime, NamespaceToken};
17
18#[derive(Clone, Debug, Serialize, Deserialize)]
25pub struct KgArchive {
26 pub format: String,
27 pub version: String,
28 pub namespace: String,
29 pub exported_at: DateTime<Utc>,
30 pub entities: Vec<ExportedEntity>,
31 pub edges: Vec<ExportedEdge>,
32}
33
34#[derive(Clone, Debug, Serialize, Deserialize)]
36pub struct ExportedEntity {
37 pub id: Uuid,
38 pub kind: String,
40 #[serde(skip_serializing_if = "Option::is_none")]
42 pub entity_type: Option<String>,
43 pub name: String,
44 #[serde(skip_serializing_if = "Option::is_none")]
45 pub description: Option<String>,
46 #[serde(skip_serializing_if = "Option::is_none")]
47 pub properties: Option<serde_json::Value>,
48 #[serde(default)]
49 pub tags: Vec<String>,
50 pub created_at: DateTime<Utc>,
51 pub updated_at: DateTime<Utc>,
52}
53
54#[derive(Clone, Debug, Serialize, Deserialize)]
56pub struct ExportedEdge {
57 #[serde(default = "Uuid::new_v4")]
62 pub edge_id: Uuid,
63 pub source: Uuid,
64 pub target: Uuid,
65 pub relation: EdgeRelation,
67 pub weight: f64,
68 #[serde(default, skip_serializing_if = "Option::is_none")]
70 pub properties: Option<serde_json::Value>,
71 #[serde(
74 default = "default_import_timestamp",
75 deserialize_with = "deserialize_import_timestamp"
76 )]
77 pub created_at: DateTime<Utc>,
78 #[serde(
81 default = "default_import_timestamp",
82 deserialize_with = "deserialize_import_timestamp"
83 )]
84 pub updated_at: DateTime<Utc>,
85}
86
87fn default_import_timestamp() -> DateTime<Utc> {
88 Utc::now()
89}
90
91fn deserialize_import_timestamp<'de, D>(deserializer: D) -> Result<DateTime<Utc>, D::Error>
92where
93 D: serde::Deserializer<'de>,
94{
95 let raw = String::deserialize(deserializer)?;
96 DateTime::parse_from_rfc3339(&raw)
97 .map(|timestamp| timestamp.with_timezone(&Utc))
98 .map_err(|error| serde::de::Error::custom(format!("timestamp must be RFC3339: {error}")))
99}
100
101#[derive(Clone, Debug, Serialize, Deserialize)]
103pub struct ImportSummary {
104 pub entities_imported: usize,
105 pub edges_imported: usize,
106 pub edges_skipped: usize,
112 #[serde(default)]
115 pub embedding_truncation: crate::retrieval::EmbeddingTruncationReport,
116}
117
118impl KhiveRuntime {
121 pub async fn export_kg(&self, token: &NamespaceToken) -> RuntimeResult<KgArchive> {
127 let ns = token.namespace().as_str().to_owned();
128
129 let entity_page = self
130 .entities(token)?
131 .query_entities(
132 &ns,
133 EntityFilter::default(),
134 PageRequest {
135 offset: 0,
136 limit: u32::MAX,
137 },
138 )
139 .await?;
140
141 let entities: Vec<ExportedEntity> = entity_page
142 .items
143 .into_iter()
144 .map(|e| {
145 let created_at =
146 DateTime::from_timestamp_micros(e.created_at).unwrap_or_else(Utc::now);
147 let updated_at =
148 DateTime::from_timestamp_micros(e.updated_at).unwrap_or_else(Utc::now);
149 ExportedEntity {
150 id: e.id,
151 kind: e.kind.to_string(),
152 entity_type: e.entity_type,
153 name: e.name,
154 description: e.description,
155 properties: e.properties,
156 tags: e.tags,
157 created_at,
158 updated_at,
159 }
160 })
161 .collect();
162
163 let source_ids: Vec<Uuid> = entities.iter().map(|e| e.id).collect();
164 let edges = if source_ids.is_empty() {
165 Vec::new()
166 } else {
167 let filter = EdgeFilter {
168 source_ids: source_ids.clone(),
169 ..Default::default()
170 };
171 let edge_page = self
172 .graph(token)?
173 .query_edges(
174 filter,
175 Vec::new(),
176 PageRequest {
177 offset: 0,
178 limit: u32::MAX,
179 },
180 )
181 .await?;
182
183 let id_set: HashSet<Uuid> = source_ids.into_iter().collect();
184 edge_page
185 .items
186 .into_iter()
187 .filter(|e| id_set.contains(&e.source_id))
188 .map(|e| ExportedEdge {
189 edge_id: e.id.into(),
190 source: e.source_id,
191 target: e.target_id,
192 relation: e.relation,
193 weight: e.weight,
194 properties: e.metadata,
195 created_at: e.created_at,
196 updated_at: e.updated_at,
197 })
198 .collect()
199 };
200
201 Ok(KgArchive {
202 format: "khive-kg".to_string(),
203 version: "0.1".to_string(),
204 namespace: ns,
205 exported_at: Utc::now(),
206 entities,
207 edges,
208 })
209 }
210
211 pub async fn export_kg_json(&self, token: &NamespaceToken) -> RuntimeResult<String> {
213 let archive = self.export_kg(token).await?;
214 serde_json::to_string(&archive).map_err(|e| RuntimeError::InvalidInput(e.to_string()))
215 }
216
217 pub async fn import_kg(
226 &self,
227 archive: &KgArchive,
228 token: &NamespaceToken,
229 ) -> RuntimeResult<ImportSummary> {
230 if archive.format != "khive-kg" {
231 return Err(RuntimeError::InvalidInput(format!(
232 "unsupported archive format {:?}; expected \"khive-kg\"",
233 archive.format
234 )));
235 }
236 if archive.version != "0.1" {
237 return Err(RuntimeError::InvalidInput(format!(
238 "unsupported archive version {:?}; supported: \"0.1\"",
239 archive.version
240 )));
241 }
242
243 let ns = token.namespace().as_str().to_owned();
244
245 for (index, entity) in archive.entities.iter().enumerate() {
249 self.validate_entity_kind(&entity.kind)?;
250 crate::secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())?;
255 if entity.name.trim().is_empty() {
256 return Err(RuntimeError::InvalidInput(format!(
257 "archive entity {index} ({}) name must be non-blank",
258 entity.id
259 )));
260 }
261 }
262 for (index, edge) in archive.edges.iter().enumerate() {
263 let record = format!("edge[{index}]");
264 crate::operations::validate_edge_weight(edge.weight)?;
265 crate::secret_gate::reject_reserved_secret_gate_property(edge.properties.as_ref())?;
270 if let Some(p) = edge.properties.as_ref() {
271 crate::secret_gate::check_json_at(p, &record, "properties")?;
272 }
273 }
274
275 let store = self.entities(token)?;
276 let mut entities_imported = 0usize;
277 let mut embedding_truncation = crate::retrieval::EmbeddingTruncationReport::default();
278 for page in archive
279 .entities
280 .chunks(crate::retrieval::EMBEDDING_BATCH_PAGE_SIZE)
281 {
282 let entities: Vec<khive_storage::entity::Entity> = page
283 .iter()
284 .map(|ee| khive_storage::entity::Entity {
285 id: ee.id,
286 namespace: ns.clone(),
287 kind: ee.kind.clone(),
288 entity_type: ee.entity_type.clone(),
289 name: ee.name.clone(),
290 description: ee.description.clone(),
291 properties: ee.properties.clone(),
292 tags: ee.tags.clone(),
293 created_at: ee.created_at.timestamp_micros(),
294 updated_at: ee.updated_at.timestamp_micros(),
295 deleted_at: None,
296 merged_into: None,
297 merge_event_id: None,
298 version: 1,
299 content_ref: None,
300 })
301 .collect();
302 let texts: Vec<String> = entities
303 .iter()
304 .map(crate::curation::entity_embedding_text)
305 .collect();
306 let mut embeddings: Vec<HashMap<String, crate::retrieval::DocumentEmbeddingOutcome>> =
307 (0..entities.len()).map(|_| HashMap::new()).collect();
308 for model_name in self.registered_embedding_model_names() {
309 match self
310 .embed_document_batch_with_model_outcomes_for_token(token, &model_name, &texts)
311 .await
312 {
313 Ok(outcomes) => {
314 for (slot, outcome) in embeddings.iter_mut().zip(outcomes) {
315 slot.insert(model_name.clone(), outcome);
316 }
317 }
318 Err(error) => {
319 tracing::warn!(
320 model = %model_name,
321 error = %error,
322 "import_kg: batch embed failed; retrying records individually"
323 );
324 }
325 }
326 }
327 for (entity, outcomes) in entities.into_iter().zip(embeddings) {
328 store.upsert_entity(entity.clone()).await?;
329 embedding_truncation.merge(
332 self.reindex_entity_with_precomputed(token, &entity, outcomes)
333 .await?,
334 );
335 entities_imported += 1;
336 }
337 }
338
339 let graph = self.graph(token)?;
343 let mut edges_imported = 0usize;
344 let mut edges_skipped = 0usize;
345 for ee in &archive.edges {
346 let source_ok = match self.get_entity(token, ee.source).await {
347 Ok(_) => true,
348 Err(RuntimeError::NotFound(_)) => false,
349 Err(e) => return Err(e),
350 };
351 if !source_ok {
352 tracing::warn!(
353 source = %ee.source,
354 target = %ee.target,
355 relation = ?ee.relation,
356 "import_kg: skipping edge — source entity not found in namespace {ns:?}"
357 );
358 edges_skipped += 1;
359 continue;
360 }
361 let target_ok = match self.get_entity(token, ee.target).await {
362 Ok(_) => true,
363 Err(RuntimeError::NotFound(_)) => false,
364 Err(e) => return Err(e),
365 };
366 if !target_ok {
367 tracing::warn!(
368 source = %ee.source,
369 target = %ee.target,
370 relation = ?ee.relation,
371 "import_kg: skipping edge — target entity not found in namespace {ns:?}"
372 );
373 edges_skipped += 1;
374 continue;
375 }
376 match self
381 .validate_edge_relation_endpoints(token, ee.source, ee.target, ee.relation)
382 .await
383 {
384 Ok(_) => {}
385 Err(e @ (RuntimeError::InvalidInput(_) | RuntimeError::NotFound(_))) => {
386 tracing::warn!(
387 source = %ee.source,
388 target = %ee.target,
389 relation = ?ee.relation,
390 error = %e,
391 "import_kg: skipping edge — endpoint contract violation in namespace {ns:?}"
392 );
393 edges_skipped += 1;
394 continue;
395 }
396 Err(e) => return Err(e),
397 }
398 let edge = khive_storage::types::Edge {
399 id: LinkId::from(ee.edge_id),
400 namespace: ns.clone(),
401 source_id: ee.source,
402 target_id: ee.target,
403 relation: ee.relation,
404 weight: ee.weight,
405 created_at: ee.created_at,
406 updated_at: ee.updated_at,
407 deleted_at: None,
408 metadata: ee.properties.clone(),
409 target_backend: None,
410 };
411 graph.upsert_edge(edge).await?;
412 edges_imported += 1;
413 }
414
415 Ok(ImportSummary {
416 entities_imported,
417 edges_imported,
418 edges_skipped,
419 embedding_truncation,
420 })
421 }
422
423 pub async fn import_kg_json(
425 &self,
426 json: &str,
427 token: &NamespaceToken,
428 ) -> RuntimeResult<ImportSummary> {
429 let archive: KgArchive =
430 serde_json::from_str(json).map_err(|e| RuntimeError::InvalidInput(e.to_string()))?;
431 self.import_kg(&archive, token).await
432 }
433}
434
435#[cfg(test)]
440mod tests {
441 use std::sync::atomic::{AtomicUsize, Ordering};
442 use std::sync::Arc;
443
444 use super::*;
445 use crate::runtime::{KhiveRuntime, NamespaceToken};
446 use crate::{EmbedderProvider, Namespace};
447 use async_trait::async_trait;
448 use khive_storage::EdgeRelation;
449 use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService, MAX_TEXT_BYTES};
450
451 const IMPORT_TEST_MODEL: &str = "all-minilm-l6-v2";
452 const IMPORT_BATCH_MODEL: &str = "import-batch-model";
453 const IMPORT_BATCH_MODEL_TWO: &str = "import-batch-model-two";
454
455 struct CountingImportService {
456 calls: Arc<AtomicUsize>,
457 reject_poison: bool,
458 }
459
460 #[async_trait]
461 impl EmbeddingService for CountingImportService {
462 async fn embed(
463 &self,
464 texts: &[String],
465 _model: EmbeddingModel,
466 ) -> std::result::Result<Vec<Vec<f32>>, EmbedError> {
467 self.calls.fetch_add(1, Ordering::SeqCst);
468 if self.reject_poison
469 && (texts.len() > 1 || texts.iter().any(|text| text.contains("poison")))
470 {
471 return Err(EmbedError::InferenceFailed("poison input".into()));
472 }
473 Ok(texts
474 .iter()
475 .map(|text| {
476 vec![
477 text.len() as f32,
478 text.bytes().map(u32::from).sum::<u32>() as f32,
479 text.as_bytes().first().copied().unwrap_or_default() as f32,
480 text.as_bytes().last().copied().unwrap_or_default() as f32,
481 ]
482 })
483 .collect())
484 }
485
486 fn supports_model(&self, _model: EmbeddingModel) -> bool {
487 true
488 }
489
490 fn name(&self) -> &'static str {
491 "import-batch-counting-service"
492 }
493 }
494
495 struct CountingImportProvider {
496 name: &'static str,
497 calls: Arc<AtomicUsize>,
498 reject_poison: bool,
499 }
500
501 #[async_trait]
502 impl EmbedderProvider for CountingImportProvider {
503 fn name(&self) -> &str {
504 self.name
505 }
506
507 fn dimensions(&self) -> usize {
508 4
509 }
510
511 async fn build(&self) -> std::result::Result<Arc<dyn EmbeddingService>, RuntimeError> {
512 Ok(Arc::new(CountingImportService {
513 calls: Arc::clone(&self.calls),
514 reject_poison: self.reject_poison,
515 }))
516 }
517 }
518
519 fn batch_archive(names: impl IntoIterator<Item = String>) -> KgArchive {
520 KgArchive {
521 format: "khive-kg".to_string(),
522 version: "0.1".to_string(),
523 namespace: "local".to_string(),
524 exported_at: Utc::now(),
525 entities: names
526 .into_iter()
527 .map(|name| ExportedEntity {
528 id: Uuid::new_v4(),
529 kind: "concept".to_string(),
530 entity_type: None,
531 name,
532 description: Some("batch body".to_string()),
533 properties: None,
534 tags: vec![],
535 created_at: Utc::now(),
536 updated_at: Utc::now(),
537 })
538 .collect(),
539 edges: vec![],
540 }
541 }
542
543 struct ImportEmbeddingService;
544
545 #[async_trait]
546 impl EmbeddingService for ImportEmbeddingService {
547 async fn embed(
548 &self,
549 texts: &[String],
550 _model: EmbeddingModel,
551 ) -> std::result::Result<Vec<Vec<f32>>, EmbedError> {
552 let dimensions = EmbeddingModel::AllMiniLmL6V2.dimensions();
553 Ok(texts.iter().map(|_| vec![1.0; dimensions]).collect())
554 }
555
556 fn supports_model(&self, _model: EmbeddingModel) -> bool {
557 true
558 }
559
560 fn name(&self) -> &'static str {
561 "kg-import-truncation-test"
562 }
563 }
564
565 struct ImportEmbeddingProvider;
566
567 #[async_trait]
568 impl EmbedderProvider for ImportEmbeddingProvider {
569 fn name(&self) -> &str {
570 IMPORT_TEST_MODEL
571 }
572
573 fn dimensions(&self) -> usize {
574 EmbeddingModel::AllMiniLmL6V2.dimensions()
575 }
576
577 async fn build(&self) -> std::result::Result<Arc<dyn EmbeddingService>, RuntimeError> {
578 Ok(Arc::new(ImportEmbeddingService))
579 }
580 }
581
582 async fn make_rt() -> KhiveRuntime {
583 KhiveRuntime::memory().expect("in-memory runtime")
584 }
585
586 fn make_rt_with_embedder() -> KhiveRuntime {
587 let runtime = KhiveRuntime::memory().expect("in-memory runtime");
588 runtime.register_embedder(ImportEmbeddingProvider);
589 runtime
590 }
591
592 #[tokio::test]
593 async fn import_batches_provider_calls_and_matches_per_record_reindex_vectors() {
594 let token = NamespaceToken::local();
595 let archive = batch_archive((0..257).map(|index| format!("Batch entity {index}")));
596 let ids: Vec<Uuid> = archive.entities.iter().map(|entity| entity.id).collect();
597 let calls = Arc::new(AtomicUsize::new(0));
598 let runtime = KhiveRuntime::memory().unwrap();
599 runtime.register_embedder(CountingImportProvider {
600 name: IMPORT_BATCH_MODEL,
601 calls: Arc::clone(&calls),
602 reject_poison: false,
603 });
604 let usage = crate::usage::UsageContext::new();
605 let summary = crate::usage::scope(usage.clone(), runtime.import_kg(&archive, &token))
606 .await
607 .unwrap();
608 assert_eq!(summary.entities_imported, ids.len());
609 assert_eq!(usage.snapshot()["embed_calls"], ids.len() as u64);
610 assert_eq!(
611 calls.load(Ordering::SeqCst),
612 2,
613 "257 records need two provider batches"
614 );
615
616 let singles = KhiveRuntime::memory().unwrap();
617 let singleton_calls = Arc::new(AtomicUsize::new(0));
618 singles.register_embedder(CountingImportProvider {
619 name: IMPORT_BATCH_MODEL,
620 calls: Arc::clone(&singleton_calls),
621 reject_poison: false,
622 });
623 for id in &ids {
624 let entity = runtime.get_entity(&token, *id).await.unwrap();
625 singles
626 .entities(&token)
627 .unwrap()
628 .upsert_entity(entity.clone())
629 .await
630 .unwrap();
631 singles.reindex_entity(&token, &entity).await.unwrap();
632 }
633 assert_eq!(singleton_calls.load(Ordering::SeqCst), ids.len());
634 let batched_vectors = runtime
635 .vectors_for_model(&token, IMPORT_BATCH_MODEL)
636 .unwrap()
637 .get_vectors(&ids, "local", "entity.body")
638 .await
639 .unwrap();
640 let singleton_vectors = singles
641 .vectors_for_model(&token, IMPORT_BATCH_MODEL)
642 .unwrap()
643 .get_vectors(&ids, "local", "entity.body")
644 .await
645 .unwrap();
646 assert_eq!(batched_vectors.len(), ids.len());
647 assert_eq!(batched_vectors, singleton_vectors);
648 }
649
650 #[tokio::test]
651 async fn import_batches_each_registered_model_once_per_page() {
652 let token = NamespaceToken::local();
653 let archive = batch_archive(["alpha", "beta", "gamma"].map(str::to_string));
654 let calls_a = Arc::new(AtomicUsize::new(0));
655 let calls_b = Arc::new(AtomicUsize::new(0));
656 let runtime = KhiveRuntime::memory().unwrap();
657 for (name, calls) in [
658 (IMPORT_BATCH_MODEL, Arc::clone(&calls_a)),
659 (IMPORT_BATCH_MODEL_TWO, Arc::clone(&calls_b)),
660 ] {
661 runtime.register_embedder(CountingImportProvider {
662 name,
663 calls,
664 reject_poison: false,
665 });
666 }
667 runtime.import_kg(&archive, &token).await.unwrap();
668 assert_eq!(calls_a.load(Ordering::SeqCst), 1);
669 assert_eq!(calls_b.load(Ordering::SeqCst), 1);
670 let ids: Vec<Uuid> = archive.entities.iter().map(|entity| entity.id).collect();
671 for model in [IMPORT_BATCH_MODEL, IMPORT_BATCH_MODEL_TWO] {
672 assert_eq!(
673 runtime
674 .vectors_for_model(&token, model)
675 .unwrap()
676 .get_vectors(&ids, "local", "entity.body")
677 .await
678 .unwrap()
679 .len(),
680 ids.len()
681 );
682 }
683 }
684
685 #[tokio::test]
686 async fn import_failed_batch_isolates_poison_without_duplicate_vector_writes() {
687 use khive_storage::types::{SqlStatement, SqlValue};
688
689 let token = NamespaceToken::local();
690 let archive = batch_archive(["good one", "poison", "good two"].map(str::to_string));
691 let ids: Vec<Uuid> = archive.entities.iter().map(|entity| entity.id).collect();
692 let calls = Arc::new(AtomicUsize::new(0));
693 let runtime = KhiveRuntime::memory().unwrap();
694 runtime.register_embedder(CountingImportProvider {
695 name: IMPORT_BATCH_MODEL,
696 calls: Arc::clone(&calls),
697 reject_poison: true,
698 });
699 let summary = runtime.import_kg(&archive, &token).await.unwrap();
700 assert_eq!(summary.entities_imported, 3);
701 assert_eq!(
702 calls.load(Ordering::SeqCst),
703 4,
704 "one failed page plus three singleton retries"
705 );
706 let vectors = runtime
707 .vectors_for_model(&token, IMPORT_BATCH_MODEL)
708 .unwrap()
709 .get_vectors(&ids, "local", "entity.body")
710 .await
711 .unwrap();
712 assert_eq!(vectors.len(), 2);
713 assert!(!vectors.contains_key(&ids[1]));
714 let mut reader = runtime.sql().reader().await.unwrap();
715 let rows = reader
716 .query_all(SqlStatement {
717 sql: "SELECT COUNT(*) FROM ann_write_log WHERE embedding_model = ?1 AND op = 'upsert'".into(),
718 params: vec![SqlValue::Text(IMPORT_BATCH_MODEL.into())],
719 label: Some("import-batch-upsert-count".into()),
720 })
721 .await
722 .unwrap();
723 assert!(matches!(&rows[0].columns[0].value, SqlValue::Integer(2)));
724 let provenance = reader
725 .query_scalar(SqlStatement {
726 sql: "SELECT COUNT(*) FROM vector_provenance WHERE model_key = ?1".into(),
727 params: vec![SqlValue::Text(crate::config::sanitize_key(
728 IMPORT_BATCH_MODEL,
729 ))],
730 label: Some("import-batch-provenance-count".into()),
731 })
732 .await
733 .unwrap();
734 assert!(matches!(provenance, Some(SqlValue::Integer(0))));
735 }
736
737 #[tokio::test]
738 async fn import_summary_reports_actual_embedding_truncation() {
739 let runtime = make_rt_with_embedder();
740 let token = NamespaceToken::local();
741 let archive = KgArchive {
742 format: "khive-kg".to_string(),
743 version: "0.1".to_string(),
744 namespace: "local".to_string(),
745 exported_at: Utc::now(),
746 entities: vec![ExportedEntity {
747 id: Uuid::new_v4(),
748 kind: "concept".to_string(),
749 entity_type: None,
750 name: "Imported long entity".to_string(),
751 description: Some("x".repeat(MAX_TEXT_BYTES + 1)),
752 properties: None,
753 tags: vec![],
754 created_at: Utc::now(),
755 updated_at: Utc::now(),
756 }],
757 edges: vec![],
758 };
759
760 let summary = runtime
761 .import_kg(&archive, &token)
762 .await
763 .expect("import with bounded embedding input");
764
765 assert_eq!(summary.entities_imported, 1);
766 assert_eq!(summary.embedding_truncation.truncated, 1);
767 assert!(summary.embedding_truncation.discarded_bytes > 0);
768 let wire = serde_json::to_value(&summary).expect("serialize import summary");
769 assert_eq!(wire["embedding_truncation"]["truncated"], 1);
770 }
771
772 #[tokio::test]
773 async fn import_rejects_whitespace_name_before_any_entity_write() {
774 let runtime = make_rt().await;
775 let token = NamespaceToken::local();
776 let valid_id = Uuid::new_v4();
777 let archive = KgArchive {
778 format: "khive-kg".to_string(),
779 version: "0.1".to_string(),
780 namespace: "local".to_string(),
781 exported_at: Utc::now(),
782 entities: vec![
783 ExportedEntity {
784 id: valid_id,
785 kind: "concept".to_string(),
786 entity_type: None,
787 name: "Must not be written".to_string(),
788 description: None,
789 properties: None,
790 tags: vec![],
791 created_at: Utc::now(),
792 updated_at: Utc::now(),
793 },
794 ExportedEntity {
795 id: Uuid::new_v4(),
796 kind: "concept".to_string(),
797 entity_type: None,
798 name: " \t\n ".to_string(),
799 description: None,
800 properties: None,
801 tags: vec![],
802 created_at: Utc::now(),
803 updated_at: Utc::now(),
804 },
805 ],
806 edges: vec![],
807 };
808
809 let err = runtime
810 .import_kg(&archive, &token)
811 .await
812 .expect_err("whitespace-only entity names must fail the whole import");
813 assert!(
814 err.to_string().contains("non-blank"),
815 "error must explain the name invariant: {err}"
816 );
817 assert!(
818 runtime.get_entity(&token, valid_id).await.is_err(),
819 "deterministic validation must finish before the first entity write"
820 );
821 }
822
823 #[tokio::test]
825 async fn roundtrip_entities_and_edges() {
826 let src = make_rt().await;
827 let tok = NamespaceToken::local();
828 let e1 = src
829 .create_entity(
830 &tok,
831 "concept",
832 None,
833 "FlashAttention",
834 Some("fast attention"),
835 None,
836 vec![],
837 )
838 .await
839 .unwrap();
840 let e2 = src
841 .create_entity(
842 &tok,
843 "concept",
844 None,
845 "FlashAttention-2",
846 None,
847 None,
848 vec![],
849 )
850 .await
851 .unwrap();
852 let e3 = src
853 .create_entity(
854 &tok,
855 "person",
856 None,
857 "Tri Dao",
858 None,
859 None,
860 vec!["author".into()],
861 )
862 .await
863 .unwrap();
864 src.link(&tok, e2.id, e1.id, EdgeRelation::Extends, 1.0, None)
865 .await
866 .unwrap();
867 src.link(&tok, e1.id, e3.id, EdgeRelation::IntroducedBy, 0.9, None)
868 .await
869 .unwrap();
870
871 let archive = src.export_kg(&tok).await.unwrap();
872 assert_eq!(archive.entities.len(), 3);
873 assert_eq!(archive.edges.len(), 2);
874 assert_eq!(archive.format, "khive-kg");
875 assert_eq!(archive.version, "0.1");
876
877 let dst = make_rt().await;
878 let summary = dst.import_kg(&archive, &tok).await.unwrap();
879 assert_eq!(summary.entities_imported, 3);
880 assert_eq!(summary.edges_imported, 2);
881
882 let got = dst.get_entity(&tok, e1.id).await.unwrap();
883 assert_eq!(got.name, "FlashAttention");
884 assert_eq!(got.description.as_deref(), Some("fast attention"));
885 }
886
887 #[tokio::test]
889 async fn json_roundtrip() {
890 let src = make_rt().await;
891 let tok = NamespaceToken::local();
892 let e1 = src
893 .create_entity(
894 &tok,
895 "concept",
896 None,
897 "LoRA",
898 Some("low-rank adaptation"),
899 Some(serde_json::json!({"year": "2021"})),
900 vec!["fine-tuning".into()],
901 )
902 .await
903 .unwrap();
904 let e2 = src
905 .create_entity(&tok, "concept", None, "QLoRA", None, None, vec![])
906 .await
907 .unwrap();
908 src.link(&tok, e2.id, e1.id, EdgeRelation::VariantOf, 0.9, None)
909 .await
910 .unwrap();
911
912 let json_str = src.export_kg_json(&tok).await.unwrap();
913 assert!(json_str.contains("khive-kg"));
914
915 let dst = make_rt().await;
916 let summary = dst.import_kg_json(&json_str, &tok).await.unwrap();
917 assert_eq!(summary.entities_imported, 2);
918 assert_eq!(summary.edges_imported, 1);
919
920 let got = dst.get_entity(&tok, e1.id).await.unwrap();
921 assert_eq!(got.tags, vec!["fine-tuning"]);
922 }
923
924 #[tokio::test]
931 async fn namespace_targeting() {
932 let src = make_rt().await;
933 let tok_a = NamespaceToken::for_namespace(Namespace::parse("a").unwrap());
934 let tok_b = NamespaceToken::for_namespace(Namespace::parse("b").unwrap());
935 src.create_entity(&tok_a, "concept", None, "Sinkhorn", None, None, vec![])
936 .await
937 .unwrap();
938
939 let archive = src.export_kg(&tok_a).await.unwrap();
940 assert_eq!(archive.namespace, "a");
941
942 let dst = make_rt().await;
943 let summary = dst.import_kg(&archive, &tok_b).await.unwrap();
944 assert_eq!(summary.entities_imported, 1);
945
946 let in_b = dst.list_entities(&tok_b, None, None, 100, 0).await.unwrap();
947 assert_eq!(in_b.len(), 1);
948 assert_eq!(in_b[0].name, "Sinkhorn");
949
950 let in_a = src.list_entities(&tok_a, None, None, 100, 0).await.unwrap();
951 assert_eq!(in_a.len(), 1);
952
953 let dst_a = dst.list_entities(&tok_a, None, None, 100, 0).await.unwrap();
954 assert_eq!(dst_a.len(), 0);
955 }
956
957 #[tokio::test]
959 async fn format_validation_rejects_wrong_format() {
960 let rt = make_rt().await;
961 let tok = NamespaceToken::local();
962 let bad = KgArchive {
963 format: "wrong".to_string(),
964 version: "0.1".to_string(),
965 namespace: "local".to_string(),
966 exported_at: Utc::now(),
967 entities: vec![],
968 edges: vec![],
969 };
970 let err = rt.import_kg(&bad, &tok).await.unwrap_err();
971 assert!(matches!(err, RuntimeError::InvalidInput(_)));
972 }
973
974 #[tokio::test]
976 async fn import_unsupported_archive_version_returns_error() {
977 let rt = make_rt().await;
978 let tok = NamespaceToken::local();
979 let bad = KgArchive {
980 format: "khive-kg".to_string(),
981 version: "999.0".to_string(),
982 namespace: "local".to_string(),
983 exported_at: Utc::now(),
984 entities: vec![],
985 edges: vec![],
986 };
987 let err = rt.import_kg(&bad, &tok).await.unwrap_err();
988 assert!(
989 matches!(err, RuntimeError::InvalidInput(_)),
990 "expected InvalidInput, got {err:?}"
991 );
992 if let RuntimeError::InvalidInput(msg) = err {
993 assert!(
994 msg.contains("999.0"),
995 "error message should mention the unsupported version, got: {msg:?}"
996 );
997 }
998 }
999
1000 #[tokio::test]
1004 async fn import_entity_with_reserved_secret_gate_property_is_rejected() {
1005 let rt = make_rt().await;
1006 let tok = NamespaceToken::local();
1007 let archive = KgArchive {
1008 format: "khive-kg".to_string(),
1009 version: "0.1".to_string(),
1010 namespace: "local".to_string(),
1011 exported_at: Utc::now(),
1012 entities: vec![ExportedEntity {
1013 id: Uuid::new_v4(),
1014 kind: "concept".to_string(),
1015 entity_type: None,
1016 name: "ReservedKeyImport".to_string(),
1017 description: None,
1018 properties: Some(serde_json::json!({
1019 "khive:secret_gate": "exempted:content-sha256-manifest-v1"
1020 })),
1021 tags: vec![],
1022 created_at: Utc::now(),
1023 updated_at: Utc::now(),
1024 }],
1025 edges: vec![],
1026 };
1027
1028 let err = rt.import_kg(&archive, &tok).await.unwrap_err();
1029 assert!(
1030 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
1031 "expected a reservation rejection, got {err:?}"
1032 );
1033
1034 let entities = rt.list_entities(&tok, None, None, 100, 0).await.unwrap();
1035 assert!(
1036 !entities.iter().any(|e| e.name == "ReservedKeyImport"),
1037 "rejected archive import must not create the entity"
1038 );
1039 }
1040
1041 #[tokio::test]
1046 async fn import_edge_with_reserved_secret_gate_property_is_rejected() {
1047 let rt = make_rt().await;
1048 let tok = NamespaceToken::local();
1049 let a = Uuid::new_v4();
1050 let b = Uuid::new_v4();
1051 let mk = |id: Uuid, name: &str| ExportedEntity {
1052 id,
1053 kind: "concept".to_string(),
1054 entity_type: None,
1055 name: name.to_string(),
1056 description: None,
1057 properties: None,
1058 tags: vec![],
1059 created_at: Utc::now(),
1060 updated_at: Utc::now(),
1061 };
1062 let archive = KgArchive {
1063 format: "khive-kg".to_string(),
1064 version: "0.1".to_string(),
1065 namespace: "local".to_string(),
1066 exported_at: Utc::now(),
1067 entities: vec![mk(a, "EdgeGateA"), mk(b, "EdgeGateB")],
1068 edges: vec![ExportedEdge {
1069 edge_id: Uuid::new_v4(),
1070 source: a,
1071 target: b,
1072 relation: EdgeRelation::Extends,
1073 weight: 0.5,
1074 properties: Some(serde_json::json!({
1075 "khive:secret_gate": "exempted:content-sha256-manifest-v1"
1076 })),
1077 created_at: Utc::now(),
1078 updated_at: Utc::now(),
1079 }],
1080 };
1081
1082 let err = rt.import_kg(&archive, &tok).await.unwrap_err();
1083 assert!(
1084 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
1085 "expected a reservation rejection, got {err:?}"
1086 );
1087
1088 let entities = rt.list_entities(&tok, None, None, 100, 0).await.unwrap();
1091 assert!(
1092 !entities.iter().any(|e| e.name == "EdgeGateA"),
1093 "rejected archive import must not create any records"
1094 );
1095 }
1096
1097 #[tokio::test]
1101 async fn import_edge_with_credential_shaped_property_is_rejected() {
1102 let rt = make_rt().await;
1103 let tok = NamespaceToken::local();
1104 let a = Uuid::new_v4();
1105 let b = Uuid::new_v4();
1106 let mk = |id: Uuid, name: &str| ExportedEntity {
1107 id,
1108 kind: "concept".to_string(),
1109 entity_type: None,
1110 name: name.to_string(),
1111 description: None,
1112 properties: None,
1113 tags: vec![],
1114 created_at: Utc::now(),
1115 updated_at: Utc::now(),
1116 };
1117 let archive = KgArchive {
1118 format: "khive-kg".to_string(),
1119 version: "0.1".to_string(),
1120 namespace: "local".to_string(),
1121 exported_at: Utc::now(),
1122 entities: vec![mk(a, "EdgeScanA"), mk(b, "EdgeScanB")],
1123 edges: vec![ExportedEdge {
1124 edge_id: Uuid::new_v4(),
1125 source: a,
1126 target: b,
1127 relation: EdgeRelation::Extends,
1128 weight: 0.5,
1129 properties: Some(serde_json::json!({
1130 "api_key": "AKIAFAKEKEY1234567890"
1131 })),
1132 created_at: Utc::now(),
1133 updated_at: Utc::now(),
1134 }],
1135 };
1136
1137 let err = rt.import_kg(&archive, &tok).await.unwrap_err();
1138 assert!(
1139 matches!(err, RuntimeError::SecretDetected(_)),
1140 "expected the blocking scanner to reject the value, got {err:?}"
1141 );
1142
1143 let entities = rt.list_entities(&tok, None, None, 100, 0).await.unwrap();
1144 assert!(
1145 !entities.iter().any(|e| e.name == "EdgeScanA"),
1146 "rejected archive import must not create any records"
1147 );
1148 }
1149
1150 #[test]
1152 fn invalid_relation_rejected_at_deserialize() {
1153 let json = r#"{
1154 "format":"khive-kg","version":"0.1","namespace":"local",
1155 "exported_at":"2026-01-01T00:00:00Z",
1156 "entities":[],
1157 "edges":[{"edge_id":"00000000-0000-0000-0000-000000000099",
1158 "source":"00000000-0000-0000-0000-000000000001",
1159 "target":"00000000-0000-0000-0000-000000000002",
1160 "relation":"related_to","weight":0.5}]
1161 }"#;
1162 let result: Result<KgArchive, _> = serde_json::from_str(json);
1163 assert!(
1164 result.is_err(),
1165 "non-canonical relation should fail to deserialize"
1166 );
1167 }
1168
1169 #[tokio::test]
1176 async fn import_edge_with_dangling_source_is_skipped() {
1177 let phantom_source = Uuid::parse_str("deadbeef-dead-4ead-dead-deadbeefcafe").unwrap();
1178
1179 let rt = make_rt().await;
1180 let tok = NamespaceToken::local();
1181 let real = rt
1182 .create_entity(&tok, "concept", None, "Real", None, None, vec![])
1183 .await
1184 .unwrap();
1185
1186 let archive = KgArchive {
1187 format: "khive-kg".to_string(),
1188 version: "0.1".to_string(),
1189 namespace: "local".to_string(),
1190 exported_at: Utc::now(),
1191 entities: vec![ExportedEntity {
1192 id: real.id,
1193 kind: "concept".to_string(),
1194 entity_type: None,
1195 name: "Real".to_string(),
1196 description: None,
1197 properties: None,
1198 tags: vec![],
1199 created_at: Utc::now(),
1200 updated_at: Utc::now(),
1201 }],
1202 edges: vec![ExportedEdge {
1203 edge_id: Uuid::new_v4(),
1204 source: phantom_source,
1205 target: real.id,
1206 relation: EdgeRelation::Extends,
1207 weight: 1.0,
1208 properties: None,
1209 created_at: Utc::now(),
1210 updated_at: Utc::now(),
1211 }],
1212 };
1213
1214 let dst = make_rt().await;
1215 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1216 assert_eq!(summary.entities_imported, 1);
1217 assert_eq!(
1218 summary.edges_imported, 0,
1219 "dangling source must not be imported"
1220 );
1221 assert_eq!(
1222 summary.edges_skipped, 1,
1223 "dangling source must be counted as skipped"
1224 );
1225 }
1226
1227 #[tokio::test]
1232 async fn import_edge_with_dangling_target_is_skipped() {
1233 let phantom_target = Uuid::parse_str("cafebabe-cafe-4abe-cafe-cafebabecafe").unwrap();
1234
1235 let rt = make_rt().await;
1236 let tok = NamespaceToken::local();
1237 let real = rt
1238 .create_entity(&tok, "concept", None, "Source", None, None, vec![])
1239 .await
1240 .unwrap();
1241
1242 let archive = KgArchive {
1243 format: "khive-kg".to_string(),
1244 version: "0.1".to_string(),
1245 namespace: "local".to_string(),
1246 exported_at: Utc::now(),
1247 entities: vec![ExportedEntity {
1248 id: real.id,
1249 kind: "concept".to_string(),
1250 entity_type: None,
1251 name: "Source".to_string(),
1252 description: None,
1253 properties: None,
1254 tags: vec![],
1255 created_at: Utc::now(),
1256 updated_at: Utc::now(),
1257 }],
1258 edges: vec![ExportedEdge {
1259 edge_id: Uuid::new_v4(),
1260 source: real.id,
1261 target: phantom_target,
1262 relation: EdgeRelation::DependsOn,
1263 weight: 0.8,
1264 properties: None,
1265 created_at: Utc::now(),
1266 updated_at: Utc::now(),
1267 }],
1268 };
1269
1270 let dst = make_rt().await;
1271 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1272 assert_eq!(summary.entities_imported, 1);
1273 assert_eq!(
1274 summary.edges_imported, 0,
1275 "dangling target must not be imported"
1276 );
1277 assert_eq!(
1278 summary.edges_skipped, 1,
1279 "dangling target must be counted as skipped"
1280 );
1281 }
1282
1283 #[tokio::test]
1288 async fn import_mixed_edges_reports_correct_counts() {
1289 let phantom = Uuid::parse_str("11111111-1111-4111-8111-111111111111").unwrap();
1290
1291 let src = make_rt().await;
1292 let tok = NamespaceToken::local();
1293 let a = src
1294 .create_entity(&tok, "concept", None, "A", None, None, vec![])
1295 .await
1296 .unwrap();
1297 let b = src
1298 .create_entity(&tok, "concept", None, "B", None, None, vec![])
1299 .await
1300 .unwrap();
1301 let c = src
1302 .create_entity(&tok, "concept", None, "C", None, None, vec![])
1303 .await
1304 .unwrap();
1305
1306 let archive = KgArchive {
1307 format: "khive-kg".to_string(),
1308 version: "0.1".to_string(),
1309 namespace: "local".to_string(),
1310 exported_at: Utc::now(),
1311 entities: vec![
1312 ExportedEntity {
1313 id: a.id,
1314 kind: "concept".to_string(),
1315 entity_type: None,
1316 name: "A".to_string(),
1317 description: None,
1318 properties: None,
1319 tags: vec![],
1320 created_at: Utc::now(),
1321 updated_at: Utc::now(),
1322 },
1323 ExportedEntity {
1324 id: b.id,
1325 kind: "concept".to_string(),
1326 entity_type: None,
1327 name: "B".to_string(),
1328 description: None,
1329 properties: None,
1330 tags: vec![],
1331 created_at: Utc::now(),
1332 updated_at: Utc::now(),
1333 },
1334 ExportedEntity {
1335 id: c.id,
1336 kind: "concept".to_string(),
1337 entity_type: None,
1338 name: "C".to_string(),
1339 description: None,
1340 properties: None,
1341 tags: vec![],
1342 created_at: Utc::now(),
1343 updated_at: Utc::now(),
1344 },
1345 ],
1346 edges: vec![
1347 ExportedEdge {
1348 edge_id: Uuid::new_v4(),
1349 source: a.id,
1350 target: b.id,
1351 relation: EdgeRelation::Extends,
1352 weight: 1.0,
1353 properties: None,
1354 created_at: Utc::now(),
1355 updated_at: Utc::now(),
1356 },
1357 ExportedEdge {
1358 edge_id: Uuid::new_v4(),
1359 source: b.id,
1360 target: c.id,
1361 relation: EdgeRelation::VariantOf,
1366 weight: 0.9,
1367 properties: None,
1368 created_at: Utc::now(),
1369 updated_at: Utc::now(),
1370 },
1371 ExportedEdge {
1372 edge_id: Uuid::new_v4(),
1373 source: a.id,
1374 target: phantom,
1375 relation: EdgeRelation::Enables,
1376 weight: 0.5,
1377 properties: None,
1378 created_at: Utc::now(),
1379 updated_at: Utc::now(),
1380 },
1381 ],
1382 };
1383
1384 let dst = make_rt().await;
1385 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1386 assert_eq!(summary.entities_imported, 3);
1387 assert_eq!(
1388 summary.edges_imported, 2,
1389 "only valid edges must be imported"
1390 );
1391 assert_eq!(
1392 summary.edges_skipped, 1,
1393 "one dangling edge must be reported"
1394 );
1395 }
1396
1397 #[tokio::test]
1399 async fn import_all_valid_edges_reports_zero_skipped() {
1400 let src = make_rt().await;
1401 let tok = NamespaceToken::local();
1402 let e1 = src
1403 .create_entity(&tok, "concept", None, "E1", None, None, vec![])
1404 .await
1405 .unwrap();
1406 let e2 = src
1407 .create_entity(&tok, "concept", None, "E2", None, None, vec![])
1408 .await
1409 .unwrap();
1410 src.link(&tok, e1.id, e2.id, EdgeRelation::VariantOf, 0.7, None)
1411 .await
1412 .unwrap();
1413
1414 let archive = src.export_kg(&tok).await.unwrap();
1415 let dst = make_rt().await;
1416 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1417 assert_eq!(summary.edges_imported, 1);
1418 assert_eq!(
1419 summary.edges_skipped, 0,
1420 "no edges should be skipped when all endpoints exist"
1421 );
1422 }
1423
1424 #[tokio::test]
1430 async fn import_edge_violating_endpoint_contract_is_skipped() {
1431 let src = make_rt().await;
1432 let tok = NamespaceToken::local();
1433 let e1 = src
1434 .create_entity(&tok, "concept", None, "E1", None, None, vec![])
1435 .await
1436 .unwrap();
1437 let e2 = src
1438 .create_entity(&tok, "concept", None, "E2", None, None, vec![])
1439 .await
1440 .unwrap();
1441
1442 let archive = KgArchive {
1443 format: "khive-kg".to_string(),
1444 version: "0.1".to_string(),
1445 namespace: "local".to_string(),
1446 exported_at: Utc::now(),
1447 entities: vec![
1448 ExportedEntity {
1449 id: e1.id,
1450 kind: "concept".to_string(),
1451 entity_type: None,
1452 name: "E1".to_string(),
1453 description: None,
1454 properties: None,
1455 tags: vec![],
1456 created_at: Utc::now(),
1457 updated_at: Utc::now(),
1458 },
1459 ExportedEntity {
1460 id: e2.id,
1461 kind: "concept".to_string(),
1462 entity_type: None,
1463 name: "E2".to_string(),
1464 description: None,
1465 properties: None,
1466 tags: vec![],
1467 created_at: Utc::now(),
1468 updated_at: Utc::now(),
1469 },
1470 ],
1471 edges: vec![ExportedEdge {
1472 edge_id: Uuid::new_v4(),
1473 source: e1.id,
1474 target: e2.id,
1475 relation: EdgeRelation::Precedes,
1476 weight: 1.0,
1477 properties: None,
1478 created_at: Utc::now(),
1479 updated_at: Utc::now(),
1480 }],
1481 };
1482
1483 let dst = make_rt().await;
1484 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1485 assert_eq!(summary.entities_imported, 2);
1486 assert_eq!(
1487 summary.edges_imported, 0,
1488 "an endpoint-contract-violating edge must not be imported"
1489 );
1490 assert_eq!(
1491 summary.edges_skipped, 1,
1492 "an endpoint-contract-violating edge must be counted as skipped"
1493 );
1494 assert!(
1495 dst.neighbors(
1496 &tok,
1497 e1.id,
1498 khive_storage::types::Direction::Out,
1499 None,
1500 None
1501 )
1502 .await
1503 .unwrap()
1504 .is_empty(),
1505 "the contract-violating edge must not exist in the destination graph"
1506 );
1507 }
1508
1509 #[tokio::test]
1513 async fn export_kg_preserves_edge_id() {
1514 let rt = make_rt().await;
1515 let tok = NamespaceToken::local();
1516 let a = rt
1517 .create_entity(&tok, "concept", None, "Alpha", None, None, vec![])
1518 .await
1519 .unwrap();
1520 let b = rt
1521 .create_entity(&tok, "concept", None, "Beta", None, None, vec![])
1522 .await
1523 .unwrap();
1524 let stored_edge = rt
1525 .link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
1526 .await
1527 .unwrap();
1528 let stored_id: Uuid = stored_edge.id.into();
1529
1530 let archive = rt.export_kg(&tok).await.unwrap();
1531 assert_eq!(archive.edges.len(), 1);
1532 assert_eq!(
1533 archive.edges[0].edge_id, stored_id,
1534 "exported edge_id must equal the LinkId returned by link"
1535 );
1536 }
1537
1538 #[tokio::test]
1540 async fn import_kg_persists_edge_id_and_timestamps() {
1541 let src = make_rt().await;
1542 let tok = NamespaceToken::local();
1543 let a = src
1544 .create_entity(&tok, "concept", None, "Alpha", None, None, vec![])
1545 .await
1546 .unwrap();
1547 let b = src
1548 .create_entity(&tok, "concept", None, "Beta", None, None, vec![])
1549 .await
1550 .unwrap();
1551 let stored_edge = src
1552 .link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
1553 .await
1554 .unwrap();
1555 let original_id: Uuid = stored_edge.id.into();
1556
1557 let expected_created = chrono::DateTime::parse_from_rfc3339("2026-03-03T00:00:00Z")
1558 .unwrap()
1559 .with_timezone(&Utc);
1560 let expected_updated = chrono::DateTime::parse_from_rfc3339("2026-04-04T00:00:00Z")
1561 .unwrap()
1562 .with_timezone(&Utc);
1563 let mut archive = src.export_kg(&tok).await.unwrap();
1564 archive.edges[0].created_at = expected_created;
1565 archive.edges[0].updated_at = expected_updated;
1566 archive.edges[0].properties = Some(serde_json::json!({"confidence": 0.95}));
1567 let dst = make_rt().await;
1568 dst.import_kg(&archive, &tok).await.unwrap();
1569
1570 let imported_edge = dst.get_edge(&tok, original_id).await.unwrap();
1571 assert!(
1572 imported_edge.is_some(),
1573 "imported edge must be retrievable by the original edge_id"
1574 );
1575 let imported_edge = imported_edge.unwrap();
1576 assert_eq!(
1577 Uuid::from(imported_edge.id),
1578 original_id,
1579 "stored edge id must equal the archive edge_id"
1580 );
1581 assert_eq!(
1582 imported_edge.created_at, expected_created,
1583 "present edge created_at must not be replaced with import time"
1584 );
1585 assert_eq!(
1586 imported_edge.updated_at, expected_updated,
1587 "present edge updated_at must not be replaced with created_at or import time"
1588 );
1589 assert_eq!(
1590 imported_edge.metadata,
1591 Some(serde_json::json!({"confidence": 0.95})),
1592 "edge properties must persist as storage metadata"
1593 );
1594
1595 let reexported = dst.export_kg(&tok).await.unwrap();
1596 assert_eq!(
1597 reexported.edges[0].properties,
1598 Some(serde_json::json!({"confidence": 0.95})),
1599 "edge metadata must re-export as portable properties"
1600 );
1601 }
1602
1603 #[tokio::test]
1608 async fn old_archive_missing_edge_id_round_trips() {
1609 let src_id = Uuid::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
1610 let tgt_id = Uuid::parse_str("00000000-0000-0000-0000-000000000002").unwrap();
1611
1612 let json = format!(
1614 r#"{{
1615 "format": "khive-kg",
1616 "version": "0.1",
1617 "namespace": "local",
1618 "exported_at": "2026-01-01T00:00:00Z",
1619 "entities": [
1620 {{"id":"{src_id}","kind":"concept","name":"SrcNode","created_at":"2026-01-01T00:00:00Z","updated_at":"2026-01-01T00:00:00Z"}},
1621 {{"id":"{tgt_id}","kind":"concept","name":"TgtNode","created_at":"2026-01-01T00:00:00Z","updated_at":"2026-01-01T00:00:00Z"}}
1622 ],
1623 "edges": [
1624 {{
1625 "source": "{src_id}",
1626 "target": "{tgt_id}",
1627 "relation": "extends",
1628 "weight": 0.9
1629 }}
1630 ]
1631 }}"#
1632 );
1633
1634 let archive: KgArchive = serde_json::from_str(&json)
1636 .expect("old archive without edge_id must deserialize successfully");
1637 assert_eq!(archive.edges.len(), 1);
1638 let generated_id = archive.edges[0].edge_id;
1639 assert_ne!(
1640 generated_id,
1641 Uuid::nil(),
1642 "missing edge_id in old archive must get a fresh non-nil UUID"
1643 );
1644
1645 let rt = make_rt().await;
1646 let tok = NamespaceToken::local();
1647 let summary = rt.import_kg(&archive, &tok).await.unwrap();
1648 assert_eq!(summary.entities_imported, 2);
1649 assert_eq!(
1650 summary.edges_imported, 1,
1651 "edge must be imported when both endpoints exist"
1652 );
1653
1654 let stored = rt.get_edge(&tok, generated_id).await.unwrap();
1655 assert!(
1656 stored.is_some(),
1657 "imported edge must be retrievable by the generated edge_id"
1658 );
1659 assert_eq!(
1660 Uuid::from(stored.unwrap().id),
1661 generated_id,
1662 "stored edge id must equal the generated edge_id"
1663 );
1664
1665 let re_archive = rt.export_kg(&tok).await.unwrap();
1666 assert_eq!(re_archive.edges.len(), 1);
1667 assert_eq!(
1668 re_archive.edges[0].edge_id, generated_id,
1669 "re-exported edge_id must equal the ID generated on first import"
1670 );
1671 }
1672
1673 #[tokio::test]
1678 async fn export_import_export_edge_id_equality() {
1679 let src = make_rt().await;
1680 let tok = NamespaceToken::local();
1681 let a = src
1682 .create_entity(&tok, "concept", None, "NodeA", None, None, vec![])
1683 .await
1684 .unwrap();
1685 let b = src
1686 .create_entity(&tok, "concept", None, "NodeB", None, None, vec![])
1687 .await
1688 .unwrap();
1689 let stored = src
1690 .link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
1691 .await
1692 .unwrap();
1693 let original_edge_id: Uuid = stored.id.into();
1694
1695 let archive1 = src.export_kg(&tok).await.unwrap();
1696 assert_eq!(archive1.edges.len(), 1);
1697 assert_eq!(
1698 archive1.edges[0].edge_id, original_edge_id,
1699 "first export must carry the stored edge_id"
1700 );
1701
1702 let dst = make_rt().await;
1703 dst.import_kg(&archive1, &tok).await.unwrap();
1704
1705 let archive2 = dst.export_kg(&tok).await.unwrap();
1706 assert_eq!(archive2.edges.len(), 1);
1707
1708 let re_edge = archive2
1709 .edges
1710 .iter()
1711 .find(|e| e.source == a.id && e.target == b.id && e.relation == EdgeRelation::Extends)
1712 .expect(
1713 "re-exported archive must contain the original edge by (source,target,relation)",
1714 );
1715 assert_eq!(
1716 re_edge.edge_id, original_edge_id,
1717 "edge_id must be identical across export → import → export"
1718 );
1719 }
1720}