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