diff --git a/core/Cargo.lock b/core/Cargo.lock index 49fc652f4..32954c205 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -1626,6 +1626,7 @@ dependencies = [ "serde_json", "sha2", "solstone-core-entity-matching", + "solstone-core-indexer-store", "solstone-core-journal-io", "unicode-normalization", "zip", diff --git a/core/crates/solstone-core-entity/Cargo.toml b/core/crates/solstone-core-entity/Cargo.toml index 3641de457..6bd690d55 100644 --- a/core/crates/solstone-core-entity/Cargo.toml +++ b/core/crates/solstone-core-entity/Cargo.toml @@ -19,6 +19,7 @@ serde.workspace = true serde_json.workspace = true sha2 = "0.10.9" solstone-core-journal-io.workspace = true +solstone-core-indexer-store.workspace = true unicode-normalization.workspace = true solstone-core-entity-matching.workspace = true zip = { version = "=2.4.2", default-features = false, features = ["deflate"] } diff --git a/core/crates/solstone-core-entity/src/fixture_tests.rs b/core/crates/solstone-core-entity/src/fixture_tests.rs index 3d2023c94..d43e35079 100644 --- a/core/crates/solstone-core-entity/src/fixture_tests.rs +++ b/core/crates/solstone-core-entity/src/fixture_tests.rs @@ -21,6 +21,9 @@ use super::test_support::{ const JOURNAL_IO_ALLOWED: &[&str] = &[ "read_json", "read_jsonl", + "day_dirs", + "iter_segments", + "PathOrDay", "read_text", "MalformedPolicy", "contained_path", @@ -44,6 +47,18 @@ const JOURNAL_IO_ALLOWED: &[&str] = &[ "StagedDirOptions", "StagedWriteError", "remove_dir_all", + "append_jsonl", + "capture_snapshot", + "restore_snapshot", + "JournalSnapshot", + "SnapshotError", + "AppendError", + "atomic_replace", + "read_bytes", + "read_json", + "read_jsonl", + "write_json", + "write_jsonl", ]; #[test] @@ -134,6 +149,8 @@ fn collect_production_sources(directory: &Path, sources: &mut Vec) { Some( "fixture_tests.rs" | "resolution_tests.rs" + | "merge_payload_tests.rs" + | "merge_tests.rs" | "store_tests.rs" | "test_support.rs" | "trust_lock_tests.rs" diff --git a/core/crates/solstone-core-entity/src/lib.rs b/core/crates/solstone-core-entity/src/lib.rs index 55d9e91da..37e792671 100644 --- a/core/crates/solstone-core-entity/src/lib.rs +++ b/core/crates/solstone-core-entity/src/lib.rs @@ -20,14 +20,16 @@ pub use store::{ EntityAmbiguityRescopeError, EntityAmbiguityRescopeReport, EntityIdentityGroupMap, EntityIdentityMap, EntityIdentityRepairError, EntityIdentityRepairGuard, EntityIdentityRepairRefusal, EntityIdentityRepairReport, EntityIdentityRepairSkip, - EntityIdentityRepairSkipReason, EntityOperationContext, EntityOperationKind, EntitySaveResult, + EntityIdentityRepairSkipReason, EntityMergeError, EntityMergeOptions, EntityMergePreview, + EntityMergeReport, EntityOperationContext, EntityOperationKind, EntitySaveResult, EntityStoreError, EntityWriteError, HistoryEvent, IdentityMapCacheLoad, IdentityMapLoser, IdentityMapLoserReason, IdentitySnapshot, PreparedHistoryEvent, PreparedHistoryOutcome, - classify_prepared_history, guard_restore_does_not_cross_merge, guard_visible_event_collision, - load_resolved_ambiguity_choice, read_ambiguities, read_entity_identity, - read_identity_group_map, read_identity_map, read_prepared_history, read_visible_history, - record_ambiguity_choice, record_ambiguity_observation, refresh_identity_map_cache, - repair_entity_identities, rescope_facet_ambiguities, save_entity_identity, + classify_prepared_history, commit_entity_merge, guard_restore_does_not_cross_merge, + guard_visible_event_collision, load_resolved_ambiguity_choice, preview_entity_merge, + read_ambiguities, read_entity_identity, read_identity_group_map, read_identity_map, + read_prepared_history, read_visible_history, record_ambiguity_choice, + record_ambiguity_observation, refresh_identity_map_cache, repair_entity_identities, + rescope_facet_ambiguities, save_entity_identity, }; pub use trust_lock::{EntityTrustLock, EntityTrustLockError, hold_entity_trust_lock}; @@ -40,6 +42,10 @@ pub(crate) use store::{ #[cfg(test)] mod fixture_tests; #[cfg(test)] +mod merge_payload_tests; +#[cfg(test)] +mod merge_tests; +#[cfg(test)] mod resolution_tests; #[cfg(test)] mod store_tests; diff --git a/core/crates/solstone-core-entity/src/merge_tests.rs b/core/crates/solstone-core-entity/src/merge_tests.rs new file mode 100644 index 000000000..6c9b16b6a --- /dev/null +++ b/core/crates/solstone-core-entity/src/merge_tests.rs @@ -0,0 +1,1182 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use super::store::merge::commit_entity_merge_with_injector; +use super::store::merge::merge_facets; +use super::store::merge::merge_voiceprints; +use super::store::merge::{dedupe_akas, dedupe_emails, dedupe_observations}; +use super::store::merge_payload::{list_entity_merge_payload_ids, load_entity_merge_payload}; +use super::store::voiceprints::{read_voiceprints_npz, write_voiceprints_npz}; +use crate::{ + EntityMergeOptions, commit_entity_merge, guard_restore_does_not_cross_merge, + hold_entity_trust_lock, preview_entity_merge, read_entity_identity, read_visible_history, + save_entity_identity, +}; +use serde_json::json; +use std::fs; +use std::path::PathBuf; +use std::sync::atomic::{AtomicU64, Ordering}; + +static NEXT_VOICEPRINT_DIRECTORY: AtomicU64 = AtomicU64::new(0); + +fn voiceprint_journal() -> PathBuf { + let path = std::env::temp_dir().join(format!( + "solstone-voiceprint-{}-{}", + std::process::id(), + NEXT_VOICEPRINT_DIRECTORY.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&path).unwrap(); + path +} +fn row(value: f32) -> Vec { + vec![value; 256] +} +fn metadata(key: &str) -> String { + format!(r#"{{"day":"d","segment_key":"s","source":"x","sentence_id":"{key}"}}"#) +} +fn write_voiceprints( + journal: &std::path::Path, + id: &str, + rows: Vec>, + metadata: Vec, +) { + let path = journal.join(format!("entities/{id}/voiceprints.npz")); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + let embeddings = rows.into_iter().flatten().collect::>(); + fs::write(path, write_voiceprints_npz(&embeddings, &metadata).unwrap()).unwrap(); +} +fn read_rows(journal: &std::path::Path, id: &str) -> super::store::voiceprints::VoiceprintArchive { + read_voiceprints_npz(&fs::read(journal.join(format!("entities/{id}/voiceprints.npz"))).unwrap()) + .unwrap() +} + +#[test] +fn voiceprints_copy_source_when_target_is_missing() { + let journal = voiceprint_journal(); + write_voiceprints(&journal, "source", vec![row(2.0)], vec![metadata("1")]); + assert_eq!( + merge_voiceprints(&journal, "source", "target") + .unwrap() + .added, + 1 + ); + assert_eq!(read_rows(&journal, "target").metadata, vec![metadata("1")]); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn voiceprints_append_without_overlap() { + let journal = voiceprint_journal(); + write_voiceprints(&journal, "source", vec![row(2.0)], vec![metadata("2")]); + write_voiceprints(&journal, "target", vec![row(3.0)], vec![metadata("1")]); + assert_eq!( + merge_voiceprints(&journal, "source", "target") + .unwrap() + .target_total, + 2 + ); + assert_eq!( + read_rows(&journal, "target").metadata, + vec![metadata("1"), metadata("2")] + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn voiceprints_skip_shared_target_key() { + let journal = voiceprint_journal(); + write_voiceprints(&journal, "source", vec![row(2.0)], vec![metadata("1")]); + write_voiceprints(&journal, "target", vec![row(3.0)], vec![metadata("1")]); + let stats = merge_voiceprints(&journal, "source", "target").unwrap(); + assert_eq!( + (stats.added, stats.skipped_duplicate, stats.target_total), + (0, 1, 1) + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn voiceprints_skip_internal_duplicate_key() { + let journal = voiceprint_journal(); + write_voiceprints( + &journal, + "source", + vec![row(2.0), row(3.0)], + vec![metadata("1"), metadata("1")], + ); + let stats = merge_voiceprints(&journal, "source", "target").unwrap(); + assert_eq!((stats.added, stats.skipped_duplicate), (1, 1)); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn voiceprints_do_nothing_without_source_file() { + let journal = voiceprint_journal(); + write_voiceprints(&journal, "target", vec![row(3.0)], vec![metadata("1")]); + assert_eq!( + merge_voiceprints(&journal, "source", "target") + .unwrap() + .target_total, + 1 + ); + assert_eq!(read_rows(&journal, "target").metadata, vec![metadata("1")]); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn voiceprints_skip_degenerate_rows_without_duplicate_count() { + let journal = voiceprint_journal(); + write_voiceprints(&journal, "source", vec![row(0.0)], vec![metadata("1")]); + let stats = merge_voiceprints(&journal, "source", "target").unwrap(); + assert_eq!( + (stats.added, stats.skipped_duplicate, stats.target_total), + (0, 0, 0) + ); + assert!(!journal.join("entities/target/voiceprints.npz").exists()); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn facets_move_relationship_and_observations() { + let journal = voiceprint_journal(); + let dir = journal.join("facets/work/entities/source"); + fs::create_dir_all(&dir).unwrap(); + fs::write( + dir.join("entity.json"), + br#"{"entity_id":"source","attached_at":"2026-01-01"}"#, + ) + .unwrap(); + fs::write( + dir.join("observations.jsonl"), + b"{\"content\":\"note\",\"observed_at\":\"x\"}\n", + ) + .unwrap(); + assert_eq!( + merge_facets(&journal, "source", "target", None) + .unwrap() + .moved_count, + 1 + ); + let relationship: serde_json::Value = serde_json::from_slice( + &fs::read(journal.join("facets/work/entities/target/entity.json")).unwrap(), + ) + .unwrap(); + assert_eq!(relationship["entity_id"], "target"); + assert_eq!( + fs::read_to_string(journal.join("facets/work/entities/target/observations.jsonl")).unwrap(), + "{\"content\":\"note\",\"observed_at\":\"x\"}\n" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn email_dedup_keeps_first_seen_order() { + assert_eq!( + dedupe_emails( + &["z@example.test".to_owned()], + &["a@example.test".to_owned()] + ), + ["z@example.test", "a@example.test"] + ); +} + +#[test] +fn alias_dedup_sorts_by_lowercase() { + assert_eq!( + dedupe_akas(&["Zulu".to_owned(), "alpha".to_owned()]), + ["alpha", "Zulu"] + ); +} + +#[test] +fn dedup_keeps_first_case_variant_and_does_not_normalize() { + assert_eq!( + dedupe_akas(&[ + "Jane".to_owned(), + "jane".to_owned(), + " Jane".to_owned(), + "Straße".to_owned(), + "STRASSE".to_owned(), + "İ".to_owned(), + "i".to_owned(), + "ẞ".to_owned(), + "ß".to_owned() + ]), + [" Jane", "i", "İ", "Jane", "STRASSE", "Straße", "ẞ"] + ); +} + +#[test] +fn dedup_survives_strasse_case_variants() { + assert_eq!( + dedupe_akas(&["Straße".to_owned(), "STRASSE".to_owned()]), + ["STRASSE", "Straße"] + ); + assert_eq!( + dedupe_emails(&["Straße".to_owned()], &["STRASSE".to_owned()]), + ["Straße", "STRASSE"] + ); +} + +#[test] +fn dedup_survives_dotted_i_variants() { + assert_eq!( + dedupe_akas(&["İstanbul".to_owned(), "istanbul".to_owned()]), + ["istanbul", "İstanbul"] + ); + assert_eq!( + dedupe_emails(&["İstanbul".to_owned()], &["istanbul".to_owned()]), + ["İstanbul", "istanbul"] + ); +} + +#[test] +fn dedup_survives_nfc_nfd_variants() { + assert_eq!( + dedupe_akas(&["café".to_owned(), "cafe\u{301}".to_owned()]), + ["cafe\u{301}", "café"] + ); + assert_eq!( + dedupe_emails(&["café".to_owned()], &["cafe\u{301}".to_owned()]), + ["café", "cafe\u{301}"] + ); +} + +#[test] +fn dedup_collapses_sharp_s_variants() { + assert_eq!(dedupe_akas(&["ẞ".to_owned(), "ß".to_owned()]), ["ẞ"]); + assert_eq!(dedupe_emails(&["ẞ".to_owned()], &["ß".to_owned()]), ["ẞ"]); +} + +#[test] +fn dedup_does_not_collapse_whitespace_padded_variants() { + assert_eq!( + dedupe_akas(&[" Jane".to_owned(), "Jane".to_owned()]), + [" Jane", "Jane"] + ); + assert_eq!( + dedupe_emails(&[" Jane".to_owned()], &["Jane".to_owned()]), + [" Jane", "Jane"] + ); +} + +#[test] +fn observation_dedup_prefers_target_entries() { + assert_eq!( + dedupe_observations( + &[json!({"content":"one","observed_at":"x"})], + &[json!({"content":"one","observed_at":"x"})] + ), + [json!({"content":"one","observed_at":"x"})] + ); +} + +#[test] +fn commit_cleanup_removes_touched_source_facet_and_discovery_cache() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let source = journal.join("facets/work/entities/source"); + let target = journal.join("facets/work/entities/target"); + fs::create_dir_all(&source).unwrap(); + fs::create_dir_all(&target).unwrap(); + fs::write(source.join("entity.json"), br#"{"entity_id":"source"}"#).unwrap(); + fs::write(target.join("entity.json"), br#"{"entity_id":"target"}"#).unwrap(); + let discovery = journal.join("awareness/discovery_clusters.json"); + fs::create_dir_all(discovery.parent().unwrap()).unwrap(); + fs::write(&discovery, b"{}").unwrap(); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!(!source.exists()); + assert!(!discovery.exists()); + fs::remove_dir_all(journal).unwrap(); +} + +fn commit_segment_merge(journal: &std::path::Path) { + for id in ["source", "target"] { + save_entity_identity( + journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + commit_entity_merge(journal, "source", "target", EntityMergeOptions::default()).unwrap(); +} + +#[test] +fn commit_rewrites_segment_speaker_labels() { + let journal = voiceprint_journal(); + let path = journal.join("chronicle/20260102/080000_300/talents/speaker_labels.json"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, br#"{"labels":[{"speaker":"source"}]}"#).unwrap(); + commit_segment_merge(&journal); + let value: serde_json::Value = serde_json::from_slice(&fs::read(&path).unwrap()).unwrap(); + assert_eq!(value["labels"][0]["speaker"], "target"); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn commit_rewrites_segment_speaker_corrections() { + let journal = voiceprint_journal(); + let path = journal.join("chronicle/20260102/080000_300/talents/speaker_corrections.json"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write( + &path, + br#"{"corrections":[{"original_speaker":"source","corrected_speaker":"other"}]}"#, + ) + .unwrap(); + commit_segment_merge(&journal); + let value: serde_json::Value = serde_json::from_slice(&fs::read(&path).unwrap()).unwrap(); + assert_eq!(value["corrections"][0]["original_speaker"], "target"); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn commit_rewrites_activity_entity_references() { + let journal = voiceprint_journal(); + let path = journal.join("facets/work/activities/20260102.jsonl"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, b"{\"id\":\"activity\",\"active_entities\":[\"source\"],\"commitments\":[{\"owner_entity_id\":\"source\",\"counterparty_entity_id\":\"someone-else\"}]}\n").unwrap(); + commit_segment_merge(&journal); + let row: serde_json::Value = serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap(); + assert_eq!(row["active_entities"][0], "target"); + assert_eq!(row["commitments"][0]["owner_entity_id"], "target"); + assert_eq!( + row["commitments"][0]["counterparty_entity_id"], + "someone-else" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn commit_remaps_other_entity_observation_relation() { + let journal = voiceprint_journal(); + let path = journal.join("facets/work/entities/other/observations.jsonl"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, b"{\"relation\":{\"target_entity_id\":\"source\"}}\n").unwrap(); + commit_segment_merge(&journal); + let row: serde_json::Value = serde_json::from_str(&fs::read_to_string(&path).unwrap()).unwrap(); + assert_eq!(row["relation"]["target_entity_id"], "target"); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn commit_payload_records_history_sequence() { + let journal = voiceprint_journal(); + commit_segment_merge(&journal); + let merge_id = list_entity_merge_payload_ids(&journal, "target") + .unwrap() + .pop() + .unwrap(); + let payload = load_entity_merge_payload(&journal, "target", &merge_id).unwrap(); + let expected = read_visible_history(&journal, "target") + .unwrap() + .last() + .unwrap() + .sequence() + .unwrap(); + assert_eq!( + payload["commit_seq"].as_i64().map(i128::from), + Some(expected) + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn commit_records_matching_payload_and_audit_counts() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + write_voiceprints(&journal, "source", vec![row(2.0)], vec![metadata("count")]); + let facet = journal.join("facets/work/entities/source"); + fs::create_dir_all(&facet).unwrap(); + fs::write(facet.join("entity.json"), br#"{"entity_id":"source"}"#).unwrap(); + let activity = journal.join("facets/work/activities/20260102.jsonl"); + fs::create_dir_all(activity.parent().unwrap()).unwrap(); + fs::write(&activity, b"{\"active_entities\":[\"source\"]}\n").unwrap(); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let audit: serde_json::Value = serde_json::from_str( + &fs::read_to_string(journal.join("logs/entity-merges.jsonl")).unwrap(), + ) + .unwrap(); + assert!(audit["ts"].as_u64().is_some()); + assert_eq!(audit["source_display_name"], "source"); + assert_eq!(audit["target_display_name"], "target"); + assert_eq!(audit["principal_transferred"], false); + assert_eq!(audit["caller"], serde_json::Value::Null); + assert!(audit.get("kind").is_none()); + assert!(audit["counts"]["facets"]["moved"].as_u64().unwrap() > 0); + assert!(audit["counts"]["voiceprints"]["added"].as_u64().unwrap() > 0); + assert!( + audit["counts"]["activities"]["records_rewritten"] + .as_u64() + .unwrap() + > 0 + ); + let merge_id = list_entity_merge_payload_ids(&journal, "target") + .unwrap() + .pop() + .unwrap(); + let payload = load_entity_merge_payload(&journal, "target", &merge_id).unwrap(); + assert_eq!(payload["result_counts"], audit["counts"]); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn facets_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let source = journal.join("facets/work/entities/source"); + fs::create_dir_all(&source).unwrap(); + fs::write(source.join("entity.json"), br#"{"entity_id":"source"}"#).unwrap(); + let result = commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase: &str| phase == "facets"), + ); + assert!(result.is_err()); + assert!(source.exists()); + assert!(journal.join("entities/source").exists()); + assert!(!journal.join("facets/work/entities/target").exists()); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!(!source.exists()); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn voiceprints_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + write_voiceprints(&journal, "source", vec![row(2.0)], vec![metadata("inject")]); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "voiceprints") + ) + .is_err() + ); + assert!(!journal.join("entities/target/voiceprints.npz").exists()); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + read_rows(&journal, "target").metadata, + vec![metadata("inject")] + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn segments_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let path = journal.join("chronicle/20260102/080000_300/talents/speaker_labels.json"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, br#"{"labels":[{"speaker":"source"}]}"#).unwrap(); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "segments") + ) + .is_err() + ); + assert_eq!( + serde_json::from_slice::(&fs::read(&path).unwrap()).unwrap()["labels"] + [0]["speaker"], + "source" + ); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + serde_json::from_slice::(&fs::read(&path).unwrap()).unwrap()["labels"] + [0]["speaker"], + "target" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn activities_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let path = journal.join("facets/work/activities/20260102.jsonl"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, b"{\"active_entities\":[\"source\"]}\n").unwrap(); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "activities") + ) + .is_err() + ); + assert_eq!( + serde_json::from_str::(&fs::read_to_string(&path).unwrap()).unwrap()["active_entities"] + [0], + "source" + ); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + serde_json::from_str::(&fs::read_to_string(&path).unwrap()).unwrap()["active_entities"] + [0], + "target" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn observation_relations_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let path = journal.join("facets/work/entities/other/observations.jsonl"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, b"{\"relation\":{\"target_entity_id\":\"source\"}}\n").unwrap(); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "observation relation remap") + ) + .is_err() + ); + assert_eq!( + serde_json::from_str::(&fs::read_to_string(&path).unwrap()).unwrap()["relation"] + ["target_entity_id"], + "source" + ); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + serde_json::from_str::(&fs::read_to_string(&path).unwrap()).unwrap()["relation"] + ["target_entity_id"], + "target" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn private_payload_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "private_payload") + ) + .is_err() + ); + assert!( + list_entity_merge_payload_ids(&journal, "target") + .unwrap() + .is_empty() + ); + assert!(journal.join("entities/source").exists()); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!( + !list_entity_merge_payload_ids(&journal, "target") + .unwrap() + .is_empty() + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn lineage_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["grandparent", "source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + commit_entity_merge( + &journal, + "grandparent", + "source", + EntityMergeOptions::default(), + ) + .unwrap(); + let descendant = list_entity_merge_payload_ids(&journal, "source") + .unwrap() + .pop() + .unwrap(); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "lineage") + ) + .is_err() + ); + assert!( + list_entity_merge_payload_ids(&journal, "source") + .unwrap() + .contains(&descendant) + ); + assert!(journal.join("entities/source").exists()); + assert!(journal.join("entities/target").exists()); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!( + list_entity_merge_payload_ids(&journal, "target") + .unwrap() + .contains(&descendant) + ); + assert!( + list_entity_merge_payload_ids(&journal, "source") + .unwrap() + .is_empty() + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn cleanup_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let source = journal.join("facets/work/entities/source"); + let target = journal.join("facets/work/entities/target"); + fs::create_dir_all(&source).unwrap(); + fs::create_dir_all(&target).unwrap(); + fs::write(source.join("entity.json"), br#"{"entity_id":"source"}"#).unwrap(); + fs::write(target.join("entity.json"), br#"{"entity_id":"target"}"#).unwrap(); + let discovery = journal.join("awareness/discovery_clusters.json"); + fs::create_dir_all(discovery.parent().unwrap()).unwrap(); + fs::write(&discovery, b"{}").unwrap(); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "cleanup") + ) + .is_err() + ); + assert!(source.exists()); + assert!(discovery.exists()); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!(!source.exists()); + assert!(!discovery.exists()); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn history_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let before = read_visible_history(&journal, "target").unwrap(); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "history") + ) + .is_err() + ); + let after = read_visible_history(&journal, "target").unwrap(); + assert_eq!(after.len(), before.len()); + assert!(after.iter().all(|event| event.value()["kind"] != "merge")); + assert!(journal.join("entities/source").exists()); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!( + read_visible_history(&journal, "target") + .unwrap() + .iter() + .any(|event| event.value()["kind"] == "merge") + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn edges_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let connection = solstone_core_indexer_store::db::open_index(&journal).unwrap(); + connection + .execute( + "INSERT INTO edges(src, dst, kind, directed, source, path, weight) VALUES ('source', 'other', 'related', 1, 'manual', 'test', 1)", + [], + ) + .unwrap(); + drop(connection); + assert!( + commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "edges") + ) + .is_err() + ); + let connection = solstone_core_indexer_store::db::open_index(&journal).unwrap(); + let source_edges: i64 = connection + .query_row( + "SELECT COUNT(*) FROM edges WHERE src = 'source'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(source_edges, 1); + drop(connection); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let connection = solstone_core_indexer_store::db::open_index(&journal).unwrap(); + let target_edges: i64 = connection + .query_row( + "SELECT COUNT(*) FROM edges WHERE src = 'target'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(target_edges, 1); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn audit_phase_injection_rolls_back_and_retry_succeeds() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let error = commit_entity_merge_with_injector( + &journal, + "source", + "target", + EntityMergeOptions::default(), + Some(&|phase| phase == "audit"), + ) + .unwrap_err(); + let merge_id = match error { + crate::EntityMergeError::Failed { report, .. } => report.merge_id, + _ => panic!("expected injected merge failure"), + }; + let audit = journal.join("logs/entity-merges.jsonl"); + assert!( + !audit.exists() + || !fs::read_to_string(&audit) + .unwrap() + .lines() + .any( + |line| serde_json::from_str::(line).unwrap()["merge_id"] + == merge_id + ) + ); + assert!(journal.join("entities/source").exists()); + let report = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!(fs::read_to_string(audit).unwrap().lines().any(|line| { + serde_json::from_str::(line).unwrap()["merge_id"] == report.merge_id + })); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_refuses_blocked_source() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"source","aka":[],"emails":[],"blocked":true}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"target","aka":[],"emails":[]}), + None, + ) + .unwrap(); + assert_eq!( + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()) + .unwrap_err() + .to_string(), + "Cannot merge blocked entity: source" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_refuses_blocked_target() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"source","aka":[],"emails":[]}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"target","aka":[],"emails":[],"blocked":true}), + None, + ) + .unwrap(); + assert_eq!( + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()) + .unwrap_err() + .to_string(), + "Cannot merge blocked entity: target" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_refuses_same_entity() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"source","aka":[],"emails":[]}), + None, + ) + .unwrap(); + assert_eq!( + commit_entity_merge(&journal, "source", "source", EntityMergeOptions::default()) + .unwrap_err() + .to_string(), + "Source and target must be different entities." + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_refuses_two_principal_entities() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[],"is_principal":true}), + None, + ) + .unwrap(); + } + assert_eq!( + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()) + .unwrap_err() + .to_string(), + "Cannot merge two principal entities." + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn committed_merge_arms_the_restore_guard() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let events = read_visible_history(&journal, "target").unwrap(); + let merge = events + .iter() + .find(|event| event.value()["kind"] == "merge") + .unwrap(); + assert_eq!( + guard_restore_does_not_cross_merge(merge, &events) + .unwrap_err() + .to_string(), + "generic identity restore cannot target a recorded merge event; use recorded-merge undo instead" + ); + let earlier = events + .iter() + .find(|event| event.value()["kind"] != "merge") + .unwrap(); + assert_eq!( + guard_restore_does_not_cross_merge(earlier, &events) + .unwrap_err() + .to_string(), + "generic identity restore cannot cross a recorded merge event; use recorded-merge undo instead" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_transfers_principal_from_source_to_target() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"source","aka":[],"emails":[],"is_principal":true}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"target","aka":[],"emails":[]}), + None, + ) + .unwrap(); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + read_entity_identity(&journal, "target") + .unwrap() + .unwrap() + .value()["is_principal"], + true + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_fills_missing_target_scalars_from_source() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"source","aka":[],"emails":[],"title":"Engineer"}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"target","aka":[],"emails":[]}), + None, + ) + .unwrap(); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + read_entity_identity(&journal, "target") + .unwrap() + .unwrap() + .value()["title"], + "Engineer" + ); + fs::remove_dir_all(journal).unwrap(); + + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"source","aka":[],"emails":[],"title":"Engineer"}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"target","aka":[],"emails":[],"title":"Director"}), + None, + ) + .unwrap(); + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert_eq!( + read_entity_identity(&journal, "target") + .unwrap() + .unwrap() + .value()["title"], + "Director" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_never_overwrites_target_name() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"Source Name","aka":[],"emails":[]}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"Target Name","aka":[],"emails":[]}), + None, + ) + .unwrap(); + commit_entity_merge( + &journal, + "source", + "target", + EntityMergeOptions { + keep_source_as_aka: true, + }, + ) + .unwrap(); + let target = read_entity_identity(&journal, "target").unwrap().unwrap(); + assert_eq!(target.value()["name"], "Target Name"); + assert!( + target.value()["aka"] + .as_array() + .unwrap() + .contains(&json!("Source Name")) + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn merge_refuses_when_third_entity_names_source_in_aka() { + let journal = voiceprint_journal(); + save_entity_identity( + &journal, + "source", + &json!({"id":"source","name":"Source Name","aka":[],"emails":[]}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "target", + &json!({"id":"target","name":"Target Name","aka":[],"emails":[]}), + None, + ) + .unwrap(); + save_entity_identity( + &journal, + "third", + &json!({"id":"third","name":"Third","aka":["source"],"emails":[]}), + None, + ) + .unwrap(); + assert_eq!( + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()) + .unwrap_err() + .to_string(), + "Cannot merge 'source': referenced in aka lists of entity ids: third" + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn preview_does_not_take_trust_lock_or_touch_index() { + let journal = voiceprint_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + let index = journal.join("indexer/journal.sqlite"); + assert!(!index.exists()); + preview_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + assert!(!index.exists()); + drop(hold_entity_trust_lock(&journal).unwrap()); + fs::remove_dir_all(journal).unwrap(); +} diff --git a/core/crates/solstone-core-entity/src/store/merge.rs b/core/crates/solstone-core-entity/src/store/merge.rs new file mode 100644 index 000000000..251e1874b --- /dev/null +++ b/core/crates/solstone-core-entity/src/store/merge.rs @@ -0,0 +1,1049 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::collections::HashSet; +use std::error::Error; +use std::fmt; +use std::path::Path; +use std::time::{SystemTime, UNIX_EPOCH}; + +use serde_json::{Value, json}; +use solstone_core_journal_io::AtomicWriteOptions; +use solstone_core_journal_io::DirEntryKind; +use solstone_core_journal_io::JournalSnapshot; +use solstone_core_journal_io::JsonWriteOptions; +use solstone_core_journal_io::MalformedPolicy; +use solstone_core_journal_io::PathOrDay; +use solstone_core_journal_io::SnapshotError; +use solstone_core_journal_io::append_jsonl; +use solstone_core_journal_io::atomic_replace; +use solstone_core_journal_io::contained_path; +use solstone_core_journal_io::day_dirs; +use solstone_core_journal_io::iter_segments; +use solstone_core_journal_io::list_dir_entries; +use solstone_core_journal_io::path_lexists; +use solstone_core_journal_io::read_bytes; +use solstone_core_journal_io::read_json; +use solstone_core_journal_io::read_jsonl; +use solstone_core_journal_io::restore_snapshot; +use solstone_core_journal_io::write_json; +use solstone_core_journal_io::write_jsonl; + +use crate::{ + EntityOperationContext, EntityOperationKind, EntityStoreError, EntityWriteError, + hold_entity_trust_lock, read_entity_identity, save_entity_identity, +}; + +use super::merge_payload::{ + MergePayloadError, list_entity_merge_payload_ids, move_entity_merge_payload, + record_entity_merge_payload, +}; +use super::merge_rollback::MergeRollback; +use super::voiceprints::{VoiceprintArchive, read_voiceprints_npz, write_voiceprints_npz}; + +const PHASES: [&str; 11] = [ + "private_payload", + "voiceprints", + "facets", + "segments", + "activities", + "lineage", + "cleanup", + "observation relation remap", + "history", + "edges", + "audit", +]; +type VoiceprintKey = (Option, Option, Option, Option); + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub struct EntityMergeOptions { + pub keep_source_as_aka: bool, +} + +#[derive(Debug, Clone, PartialEq)] +pub struct EntityMergePreview { + pub source_id: String, + pub target_id: String, + pub target_identity: Value, + pub aliases_added: usize, + pub emails_added: usize, +} + +#[derive(Debug, Clone, PartialEq)] +pub struct EntityMergeReport { + pub merge_id: String, + pub source_id: String, + pub target_id: String, + pub completed_phases: Vec, + pub aliases_added: usize, + pub emails_added: usize, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub(crate) struct VoiceprintMergeStats { + pub added: usize, + pub skipped_duplicate: usize, + pub target_total: usize, +} +#[derive(Debug, Clone, PartialEq, Eq, Default)] +pub(crate) struct FacetMergeStats { + pub moved_count: usize, + pub merged_count: usize, + pub touched_facets: Vec, +} +#[allow(dead_code)] +#[derive(Debug, Default)] +struct MergeStats { + voiceprints_added: usize, + voiceprints_skipped_duplicate: usize, + voiceprints_target_total: usize, + facets_moved: usize, + facets_merged: usize, + facets_observations_appended: usize, + segments_labels_rewritten: usize, + segments_corrections_rewritten: usize, + segments_files_scanned: usize, + activities_records_rewritten: usize, + activities_fields_rewritten: usize, + activities_files_scanned: usize, + activities_files_rewritten: usize, + observation_relations_rewritten: usize, + edges_rows_folded: usize, + edges_self_edges_dropped: usize, +} +fn audit_counts( + stats: &MergeStats, + akas_added: usize, + emails_added: usize, + principal_transferred: bool, +) -> Value { + json!({"identity":{"akas_added":akas_added,"emails_added":emails_added,"principal_transferred":principal_transferred},"voiceprints":{"added":stats.voiceprints_added,"skipped_duplicate":stats.voiceprints_skipped_duplicate,"target_total":stats.voiceprints_target_total},"facets":{"moved":stats.facets_moved,"merged":stats.facets_merged,"observations_appended":stats.facets_observations_appended,"observation_relations_rewritten":stats.observation_relations_rewritten},"segments":{"labels_rewritten":stats.segments_labels_rewritten,"corrections_rewritten":stats.segments_corrections_rewritten,"files_scanned":stats.segments_files_scanned,"errors":0},"activities":{"records_rewritten":stats.activities_records_rewritten,"fields_rewritten":stats.activities_fields_rewritten,"files_scanned":stats.activities_files_scanned,"files_rewritten":stats.activities_files_rewritten,"errors":0},"edges":{"rows_folded":stats.edges_rows_folded,"self_edges_dropped":stats.edges_self_edges_dropped,"error":null}}) +} + +#[derive(Debug)] +pub enum EntityMergeError { + Refused(String), + Read(EntityStoreError), + Write(EntityWriteError), + Payload(MergePayloadError), + Snapshot(SnapshotError), + Index(solstone_core_indexer_store::StoreError), + Audit(solstone_core_journal_io::AppendError), + Failed { + failed_phase: String, + report: Box, + rollback_error: Option, + }, +} +impl fmt::Display for EntityMergeError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Refused(message) => f.write_str(message), + Self::Read(error) => error.fmt(f), + Self::Write(error) => error.fmt(f), + Self::Payload(error) => error.fmt(f), + Self::Snapshot(error) => error.fmt(f), + Self::Index(error) => error.fmt(f), + Self::Audit(error) => error.fmt(f), + Self::Failed { failed_phase, .. } => { + write!(f, "entity merge failed during {failed_phase}") + } + } + } +} +impl Error for EntityMergeError {} +impl From for EntityMergeError { + fn from(error: EntityStoreError) -> Self { + Self::Read(error) + } +} +impl From for EntityMergeError { + fn from(error: EntityWriteError) -> Self { + Self::Write(error) + } +} +impl From for EntityMergeError { + fn from(error: MergePayloadError) -> Self { + Self::Payload(error) + } +} +impl From for EntityMergeError { + fn from(error: SnapshotError) -> Self { + Self::Snapshot(error) + } +} + +pub fn preview_entity_merge( + journal: &Path, + source_id: &str, + target_id: &str, + options: EntityMergeOptions, +) -> Result { + let plan = plan_merge(journal, source_id, target_id, options)?; + Ok(EntityMergePreview { + source_id: source_id.to_owned(), + target_id: target_id.to_owned(), + target_identity: plan.target_after, + aliases_added: plan.aliases_added, + emails_added: plan.emails_added, + }) +} + +pub fn commit_entity_merge( + journal: &Path, + source_id: &str, + target_id: &str, + options: EntityMergeOptions, +) -> Result { + commit_entity_merge_with_injector(journal, source_id, target_id, options, None) +} + +pub(crate) fn commit_entity_merge_with_injector( + journal: &Path, + source_id: &str, + target_id: &str, + options: EntityMergeOptions, + injector: Option<&dyn Fn(&str) -> bool>, +) -> Result { + let _trust = hold_entity_trust_lock(journal).map_err(EntityWriteError::TrustLock)?; + let plan = plan_merge(journal, source_id, target_id, options)?; + let merge_id = format!( + "em_{:x}", + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + .as_nanos() + ); + let mut rollback = MergeRollback::default(); + for path in [ + format!("entities/{source_id}"), + format!("entities/{target_id}"), + format!("entities/{target_id}/history/private/{merge_id}.json"), + "indexer".to_owned(), + ] { + rollback.capture(journal, &path)?; + } + let mut report = EntityMergeReport { + merge_id: merge_id.clone(), + source_id: source_id.to_owned(), + target_id: target_id.to_owned(), + completed_phases: Vec::new(), + aliases_added: plan.aliases_added, + emails_added: plan.emails_added, + }; + let mut payload = payload_for_merge(&merge_id, source_id, target_id, &plan); + let mut touched_facets = Vec::new(); + let mut stats = MergeStats::default(); + for phase in PHASES { + let result = match phase { + "private_payload" => record_entity_merge_payload(journal, target_id, &merge_id, &payload).map(|_| ()).map_err(Into::into), + "voiceprints" => merge_voiceprints(journal, source_id, target_id).map(|result| { stats.voiceprints_added=result.added; stats.voiceprints_skipped_duplicate=result.skipped_duplicate; stats.voiceprints_target_total=result.target_total; }), + "facets" => merge_facets(journal, source_id, target_id, Some(&mut rollback)).map(|result| { stats.facets_moved=result.moved_count; stats.facets_merged=result.merged_count; touched_facets = result.touched_facets; }), + "history" => save_entity_identity(journal, target_id, &plan.target_after, Some(&EntityOperationContext { kind: EntityOperationKind::Merge, caller: Value::Null, actor: Value::Null, metadata: json!({"merge_id": merge_id, "source_id": source_id, "target_id": target_id}) })).map_err(EntityMergeError::Write).and_then(|saved| { let sequence = saved.event.and_then(|event| event.get("seq").cloned()).ok_or_else(|| EntityMergeError::Refused("merge history event was not written".to_owned()))?; payload["commit_seq"] = sequence; record_entity_merge_payload(journal, target_id, &merge_id, &payload)?; Ok(()) }), + "cleanup" => cleanup_merge(journal, source_id, &touched_facets, Some(&mut rollback)), + "edges" => solstone_core_indexer_store::merge::fold_entity_edges_for_recorded_merge(journal, source_id, target_id).map_err(EntityMergeError::Index).and_then(|result| { stats.edges_rows_folded=result.rows_folded; stats.edges_self_edges_dropped=result.self_edges_dropped; payload["result_counts"] = audit_counts(&stats, plan.aliases_added, plan.emails_added, plan.principal_transferred); record_entity_merge_payload(journal, target_id, &merge_id, &payload)?; Ok(()) }), + "audit" => { let path = contained_path(journal, "logs/entity-merges.jsonl").map_err(|error| EntityMergeError::Refused(error.to_string()))?; rollback.capture(journal, "logs/entity-merges.jsonl")?; let ts = u64::try_from(SystemTime::now().duration_since(UNIX_EPOCH).map_err(|error| EntityMergeError::Refused(error.to_string()))?.as_millis()).map_err(|_| EntityMergeError::Refused("merge audit timestamp exceeds u64".to_owned()))?; append_jsonl(path, &json!({"ts":ts,"merge_id":merge_id,"source_id":source_id,"source_display_name":plan.source_display_name,"target_id":target_id,"target_display_name":plan.target_display_name,"principal_transferred":plan.principal_transferred,"counts":audit_counts(&stats, plan.aliases_added, plan.emails_added, plan.principal_transferred),"caller":Value::Null})).map_err(EntityMergeError::Audit) }, + "segments" => merge_segment_labels(journal, source_id, target_id, Some(&mut rollback)).map(|result| { stats.segments_labels_rewritten=result.labels_rewritten; stats.segments_corrections_rewritten=result.corrections_rewritten; stats.segments_files_scanned=result.files_scanned; }), + "activities" => merge_activities(journal, source_id, target_id, Some(&mut rollback)).map(|result| { stats.activities_records_rewritten=result.records_rewritten; stats.activities_fields_rewritten=result.fields_rewritten; stats.activities_files_scanned=result.files_scanned; stats.activities_files_rewritten=result.files_rewritten; }), + "observation relation remap" => merge_observation_relations(journal, source_id, target_id, Some(&mut rollback)).map(|result| { stats.observation_relations_rewritten=result.rows_rewritten; }), + "lineage" => rebase_lineage(journal, source_id, target_id, &plan.target_after).and_then(|ids| { if !ids.is_empty() { payload["manifest"]["rebased_merge_ids"] = Value::Array(ids.into_iter().map(Value::String).collect()); record_entity_merge_payload(journal, target_id, &merge_id, &payload)?; } Ok(()) }), + _ => unreachable!("merge phase list is fixed"), + }; + if let Err(error) = result { + let rollback_error = rollback + .restore(journal) + .err() + .map(|rollback| rollback.to_string()); + return Err(EntityMergeError::Failed { + failed_phase: phase.to_owned(), + report: Box::new(report), + rollback_error: rollback_error.or_else(|| Some(error.to_string())), + }); + } + if injector.is_some_and(|injector| injector(phase)) { + let rollback_error = rollback + .restore(journal) + .err() + .map(|error| error.to_string()); + return Err(EntityMergeError::Failed { + failed_phase: phase.to_owned(), + report: Box::new(report), + rollback_error: rollback_error + .or_else(|| Some(format!("injected failure after phase {phase}"))), + }); + } + report.completed_phases.push(phase.to_owned()); + } + Ok(report) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub(crate) struct ObservationRelationMergeStats { + pub rows_rewritten: usize, +} +pub(crate) fn merge_observation_relations( + journal: &Path, + source_id: &str, + target_id: &str, + mut rollback: Option<&mut MergeRollback>, +) -> Result { + let facets = contained_path(journal, "facets") + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut stats = ObservationRelationMergeStats::default(); + for facet in + list_dir_entries(&facets).map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + if facet.kind != DirEntryKind::Directory { + continue; + } + let entities = facet.path.join("entities"); + for entity in list_dir_entries(&entities) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + if entity.kind != DirEntryKind::Directory { + continue; + } + let path = entity.path.join("observations.jsonl"); + if !path_lexists(&path).map_err(|error| EntityMergeError::Refused(error.to_string()))? { + continue; + } + let mut rows: Vec = read_jsonl(&path, Vec::new(), MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut changed = false; + for row in &mut rows { + if let Some(relation) = row.get_mut("relation").and_then(Value::as_object_mut) + && relation.get("target_entity_id").and_then(Value::as_str) == Some(source_id) + { + relation.insert( + "target_entity_id".to_owned(), + Value::String(target_id.to_owned()), + ); + stats.rows_rewritten += 1; + changed = true; + } + } + if changed { + capture_rollback_file(&mut rollback, journal, &path)?; + write_jsonl(&path, rows, AtomicWriteOptions::default()) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + } + } + } + Ok(stats) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub(crate) struct ActivityMergeStats { + pub files_scanned: usize, + pub files_rewritten: usize, + pub records_rewritten: usize, + pub fields_rewritten: usize, +} +pub(crate) fn merge_activities( + journal: &Path, + source_id: &str, + target_id: &str, + mut rollback: Option<&mut MergeRollback>, +) -> Result { + let facets = contained_path(journal, "facets") + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut stats = ActivityMergeStats::default(); + for facet in + list_dir_entries(&facets).map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + if facet.kind != DirEntryKind::Directory { + continue; + } + let activities = facet.path.join("activities"); + for file in list_dir_entries(&activities) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + if file.kind != DirEntryKind::File + || file.path.extension().and_then(|value| value.to_str()) != Some("jsonl") + { + continue; + } + stats.files_scanned += 1; + let mut rows: Vec = + read_jsonl(&file.path, Vec::new(), MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut file_changed = false; + for row in &mut rows { + let mut changed = false; + if let Some(object) = row.as_object_mut() { + if let Some(active) = object + .get_mut("active_entities") + .and_then(Value::as_array_mut) + { + for value in active { + if value.as_str() == Some(source_id) { + *value = Value::String(target_id.to_owned()); + stats.fields_rewritten += 1; + changed = true; + } + } + } + for (field, keys) in [ + ("participation", &["entity_id"][..]), + ( + "commitments", + &["owner_entity_id", "counterparty_entity_id"][..], + ), + ( + "closures", + &["owner_entity_id", "counterparty_entity_id"][..], + ), + ( + "decisions", + &["owner_entity_id", "counterparty_entity_id"][..], + ), + ("relations", &["from_entity_id", "to_entity_id"][..]), + ] { + if let Some(items) = object.get_mut(field).and_then(Value::as_array_mut) { + for item in items { + if let Some(item) = item.as_object_mut() { + for key in keys { + if item.get(*key).and_then(Value::as_str) == Some(source_id) + { + item.insert( + (*key).to_owned(), + Value::String(target_id.to_owned()), + ); + stats.fields_rewritten += 1; + changed = true; + } + } + } + } + } + } + } + if changed { + stats.records_rewritten += 1; + file_changed = true; + } + } + if file_changed { + capture_rollback_file(&mut rollback, journal, &file.path)?; + write_jsonl(&file.path, rows, AtomicWriteOptions::default()) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + stats.files_rewritten += 1; + } + } + } + Ok(stats) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub(crate) struct SegmentMergeStats { + pub files_scanned: usize, + pub labels_rewritten: usize, + pub corrections_rewritten: usize, +} + +pub(crate) fn merge_segment_labels( + journal: &Path, + source_id: &str, + target_id: &str, + mut rollback: Option<&mut MergeRollback>, +) -> Result { + let mut stats = SegmentMergeStats::default(); + for day in day_dirs(journal) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + .into_values() + { + for segment in iter_segments(journal, PathOrDay::Directory(&day)) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + let path = segment.path.join("talents/speaker_labels.json"); + if path_lexists(&path).map_err(|error| EntityMergeError::Refused(error.to_string()))? { + stats.files_scanned += 1; + let raw = read_bytes(&path, Vec::new()) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if !raw + .windows(source_id.len()) + .any(|bytes| bytes == source_id.as_bytes()) + { + continue; + } + let mut value: Value = read_json(&path, Value::Null, MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut changed = false; + if let Some(labels) = value.get_mut("labels").and_then(Value::as_array_mut) { + for label in labels { + if let Some(object) = label.as_object_mut() + && object.get("speaker").and_then(Value::as_str) == Some(source_id) + { + object + .insert("speaker".to_owned(), Value::String(target_id.to_owned())); + changed = true; + } + } + } + if changed { + capture_rollback_file(&mut rollback, journal, &path)?; + write_json( + &path, + &value, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + stats.labels_rewritten += 1; + } + } + let path = segment.path.join("talents/speaker_corrections.json"); + if !path_lexists(&path).map_err(|error| EntityMergeError::Refused(error.to_string()))? { + continue; + } + stats.files_scanned += 1; + let raw = read_bytes(&path, Vec::new()) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if !raw + .windows(source_id.len()) + .any(|bytes| bytes == source_id.as_bytes()) + { + continue; + } + let mut value: Value = read_json(&path, Value::Null, MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut changed = false; + if let Some(corrections) = value.get_mut("corrections").and_then(Value::as_array_mut) { + for correction in corrections { + if let Some(object) = correction.as_object_mut() { + for field in ["original_speaker", "corrected_speaker"] { + if object.get(field).and_then(Value::as_str) == Some(source_id) { + object + .insert(field.to_owned(), Value::String(target_id.to_owned())); + changed = true; + } + } + } + } + } + if changed { + capture_rollback_file(&mut rollback, journal, &path)?; + write_json( + &path, + &value, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + stats.corrections_rewritten += 1; + } + } + } + Ok(stats) +} + +fn capture_rollback_file( + rollback: &mut Option<&mut MergeRollback>, + journal: &Path, + path: &Path, +) -> Result<(), EntityMergeError> { + let Some(rollback) = rollback.as_deref_mut() else { + return Ok(()); + }; + let relative = path + .strip_prefix(journal) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + .to_str() + .ok_or_else(|| { + EntityMergeError::Refused(format!("path is not UTF-8: {}", path.display())) + })?; + rollback.capture(journal, relative)?; + Ok(()) +} + +pub(crate) fn merge_facets( + journal: &Path, + source_id: &str, + target_id: &str, + mut rollback: Option<&mut MergeRollback>, +) -> Result { + let facets = contained_path(journal, "facets") + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut stats = FacetMergeStats::default(); + for entry in + list_dir_entries(&facets).map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + if entry.kind != DirEntryKind::Directory { + continue; + } + let facet = entry.name.to_string_lossy(); + let source_rel = contained_path( + journal, + &format!("facets/{facet}/entities/{source_id}/entity.json"), + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if !path_lexists(&source_rel) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + continue; + } + let target_directory = format!("facets/{facet}/entities/{target_id}"); + let target_rel = contained_path(journal, &format!("{target_directory}/entity.json")) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if let Some(rollback) = rollback.as_deref_mut() { + rollback.capture(journal, &target_directory)?; + } + let source: Value = read_json(&source_rel, Value::Null, MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let source_obs_path = source_rel.parent().unwrap().join("observations.jsonl"); + let source_obs: Vec = + read_jsonl(&source_obs_path, Vec::new(), MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if !path_lexists(&target_rel) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + let mut moved = source; + moved + .as_object_mut() + .ok_or_else(|| { + EntityMergeError::Refused("facet relationship is not an object".to_owned()) + })? + .insert("entity_id".to_owned(), Value::String(target_id.to_owned())); + write_json( + &target_rel, + &moved, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if path_lexists(&source_obs_path) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + write_jsonl( + target_rel.parent().unwrap().join("observations.jsonl"), + source_obs, + AtomicWriteOptions::default(), + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + } + stats.moved_count += 1; + stats.touched_facets.push(facet.into_owned()); + } else { + let mut target: Value = read_json(&target_rel, Value::Null, MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let target_obs_path = target_rel.parent().unwrap().join("observations.jsonl"); + let target_obs: Vec = + read_jsonl(&target_obs_path, Vec::new(), MalformedPolicy::Raise) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + merge_facet_scalars(&source, &mut target); + write_json( + &target_rel, + &target, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + write_jsonl( + target_obs_path, + dedupe_observations(&source_obs, &target_obs), + AtomicWriteOptions::default(), + ) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + stats.merged_count += 1; + stats.touched_facets.push(facet.into_owned()); + } + } + Ok(stats) +} +fn cleanup_merge( + journal: &Path, + source_id: &str, + facets: &[String], + mut rollback: Option<&mut MergeRollback>, +) -> Result<(), EntityMergeError> { + for facet in facets { + let path = format!("facets/{facet}/entities/{source_id}"); + if let Some(rollback) = rollback.as_deref_mut() { + rollback.capture(journal, &path)?; + } + restore_snapshot(journal, &JournalSnapshot::Missing { path })?; + } + let discovery = contained_path(journal, "awareness/discovery_clusters.json") + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if path_lexists(&discovery).map_err(|error| EntityMergeError::Refused(error.to_string()))? { + if let Some(rollback) = rollback { + rollback.capture(journal, "awareness/discovery_clusters.json")?; + } + restore_snapshot( + journal, + &JournalSnapshot::Missing { + path: "awareness/discovery_clusters.json".to_owned(), + }, + )?; + } + restore_snapshot( + journal, + &JournalSnapshot::Missing { + path: format!("entities/{source_id}"), + }, + ) + .map_err(Into::into) +} + +fn rebase_lineage( + journal: &Path, + source_id: &str, + target_id: &str, + target: &Value, +) -> Result, EntityMergeError> { + let mut rebased = Vec::new(); + for merge_id in list_entity_merge_payload_ids(journal, source_id)? { + let (payload, private_payload) = + move_entity_merge_payload(journal, source_id, target_id, &merge_id, Some(source_id))?; + let descendant_source = payload + .get("source_id") + .and_then(Value::as_str) + .unwrap_or_default(); + save_entity_identity( + journal, + target_id, + target, + Some(&EntityOperationContext { + kind: EntityOperationKind::Merge, + caller: Value::Null, + actor: Value::Null, + metadata: json!({"merge_id":merge_id,"source_id":descendant_source,"target_id":target_id,"rebased_from_entity_id":source_id,"private_payload":private_payload}), + }), + )?; + rebased.push(merge_id); + } + Ok(rebased) +} +fn merge_facet_scalars(source: &Value, target: &mut Value) { + let source = source.as_object().expect("facet relationship object"); + let target = target.as_object_mut().expect("facet relationship object"); + for (field, earlier) in [ + ("attached_at", true), + ("updated_at", false), + ("last_seen", false), + ] { + if let Some(value) = source.get(field).filter(|value| !is_blank(Some(value))) { + let replace = is_blank(target.get(field)) + || target.get(field).is_some_and(|existing| { + if earlier { + value.as_str() < existing.as_str() + } else { + value.as_str() > existing.as_str() + } + }); + if replace { + target.insert(field.to_owned(), value.clone()); + } + } + } + if is_blank(target.get("description")) && !is_blank(source.get("description")) { + target.insert("description".to_owned(), source["description"].clone()); + } +} + +pub(crate) fn merge_voiceprints( + journal: &Path, + source_id: &str, + target_id: &str, +) -> Result { + let source_path = contained_path(journal, &format!("entities/{source_id}/voiceprints.npz")) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let target_path = contained_path(journal, &format!("entities/{target_id}/voiceprints.npz")) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + if !path_lexists(&source_path).map_err(|error| EntityMergeError::Refused(error.to_string()))? { + let target_total = load_voiceprints(&target_path)?.map_or(0, |archive| archive.rows); + return Ok(VoiceprintMergeStats { + target_total, + ..VoiceprintMergeStats::default() + }); + } + let source = load_voiceprints(&source_path)?.expect("existing source voiceprints"); + let mut target = load_voiceprints(&target_path)?.unwrap_or(VoiceprintArchive { + embeddings: Vec::new(), + rows: 0, + metadata: Vec::new(), + }); + let mut existing = target + .metadata + .iter() + .map(|metadata| voiceprint_key(metadata)) + .collect::, _>>()?; + let mut stats = VoiceprintMergeStats::default(); + for (embedding, metadata) in source.embeddings.chunks_exact(256).zip(&source.metadata) { + let key = voiceprint_key(metadata)?; + if !existing.insert(key) { + stats.skipped_duplicate += 1; + continue; + } + let norm = embedding + .iter() + .map(|value| value * value) + .sum::() + .sqrt(); + if norm > 0.0 { + target + .embeddings + .extend(embedding.iter().map(|value| value / norm)); + target.metadata.push(metadata.clone()); + stats.added += 1; + } + } + target.rows = target.metadata.len(); + stats.target_total = target.rows; + if stats.added > 0 { + let bytes = write_voiceprints_npz(&target.embeddings, &target.metadata) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + atomic_replace(&target_path, &bytes, AtomicWriteOptions::default()) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + } + Ok(stats) +} + +fn load_voiceprints(path: &Path) -> Result, EntityMergeError> { + if !path_lexists(path).map_err(|error| EntityMergeError::Refused(error.to_string()))? { + return Ok(None); + } + let bytes = read_bytes(path, Vec::new()) + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + read_voiceprints_npz(&bytes) + .map(Some) + .map_err(|error| EntityMergeError::Refused(error.to_string())) +} + +fn voiceprint_key(metadata: &str) -> Result { + let object: Value = serde_json::from_str(metadata).map_err(|error| { + EntityMergeError::Refused(format!("invalid voiceprint metadata: {error}")) + })?; + Ok(( + object.get("day").cloned(), + object.get("segment_key").cloned(), + object.get("source").cloned(), + object.get("sentence_id").cloned(), + )) +} + +struct MergePlan { + target_before: Value, + target_after: Value, + aliases_added: usize, + emails_added: usize, + source_display_name: String, + target_display_name: String, + principal_transferred: bool, +} + +fn plan_merge( + journal: &Path, + source_id: &str, + target_id: &str, + options: EntityMergeOptions, +) -> Result { + if source_id == target_id { + return Err(EntityMergeError::Refused( + "Source and target must be different entities.".to_owned(), + )); + } + let source = read_entity_identity(journal, source_id)? + .ok_or_else(|| EntityMergeError::Refused(format!("Source entity not found: {source_id}")))? + .value() + .clone(); + let target = read_entity_identity(journal, target_id)? + .ok_or_else(|| EntityMergeError::Refused(format!("Target entity not found: {target_id}")))? + .value() + .clone(); + if source.get("blocked").and_then(Value::as_bool) == Some(true) { + return Err(EntityMergeError::Refused(format!( + "Cannot merge blocked entity: {source_id}" + ))); + } + if target.get("blocked").and_then(Value::as_bool) == Some(true) { + return Err(EntityMergeError::Refused(format!( + "Cannot merge blocked entity: {target_id}" + ))); + } + if source.get("is_principal").and_then(Value::as_bool) == Some(true) + && target.get("is_principal").and_then(Value::as_bool) == Some(true) + { + return Err(EntityMergeError::Refused( + "Cannot merge two principal entities.".to_owned(), + )); + } + check_aka_cross_references( + journal, + source_id, + source + .get("name") + .and_then(Value::as_str) + .unwrap_or_default(), + target_id, + )?; + let mut after = target.clone(); + let object = after.as_object_mut().ok_or_else(|| { + EntityMergeError::Refused(format!("Target entity not found: {target_id}")) + })?; + let aliases_before = values(&target, "aka"); + let mut alias_values = aliases_before.clone(); + alias_values.extend(values(&source, "aka")); + if options.keep_source_as_aka + && let Some(name) = source.get("name").and_then(Value::as_str) + { + alias_values.push(name.to_owned()); + } + let aliases = dedupe_akas(&alias_values); + object.insert( + "aka".to_owned(), + Value::Array(aliases.iter().cloned().map(Value::String).collect()), + ); + let emails_before = values(&target, "emails"); + let emails = dedupe_emails(&emails_before, &values(&source, "emails")); + object.insert( + "emails".to_owned(), + Value::Array(emails.iter().cloned().map(Value::String).collect()), + ); + for (field, value) in source.as_object().expect("identity object") { + if ![ + "id", + "name", + "aka", + "emails", + "created_at", + "updated_at", + "merged_into", + "blocked", + "is_principal", + ] + .contains(&field.as_str()) + && is_blank(object.get(field)) + && !is_blank(Some(value)) + { + object.insert(field.clone(), value.clone()); + } + } + if source.get("is_principal").and_then(Value::as_bool) == Some(true) { + object.insert("is_principal".to_owned(), Value::Bool(true)); + } + let source_display_name = source + .get("name") + .and_then(Value::as_str) + .unwrap_or(source_id) + .to_owned(); + let target_display_name = after + .get("name") + .and_then(Value::as_str) + .unwrap_or(target_id) + .to_owned(); + Ok(MergePlan { + target_before: target, + target_after: after, + aliases_added: aliases.len().saturating_sub(aliases_before.len()), + emails_added: emails.len().saturating_sub(emails_before.len()), + source_display_name, + target_display_name, + principal_transferred: source.get("is_principal").and_then(Value::as_bool) == Some(true), + }) +} + +pub(crate) fn dedupe_akas(values: &[String]) -> Vec { + let mut seen = HashSet::new(); + let mut output = Vec::new(); + for value in values { + if seen.insert(value.to_lowercase()) { + output.push(value.clone()); + } + } + output.sort_by_key(|value| value.to_lowercase()); + output +} +pub(crate) fn dedupe_emails(target_values: &[String], source_values: &[String]) -> Vec { + let mut seen = HashSet::new(); + target_values + .iter() + .chain(source_values) + .filter(|value| seen.insert(value.to_lowercase())) + .cloned() + .collect() +} +#[allow(dead_code)] +pub(crate) fn dedupe_observations(source: &[Value], target: &[Value]) -> Vec { + let mut seen = HashSet::new(); + let mut result = Vec::new(); + for value in target.iter().chain(source) { + let key = ( + value + .get("content") + .and_then(Value::as_str) + .unwrap_or_default() + .to_owned(), + value.get("observed_at").cloned().unwrap_or(Value::Null), + ); + if seen.insert(key) { + result.push(value.clone()); + } + } + result +} +fn values(value: &Value, field: &str) -> Vec { + value + .get(field) + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(Value::as_str) + .map(str::to_owned) + .collect() +} +fn is_blank(value: Option<&Value>) -> bool { + value.is_none_or(|value| value.is_null() || value.as_str().is_some_and(str::is_empty)) +} +fn check_aka_cross_references( + journal: &Path, + source_id: &str, + source_name: &str, + target_id: &str, +) -> Result<(), EntityMergeError> { + let directory = contained_path(journal, "entities") + .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + let mut ids = Vec::new(); + for entry in list_dir_entries(&directory) + .map_err(|error| EntityMergeError::Refused(error.to_string()))? + { + if entry.kind != DirEntryKind::Directory { + continue; + } + let id = entry.name.to_string_lossy(); + if id == source_id || id == target_id { + continue; + } + if let Some(identity) = read_entity_identity(journal, &id)? + && values(identity.value(), "aka") + .iter() + .any(|aka| aka == source_id || aka == source_name) + { + ids.push(id.into_owned()); + } + } + if ids.is_empty() { + Ok(()) + } else { + Err(EntityMergeError::Refused(format!( + "Cannot merge '{source_id}': referenced in aka lists of entity ids: {}", + ids.join(", ") + ))) + } +} +fn payload_for_merge(merge_id: &str, source_id: &str, target_id: &str, plan: &MergePlan) -> Value { + json!({"schema_version":1,"merge_id":merge_id,"source_id":source_id,"target_id":target_id,"commit_seq":null,"source_state":{"snapshots":[{"rel":format!("entities/{source_id}"),"files":[]}]},"result_counts":{},"manifest":{"identity":{"target_before":plan.target_before,"aka_support":[],"email_support":[],"scalar_support":[]},"voiceprints":{"support":[]},"facets":{"entries":[]},"segments":{"entries":[]},"activities":{"entries":[]},"observation_relations":{"entries":[]},"rebased_merge_ids":[]}}) +} diff --git a/core/crates/solstone-core-entity/src/store/mod.rs b/core/crates/solstone-core-entity/src/store/mod.rs index 10b904e18..ce3fded99 100644 --- a/core/crates/solstone-core-entity/src/store/mod.rs +++ b/core/crates/solstone-core-entity/src/store/mod.rs @@ -8,6 +8,9 @@ mod error; mod history; mod identity; mod map; +pub(crate) mod merge; +pub(crate) mod merge_payload; +mod merge_rollback; mod paths; mod reconcile; mod repair; @@ -29,6 +32,10 @@ pub use map::{ EntityIdentityGroupMap, EntityIdentityMap, IdentityMapLoser, IdentityMapLoserReason, read_identity_group_map, read_identity_map, }; +pub use merge::{ + EntityMergeError, EntityMergeOptions, EntityMergePreview, EntityMergeReport, + commit_entity_merge, preview_entity_merge, +}; pub use reconcile::{PreparedHistoryOutcome, classify_prepared_history}; pub use repair::{ EntityIdentityRepairError, EntityIdentityRepairGuard, EntityIdentityRepairRefusal,