diff --git a/core/crates/solstone-core-entity/src/store/merge.rs b/core/crates/solstone-core-entity/src/store/merge.rs index 962018685..7af91baf8 100644 --- a/core/crates/solstone-core-entity/src/store/merge.rs +++ b/core/crates/solstone-core-entity/src/store/merge.rs @@ -1177,6 +1177,10 @@ fn payload_for_merge( journal, &format!("entities/{source_id}"), )?]; + let target_voiceprints = snapshot_payload(&capture_snapshot( + journal, + &format!("entities/{target_id}/voiceprints.npz"), + )?); let facets = contained_path(journal, "facets") .map_err(|error| EntityMergeError::Refused(error.to_string()))?; for entry in @@ -1194,7 +1198,7 @@ fn payload_for_merge( } } Ok( - json!({"schema_version":1,"merge_id":merge_id,"source_id":source_id,"target_id":target_id,"commit_seq":null,"source_state":{"identity":plan.source_before,"snapshots":snapshots},"result_counts":{},"manifest":{"identity":{"target_before":plan.target_before,"aka_support":plan.aka_support,"email_support":plan.email_support,"scalar_support":plan.scalar_support},"voiceprints":{"support":[]},"facets":{"entries":[]},"segments":{"entries":[]},"activities":{"entries":[]},"observation_relations":{"entries":[]},"rebased_merge_ids":[]}}), + json!({"schema_version":1,"merge_id":merge_id,"source_id":source_id,"target_id":target_id,"commit_seq":null,"source_state":{"identity":plan.source_before,"snapshots":snapshots},"result_counts":{},"manifest":{"identity":{"target_before":plan.target_before,"aka_support":plan.aka_support,"email_support":plan.email_support,"scalar_support":plan.scalar_support},"voiceprints":{"support":[],"target_before":target_voiceprints},"facets":{"entries":[]},"segments":{"entries":[]},"activities":{"entries":[]},"observation_relations":{"entries":[]},"rebased_merge_ids":[]}}), ) } diff --git a/core/crates/solstone-core-entity/src/store/undo.rs b/core/crates/solstone-core-entity/src/store/undo.rs index 04b6daabf..86427faca 100644 --- a/core/crates/solstone-core-entity/src/store/undo.rs +++ b/core/crates/solstone-core-entity/src/store/undo.rs @@ -142,6 +142,8 @@ pub(crate) fn undo_entity_merge_with_injector( for snapshot in &source_snapshots { restore_snapshot(journal, snapshot)?; } + phase = "voiceprints"; + undo_voiceprints(journal, &target_id, &payload, &mut rollback)?; phase = "segments"; undo_segments(journal, &payload, &mut rollback)?; inject_failure(injector, phase)?; @@ -666,6 +668,31 @@ fn undo_observation_relations( Ok(()) } +fn undo_voiceprints( + journal: &Path, + target_id: &str, + payload: &Value, + rollback: &mut MergeRollback, +) -> Result<(), EntityUndoError> { + let snapshot = payload + .get("manifest") + .and_then(|manifest| manifest.get("voiceprints")) + .and_then(|voiceprints| voiceprints.get("target_before")) + .ok_or_else(|| { + EntityUndoError::Refused("merge payload voiceprints missing target_before".to_owned()) + })?; + let snapshot = snapshot_from_payload(snapshot)?; + let path = format!("entities/{target_id}/voiceprints.npz"); + if snapshot_path(&snapshot) != path { + return Err(EntityUndoError::Refused( + "merge payload voiceprints snapshot path does not match target".to_owned(), + )); + } + rollback.capture(journal, &path)?; + restore_snapshot(journal, &snapshot)?; + Ok(()) +} + fn undo_segments( journal: &Path, payload: &Value, diff --git a/core/crates/solstone-core-entity/src/undo_tests.rs b/core/crates/solstone-core-entity/src/undo_tests.rs index e14df961a..5cbe118a9 100644 --- a/core/crates/solstone-core-entity/src/undo_tests.rs +++ b/core/crates/solstone-core-entity/src/undo_tests.rs @@ -8,8 +8,10 @@ use std::path::PathBuf; use std::sync::atomic::{AtomicU64, Ordering}; use serde_json::{Value, json}; +use solstone_core_journal_io::{AtomicWriteOptions, JsonWriteOptions, write_json, write_jsonl}; use super::store::undo_entity_merge_with_injector; +use super::store::voiceprints::write_voiceprints_npz; use crate::{ EntityMergeOptions, commit_entity_merge, guard_restore_does_not_cross_merge, read_entity_identity, read_visible_history, save_entity_identity, undo_entity_merge, @@ -59,6 +61,53 @@ fn journal_tree(journal: &std::path::Path) -> Vec<(String, Vec)> { files } +fn comparable_journal_tree(journal: &std::path::Path, target_id: &str) -> Vec<(String, Vec)> { + fn excluded(relative: &str, target_id: &str) -> bool { + relative.ends_with(".lock") + || relative == "indexer/" + || relative.starts_with("indexer/") + || relative == format!("entities/{target_id}/history/") + || relative.starts_with(&format!("entities/{target_id}/history/")) + || relative == "logs/entity-merges.jsonl" + || relative == "awareness/discovery_clusters.json" + } + + fn collect( + root: &std::path::Path, + directory: &std::path::Path, + target_id: &str, + files: &mut Vec<(String, Vec)>, + ) { + let mut entries = fs::read_dir(directory) + .unwrap() + .map(|entry| entry.unwrap()) + .collect::>(); + entries.sort_by_key(|entry| entry.file_name()); + for entry in entries { + let path = entry.path(); + let relative = path + .strip_prefix(root) + .unwrap() + .to_string_lossy() + .into_owned(); + let relative_directory = format!("{relative}/"); + if excluded(&relative, target_id) || excluded(&relative_directory, target_id) { + continue; + } + if entry.file_type().unwrap().is_dir() { + files.push((relative_directory, Vec::new())); + collect(root, &path, target_id, files); + } else if entry.file_type().unwrap().is_file() { + files.push((relative, fs::read(path).unwrap())); + } + } + } + + let mut files = Vec::new(); + collect(journal, journal, target_id, &mut files); + files +} + #[test] fn undo_reverts_target_identity_restores_source_and_removes_payload() { let journal = undo_journal(); @@ -286,6 +335,93 @@ fn undo_restores_file_modes_not_directory_modes() { fs::remove_dir_all(journal).unwrap(); } +#[test] +fn undo_restores_byte_identical_journal_outside_excluded_paths() { + let journal = undo_journal(); + for id in ["source", "target", "other"] { + save_entity_identity( + &journal, + id, + &json!({"id":id,"name":id,"aka":[],"emails":[]}), + None, + ) + .unwrap(); + } + fs::create_dir_all(journal.join("logs")).unwrap(); + let discovery = journal.join("awareness/discovery_clusters.json"); + fs::create_dir_all(discovery.parent().unwrap()).unwrap(); + fs::write(&discovery, b"{\"clusters\":[]}").unwrap(); + + let source_facet = journal.join("facets/work/entities/source"); + fs::create_dir_all(&source_facet).unwrap(); + write_json( + source_facet.join("entity.json"), + &json!({"entity_id":"source","description":"source relationship"}), + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .unwrap(); + let labels = journal.join("chronicle/20260102/080000_300/talents/speaker_labels.json"); + fs::create_dir_all(labels.parent().unwrap()).unwrap(); + write_json( + &labels, + &json!({"labels":[{"speaker":"source"}]}), + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .unwrap(); + let activity = journal.join("facets/work/activities/20260102.jsonl"); + fs::create_dir_all(activity.parent().unwrap()).unwrap(); + write_jsonl( + &activity, + vec![json!({"id":"activity","active_entities":["source"]})], + AtomicWriteOptions::default(), + ) + .unwrap(); + let observations = journal.join("facets/work/entities/other/observations.jsonl"); + fs::create_dir_all(observations.parent().unwrap()).unwrap(); + write_jsonl( + &observations, + vec![json!({"observed_at":1,"relation":{"kind":"works-with","target_entity_id":"source","target_name":"source"}})], + AtomicWriteOptions::default(), + ) + .unwrap(); + let source_voiceprints = journal.join("entities/source/voiceprints.npz"); + fs::write( + &source_voiceprints, + write_voiceprints_npz( + &[2.0; 256], + &[ + "{\"day\":\"d\",\"segment_key\":\"s\",\"source\":\"x\",\"sentence_id\":\"1\"}" + .to_owned(), + ], + ) + .unwrap(), + ) + .unwrap(); + + solstone_core_indexer_store::scan::rebuild_edges(&journal).unwrap(); + let index_before = solstone_core_indexer_store::merge::fingerprint_edge_rows(&journal).unwrap(); + let tree_before = comparable_journal_tree(&journal, "target"); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + + assert_eq!(comparable_journal_tree(&journal, "target"), tree_before); + assert_eq!( + solstone_core_indexer_store::merge::fingerprint_edge_rows(&journal).unwrap(), + index_before + ); + fs::remove_dir_all(journal).unwrap(); +} + #[test] fn undo_removes_moved_target_facet_relationship() { let journal = undo_journal();