1use std::collections::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 ee in &archive.entities {
279 let created_micros = ee.created_at.timestamp_micros();
280 let updated_micros = ee.updated_at.timestamp_micros();
281 let entity = khive_storage::entity::Entity {
282 id: ee.id,
283 namespace: ns.clone(),
284 kind: ee.kind.clone(),
285 entity_type: ee.entity_type.clone(),
286 name: ee.name.clone(),
287 description: ee.description.clone(),
288 properties: ee.properties.clone(),
289 tags: ee.tags.clone(),
290 created_at: created_micros,
291 updated_at: updated_micros,
292 deleted_at: None,
293 merged_into: None,
294 merge_event_id: None,
295 version: 1,
296 content_ref: None,
297 };
298 store.upsert_entity(entity.clone()).await?;
299 embedding_truncation.merge(self.reindex_entity(token, &entity).await?);
301 entities_imported += 1;
302 }
303
304 let graph = self.graph(token)?;
308 let mut edges_imported = 0usize;
309 let mut edges_skipped = 0usize;
310 for ee in &archive.edges {
311 let source_ok = match self.get_entity(token, ee.source).await {
312 Ok(_) => true,
313 Err(RuntimeError::NotFound(_)) => false,
314 Err(e) => return Err(e),
315 };
316 if !source_ok {
317 tracing::warn!(
318 source = %ee.source,
319 target = %ee.target,
320 relation = ?ee.relation,
321 "import_kg: skipping edge — source entity not found in namespace {ns:?}"
322 );
323 edges_skipped += 1;
324 continue;
325 }
326 let target_ok = match self.get_entity(token, ee.target).await {
327 Ok(_) => true,
328 Err(RuntimeError::NotFound(_)) => false,
329 Err(e) => return Err(e),
330 };
331 if !target_ok {
332 tracing::warn!(
333 source = %ee.source,
334 target = %ee.target,
335 relation = ?ee.relation,
336 "import_kg: skipping edge — target entity not found in namespace {ns:?}"
337 );
338 edges_skipped += 1;
339 continue;
340 }
341 match self
346 .validate_edge_relation_endpoints(token, ee.source, ee.target, ee.relation)
347 .await
348 {
349 Ok(()) => {}
350 Err(e @ (RuntimeError::InvalidInput(_) | RuntimeError::NotFound(_))) => {
351 tracing::warn!(
352 source = %ee.source,
353 target = %ee.target,
354 relation = ?ee.relation,
355 error = %e,
356 "import_kg: skipping edge — endpoint contract violation in namespace {ns:?}"
357 );
358 edges_skipped += 1;
359 continue;
360 }
361 Err(e) => return Err(e),
362 }
363 let edge = khive_storage::types::Edge {
364 id: LinkId::from(ee.edge_id),
365 namespace: ns.clone(),
366 source_id: ee.source,
367 target_id: ee.target,
368 relation: ee.relation,
369 weight: ee.weight,
370 created_at: ee.created_at,
371 updated_at: ee.updated_at,
372 deleted_at: None,
373 metadata: ee.properties.clone(),
374 target_backend: None,
375 };
376 graph.upsert_edge(edge).await?;
377 edges_imported += 1;
378 }
379
380 Ok(ImportSummary {
381 entities_imported,
382 edges_imported,
383 edges_skipped,
384 embedding_truncation,
385 })
386 }
387
388 pub async fn import_kg_json(
390 &self,
391 json: &str,
392 token: &NamespaceToken,
393 ) -> RuntimeResult<ImportSummary> {
394 let archive: KgArchive =
395 serde_json::from_str(json).map_err(|e| RuntimeError::InvalidInput(e.to_string()))?;
396 self.import_kg(&archive, token).await
397 }
398}
399
400#[cfg(test)]
405mod tests {
406 use std::sync::Arc;
407
408 use super::*;
409 use crate::runtime::{KhiveRuntime, NamespaceToken};
410 use crate::{EmbedderProvider, Namespace};
411 use async_trait::async_trait;
412 use khive_storage::EdgeRelation;
413 use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService, MAX_TEXT_BYTES};
414
415 const IMPORT_TEST_MODEL: &str = "all-minilm-l6-v2";
416
417 struct ImportEmbeddingService;
418
419 #[async_trait]
420 impl EmbeddingService for ImportEmbeddingService {
421 async fn embed(
422 &self,
423 texts: &[String],
424 _model: EmbeddingModel,
425 ) -> std::result::Result<Vec<Vec<f32>>, EmbedError> {
426 let dimensions = EmbeddingModel::AllMiniLmL6V2.dimensions();
427 Ok(texts.iter().map(|_| vec![1.0; dimensions]).collect())
428 }
429
430 fn supports_model(&self, _model: EmbeddingModel) -> bool {
431 true
432 }
433
434 fn name(&self) -> &'static str {
435 "kg-import-truncation-test"
436 }
437 }
438
439 struct ImportEmbeddingProvider;
440
441 #[async_trait]
442 impl EmbedderProvider for ImportEmbeddingProvider {
443 fn name(&self) -> &str {
444 IMPORT_TEST_MODEL
445 }
446
447 fn dimensions(&self) -> usize {
448 EmbeddingModel::AllMiniLmL6V2.dimensions()
449 }
450
451 async fn build(&self) -> std::result::Result<Arc<dyn EmbeddingService>, RuntimeError> {
452 Ok(Arc::new(ImportEmbeddingService))
453 }
454 }
455
456 async fn make_rt() -> KhiveRuntime {
457 KhiveRuntime::memory().expect("in-memory runtime")
458 }
459
460 fn make_rt_with_embedder() -> KhiveRuntime {
461 let runtime = KhiveRuntime::memory().expect("in-memory runtime");
462 runtime.register_embedder(ImportEmbeddingProvider);
463 runtime
464 }
465
466 #[tokio::test]
467 async fn import_summary_reports_actual_embedding_truncation() {
468 let runtime = make_rt_with_embedder();
469 let token = NamespaceToken::local();
470 let archive = KgArchive {
471 format: "khive-kg".to_string(),
472 version: "0.1".to_string(),
473 namespace: "local".to_string(),
474 exported_at: Utc::now(),
475 entities: vec![ExportedEntity {
476 id: Uuid::new_v4(),
477 kind: "concept".to_string(),
478 entity_type: None,
479 name: "Imported long entity".to_string(),
480 description: Some("x".repeat(MAX_TEXT_BYTES + 1)),
481 properties: None,
482 tags: vec![],
483 created_at: Utc::now(),
484 updated_at: Utc::now(),
485 }],
486 edges: vec![],
487 };
488
489 let summary = runtime
490 .import_kg(&archive, &token)
491 .await
492 .expect("import with bounded embedding input");
493
494 assert_eq!(summary.entities_imported, 1);
495 assert_eq!(summary.embedding_truncation.truncated, 1);
496 assert!(summary.embedding_truncation.discarded_bytes > 0);
497 let wire = serde_json::to_value(&summary).expect("serialize import summary");
498 assert_eq!(wire["embedding_truncation"]["truncated"], 1);
499 }
500
501 #[tokio::test]
502 async fn import_rejects_whitespace_name_before_any_entity_write() {
503 let runtime = make_rt().await;
504 let token = NamespaceToken::local();
505 let valid_id = Uuid::new_v4();
506 let archive = KgArchive {
507 format: "khive-kg".to_string(),
508 version: "0.1".to_string(),
509 namespace: "local".to_string(),
510 exported_at: Utc::now(),
511 entities: vec![
512 ExportedEntity {
513 id: valid_id,
514 kind: "concept".to_string(),
515 entity_type: None,
516 name: "Must not be written".to_string(),
517 description: None,
518 properties: None,
519 tags: vec![],
520 created_at: Utc::now(),
521 updated_at: Utc::now(),
522 },
523 ExportedEntity {
524 id: Uuid::new_v4(),
525 kind: "concept".to_string(),
526 entity_type: None,
527 name: " \t\n ".to_string(),
528 description: None,
529 properties: None,
530 tags: vec![],
531 created_at: Utc::now(),
532 updated_at: Utc::now(),
533 },
534 ],
535 edges: vec![],
536 };
537
538 let err = runtime
539 .import_kg(&archive, &token)
540 .await
541 .expect_err("whitespace-only entity names must fail the whole import");
542 assert!(
543 err.to_string().contains("non-blank"),
544 "error must explain the name invariant: {err}"
545 );
546 assert!(
547 runtime.get_entity(&token, valid_id).await.is_err(),
548 "deterministic validation must finish before the first entity write"
549 );
550 }
551
552 #[tokio::test]
554 async fn roundtrip_entities_and_edges() {
555 let src = make_rt().await;
556 let tok = NamespaceToken::local();
557 let e1 = src
558 .create_entity(
559 &tok,
560 "concept",
561 None,
562 "FlashAttention",
563 Some("fast attention"),
564 None,
565 vec![],
566 )
567 .await
568 .unwrap();
569 let e2 = src
570 .create_entity(
571 &tok,
572 "concept",
573 None,
574 "FlashAttention-2",
575 None,
576 None,
577 vec![],
578 )
579 .await
580 .unwrap();
581 let e3 = src
582 .create_entity(
583 &tok,
584 "person",
585 None,
586 "Tri Dao",
587 None,
588 None,
589 vec!["author".into()],
590 )
591 .await
592 .unwrap();
593 src.link(&tok, e2.id, e1.id, EdgeRelation::Extends, 1.0, None)
594 .await
595 .unwrap();
596 src.link(&tok, e1.id, e3.id, EdgeRelation::IntroducedBy, 0.9, None)
597 .await
598 .unwrap();
599
600 let archive = src.export_kg(&tok).await.unwrap();
601 assert_eq!(archive.entities.len(), 3);
602 assert_eq!(archive.edges.len(), 2);
603 assert_eq!(archive.format, "khive-kg");
604 assert_eq!(archive.version, "0.1");
605
606 let dst = make_rt().await;
607 let summary = dst.import_kg(&archive, &tok).await.unwrap();
608 assert_eq!(summary.entities_imported, 3);
609 assert_eq!(summary.edges_imported, 2);
610
611 let got = dst.get_entity(&tok, e1.id).await.unwrap();
612 assert_eq!(got.name, "FlashAttention");
613 assert_eq!(got.description.as_deref(), Some("fast attention"));
614 }
615
616 #[tokio::test]
618 async fn json_roundtrip() {
619 let src = make_rt().await;
620 let tok = NamespaceToken::local();
621 let e1 = src
622 .create_entity(
623 &tok,
624 "concept",
625 None,
626 "LoRA",
627 Some("low-rank adaptation"),
628 Some(serde_json::json!({"year": "2021"})),
629 vec!["fine-tuning".into()],
630 )
631 .await
632 .unwrap();
633 let e2 = src
634 .create_entity(&tok, "concept", None, "QLoRA", None, None, vec![])
635 .await
636 .unwrap();
637 src.link(&tok, e2.id, e1.id, EdgeRelation::VariantOf, 0.9, None)
638 .await
639 .unwrap();
640
641 let json_str = src.export_kg_json(&tok).await.unwrap();
642 assert!(json_str.contains("khive-kg"));
643
644 let dst = make_rt().await;
645 let summary = dst.import_kg_json(&json_str, &tok).await.unwrap();
646 assert_eq!(summary.entities_imported, 2);
647 assert_eq!(summary.edges_imported, 1);
648
649 let got = dst.get_entity(&tok, e1.id).await.unwrap();
650 assert_eq!(got.tags, vec!["fine-tuning"]);
651 }
652
653 #[tokio::test]
660 async fn namespace_targeting() {
661 let src = make_rt().await;
662 let tok_a = NamespaceToken::for_namespace(Namespace::parse("a").unwrap());
663 let tok_b = NamespaceToken::for_namespace(Namespace::parse("b").unwrap());
664 src.create_entity(&tok_a, "concept", None, "Sinkhorn", None, None, vec![])
665 .await
666 .unwrap();
667
668 let archive = src.export_kg(&tok_a).await.unwrap();
669 assert_eq!(archive.namespace, "a");
670
671 let dst = make_rt().await;
672 let summary = dst.import_kg(&archive, &tok_b).await.unwrap();
673 assert_eq!(summary.entities_imported, 1);
674
675 let in_b = dst.list_entities(&tok_b, None, None, 100, 0).await.unwrap();
676 assert_eq!(in_b.len(), 1);
677 assert_eq!(in_b[0].name, "Sinkhorn");
678
679 let in_a = src.list_entities(&tok_a, None, None, 100, 0).await.unwrap();
680 assert_eq!(in_a.len(), 1);
681
682 let dst_a = dst.list_entities(&tok_a, None, None, 100, 0).await.unwrap();
683 assert_eq!(dst_a.len(), 0);
684 }
685
686 #[tokio::test]
688 async fn format_validation_rejects_wrong_format() {
689 let rt = make_rt().await;
690 let tok = NamespaceToken::local();
691 let bad = KgArchive {
692 format: "wrong".to_string(),
693 version: "0.1".to_string(),
694 namespace: "local".to_string(),
695 exported_at: Utc::now(),
696 entities: vec![],
697 edges: vec![],
698 };
699 let err = rt.import_kg(&bad, &tok).await.unwrap_err();
700 assert!(matches!(err, RuntimeError::InvalidInput(_)));
701 }
702
703 #[tokio::test]
705 async fn import_unsupported_archive_version_returns_error() {
706 let rt = make_rt().await;
707 let tok = NamespaceToken::local();
708 let bad = KgArchive {
709 format: "khive-kg".to_string(),
710 version: "999.0".to_string(),
711 namespace: "local".to_string(),
712 exported_at: Utc::now(),
713 entities: vec![],
714 edges: vec![],
715 };
716 let err = rt.import_kg(&bad, &tok).await.unwrap_err();
717 assert!(
718 matches!(err, RuntimeError::InvalidInput(_)),
719 "expected InvalidInput, got {err:?}"
720 );
721 if let RuntimeError::InvalidInput(msg) = err {
722 assert!(
723 msg.contains("999.0"),
724 "error message should mention the unsupported version, got: {msg:?}"
725 );
726 }
727 }
728
729 #[tokio::test]
733 async fn import_entity_with_reserved_secret_gate_property_is_rejected() {
734 let rt = make_rt().await;
735 let tok = NamespaceToken::local();
736 let archive = KgArchive {
737 format: "khive-kg".to_string(),
738 version: "0.1".to_string(),
739 namespace: "local".to_string(),
740 exported_at: Utc::now(),
741 entities: vec![ExportedEntity {
742 id: Uuid::new_v4(),
743 kind: "concept".to_string(),
744 entity_type: None,
745 name: "ReservedKeyImport".to_string(),
746 description: None,
747 properties: Some(serde_json::json!({
748 "khive:secret_gate": "exempted:content-sha256-manifest-v1"
749 })),
750 tags: vec![],
751 created_at: Utc::now(),
752 updated_at: Utc::now(),
753 }],
754 edges: vec![],
755 };
756
757 let err = rt.import_kg(&archive, &tok).await.unwrap_err();
758 assert!(
759 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
760 "expected a reservation rejection, got {err:?}"
761 );
762
763 let entities = rt.list_entities(&tok, None, None, 100, 0).await.unwrap();
764 assert!(
765 !entities.iter().any(|e| e.name == "ReservedKeyImport"),
766 "rejected archive import must not create the entity"
767 );
768 }
769
770 #[tokio::test]
775 async fn import_edge_with_reserved_secret_gate_property_is_rejected() {
776 let rt = make_rt().await;
777 let tok = NamespaceToken::local();
778 let a = Uuid::new_v4();
779 let b = Uuid::new_v4();
780 let mk = |id: Uuid, name: &str| ExportedEntity {
781 id,
782 kind: "concept".to_string(),
783 entity_type: None,
784 name: name.to_string(),
785 description: None,
786 properties: None,
787 tags: vec![],
788 created_at: Utc::now(),
789 updated_at: Utc::now(),
790 };
791 let archive = KgArchive {
792 format: "khive-kg".to_string(),
793 version: "0.1".to_string(),
794 namespace: "local".to_string(),
795 exported_at: Utc::now(),
796 entities: vec![mk(a, "EdgeGateA"), mk(b, "EdgeGateB")],
797 edges: vec![ExportedEdge {
798 edge_id: Uuid::new_v4(),
799 source: a,
800 target: b,
801 relation: EdgeRelation::Extends,
802 weight: 0.5,
803 properties: Some(serde_json::json!({
804 "khive:secret_gate": "exempted:content-sha256-manifest-v1"
805 })),
806 created_at: Utc::now(),
807 updated_at: Utc::now(),
808 }],
809 };
810
811 let err = rt.import_kg(&archive, &tok).await.unwrap_err();
812 assert!(
813 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
814 "expected a reservation rejection, got {err:?}"
815 );
816
817 let entities = rt.list_entities(&tok, None, None, 100, 0).await.unwrap();
820 assert!(
821 !entities.iter().any(|e| e.name == "EdgeGateA"),
822 "rejected archive import must not create any records"
823 );
824 }
825
826 #[tokio::test]
830 async fn import_edge_with_credential_shaped_property_is_rejected() {
831 let rt = make_rt().await;
832 let tok = NamespaceToken::local();
833 let a = Uuid::new_v4();
834 let b = Uuid::new_v4();
835 let mk = |id: Uuid, name: &str| ExportedEntity {
836 id,
837 kind: "concept".to_string(),
838 entity_type: None,
839 name: name.to_string(),
840 description: None,
841 properties: None,
842 tags: vec![],
843 created_at: Utc::now(),
844 updated_at: Utc::now(),
845 };
846 let archive = KgArchive {
847 format: "khive-kg".to_string(),
848 version: "0.1".to_string(),
849 namespace: "local".to_string(),
850 exported_at: Utc::now(),
851 entities: vec![mk(a, "EdgeScanA"), mk(b, "EdgeScanB")],
852 edges: vec![ExportedEdge {
853 edge_id: Uuid::new_v4(),
854 source: a,
855 target: b,
856 relation: EdgeRelation::Extends,
857 weight: 0.5,
858 properties: Some(serde_json::json!({
859 "api_key": "AKIAFAKEKEY1234567890"
860 })),
861 created_at: Utc::now(),
862 updated_at: Utc::now(),
863 }],
864 };
865
866 let err = rt.import_kg(&archive, &tok).await.unwrap_err();
867 assert!(
868 matches!(err, RuntimeError::SecretDetected(_)),
869 "expected the blocking scanner to reject the value, got {err:?}"
870 );
871
872 let entities = rt.list_entities(&tok, None, None, 100, 0).await.unwrap();
873 assert!(
874 !entities.iter().any(|e| e.name == "EdgeScanA"),
875 "rejected archive import must not create any records"
876 );
877 }
878
879 #[test]
881 fn invalid_relation_rejected_at_deserialize() {
882 let json = r#"{
883 "format":"khive-kg","version":"0.1","namespace":"local",
884 "exported_at":"2026-01-01T00:00:00Z",
885 "entities":[],
886 "edges":[{"edge_id":"00000000-0000-0000-0000-000000000099",
887 "source":"00000000-0000-0000-0000-000000000001",
888 "target":"00000000-0000-0000-0000-000000000002",
889 "relation":"related_to","weight":0.5}]
890 }"#;
891 let result: Result<KgArchive, _> = serde_json::from_str(json);
892 assert!(
893 result.is_err(),
894 "non-canonical relation should fail to deserialize"
895 );
896 }
897
898 #[tokio::test]
905 async fn import_edge_with_dangling_source_is_skipped() {
906 let phantom_source = Uuid::parse_str("deadbeef-dead-4ead-dead-deadbeefcafe").unwrap();
907
908 let rt = make_rt().await;
909 let tok = NamespaceToken::local();
910 let real = rt
911 .create_entity(&tok, "concept", None, "Real", None, None, vec![])
912 .await
913 .unwrap();
914
915 let archive = KgArchive {
916 format: "khive-kg".to_string(),
917 version: "0.1".to_string(),
918 namespace: "local".to_string(),
919 exported_at: Utc::now(),
920 entities: vec![ExportedEntity {
921 id: real.id,
922 kind: "concept".to_string(),
923 entity_type: None,
924 name: "Real".to_string(),
925 description: None,
926 properties: None,
927 tags: vec![],
928 created_at: Utc::now(),
929 updated_at: Utc::now(),
930 }],
931 edges: vec![ExportedEdge {
932 edge_id: Uuid::new_v4(),
933 source: phantom_source,
934 target: real.id,
935 relation: EdgeRelation::Extends,
936 weight: 1.0,
937 properties: None,
938 created_at: Utc::now(),
939 updated_at: Utc::now(),
940 }],
941 };
942
943 let dst = make_rt().await;
944 let summary = dst.import_kg(&archive, &tok).await.unwrap();
945 assert_eq!(summary.entities_imported, 1);
946 assert_eq!(
947 summary.edges_imported, 0,
948 "dangling source must not be imported"
949 );
950 assert_eq!(
951 summary.edges_skipped, 1,
952 "dangling source must be counted as skipped"
953 );
954 }
955
956 #[tokio::test]
961 async fn import_edge_with_dangling_target_is_skipped() {
962 let phantom_target = Uuid::parse_str("cafebabe-cafe-4abe-cafe-cafebabecafe").unwrap();
963
964 let rt = make_rt().await;
965 let tok = NamespaceToken::local();
966 let real = rt
967 .create_entity(&tok, "concept", None, "Source", None, None, vec![])
968 .await
969 .unwrap();
970
971 let archive = KgArchive {
972 format: "khive-kg".to_string(),
973 version: "0.1".to_string(),
974 namespace: "local".to_string(),
975 exported_at: Utc::now(),
976 entities: vec![ExportedEntity {
977 id: real.id,
978 kind: "concept".to_string(),
979 entity_type: None,
980 name: "Source".to_string(),
981 description: None,
982 properties: None,
983 tags: vec![],
984 created_at: Utc::now(),
985 updated_at: Utc::now(),
986 }],
987 edges: vec![ExportedEdge {
988 edge_id: Uuid::new_v4(),
989 source: real.id,
990 target: phantom_target,
991 relation: EdgeRelation::DependsOn,
992 weight: 0.8,
993 properties: None,
994 created_at: Utc::now(),
995 updated_at: Utc::now(),
996 }],
997 };
998
999 let dst = make_rt().await;
1000 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1001 assert_eq!(summary.entities_imported, 1);
1002 assert_eq!(
1003 summary.edges_imported, 0,
1004 "dangling target must not be imported"
1005 );
1006 assert_eq!(
1007 summary.edges_skipped, 1,
1008 "dangling target must be counted as skipped"
1009 );
1010 }
1011
1012 #[tokio::test]
1017 async fn import_mixed_edges_reports_correct_counts() {
1018 let phantom = Uuid::parse_str("11111111-1111-4111-8111-111111111111").unwrap();
1019
1020 let src = make_rt().await;
1021 let tok = NamespaceToken::local();
1022 let a = src
1023 .create_entity(&tok, "concept", None, "A", None, None, vec![])
1024 .await
1025 .unwrap();
1026 let b = src
1027 .create_entity(&tok, "concept", None, "B", None, None, vec![])
1028 .await
1029 .unwrap();
1030 let c = src
1031 .create_entity(&tok, "concept", None, "C", None, None, vec![])
1032 .await
1033 .unwrap();
1034
1035 let archive = KgArchive {
1036 format: "khive-kg".to_string(),
1037 version: "0.1".to_string(),
1038 namespace: "local".to_string(),
1039 exported_at: Utc::now(),
1040 entities: vec![
1041 ExportedEntity {
1042 id: a.id,
1043 kind: "concept".to_string(),
1044 entity_type: None,
1045 name: "A".to_string(),
1046 description: None,
1047 properties: None,
1048 tags: vec![],
1049 created_at: Utc::now(),
1050 updated_at: Utc::now(),
1051 },
1052 ExportedEntity {
1053 id: b.id,
1054 kind: "concept".to_string(),
1055 entity_type: None,
1056 name: "B".to_string(),
1057 description: None,
1058 properties: None,
1059 tags: vec![],
1060 created_at: Utc::now(),
1061 updated_at: Utc::now(),
1062 },
1063 ExportedEntity {
1064 id: c.id,
1065 kind: "concept".to_string(),
1066 entity_type: None,
1067 name: "C".to_string(),
1068 description: None,
1069 properties: None,
1070 tags: vec![],
1071 created_at: Utc::now(),
1072 updated_at: Utc::now(),
1073 },
1074 ],
1075 edges: vec![
1076 ExportedEdge {
1077 edge_id: Uuid::new_v4(),
1078 source: a.id,
1079 target: b.id,
1080 relation: EdgeRelation::Extends,
1081 weight: 1.0,
1082 properties: None,
1083 created_at: Utc::now(),
1084 updated_at: Utc::now(),
1085 },
1086 ExportedEdge {
1087 edge_id: Uuid::new_v4(),
1088 source: b.id,
1089 target: c.id,
1090 relation: EdgeRelation::VariantOf,
1095 weight: 0.9,
1096 properties: None,
1097 created_at: Utc::now(),
1098 updated_at: Utc::now(),
1099 },
1100 ExportedEdge {
1101 edge_id: Uuid::new_v4(),
1102 source: a.id,
1103 target: phantom,
1104 relation: EdgeRelation::Enables,
1105 weight: 0.5,
1106 properties: None,
1107 created_at: Utc::now(),
1108 updated_at: Utc::now(),
1109 },
1110 ],
1111 };
1112
1113 let dst = make_rt().await;
1114 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1115 assert_eq!(summary.entities_imported, 3);
1116 assert_eq!(
1117 summary.edges_imported, 2,
1118 "only valid edges must be imported"
1119 );
1120 assert_eq!(
1121 summary.edges_skipped, 1,
1122 "one dangling edge must be reported"
1123 );
1124 }
1125
1126 #[tokio::test]
1128 async fn import_all_valid_edges_reports_zero_skipped() {
1129 let src = make_rt().await;
1130 let tok = NamespaceToken::local();
1131 let e1 = src
1132 .create_entity(&tok, "concept", None, "E1", None, None, vec![])
1133 .await
1134 .unwrap();
1135 let e2 = src
1136 .create_entity(&tok, "concept", None, "E2", None, None, vec![])
1137 .await
1138 .unwrap();
1139 src.link(&tok, e1.id, e2.id, EdgeRelation::VariantOf, 0.7, None)
1140 .await
1141 .unwrap();
1142
1143 let archive = src.export_kg(&tok).await.unwrap();
1144 let dst = make_rt().await;
1145 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1146 assert_eq!(summary.edges_imported, 1);
1147 assert_eq!(
1148 summary.edges_skipped, 0,
1149 "no edges should be skipped when all endpoints exist"
1150 );
1151 }
1152
1153 #[tokio::test]
1159 async fn import_edge_violating_endpoint_contract_is_skipped() {
1160 let src = make_rt().await;
1161 let tok = NamespaceToken::local();
1162 let e1 = src
1163 .create_entity(&tok, "concept", None, "E1", None, None, vec![])
1164 .await
1165 .unwrap();
1166 let e2 = src
1167 .create_entity(&tok, "concept", None, "E2", None, None, vec![])
1168 .await
1169 .unwrap();
1170
1171 let archive = KgArchive {
1172 format: "khive-kg".to_string(),
1173 version: "0.1".to_string(),
1174 namespace: "local".to_string(),
1175 exported_at: Utc::now(),
1176 entities: vec![
1177 ExportedEntity {
1178 id: e1.id,
1179 kind: "concept".to_string(),
1180 entity_type: None,
1181 name: "E1".to_string(),
1182 description: None,
1183 properties: None,
1184 tags: vec![],
1185 created_at: Utc::now(),
1186 updated_at: Utc::now(),
1187 },
1188 ExportedEntity {
1189 id: e2.id,
1190 kind: "concept".to_string(),
1191 entity_type: None,
1192 name: "E2".to_string(),
1193 description: None,
1194 properties: None,
1195 tags: vec![],
1196 created_at: Utc::now(),
1197 updated_at: Utc::now(),
1198 },
1199 ],
1200 edges: vec![ExportedEdge {
1201 edge_id: Uuid::new_v4(),
1202 source: e1.id,
1203 target: e2.id,
1204 relation: EdgeRelation::Precedes,
1205 weight: 1.0,
1206 properties: None,
1207 created_at: Utc::now(),
1208 updated_at: Utc::now(),
1209 }],
1210 };
1211
1212 let dst = make_rt().await;
1213 let summary = dst.import_kg(&archive, &tok).await.unwrap();
1214 assert_eq!(summary.entities_imported, 2);
1215 assert_eq!(
1216 summary.edges_imported, 0,
1217 "an endpoint-contract-violating edge must not be imported"
1218 );
1219 assert_eq!(
1220 summary.edges_skipped, 1,
1221 "an endpoint-contract-violating edge must be counted as skipped"
1222 );
1223 assert!(
1224 dst.neighbors(
1225 &tok,
1226 e1.id,
1227 khive_storage::types::Direction::Out,
1228 None,
1229 None
1230 )
1231 .await
1232 .unwrap()
1233 .is_empty(),
1234 "the contract-violating edge must not exist in the destination graph"
1235 );
1236 }
1237
1238 #[tokio::test]
1242 async fn export_kg_preserves_edge_id() {
1243 let rt = make_rt().await;
1244 let tok = NamespaceToken::local();
1245 let a = rt
1246 .create_entity(&tok, "concept", None, "Alpha", None, None, vec![])
1247 .await
1248 .unwrap();
1249 let b = rt
1250 .create_entity(&tok, "concept", None, "Beta", None, None, vec![])
1251 .await
1252 .unwrap();
1253 let stored_edge = rt
1254 .link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
1255 .await
1256 .unwrap();
1257 let stored_id: Uuid = stored_edge.id.into();
1258
1259 let archive = rt.export_kg(&tok).await.unwrap();
1260 assert_eq!(archive.edges.len(), 1);
1261 assert_eq!(
1262 archive.edges[0].edge_id, stored_id,
1263 "exported edge_id must equal the LinkId returned by link"
1264 );
1265 }
1266
1267 #[tokio::test]
1269 async fn import_kg_persists_edge_id_and_timestamps() {
1270 let src = make_rt().await;
1271 let tok = NamespaceToken::local();
1272 let a = src
1273 .create_entity(&tok, "concept", None, "Alpha", None, None, vec![])
1274 .await
1275 .unwrap();
1276 let b = src
1277 .create_entity(&tok, "concept", None, "Beta", None, None, vec![])
1278 .await
1279 .unwrap();
1280 let stored_edge = src
1281 .link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
1282 .await
1283 .unwrap();
1284 let original_id: Uuid = stored_edge.id.into();
1285
1286 let expected_created = chrono::DateTime::parse_from_rfc3339("2026-03-03T00:00:00Z")
1287 .unwrap()
1288 .with_timezone(&Utc);
1289 let expected_updated = chrono::DateTime::parse_from_rfc3339("2026-04-04T00:00:00Z")
1290 .unwrap()
1291 .with_timezone(&Utc);
1292 let mut archive = src.export_kg(&tok).await.unwrap();
1293 archive.edges[0].created_at = expected_created;
1294 archive.edges[0].updated_at = expected_updated;
1295 archive.edges[0].properties = Some(serde_json::json!({"confidence": 0.95}));
1296 let dst = make_rt().await;
1297 dst.import_kg(&archive, &tok).await.unwrap();
1298
1299 let imported_edge = dst.get_edge(&tok, original_id).await.unwrap();
1300 assert!(
1301 imported_edge.is_some(),
1302 "imported edge must be retrievable by the original edge_id"
1303 );
1304 let imported_edge = imported_edge.unwrap();
1305 assert_eq!(
1306 Uuid::from(imported_edge.id),
1307 original_id,
1308 "stored edge id must equal the archive edge_id"
1309 );
1310 assert_eq!(
1311 imported_edge.created_at, expected_created,
1312 "present edge created_at must not be replaced with import time"
1313 );
1314 assert_eq!(
1315 imported_edge.updated_at, expected_updated,
1316 "present edge updated_at must not be replaced with created_at or import time"
1317 );
1318 assert_eq!(
1319 imported_edge.metadata,
1320 Some(serde_json::json!({"confidence": 0.95})),
1321 "edge properties must persist as storage metadata"
1322 );
1323
1324 let reexported = dst.export_kg(&tok).await.unwrap();
1325 assert_eq!(
1326 reexported.edges[0].properties,
1327 Some(serde_json::json!({"confidence": 0.95})),
1328 "edge metadata must re-export as portable properties"
1329 );
1330 }
1331
1332 #[tokio::test]
1337 async fn old_archive_missing_edge_id_round_trips() {
1338 let src_id = Uuid::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
1339 let tgt_id = Uuid::parse_str("00000000-0000-0000-0000-000000000002").unwrap();
1340
1341 let json = format!(
1343 r#"{{
1344 "format": "khive-kg",
1345 "version": "0.1",
1346 "namespace": "local",
1347 "exported_at": "2026-01-01T00:00:00Z",
1348 "entities": [
1349 {{"id":"{src_id}","kind":"concept","name":"SrcNode","created_at":"2026-01-01T00:00:00Z","updated_at":"2026-01-01T00:00:00Z"}},
1350 {{"id":"{tgt_id}","kind":"concept","name":"TgtNode","created_at":"2026-01-01T00:00:00Z","updated_at":"2026-01-01T00:00:00Z"}}
1351 ],
1352 "edges": [
1353 {{
1354 "source": "{src_id}",
1355 "target": "{tgt_id}",
1356 "relation": "extends",
1357 "weight": 0.9
1358 }}
1359 ]
1360 }}"#
1361 );
1362
1363 let archive: KgArchive = serde_json::from_str(&json)
1365 .expect("old archive without edge_id must deserialize successfully");
1366 assert_eq!(archive.edges.len(), 1);
1367 let generated_id = archive.edges[0].edge_id;
1368 assert_ne!(
1369 generated_id,
1370 Uuid::nil(),
1371 "missing edge_id in old archive must get a fresh non-nil UUID"
1372 );
1373
1374 let rt = make_rt().await;
1375 let tok = NamespaceToken::local();
1376 let summary = rt.import_kg(&archive, &tok).await.unwrap();
1377 assert_eq!(summary.entities_imported, 2);
1378 assert_eq!(
1379 summary.edges_imported, 1,
1380 "edge must be imported when both endpoints exist"
1381 );
1382
1383 let stored = rt.get_edge(&tok, generated_id).await.unwrap();
1384 assert!(
1385 stored.is_some(),
1386 "imported edge must be retrievable by the generated edge_id"
1387 );
1388 assert_eq!(
1389 Uuid::from(stored.unwrap().id),
1390 generated_id,
1391 "stored edge id must equal the generated edge_id"
1392 );
1393
1394 let re_archive = rt.export_kg(&tok).await.unwrap();
1395 assert_eq!(re_archive.edges.len(), 1);
1396 assert_eq!(
1397 re_archive.edges[0].edge_id, generated_id,
1398 "re-exported edge_id must equal the ID generated on first import"
1399 );
1400 }
1401
1402 #[tokio::test]
1407 async fn export_import_export_edge_id_equality() {
1408 let src = make_rt().await;
1409 let tok = NamespaceToken::local();
1410 let a = src
1411 .create_entity(&tok, "concept", None, "NodeA", None, None, vec![])
1412 .await
1413 .unwrap();
1414 let b = src
1415 .create_entity(&tok, "concept", None, "NodeB", None, None, vec![])
1416 .await
1417 .unwrap();
1418 let stored = src
1419 .link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
1420 .await
1421 .unwrap();
1422 let original_edge_id: Uuid = stored.id.into();
1423
1424 let archive1 = src.export_kg(&tok).await.unwrap();
1425 assert_eq!(archive1.edges.len(), 1);
1426 assert_eq!(
1427 archive1.edges[0].edge_id, original_edge_id,
1428 "first export must carry the stored edge_id"
1429 );
1430
1431 let dst = make_rt().await;
1432 dst.import_kg(&archive1, &tok).await.unwrap();
1433
1434 let archive2 = dst.export_kg(&tok).await.unwrap();
1435 assert_eq!(archive2.edges.len(), 1);
1436
1437 let re_edge = archive2
1438 .edges
1439 .iter()
1440 .find(|e| e.source == a.id && e.target == b.id && e.relation == EdgeRelation::Extends)
1441 .expect(
1442 "re-exported archive must contain the original edge by (source,target,relation)",
1443 );
1444 assert_eq!(
1445 re_edge.edge_id, original_edge_id,
1446 "edge_id must be identical across export → import → export"
1447 );
1448 }
1449}