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 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
19/// Parse an optional archive timestamp, using `fallback` only when absent.
20///
21/// Present values must be RFC3339 and are converted to UTC. Invalid values keep
22/// their parser error with the original timestamp as context.
23pub 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// ── Archive types ─────────────────────────────────────────────────────────────
36
37/// Portable JSON archive of a namespace-scoped knowledge graph.
38///
39/// The `format` field is always `"khive-kg"`. The `version` field identifies
40/// the serialization schema; parsers should reject unknown versions.
41#[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/// An entity record in the portable archive.
52#[derive(Clone, Debug, Serialize, Deserialize)]
53pub struct ExportedEntity {
54    pub id: Uuid,
55    /// Pack-owned kind string (e.g. `"concept"`, `"person"`).
56    pub kind: String,
57    /// Pack-governed subtype token (e.g. `"paper"`, `"snapshot"`).
58    #[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/// A directed edge record in the portable archive.
72#[derive(Clone, Debug, Serialize, Deserialize)]
73pub struct ExportedEdge {
74    /// Stable edge identity across export/import cycles.
75    ///
76    /// Old archives (pre-0.2) omit this field. `serde(default)` assigns a fresh
77    /// UUID on import so backward-compatible archives are accepted as-is.
78    #[serde(default = "Uuid::new_v4")]
79    pub edge_id: Uuid,
80    pub source: Uuid,
81    pub target: Uuid,
82    /// One of the canonical edge relations (closed enum).
83    pub relation: EdgeRelation,
84    pub weight: f64,
85    /// Portable edge metadata, named `properties` in the ADR-020 wire shape.
86    #[serde(default, skip_serializing_if = "Option::is_none")]
87    pub properties: Option<serde_json::Value>,
88    /// Edge creation time. Older archives may omit it; those imports use the
89    /// time at which the archive is decoded.
90    #[serde(
91        default = "default_import_timestamp",
92        deserialize_with = "deserialize_import_timestamp"
93    )]
94    pub created_at: DateTime<Utc>,
95    /// Edge last-update time. Kept independent from `created_at` so imported
96    /// provenance is not collapsed during a round trip.
97    #[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/// Outcome of a successful import operation.
119#[derive(Clone, Debug, Serialize, Deserialize)]
120pub struct ImportSummary {
121    pub entities_imported: usize,
122    pub edges_imported: usize,
123    /// Number of edges that were skipped because one or both endpoint UUIDs
124    /// were not found in the target namespace after entity import.
125    ///
126    /// A non-zero value indicates the archive contained dangling edges (edges
127    /// referencing entities not present in the archive or the existing graph).
128    pub edges_skipped: usize,
129    /// Aggregate actual outcomes from reindexing imported entities. Source
130    /// text remains complete in the entity store and FTS when this is non-zero.
131    #[serde(default)]
132    pub embedding_truncation: crate::retrieval::EmbeddingTruncationReport,
133}
134
135// ── KhiveRuntime impl ─────────────────────────────────────────────────────────
136
137impl KhiveRuntime {
138    /// Export all entities and edges in a namespace to a portable JSON archive.
139    ///
140    /// Edge collection: all entity IDs in the namespace are gathered first;
141    /// `query_edges` is then called with those IDs as `source_ids`. This
142    /// captures every edge whose source entity belongs to the namespace.
143    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    /// Export to a JSON string (convenience wrapper around `export_kg`).
229    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    /// Import an archive into `target_namespace`.
235    ///
236    /// If `target_namespace` is `None`, the archive's own namespace is used.
237    ///
238    /// - Entities: upserted by ID; existing records are overwritten.
239    /// - Edges: upserted; existing records are overwritten.
240    /// - Validation: `format != "khive-kg"` or unsupported version → `InvalidInput`.
241    ///   Invalid edge relations are caught at JSON deserialization time.
242    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        // Complete deterministic validation before opening a store or issuing
263        // the first write. Endpoint existence and namespace checks remain at
264        // write time because they depend on mutable target state.
265        for (index, entity) in archive.entities.iter().enumerate() {
266            self.validate_entity_kind(&entity.kind)?;
267            // Archive content is caller-controlled input: the runtime-owned
268            // `khive:secret_gate` property key is reservation-only on import,
269            // exactly as on every other properties-bearing write path
270            // (ADR-115 Amendment 1 §3). Import never consumes exemptions.
271            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            // Edge properties are caller-controlled input: the runtime-owned
283            // `khive:secret_gate` key is reservation-only on import, and edge
284            // metadata is in ADR-115 Amendment 1 §3's unchanged blocking-scanner
285            // class, so credential-shaped values are rejected here as well.
286            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                // Publish each record through the same guarded FTS/vector writer
344                // as an ordinary reindex, including its best-effort model failures.
345                embedding_truncation.merge(
346                    self.reindex_entity_with_precomputed(token, &entity, outcomes)
347                        .await?,
348                );
349                entities_imported += 1;
350            }
351        }
352
353        // Untrusted archives may reference entities absent from the target namespace;
354        // check both endpoints and skip edges with a missing source or target to avoid
355        // dangling references in the graph store.
356        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            // Only a contract verdict (InvalidInput) or an endpoint-resolution
391            // miss (NotFound) is a property of the edge being imported; any
392            // other error is an operational failure of the check itself and
393            // must abort the import rather than silently discard a valid edge.
394            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    /// Import from a JSON string (convenience wrapper around `import_kg`).
438    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// ── Tests ─────────────────────────────────────────────────────────────────────
450
451// Kept inline: these tests exercise round-trip invariants over private encoding
452// helpers that would otherwise need to be made pub to test from tests/.
453#[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    /// 1. Roundtrip: 3 entities + 2 edges survive export → import on a fresh runtime.
861    #[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    /// 2. JSON roundtrip: export_kg_json → import_kg_json produces equivalent state.
925    #[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    /// 3. Namespace targeting: export from namespace "a", import into namespace "b" on a
962    ///    fresh runtime — entities land in "b", and the source runtime's "a" is unaffected.
963    ///
964    ///    Note: source and destination are separate runtimes (separate in-memory DBs).
965    ///    Same-DB cross-namespace copy is not a portability use case — portability is about
966    ///    moving graphs between instances, not between namespaces within one instance.
967    #[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    /// 4. Format validation: wrong `format` field → InvalidInput.
995    #[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    /// 5. Unsupported archive version → InvalidInput.
1012    #[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    /// ADR-115 Amendment 1 §3: an archive entity carrying the reserved
1038    /// `khive:secret_gate` property key must be rejected before any upsert —
1039    /// import is a caller-controlled write path and stays reservation-only.
1040    #[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    /// ADR-115 Amendment 1 §3: an archive edge carrying the reserved
1079    /// `khive:secret_gate` property key is rejected in the same deterministic
1080    /// pre-write pass as entities — edge metadata is a properties-bearing
1081    /// write path.
1082    #[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        // The rejection happens before any write, so the archive's entities
1126        // must not have been created either.
1127        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    /// ADR-115 Amendment 1 §3: edge metadata is in the unchanged
1135    /// blocking-scanner class, so a credential-shaped value in edge
1136    /// properties fails the import in the same pre-write pass.
1137    #[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    /// 6. Invalid relation in archive → InvalidInput.
1188    #[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    // ── Dangling-edge validation tests ────────────────────────────────────────
1207
1208    /// 6. Edge with dangling source (source UUID not in entity table) is skipped.
1209    ///
1210    /// The archive has one entity + one edge whose source is a phantom UUID.
1211    /// Import succeeds, entities_imported=1, edges_imported=0, edges_skipped=1.
1212    #[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    /// 7. Edge with dangling target (target UUID not in entity table) is skipped.
1265    ///
1266    /// The archive has one entity + one edge whose target is a phantom UUID.
1267    /// Import succeeds, entities_imported=1, edges_imported=0, edges_skipped=1.
1268    #[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    /// 8. Mixed batch: some valid edges and some dangling edges — correct counts reported.
1321    ///
1322    /// Archive has 3 entities, 2 valid edges, and 1 dangling edge (phantom target).
1323    /// Import succeeds with edges_imported=2, edges_skipped=1.
1324    #[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                    // #1235: concept->concept is not in depends_on's endpoint
1399                    // contract; use variant_of (concept->concept is valid) so
1400                    // this fixture exercises "one dangling edge skipped" only,
1401                    // not an endpoint-contract violation this test isn't about.
1402                    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    /// 9. All-valid edges produce edges_skipped=0 (no regression on the happy path).
1435    #[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    /// #1235: an edge whose (source kind, relation, target kind) triple violates the
1462    /// endpoint contract (here: `precedes` between two `concept` entities — `precedes`
1463    /// is restricted to document/dataset/artifact/service/project pairs) must be
1464    /// skipped by import the same way a dangling endpoint is, not written straight
1465    /// through with `graph.upsert_edge`.
1466    #[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    // ── edge_id contract tests ────────────────────────────────────────────────
1547
1548    /// 10. export_kg sets edge_id in the archive to the LinkId returned by link.
1549    #[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    /// 11. import_kg writes the archive edge identity and timestamps exactly.
1576    #[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    /// 12. Old archive (no edge_id field) deserializes, imports, and re-exports with the
1641    ///     same generated UUID — proving the generated ID survives the full round trip.
1642    ///
1643    ///     The fixture includes two entities so the edge is not skipped during import.
1644    #[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        // Simulate a pre-0.2 archive JSON where the edge lacks an edge_id field.
1650        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        // serde(default) must assign a fresh non-nil UUID when edge_id is absent.
1672        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    /// 13. Explicit export → import → export equality: the edge_id is unchanged across
1711    ///     a full round trip when the source archive already contains an edge_id.
1712    ///
1713    ///     Verifies by (source, target, relation) key that re-export emits the original ID.
1714    #[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}