Skip to main content

khive_runtime/
portability.rs

1//! KG export / import — portable JSON archive for namespace-scoped knowledge graphs.
2//!
3//! Embeddings are excluded (regenerable from text + model). Edges are collected by
4//! querying all entity IDs in the namespace first, then fetching incident edges.
5
6use 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// ── Archive types ─────────────────────────────────────────────────────────────
19
20/// Portable JSON archive of a namespace-scoped knowledge graph.
21///
22/// The `format` field is always `"khive-kg"`. The `version` field identifies
23/// the serialization schema; parsers should reject unknown versions.
24#[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/// An entity record in the portable archive.
35#[derive(Clone, Debug, Serialize, Deserialize)]
36pub struct ExportedEntity {
37    pub id: Uuid,
38    /// Pack-owned kind string (e.g. `"concept"`, `"person"`).
39    pub kind: String,
40    /// Pack-governed subtype token (e.g. `"paper"`, `"snapshot"`).
41    #[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/// A directed edge record in the portable archive.
55#[derive(Clone, Debug, Serialize, Deserialize)]
56pub struct ExportedEdge {
57    /// Stable edge identity across export/import cycles.
58    ///
59    /// Old archives (pre-0.2) omit this field. `serde(default)` assigns a fresh
60    /// UUID on import so backward-compatible archives are accepted as-is.
61    #[serde(default = "Uuid::new_v4")]
62    pub edge_id: Uuid,
63    pub source: Uuid,
64    pub target: Uuid,
65    /// One of the canonical edge relations (closed enum).
66    pub relation: EdgeRelation,
67    pub weight: f64,
68    /// Portable edge metadata, named `properties` in the ADR-020 wire shape.
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub properties: Option<serde_json::Value>,
71    /// Edge creation time. Older archives may omit it; those imports use the
72    /// time at which the archive is decoded.
73    #[serde(
74        default = "default_import_timestamp",
75        deserialize_with = "deserialize_import_timestamp"
76    )]
77    pub created_at: DateTime<Utc>,
78    /// Edge last-update time. Kept independent from `created_at` so imported
79    /// provenance is not collapsed during a round trip.
80    #[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/// Outcome of a successful import operation.
102#[derive(Clone, Debug, Serialize, Deserialize)]
103pub struct ImportSummary {
104    pub entities_imported: usize,
105    pub edges_imported: usize,
106    /// Number of edges that were skipped because one or both endpoint UUIDs
107    /// were not found in the target namespace after entity import.
108    ///
109    /// A non-zero value indicates the archive contained dangling edges (edges
110    /// referencing entities not present in the archive or the existing graph).
111    pub edges_skipped: usize,
112    /// Aggregate actual outcomes from reindexing imported entities. Source
113    /// text remains complete in the entity store and FTS when this is non-zero.
114    #[serde(default)]
115    pub embedding_truncation: crate::retrieval::EmbeddingTruncationReport,
116}
117
118// ── KhiveRuntime impl ─────────────────────────────────────────────────────────
119
120impl KhiveRuntime {
121    /// Export all entities and edges in a namespace to a portable JSON archive.
122    ///
123    /// Edge collection: all entity IDs in the namespace are gathered first;
124    /// `query_edges` is then called with those IDs as `source_ids`. This
125    /// captures every edge whose source entity belongs to the namespace.
126    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    /// Export to a JSON string (convenience wrapper around `export_kg`).
212    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    /// Import an archive into `target_namespace`.
218    ///
219    /// If `target_namespace` is `None`, the archive's own namespace is used.
220    ///
221    /// - Entities: upserted by ID; existing records are overwritten.
222    /// - Edges: upserted; existing records are overwritten.
223    /// - Validation: `format != "khive-kg"` or unsupported version → `InvalidInput`.
224    ///   Invalid edge relations are caught at JSON deserialization time.
225    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        // Complete deterministic validation before opening a store or issuing
246        // the first write. Endpoint existence and namespace checks remain at
247        // write time because they depend on mutable target state.
248        for (index, entity) in archive.entities.iter().enumerate() {
249            self.validate_entity_kind(&entity.kind)?;
250            // Archive content is caller-controlled input: the runtime-owned
251            // `khive:secret_gate` property key is reservation-only on import,
252            // exactly as on every other properties-bearing write path
253            // (ADR-115 Amendment 1 §3). Import never consumes exemptions.
254            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            // Edge properties are caller-controlled input: the runtime-owned
266            // `khive:secret_gate` key is reservation-only on import, and edge
267            // metadata is in ADR-115 Amendment 1 §3's unchanged blocking-scanner
268            // class, so credential-shaped values are rejected here as well.
269            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                // Publish each record through the same guarded FTS/vector writer
330                // as an ordinary reindex, including its best-effort model failures.
331                embedding_truncation.merge(
332                    self.reindex_entity_with_precomputed(token, &entity, outcomes)
333                        .await?,
334                );
335                entities_imported += 1;
336            }
337        }
338
339        // Untrusted archives may reference entities absent from the target namespace;
340        // check both endpoints and skip edges with a missing source or target to avoid
341        // dangling references in the graph store.
342        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            // Only a contract verdict (InvalidInput) or an endpoint-resolution
377            // miss (NotFound) is a property of the edge being imported; any
378            // other error is an operational failure of the check itself and
379            // must abort the import rather than silently discard a valid edge.
380            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    /// Import from a JSON string (convenience wrapper around `import_kg`).
424    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// ── Tests ─────────────────────────────────────────────────────────────────────
436
437// Kept inline: these tests exercise round-trip invariants over private encoding
438// helpers that would otherwise need to be made pub to test from tests/.
439#[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    /// 1. Roundtrip: 3 entities + 2 edges survive export → import on a fresh runtime.
824    #[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    /// 2. JSON roundtrip: export_kg_json → import_kg_json produces equivalent state.
888    #[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    /// 3. Namespace targeting: export from namespace "a", import into namespace "b" on a
925    ///    fresh runtime — entities land in "b", and the source runtime's "a" is unaffected.
926    ///
927    ///    Note: source and destination are separate runtimes (separate in-memory DBs).
928    ///    Same-DB cross-namespace copy is not a portability use case — portability is about
929    ///    moving graphs between instances, not between namespaces within one instance.
930    #[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    /// 4. Format validation: wrong `format` field → InvalidInput.
958    #[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    /// 5. Unsupported archive version → InvalidInput.
975    #[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    /// ADR-115 Amendment 1 §3: an archive entity carrying the reserved
1001    /// `khive:secret_gate` property key must be rejected before any upsert —
1002    /// import is a caller-controlled write path and stays reservation-only.
1003    #[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    /// ADR-115 Amendment 1 §3: an archive edge carrying the reserved
1042    /// `khive:secret_gate` property key is rejected in the same deterministic
1043    /// pre-write pass as entities — edge metadata is a properties-bearing
1044    /// write path.
1045    #[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        // The rejection happens before any write, so the archive's entities
1089        // must not have been created either.
1090        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    /// ADR-115 Amendment 1 §3: edge metadata is in the unchanged
1098    /// blocking-scanner class, so a credential-shaped value in edge
1099    /// properties fails the import in the same pre-write pass.
1100    #[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    /// 6. Invalid relation in archive → InvalidInput.
1151    #[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    // ── Dangling-edge validation tests ────────────────────────────────────────
1170
1171    /// 6. Edge with dangling source (source UUID not in entity table) is skipped.
1172    ///
1173    /// The archive has one entity + one edge whose source is a phantom UUID.
1174    /// Import succeeds, entities_imported=1, edges_imported=0, edges_skipped=1.
1175    #[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    /// 7. Edge with dangling target (target UUID not in entity table) is skipped.
1228    ///
1229    /// The archive has one entity + one edge whose target is a phantom UUID.
1230    /// Import succeeds, entities_imported=1, edges_imported=0, edges_skipped=1.
1231    #[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    /// 8. Mixed batch: some valid edges and some dangling edges — correct counts reported.
1284    ///
1285    /// Archive has 3 entities, 2 valid edges, and 1 dangling edge (phantom target).
1286    /// Import succeeds with edges_imported=2, edges_skipped=1.
1287    #[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                    // #1235: concept->concept is not in depends_on's endpoint
1362                    // contract; use variant_of (concept->concept is valid) so
1363                    // this fixture exercises "one dangling edge skipped" only,
1364                    // not an endpoint-contract violation this test isn't about.
1365                    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    /// 9. All-valid edges produce edges_skipped=0 (no regression on the happy path).
1398    #[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    /// #1235: an edge whose (source kind, relation, target kind) triple violates the
1425    /// endpoint contract (here: `precedes` between two `concept` entities — `precedes`
1426    /// is restricted to document/dataset/artifact/service/project pairs) must be
1427    /// skipped by import the same way a dangling endpoint is, not written straight
1428    /// through with `graph.upsert_edge`.
1429    #[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    // ── edge_id contract tests ────────────────────────────────────────────────
1510
1511    /// 10. export_kg sets edge_id in the archive to the LinkId returned by link.
1512    #[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    /// 11. import_kg writes the archive edge identity and timestamps exactly.
1539    #[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    /// 12. Old archive (no edge_id field) deserializes, imports, and re-exports with the
1604    ///     same generated UUID — proving the generated ID survives the full round trip.
1605    ///
1606    ///     The fixture includes two entities so the edge is not skipped during import.
1607    #[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        // Simulate a pre-0.2 archive JSON where the edge lacks an edge_id field.
1613        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        // serde(default) must assign a fresh non-nil UUID when edge_id is absent.
1635        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    /// 13. Explicit export → import → export equality: the edge_id is unchanged across
1674    ///     a full round trip when the source archive already contains an edge_id.
1675    ///
1676    ///     Verifies by (source, target, relation) key that re-export emits the original ID.
1677    #[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}