diff --git a/core/crates/solstone-core-entity/src/fixture_tests.rs b/core/crates/solstone-core-entity/src/fixture_tests.rs index d43e35079..eacae0232 100644 --- a/core/crates/solstone-core-entity/src/fixture_tests.rs +++ b/core/crates/solstone-core-entity/src/fixture_tests.rs @@ -51,6 +51,7 @@ const JOURNAL_IO_ALLOWED: &[&str] = &[ "capture_snapshot", "restore_snapshot", "JournalSnapshot", + "SnapshotDirectory", "SnapshotError", "AppendError", "atomic_replace", @@ -154,6 +155,7 @@ fn collect_production_sources(directory: &Path, sources: &mut Vec) { | "store_tests.rs" | "test_support.rs" | "trust_lock_tests.rs" + | "undo_tests.rs" ) ) { diff --git a/core/crates/solstone-core-entity/src/lib.rs b/core/crates/solstone-core-entity/src/lib.rs index 37e792671..a33120dff 100644 --- a/core/crates/solstone-core-entity/src/lib.rs +++ b/core/crates/solstone-core-entity/src/lib.rs @@ -22,14 +22,14 @@ pub use store::{ EntityIdentityRepairRefusal, EntityIdentityRepairReport, EntityIdentityRepairSkip, EntityIdentityRepairSkipReason, EntityMergeError, EntityMergeOptions, EntityMergePreview, EntityMergeReport, EntityOperationContext, EntityOperationKind, EntitySaveResult, - EntityStoreError, EntityWriteError, HistoryEvent, IdentityMapCacheLoad, IdentityMapLoser, - IdentityMapLoserReason, IdentitySnapshot, PreparedHistoryEvent, PreparedHistoryOutcome, - 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, + EntityStoreError, EntityUndoError, EntityUndoReport, EntityWriteError, HistoryEvent, + IdentityMapCacheLoad, IdentityMapLoser, IdentityMapLoserReason, IdentitySnapshot, + PreparedHistoryEvent, PreparedHistoryOutcome, 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, undo_entity_merge, }; pub use trust_lock::{EntityTrustLock, EntityTrustLockError, hold_entity_trust_lock}; @@ -53,3 +53,5 @@ mod store_tests; mod test_support; #[cfg(test)] mod trust_lock_tests; +#[cfg(test)] +mod undo_tests; diff --git a/core/crates/solstone-core-entity/src/store/mod.rs b/core/crates/solstone-core-entity/src/store/mod.rs index ce3fded99..efb49b0a7 100644 --- a/core/crates/solstone-core-entity/src/store/mod.rs +++ b/core/crates/solstone-core-entity/src/store/mod.rs @@ -14,6 +14,7 @@ mod merge_rollback; mod paths; mod reconcile; mod repair; +mod undo; #[allow(dead_code)] pub(crate) mod voiceprints; mod write; @@ -42,6 +43,7 @@ pub use repair::{ EntityIdentityRepairReport, EntityIdentityRepairSkip, EntityIdentityRepairSkipReason, repair_entity_identities, }; +pub use undo::{EntityUndoError, EntityUndoReport, undo_entity_merge}; pub use write::{ AmbiguityChoiceEntity, AmbiguityChoiceRequest, AmbiguityObservation, EntityOperationContext, EntityOperationKind, EntitySaveResult, EntityWriteError, IdentityMapCacheLoad, @@ -52,6 +54,8 @@ pub use write::{ #[cfg(test)] pub(crate) use repair::set_repair_identity_write_failure_on_attempt; #[cfg(test)] +pub(crate) use undo::undo_entity_merge_with_injector; +#[cfg(test)] pub(crate) use write::{ save_entity_identity_with_timeout, set_forced_identity_write_failure, write_history_event_json_for_test, diff --git a/core/crates/solstone-core-entity/src/store/undo.rs b/core/crates/solstone-core-entity/src/store/undo.rs new file mode 100644 index 000000000..f4d06ed69 --- /dev/null +++ b/core/crates/solstone-core-entity/src/store/undo.rs @@ -0,0 +1,591 @@ +// 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 serde_json::Value; +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::SnapshotDirectory; +use solstone_core_journal_io::SnapshotError; +use solstone_core_journal_io::contained_path; +use solstone_core_journal_io::list_dir_entries; +use solstone_core_journal_io::path_lexists; +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, EntityWriteError, hold_entity_trust_lock, + save_entity_identity, +}; + +use super::merge_payload::{ + MergePayloadError, load_entity_merge_payload, remove_entity_merge_payload, +}; +use super::merge_rollback::MergeRollback; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct EntityUndoReport { + pub merge_id: String, + pub source_id: String, + pub target_id: String, +} + +#[derive(Debug)] +pub enum EntityUndoError { + Refused(String), + Payload(MergePayloadError), + Write(EntityWriteError), + Snapshot(SnapshotError), + Index(solstone_core_indexer_store::StoreError), + Failed { + failed_phase: String, + report: Box, + rollback_error: Option, + }, +} + +impl fmt::Display for EntityUndoError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Refused(message) => formatter.write_str(message), + Self::Payload(error) => error.fmt(formatter), + Self::Write(error) => error.fmt(formatter), + Self::Snapshot(error) => error.fmt(formatter), + Self::Index(error) => error.fmt(formatter), + Self::Failed { failed_phase, .. } => { + write!(formatter, "entity merge undo failed during {failed_phase}") + } + } + } +} + +impl Error for EntityUndoError {} + +impl From for EntityUndoError { + fn from(error: MergePayloadError) -> Self { + Self::Payload(error) + } +} + +impl From for EntityUndoError { + fn from(error: EntityWriteError) -> Self { + Self::Write(error) + } +} + +impl From for EntityUndoError { + fn from(error: SnapshotError) -> Self { + Self::Snapshot(error) + } +} + +pub fn undo_entity_merge( + journal: &Path, + merge_id: &str, + caller: Value, +) -> Result { + undo_entity_merge_with_injector(journal, merge_id, caller, None) +} + +pub(crate) fn undo_entity_merge_with_injector( + journal: &Path, + merge_id: &str, + caller: Value, + injector: Option<&dyn Fn(&str) -> bool>, +) -> Result { + let target_id = find_payload_holder(journal, merge_id)?; + let payload = load_entity_merge_payload(journal, &target_id, merge_id)?; + let source_id = payload + .get("source_id") + .and_then(Value::as_str) + .ok_or_else(|| { + EntityUndoError::Refused("merge payload missing source entity id".to_owned()) + })? + .to_owned(); + let snapshot_paths = snapshot_paths(&payload)?; + let source_identity = payload + .get("source_state") + .and_then(Value::as_object) + .and_then(|source_state| source_state.get("identity")) + .filter(|identity| identity.is_object()) + .cloned() + .ok_or_else(|| { + EntityUndoError::Refused("merge payload source_state is missing identity".to_owned()) + })?; + let target_before = payload + .get("manifest") + .and_then(Value::as_object) + .and_then(|manifest| manifest.get("identity")) + .and_then(Value::as_object) + .and_then(|identity| identity.get("target_before")) + .filter(|identity| identity.is_object()) + .cloned() + .ok_or_else(|| { + EntityUndoError::Refused("merge payload identity missing target_before".to_owned()) + })?; + let report = EntityUndoReport { + merge_id: merge_id.to_owned(), + source_id: source_id.clone(), + target_id: target_id.clone(), + }; + let _trust = hold_entity_trust_lock(journal).map_err(|error| EntityUndoError::Failed { + failed_phase: "trust_lock".to_owned(), + report: Box::new(report.clone()), + rollback_error: Some(error.to_string()), + })?; + let mut rollback = MergeRollback::default(); + let mut phase = "snapshot"; + let result: Result<(), EntityUndoError> = (|| { + let mut paths = HashSet::from([ + format!("entities/{target_id}"), + format!("entities/{source_id}"), + "indexer".to_owned(), + ]); + paths.extend(snapshot_paths.iter().cloned()); + for path in paths { + rollback.capture(journal, &path)?; + } + phase = "source_state"; + for path in &snapshot_paths { + restore_snapshot( + journal, + &JournalSnapshot::Directory(SnapshotDirectory { + path: path.clone(), + entries: Vec::new(), + }), + )?; + } + save_entity_identity(journal, &source_id, &source_identity, None)?; + phase = "segments"; + undo_segments(journal, &payload, &mut rollback)?; + inject_failure(injector, phase)?; + phase = "activities"; + undo_activities(journal, &payload, &mut rollback)?; + inject_failure(injector, phase)?; + phase = "observations"; + undo_observation_relations(journal, &payload, &mut rollback)?; + inject_failure(injector, phase)?; + phase = "facets"; + undo_facets(journal, &target_id, &payload, &mut rollback)?; + inject_failure(injector, phase)?; + phase = "identity"; + save_entity_identity( + journal, + &target_id, + &target_before, + Some(&EntityOperationContext { + kind: EntityOperationKind::MergeUndo, + caller, + actor: Value::Null, + metadata: serde_json::json!({ + "undo_of": merge_id, + "source_id": source_id, + "target_id": target_id, + }), + }), + )?; + phase = "private_payload"; + remove_entity_merge_payload(journal, &target_id, merge_id)?; + phase = "edges"; + solstone_core_indexer_store::merge::rebuild_edges_for_recorded_merge_undo(journal) + .map_err(EntityUndoError::Index)?; + Ok(()) + })(); + if let Err(error) = result { + let rollback_error = rollback + .restore(journal) + .err() + .map(|rollback_error| rollback_error.to_string()); + return Err(EntityUndoError::Failed { + failed_phase: phase.to_owned(), + report: Box::new(report), + rollback_error: rollback_error.or_else(|| Some(error.to_string())), + }); + } + Ok(report) +} + +fn inject_failure( + injector: Option<&dyn Fn(&str) -> bool>, + phase: &str, +) -> Result<(), EntityUndoError> { + if injector.is_some_and(|injector| injector(phase)) { + return Err(EntityUndoError::Refused(format!( + "injected failure after phase {phase}" + ))); + } + Ok(()) +} + +fn undo_facets( + journal: &Path, + target_id: &str, + payload: &Value, + rollback: &mut MergeRollback, +) -> Result<(), EntityUndoError> { + let entries = payload + .get("manifest") + .and_then(Value::as_object) + .and_then(|manifest| manifest.get("facets")) + .and_then(Value::as_object) + .and_then(|facets| facets.get("entries")) + .and_then(Value::as_array) + .ok_or_else(|| { + EntityUndoError::Refused("merge payload facets entries are missing".to_owned()) + })?; + for entry in entries { + let facet = entry + .get("facet") + .and_then(Value::as_str) + .filter(|facet| !facet.is_empty()) + .ok_or_else(|| { + EntityUndoError::Refused("merge payload facet entry is missing facet".to_owned()) + })?; + let kind = entry.get("kind").and_then(Value::as_str).ok_or_else(|| { + EntityUndoError::Refused("merge payload facet entry is missing kind".to_owned()) + })?; + let directory = format!("facets/{facet}/entities/{target_id}"); + match kind { + "move" => { + rollback.capture(journal, &directory)?; + restore_snapshot(journal, &JournalSnapshot::Missing { path: directory })?; + } + "merge" => { + let target_before = entry + .get("target_before") + .filter(|value| value.is_object()) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload facet entry is missing target_before".to_owned(), + ) + })?; + let path = format!("{directory}/entity.json"); + let destination = contained_path(journal, &path) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + rollback.capture(journal, &path)?; + write_json( + destination, + target_before, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let observations_before = entry + .get("target_observations_before") + .and_then(Value::as_array) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload facet entry is missing target_observations_before" + .to_owned(), + ) + })?; + let observations_existed = entry + .get("target_observations_existed") + .and_then(Value::as_bool) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload facet entry is missing target_observations_existed" + .to_owned(), + ) + })?; + let observations_path = format!("{directory}/observations.jsonl"); + rollback.capture(journal, &observations_path)?; + if observations_existed { + let destination = contained_path(journal, &observations_path) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + write_jsonl( + destination, + observations_before.to_vec(), + AtomicWriteOptions::default(), + ) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + } else { + restore_snapshot( + journal, + &JournalSnapshot::Missing { + path: observations_path, + }, + )?; + } + } + _ => { + return Err(EntityUndoError::Refused(format!( + "merge payload facet entry has unknown kind: {kind}" + ))); + } + } + } + Ok(()) +} + +fn undo_observation_relations( + journal: &Path, + payload: &Value, + rollback: &mut MergeRollback, +) -> Result<(), EntityUndoError> { + let entries = payload + .get("manifest") + .and_then(Value::as_object) + .and_then(|manifest| manifest.get("observation_relations")) + .and_then(Value::as_object) + .and_then(|relations| relations.get("entries")) + .and_then(Value::as_array) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload observation relation entries are missing".to_owned(), + ) + })?; + for entry in entries { + let path = entry.get("path").and_then(Value::as_str).ok_or_else(|| { + EntityUndoError::Refused( + "merge payload observation relation entry is missing path".to_owned(), + ) + })?; + let row_index = entry + .get("row_index") + .and_then(Value::as_u64) + .and_then(|index| usize::try_from(index).ok()) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload observation relation entry is missing row_index".to_owned(), + ) + })?; + let target_before = entry + .get("target_before") + .and_then(Value::as_str) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload observation relation entry is missing target_before".to_owned(), + ) + })?; + let destination = contained_path(journal, path) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let mut rows: Vec = solstone_core_journal_io::read_jsonl( + &destination, + Vec::new(), + solstone_core_journal_io::MalformedPolicy::Raise, + ) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let relation = rows + .get_mut(row_index) + .and_then(|row| row.get_mut("relation")) + .and_then(Value::as_object_mut) + .ok_or_else(|| { + EntityUndoError::Refused(format!( + "merge payload observation relation row is missing: {path}:{row_index}" + )) + })?; + relation.insert( + "target_entity_id".to_owned(), + Value::String(target_before.to_owned()), + ); + rollback.capture(journal, path)?; + write_jsonl(destination, rows, AtomicWriteOptions::default()) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + } + Ok(()) +} + +fn undo_segments( + journal: &Path, + payload: &Value, + rollback: &mut MergeRollback, +) -> Result<(), EntityUndoError> { + let entries = manifest_entries(payload, "segments")?; + for entry in entries { + let path = entry_path(entry, "segment")?; + let section = entry + .get("section") + .and_then(Value::as_str) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload segment entry is missing section".to_owned(), + ) + })?; + let row_index = entry_index(entry, "row_index", "segment")?; + let field = entry.get("field").and_then(Value::as_str).ok_or_else(|| { + EntityUndoError::Refused("merge payload segment entry is missing field".to_owned()) + })?; + let before = entry.get("before").cloned().ok_or_else(|| { + EntityUndoError::Refused("merge payload segment entry is missing before".to_owned()) + })?; + let destination = contained_path(journal, path) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let mut value: Value = read_json(&destination, Value::Null, MalformedPolicy::Raise) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let object = value + .get_mut(section) + .and_then(Value::as_array_mut) + .and_then(|rows| rows.get_mut(row_index)) + .and_then(Value::as_object_mut) + .ok_or_else(|| { + EntityUndoError::Refused(format!( + "merge payload segment row is missing: {path}:{section}:{row_index}" + )) + })?; + object.insert(field.to_owned(), before); + rollback.capture(journal, path)?; + write_json( + destination, + &value, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + } + Ok(()) +} + +fn undo_activities( + journal: &Path, + payload: &Value, + rollback: &mut MergeRollback, +) -> Result<(), EntityUndoError> { + let entries = manifest_entries(payload, "activities")?; + for entry in entries { + let path = entry_path(entry, "activity")?; + let row_index = entry_index(entry, "row_index", "activity")?; + let container = entry + .get("container") + .and_then(Value::as_str) + .ok_or_else(|| { + EntityUndoError::Refused( + "merge payload activity entry is missing container".to_owned(), + ) + })?; + let item_index = entry_index(entry, "item_index", "activity")?; + let before = entry.get("before").cloned().ok_or_else(|| { + EntityUndoError::Refused("merge payload activity entry is missing before".to_owned()) + })?; + let destination = contained_path(journal, path) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let mut rows: Vec = read_jsonl(&destination, Vec::new(), MalformedPolicy::Raise) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + let row = rows + .get_mut(row_index) + .and_then(Value::as_object_mut) + .ok_or_else(|| { + EntityUndoError::Refused(format!( + "merge payload activity row is missing: {path}:{row_index}" + )) + })?; + if let Some(field) = entry.get("field").and_then(Value::as_str) { + let object = row + .get_mut(container) + .and_then(Value::as_array_mut) + .and_then(|items| items.get_mut(item_index)) + .and_then(Value::as_object_mut) + .ok_or_else(|| { + EntityUndoError::Refused(format!( + "merge payload activity field is missing: {path}:{row_index}:{container}:{item_index}" + )) + })?; + object.insert(field.to_owned(), before); + } else { + let value = row + .get_mut(container) + .and_then(Value::as_array_mut) + .and_then(|items| items.get_mut(item_index)) + .ok_or_else(|| { + EntityUndoError::Refused(format!( + "merge payload activity item is missing: {path}:{row_index}:{container}:{item_index}" + )) + })?; + *value = before; + } + rollback.capture(journal, path)?; + write_jsonl(destination, rows, AtomicWriteOptions::default()) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + } + Ok(()) +} + +fn manifest_entries<'a>(payload: &'a Value, section: &str) -> Result<&'a [Value], EntityUndoError> { + payload + .get("manifest") + .and_then(Value::as_object) + .and_then(|manifest| manifest.get(section)) + .and_then(Value::as_object) + .and_then(|section| section.get("entries")) + .and_then(Value::as_array) + .map(Vec::as_slice) + .ok_or_else(|| { + EntityUndoError::Refused(format!("merge payload {section} entries are missing")) + }) +} + +fn entry_path<'a>(entry: &'a Value, section: &str) -> Result<&'a str, EntityUndoError> { + entry.get("path").and_then(Value::as_str).ok_or_else(|| { + EntityUndoError::Refused(format!("merge payload {section} entry is missing path")) + }) +} + +fn entry_index(entry: &Value, field: &str, section: &str) -> Result { + entry + .get(field) + .and_then(Value::as_u64) + .and_then(|index| usize::try_from(index).ok()) + .ok_or_else(|| { + EntityUndoError::Refused(format!("merge payload {section} entry is missing {field}")) + }) +} + +fn find_payload_holder(journal: &Path, merge_id: &str) -> Result { + let entities = contained_path(journal, "entities") + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + for entry in + list_dir_entries(&entities).map_err(|error| EntityUndoError::Refused(error.to_string()))? + { + if entry.kind != DirEntryKind::Directory { + continue; + } + let entity_id = entry.name.to_string_lossy().into_owned(); + let path = contained_path( + journal, + &format!("entities/{entity_id}/history/private/{merge_id}.json"), + ) + .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + if path_lexists(&path).map_err(|error| EntityUndoError::Refused(error.to_string()))? { + return Ok(entity_id); + } + } + Err(EntityUndoError::Refused(format!( + "private merge payload not found: {merge_id}" + ))) +} + +fn snapshot_paths(payload: &Value) -> Result, EntityUndoError> { + payload + .get("source_state") + .and_then(Value::as_object) + .and_then(|source_state| source_state.get("snapshots")) + .and_then(Value::as_array) + .ok_or_else(|| EntityUndoError::Refused("merge payload missing snapshots".to_owned()))? + .iter() + .map(|snapshot| { + snapshot + .get("rel") + .and_then(Value::as_str) + .map(str::to_owned) + .ok_or_else(|| { + EntityUndoError::Refused("manifest snapshot missing relative path".to_owned()) + }) + }) + .collect() +} diff --git a/core/crates/solstone-core-entity/src/undo_tests.rs b/core/crates/solstone-core-entity/src/undo_tests.rs new file mode 100644 index 000000000..17ac6b43f --- /dev/null +++ b/core/crates/solstone-core-entity/src/undo_tests.rs @@ -0,0 +1,465 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::fs; +use std::path::PathBuf; +use std::sync::atomic::{AtomicU64, Ordering}; + +use serde_json::{Value, json}; + +use super::store::undo_entity_merge_with_injector; +use crate::{ + EntityMergeOptions, commit_entity_merge, read_entity_identity, save_entity_identity, + undo_entity_merge, +}; + +static NEXT_UNDO_DIRECTORY: AtomicU64 = AtomicU64::new(0); + +fn undo_journal() -> PathBuf { + let path = std::env::temp_dir().join(format!( + "solstone-undo-{}-{}", + std::process::id(), + NEXT_UNDO_DIRECTORY.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&path).unwrap(); + path +} + +#[test] +fn undo_reverts_target_identity_restores_source_and_removes_payload() { + let journal = undo_journal(); + let source = json!({"id":"source","name":"Source","aka":[],"emails":[],"title":"Engineer"}); + let target = json!({"id":"target","name":"Target","aka":[],"emails":[],"title":"Director"}); + save_entity_identity(&journal, "source", &source, None).unwrap(); + save_entity_identity(&journal, "target", &target, None).unwrap(); + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + assert_eq!( + read_entity_identity(&journal, "target") + .unwrap() + .unwrap() + .value(), + &target + ); + assert!(journal.join("entities/source/entity.json").exists()); + assert!(read_entity_identity(&journal, "source").unwrap().is_some()); + assert!( + !journal + .join(format!( + "entities/target/history/private/{}.json", + merge.merge_id + )) + .exists() + ); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn undo_removes_moved_target_facet_relationship() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let source_facet = journal.join("facets/work/entities/source"); + fs::create_dir_all(&source_facet).unwrap(); + fs::write( + source_facet.join("entity.json"), + br#"{"entity_id":"source"}"#, + ) + .unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let target_facet = journal.join("facets/work/entities/target"); + assert!(target_facet.exists()); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + assert!(!target_facet.exists()); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn undo_restores_merged_target_facet_relationship() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let source_facet = journal.join("facets/work/entities/source"); + let target_facet = journal.join("facets/work/entities/target"); + fs::create_dir_all(&source_facet).unwrap(); + fs::create_dir_all(&target_facet).unwrap(); + fs::write( + source_facet.join("entity.json"), + br#"{"entity_id":"source","attached_at":"2026-01-01"}"#, + ) + .unwrap(); + let target_before = json!({ + "entity_id": "target", + "attached_at": "2026-02-01", + "description": "target description" + }); + fs::write( + target_facet.join("entity.json"), + serde_json::to_vec(&target_before).unwrap(), + ) + .unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + let restored: Value = + serde_json::from_slice(&fs::read(target_facet.join("entity.json")).unwrap()).unwrap(); + assert_eq!(restored, target_before); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn undo_restores_target_observations_from_merged_facet() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let source_facet = journal.join("facets/work/entities/source"); + let target_facet = journal.join("facets/work/entities/target"); + fs::create_dir_all(&source_facet).unwrap(); + fs::create_dir_all(&target_facet).unwrap(); + fs::write( + source_facet.join("entity.json"), + br#"{"entity_id":"source"}"#, + ) + .unwrap(); + fs::write( + target_facet.join("entity.json"), + br#"{"entity_id":"target"}"#, + ) + .unwrap(); + let source_observation = json!({"content": "source", "observed_at": "2026-01-01"}); + let target_observation = json!({"content": "target", "observed_at": "2026-01-02"}); + fs::write( + source_facet.join("observations.jsonl"), + format!("{}\n", source_observation), + ) + .unwrap(); + fs::write( + target_facet.join("observations.jsonl"), + format!("{}\n", target_observation), + ) + .unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + let restored = fs::read_to_string(target_facet.join("observations.jsonl")) + .unwrap() + .lines() + .map(|line| serde_json::from_str::(line).unwrap()) + .collect::>(); + assert_eq!(restored, vec![target_observation]); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn undo_restores_observation_relation_target() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let observations = journal.join("facets/work/entities/other/observations.jsonl"); + fs::create_dir_all(observations.parent().unwrap()).unwrap(); + fs::write( + &observations, + b"{\"content\":\"note\",\"observed_at\":\"2026-01-01\",\"relation\":{\"kind\":\"works-with\",\"target_entity_id\":\"source\"}}\n", + ) + .unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let remapped: Value = + serde_json::from_str(&fs::read_to_string(&observations).unwrap()).unwrap(); + assert_eq!(remapped["relation"]["target_entity_id"], "target"); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + let restored: Value = + serde_json::from_str(&fs::read_to_string(&observations).unwrap()).unwrap(); + assert_eq!(restored["relation"]["target_entity_id"], "source"); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn undo_restores_segment_speaker_label() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let labels = journal.join("chronicle/20260102/080000_300/talents/speaker_labels.json"); + fs::create_dir_all(labels.parent().unwrap()).unwrap(); + fs::write(&labels, br#"{"labels":[{"speaker":"source"}]}"#).unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let rewritten: Value = serde_json::from_slice(&fs::read(&labels).unwrap()).unwrap(); + assert_eq!(rewritten["labels"][0]["speaker"], "target"); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + let restored: Value = serde_json::from_slice(&fs::read(&labels).unwrap()).unwrap(); + assert_eq!(restored["labels"][0]["speaker"], "source"); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn undo_restores_activity_active_entity() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let activities = journal.join("facets/work/activities/20260102.jsonl"); + fs::create_dir_all(activities.parent().unwrap()).unwrap(); + fs::write( + &activities, + b"{\"id\":\"activity\",\"active_entities\":[\"source\"]}\n", + ) + .unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let rewritten: Value = serde_json::from_str(&fs::read_to_string(&activities).unwrap()).unwrap(); + assert_eq!(rewritten["active_entities"][0], "target"); + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + let restored: Value = serde_json::from_str(&fs::read_to_string(&activities).unwrap()).unwrap(); + assert_eq!(restored["active_entities"][0], "source"); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn facets_undo_injection_rolls_back_and_retry_succeeds() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let moved_source = journal.join("facets/moved/entities/source"); + fs::create_dir_all(&moved_source).unwrap(); + fs::write( + moved_source.join("entity.json"), + br#"{"entity_id":"source"}"#, + ) + .unwrap(); + let merged_source = journal.join("facets/merged/entities/source"); + let merged_target = journal.join("facets/merged/entities/target"); + fs::create_dir_all(&merged_source).unwrap(); + fs::create_dir_all(&merged_target).unwrap(); + fs::write( + merged_source.join("entity.json"), + br#"{"entity_id":"source","attached_at":"2026-01-01"}"#, + ) + .unwrap(); + let merged_before = json!({"entity_id":"target","attached_at":"2026-02-01"}); + fs::write( + merged_target.join("entity.json"), + serde_json::to_vec(&merged_before).unwrap(), + ) + .unwrap(); + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let moved_target = journal.join("facets/moved/entities/target/entity.json"); + let moved_after_merge = fs::read(&moved_target).unwrap(); + let merged_after_merge = fs::read(merged_target.join("entity.json")).unwrap(); + assert!( + undo_entity_merge_with_injector( + &journal, + &merge.merge_id, + Value::Null, + Some(&|phase| phase == "facets"), + ) + .is_err() + ); + assert_eq!(fs::read(&moved_target).unwrap(), moved_after_merge); + assert_eq!( + fs::read(merged_target.join("entity.json")).unwrap(), + merged_after_merge + ); + + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + assert!(!moved_target.exists()); + let restored: Value = + serde_json::from_slice(&fs::read(merged_target.join("entity.json")).unwrap()).unwrap(); + assert_eq!(restored, merged_before); + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn observations_undo_injection_rolls_back_and_retry_succeeds() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let first = journal.join("facets/one/entities/other/observations.jsonl"); + let second = journal.join("facets/two/entities/other/observations.jsonl"); + for path in [&first, &second] { + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write( + path, + b"{\"relation\":{\"kind\":\"works-with\",\"target_entity_id\":\"source\"}}\n", + ) + .unwrap(); + } + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let first_after_merge = fs::read(&first).unwrap(); + let second_after_merge = fs::read(&second).unwrap(); + assert!( + undo_entity_merge_with_injector( + &journal, + &merge.merge_id, + Value::Null, + Some(&|phase| phase == "observations"), + ) + .is_err() + ); + assert_eq!(fs::read(&first).unwrap(), first_after_merge); + assert_eq!(fs::read(&second).unwrap(), second_after_merge); + + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + for path in [&first, &second] { + let restored: Value = serde_json::from_slice(&fs::read(path).unwrap()).unwrap(); + assert_eq!(restored["relation"]["target_entity_id"], "source"); + } + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn segments_undo_injection_rolls_back_and_retry_succeeds() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let first = journal.join("chronicle/20260102/080000_300/talents/speaker_labels.json"); + let second = journal.join("chronicle/20260102/090000_300/talents/speaker_labels.json"); + for path in [&first, &second] { + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(path, br#"{"labels":[{"speaker":"source"}]}"#).unwrap(); + } + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let first_after_merge = fs::read(&first).unwrap(); + let second_after_merge = fs::read(&second).unwrap(); + assert!( + undo_entity_merge_with_injector( + &journal, + &merge.merge_id, + Value::Null, + Some(&|phase| phase == "segments"), + ) + .is_err() + ); + assert_eq!(fs::read(&first).unwrap(), first_after_merge); + assert_eq!(fs::read(&second).unwrap(), second_after_merge); + + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + for path in [&first, &second] { + let restored: Value = serde_json::from_slice(&fs::read(path).unwrap()).unwrap(); + assert_eq!(restored["labels"][0]["speaker"], "source"); + } + fs::remove_dir_all(journal).unwrap(); +} + +#[test] +fn activities_undo_injection_rolls_back_and_retry_succeeds() { + let journal = undo_journal(); + for id in ["source", "target"] { + save_entity_identity( + &journal, + id, + &json!({"id": id, "name": id, "aka": [], "emails": []}), + None, + ) + .unwrap(); + } + let first = journal.join("facets/work/activities/20260102.jsonl"); + let second = journal.join("facets/work/activities/20260103.jsonl"); + for path in [&first, &second] { + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write( + path, + b"{\"id\":\"activity\",\"active_entities\":[\"source\"]}\n", + ) + .unwrap(); + } + + let merge = + commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); + let first_after_merge = fs::read(&first).unwrap(); + let second_after_merge = fs::read(&second).unwrap(); + assert!( + undo_entity_merge_with_injector( + &journal, + &merge.merge_id, + Value::Null, + Some(&|phase| phase == "activities"), + ) + .is_err() + ); + assert_eq!(fs::read(&first).unwrap(), first_after_merge); + assert_eq!(fs::read(&second).unwrap(), second_after_merge); + + undo_entity_merge(&journal, &merge.merge_id, Value::Null).unwrap(); + for path in [&first, &second] { + let restored: Value = serde_json::from_slice(&fs::read(path).unwrap()).unwrap(); + assert_eq!(restored["active_entities"][0], "source"); + } + fs::remove_dir_all(journal).unwrap(); +}