From 86521bb35596ed5c05f95915b89cdbdf9ab514cb Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Wed, 5 Aug 2026 15:49:41 -0600 Subject: [PATCH] Record complete entity merge metadata Track appended facet observations and per-row voiceprint support in merge payloads and audit results. Surface rollback causes in merge and undo errors while removing obsolete lints and duplicate I/O allowlist entries. --- .../solstone-core-entity/src/fixture_tests.rs | 3 -- .../solstone-core-entity/src/merge_tests.rs | 43 +++++++++++++++- .../solstone-core-entity/src/store/merge.rs | 49 ++++++++++++++++--- .../src/store/merge_payload.rs | 4 -- .../solstone-core-entity/src/store/mod.rs | 1 - .../solstone-core-entity/src/store/undo.rs | 10 ++++ .../solstone-core-entity/src/undo_tests.rs | 18 ++++--- 7 files changed, 104 insertions(+), 24 deletions(-) diff --git a/core/crates/solstone-core-entity/src/fixture_tests.rs b/core/crates/solstone-core-entity/src/fixture_tests.rs index b4a2149be..e79a83ac2 100644 --- a/core/crates/solstone-core-entity/src/fixture_tests.rs +++ b/core/crates/solstone-core-entity/src/fixture_tests.rs @@ -57,9 +57,6 @@ const JOURNAL_IO_ALLOWED: &[&str] = &[ "AppendError", "atomic_replace", "read_bytes", - "read_json", - "read_jsonl", - "write_json", "write_jsonl", ]; diff --git a/core/crates/solstone-core-entity/src/merge_tests.rs b/core/crates/solstone-core-entity/src/merge_tests.rs index 785002777..fa000d54b 100644 --- a/core/crates/solstone-core-entity/src/merge_tests.rs +++ b/core/crates/solstone-core-entity/src/merge_tests.rs @@ -524,6 +524,23 @@ fn commit_records_matching_payload_and_audit_counts() { 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(); + fs::write( + facet.join("observations.jsonl"), + b"{\"content\":\"source observation\",\"observed_at\":\"2026-01-01\"}\n", + ) + .unwrap(); + let target_facet = journal.join("facets/work/entities/target"); + fs::create_dir_all(&target_facet).unwrap(); + fs::write( + target_facet.join("entity.json"), + br#"{"entity_id":"target"}"#, + ) + .unwrap(); + fs::write( + target_facet.join("observations.jsonl"), + b"{\"content\":\"target observation\",\"observed_at\":\"2026-01-01\"}\n", + ) + .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(); @@ -538,7 +555,8 @@ fn commit_records_matching_payload_and_audit_counts() { 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"]["facets"]["merged"].as_u64().unwrap() > 0); + assert_eq!(audit["counts"]["facets"]["observations_appended"], 1); assert!(audit["counts"]["voiceprints"]["added"].as_u64().unwrap() > 0); assert!( audit["counts"]["activities"]["records_rewritten"] @@ -552,6 +570,25 @@ fn commit_records_matching_payload_and_audit_counts() { .unwrap(); let payload = load_entity_merge_payload(&journal, "target", &merge_id).unwrap(); assert_eq!(payload["result_counts"], audit["counts"]); + assert_eq!( + payload["manifest"]["voiceprints"]["support"] + .as_array() + .unwrap() + .len(), + 1 + ); + assert_eq!( + payload["manifest"]["voiceprints"]["support"][0]["key"]["sentence_id"], + "count" + ); + assert_eq!( + payload["manifest"]["voiceprints"]["support"][0]["target_preexisting"], + false + ); + assert_eq!( + payload["manifest"]["voiceprints"]["support"][0]["added"], + true + ); fs::remove_dir_all(journal).unwrap(); } @@ -1011,6 +1048,10 @@ fn audit_phase_injection_rolls_back_and_retry_succeeds() { Some(&|phase, artifact_index| phase == "audit" && artifact_index == 0), ) .unwrap_err(); + assert_eq!( + error.to_string(), + "entity merge failed during audit: injected failure after audit artifact 0" + ); let merge_id = match error { crate::EntityMergeError::Failed { report, .. } => report.merge_id, _ => panic!("expected injected merge failure"), diff --git a/core/crates/solstone-core-entity/src/store/merge.rs b/core/crates/solstone-core-entity/src/store/merge.rs index 588ba025e..6476850b2 100644 --- a/core/crates/solstone-core-entity/src/store/merge.rs +++ b/core/crates/solstone-core-entity/src/store/merge.rs @@ -82,20 +82,21 @@ pub struct EntityMergeReport { pub emails_added: usize, } -#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +#[derive(Debug, Clone, PartialEq, Default)] pub(crate) struct VoiceprintMergeStats { pub added: usize, pub skipped_duplicate: usize, pub target_total: usize, + pub support: Vec, } #[derive(Debug, Clone, PartialEq, Eq, Default)] pub(crate) struct FacetMergeStats { pub moved_count: usize, pub merged_count: usize, + pub observations_appended: usize, pub touched_facets: Vec, pub entries: Vec, } -#[allow(dead_code)] #[derive(Debug, Default)] struct MergeStats { voiceprints_added: usize, @@ -149,6 +150,13 @@ impl fmt::Display for EntityMergeError { Self::Snapshot(error) => error.fmt(f), Self::Index(error) => error.fmt(f), Self::Audit(error) => error.fmt(f), + Self::Failed { + failed_phase, + rollback_error: Some(error), + .. + } => { + write!(f, "entity merge failed during {failed_phase}: {error}") + } Self::Failed { failed_phase, .. } => { write!(f, "entity merge failed during {failed_phase}") } @@ -241,8 +249,8 @@ pub(crate) fn commit_entity_merge_with_injector( 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), 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); }), + "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; payload["manifest"]["voiceprints"]["support"] = Value::Array(result.support); }), + "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; stats.facets_observations_appended=result.observations_appended; 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(()) }), @@ -768,9 +776,12 @@ pub(crate) fn merge_facets( .map_err(|error| EntityMergeError::Refused(error.to_string()))?; inject_failure(injector, "facets", artifact_index)?; artifact_index += 1; + let merged_observations = dedupe_observations(&source_obs, &target_obs); + stats.observations_appended += + merged_observations.len().saturating_sub(target_obs.len()); write_jsonl( target_obs_path, - dedupe_observations(&source_obs, &target_obs), + merged_observations, AtomicWriteOptions::default(), ) .map_err(|error| EntityMergeError::Refused(error.to_string()))?; @@ -907,11 +918,16 @@ pub(crate) fn merge_voiceprints( .iter() .map(|metadata| voiceprint_key(metadata)) .collect::, _>>()?; + let target_existing = existing.clone(); 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) { + let target_preexisting = target_existing.contains(&key); + if !existing.insert(key.clone()) { stats.skipped_duplicate += 1; + stats + .support + .push(voiceprint_support_entry(&key, target_preexisting, false)); continue; } let norm = embedding @@ -925,6 +941,13 @@ pub(crate) fn merge_voiceprints( .extend(embedding.iter().map(|value| value / norm)); target.metadata.push(metadata.clone()); stats.added += 1; + stats + .support + .push(voiceprint_support_entry(&key, target_preexisting, true)); + } else { + stats + .support + .push(voiceprint_support_entry(&key, target_preexisting, false)); } } target.rows = target.metadata.len(); @@ -938,6 +961,19 @@ pub(crate) fn merge_voiceprints( Ok(stats) } +fn voiceprint_support_entry(key: &VoiceprintKey, target_preexisting: bool, added: bool) -> Value { + json!({ + "key": { + "day": &key.0, + "segment_key": &key.1, + "source": &key.2, + "sentence_id": &key.3, + }, + "target_preexisting": target_preexisting, + "added": added, + }) +} + fn load_voiceprints(path: &Path) -> Result, EntityMergeError> { if !path_lexists(path).map_err(|error| EntityMergeError::Refused(error.to_string()))? { return Ok(None); @@ -1140,7 +1176,6 @@ pub(crate) fn dedupe_emails(target_values: &[String], source_values: &[String]) .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(); diff --git a/core/crates/solstone-core-entity/src/store/merge_payload.rs b/core/crates/solstone-core-entity/src/store/merge_payload.rs index 2803b2176..9da3c3fe2 100644 --- a/core/crates/solstone-core-entity/src/store/merge_payload.rs +++ b/core/crates/solstone-core-entity/src/store/merge_payload.rs @@ -65,7 +65,6 @@ pub(crate) fn record_entity_merge_payload( Ok(payload_relative_path(entity_id, merge_id)) } -#[allow(dead_code)] pub(crate) fn load_entity_merge_payload( journal: &Path, entity_id: &str, @@ -93,7 +92,6 @@ pub(crate) fn load_entity_merge_payload( Ok(payload) } -#[allow(dead_code)] pub(crate) fn move_entity_merge_payload( journal: &Path, source_id: &str, @@ -119,7 +117,6 @@ pub(crate) fn move_entity_merge_payload( Ok((payload, target_rel)) } -#[allow(dead_code)] pub(crate) fn remove_entity_merge_payload( journal: &Path, entity_id: &str, @@ -154,7 +151,6 @@ fn payload_relative_path(entity_id: &str, merge_id: &str) -> String { format!("entities/{entity_id}/history/private/{merge_id}.json") } -#[allow(dead_code)] pub(crate) fn list_entity_merge_payload_ids( journal: &Path, entity_id: &str, diff --git a/core/crates/solstone-core-entity/src/store/mod.rs b/core/crates/solstone-core-entity/src/store/mod.rs index efb49b0a7..3cd674b57 100644 --- a/core/crates/solstone-core-entity/src/store/mod.rs +++ b/core/crates/solstone-core-entity/src/store/mod.rs @@ -15,7 +15,6 @@ mod paths; mod reconcile; mod repair; mod undo; -#[allow(dead_code)] pub(crate) mod voiceprints; mod write; diff --git a/core/crates/solstone-core-entity/src/store/undo.rs b/core/crates/solstone-core-entity/src/store/undo.rs index eb466a3e0..2942e4df8 100644 --- a/core/crates/solstone-core-entity/src/store/undo.rs +++ b/core/crates/solstone-core-entity/src/store/undo.rs @@ -67,6 +67,16 @@ impl fmt::Display for EntityUndoError { Self::Write(error) => error.fmt(formatter), Self::Snapshot(error) => error.fmt(formatter), Self::Index(error) => error.fmt(formatter), + Self::Failed { + failed_phase, + rollback_error: Some(error), + .. + } => { + write!( + formatter, + "entity merge undo failed during {failed_phase}: {error}" + ) + } Self::Failed { failed_phase, .. } => { write!(formatter, "entity merge undo failed during {failed_phase}") } diff --git a/core/crates/solstone-core-entity/src/undo_tests.rs b/core/crates/solstone-core-entity/src/undo_tests.rs index 5a8f648c3..722e2d1d8 100644 --- a/core/crates/solstone-core-entity/src/undo_tests.rs +++ b/core/crates/solstone-core-entity/src/undo_tests.rs @@ -754,14 +754,16 @@ fn facets_undo_injection_rolls_back_and_retry_succeeds() { 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, artifact_index| phase == "facets" && artifact_index == 0), - ) - .is_err() + let error = undo_entity_merge_with_injector( + &journal, + &merge.merge_id, + Value::Null, + Some(&|phase, artifact_index| phase == "facets" && artifact_index == 0), + ) + .unwrap_err(); + assert_eq!( + error.to_string(), + "entity merge undo failed during facets: injected failure after facets artifact 0" ); assert_eq!(fs::read(&moved_target).unwrap(), moved_after_merge); assert_eq!( -- 2.51.2