1use std::collections::HashSet;
10
11use serde::{Deserialize, Serialize};
12use serde_json::Value;
13use uuid::Uuid;
14
15use khive_db::SqliteError;
16use khive_storage::note::Note;
17use khive_storage::types::{EdgeFilter, TextDocument};
18use khive_storage::{EdgeRelation, Entity, SubstrateKind};
19use khive_types::EventKind;
20use rusqlite::OptionalExtension;
21
22use crate::error::{RuntimeError, RuntimeResult};
23use crate::operations::canonical_edge_endpoints;
24use crate::runtime::{KhiveRuntime, NamespaceToken};
25
26#[derive(Clone, Debug, Default)]
48pub struct EntityPatch {
49 pub name: Option<String>,
50 pub description: Option<Option<String>>,
51 pub properties: Option<Value>,
52 pub tags: Option<Vec<String>>,
53}
54
55#[derive(Clone, Copy, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
57#[serde(rename_all = "snake_case")]
58pub enum EntityDedupMergePolicy {
59 #[default]
62 PreferInto,
63 PreferFrom,
65 Union,
67}
68
69#[derive(Clone, Copy, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
71#[serde(rename_all = "snake_case")]
72pub enum ContentMergeStrategy {
73 #[default]
74 Append,
75 PreferInto,
76 PreferFrom,
77}
78
79#[derive(Clone, Debug, Serialize, Deserialize)]
81pub struct MergeSummary {
82 pub kept_id: Uuid,
83 pub removed_id: Uuid,
84 pub edges_rewired: usize,
85 pub properties_merged: usize,
86 pub tags_unioned: usize,
87 pub content_appended: bool,
88 pub dry_run: bool,
89}
90
91#[derive(Clone, Debug, Default)]
96pub struct EdgePatch {
97 pub relation: Option<EdgeRelation>,
98 pub weight: Option<f64>,
99 pub properties: Option<Value>,
100}
101
102#[derive(Clone, Debug, Default)]
109pub struct NotePatch {
110 pub name: Option<Option<String>>,
111 pub content: Option<String>,
112 pub salience: Option<Option<f64>>,
113 pub decay_factor: Option<Option<f64>>,
114 pub properties: Option<Value>,
115 pub(crate) kind_status: Option<String>,
116}
117
118impl NotePatch {
119 pub fn new(
122 name: Option<Option<String>>,
123 content: Option<String>,
124 salience: Option<Option<f64>>,
125 decay_factor: Option<Option<f64>>,
126 properties: Option<Value>,
127 ) -> Self {
128 Self {
129 name,
130 content,
131 salience,
132 decay_factor,
133 properties,
134 kind_status: None,
135 }
136 }
137}
138
139#[derive(Clone, Debug, Default)]
141pub struct EdgeListFilter {
142 pub source_id: Option<Uuid>,
143 pub target_id: Option<Uuid>,
144 pub relations: Vec<EdgeRelation>,
146 pub min_weight: Option<f64>,
147 pub max_weight: Option<f64>,
148}
149
150impl From<EdgeListFilter> for EdgeFilter {
151 fn from(f: EdgeListFilter) -> Self {
152 EdgeFilter {
153 source_ids: f.source_id.into_iter().collect(),
154 target_ids: f.target_id.into_iter().collect(),
155 relations: f.relations,
156 min_weight: f.min_weight,
157 max_weight: f.max_weight,
158 ..Default::default()
159 }
160 }
161}
162
163#[allow(dead_code)]
171struct EdgeRow {
172 id: Uuid,
173 source_id: Uuid,
174 target_id: Uuid,
175 relation: String,
176 weight: f64,
177 created_at: i64,
178 updated_at: i64,
179 deleted_at: Option<i64>,
180 target_backend: Option<String>,
181 metadata: Option<String>,
182}
183
184impl KhiveRuntime {
189 pub(crate) async fn prepare_update_entity(
200 &self,
201 token: &NamespaceToken,
202 id: Uuid,
203 patch: EntityPatch,
204 ) -> RuntimeResult<(Entity, bool, Vec<&'static str>)> {
205 if let Some(ref name) = patch.name {
206 crate::secret_gate::check(name)?;
207 }
208 if let Some(Some(ref desc)) = patch.description {
209 crate::secret_gate::check(desc)?;
210 }
211 if let Some(ref props) = patch.properties {
212 crate::secret_gate::check_json(props)?;
213 }
214 if let Some(ref tags) = patch.tags {
215 crate::secret_gate::check_tags(tags)?;
216 }
217 let store = self.entities(token)?;
218 let mut entity = store
219 .get_entity(id)
220 .await?
221 .ok_or_else(|| RuntimeError::NotFound(format!("entity {id}")))?;
222
223 let mut text_changed = false;
224 let mut changed_fields: Vec<&'static str> = Vec::new();
225
226 if let Some(name) = patch.name {
227 text_changed |= entity.name != name;
228 entity.name = name;
229 changed_fields.push("name");
230 }
231 if let Some(desc_patch) = patch.description {
232 text_changed |= entity.description != desc_patch;
233 entity.description = desc_patch;
234 changed_fields.push("description");
235 }
236 if let Some(props) = patch.properties {
237 let (merged, _) = merge_properties(
238 &entity.properties,
239 &Some(props),
240 EntityDedupMergePolicy::PreferFrom,
241 );
242 entity.properties = merged;
243 changed_fields.push("properties");
244 }
245 if let Some(tags) = patch.tags {
246 entity.tags = tags;
247 changed_fields.push("tags");
248 }
249
250 entity.updated_at = chrono::Utc::now().timestamp_micros();
251 Ok((entity, text_changed, changed_fields))
252 }
253
254 pub async fn update_entity(
255 &self,
256 token: &NamespaceToken,
257 id: Uuid,
258 patch: EntityPatch,
259 ) -> RuntimeResult<Entity> {
260 let (entity, text_changed, changed_fields) =
261 self.prepare_update_entity(token, id, patch).await?;
262
263 let store = self.entities(token)?;
264 store.upsert_entity(entity.clone()).await?;
265
266 if text_changed {
267 self.reindex_entity(token, &entity).await?;
268 }
269
270 let event_store = self.events(token)?;
271 let event = khive_storage::event::Event::new(
272 entity.namespace.clone(),
273 "update",
274 EventKind::EntityUpdated,
275 SubstrateKind::Entity,
276 "",
277 )
278 .with_target(entity.id)
279 .with_payload(serde_json::json!({
280 "id": entity.id,
281 "namespace": entity.namespace,
282 "changed_fields": changed_fields,
283 }));
284 event_store.append_event(event).await.map_err(|e| {
285 RuntimeError::Internal(format!("update_entity: event store write failed: {e}"))
286 })?;
287
288 Ok(entity)
289 }
290
291 pub async fn merge_entity(
304 &self,
305 token: &NamespaceToken,
306 into_id: Uuid,
307 from_id: Uuid,
308 strategy: EntityDedupMergePolicy,
309 content_strategy: ContentMergeStrategy,
310 dry_run: bool,
311 ) -> RuntimeResult<MergeSummary> {
312 self.merge_entity_with_reason(
313 token,
314 into_id,
315 from_id,
316 strategy,
317 content_strategy,
318 dry_run,
319 None,
320 )
321 .await
322 }
323
324 #[allow(clippy::too_many_arguments)]
328 pub async fn merge_entity_with_reason(
329 &self,
330 token: &NamespaceToken,
331 into_id: Uuid,
332 from_id: Uuid,
333 strategy: EntityDedupMergePolicy,
334 content_strategy: ContentMergeStrategy,
335 dry_run: bool,
336 reason: Option<String>,
337 ) -> RuntimeResult<MergeSummary> {
338 if let Some(reason) = reason.as_deref() {
339 crate::secret_gate::check(reason)?;
340 }
341 if into_id == from_id {
342 return Err(RuntimeError::InvalidInput(
343 "cannot merge an entity into itself".into(),
344 ));
345 }
346 {
349 let into_entity = self.get_entity(token, into_id).await?;
350 let from_entity = self.get_entity(token, from_id).await?;
351 if into_entity.kind != from_entity.kind {
352 return Err(RuntimeError::InvalidInput(format!(
353 "cannot merge entities of different kinds: into={} ({}), from={} ({}); \
354 merge requires both entities to share the same kind",
355 into_id, into_entity.kind, from_id, from_entity.kind
356 )));
357 }
358 }
359 let ns = token.namespace().as_str().to_owned();
360 let fts_table = "fts_entities".to_string();
361 let vec_tables: Vec<String> = self
362 .registered_embedding_model_names()
363 .iter()
364 .map(|name| format!("vec_{}", crate::config::sanitize_key(name)))
365 .collect();
366
367 let _ = self.entities(token)?;
369 let _ = self.graph(token)?;
370 let _ = self.text(token)?;
371 for model_name in &self.registered_embedding_model_names() {
374 let _ = self.vectors_for_model(token, model_name)?;
375 }
376
377 let pool = self.backend().pool_arc();
378 let writer_task = pool.writer_task_handle().ok().flatten();
382
383 let (summary, updated_entity) = if let Some(writer_task) = writer_task {
384 writer_task
385 .send(move |conn| {
386 merge_entity_sql(
387 conn,
388 ns,
389 fts_table,
390 vec_tables,
391 into_id,
392 from_id,
393 strategy,
394 content_strategy,
395 dry_run,
396 )
397 .map_err(|e| {
398 khive_storage::StorageError::driver(
399 khive_storage::StorageCapability::Entities,
400 "merge_entity",
401 e,
402 )
403 })
404 })
405 .await
406 .map_err(RuntimeError::Storage)?
407 } else {
408 tokio::task::spawn_blocking(move || {
409 let guard = pool.writer()?;
410 guard.transaction(|conn| {
411 merge_entity_sql(
412 conn,
413 ns,
414 fts_table,
415 vec_tables,
416 into_id,
417 from_id,
418 strategy,
419 content_strategy,
420 dry_run,
421 )
422 })
423 })
424 .await
425 .map_err(|e| RuntimeError::Internal(e.to_string()))??
426 };
427
428 if !dry_run && !self.registered_embedding_model_names().is_empty() {
431 self.reindex_entity(token, &updated_entity).await?;
432 }
433
434 if !dry_run {
436 let event_store = self.events(token)?;
437 let policy_str = match strategy {
440 EntityDedupMergePolicy::PreferInto => "prefer_into",
441 EntityDedupMergePolicy::PreferFrom => "prefer_from",
442 EntityDedupMergePolicy::Union => "union",
443 };
444 let mut payload = serde_json::json!({
445 "into_id": summary.kept_id,
446 "from_id": summary.removed_id,
447 "policy": policy_str,
448 "content_strategy": format!("{:?}", content_strategy),
449 "edges_rewired": summary.edges_rewired,
450 });
451 if let Some(reason) = reason {
452 payload["reason"] = serde_json::Value::String(reason);
453 }
454 let event = khive_storage::event::Event::new(
455 updated_entity.namespace.clone(),
456 "merge",
457 EventKind::EntityMerged,
458 SubstrateKind::Entity,
459 "",
460 )
461 .with_target(summary.kept_id)
462 .with_payload(payload);
463 event_store.append_event(event).await.map_err(|e| {
464 RuntimeError::Internal(format!("merge_entity: event store write failed: {e}"))
465 })?;
466 }
467
468 Ok(summary)
469 }
470
471 pub(crate) async fn reindex_entity(
486 &self,
487 token: &NamespaceToken,
488 entity: &Entity,
489 ) -> RuntimeResult<()> {
490 let ns = entity.namespace.clone();
492 let doc = entity_fts_document(entity);
493 let embed_body = doc.body.clone();
494 self.text(token)?.upsert_document(doc).await?;
495
496 let embed_model_names = self.registered_embedding_model_names();
497 for model_name in &embed_model_names {
498 match self
499 .embed_document_with_model(model_name, &embed_body)
500 .await
501 {
502 Ok(vector) => match self.vectors_for_model(token, model_name) {
503 Ok(vs) => {
504 if let Err(e) = vs
505 .insert(
506 entity.id,
507 SubstrateKind::Entity,
508 &ns,
509 "entity.body",
510 vec![vector],
511 )
512 .await
513 {
514 tracing::warn!(
515 model = model_name,
516 id = %entity.id,
517 "reindex_entity: vector insert failed, skipping model: {e}"
518 );
519 }
520 }
521 Err(e) => {
522 tracing::warn!(
523 model = model_name,
524 id = %entity.id,
525 "reindex_entity: could not access vector store for model, skipping: {e}"
526 );
527 }
528 },
529 Err(e) => {
530 tracing::warn!(
531 model = model_name,
532 id = %entity.id,
533 "reindex_entity: embed failed for model, skipping: {e}"
534 );
535 }
536 }
537 }
538
539 Ok(())
540 }
541
542 pub(crate) async fn remove_from_indexes(
544 &self,
545 token: &NamespaceToken,
546 id: Uuid,
547 ) -> RuntimeResult<()> {
548 let ns = token.namespace().as_str().to_owned();
549 self.text(token)?.delete_document(&ns, id).await?;
550 for model_name in self.registered_embedding_model_names() {
551 self.vectors_for_model(token, &model_name)?
552 .delete(id)
553 .await?;
554 }
555 Ok(())
556 }
557
558 pub(crate) async fn reindex_note(
562 &self,
563 token: &NamespaceToken,
564 note: &khive_storage::note::Note,
565 ) -> RuntimeResult<()> {
566 self.text_for_notes(token)?
567 .upsert_document(note_fts_document(note))
568 .await?;
569
570 let ns = note.namespace.clone();
571 let embed_model_names = self.registered_embedding_model_names();
572 for model_name in &embed_model_names {
573 match self
574 .embed_document_with_model(model_name, ¬e.content)
575 .await
576 {
577 Ok(vector) => match self.vectors_for_model(token, model_name) {
578 Ok(vs) => {
579 if let Err(e) = vs
580 .insert(
581 note.id,
582 SubstrateKind::Note,
583 &ns,
584 "note.content",
585 vec![vector],
586 )
587 .await
588 {
589 tracing::warn!(
590 model = model_name,
591 id = %note.id,
592 "reindex_note: vector insert failed, skipping model: {e}"
593 );
594 }
595 }
596 Err(e) => {
597 tracing::warn!(
598 model = model_name,
599 id = %note.id,
600 "reindex_note: could not access vector store for model, skipping: {e}"
601 );
602 }
603 },
604 Err(e) => {
605 tracing::warn!(
606 model = model_name,
607 id = %note.id,
608 "reindex_note: embed failed for model, skipping: {e}"
609 );
610 }
611 }
612 }
613 Ok(())
614 }
615
616 pub(crate) async fn prepare_update_note(
620 &self,
621 token: &NamespaceToken,
622 id: Uuid,
623 patch: NotePatch,
624 ) -> RuntimeResult<(khive_storage::note::Note, bool)> {
625 if let Some(ref content) = patch.content {
626 crate::secret_gate::check(content)?;
627 }
628 if let Some(Some(ref name)) = patch.name {
629 crate::secret_gate::check(name)?;
630 }
631 if let Some(ref props) = patch.properties {
632 crate::secret_gate::check_json(props)?;
633 }
634 let store = self.notes(token)?;
635 let mut note = store
636 .get_note(id)
637 .await?
638 .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))?;
639
640 let mut text_changed = false;
641
642 if let Some(name_patch) = patch.name {
643 text_changed |= note.name != name_patch;
644 note.name = name_patch;
645 }
646 if let Some(content) = patch.content {
647 text_changed |= note.content != content;
648 note.content = content;
649 }
650 if let Some(salience_patch) = patch.salience {
651 if let Some(s) = salience_patch {
653 if !s.is_finite() || !(0.0..=1.0).contains(&s) {
654 return Err(crate::RuntimeError::InvalidInput(format!(
655 "salience must be a finite value in [0.0, 1.0]; got {s}"
656 )));
657 }
658 }
659 note.salience = salience_patch;
660 }
661 if let Some(decay_patch) = patch.decay_factor {
662 if let Some(d) = decay_patch {
664 if !d.is_finite() || d < 0.0 {
665 return Err(crate::RuntimeError::InvalidInput(format!(
666 "decay_factor must be a finite value >= 0.0; got {d}"
667 )));
668 }
669 }
670 note.decay_factor = decay_patch;
671 }
672 if let Some(props) = patch.properties {
673 let (merged, _) = merge_properties(
674 ¬e.properties,
675 &Some(props),
676 EntityDedupMergePolicy::PreferFrom,
677 );
678 note.properties = merged;
679 }
680 if let Some(status) = patch.kind_status {
681 note.status = status;
682 }
683
684 note.updated_at = chrono::Utc::now().timestamp_micros();
685 Ok((note, text_changed))
686 }
687
688 pub async fn update_note(
690 &self,
691 token: &NamespaceToken,
692 id: Uuid,
693 patch: NotePatch,
694 ) -> RuntimeResult<khive_storage::note::Note> {
695 let (note, text_changed) = self.prepare_update_note(token, id, patch).await?;
696
697 let store = self.notes(token)?;
698 store.upsert_note(note.clone()).await?;
699
700 if text_changed {
701 self.reindex_note(token, ¬e).await?;
702 self.fire_note_mutation_hook(¬e.kind, note.id).await;
706 }
707
708 Ok(note)
709 }
710
711 pub async fn merge_note(
720 &self,
721 token: &NamespaceToken,
722 into_id: Uuid,
723 from_id: Uuid,
724 strategy: EntityDedupMergePolicy,
725 content_strategy: ContentMergeStrategy,
726 dry_run: bool,
727 ) -> RuntimeResult<MergeSummary> {
728 self.merge_note_with_reason(
729 token,
730 into_id,
731 from_id,
732 strategy,
733 content_strategy,
734 dry_run,
735 None,
736 )
737 .await
738 }
739
740 #[allow(clippy::too_many_arguments)]
744 pub async fn merge_note_with_reason(
745 &self,
746 token: &NamespaceToken,
747 into_id: Uuid,
748 from_id: Uuid,
749 strategy: EntityDedupMergePolicy,
750 content_strategy: ContentMergeStrategy,
751 dry_run: bool,
752 reason: Option<String>,
753 ) -> RuntimeResult<MergeSummary> {
754 if let Some(reason) = reason.as_deref() {
755 crate::secret_gate::check(reason)?;
756 }
757 if into_id == from_id {
758 return Err(RuntimeError::InvalidInput(
759 "cannot merge a note into itself".into(),
760 ));
761 }
762 let ns = token.namespace().as_str().to_string();
763 let fts_table = "fts_notes".to_string();
764 let vec_tables: Vec<String> = self
765 .registered_embedding_model_names()
766 .iter()
767 .map(|name| format!("vec_{}", crate::config::sanitize_key(name)))
768 .collect();
769
770 let note_store = self.notes(token)?;
771 let into_note = note_store
772 .get_note(into_id)
773 .await?
774 .ok_or_else(|| RuntimeError::NotFound("not found in this namespace".into()))?;
775 Self::ensure_namespace(&into_note.namespace, &ns)?;
776
777 let from_note = note_store
778 .get_note(from_id)
779 .await?
780 .ok_or_else(|| RuntimeError::NotFound("not found in this namespace".into()))?;
781 Self::ensure_namespace(&from_note.namespace, &ns)?;
782
783 let _ = self.graph(token)?;
784 let _ = self.text_for_notes(token)?;
785 for model_name in &self.registered_embedding_model_names() {
786 let _ = self.vectors_for_model(token, model_name)?;
787 }
788
789 let pool = self.backend().pool_arc();
790 let writer_task = pool.writer_task_handle().ok().flatten();
791
792 let (summary, updated_note) = if let Some(writer_task) = writer_task {
793 writer_task
794 .send(move |conn| {
795 merge_note_sql(
796 conn,
797 ns,
798 fts_table,
799 vec_tables,
800 into_id,
801 from_id,
802 strategy,
803 content_strategy,
804 dry_run,
805 )
806 .map_err(|e| {
807 khive_storage::StorageError::driver(
808 khive_storage::StorageCapability::Notes,
809 "merge_note",
810 e,
811 )
812 })
813 })
814 .await
815 .map_err(RuntimeError::Storage)?
816 } else {
817 tokio::task::spawn_blocking(move || {
818 let guard = pool.writer()?;
819 guard.transaction(|conn| {
820 merge_note_sql(
821 conn,
822 ns,
823 fts_table,
824 vec_tables,
825 into_id,
826 from_id,
827 strategy,
828 content_strategy,
829 dry_run,
830 )
831 })
832 })
833 .await
834 .map_err(|e| RuntimeError::Internal(e.to_string()))??
835 };
836
837 if !dry_run && !self.registered_embedding_model_names().is_empty() {
838 self.reindex_note(token, &updated_note).await?;
839 self.fire_note_mutation_hook(&updated_note.kind, updated_note.id)
843 .await;
844 }
845
846 if !dry_run {
848 let event_store = self.events(token)?;
849 let policy_str = match strategy {
852 EntityDedupMergePolicy::PreferInto => "prefer_into",
853 EntityDedupMergePolicy::PreferFrom => "prefer_from",
854 EntityDedupMergePolicy::Union => "union",
855 };
856 let mut payload = serde_json::json!({
857 "into_id": summary.kept_id,
858 "from_id": summary.removed_id,
859 "policy": policy_str,
860 "content_strategy": format!("{:?}", content_strategy),
861 "edges_rewired": summary.edges_rewired,
862 });
863 if let Some(reason) = reason {
864 payload["reason"] = serde_json::Value::String(reason);
865 }
866 let event = khive_storage::event::Event::new(
867 updated_note.namespace.clone(),
868 "merge",
869 EventKind::NoteMerged,
870 SubstrateKind::Note,
871 "",
872 )
873 .with_target(summary.kept_id)
874 .with_payload(payload);
875 event_store.append_event(event).await.map_err(|e| {
876 RuntimeError::Internal(format!("merge_note: event store write failed: {e}"))
877 })?;
878 }
879
880 Ok(summary)
881 }
882}
883
884pub fn entity_fts_document(entity: &Entity) -> TextDocument {
901 let body = match &entity.description {
902 Some(d) if !d.is_empty() => format!("{} {}", entity.name, d),
903 _ => entity.name.clone(),
904 };
905 let updated_at =
906 chrono::DateTime::from_timestamp_micros(entity.updated_at).unwrap_or_else(chrono::Utc::now);
907 TextDocument {
908 subject_id: entity.id,
909 kind: SubstrateKind::Entity,
910 title: Some(entity.name.clone()),
911 body,
912 tags: entity.tags.clone(),
913 namespace: entity.namespace.clone(),
914 metadata: entity.properties.clone(),
915 updated_at,
916 }
917}
918
919pub fn note_fts_document(note: &Note) -> TextDocument {
932 let body = match ¬e.name {
933 Some(n) => format!("{n} {}", note.content),
934 None => note.content.clone(),
935 };
936 let updated_at =
937 chrono::DateTime::from_timestamp_micros(note.updated_at).unwrap_or_else(chrono::Utc::now);
938 TextDocument {
939 subject_id: note.id,
940 kind: SubstrateKind::Note,
941 title: note.name.clone(),
942 body,
943 tags: vec![],
944 namespace: note.namespace.clone(),
945 metadata: note.properties.clone(),
946 updated_at,
947 }
948}
949
950pub(crate) struct NoteFtsScalars {
956 pub title: String,
959 pub body: String,
960 pub tags: String,
962 pub metadata: Option<String>,
964 pub updated_at_micros: i64,
966}
967
968pub(crate) fn note_fts_scalars(note: &Note) -> NoteFtsScalars {
973 let doc = note_fts_document(note);
974 NoteFtsScalars {
975 title: doc.title.unwrap_or_default(),
976 body: doc.body,
977 tags: "[]".to_string(),
978 metadata: doc
979 .metadata
980 .as_ref()
981 .map(|v| serde_json::to_string(v).unwrap_or_default()),
982 updated_at_micros: doc.updated_at.timestamp_micros(),
983 }
984}
985
986fn read_merge_entity(
992 conn: &rusqlite::Connection,
993 id: Uuid,
994 namespace: &str,
995) -> Result<Entity, SqliteError> {
996 let id_str = id.to_string();
997 let mut stmt = conn.prepare(
998 "SELECT id, namespace, kind, entity_type, name, description, properties, tags, \
999 created_at, updated_at, deleted_at, merged_into, merge_event_id, content_ref \
1000 FROM entities WHERE id = ?1 AND deleted_at IS NULL",
1001 )?;
1002 let mut rows = stmt.query(rusqlite::params![id_str])?;
1003 let row = rows
1004 .next()?
1005 .ok_or_else(|| SqliteError::InvalidData(format!("entity {id} not found")))?;
1006
1007 let id_s: String = row.get(0)?;
1008 let ns: String = row.get(1)?;
1009 let kind: String = row.get(2)?;
1010 let entity_type: Option<String> = row.get(3)?;
1011 let name: String = row.get(4)?;
1012 let description: Option<String> = row.get(5)?;
1013 let properties_str: Option<String> = row.get(6)?;
1014 let tags_str: String = row.get(7)?;
1015 let created_at: i64 = row.get(8)?;
1016 let updated_at: i64 = row.get(9)?;
1017 let deleted_at: Option<i64> = row.get(10)?;
1018 let merged_into_str: Option<String> = row.get(11)?;
1019 let merge_event_id_str: Option<String> = row.get(12)?;
1020 let content_ref: Option<String> = row.get(13)?;
1021
1022 if ns != namespace {
1023 return Err(SqliteError::InvalidData(format!(
1024 "entity {id} belongs to namespace '{ns}', not '{namespace}'"
1025 )));
1026 }
1027
1028 let entity_id = Uuid::parse_str(&id_s).map_err(|e| SqliteError::InvalidData(e.to_string()))?;
1029 let properties: Option<Value> = properties_str
1030 .map(|s| {
1031 serde_json::from_str::<Value>(&s).map_err(|e| SqliteError::InvalidData(e.to_string()))
1032 })
1033 .transpose()?;
1034 let tags: Vec<String> =
1035 serde_json::from_str(&tags_str).map_err(|e| SqliteError::InvalidData(e.to_string()))?;
1036 let merged_into = merged_into_str
1037 .as_deref()
1038 .map(Uuid::parse_str)
1039 .transpose()
1040 .map_err(|e| SqliteError::InvalidData(e.to_string()))?;
1041 let merge_event_id = merge_event_id_str
1042 .as_deref()
1043 .map(Uuid::parse_str)
1044 .transpose()
1045 .map_err(|e| SqliteError::InvalidData(e.to_string()))?;
1046
1047 Ok(Entity {
1048 id: entity_id,
1049 namespace: ns,
1050 kind,
1051 entity_type,
1052 name,
1053 description,
1054 properties,
1055 tags,
1056 created_at,
1057 updated_at,
1058 deleted_at,
1059 merged_into,
1060 merge_event_id,
1061 content_ref,
1062 })
1063}
1064
1065#[allow(clippy::too_many_arguments)]
1076fn merge_entity_sql(
1077 conn: &rusqlite::Connection,
1078 namespace: String,
1079 fts_table: String,
1080 vec_tables: Vec<String>,
1081 into_id: Uuid,
1082 from_id: Uuid,
1083 strategy: EntityDedupMergePolicy,
1084 content_strategy: ContentMergeStrategy,
1085 dry_run: bool,
1086) -> Result<(MergeSummary, Entity), SqliteError> {
1087 let into_entity = read_merge_entity(conn, into_id, &namespace)?;
1088 let from_entity = read_merge_entity(conn, from_id, &namespace)?;
1089
1090 let parse_id =
1092 |s: String| Uuid::parse_str(&s).map_err(|e| SqliteError::InvalidData(e.to_string()));
1093
1094 let from_str = from_id.to_string();
1095
1096 let mut outbound: Vec<EdgeRow> = Vec::new();
1097 {
1098 let mut stmt = conn.prepare(
1099 "SELECT id, source_id, target_id, relation, weight, created_at, \
1100 updated_at, deleted_at, target_backend, metadata \
1101 FROM graph_edges WHERE namespace = ?1 AND source_id = ?2",
1102 )?;
1103 let mut rows = stmt.query(rusqlite::params![&namespace, &from_str])?;
1104 while let Some(row) = rows.next()? {
1105 outbound.push(EdgeRow {
1106 id: parse_id(row.get(0)?)?,
1107 source_id: parse_id(row.get(1)?)?,
1108 target_id: parse_id(row.get(2)?)?,
1109 relation: row.get(3)?,
1110 weight: row.get(4)?,
1111 created_at: row.get(5)?,
1112 updated_at: row.get(6)?,
1113 deleted_at: row.get(7)?,
1114 target_backend: row.get(8)?,
1115 metadata: row.get(9)?,
1116 });
1117 }
1118 }
1119
1120 let mut inbound: Vec<EdgeRow> = Vec::new();
1121 {
1122 let mut stmt = conn.prepare(
1123 "SELECT id, source_id, target_id, relation, weight, created_at, \
1124 updated_at, deleted_at, target_backend, metadata \
1125 FROM graph_edges WHERE namespace = ?1 AND target_id = ?2",
1126 )?;
1127 let mut rows = stmt.query(rusqlite::params![&namespace, &from_str])?;
1128 while let Some(row) = rows.next()? {
1129 inbound.push(EdgeRow {
1130 id: parse_id(row.get(0)?)?,
1131 source_id: parse_id(row.get(1)?)?,
1132 target_id: parse_id(row.get(2)?)?,
1133 relation: row.get(3)?,
1134 weight: row.get(4)?,
1135 created_at: row.get(5)?,
1136 updated_at: row.get(6)?,
1137 deleted_at: row.get(7)?,
1138 target_backend: row.get(8)?,
1139 metadata: row.get(9)?,
1140 });
1141 }
1142 }
1143
1144 let mut seen: HashSet<Uuid> = HashSet::new();
1146 let mut all_edges: Vec<EdgeRow> = Vec::new();
1147 for edge in outbound.into_iter().chain(inbound) {
1148 if seen.insert(edge.id) {
1149 all_edges.push(edge);
1150 }
1151 }
1152
1153 let (merged_props, properties_merged) =
1155 merge_properties(&into_entity.properties, &from_entity.properties, strategy);
1156 let merged_name = merge_string_field(&into_entity.name, &from_entity.name, strategy);
1157 let (merged_description, content_appended) = match content_strategy {
1158 ContentMergeStrategy::Append => {
1159 let into_desc = into_entity.description.as_deref().unwrap_or("");
1160 let from_desc = from_entity.description.as_deref().unwrap_or("");
1161 if from_desc.is_empty() {
1162 (into_entity.description.clone(), false)
1163 } else if into_desc.is_empty() {
1164 (from_entity.description.clone(), true)
1165 } else {
1166 (Some(format!("{}\n\n---\n\n{}", into_desc, from_desc)), true)
1167 }
1168 }
1169 ContentMergeStrategy::PreferInto => (into_entity.description.clone(), false),
1173 ContentMergeStrategy::PreferFrom => (from_entity.description.clone(), false),
1174 };
1175 let (merged_tags, tags_unioned) = union_tags(&into_entity.tags, &from_entity.tags);
1176
1177 let now = chrono::Utc::now().timestamp_micros();
1178 let into_str = into_id.to_string();
1179 let props_str = merged_props
1180 .as_ref()
1181 .map(|v| serde_json::to_string(v).unwrap_or_default());
1182 let tags_json = serde_json::to_string(&merged_tags).unwrap_or_else(|_| "[]".to_string());
1183
1184 let mut edges_rewired = 0usize;
1187 for edge in all_edges {
1188 let raw_src = if edge.source_id == from_id {
1189 into_id
1190 } else {
1191 edge.source_id
1192 };
1193 let raw_tgt = if edge.target_id == from_id {
1194 into_id
1195 } else {
1196 edge.target_id
1197 };
1198 let (new_src, new_tgt) = match edge.relation.parse::<EdgeRelation>() {
1201 Ok(rel) => canonical_edge_endpoints(rel, raw_src, raw_tgt),
1202 Err(_) => (raw_src, raw_tgt),
1203 };
1204
1205 if new_src == new_tgt {
1206 if !dry_run {
1207 conn.execute(
1208 "DELETE FROM graph_edges WHERE namespace = ?1 AND id = ?2",
1209 rusqlite::params![&namespace, edge.id.to_string()],
1210 )?;
1211 }
1212 continue;
1213 }
1214
1215 if dry_run {
1216 edges_rewired += 1;
1218 continue;
1219 }
1220
1221 let now_ts = chrono::Utc::now().timestamp();
1222 let conflict_id: Option<String> = {
1228 let conflict_src = new_src.to_string();
1229 let conflict_tgt = new_tgt.to_string();
1230 conn.query_row(
1231 khive_db::stores::graph::EDGE_SYMMETRIC_CONFLICT_PROBE_SQL,
1232 rusqlite::params![
1233 &namespace,
1234 &conflict_src,
1235 &conflict_tgt,
1236 &edge.relation,
1237 edge.id.to_string(),
1238 ],
1239 |row| row.get(0),
1240 )
1241 .optional()
1242 .map_err(SqliteError::Rusqlite)?
1243 };
1244
1245 let changed = if let Some(existing_id) = conflict_id {
1246 conn.execute(
1248 khive_db::stores::graph::EDGE_SYMMETRIC_DELETE_NONCANONICAL_SQL,
1249 rusqlite::params![&namespace, edge.id.to_string()],
1250 )?;
1251 conn.execute(
1252 khive_db::stores::graph::EDGE_SYMMETRIC_REFRESH_CANONICAL_SQL,
1253 rusqlite::params![
1254 edge.weight,
1255 now_ts,
1256 edge.target_backend,
1257 edge.metadata,
1258 &namespace,
1259 &existing_id,
1260 ],
1261 )?
1262 } else {
1263 conn.execute(
1264 "UPDATE graph_edges SET \
1265 source_id = ?1, target_id = ?2, updated_at = ?3 \
1266 WHERE namespace = ?4 AND id = ?5",
1267 rusqlite::params![
1268 new_src.to_string(),
1269 new_tgt.to_string(),
1270 now_ts,
1271 &namespace,
1272 edge.id.to_string(),
1273 ],
1274 )?
1275 };
1276 if changed > 0 {
1277 edges_rewired += 1;
1278 }
1279 }
1280
1281 if !dry_run {
1282 conn.execute(
1283 "INSERT OR REPLACE INTO entities \
1284 (id, namespace, kind, name, description, properties, tags, \
1285 created_at, updated_at, deleted_at, merged_into, merge_event_id) \
1286 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
1287 rusqlite::params![
1288 &into_str,
1289 &namespace,
1290 &into_entity.kind,
1291 &merged_name,
1292 &merged_description,
1293 &props_str,
1294 &tags_json,
1295 into_entity.created_at,
1296 now,
1297 into_entity.deleted_at,
1298 Option::<String>::None,
1299 Option::<String>::None,
1300 ],
1301 )?;
1302
1303 let fts_body = match &merged_description {
1307 Some(d) if !d.is_empty() => format!("{} {}", merged_name, d),
1308 _ => merged_name.clone(),
1309 };
1310 let kind_str = SubstrateKind::Entity.to_string();
1311
1312 conn.execute(
1313 &format!(
1314 "DELETE FROM {} WHERE namespace = ?1 AND subject_id = ?2",
1315 fts_table
1316 ),
1317 rusqlite::params![&namespace, &into_str],
1318 )?;
1319 conn.execute(
1320 &format!(
1321 "INSERT INTO {} \
1322 (subject_id, kind, title, body, tags, namespace, metadata, updated_at) \
1323 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
1324 fts_table
1325 ),
1326 rusqlite::params![
1327 &into_str,
1328 &kind_str,
1329 &merged_name,
1330 &fts_body,
1331 &tags_json,
1332 &namespace,
1333 &props_str,
1334 now,
1335 ],
1336 )?;
1337
1338 conn.execute(
1339 &format!(
1340 "DELETE FROM {} WHERE namespace = ?1 AND subject_id = ?2",
1341 fts_table
1342 ),
1343 rusqlite::params![&namespace, &from_str],
1344 )?;
1345
1346 khive_db::stores::vectors::delete_subject_from_vector_tables(
1347 conn,
1348 &vec_tables,
1349 from_id,
1350 &namespace,
1351 )?;
1352
1353 let merge_event_id = Uuid::new_v4();
1354 conn.execute(
1355 "UPDATE entities \
1356 SET deleted_at = ?1, merged_into = ?2, merge_event_id = ?3, updated_at = ?1 \
1357 WHERE namespace = ?4 AND id = ?5 AND deleted_at IS NULL",
1358 rusqlite::params![
1359 now,
1360 into_str,
1361 merge_event_id.to_string(),
1362 &namespace,
1363 &from_str,
1364 ],
1365 )?;
1366 }
1367
1368 let updated_entity = Entity {
1369 id: into_id,
1370 namespace,
1371 kind: into_entity.kind,
1372 entity_type: into_entity.entity_type,
1373 name: merged_name,
1374 description: merged_description,
1375 properties: merged_props,
1376 tags: merged_tags,
1377 created_at: into_entity.created_at,
1378 updated_at: now,
1379 deleted_at: into_entity.deleted_at,
1380 merged_into: None,
1381 merge_event_id: None,
1382 content_ref: into_entity.content_ref,
1383 };
1384
1385 Ok((
1386 MergeSummary {
1387 kept_id: into_id,
1388 removed_id: from_id,
1389 edges_rewired,
1390 properties_merged,
1391 tags_unioned,
1392 content_appended,
1393 dry_run,
1394 },
1395 updated_entity,
1396 ))
1397}
1398
1399fn read_merge_note(
1405 conn: &rusqlite::Connection,
1406 id: Uuid,
1407 namespace: &str,
1408) -> Result<khive_storage::note::Note, SqliteError> {
1409 use khive_storage::note::Note;
1410 let id_str = id.to_string();
1411 let mut stmt = conn.prepare(
1412 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
1413 expires_at, properties, created_at, updated_at, deleted_at \
1414 FROM notes WHERE id = ?1 AND deleted_at IS NULL",
1415 )?;
1416 let mut rows = stmt.query(rusqlite::params![id_str])?;
1417 let row = rows
1418 .next()?
1419 .ok_or_else(|| SqliteError::InvalidData(format!("note {id} not found")))?;
1420
1421 let id_s: String = row.get(0)?;
1422 let ns: String = row.get(1)?;
1423 let kind: String = row.get(2)?;
1424 let status: String = row.get(3)?;
1425 let name: Option<String> = row.get(4)?;
1426 let content: String = row.get(5)?;
1427 let salience: Option<f64> = row.get(6)?;
1428 let decay_factor: Option<f64> = row.get(7)?;
1429 let expires_at: Option<i64> = row.get(8)?;
1430 let properties_str: Option<String> = row.get(9)?;
1431 let created_at: i64 = row.get(10)?;
1432 let updated_at: i64 = row.get(11)?;
1433 let deleted_at: Option<i64> = row.get(12)?;
1434
1435 if ns != namespace {
1436 return Err(SqliteError::InvalidData(format!(
1437 "note {id} belongs to namespace '{ns}', not '{namespace}'"
1438 )));
1439 }
1440
1441 let note_id = Uuid::parse_str(&id_s).map_err(|e| SqliteError::InvalidData(e.to_string()))?;
1442 let properties: Option<serde_json::Value> = properties_str
1443 .map(|s| serde_json::from_str(&s).map_err(|e| SqliteError::InvalidData(e.to_string())))
1444 .transpose()?;
1445
1446 Ok(Note {
1447 id: note_id,
1448 namespace: ns,
1449 kind,
1450 status,
1451 name,
1452 content,
1453 salience,
1454 decay_factor,
1455 expires_at,
1456 properties,
1457 created_at,
1458 updated_at,
1459 deleted_at,
1460 })
1461}
1462
1463fn max_option_f64(a: Option<f64>, b: Option<f64>) -> Option<f64> {
1464 match (a, b) {
1465 (Some(x), Some(y)) => Some(x.max(y)),
1466 (Some(x), None) => Some(x),
1467 (None, Some(y)) => Some(y),
1468 (None, None) => None,
1469 }
1470}
1471
1472fn append_merge_history(props: Option<Value>, entry: Value) -> Result<Option<Value>, SqliteError> {
1473 use serde_json::{json, Map};
1474 let mut obj: Map<String, Value> = match props {
1475 Some(Value::Object(m)) => m,
1476 Some(other) => {
1477 let mut m = Map::new();
1478 m.insert("_value".into(), other);
1479 m
1480 }
1481 None => Map::new(),
1482 };
1483 let history = obj
1484 .entry("_merge_history".to_string())
1485 .or_insert_with(|| json!([]));
1486 if let Value::Array(arr) = history {
1487 arr.push(entry);
1488 }
1489 Ok(Some(Value::Object(obj)))
1490}
1491
1492#[allow(clippy::too_many_arguments)]
1502fn merge_note_sql(
1503 conn: &rusqlite::Connection,
1504 namespace: String,
1505 fts_table: String,
1506 vec_tables: Vec<String>,
1507 into_id: Uuid,
1508 from_id: Uuid,
1509 strategy: EntityDedupMergePolicy,
1510 content_strategy: ContentMergeStrategy,
1511 dry_run: bool,
1512) -> Result<(MergeSummary, khive_storage::note::Note), SqliteError> {
1513 let into_note = read_merge_note(conn, into_id, &namespace)?;
1514 let from_note = read_merge_note(conn, from_id, &namespace)?;
1515
1516 if into_note.kind != from_note.kind {
1517 return Err(SqliteError::InvalidData(format!(
1518 "cannot merge notes of different kinds: {} vs {}",
1519 into_note.kind, from_note.kind
1520 )));
1521 }
1522
1523 let now = chrono::Utc::now().timestamp_micros();
1524 let into_str = into_id.to_string();
1525 let from_str = from_id.to_string();
1526
1527 let parse_id =
1529 |s: String| Uuid::parse_str(&s).map_err(|e| SqliteError::InvalidData(e.to_string()));
1530
1531 let mut outbound: Vec<EdgeRow> = Vec::new();
1532 {
1533 let mut stmt = conn.prepare(
1534 "SELECT id, source_id, target_id, relation, weight, created_at, updated_at, deleted_at, target_backend, metadata \
1535 FROM graph_edges WHERE namespace = ?1 AND source_id = ?2",
1536 )?;
1537 let mut rows = stmt.query(rusqlite::params![&namespace, &from_str])?;
1538 while let Some(row) = rows.next()? {
1539 outbound.push(EdgeRow {
1540 id: parse_id(row.get(0)?)?,
1541 source_id: parse_id(row.get(1)?)?,
1542 target_id: parse_id(row.get(2)?)?,
1543 relation: row.get(3)?,
1544 weight: row.get(4)?,
1545 created_at: row.get(5)?,
1546 updated_at: row.get(6)?,
1547 deleted_at: row.get(7)?,
1548 target_backend: row.get(8)?,
1549 metadata: row.get(9)?,
1550 });
1551 }
1552 }
1553 let mut inbound: Vec<EdgeRow> = Vec::new();
1554 {
1555 let mut stmt = conn.prepare(
1556 "SELECT id, source_id, target_id, relation, weight, created_at, updated_at, deleted_at, target_backend, metadata \
1557 FROM graph_edges WHERE namespace = ?1 AND target_id = ?2",
1558 )?;
1559 let mut rows = stmt.query(rusqlite::params![&namespace, &from_str])?;
1560 while let Some(row) = rows.next()? {
1561 inbound.push(EdgeRow {
1562 id: parse_id(row.get(0)?)?,
1563 source_id: parse_id(row.get(1)?)?,
1564 target_id: parse_id(row.get(2)?)?,
1565 relation: row.get(3)?,
1566 weight: row.get(4)?,
1567 created_at: row.get(5)?,
1568 updated_at: row.get(6)?,
1569 deleted_at: row.get(7)?,
1570 target_backend: row.get(8)?,
1571 metadata: row.get(9)?,
1572 });
1573 }
1574 }
1575 let mut seen: HashSet<Uuid> = HashSet::new();
1576 let mut all_edges: Vec<EdgeRow> = Vec::new();
1577 for edge in outbound.into_iter().chain(inbound) {
1578 if seen.insert(edge.id) {
1579 all_edges.push(edge);
1580 }
1581 }
1582
1583 let (merged_content, content_appended) = match content_strategy {
1585 ContentMergeStrategy::Append => {
1586 if from_note.content.is_empty() {
1587 (into_note.content.clone(), false)
1588 } else {
1589 (
1590 format!("{}\n\n---\n\n{}", into_note.content, from_note.content),
1591 true,
1592 )
1593 }
1594 }
1595 ContentMergeStrategy::PreferInto => (into_note.content.clone(), false),
1596 ContentMergeStrategy::PreferFrom => (from_note.content.clone(), false),
1597 };
1598
1599 let merged_name = match strategy {
1600 EntityDedupMergePolicy::PreferFrom => from_note.name.clone().or(into_note.name.clone()),
1601 _ => into_note.name.clone().or(from_note.name.clone()),
1602 };
1603
1604 let (merged_props, properties_merged) =
1605 merge_properties(&into_note.properties, &from_note.properties, strategy);
1606
1607 let merge_history_entry = serde_json::json!({
1608 "merged_from": from_id.to_string(),
1609 "merged_at": now,
1610 "strategy": format!("{:?}", strategy),
1611 "content_strategy": format!("{:?}", content_strategy),
1612 });
1613 let merged_props = append_merge_history(merged_props, merge_history_entry)?;
1614
1615 let merged_salience = max_option_f64(into_note.salience, from_note.salience);
1616 let merged_expires_at = match (into_note.expires_at, from_note.expires_at) {
1617 (Some(a), Some(b)) => Some(a.max(b)),
1618 (Some(a), None) => Some(a),
1619 (None, Some(b)) => Some(b),
1620 (None, None) => None,
1621 };
1622
1623 let props_str = merged_props
1624 .as_ref()
1625 .map(|v| serde_json::to_string(v).unwrap_or_default());
1626
1627 let mut edges_rewired = 0usize;
1628 if !dry_run {
1629 for edge in all_edges {
1630 let raw_src = if edge.source_id == from_id {
1631 into_id
1632 } else {
1633 edge.source_id
1634 };
1635 let raw_tgt = if edge.target_id == from_id {
1636 into_id
1637 } else {
1638 edge.target_id
1639 };
1640 let (new_src, new_tgt) = match edge.relation.parse::<EdgeRelation>() {
1642 Ok(rel) => canonical_edge_endpoints(rel, raw_src, raw_tgt),
1643 Err(_) => (raw_src, raw_tgt),
1644 };
1645 if new_src == new_tgt {
1646 conn.execute(
1647 "DELETE FROM graph_edges WHERE namespace = ?1 AND id = ?2",
1648 rusqlite::params![&namespace, edge.id.to_string()],
1649 )?;
1650 continue;
1651 }
1652 let now_ts = chrono::Utc::now().timestamp();
1653 let conflict_id: Option<String> = {
1654 let conflict_src = new_src.to_string();
1655 let conflict_tgt = new_tgt.to_string();
1656 conn.query_row(
1657 khive_db::stores::graph::EDGE_SYMMETRIC_CONFLICT_PROBE_SQL,
1658 rusqlite::params![
1659 &namespace,
1660 &conflict_src,
1661 &conflict_tgt,
1662 &edge.relation,
1663 edge.id.to_string(),
1664 ],
1665 |row| row.get(0),
1666 )
1667 .optional()
1668 .map_err(SqliteError::Rusqlite)?
1669 };
1670
1671 let changed = if let Some(existing_id) = conflict_id {
1672 conn.execute(
1673 khive_db::stores::graph::EDGE_SYMMETRIC_DELETE_NONCANONICAL_SQL,
1674 rusqlite::params![&namespace, edge.id.to_string()],
1675 )?;
1676 conn.execute(
1677 khive_db::stores::graph::EDGE_SYMMETRIC_REFRESH_CANONICAL_SQL,
1678 rusqlite::params![
1679 edge.weight,
1680 now_ts,
1681 edge.target_backend,
1682 edge.metadata,
1683 &namespace,
1684 &existing_id,
1685 ],
1686 )?
1687 } else {
1688 conn.execute(
1689 "UPDATE graph_edges SET \
1690 source_id = ?1, target_id = ?2, updated_at = ?3 \
1691 WHERE namespace = ?4 AND id = ?5",
1692 rusqlite::params![
1693 new_src.to_string(),
1694 new_tgt.to_string(),
1695 now_ts,
1696 &namespace,
1697 edge.id.to_string(),
1698 ],
1699 )?
1700 };
1701 if changed > 0 {
1702 edges_rewired += 1;
1703 }
1704 }
1705
1706 conn.prepare_cached(khive_db::stores::note::NOTE_UPSERT_SQL)?
1707 .execute(rusqlite::params![
1708 &into_str,
1709 &namespace,
1710 &into_note.kind,
1711 &into_note.status,
1712 &merged_name,
1713 &merged_content,
1714 merged_salience,
1715 into_note.decay_factor,
1716 merged_expires_at,
1717 &props_str,
1718 into_note.created_at,
1719 now,
1720 into_note.deleted_at,
1721 ])?;
1722
1723 conn.execute(
1724 &format!(
1725 "DELETE FROM {} WHERE namespace = ?1 AND subject_id = ?2",
1726 fts_table
1727 ),
1728 rusqlite::params![&namespace, &into_str],
1729 )?;
1730 let fts_merged = {
1735 let mut merged_note = Note::new(&namespace, &*into_note.kind, &*merged_content);
1736 merged_note.id = into_id;
1737 merged_note.name = merged_name.clone();
1738 merged_note.properties = merged_props.clone();
1739 merged_note.updated_at = now;
1740 note_fts_scalars(&merged_note)
1741 };
1742 conn.execute(
1743 &format!(
1744 "INSERT INTO {} \
1745 (subject_id, kind, title, body, tags, namespace, metadata, updated_at) \
1746 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
1747 fts_table
1748 ),
1749 rusqlite::params![
1750 &into_str,
1751 SubstrateKind::Note.to_string(),
1752 &fts_merged.title,
1753 &fts_merged.body,
1754 &fts_merged.tags,
1755 &namespace,
1756 &fts_merged.metadata,
1757 fts_merged.updated_at_micros,
1758 ],
1759 )?;
1760
1761 conn.execute(
1762 &format!(
1763 "DELETE FROM {} WHERE namespace = ?1 AND subject_id = ?2",
1764 fts_table
1765 ),
1766 rusqlite::params![&namespace, &from_str],
1767 )?;
1768
1769 khive_db::stores::vectors::delete_subject_from_vector_tables(
1770 conn,
1771 &vec_tables,
1772 from_id,
1773 &namespace,
1774 )?;
1775
1776 conn.execute(
1777 "UPDATE notes SET status = 'deleted', deleted_at = ?1, updated_at = ?1 \
1778 WHERE namespace = ?2 AND id = ?3 AND deleted_at IS NULL",
1779 rusqlite::params![now, &namespace, &from_str],
1780 )?;
1781 }
1782
1783 let updated_note = khive_storage::note::Note {
1784 id: into_id,
1785 namespace: namespace.clone(),
1786 kind: into_note.kind.clone(),
1787 status: into_note.status.clone(),
1788 name: merged_name,
1789 content: merged_content,
1790 salience: merged_salience,
1791 decay_factor: into_note.decay_factor,
1792 expires_at: merged_expires_at,
1793 properties: merged_props,
1794 created_at: into_note.created_at,
1795 updated_at: now,
1796 deleted_at: into_note.deleted_at,
1797 };
1798
1799 Ok((
1800 MergeSummary {
1801 kept_id: into_id,
1802 removed_id: from_id,
1803 edges_rewired,
1804 properties_merged,
1805 tags_unioned: 0,
1806 content_appended,
1807 dry_run,
1808 },
1809 updated_note,
1810 ))
1811}
1812
1813pub(crate) fn merge_string_field(
1820 into: &str,
1821 from: &str,
1822 strategy: EntityDedupMergePolicy,
1823) -> String {
1824 match strategy {
1825 EntityDedupMergePolicy::PreferInto | EntityDedupMergePolicy::Union => into.to_string(),
1826 EntityDedupMergePolicy::PreferFrom => from.to_string(),
1827 }
1828}
1829
1830pub(crate) fn merge_properties(
1835 into: &Option<Value>,
1836 from: &Option<Value>,
1837 strategy: EntityDedupMergePolicy,
1838) -> (Option<Value>, usize) {
1839 match (into, from) {
1840 (None, None) => (None, 0),
1841 (Some(a), None) => (Some(a.clone()), 0),
1842 (None, Some(b)) => {
1843 let count = if let Value::Object(m) = b { m.len() } else { 1 };
1844 (Some(b.clone()), count)
1845 }
1846 (Some(into_val), Some(from_val)) => {
1847 let (merged, added) = merge_json(into_val, from_val, strategy);
1848 (Some(merged), added)
1849 }
1850 }
1851}
1852
1853fn merge_json(into: &Value, from: &Value, strategy: EntityDedupMergePolicy) -> (Value, usize) {
1855 match (into, from, strategy) {
1856 (Value::Object(a), Value::Object(b), EntityDedupMergePolicy::Union) => {
1857 let mut result = a.clone();
1858 let mut added = 0usize;
1859 for (k, v_from) in b {
1860 if let Some(v_into) = a.get(k) {
1861 let (merged, sub_added) =
1862 merge_json(v_into, v_from, EntityDedupMergePolicy::Union);
1863 result.insert(k.clone(), merged);
1864 added += sub_added;
1865 } else {
1866 result.insert(k.clone(), v_from.clone());
1867 added += 1;
1868 }
1869 }
1870 (Value::Object(result), added)
1871 }
1872 (Value::Object(a), Value::Object(b), EntityDedupMergePolicy::PreferInto) => {
1873 let mut result = a.clone();
1874 let mut added = 0usize;
1875 for (k, v) in b {
1876 if !a.contains_key(k) {
1877 result.insert(k.clone(), v.clone());
1878 added += 1;
1879 }
1880 }
1881 (Value::Object(result), added)
1882 }
1883 (Value::Object(a), Value::Object(b), EntityDedupMergePolicy::PreferFrom) => {
1884 let mut result = a.clone();
1885 let mut added = 0usize;
1886 for (k, v) in b {
1887 result.insert(k.clone(), v.clone());
1888 if !a.contains_key(k) {
1889 added += 1;
1890 }
1891 }
1892 (Value::Object(result), added)
1893 }
1894 (_into_val, from_val, EntityDedupMergePolicy::PreferFrom) => (from_val.clone(), 1),
1896 _ => (into.clone(), 0),
1897 }
1898}
1899
1900pub(crate) fn union_tags(into: &[String], from: &[String]) -> (Vec<String>, usize) {
1903 let mut seen: HashSet<&str> = into.iter().map(|s| s.as_str()).collect();
1904 let mut result: Vec<String> = into.to_vec();
1905 let mut added = 0usize;
1906 for tag in from {
1907 if seen.insert(tag.as_str()) {
1908 result.push(tag.clone());
1909 added += 1;
1910 }
1911 }
1912 (result, added)
1913}
1914
1915#[cfg(test)]
1924mod tests {
1925 use super::*;
1926 use crate::runtime::{KhiveRuntime, NamespaceToken};
1927 use khive_storage::types::{Direction, TextFilter, TextQueryMode, TextSearchRequest};
1928
1929 fn rt() -> KhiveRuntime {
1930 KhiveRuntime::memory().unwrap()
1931 }
1932
1933 fn secret_shaped_reason() -> String {
1934 const ALPHANUMERIC: &[u8] =
1935 b"0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz";
1936 let candidate: String = (0..48)
1937 .map(|index| char::from(ALPHANUMERIC[(index * 17 + 11) % ALPHANUMERIC.len()]))
1938 .collect();
1939 format!("secret value: {candidate}")
1940 }
1941
1942 async fn fts_hit(rt: &KhiveRuntime, token: &NamespaceToken, query: &str) -> Vec<Uuid> {
1944 let ns = token.namespace().as_str().to_string();
1945 rt.text(token)
1946 .unwrap()
1947 .search(TextSearchRequest {
1948 query: query.to_string(),
1949 mode: TextQueryMode::Plain,
1950 filter: Some(TextFilter {
1951 namespaces: vec![ns],
1952 ..Default::default()
1953 }),
1954 top_k: 50,
1955 snippet_chars: 100,
1956 })
1957 .await
1958 .unwrap()
1959 .into_iter()
1960 .map(|h| h.subject_id)
1961 .collect()
1962 }
1963
1964 #[tokio::test]
1965 async fn update_entity_patch_changes_only_specified_fields() {
1966 let rt = rt();
1967 let tok = NamespaceToken::local();
1968 let entity = rt
1969 .create_entity(
1970 &tok,
1971 "concept",
1972 None,
1973 "OriginalName",
1974 Some("orig desc"),
1975 Some(serde_json::json!({"k":"v"})),
1976 vec![],
1977 )
1978 .await
1979 .unwrap();
1980
1981 let updated = rt
1982 .update_entity(
1983 &tok,
1984 entity.id,
1985 EntityPatch {
1986 description: Some(Some("new desc".to_string())),
1987 ..Default::default()
1988 },
1989 )
1990 .await
1991 .unwrap();
1992
1993 assert_eq!(updated.name, "OriginalName");
1994 assert_eq!(updated.description.as_deref(), Some("new desc"));
1995 assert_eq!(updated.properties, Some(serde_json::json!({"k":"v"})));
1996 }
1997
1998 #[tokio::test]
1999 async fn update_entity_clear_description_with_some_none() {
2000 let rt = rt();
2001 let tok = NamespaceToken::local();
2002 let entity = rt
2003 .create_entity(
2004 &tok,
2005 "concept",
2006 None,
2007 "ClearDesc",
2008 Some("has description"),
2009 None,
2010 vec![],
2011 )
2012 .await
2013 .unwrap();
2014
2015 let updated = rt
2016 .update_entity(
2017 &tok,
2018 entity.id,
2019 EntityPatch {
2020 description: Some(None),
2021 ..Default::default()
2022 },
2023 )
2024 .await
2025 .unwrap();
2026
2027 assert!(
2028 updated.description.is_none(),
2029 "description should be cleared"
2030 );
2031 }
2032
2033 #[tokio::test]
2034 async fn update_entity_reindexes_when_name_changes() {
2035 let rt = rt();
2036 let tok = NamespaceToken::local();
2037 let entity = rt
2038 .create_entity(&tok, "concept", None, "OldName", None, None, vec![])
2039 .await
2040 .unwrap();
2041
2042 let hits_before = fts_hit(&rt, &tok, "OldName").await;
2043 assert!(
2044 hits_before.contains(&entity.id),
2045 "entity should be findable by old name"
2046 );
2047
2048 rt.update_entity(
2049 &tok,
2050 entity.id,
2051 EntityPatch {
2052 name: Some("NewName".to_string()),
2053 ..Default::default()
2054 },
2055 )
2056 .await
2057 .unwrap();
2058
2059 let hits_old = fts_hit(&rt, &tok, "OldName").await;
2060 let hits_new = fts_hit(&rt, &tok, "NewName").await;
2061
2062 assert!(
2063 !hits_old.contains(&entity.id),
2064 "old name should no longer match after rename"
2065 );
2066 assert!(
2067 hits_new.contains(&entity.id),
2068 "new name should be findable after rename"
2069 );
2070 }
2071
2072 #[tokio::test]
2073 async fn update_entity_properties_merges_preserving_existing_keys() {
2074 let rt = rt();
2075 let tok = NamespaceToken::local();
2076 let entity = rt
2077 .create_entity(
2078 &tok,
2079 "concept",
2080 None,
2081 "MergeProps",
2082 None,
2083 Some(serde_json::json!({
2084 "domain": "inference",
2085 "repo": "lattice",
2086 "status": "researched",
2087 })),
2088 vec![],
2089 )
2090 .await
2091 .unwrap();
2092
2093 let updated = rt
2094 .update_entity(
2095 &tok,
2096 entity.id,
2097 EntityPatch {
2098 properties: Some(serde_json::json!({"status": "implemented"})),
2099 ..Default::default()
2100 },
2101 )
2102 .await
2103 .unwrap();
2104
2105 let props = updated.properties.expect("properties should remain set");
2106 assert_eq!(props["domain"], "inference", "domain key must be preserved");
2107 assert_eq!(props["repo"], "lattice", "repo key must be preserved");
2108 assert_eq!(
2109 props["status"], "implemented",
2110 "status key must be updated by patch"
2111 );
2112 }
2113
2114 #[tokio::test]
2115 async fn update_entity_skips_reindex_when_only_properties_change() {
2116 let rt = rt();
2117 let tok = NamespaceToken::local();
2118 let entity = rt
2119 .create_entity(&tok, "concept", None, "StableIndexed", None, None, vec![])
2120 .await
2121 .unwrap();
2122
2123 let hits_before = fts_hit(&rt, &tok, "StableIndexed").await;
2124 assert!(hits_before.contains(&entity.id));
2125
2126 rt.update_entity(
2127 &tok,
2128 entity.id,
2129 EntityPatch {
2130 properties: Some(serde_json::json!({"new": "prop"})),
2131 ..Default::default()
2132 },
2133 )
2134 .await
2135 .unwrap();
2136
2137 let hits_after = fts_hit(&rt, &tok, "StableIndexed").await;
2138 assert!(
2139 hits_after.contains(&entity.id),
2140 "still findable after props-only patch"
2141 );
2142 }
2143
2144 #[tokio::test]
2145 async fn merge_entity_rewires_edges() {
2146 let rt = rt();
2147 let tok = NamespaceToken::local();
2148 let a = rt
2149 .create_entity(&tok, "concept", None, "A", None, None, vec![])
2150 .await
2151 .unwrap();
2152 let b = rt
2153 .create_entity(&tok, "concept", None, "B", None, None, vec![])
2154 .await
2155 .unwrap();
2156 let c = rt
2157 .create_entity(&tok, "concept", None, "C", None, None, vec![])
2158 .await
2159 .unwrap();
2160 let d = rt
2161 .create_entity(&tok, "concept", None, "D", None, None, vec![])
2162 .await
2163 .unwrap();
2164
2165 rt.link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
2167 .await
2168 .unwrap();
2169 rt.link(&tok, c.id, b.id, EdgeRelation::Extends, 1.0, None)
2170 .await
2171 .unwrap();
2172
2173 let summary = rt
2174 .merge_entity_with_reason(
2175 &tok,
2176 d.id,
2177 b.id,
2178 EntityDedupMergePolicy::PreferInto,
2179 ContentMergeStrategy::Append,
2180 false,
2181 None,
2182 )
2183 .await
2184 .unwrap();
2185
2186 assert_eq!(summary.kept_id, d.id);
2187 assert_eq!(summary.removed_id, b.id);
2188 assert_eq!(summary.edges_rewired, 2);
2189
2190 let a_neighbors = rt
2191 .neighbors(&tok, a.id, Direction::Out, None, None)
2192 .await
2193 .unwrap();
2194 assert_eq!(a_neighbors.len(), 1);
2195 assert_eq!(a_neighbors[0].node_id, d.id);
2196
2197 let c_neighbors = rt
2198 .neighbors(&tok, c.id, Direction::Out, None, None)
2199 .await
2200 .unwrap();
2201 assert_eq!(c_neighbors.len(), 1);
2202 assert_eq!(c_neighbors[0].node_id, d.id);
2203 }
2204
2205 #[tokio::test]
2206 async fn merge_entity_self_merge_rejected() {
2207 let rt = rt();
2208 let tok = NamespaceToken::local();
2209 let a = rt
2210 .create_entity(&tok, "concept", None, "A", None, None, vec![])
2211 .await
2212 .unwrap();
2213 let err = rt
2214 .merge_entity_with_reason(
2215 &tok,
2216 a.id,
2217 a.id,
2218 EntityDedupMergePolicy::PreferInto,
2219 ContentMergeStrategy::Append,
2220 false,
2221 None,
2222 )
2223 .await
2224 .unwrap_err();
2225 assert!(
2226 format!("{err:?}").contains("cannot merge an entity into itself"),
2227 "expected self-merge rejection, got: {err:?}"
2228 );
2229 }
2230
2231 #[tokio::test]
2232 async fn merge_entity_prefer_into_strategy() {
2233 let rt = rt();
2234 let tok = NamespaceToken::local();
2235 let into = rt
2236 .create_entity(
2237 &tok,
2238 "concept",
2239 None,
2240 "Into",
2241 None,
2242 Some(serde_json::json!({"a": 1})),
2243 vec![],
2244 )
2245 .await
2246 .unwrap();
2247 let from = rt
2248 .create_entity(
2249 &tok,
2250 "concept",
2251 None,
2252 "From",
2253 None,
2254 Some(serde_json::json!({"a": 2, "b": 3})),
2255 vec![],
2256 )
2257 .await
2258 .unwrap();
2259
2260 rt.merge_entity_with_reason(
2261 &tok,
2262 into.id,
2263 from.id,
2264 EntityDedupMergePolicy::PreferInto,
2265 ContentMergeStrategy::Append,
2266 false,
2267 None,
2268 )
2269 .await
2270 .unwrap();
2271
2272 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2273 let props = kept.properties.unwrap();
2274 assert_eq!(props["a"], 1);
2276 assert_eq!(props["b"], 3);
2277 }
2278
2279 #[tokio::test]
2280 async fn merge_entity_prefer_from_strategy() {
2281 let rt = rt();
2282 let tok = NamespaceToken::local();
2283 let into = rt
2284 .create_entity(
2285 &tok,
2286 "concept",
2287 None,
2288 "Into",
2289 None,
2290 Some(serde_json::json!({"a": 1})),
2291 vec![],
2292 )
2293 .await
2294 .unwrap();
2295 let from = rt
2296 .create_entity(
2297 &tok,
2298 "concept",
2299 None,
2300 "From",
2301 None,
2302 Some(serde_json::json!({"a": 2, "b": 3})),
2303 vec![],
2304 )
2305 .await
2306 .unwrap();
2307
2308 rt.merge_entity_with_reason(
2309 &tok,
2310 into.id,
2311 from.id,
2312 EntityDedupMergePolicy::PreferFrom,
2313 ContentMergeStrategy::Append,
2314 false,
2315 None,
2316 )
2317 .await
2318 .unwrap();
2319
2320 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2321 let props = kept.properties.unwrap();
2322 assert_eq!(props["a"], 2);
2324 assert_eq!(props["b"], 3);
2325 }
2326
2327 #[tokio::test]
2328 async fn merge_entity_union_strategy() {
2329 let rt = rt();
2330 let tok = NamespaceToken::local();
2331 let into = rt
2332 .create_entity(
2333 &tok,
2334 "concept",
2335 None,
2336 "Into",
2337 None,
2338 Some(serde_json::json!({"a": 1})),
2339 vec![],
2340 )
2341 .await
2342 .unwrap();
2343 let from = rt
2344 .create_entity(
2345 &tok,
2346 "concept",
2347 None,
2348 "From",
2349 None,
2350 Some(serde_json::json!({"a": 2, "b": 3})),
2351 vec![],
2352 )
2353 .await
2354 .unwrap();
2355
2356 rt.merge_entity_with_reason(
2357 &tok,
2358 into.id,
2359 from.id,
2360 EntityDedupMergePolicy::Union,
2361 ContentMergeStrategy::Append,
2362 false,
2363 None,
2364 )
2365 .await
2366 .unwrap();
2367
2368 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2369 let props = kept.properties.unwrap();
2370 assert_eq!(props["a"], 1);
2372 assert_eq!(props["b"], 3);
2373 }
2374
2375 #[tokio::test]
2376 async fn merge_entity_unions_tags() {
2377 let rt = rt();
2378 let tok = NamespaceToken::local();
2379 let into = rt
2380 .create_entity(
2381 &tok,
2382 "concept",
2383 None,
2384 "Into",
2385 None,
2386 None,
2387 vec!["x".to_string(), "y".to_string()],
2388 )
2389 .await
2390 .unwrap();
2391 let from = rt
2392 .create_entity(
2393 &tok,
2394 "concept",
2395 None,
2396 "From",
2397 None,
2398 None,
2399 vec!["y".to_string(), "z".to_string()],
2400 )
2401 .await
2402 .unwrap();
2403
2404 rt.merge_entity_with_reason(
2405 &tok,
2406 into.id,
2407 from.id,
2408 EntityDedupMergePolicy::PreferInto,
2409 ContentMergeStrategy::Append,
2410 false,
2411 None,
2412 )
2413 .await
2414 .unwrap();
2415
2416 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2417 let mut tags = kept.tags.clone();
2418 tags.sort();
2419 assert_eq!(tags, vec!["x", "y", "z"]);
2420 }
2421
2422 #[tokio::test]
2423 async fn merge_entity_drops_self_loops() {
2424 let rt = rt();
2425 let tok = NamespaceToken::local();
2426 let a = rt
2427 .create_entity(&tok, "concept", None, "A", None, None, vec![])
2428 .await
2429 .unwrap();
2430 let b = rt
2431 .create_entity(&tok, "concept", None, "B", None, None, vec![])
2432 .await
2433 .unwrap();
2434
2435 rt.link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
2437 .await
2438 .unwrap();
2439
2440 let summary = rt
2441 .merge_entity_with_reason(
2442 &tok,
2443 a.id,
2444 b.id,
2445 EntityDedupMergePolicy::PreferInto,
2446 ContentMergeStrategy::Append,
2447 false,
2448 None,
2449 )
2450 .await
2451 .unwrap();
2452
2453 assert_eq!(
2454 summary.edges_rewired, 0,
2455 "self-loop should be dropped, not rewired"
2456 );
2457
2458 let a_out = rt
2459 .neighbors(&tok, a.id, Direction::Out, None, None)
2460 .await
2461 .unwrap();
2462 assert!(a_out.is_empty(), "no self-loop should remain");
2463 }
2464
2465 #[tokio::test]
2468 async fn merge_entity_append_strategy_concatenates_descriptions() {
2469 let rt = rt();
2470 let tok = NamespaceToken::local();
2471 let into = rt
2472 .create_entity(&tok, "concept", None, "Into", Some("desc A"), None, vec![])
2473 .await
2474 .unwrap();
2475 let from = rt
2476 .create_entity(&tok, "concept", None, "From", Some("desc B"), None, vec![])
2477 .await
2478 .unwrap();
2479
2480 let summary = rt
2481 .merge_entity_with_reason(
2482 &tok,
2483 into.id,
2484 from.id,
2485 EntityDedupMergePolicy::PreferInto,
2486 ContentMergeStrategy::Append,
2487 false,
2488 None,
2489 )
2490 .await
2491 .unwrap();
2492
2493 assert!(
2494 summary.content_appended,
2495 "append strategy with two non-empty descriptions must report content_appended=true"
2496 );
2497 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2498 assert_eq!(kept.description.as_deref(), Some("desc A\n\n---\n\ndesc B"));
2499 }
2500
2501 #[tokio::test]
2502 async fn merge_entity_append_strategy_from_empty_is_noop() {
2503 let rt = rt();
2504 let tok = NamespaceToken::local();
2505 let into = rt
2506 .create_entity(&tok, "concept", None, "Into", Some("desc A"), None, vec![])
2507 .await
2508 .unwrap();
2509 let from = rt
2510 .create_entity(&tok, "concept", None, "From", None, None, vec![])
2511 .await
2512 .unwrap();
2513
2514 let summary = rt
2515 .merge_entity_with_reason(
2516 &tok,
2517 into.id,
2518 from.id,
2519 EntityDedupMergePolicy::PreferInto,
2520 ContentMergeStrategy::Append,
2521 false,
2522 None,
2523 )
2524 .await
2525 .unwrap();
2526
2527 assert!(
2528 !summary.content_appended,
2529 "from's empty description means nothing was appended"
2530 );
2531 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2532 assert_eq!(kept.description.as_deref(), Some("desc A"));
2533 }
2534
2535 #[tokio::test]
2536 async fn merge_entity_append_strategy_into_empty_takes_from() {
2537 let rt = rt();
2538 let tok = NamespaceToken::local();
2539 let into = rt
2540 .create_entity(&tok, "concept", None, "Into", None, None, vec![])
2541 .await
2542 .unwrap();
2543 let from = rt
2544 .create_entity(&tok, "concept", None, "From", Some("desc B"), None, vec![])
2545 .await
2546 .unwrap();
2547
2548 let summary = rt
2549 .merge_entity_with_reason(
2550 &tok,
2551 into.id,
2552 from.id,
2553 EntityDedupMergePolicy::PreferInto,
2554 ContentMergeStrategy::Append,
2555 false,
2556 None,
2557 )
2558 .await
2559 .unwrap();
2560
2561 assert!(
2562 summary.content_appended,
2563 "taking from's description into an empty into is real content preservation"
2564 );
2565 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2566 assert_eq!(kept.description.as_deref(), Some("desc B"));
2567 }
2568
2569 #[tokio::test]
2570 async fn merge_entity_prefer_into_strategy_still_discards_explicitly() {
2571 let rt = rt();
2572 let tok = NamespaceToken::local();
2573 let into = rt
2574 .create_entity(&tok, "concept", None, "Into", Some("desc A"), None, vec![])
2575 .await
2576 .unwrap();
2577 let from = rt
2578 .create_entity(&tok, "concept", None, "From", Some("desc B"), None, vec![])
2579 .await
2580 .unwrap();
2581
2582 let summary = rt
2583 .merge_entity_with_reason(
2584 &tok,
2585 into.id,
2586 from.id,
2587 EntityDedupMergePolicy::PreferInto,
2588 ContentMergeStrategy::PreferInto,
2589 false,
2590 None,
2591 )
2592 .await
2593 .unwrap();
2594
2595 assert!(
2596 !summary.content_appended,
2597 "explicit PreferInto opt-out must not report an append"
2598 );
2599 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2600 assert_eq!(
2601 kept.description.as_deref(),
2602 Some("desc A"),
2603 "explicit PreferInto opt-out keeps the old discard behavior"
2604 );
2605 }
2606
2607 #[tokio::test]
2612 async fn merge_entity_prefer_from_content_strategy_wins_over_default_entity_policy() {
2613 let rt = rt();
2614 let tok = NamespaceToken::local();
2615 let into = rt
2616 .create_entity(&tok, "concept", None, "Into", Some("desc A"), None, vec![])
2617 .await
2618 .unwrap();
2619 let from = rt
2620 .create_entity(&tok, "concept", None, "From", Some("desc B"), None, vec![])
2621 .await
2622 .unwrap();
2623
2624 let summary = rt
2625 .merge_entity_with_reason(
2626 &tok,
2627 into.id,
2628 from.id,
2629 EntityDedupMergePolicy::PreferInto,
2630 ContentMergeStrategy::PreferFrom,
2631 false,
2632 None,
2633 )
2634 .await
2635 .unwrap();
2636
2637 assert!(
2638 !summary.content_appended,
2639 "explicit PreferFrom is not an append"
2640 );
2641 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2642 assert_eq!(
2643 kept.description.as_deref(),
2644 Some("desc B"),
2645 "content_strategy=prefer_from must win over the default prefer_into entity policy"
2646 );
2647 }
2648
2649 #[tokio::test]
2650 async fn merge_entity_dry_run_previews_append() {
2651 let rt = rt();
2652 let tok = NamespaceToken::local();
2653 let into = rt
2654 .create_entity(&tok, "concept", None, "Into", Some("desc A"), None, vec![])
2655 .await
2656 .unwrap();
2657 let from = rt
2658 .create_entity(&tok, "concept", None, "From", Some("desc B"), None, vec![])
2659 .await
2660 .unwrap();
2661
2662 let summary = rt
2663 .merge_entity_with_reason(
2664 &tok,
2665 into.id,
2666 from.id,
2667 EntityDedupMergePolicy::PreferInto,
2668 ContentMergeStrategy::Append,
2669 true,
2670 None,
2671 )
2672 .await
2673 .unwrap();
2674
2675 assert!(summary.dry_run);
2676 assert!(
2677 summary.content_appended,
2678 "dry-run must preview the append outcome without writing"
2679 );
2680 let kept = rt.get_entity(&tok, into.id).await.unwrap();
2681 assert_eq!(
2682 kept.description.as_deref(),
2683 Some("desc A"),
2684 "dry_run=true must not mutate the into entity's description"
2685 );
2686 }
2687
2688 #[tokio::test]
2691 async fn merge_entity_dry_run_predicts_edges_rewired_without_writing() {
2692 use khive_storage::EdgeRelation;
2693
2694 let rt = rt();
2695 let tok = NamespaceToken::local();
2696 let a = rt
2697 .create_entity(&tok, "concept", None, "A", None, None, vec![])
2698 .await
2699 .unwrap();
2700 let into = rt
2701 .create_entity(&tok, "concept", None, "Into", None, None, vec![])
2702 .await
2703 .unwrap();
2704 let from = rt
2705 .create_entity(&tok, "concept", None, "From", None, None, vec![])
2706 .await
2707 .unwrap();
2708
2709 rt.link(&tok, a.id, from.id, EdgeRelation::Extends, 1.0, None)
2710 .await
2711 .unwrap();
2712
2713 let summary = rt
2714 .merge_entity_with_reason(
2715 &tok,
2716 into.id,
2717 from.id,
2718 EntityDedupMergePolicy::PreferInto,
2719 ContentMergeStrategy::Append,
2720 true,
2721 None,
2722 )
2723 .await
2724 .unwrap();
2725
2726 assert!(summary.dry_run);
2727 assert_eq!(
2728 summary.edges_rewired, 1,
2729 "dry-run must predict the edge that would be rewired, not report zero"
2730 );
2731
2732 let a_neighbors = rt
2733 .neighbors(&tok, a.id, Direction::Out, None, None)
2734 .await
2735 .unwrap();
2736 assert_eq!(a_neighbors.len(), 1);
2737 assert_eq!(
2738 a_neighbors[0].node_id, from.id,
2739 "dry_run=true must not rewire any edges"
2740 );
2741
2742 let events = rt
2743 .events(&tok)
2744 .unwrap()
2745 .query_events(
2746 khive_storage::EventFilter {
2747 kinds: vec![EventKind::EntityMerged],
2748 ..Default::default()
2749 },
2750 khive_storage::types::PageRequest {
2751 offset: 0,
2752 limit: 10,
2753 },
2754 )
2755 .await
2756 .unwrap();
2757 assert!(
2758 events.items.is_empty(),
2759 "dry_run=true must not append an EntityMerged event"
2760 );
2761 }
2762
2763 #[tokio::test]
2767 async fn merge_entity_event_reason_present_when_supplied_absent_when_not() {
2768 let rt = rt();
2769 let tok = NamespaceToken::local();
2770
2771 let into_a = rt
2772 .create_entity(&tok, "concept", None, "IntoA", None, None, vec![])
2773 .await
2774 .unwrap();
2775 let from_a = rt
2776 .create_entity(&tok, "concept", None, "FromA", None, None, vec![])
2777 .await
2778 .unwrap();
2779 rt.merge_entity_with_reason(
2780 &tok,
2781 into_a.id,
2782 from_a.id,
2783 EntityDedupMergePolicy::PreferInto,
2784 ContentMergeStrategy::Append,
2785 false,
2786 Some("duplicate".to_string()),
2787 )
2788 .await
2789 .unwrap();
2790
2791 let into_b = rt
2792 .create_entity(&tok, "concept", None, "IntoB", None, None, vec![])
2793 .await
2794 .unwrap();
2795 let from_b = rt
2796 .create_entity(&tok, "concept", None, "FromB", None, None, vec![])
2797 .await
2798 .unwrap();
2799 rt.merge_entity_with_reason(
2800 &tok,
2801 into_b.id,
2802 from_b.id,
2803 EntityDedupMergePolicy::PreferInto,
2804 ContentMergeStrategy::Append,
2805 false,
2806 None,
2807 )
2808 .await
2809 .unwrap();
2810
2811 let events = rt
2812 .events(&tok)
2813 .unwrap()
2814 .query_events(
2815 khive_storage::EventFilter {
2816 kinds: vec![EventKind::EntityMerged],
2817 ..Default::default()
2818 },
2819 khive_storage::types::PageRequest {
2820 offset: 0,
2821 limit: 10,
2822 },
2823 )
2824 .await
2825 .unwrap();
2826 assert_eq!(events.items.len(), 2);
2827
2828 let with_reason = events
2829 .items
2830 .iter()
2831 .find(|e| {
2832 e.payload.get("from_id").and_then(|v| v.as_str())
2833 == Some(from_a.id.to_string()).as_deref()
2834 })
2835 .expect("event for the reasoned merge must exist");
2836 assert_eq!(
2837 with_reason.payload.get("reason").and_then(|v| v.as_str()),
2838 Some("duplicate"),
2839 "reason must be threaded verbatim into the payload when supplied"
2840 );
2841
2842 let without_reason = events
2843 .items
2844 .iter()
2845 .find(|e| {
2846 e.payload.get("from_id").and_then(|v| v.as_str())
2847 == Some(from_b.id.to_string()).as_deref()
2848 })
2849 .expect("event for the reasonless merge must exist");
2850 assert!(
2851 without_reason.payload.get("reason").is_none(),
2852 "reason key must be absent (never null) when the caller omits it, got: {:?}",
2853 without_reason.payload
2854 );
2855 }
2856
2857 #[tokio::test]
2858 async fn merge_entity_with_reason_preserves_an_explicit_empty_reason() {
2859 let rt = rt();
2860 let tok = NamespaceToken::local();
2861 let into = rt
2862 .create_entity(&tok, "concept", None, "Into", None, None, vec![])
2863 .await
2864 .unwrap();
2865 let from = rt
2866 .create_entity(&tok, "concept", None, "From", None, None, vec![])
2867 .await
2868 .unwrap();
2869
2870 rt.merge_entity_with_reason(
2871 &tok,
2872 into.id,
2873 from.id,
2874 EntityDedupMergePolicy::PreferInto,
2875 ContentMergeStrategy::Append,
2876 false,
2877 Some(String::new()),
2878 )
2879 .await
2880 .unwrap();
2881
2882 let events = rt
2883 .events(&tok)
2884 .unwrap()
2885 .query_events(
2886 khive_storage::EventFilter {
2887 kinds: vec![EventKind::EntityMerged],
2888 ..Default::default()
2889 },
2890 khive_storage::types::PageRequest {
2891 offset: 0,
2892 limit: 10,
2893 },
2894 )
2895 .await
2896 .unwrap();
2897 assert_eq!(events.items.len(), 1);
2898 assert_eq!(
2899 events.items[0].payload.get("reason"),
2900 Some(&Value::String(String::new()))
2901 );
2902 }
2903
2904 #[tokio::test]
2905 async fn merge_entity_with_reason_rejects_secrets_before_reads_or_writes() {
2906 let rt = rt();
2907 let tok = NamespaceToken::local();
2908 let into = rt
2909 .create_entity(&tok, "concept", None, "Into", None, None, vec![])
2910 .await
2911 .unwrap();
2912 let from = rt
2913 .create_entity(&tok, "concept", None, "From", None, None, vec![])
2914 .await
2915 .unwrap();
2916 let secret = secret_shaped_reason();
2917
2918 let error = rt
2919 .merge_entity_with_reason(
2920 &tok,
2921 into.id,
2922 from.id,
2923 EntityDedupMergePolicy::PreferInto,
2924 ContentMergeStrategy::Append,
2925 false,
2926 Some(secret),
2927 )
2928 .await
2929 .unwrap_err();
2930
2931 assert!(matches!(error, RuntimeError::SecretDetected(_)));
2932 assert_eq!(rt.get_entity(&tok, into.id).await.unwrap().id, into.id);
2933 assert_eq!(rt.get_entity(&tok, from.id).await.unwrap().id, from.id);
2934 let event_count = rt
2935 .events(&tok)
2936 .unwrap()
2937 .count_events(khive_storage::EventFilter {
2938 kinds: vec![EventKind::EntityMerged],
2939 ..Default::default()
2940 })
2941 .await
2942 .unwrap();
2943 assert_eq!(event_count, 0);
2944 }
2945
2946 #[tokio::test]
2949 async fn merge_note_emits_exactly_one_note_merged_event_with_kept_and_absorbed_ids() {
2950 let rt = rt();
2951 let tok = NamespaceToken::local();
2952 let into = rt
2953 .create_note(&tok, "observation", None, "into note", None, None, vec![])
2954 .await
2955 .unwrap();
2956 let from = rt
2957 .create_note(&tok, "observation", None, "from note", None, None, vec![])
2958 .await
2959 .unwrap();
2960
2961 let summary = rt
2962 .merge_note_with_reason(
2963 &tok,
2964 into.id,
2965 from.id,
2966 EntityDedupMergePolicy::PreferInto,
2967 ContentMergeStrategy::Append,
2968 false,
2969 Some("duplicate".to_string()),
2970 )
2971 .await
2972 .unwrap();
2973
2974 let events = rt
2975 .events(&tok)
2976 .unwrap()
2977 .query_events(
2978 khive_storage::EventFilter {
2979 kinds: vec![EventKind::NoteMerged],
2980 ..Default::default()
2981 },
2982 khive_storage::types::PageRequest {
2983 offset: 0,
2984 limit: 10,
2985 },
2986 )
2987 .await
2988 .unwrap();
2989 assert_eq!(
2990 events.items.len(),
2991 1,
2992 "merge_note must emit exactly one NoteMerged event"
2993 );
2994
2995 let payload = &events.items[0].payload;
2996 assert_eq!(
2997 payload.get("into_id").and_then(|v| v.as_str()),
2998 Some(summary.kept_id.to_string()).as_deref()
2999 );
3000 assert_eq!(
3001 payload.get("from_id").and_then(|v| v.as_str()),
3002 Some(summary.removed_id.to_string()).as_deref()
3003 );
3004 assert_eq!(
3005 payload.get("reason").and_then(|v| v.as_str()),
3006 Some("duplicate")
3007 );
3008 }
3009
3010 #[tokio::test]
3011 async fn merge_note_with_reason_preserves_an_explicit_empty_reason() {
3012 let rt = rt();
3013 let tok = NamespaceToken::local();
3014 let into = rt
3015 .create_note(&tok, "observation", None, "into note", None, None, vec![])
3016 .await
3017 .unwrap();
3018 let from = rt
3019 .create_note(&tok, "observation", None, "from note", None, None, vec![])
3020 .await
3021 .unwrap();
3022
3023 rt.merge_note_with_reason(
3024 &tok,
3025 into.id,
3026 from.id,
3027 EntityDedupMergePolicy::PreferInto,
3028 ContentMergeStrategy::Append,
3029 false,
3030 Some(String::new()),
3031 )
3032 .await
3033 .unwrap();
3034
3035 let events = rt
3036 .events(&tok)
3037 .unwrap()
3038 .query_events(
3039 khive_storage::EventFilter {
3040 kinds: vec![EventKind::NoteMerged],
3041 ..Default::default()
3042 },
3043 khive_storage::types::PageRequest {
3044 offset: 0,
3045 limit: 10,
3046 },
3047 )
3048 .await
3049 .unwrap();
3050 assert_eq!(events.items.len(), 1);
3051 assert_eq!(
3052 events.items[0].payload.get("reason"),
3053 Some(&Value::String(String::new()))
3054 );
3055 }
3056
3057 #[tokio::test]
3058 async fn merge_note_with_reason_rejects_secrets_before_reads_or_writes() {
3059 let rt = rt();
3060 let tok = NamespaceToken::local();
3061 let into = rt
3062 .create_note(&tok, "observation", None, "into note", None, None, vec![])
3063 .await
3064 .unwrap();
3065 let from = rt
3066 .create_note(&tok, "observation", None, "from note", None, None, vec![])
3067 .await
3068 .unwrap();
3069 let secret = secret_shaped_reason();
3070
3071 let error = rt
3072 .merge_note_with_reason(
3073 &tok,
3074 into.id,
3075 from.id,
3076 EntityDedupMergePolicy::PreferInto,
3077 ContentMergeStrategy::Append,
3078 false,
3079 Some(secret),
3080 )
3081 .await
3082 .unwrap_err();
3083
3084 assert!(matches!(error, RuntimeError::SecretDetected(_)));
3085 let note_store = rt.notes(&tok).unwrap();
3086 assert_eq!(
3087 note_store.get_note(into.id).await.unwrap().unwrap().id,
3088 into.id
3089 );
3090 assert_eq!(
3091 note_store.get_note(from.id).await.unwrap().unwrap().id,
3092 from.id
3093 );
3094 let event_count = rt
3095 .events(&tok)
3096 .unwrap()
3097 .count_events(khive_storage::EventFilter {
3098 kinds: vec![EventKind::NoteMerged],
3099 ..Default::default()
3100 })
3101 .await
3102 .unwrap();
3103 assert_eq!(event_count, 0);
3104 }
3105
3106 #[tokio::test]
3107 async fn legacy_merge_methods_remain_source_compatible() {
3108 let rt = rt();
3109 let tok = NamespaceToken::local();
3110 let into_entity = rt
3111 .create_entity(&tok, "concept", None, "Entity A", None, None, vec![])
3112 .await
3113 .unwrap();
3114 let from_entity = rt
3115 .create_entity(&tok, "concept", None, "Entity B", None, None, vec![])
3116 .await
3117 .unwrap();
3118 let into_note = rt
3119 .create_note(&tok, "observation", None, "note A", None, None, vec![])
3120 .await
3121 .unwrap();
3122 let from_note = rt
3123 .create_note(&tok, "observation", None, "note B", None, None, vec![])
3124 .await
3125 .unwrap();
3126
3127 rt.merge_entity(
3128 &tok,
3129 into_entity.id,
3130 from_entity.id,
3131 EntityDedupMergePolicy::PreferInto,
3132 ContentMergeStrategy::Append,
3133 false,
3134 )
3135 .await
3136 .unwrap();
3137 rt.merge_note(
3138 &tok,
3139 into_note.id,
3140 from_note.id,
3141 EntityDedupMergePolicy::PreferInto,
3142 ContentMergeStrategy::Append,
3143 false,
3144 )
3145 .await
3146 .unwrap();
3147 }
3148
3149 #[test]
3152 fn union_tags_deduplicates() {
3153 let (tags, added) = union_tags(
3154 &["x".to_string(), "y".to_string()],
3155 &["y".to_string(), "z".to_string()],
3156 );
3157 let mut sorted = tags.clone();
3158 sorted.sort();
3159 assert_eq!(sorted, vec!["x", "y", "z"]);
3160 assert_eq!(added, 1);
3161 }
3162
3163 #[test]
3164 fn merge_properties_prefer_into_fills_missing_keys() {
3165 let a = serde_json::json!({"a": 1});
3166 let b = serde_json::json!({"a": 99, "b": 2});
3167 let (merged, added) =
3168 merge_properties(&Some(a), &Some(b), EntityDedupMergePolicy::PreferInto);
3169 let m = merged.unwrap();
3170 assert_eq!(m["a"], 1);
3171 assert_eq!(m["b"], 2);
3172 assert_eq!(added, 1);
3173 }
3174
3175 #[tokio::test]
3178 async fn merge_entity_tombstones_source_with_provenance() {
3179 let rt = rt();
3180 let tok = NamespaceToken::local();
3181 let into = rt
3182 .create_entity(&tok, "concept", None, "Into", None, None, vec![])
3183 .await
3184 .unwrap();
3185 let from = rt
3186 .create_entity(&tok, "concept", None, "From", None, None, vec![])
3187 .await
3188 .unwrap();
3189 let from_id = from.id;
3190
3191 rt.merge_entity_with_reason(
3192 &tok,
3193 into.id,
3194 from_id,
3195 EntityDedupMergePolicy::PreferInto,
3196 ContentMergeStrategy::Append,
3197 false,
3198 None,
3199 )
3200 .await
3201 .unwrap();
3202
3203 assert!(
3204 rt.get_entity(&tok, from_id).await.is_err(),
3205 "tombstoned source should not be returned by get_entity"
3206 );
3207
3208 let pool = rt.backend().pool_arc();
3209 let (deleted_at, merged_into): (Option<i64>, Option<String>) =
3210 tokio::task::spawn_blocking(move || {
3211 let guard = pool.writer().unwrap();
3212 guard
3213 .conn()
3214 .query_row(
3215 "SELECT deleted_at, merged_into FROM entities WHERE id = ?1",
3216 [from_id.to_string()],
3217 |row| Ok((row.get(0)?, row.get(1)?)),
3218 )
3219 .unwrap()
3220 })
3221 .await
3222 .unwrap();
3223 assert!(
3224 deleted_at.is_some(),
3225 "tombstoned entity must have deleted_at set"
3226 );
3227 assert_eq!(
3228 merged_into.as_deref(),
3229 Some(into.id.to_string().as_str()),
3230 "merged_into must point to into_id"
3231 );
3232 }
3233
3234 #[tokio::test]
3235 async fn merge_note_same_kind_appends_content() {
3236 let rt = rt();
3237 let tok = NamespaceToken::local();
3238 let into = rt
3239 .create_note(
3240 &tok,
3241 "observation",
3242 None,
3243 "Into content",
3244 None,
3245 None,
3246 vec![],
3247 )
3248 .await
3249 .unwrap();
3250 let from = rt
3251 .create_note(
3252 &tok,
3253 "observation",
3254 None,
3255 "From content",
3256 None,
3257 None,
3258 vec![],
3259 )
3260 .await
3261 .unwrap();
3262 let from_id = from.id;
3263
3264 let summary = rt
3265 .merge_note_with_reason(
3266 &tok,
3267 into.id,
3268 from_id,
3269 EntityDedupMergePolicy::PreferInto,
3270 ContentMergeStrategy::Append,
3271 false,
3272 None,
3273 )
3274 .await
3275 .unwrap();
3276
3277 assert_eq!(summary.kept_id, into.id);
3278 assert_eq!(summary.removed_id, from_id);
3279 assert!(summary.content_appended);
3280 assert!(!summary.dry_run);
3281
3282 let from_store = rt.notes(&tok).unwrap();
3283 assert!(
3284 from_store.get_note(from_id).await.unwrap().is_none(),
3285 "merged-from note should be soft-deleted"
3286 );
3287 }
3288
3289 #[tokio::test]
3292 async fn merge_note_survives_shared_edge_to_third_party() {
3293 use khive_storage::EdgeRelation;
3294 let rt = rt();
3295 let tok = NamespaceToken::local();
3296
3297 let into = rt
3298 .create_note(&tok, "observation", None, "Into", None, None, vec![])
3299 .await
3300 .unwrap();
3301 let from = rt
3302 .create_note(&tok, "observation", None, "From", None, None, vec![])
3303 .await
3304 .unwrap();
3305 let shared = rt
3306 .create_entity(&tok, "concept", None, "Shared", None, None, vec![])
3307 .await
3308 .unwrap();
3309
3310 rt.link(&tok, into.id, shared.id, EdgeRelation::Annotates, 1.0, None)
3314 .await
3315 .unwrap();
3316 rt.link(&tok, from.id, shared.id, EdgeRelation::Annotates, 1.0, None)
3317 .await
3318 .unwrap();
3319
3320 let summary = rt
3321 .merge_note_with_reason(
3322 &tok,
3323 into.id,
3324 from.id,
3325 EntityDedupMergePolicy::PreferInto,
3326 ContentMergeStrategy::Append,
3327 false,
3328 None,
3329 )
3330 .await
3331 .expect("merge must succeed even when both notes annotate the same entity");
3332
3333 assert_eq!(summary.kept_id, into.id);
3334 assert_eq!(summary.removed_id, from.id);
3335
3336 let into_edges = rt
3337 .list_edges(
3338 &tok,
3339 crate::EdgeListFilter {
3340 source_id: Some(into.id),
3341 target_id: Some(shared.id),
3342 relations: vec![EdgeRelation::Annotates],
3343 ..Default::default()
3344 },
3345 10,
3346 0,
3347 )
3348 .await
3349 .unwrap();
3350 assert_eq!(
3351 into_edges.len(),
3352 1,
3353 "exactly one live into→shared annotates edge must exist after merge; got: {into_edges:?}"
3354 );
3355 }
3356
3357 #[tokio::test]
3358 async fn merge_note_different_kinds_rejected() {
3359 let rt = rt();
3360 let tok = NamespaceToken::local();
3361 let into = rt
3362 .create_note(&tok, "observation", None, "Into", None, None, vec![])
3363 .await
3364 .unwrap();
3365 let from = rt
3366 .create_note(&tok, "decision", None, "From", None, None, vec![])
3367 .await
3368 .unwrap();
3369
3370 let result = rt
3371 .merge_note_with_reason(
3372 &tok,
3373 into.id,
3374 from.id,
3375 EntityDedupMergePolicy::PreferInto,
3376 ContentMergeStrategy::Append,
3377 false,
3378 None,
3379 )
3380 .await;
3381 assert!(result.is_err(), "merging different note kinds must fail");
3382 }
3383
3384 #[tokio::test]
3385 async fn merge_note_dry_run_leaves_notes_unchanged() {
3386 let rt = rt();
3387 let tok = NamespaceToken::local();
3388 let into = rt
3389 .create_note(
3390 &tok,
3391 "observation",
3392 None,
3393 "Into content",
3394 None,
3395 None,
3396 vec![],
3397 )
3398 .await
3399 .unwrap();
3400 let from = rt
3401 .create_note(
3402 &tok,
3403 "observation",
3404 None,
3405 "From content",
3406 None,
3407 None,
3408 vec![],
3409 )
3410 .await
3411 .unwrap();
3412 let into_id = into.id;
3413 let from_id = from.id;
3414
3415 let summary = rt
3416 .merge_note_with_reason(
3417 &tok,
3418 into_id,
3419 from_id,
3420 EntityDedupMergePolicy::PreferInto,
3421 ContentMergeStrategy::Append,
3422 true,
3423 None,
3424 )
3425 .await
3426 .unwrap();
3427
3428 assert!(summary.dry_run);
3429
3430 let store = rt.notes(&tok).unwrap();
3431 let into_after = store.get_note(into_id).await.unwrap().unwrap();
3432 let from_after = store.get_note(from_id).await.unwrap().unwrap();
3433 assert_eq!(
3434 into_after.content, "Into content",
3435 "dry_run must not mutate into-note"
3436 );
3437 assert_eq!(
3438 from_after.content, "From content",
3439 "dry_run must not mutate from-note"
3440 );
3441
3442 let events = rt
3443 .events(&tok)
3444 .unwrap()
3445 .query_events(
3446 khive_storage::EventFilter {
3447 kinds: vec![EventKind::NoteMerged],
3448 ..Default::default()
3449 },
3450 khive_storage::types::PageRequest {
3451 offset: 0,
3452 limit: 10,
3453 },
3454 )
3455 .await
3456 .unwrap();
3457 assert!(
3458 events.items.is_empty(),
3459 "dry_run=true must not append a NoteMerged event"
3460 );
3461 }
3462
3463 #[tokio::test]
3468 async fn merge_nameless_notes_fts_document_is_parity_correct() {
3469 use khive_storage::types::TextSearchRequest;
3470
3471 let rt = rt(); let tok = NamespaceToken::local();
3473
3474 let into = rt
3475 .create_note(
3476 &tok,
3477 "observation",
3478 None,
3479 "intosentinelzxq body",
3480 None,
3481 Some(serde_json::json!({"src": "into"})),
3482 vec![],
3483 )
3484 .await
3485 .expect("create into-note");
3486 let from = rt
3487 .create_note(
3488 &tok,
3489 "observation",
3490 None,
3491 "fromsentinelzxq body",
3492 None,
3493 None,
3494 vec![],
3495 )
3496 .await
3497 .expect("create from-note");
3498
3499 let into_id = into.id;
3500 let from_id = from.id;
3501
3502 rt.merge_note_with_reason(
3503 &tok,
3504 into_id,
3505 from_id,
3506 EntityDedupMergePolicy::PreferInto,
3507 ContentMergeStrategy::Append,
3508 false,
3509 None,
3510 )
3511 .await
3512 .expect("merge_note must succeed");
3513
3514 let note_store = rt.notes(&tok).expect("note store");
3515 let merged_note = note_store
3516 .get_note(into_id)
3517 .await
3518 .expect("get_note")
3519 .expect("merged note must exist");
3520
3521 let expected = note_fts_document(&merged_note);
3522
3523 let fts = rt.text_for_notes(&tok).expect("FTS store");
3524 let stored = fts
3525 .get_document("local", into_id)
3526 .await
3527 .expect("get_document must not error")
3528 .expect("FTS document must exist after merge");
3529
3530 assert_eq!(stored.subject_id, expected.subject_id, "subject_id");
3531 assert_eq!(
3532 stored.title, expected.title,
3533 "title (None for nameless note)"
3534 );
3535 assert_eq!(stored.body, expected.body, "body");
3536 assert_eq!(stored.namespace, expected.namespace, "namespace");
3537 assert_eq!(stored.kind, expected.kind, "kind");
3538
3539 assert!(
3540 stored.title.is_none(),
3541 "nameless merged note must have title=None in FTS (was NULL before fix)"
3542 );
3543
3544 let hits = fts
3546 .search(TextSearchRequest {
3547 query: "intosentinelzxq".to_string(),
3548 mode: khive_storage::types::TextQueryMode::Plain,
3549 filter: None,
3550 top_k: 10,
3551 snippet_chars: 0,
3552 })
3553 .await
3554 .expect("search");
3555 assert!(
3556 hits.iter().any(|h| h.subject_id == into_id),
3557 "merged note must be searchable by into-note content"
3558 );
3559 }
3560
3561 #[tokio::test]
3562 async fn update_edge_updates_properties() {
3563 use khive_storage::EdgeRelation;
3564 let rt = rt();
3565 let tok = NamespaceToken::local();
3566 let a = rt
3567 .create_entity(&tok, "concept", None, "A", None, None, vec![])
3568 .await
3569 .unwrap();
3570 let b = rt
3571 .create_entity(&tok, "concept", None, "B", None, None, vec![])
3572 .await
3573 .unwrap();
3574 let edge = rt
3575 .link(&tok, a.id, b.id, EdgeRelation::Extends, 0.5, None)
3576 .await
3577 .unwrap();
3578 let edge_id: Uuid = edge.id.into();
3579
3580 let updated = rt
3581 .update_edge(
3582 &tok,
3583 edge_id,
3584 EdgePatch {
3585 properties: Some(serde_json::json!({"source": "manual"})),
3586 ..Default::default()
3587 },
3588 )
3589 .await
3590 .unwrap();
3591
3592 assert_eq!(updated.metadata.as_ref().unwrap()["source"], "manual");
3593 assert!((updated.weight - 0.5).abs() < 0.001, "weight unchanged");
3594 }
3595
3596 #[tokio::test]
3600 async fn merge_entity_survives_shared_edge_to_third_party() {
3601 use khive_storage::EdgeRelation;
3602 let rt = rt();
3603 let tok = NamespaceToken::local();
3604
3605 let a = rt
3608 .create_entity(&tok, "concept", None, "A", None, None, vec![])
3609 .await
3610 .unwrap();
3611 let b = rt
3612 .create_entity(&tok, "concept", None, "B", None, None, vec![])
3613 .await
3614 .unwrap();
3615 let shared = rt
3616 .create_entity(&tok, "concept", None, "Shared", None, None, vec![])
3617 .await
3618 .unwrap();
3619
3620 rt.link(&tok, a.id, shared.id, EdgeRelation::Extends, 1.0, None)
3623 .await
3624 .unwrap();
3625 rt.link(&tok, b.id, shared.id, EdgeRelation::Extends, 1.0, None)
3626 .await
3627 .unwrap();
3628
3629 let summary = rt
3630 .merge_entity_with_reason(
3631 &tok,
3632 a.id,
3633 b.id,
3634 crate::EntityDedupMergePolicy::PreferInto,
3635 ContentMergeStrategy::Append,
3636 false,
3637 None,
3638 )
3639 .await
3640 .expect(
3641 "C1: merge must succeed even when both entities share an edge to a third party",
3642 );
3643
3644 assert_eq!(summary.kept_id, a.id);
3645 assert_eq!(summary.removed_id, b.id);
3646 let a_edges = rt
3651 .list_edges(
3652 &tok,
3653 crate::EdgeListFilter {
3654 source_id: Some(a.id),
3655 target_id: Some(shared.id),
3656 relations: vec![EdgeRelation::Extends],
3657 ..Default::default()
3658 },
3659 10,
3660 0,
3661 )
3662 .await
3663 .unwrap();
3664 assert_eq!(
3665 a_edges.len(),
3666 1,
3667 "C1: exactly one live A→shared Extends edge must exist after merge; got: {a_edges:?}"
3668 );
3669
3670 let b_after = rt.entities(&tok).unwrap().get_entity(b.id).await.unwrap();
3672 assert!(
3673 b_after.is_none(),
3674 "C3: from_entity must be tombstoned (get_entity returns None for deleted) after merge; got: {b_after:?}"
3675 );
3676 }
3677
3678 #[tokio::test]
3682 async fn merge_entity_cross_kind_rejected_at_runtime() {
3683 let rt = rt();
3684 let tok = NamespaceToken::local();
3685
3686 let concept = rt
3687 .create_entity(&tok, "concept", None, "H2Concept", None, None, vec![])
3688 .await
3689 .unwrap();
3690 let project = rt
3691 .create_entity(&tok, "project", None, "H2Project", None, None, vec![])
3692 .await
3693 .unwrap();
3694
3695 let err = rt
3696 .merge_entity_with_reason(
3697 &tok,
3698 concept.id,
3699 project.id,
3700 crate::EntityDedupMergePolicy::PreferInto,
3701 ContentMergeStrategy::Append,
3702 false,
3703 None,
3704 )
3705 .await
3706 .expect_err("H2: cross-kind merge must be rejected by runtime");
3707 assert!(
3708 matches!(err, crate::RuntimeError::InvalidInput(_)),
3709 "H2: expected InvalidInput, got: {err:?}"
3710 );
3711
3712 let concept_after = rt.get_entity(&tok, concept.id).await;
3713 let project_after = rt.get_entity(&tok, project.id).await;
3714 assert!(
3715 concept_after.is_ok(),
3716 "H2: concept must remain live after rejected merge; got: {concept_after:?}"
3717 );
3718 assert!(
3719 project_after.is_ok(),
3720 "H2: project must remain live after rejected merge; got: {project_after:?}"
3721 );
3722 }
3723
3724 #[tokio::test]
3726 async fn merge_entity_same_kind_succeeds() {
3727 let rt = rt();
3728 let tok = NamespaceToken::local();
3729
3730 let c1 = rt
3731 .create_entity(&tok, "concept", None, "Concept1", None, None, vec![])
3732 .await
3733 .unwrap();
3734 let c2 = rt
3735 .create_entity(&tok, "concept", None, "Concept2", None, None, vec![])
3736 .await
3737 .unwrap();
3738
3739 let summary = rt
3740 .merge_entity_with_reason(
3741 &tok,
3742 c1.id,
3743 c2.id,
3744 crate::EntityDedupMergePolicy::PreferInto,
3745 ContentMergeStrategy::Append,
3746 false,
3747 None,
3748 )
3749 .await
3750 .expect("same-kind merge must succeed");
3751 assert_eq!(summary.kept_id, c1.id);
3752 assert_eq!(summary.removed_id, c2.id);
3753
3754 let c2_after = rt.entities(&tok).unwrap().get_entity(c2.id).await.unwrap();
3755 assert!(c2_after.is_none(), "from_entity must be tombstoned");
3756 }
3757
3758 #[tokio::test]
3761 async fn merge_note_cross_namespace_either_id_returns_not_found() {
3762 use crate::error::RuntimeError;
3763 use crate::Namespace;
3764
3765 let rt = rt();
3766 let ns_a = NamespaceToken::for_namespace(Namespace::parse("ns-a").unwrap());
3767 let ns_b = NamespaceToken::for_namespace(Namespace::parse("ns-b").unwrap());
3768
3769 let into_a = rt
3770 .create_note(&ns_a, "observation", None, "Into A", None, None, vec![])
3771 .await
3772 .unwrap();
3773 let from_a = rt
3774 .create_note(&ns_a, "observation", None, "From A", None, None, vec![])
3775 .await
3776 .unwrap();
3777 let note_b = rt
3778 .create_note(&ns_b, "observation", None, "Note B", None, None, vec![])
3779 .await
3780 .unwrap();
3781
3782 let foreign_into = rt
3784 .merge_note_with_reason(
3785 &ns_a,
3786 note_b.id,
3787 from_a.id,
3788 EntityDedupMergePolicy::PreferInto,
3789 ContentMergeStrategy::Append,
3790 false,
3791 None,
3792 )
3793 .await;
3794 assert!(
3795 matches!(foreign_into, Err(RuntimeError::NotFound(_))),
3796 "foreign into_id must be denied before merge, got {foreign_into:?}"
3797 );
3798
3799 let foreign_from = rt
3801 .merge_note_with_reason(
3802 &ns_a,
3803 into_a.id,
3804 note_b.id,
3805 EntityDedupMergePolicy::PreferInto,
3806 ContentMergeStrategy::Append,
3807 false,
3808 None,
3809 )
3810 .await;
3811 assert!(
3812 matches!(foreign_from, Err(RuntimeError::NotFound(_))),
3813 "foreign from_id must be denied before merge, got {foreign_from:?}"
3814 );
3815 }
3816
3817 #[tokio::test]
3820 async fn update_entity_cross_namespace_succeeds() {
3821 use crate::Namespace;
3822
3823 let rt = rt();
3824 let ns_a = NamespaceToken::for_namespace(Namespace::parse("ns-a").unwrap());
3825 let ns_b = NamespaceToken::for_namespace(Namespace::parse("ns-b").unwrap());
3826
3827 let entity = rt
3828 .create_entity(
3829 &ns_a,
3830 "concept",
3831 None,
3832 "Alpha",
3833 Some("original"),
3834 None,
3835 vec![],
3836 )
3837 .await
3838 .unwrap();
3839
3840 let result = rt
3841 .update_entity(
3842 &ns_b,
3843 entity.id,
3844 EntityPatch {
3845 name: Some("Updated".into()),
3846 ..Default::default()
3847 },
3848 )
3849 .await;
3850
3851 assert!(
3852 result.is_ok(),
3853 "cross-namespace update must succeed in shared-brain OSS; got {result:?}"
3854 );
3855 assert_eq!(result.unwrap().name, "Updated");
3856 }
3857
3858 #[tokio::test]
3862 async fn merge_entity_cross_namespace_ids_fail_at_sql_layer() {
3863 use crate::Namespace;
3864
3865 let rt = rt();
3866 let ns_a = NamespaceToken::for_namespace(Namespace::parse("ns-a").unwrap());
3867 let ns_b = NamespaceToken::for_namespace(Namespace::parse("ns-b").unwrap());
3868
3869 let into_a = rt
3870 .create_entity(&ns_a, "concept", None, "Into A", None, None, vec![])
3871 .await
3872 .unwrap();
3873 let from_a = rt
3874 .create_entity(&ns_a, "concept", None, "From A", None, None, vec![])
3875 .await
3876 .unwrap();
3877 let foreign_b = rt
3878 .create_entity(&ns_b, "concept", None, "Foreign B", None, None, vec![])
3879 .await
3880 .unwrap();
3881
3882 let foreign_into = rt
3884 .merge_entity_with_reason(
3885 &ns_a,
3886 foreign_b.id,
3887 from_a.id,
3888 EntityDedupMergePolicy::PreferInto,
3889 ContentMergeStrategy::Append,
3890 false,
3891 None,
3892 )
3893 .await;
3894 assert!(
3895 foreign_into.is_err(),
3896 "cross-namespace into_id must still fail at SQL layer; got {foreign_into:?}"
3897 );
3898
3899 let foreign_from = rt
3901 .merge_entity_with_reason(
3902 &ns_a,
3903 into_a.id,
3904 foreign_b.id,
3905 EntityDedupMergePolicy::PreferInto,
3906 ContentMergeStrategy::Append,
3907 false,
3908 None,
3909 )
3910 .await;
3911 assert!(
3912 foreign_from.is_err(),
3913 "cross-namespace from_id must still fail at SQL layer; got {foreign_from:?}"
3914 );
3915
3916 assert!(rt.get_entity(&ns_a, into_a.id).await.is_ok());
3918 assert!(rt.get_entity(&ns_a, from_a.id).await.is_ok());
3919 assert!(rt.get_entity(&ns_b, foreign_b.id).await.is_ok());
3920 }
3921
3922 #[test]
3925 fn entity_fts_document_with_description() {
3926 let mut entity = Entity::new("local", "concept", "MyEntity");
3927 entity = entity.with_description("some description text");
3928 let doc = entity_fts_document(&entity);
3929 assert_eq!(doc.subject_id, entity.id);
3930 assert_eq!(doc.namespace, "local");
3931 assert_eq!(doc.title.as_deref(), Some("MyEntity"));
3932 assert_eq!(doc.body, "MyEntity some description text");
3933 assert_eq!(doc.kind, khive_types::SubstrateKind::Entity);
3934 }
3935
3936 #[test]
3937 fn entity_fts_document_without_description() {
3938 let entity = Entity::new("local", "concept", "NameOnly");
3939 let doc = entity_fts_document(&entity);
3940 assert_eq!(doc.title.as_deref(), Some("NameOnly"));
3941 assert_eq!(doc.body, "NameOnly");
3942 }
3943
3944 #[test]
3945 fn entity_fts_document_empty_description_uses_name_only() {
3946 let mut entity = Entity::new("local", "concept", "TitleOnly");
3947 entity = entity.with_description("");
3948 let doc = entity_fts_document(&entity);
3949 assert_eq!(
3950 doc.body, "TitleOnly",
3951 "empty description must not be appended"
3952 );
3953 }
3954
3955 #[tokio::test]
3959 async fn entity_fts_document_matches_runtime_create_path() {
3960 let rt = rt();
3961 let tok = NamespaceToken::local();
3962
3963 let entity = rt
3964 .create_entity(
3965 &tok,
3966 "concept",
3967 None,
3968 "CrossPathTitle",
3969 Some("cross path description body"),
3970 Some(serde_json::json!({"key": "val"})),
3971 vec!["tag1".to_string()],
3972 )
3973 .await
3974 .expect("create_entity");
3975
3976 let fts = rt.text(&tok).expect("FTS store");
3977 let stored = fts
3978 .get_document("local", entity.id)
3979 .await
3980 .expect("get_document")
3981 .expect("document must exist after create_entity");
3982
3983 let expected = entity_fts_document(&entity);
3984
3985 assert_eq!(stored.subject_id, expected.subject_id, "subject_id");
3986 assert_eq!(stored.kind, expected.kind, "kind");
3987 assert_eq!(stored.title, expected.title, "title");
3988 assert_eq!(stored.body, expected.body, "body");
3989 assert_eq!(stored.namespace, expected.namespace, "namespace");
3990 }
3991
3992 #[tokio::test]
3995 async fn entity_fts_document_matches_runtime_update_path() {
3996 let rt = rt();
3997 let tok = NamespaceToken::local();
3998
3999 let entity = rt
4000 .create_entity(
4001 &tok,
4002 "concept",
4003 None,
4004 "OldName",
4005 Some("old desc"),
4006 None,
4007 vec![],
4008 )
4009 .await
4010 .expect("create_entity");
4011
4012 let updated = rt
4013 .update_entity(
4014 &tok,
4015 entity.id,
4016 EntityPatch {
4017 name: Some("NewName".to_string()),
4018 description: Some(Some("new desc".to_string())),
4019 ..Default::default()
4020 },
4021 )
4022 .await
4023 .expect("update_entity");
4024
4025 let fts = rt.text(&tok).expect("FTS store");
4026 let stored = fts
4027 .get_document("local", updated.id)
4028 .await
4029 .expect("get_document")
4030 .expect("document must exist after update_entity");
4031
4032 let expected = entity_fts_document(&updated);
4033
4034 assert_eq!(stored.title, expected.title, "title after update");
4035 assert_eq!(stored.body, expected.body, "body after update");
4036 }
4037
4038 struct MergeTestVecService {
4044 dims: usize,
4045 }
4046
4047 #[async_trait::async_trait]
4048 impl lattice_embed::EmbeddingService for MergeTestVecService {
4049 async fn embed(
4050 &self,
4051 texts: &[String],
4052 _model: lattice_embed::EmbeddingModel,
4053 ) -> std::result::Result<Vec<Vec<f32>>, lattice_embed::EmbedError> {
4054 Ok(texts.iter().map(|_| vec![1.0_f32; self.dims]).collect())
4055 }
4056
4057 fn supports_model(&self, _model: lattice_embed::EmbeddingModel) -> bool {
4058 true
4059 }
4060
4061 fn name(&self) -> &'static str {
4062 "merge-test-const-vec"
4063 }
4064 }
4065
4066 struct MergeTestVecProvider {
4067 provider_name: String,
4068 dims: usize,
4069 }
4070
4071 impl MergeTestVecProvider {
4072 fn new(name: &str, dims: usize) -> Self {
4073 Self {
4074 provider_name: name.to_owned(),
4075 dims,
4076 }
4077 }
4078 }
4079
4080 #[async_trait::async_trait]
4081 impl crate::embedder_registry::EmbedderProvider for MergeTestVecProvider {
4082 fn name(&self) -> &str {
4083 &self.provider_name
4084 }
4085
4086 fn dimensions(&self) -> usize {
4087 self.dims
4088 }
4089
4090 async fn build(
4091 &self,
4092 ) -> crate::error::RuntimeResult<std::sync::Arc<dyn lattice_embed::EmbeddingService>>
4093 {
4094 Ok(std::sync::Arc::new(MergeTestVecService { dims: self.dims }))
4095 }
4096 }
4097
4098 #[tokio::test]
4104 async fn merge_entity_clears_vectors_from_all_registered_models() {
4105 const DIMS: usize = 4;
4106 let rt = KhiveRuntime::memory().unwrap();
4107 rt.register_embedder(MergeTestVecProvider::new("merge-vec-a", DIMS));
4108 rt.register_embedder(MergeTestVecProvider::new("merge-vec-b", DIMS));
4109
4110 let ns_str = "merge-entity-vec-cleanup";
4111 let ns = crate::Namespace::parse(ns_str).unwrap();
4112 let tok = NamespaceToken::for_namespace(ns);
4113
4114 let into_e = rt
4115 .create_entity(
4116 &tok,
4117 "concept",
4118 None,
4119 "IntoVecEntity",
4120 Some("desc a"),
4121 None,
4122 vec![],
4123 )
4124 .await
4125 .expect("create into");
4126 let from_e = rt
4127 .create_entity(
4128 &tok,
4129 "concept",
4130 None,
4131 "FromVecEntity",
4132 Some("desc b"),
4133 None,
4134 vec![],
4135 )
4136 .await
4137 .expect("create from");
4138
4139 let vs_a = rt.vectors_for_model(&tok, "merge-vec-a").unwrap();
4141 let vs_b = rt.vectors_for_model(&tok, "merge-vec-b").unwrap();
4142 use khive_storage::types::VectorSearchRequest;
4143 let query = vec![1.0_f32; DIMS];
4144 let pre_a = vs_a
4145 .search(VectorSearchRequest {
4146 query_vectors: vec![query.clone()],
4147 top_k: 100,
4148 namespace: Some(ns_str.to_string()),
4149 kind: Some(khive_types::SubstrateKind::Entity),
4150 embedding_model: Some("merge-vec-a".to_string()),
4151 filter: None,
4152 backend_hints: None,
4153 })
4154 .await
4155 .unwrap();
4156 assert!(
4157 pre_a.iter().any(|h| h.subject_id == into_e.id)
4158 && pre_a.iter().any(|h| h.subject_id == from_e.id),
4159 "both entities must be in model-a before merge; got {pre_a:?}"
4160 );
4161
4162 let pre_b = vs_b
4165 .search(VectorSearchRequest {
4166 query_vectors: vec![query.clone()],
4167 top_k: 100,
4168 namespace: Some(ns_str.to_string()),
4169 kind: Some(khive_types::SubstrateKind::Entity),
4170 embedding_model: Some("merge-vec-b".to_string()),
4171 filter: None,
4172 backend_hints: None,
4173 })
4174 .await
4175 .unwrap();
4176 assert!(
4177 pre_b.iter().any(|h| h.subject_id == into_e.id)
4178 && pre_b.iter().any(|h| h.subject_id == from_e.id),
4179 "both entities must be in model-b before merge; got {pre_b:?}"
4180 );
4181
4182 rt.merge_entity_with_reason(
4183 &tok,
4184 into_e.id,
4185 from_e.id,
4186 EntityDedupMergePolicy::PreferInto,
4187 ContentMergeStrategy::Append,
4188 false,
4189 None,
4190 )
4191 .await
4192 .expect("merge_entity");
4193
4194 let post_a = vs_a
4195 .search(VectorSearchRequest {
4196 query_vectors: vec![query.clone()],
4197 top_k: 100,
4198 namespace: Some(ns_str.to_string()),
4199 kind: Some(khive_types::SubstrateKind::Entity),
4200 embedding_model: Some("merge-vec-a".to_string()),
4201 filter: None,
4202 backend_hints: None,
4203 })
4204 .await
4205 .unwrap();
4206 let from_ids_a: Vec<_> = post_a
4207 .iter()
4208 .filter(|h| h.subject_id == from_e.id)
4209 .collect();
4210 assert!(
4211 from_ids_a.is_empty(),
4212 "from_id must have no vectors in model-a after merge; got {from_ids_a:?}"
4213 );
4214
4215 let post_b = vs_b
4216 .search(VectorSearchRequest {
4217 query_vectors: vec![query],
4218 top_k: 100,
4219 namespace: Some(ns_str.to_string()),
4220 kind: Some(khive_types::SubstrateKind::Entity),
4221 embedding_model: Some("merge-vec-b".to_string()),
4222 filter: None,
4223 backend_hints: None,
4224 })
4225 .await
4226 .unwrap();
4227 let from_ids_b: Vec<_> = post_b
4228 .iter()
4229 .filter(|h| h.subject_id == from_e.id)
4230 .collect();
4231 assert!(
4232 from_ids_b.is_empty(),
4233 "from_id must have no vectors in model-b after merge; got {from_ids_b:?}"
4234 );
4235 }
4236
4237 #[tokio::test]
4243 async fn merge_note_clears_vectors_from_all_registered_models() {
4244 const DIMS: usize = 4;
4245 let rt = KhiveRuntime::memory().unwrap();
4246 rt.register_embedder(MergeTestVecProvider::new("merge-note-vec-a", DIMS));
4247 rt.register_embedder(MergeTestVecProvider::new("merge-note-vec-b", DIMS));
4248
4249 let ns_str = "merge-note-vec-cleanup";
4250 let ns = crate::Namespace::parse(ns_str).unwrap();
4251 let tok = NamespaceToken::for_namespace(ns);
4252
4253 let into_n = rt
4254 .create_note(
4255 &tok,
4256 "observation",
4257 None,
4258 "IntoVecNote content",
4259 None,
4260 None,
4261 vec![],
4262 )
4263 .await
4264 .expect("create into note");
4265 let from_n = rt
4266 .create_note(
4267 &tok,
4268 "observation",
4269 None,
4270 "FromVecNote content",
4271 None,
4272 None,
4273 vec![],
4274 )
4275 .await
4276 .expect("create from note");
4277
4278 let vs_a = rt.vectors_for_model(&tok, "merge-note-vec-a").unwrap();
4279 let vs_b = rt.vectors_for_model(&tok, "merge-note-vec-b").unwrap();
4280 use khive_storage::types::VectorSearchRequest;
4281 let query = vec![1.0_f32; DIMS];
4282
4283 let pre_a = vs_a
4284 .search(VectorSearchRequest {
4285 query_vectors: vec![query.clone()],
4286 top_k: 100,
4287 namespace: Some(ns_str.to_string()),
4288 kind: Some(khive_types::SubstrateKind::Note),
4289 embedding_model: Some("merge-note-vec-a".to_string()),
4290 filter: None,
4291 backend_hints: None,
4292 })
4293 .await
4294 .unwrap();
4295 assert!(
4296 pre_a.iter().any(|h| h.subject_id == into_n.id)
4297 && pre_a.iter().any(|h| h.subject_id == from_n.id),
4298 "both notes must be in model-a before merge; got {pre_a:?}"
4299 );
4300
4301 let pre_b = vs_b
4304 .search(VectorSearchRequest {
4305 query_vectors: vec![query.clone()],
4306 top_k: 100,
4307 namespace: Some(ns_str.to_string()),
4308 kind: Some(khive_types::SubstrateKind::Note),
4309 embedding_model: Some("merge-note-vec-b".to_string()),
4310 filter: None,
4311 backend_hints: None,
4312 })
4313 .await
4314 .unwrap();
4315 assert!(
4316 pre_b.iter().any(|h| h.subject_id == into_n.id)
4317 && pre_b.iter().any(|h| h.subject_id == from_n.id),
4318 "both notes must be in model-b before merge; got {pre_b:?}"
4319 );
4320
4321 rt.merge_note_with_reason(
4322 &tok,
4323 into_n.id,
4324 from_n.id,
4325 EntityDedupMergePolicy::PreferInto,
4326 ContentMergeStrategy::PreferInto,
4327 false,
4328 None,
4329 )
4330 .await
4331 .expect("merge_note");
4332
4333 let post_a = vs_a
4334 .search(VectorSearchRequest {
4335 query_vectors: vec![query.clone()],
4336 top_k: 100,
4337 namespace: Some(ns_str.to_string()),
4338 kind: Some(khive_types::SubstrateKind::Note),
4339 embedding_model: Some("merge-note-vec-a".to_string()),
4340 filter: None,
4341 backend_hints: None,
4342 })
4343 .await
4344 .unwrap();
4345 let from_ids_a: Vec<_> = post_a
4346 .iter()
4347 .filter(|h| h.subject_id == from_n.id)
4348 .collect();
4349 assert!(
4350 from_ids_a.is_empty(),
4351 "from_id must have no vectors in model-a after merge; got {from_ids_a:?}"
4352 );
4353
4354 let post_b = vs_b
4355 .search(VectorSearchRequest {
4356 query_vectors: vec![query],
4357 top_k: 100,
4358 namespace: Some(ns_str.to_string()),
4359 kind: Some(khive_types::SubstrateKind::Note),
4360 embedding_model: Some("merge-note-vec-b".to_string()),
4361 filter: None,
4362 backend_hints: None,
4363 })
4364 .await
4365 .unwrap();
4366 let from_ids_b: Vec<_> = post_b
4367 .iter()
4368 .filter(|h| h.subject_id == from_n.id)
4369 .collect();
4370 assert!(
4371 from_ids_b.is_empty(),
4372 "from_id must have no vectors in model-b after merge; got {from_ids_b:?}"
4373 );
4374 }
4375
4376 #[tokio::test]
4379 async fn entity_fts_document_matches_runtime_merge_path() {
4380 let rt = rt();
4381 let tok = NamespaceToken::local();
4382
4383 let into_e = rt
4384 .create_entity(
4385 &tok,
4386 "concept",
4387 None,
4388 "IntoEntity",
4389 Some("into desc"),
4390 None,
4391 vec![],
4392 )
4393 .await
4394 .expect("create into");
4395 let from_e = rt
4396 .create_entity(
4397 &tok,
4398 "concept",
4399 None,
4400 "FromEntity",
4401 Some("from desc"),
4402 None,
4403 vec![],
4404 )
4405 .await
4406 .expect("create from");
4407
4408 let summary = rt
4409 .merge_entity_with_reason(
4410 &tok,
4411 into_e.id,
4412 from_e.id,
4413 EntityDedupMergePolicy::PreferInto,
4414 ContentMergeStrategy::Append,
4415 false,
4416 None,
4417 )
4418 .await
4419 .expect("merge_entity");
4420
4421 let kept = rt
4422 .get_entity(&tok, summary.kept_id)
4423 .await
4424 .expect("get kept");
4425
4426 let fts = rt.text(&tok).expect("FTS store");
4427 let stored = fts
4428 .get_document("local", kept.id)
4429 .await
4430 .expect("get_document")
4431 .expect("FTS document must exist for kept entity after merge");
4432
4433 let expected = entity_fts_document(&kept);
4434
4435 assert_eq!(stored.title, expected.title, "title after merge");
4436 assert_eq!(stored.body, expected.body, "body after merge");
4437 }
4438}