From e991c22f76856e09e1efbeb4f511f531258cbddf Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Wed, 5 Aug 2026 15:44:15 -0600 Subject: [PATCH] Test merge rollback within owner phases Inject failures after each owner-store artifact write so merge and undo rollback coverage exercises partial phase mutations across facets, segments, activities, and observation relation remapping. --- .../solstone-core-entity/src/merge_tests.rs | 148 +++++++++++------- .../solstone-core-entity/src/store/merge.rs | 56 ++++++- .../solstone-core-entity/src/store/undo.rs | 43 +++-- .../solstone-core-entity/src/undo_tests.rs | 8 +- 4 files changed, 173 insertions(+), 82 deletions(-) diff --git a/core/crates/solstone-core-entity/src/merge_tests.rs b/core/crates/solstone-core-entity/src/merge_tests.rs index c2e0fc266..785002777 100644 --- a/core/crates/solstone-core-entity/src/merge_tests.rs +++ b/core/crates/solstone-core-entity/src/merge_tests.rs @@ -152,7 +152,7 @@ fn facets_move_relationship_and_observations() { ) .unwrap(); assert_eq!( - merge_facets(&journal, "source", "target", None) + merge_facets(&journal, "source", "target", None, None) .unwrap() .moved_count, 1 @@ -567,22 +567,28 @@ fn facets_phase_injection_rolls_back_and_retry_succeeds() { ) .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 sources = [ + journal.join("facets/work/entities/source"), + journal.join("facets/personal/entities/source"), + ]; + for source in &sources { + 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"), + Some(&|phase: &str, artifact_index| phase == "facets" && artifact_index == 0), ); assert!(result.is_err()); - assert!(source.exists()); + assert!(sources.iter().all(|source| source.exists())); assert!(journal.join("entities/source").exists()); assert!(!journal.join("facets/work/entities/target").exists()); + assert!(!journal.join("facets/personal/entities/target").exists()); commit_entity_merge(&journal, "source", "target", EntityMergeOptions::default()).unwrap(); - assert!(!source.exists()); + assert!(sources.iter().all(|source| !source.exists())); fs::remove_dir_all(journal).unwrap(); } @@ -605,7 +611,7 @@ fn voiceprints_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "voiceprints") + Some(&|phase, artifact_index| phase == "voiceprints" && artifact_index == 0) ) .is_err() ); @@ -630,30 +636,39 @@ fn segments_phase_injection_rolls_back_and_retry_succeeds() { ) .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(); + let paths = [ + journal.join("chronicle/20260102/080000_300/talents/speaker_labels.json"), + journal.join("chronicle/20260102/090000_300/talents/speaker_labels.json"), + ]; + for path in &paths { + 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") + Some(&|phase, artifact_index| phase == "segments" && artifact_index == 0) ) .is_err() ); - assert_eq!( - serde_json::from_slice::(&fs::read(&path).unwrap()).unwrap()["labels"] - [0]["speaker"], - "source" - ); + for path in &paths { + 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" - ); + for path in &paths { + assert_eq!( + serde_json::from_slice::(&fs::read(path).unwrap()).unwrap()["labels"] + [0]["speaker"], + "target" + ); + } fs::remove_dir_all(journal).unwrap(); } @@ -669,30 +684,39 @@ fn activities_phase_injection_rolls_back_and_retry_succeeds() { ) .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(); + let paths = [ + journal.join("facets/work/activities/20260102.jsonl"), + journal.join("facets/work/activities/20260103.jsonl"), + ]; + for path in &paths { + 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") + Some(&|phase, artifact_index| phase == "activities" && artifact_index == 0) ) .is_err() ); - assert_eq!( - serde_json::from_str::(&fs::read_to_string(&path).unwrap()).unwrap()["active_entities"] - [0], - "source" - ); + for path in &paths { + 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" - ); + for path in &paths { + assert_eq!( + serde_json::from_str::(&fs::read_to_string(path).unwrap()).unwrap() + ["active_entities"][0], + "target" + ); + } fs::remove_dir_all(journal).unwrap(); } @@ -708,30 +732,42 @@ fn observation_relations_phase_injection_rolls_back_and_retry_succeeds() { ) .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(); + let paths = [ + journal.join("facets/work/entities/other-one/observations.jsonl"), + journal.join("facets/work/entities/other-two/observations.jsonl"), + ]; + for path in &paths { + 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") + Some( + &|phase, artifact_index| phase == "observation relation remap" + && artifact_index == 0 + ) ) .is_err() ); - assert_eq!( - serde_json::from_str::(&fs::read_to_string(&path).unwrap()).unwrap()["relation"] - ["target_entity_id"], - "source" - ); + for path in &paths { + 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" - ); + for path in &paths { + assert_eq!( + serde_json::from_str::(&fs::read_to_string(path).unwrap()).unwrap() + ["relation"]["target_entity_id"], + "target" + ); + } fs::remove_dir_all(journal).unwrap(); } @@ -753,7 +789,7 @@ fn private_payload_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "private_payload") + Some(&|phase, artifact_index| phase == "private_payload" && artifact_index == 0) ) .is_err() ); @@ -801,7 +837,7 @@ fn lineage_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "lineage") + Some(&|phase, artifact_index| phase == "lineage" && artifact_index == 0) ) .is_err() ); @@ -853,7 +889,7 @@ fn cleanup_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "cleanup") + Some(&|phase, artifact_index| phase == "cleanup" && artifact_index == 0) ) .is_err() ); @@ -884,7 +920,7 @@ fn history_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "history") + Some(&|phase, artifact_index| phase == "history" && artifact_index == 0) ) .is_err() ); @@ -928,7 +964,7 @@ fn edges_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "edges") + Some(&|phase, artifact_index| phase == "edges" && artifact_index == 0) ) .is_err() ); @@ -972,7 +1008,7 @@ fn audit_phase_injection_rolls_back_and_retry_succeeds() { "source", "target", EntityMergeOptions::default(), - Some(&|phase| phase == "audit"), + Some(&|phase, artifact_index| phase == "audit" && artifact_index == 0), ) .unwrap_err(); let merge_id = match error { diff --git a/core/crates/solstone-core-entity/src/store/merge.rs b/core/crates/solstone-core-entity/src/store/merge.rs index 7af91baf8..588ba025e 100644 --- a/core/crates/solstone-core-entity/src/store/merge.rs +++ b/core/crates/solstone-core-entity/src/store/merge.rs @@ -56,6 +56,7 @@ const PHASES: [&str; 11] = [ "audit", ]; type VoiceprintKey = (Option, Option, Option, Option); +type FailureInjector = dyn Fn(&str, usize) -> bool; #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub struct EntityMergeOptions { @@ -206,7 +207,7 @@ pub(crate) fn commit_entity_merge_with_injector( source_id: &str, target_id: &str, options: EntityMergeOptions, - injector: Option<&dyn Fn(&str) -> bool>, + injector: Option<&FailureInjector>, ) -> Result { let _trust = hold_entity_trust_lock(journal).map_err(EntityWriteError::TrustLock)?; let plan = plan_merge(journal, source_id, target_id, options)?; @@ -241,14 +242,14 @@ pub(crate) fn commit_entity_merge_with_injector( 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; payload["manifest"]["facets"]["entries"] = Value::Array(result.entries); }), + "facets" => merge_facets(journal, source_id, target_id, Some(&mut rollback), injector).map(|result| { stats.facets_moved=result.moved_count; stats.facets_merged=result.merged_count; touched_facets = result.touched_facets; payload["manifest"]["facets"]["entries"] = Value::Array(result.entries); }), "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; payload["manifest"]["segments"]["entries"] = Value::Array(result.entries); }), - "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; payload["manifest"]["activities"]["entries"] = Value::Array(result.entries); }), - "observation relation remap" => merge_observation_relations(journal, source_id, target_id, Some(&mut rollback)).map(|result| { stats.observation_relations_rewritten=result.rows_rewritten; payload["manifest"]["observation_relations"]["entries"] = Value::Array(result.entries); }), + "segments" => merge_segment_labels(journal, source_id, target_id, Some(&mut rollback), injector).map(|result| { stats.segments_labels_rewritten=result.labels_rewritten; stats.segments_corrections_rewritten=result.corrections_rewritten; stats.segments_files_scanned=result.files_scanned; payload["manifest"]["segments"]["entries"] = Value::Array(result.entries); }), + "activities" => merge_activities(journal, source_id, target_id, Some(&mut rollback), injector).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; payload["manifest"]["activities"]["entries"] = Value::Array(result.entries); }), + "observation relation remap" => merge_observation_relations(journal, source_id, target_id, Some(&mut rollback), injector).map(|result| { stats.observation_relations_rewritten=result.rows_rewritten; payload["manifest"]["observation_relations"]["entries"] = Value::Array(result.entries); }), "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"), }; @@ -263,7 +264,11 @@ pub(crate) fn commit_entity_merge_with_injector( rollback_error: rollback_error.or_else(|| Some(error.to_string())), }); } - if injector.is_some_and(|injector| injector(phase)) { + if !matches!( + phase, + "facets" | "segments" | "activities" | "observation relation remap" + ) && injector.is_some_and(|injector| injector(phase, 0)) + { let rollback_error = rollback .restore(journal) .err() @@ -272,7 +277,7 @@ pub(crate) fn commit_entity_merge_with_injector( failed_phase: phase.to_owned(), report: Box::new(report), rollback_error: rollback_error - .or_else(|| Some(format!("injected failure after phase {phase}"))), + .or_else(|| Some(format!("injected failure after {phase} artifact 0"))), }); } report.completed_phases.push(phase.to_owned()); @@ -280,6 +285,19 @@ pub(crate) fn commit_entity_merge_with_injector( Ok(report) } +fn inject_failure( + injector: Option<&FailureInjector>, + phase: &str, + artifact_index: usize, +) -> Result<(), EntityMergeError> { + if injector.is_some_and(|injector| injector(phase, artifact_index)) { + return Err(EntityMergeError::Refused(format!( + "injected failure after {phase} artifact {artifact_index}" + ))); + } + Ok(()) +} + #[derive(Debug, Clone, PartialEq, Default)] pub(crate) struct ObservationRelationMergeStats { pub rows_rewritten: usize, @@ -290,10 +308,12 @@ pub(crate) fn merge_observation_relations( source_id: &str, target_id: &str, mut rollback: Option<&mut MergeRollback>, + injector: Option<&FailureInjector>, ) -> Result { let facets = contained_path(journal, "facets") .map_err(|error| EntityMergeError::Refused(error.to_string()))?; let mut stats = ObservationRelationMergeStats::default(); + let mut artifact_index = 0; for facet in list_dir_entries(&facets).map_err(|error| EntityMergeError::Refused(error.to_string()))? { @@ -343,6 +363,8 @@ pub(crate) fn merge_observation_relations( capture_rollback_file(&mut rollback, journal, &path)?; write_jsonl(&path, rows, AtomicWriteOptions::default()) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "observation relation remap", artifact_index)?; + artifact_index += 1; } } } @@ -362,10 +384,12 @@ pub(crate) fn merge_activities( source_id: &str, target_id: &str, mut rollback: Option<&mut MergeRollback>, + injector: Option<&FailureInjector>, ) -> Result { let facets = contained_path(journal, "facets") .map_err(|error| EntityMergeError::Refused(error.to_string()))?; let mut stats = ActivityMergeStats::default(); + let mut artifact_index = 0; for facet in list_dir_entries(&facets).map_err(|error| EntityMergeError::Refused(error.to_string()))? { @@ -470,6 +494,8 @@ pub(crate) fn merge_activities( capture_rollback_file(&mut rollback, journal, &file.path)?; write_jsonl(&file.path, rows, AtomicWriteOptions::default()) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "activities", artifact_index)?; + artifact_index += 1; stats.files_rewritten += 1; } } @@ -490,8 +516,10 @@ pub(crate) fn merge_segment_labels( source_id: &str, target_id: &str, mut rollback: Option<&mut MergeRollback>, + injector: Option<&FailureInjector>, ) -> Result { let mut stats = SegmentMergeStats::default(); + let mut artifact_index = 0; for day in day_dirs(journal) .map_err(|error| EntityMergeError::Refused(error.to_string()))? .into_values() @@ -551,6 +579,8 @@ pub(crate) fn merge_segment_labels( }, ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "segments", artifact_index)?; + artifact_index += 1; stats.labels_rewritten += 1; } } @@ -610,6 +640,8 @@ pub(crate) fn merge_segment_labels( }, ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "segments", artifact_index)?; + artifact_index += 1; stats.corrections_rewritten += 1; } } @@ -641,10 +673,12 @@ pub(crate) fn merge_facets( source_id: &str, target_id: &str, mut rollback: Option<&mut MergeRollback>, + injector: Option<&FailureInjector>, ) -> Result { let facets = contained_path(journal, "facets") .map_err(|error| EntityMergeError::Refused(error.to_string()))?; let mut stats = FacetMergeStats::default(); + let mut artifact_index = 0; for entry in list_dir_entries(&facets).map_err(|error| EntityMergeError::Refused(error.to_string()))? { @@ -694,6 +728,8 @@ pub(crate) fn merge_facets( }, ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; if path_lexists(&source_obs_path) .map_err(|error| EntityMergeError::Refused(error.to_string()))? { @@ -703,6 +739,8 @@ pub(crate) fn merge_facets( AtomicWriteOptions::default(), ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; } stats.moved_count += 1; stats.touched_facets.push(facet.clone()); @@ -728,12 +766,16 @@ pub(crate) fn merge_facets( }, ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; write_jsonl( target_obs_path, dedupe_observations(&source_obs, &target_obs), AtomicWriteOptions::default(), ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; stats.merged_count += 1; stats.touched_facets.push(facet.clone()); stats.entries.push(json!({ diff --git a/core/crates/solstone-core-entity/src/store/undo.rs b/core/crates/solstone-core-entity/src/store/undo.rs index de5e82afe..eb466a3e0 100644 --- a/core/crates/solstone-core-entity/src/store/undo.rs +++ b/core/crates/solstone-core-entity/src/store/undo.rs @@ -36,6 +36,8 @@ use super::merge_payload::{ }; use super::merge_rollback::MergeRollback; +type FailureInjector = dyn Fn(&str, usize) -> bool; + #[derive(Debug, Clone, PartialEq, Eq)] pub struct EntityUndoReport { pub merge_id: String, @@ -104,7 +106,7 @@ pub(crate) fn undo_entity_merge_with_injector( journal: &Path, merge_id: &str, caller: Value, - injector: Option<&dyn Fn(&str) -> bool>, + injector: Option<&FailureInjector>, ) -> Result { let target_id = find_payload_holder(journal, merge_id)?; let payload = load_entity_merge_payload(journal, &target_id, merge_id)?; @@ -146,17 +148,13 @@ pub(crate) fn undo_entity_merge_with_injector( phase = "voiceprints"; undo_voiceprints(journal, &target_id, &payload, &mut rollback)?; phase = "segments"; - undo_segments(journal, &payload, &mut rollback)?; - inject_failure(injector, phase)?; + undo_segments(journal, &payload, &mut rollback, injector)?; phase = "activities"; - undo_activities(journal, &payload, &mut rollback)?; - inject_failure(injector, phase)?; + undo_activities(journal, &payload, &mut rollback, injector)?; phase = "observations"; - undo_observation_relations(journal, &payload, &mut rollback)?; - inject_failure(injector, phase)?; + undo_observation_relations(journal, &payload, &mut rollback, injector)?; phase = "facets"; - undo_facets(journal, &target_id, &payload, &mut rollback)?; - inject_failure(injector, phase)?; + undo_facets(journal, &target_id, &payload, &mut rollback, injector)?; phase = "lineage"; undo_rebased_payloads(journal, &source_id, &target_id, &payload, &mut rollback)?; phase = "identity"; @@ -488,12 +486,13 @@ fn is_missing_value(value: Option<&Value>) -> bool { } fn inject_failure( - injector: Option<&dyn Fn(&str) -> bool>, + injector: Option<&FailureInjector>, phase: &str, + artifact_index: usize, ) -> Result<(), EntityUndoError> { - if injector.is_some_and(|injector| injector(phase)) { + if injector.is_some_and(|injector| injector(phase, artifact_index)) { return Err(EntityUndoError::Refused(format!( - "injected failure after phase {phase}" + "injected failure after {phase} artifact {artifact_index}" ))); } Ok(()) @@ -504,6 +503,7 @@ fn undo_facets( target_id: &str, payload: &Value, rollback: &mut MergeRollback, + injector: Option<&FailureInjector>, ) -> Result<(), EntityUndoError> { let entries = payload .get("manifest") @@ -515,6 +515,7 @@ fn undo_facets( .ok_or_else(|| { EntityUndoError::Refused("merge payload facets entries are missing".to_owned()) })?; + let mut artifact_index = 0; for entry in entries { let facet = entry .get("facet") @@ -531,6 +532,8 @@ fn undo_facets( "move" => { rollback.capture(journal, &directory)?; restore_snapshot(journal, &JournalSnapshot::Missing { path: directory })?; + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; } "merge" => { let target_before = entry @@ -555,6 +558,8 @@ fn undo_facets( }, ) .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; let observations_before = entry .get("target_observations_before") .and_then(Value::as_array) @@ -592,6 +597,8 @@ fn undo_facets( }, )?; } + inject_failure(injector, "facets", artifact_index)?; + artifact_index += 1; } _ => { return Err(EntityUndoError::Refused(format!( @@ -607,6 +614,7 @@ fn undo_observation_relations( journal: &Path, payload: &Value, rollback: &mut MergeRollback, + injector: Option<&FailureInjector>, ) -> Result<(), EntityUndoError> { let entries = payload .get("manifest") @@ -620,7 +628,7 @@ fn undo_observation_relations( "merge payload observation relation entries are missing".to_owned(), ) })?; - for entry in entries { + for (artifact_index, entry) in entries.iter().enumerate() { 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(), @@ -667,6 +675,7 @@ fn undo_observation_relations( rollback.capture(journal, path)?; write_jsonl(destination, rows, AtomicWriteOptions::default()) .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + inject_failure(injector, "observations", artifact_index)?; } Ok(()) } @@ -733,9 +742,10 @@ fn undo_segments( journal: &Path, payload: &Value, rollback: &mut MergeRollback, + injector: Option<&FailureInjector>, ) -> Result<(), EntityUndoError> { let entries = manifest_entries(payload, "segments")?; - for entry in entries { + for (artifact_index, entry) in entries.iter().enumerate() { let path = entry_path(entry, "segment")?; let section = entry .get("section") @@ -778,6 +788,7 @@ fn undo_segments( }, ) .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + inject_failure(injector, "segments", artifact_index)?; } Ok(()) } @@ -786,9 +797,10 @@ fn undo_activities( journal: &Path, payload: &Value, rollback: &mut MergeRollback, + injector: Option<&FailureInjector>, ) -> Result<(), EntityUndoError> { let entries = manifest_entries(payload, "activities")?; - for entry in entries { + for (artifact_index, entry) in entries.iter().enumerate() { let path = entry_path(entry, "activity")?; let row_index = entry_index(entry, "row_index", "activity")?; let container = entry @@ -842,6 +854,7 @@ fn undo_activities( rollback.capture(journal, path)?; write_jsonl(destination, rows, AtomicWriteOptions::default()) .map_err(|error| EntityUndoError::Refused(error.to_string()))?; + inject_failure(injector, "activities", artifact_index)?; } Ok(()) } diff --git a/core/crates/solstone-core-entity/src/undo_tests.rs b/core/crates/solstone-core-entity/src/undo_tests.rs index 71a80f287..5a8f648c3 100644 --- a/core/crates/solstone-core-entity/src/undo_tests.rs +++ b/core/crates/solstone-core-entity/src/undo_tests.rs @@ -759,7 +759,7 @@ fn facets_undo_injection_rolls_back_and_retry_succeeds() { &journal, &merge.merge_id, Value::Null, - Some(&|phase| phase == "facets"), + Some(&|phase, artifact_index| phase == "facets" && artifact_index == 0), ) .is_err() ); @@ -809,7 +809,7 @@ fn observations_undo_injection_rolls_back_and_retry_succeeds() { &journal, &merge.merge_id, Value::Null, - Some(&|phase| phase == "observations"), + Some(&|phase, artifact_index| phase == "observations" && artifact_index == 0), ) .is_err() ); @@ -852,7 +852,7 @@ fn segments_undo_injection_rolls_back_and_retry_succeeds() { &journal, &merge.merge_id, Value::Null, - Some(&|phase| phase == "segments"), + Some(&|phase, artifact_index| phase == "segments" && artifact_index == 0), ) .is_err() ); @@ -899,7 +899,7 @@ fn activities_undo_injection_rolls_back_and_retry_succeeds() { &journal, &merge.merge_id, Value::Null, - Some(&|phase| phase == "activities"), + Some(&|phase, artifact_index| phase == "activities" && artifact_index == 0), ) .is_err() ); -- 2.51.2