From 9d0bacbb9e7f04cef790bc05a26f355ecaf9eacb Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Wed, 5 Aug 2026 14:40:10 -0600 Subject: [PATCH] Add merge payload validation Persist, validate, rebase, and list private merge payloads with rollback snapshot support. --- .../src/merge_payload_tests.rs | 438 ++++++++++++++++++ .../src/store/merge_payload.rs | 366 +++++++++++++++ .../src/store/merge_rollback.rs | 28 ++ 3 files changed, 832 insertions(+) create mode 100644 core/crates/solstone-core-entity/src/merge_payload_tests.rs create mode 100644 core/crates/solstone-core-entity/src/store/merge_payload.rs create mode 100644 core/crates/solstone-core-entity/src/store/merge_rollback.rs diff --git a/core/crates/solstone-core-entity/src/merge_payload_tests.rs b/core/crates/solstone-core-entity/src/merge_payload_tests.rs new file mode 100644 index 000000000..cf1308112 --- /dev/null +++ b/core/crates/solstone-core-entity/src/merge_payload_tests.rs @@ -0,0 +1,438 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::fs; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU64, Ordering}; + +use serde_json::{Value, json}; + +use super::store::merge_payload::{ + load_entity_merge_payload, move_entity_merge_payload, record_entity_merge_payload, + validate_merge_payload, +}; + +static NEXT_DIRECTORY: AtomicU64 = AtomicU64::new(0); + +struct TempDir(PathBuf); +impl TempDir { + fn new() -> Self { + let path = std::env::temp_dir().join(format!( + "solstone-merge-payload-{}-{}", + std::process::id(), + NEXT_DIRECTORY.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&path).unwrap(); + Self(path) + } +} +impl Drop for TempDir { + fn drop(&mut self) { + fs::remove_dir_all(&self.0).unwrap(); + } +} + +fn payload() -> Value { + json!({"source_id":"source","target_id":"target","source_state":{"snapshots":[]},"manifest":{"identity":{"aka_support":[],"email_support":[],"scalar_support":[],"target_before":{}},"voiceprints":{"support":[]},"facets":{"entries":[]},"segments":{"entries":[]},"activities":{"entries":[]},"observation_relations":{"entries":[]},"rebased_merge_ids":[]}}) +} + +fn message(value: Value) -> String { + validate_merge_payload(Path::new("/tmp/journal"), &value) + .unwrap_err() + .to_string() +} + +#[test] +fn validator_reports_source_identity_message() { + let mut value = payload(); + value.as_object_mut().unwrap().remove("source_id"); + assert_eq!(message(value), "merge payload missing source entity id"); +} + +#[test] +fn validator_reports_target_identity_message() { + let mut value = payload(); + value.as_object_mut().unwrap().remove("target_id"); + assert_eq!(message(value), "merge payload missing target entity id"); +} + +#[test] +fn validator_reports_source_state_messages() { + let mut absent = payload(); + absent.as_object_mut().unwrap().remove("source_state"); + assert_eq!(message(absent), "merge payload missing source_state"); + let mut non_object = payload(); + non_object["source_state"] = json!(false); + assert_eq!( + message(non_object), + "merge payload source_state is not an object" + ); + let mut missing = payload(); + missing["source_state"] + .as_object_mut() + .unwrap() + .remove("snapshots"); + assert_eq!(message(missing), "merge payload missing snapshots"); + let mut non_list = payload(); + non_list["source_state"]["snapshots"] = json!(false); + assert_eq!(message(non_list), "merge payload snapshots is not a list"); +} + +#[test] +fn well_formed_payload_validates() { + validate_merge_payload(Path::new("/tmp/journal"), &payload()).unwrap(); +} + +#[test] +fn validator_reports_non_object_payload() { + assert_eq!(message(json!(false)), "merge payload is not an object"); +} + +#[test] +fn validator_reports_snapshot_messages() { + let mut non_object = payload(); + non_object["source_state"]["snapshots"] = json!([false]); + assert_eq!( + message(non_object), + "merge payload snapshot is not an object" + ); + let mut missing_rel = payload(); + missing_rel["source_state"]["snapshots"] = json!([{}]); + assert_eq!( + message(missing_rel), + "manifest snapshot missing relative path" + ); + let mut files_not_list = payload(); + files_not_list["source_state"]["snapshots"] = json!([{"rel":"entities/source","files":false}]); + assert_eq!( + message(files_not_list), + "manifest snapshot files is not a list" + ); + let mut file_not_object = payload(); + file_not_object["source_state"]["snapshots"] = + json!([{"rel":"entities/source","files":[false]}]); + assert_eq!( + message(file_not_object), + "manifest snapshot file is not an object" + ); + let mut missing_file_rel = payload(); + missing_file_rel["source_state"]["snapshots"] = json!([{"rel":"entities/source","files":[{}]}]); + assert_eq!( + message(missing_file_rel), + "manifest snapshot file missing relative path" + ); +} + +#[test] +fn validator_reports_manifest_messages() { + let mut absent = payload(); + absent.as_object_mut().unwrap().remove("manifest"); + assert_eq!(message(absent), "merge payload missing manifest"); + let mut non_object = payload(); + non_object["manifest"] = json!(false); + assert_eq!( + message(non_object), + "merge payload manifest is not an object" + ); + let mut missing = payload(); + missing["manifest"] + .as_object_mut() + .unwrap() + .remove("identity"); + assert_eq!(message(missing), "merge payload missing identity manifest"); + let mut non_object = payload(); + non_object["manifest"]["identity"] = json!(false); + assert_eq!( + message(non_object), + "merge payload identity manifest is not an object" + ); +} + +#[test] +fn validator_reports_identity_support_messages() { + for field in ["aka_support", "email_support", "scalar_support"] { + let mut absent = payload(); + absent["manifest"]["identity"] + .as_object_mut() + .unwrap() + .remove(field); + assert_eq!( + message(absent), + format!("merge payload identity missing {field}") + ); + let mut non_list = payload(); + non_list["manifest"]["identity"][field] = json!(false); + assert_eq!( + message(non_list), + format!("merge payload identity {field} is not a list") + ); + } +} + +#[test] +fn validator_reports_identity_target_before_messages() { + let mut absent = payload(); + absent["manifest"]["identity"] + .as_object_mut() + .unwrap() + .remove("target_before"); + assert_eq!( + message(absent), + "merge payload identity missing target_before" + ); + let mut non_object = payload(); + non_object["manifest"]["identity"]["target_before"] = json!(false); + assert_eq!( + message(non_object), + "merge payload identity target_before is not an object" + ); +} + +#[test] +fn validator_reports_voiceprint_messages() { + let mut absent = payload(); + absent["manifest"] + .as_object_mut() + .unwrap() + .remove("voiceprints"); + assert_eq!( + message(absent), + "merge payload missing voiceprints manifest" + ); + let mut non_object = payload(); + non_object["manifest"]["voiceprints"] = json!(false); + assert_eq!( + message(non_object), + "merge payload voiceprints manifest is not an object" + ); + let mut support = payload(); + support["manifest"]["voiceprints"] + .as_object_mut() + .unwrap() + .remove("support"); + assert_eq!( + message(support), + "merge payload voiceprints missing support" + ); + let mut non_list = payload(); + non_list["manifest"]["voiceprints"]["support"] = json!(false); + assert_eq!( + message(non_list), + "merge payload voiceprints support is not a list" + ); +} + +#[test] +fn validator_reports_facet_messages() { + let mut absent = payload(); + absent["manifest"].as_object_mut().unwrap().remove("facets"); + assert_eq!(message(absent), "merge payload missing facets manifest"); + let mut non_object = payload(); + non_object["manifest"]["facets"] = json!(false); + assert_eq!( + message(non_object), + "merge payload facets manifest is not an object" + ); + let mut missing_entries = payload(); + missing_entries["manifest"]["facets"] + .as_object_mut() + .unwrap() + .remove("entries"); + assert_eq!( + message(missing_entries), + "merge payload facets missing entries" + ); + let mut non_list = payload(); + non_list["manifest"]["facets"]["entries"] = json!(false); + assert_eq!( + message(non_list), + "merge payload facets entries is not a list" + ); + let mut non_object_entry = payload(); + non_object_entry["manifest"]["facets"]["entries"] = json!([false]); + assert_eq!( + message(non_object_entry), + "merge payload facet entry is not an object" + ); + let mut missing_name = payload(); + missing_name["manifest"]["facets"]["entries"] = json!([{}]); + assert_eq!( + message(missing_name), + "manifest facet entry missing facet name" + ); + let mut empty_name = payload(); + empty_name["manifest"]["facets"]["entries"] = json!([{"facet":""}]); + assert_eq!( + message(empty_name), + "manifest facet entry missing facet name" + ); +} + +#[test] +fn validator_reports_path_section_messages() { + for section in ["segments", "activities", "observation_relations"] { + let mut absent = payload(); + absent["manifest"].as_object_mut().unwrap().remove(section); + assert_eq!( + message(absent), + format!("merge payload missing {section} manifest") + ); + let mut non_object = payload(); + non_object["manifest"][section] = json!(false); + assert_eq!( + message(non_object), + format!("merge payload {section} manifest is not an object") + ); + let mut missing_entries = payload(); + missing_entries["manifest"][section] + .as_object_mut() + .unwrap() + .remove("entries"); + assert_eq!( + message(missing_entries), + format!("merge payload {section} missing entries") + ); + let mut non_list = payload(); + non_list["manifest"][section]["entries"] = json!(false); + assert_eq!( + message(non_list), + format!("merge payload {section} entries is not a list") + ); + let mut non_object_entry = payload(); + non_object_entry["manifest"][section]["entries"] = json!([false]); + assert_eq!( + message(non_object_entry), + format!("merge payload {section} entry is not an object") + ); + let mut missing_path = payload(); + missing_path["manifest"][section]["entries"] = json!([{}]); + assert_eq!( + message(missing_path), + format!("manifest {section} entry missing path") + ); + } +} + +#[test] +fn validator_reports_contained_path_messages() { + let expected = "journal path contains invalid component"; + let mut source = payload(); + source["source_id"] = json!("../outside"); + assert_eq!(message(source), expected); + let mut target = payload(); + target["target_id"] = json!("../outside"); + assert_eq!(message(target), expected); + let mut snapshot = payload(); + snapshot["source_state"]["snapshots"] = json!([{"rel":"../outside","files":[]}]); + assert_eq!(message(snapshot), expected); + let mut snapshot_file = payload(); + snapshot_file["source_state"]["snapshots"] = + json!([{"rel":"entities/source","files":[{"rel":"../../outside"}]}]); + assert_eq!(message(snapshot_file), expected); + let mut facet = payload(); + facet["manifest"]["facets"]["entries"] = json!([{"facet":"../outside"}]); + assert_eq!(message(facet), expected); + for section in ["segments", "activities", "observation_relations"] { + let mut value = payload(); + value["manifest"][section]["entries"] = json!([{"path":"../outside"}]); + assert_eq!(message(value), expected); + } +} + +#[test] +fn validator_reports_rebased_messages_from_manifest() { + let mut absent = payload(); + absent["manifest"] + .as_object_mut() + .unwrap() + .remove("rebased_merge_ids"); + assert_eq!(message(absent), "merge payload missing rebased_merge_ids"); + let mut non_list = payload(); + non_list["manifest"]["rebased_merge_ids"] = json!(false); + assert_eq!( + message(non_list), + "merge payload rebased_merge_ids is not a list" + ); +} + +#[test] +fn load_revalidates_tampered_payload() { + let directory = TempDir::new(); + let value = payload(); + record_entity_merge_payload(&directory.0, "source", "em_1", &value).unwrap(); + let path = directory + .0 + .join("entities/source/history/private/em_1.json"); + let mut tampered: Value = serde_json::from_slice(&fs::read(&path).unwrap()).unwrap(); + tampered["manifest"] + .as_object_mut() + .unwrap() + .remove("rebased_merge_ids"); + fs::write(&path, serde_json::to_vec(&tampered).unwrap()).unwrap(); + assert_eq!( + load_entity_merge_payload(&directory.0, "source", "em_1") + .unwrap_err() + .to_string(), + "invalid private merge payload for source:em_1: merge payload missing rebased_merge_ids" + ); +} + +#[test] +fn load_distinguishes_missing_and_non_object_payloads() { + let directory = TempDir::new(); + assert_eq!( + load_entity_merge_payload(&directory.0, "source", "missing") + .unwrap_err() + .to_string(), + "missing private merge payload for source: missing" + ); + let path = directory + .0 + .join("entities/source/history/private/non-object.json"); + fs::create_dir_all(path.parent().unwrap()).unwrap(); + fs::write(&path, b"[]").unwrap(); + assert_eq!( + load_entity_merge_payload(&directory.0, "source", "non-object") + .unwrap_err() + .to_string(), + format!("private merge payload is not an object: {}", path.display()) + ); +} + +#[test] +fn move_rebases_and_removes_only_a_distinct_source() { + let directory = TempDir::new(); + let value = payload(); + record_entity_merge_payload(&directory.0, "source", "em_1", &value).unwrap(); + let (moved, rel) = + move_entity_merge_payload(&directory.0, "source", "target", "em_1", Some("ancestor")) + .unwrap(); + assert_eq!(rel, "entities/target/history/private/em_1.json"); + assert_eq!(moved["target_id"], "target"); + assert_eq!(moved["rebased_from_entity_id"], "ancestor"); + assert!( + !directory + .0 + .join("entities/source/history/private/em_1.json") + .exists() + ); + assert!( + directory + .0 + .join("entities/target/history/private/em_1.json") + .is_file() + ); + let (_, _) = move_entity_merge_payload(&directory.0, "target", "target", "em_1", None).unwrap(); + let same = load_entity_merge_payload(&directory.0, "target", "em_1").unwrap(); + assert!(same.get("rebased_from_entity_id").is_some()); + record_entity_merge_payload(&directory.0, "target", "em_2", &payload()).unwrap(); + let (same, _) = + move_entity_merge_payload(&directory.0, "target", "target", "em_2", None).unwrap(); + assert!(same.get("rebased_from_entity_id").is_none()); + assert!( + directory + .0 + .join("entities/target/history/private/em_2.json") + .is_file() + ); +} diff --git a/core/crates/solstone-core-entity/src/store/merge_payload.rs b/core/crates/solstone-core-entity/src/store/merge_payload.rs new file mode 100644 index 000000000..4caf7628f --- /dev/null +++ b/core/crates/solstone-core-entity/src/store/merge_payload.rs @@ -0,0 +1,366 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::error::Error; +use std::fmt; +use std::path::Path; + +use serde_json::Value; +use solstone_core_journal_io::AtomicWriteError; +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::PathError; +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::restore_snapshot; +use solstone_core_journal_io::write_json; + +#[derive(Debug)] +pub enum MergePayloadError { + Path(PathError), + Read(solstone_core_journal_io::ReadError), + Write(AtomicWriteError), + Snapshot(SnapshotError), + Invalid(String), +} + +impl fmt::Display for MergePayloadError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Path(error) => error.fmt(formatter), + Self::Read(error) => error.fmt(formatter), + Self::Write(error) => error.fmt(formatter), + Self::Snapshot(error) => error.fmt(formatter), + Self::Invalid(message) => formatter.write_str(message), + } + } +} +impl Error for MergePayloadError {} + +pub(crate) fn record_entity_merge_payload( + journal: &Path, + entity_id: &str, + merge_id: &str, + payload: &Value, +) -> Result { + validate_merge_payload(journal, payload)?; + let path = payload_path(journal, entity_id, merge_id)?; + write_json( + &path, + payload, + JsonWriteOptions { + indent: Some(2), + sort_keys: false, + mode: None, + }, + ) + .map_err(MergePayloadError::Write)?; + Ok(payload_relative_path(entity_id, merge_id)) +} + +#[allow(dead_code)] +pub(crate) fn load_entity_merge_payload( + journal: &Path, + entity_id: &str, + merge_id: &str, +) -> Result { + let path = payload_path(journal, entity_id, merge_id)?; + if !path_lexists(&path).map_err(MergePayloadError::Path)? { + return Err(invalid(&format!( + "missing private merge payload for {entity_id}: {merge_id}" + ))); + } + let payload = + read_json(&path, Value::Null, MalformedPolicy::Raise).map_err(MergePayloadError::Read)?; + if !payload.is_object() { + return Err(invalid(&format!( + "private merge payload is not an object: {}", + path.display() + ))); + } + validate_merge_payload(journal, &payload).map_err(|error| { + MergePayloadError::Invalid(format!( + "invalid private merge payload for {entity_id}:{merge_id}: {error}" + )) + })?; + Ok(payload) +} + +#[allow(dead_code)] +pub(crate) fn move_entity_merge_payload( + journal: &Path, + source_id: &str, + target_id: &str, + merge_id: &str, + rebased_from_entity_id: Option<&str>, +) -> Result<(Value, String), MergePayloadError> { + let mut payload = load_entity_merge_payload(journal, source_id, merge_id)?; + payload + .as_object_mut() + .ok_or_else(|| MergePayloadError::Invalid("merge payload is not an object".to_owned()))? + .insert("target_id".to_owned(), Value::String(target_id.to_owned())); + if let Some(rebased_from_entity_id) = rebased_from_entity_id { + payload.as_object_mut().expect("validated object").insert( + "rebased_from_entity_id".to_owned(), + Value::String(rebased_from_entity_id.to_owned()), + ); + } + let target_rel = record_entity_merge_payload(journal, target_id, merge_id, &payload)?; + if source_id != target_id { + remove_entity_merge_payload(journal, source_id, merge_id)?; + } + Ok((payload, target_rel)) +} + +#[allow(dead_code)] +pub(crate) fn remove_entity_merge_payload( + journal: &Path, + entity_id: &str, + merge_id: &str, +) -> Result<(), MergePayloadError> { + let path = payload_path(journal, entity_id, merge_id)?; + if path_lexists(&path).map_err(MergePayloadError::Path)? { + restore_snapshot( + journal, + &JournalSnapshot::Missing { + path: payload_relative_path(entity_id, merge_id), + }, + ) + .map_err(MergePayloadError::Snapshot)?; + } + Ok(()) +} + +fn payload_path( + journal: &Path, + entity_id: &str, + merge_id: &str, +) -> Result { + contained_path( + journal, + &format!("entities/{entity_id}/history/private/{merge_id}.json"), + ) + .map_err(MergePayloadError::Path) +} + +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, +) -> Result, MergePayloadError> { + let directory = contained_path(journal, &format!("entities/{entity_id}/history/private")) + .map_err(MergePayloadError::Path)?; + Ok(list_dir_entries(&directory) + .map_err(MergePayloadError::Path)? + .into_iter() + .filter(|entry| entry.kind == DirEntryKind::File) + .filter_map(|entry| { + entry + .name + .to_str() + .and_then(|name| name.strip_suffix(".json")) + .map(ToOwned::to_owned) + }) + .collect()) +} + +pub(crate) fn validate_merge_payload( + journal: &Path, + payload: &Value, +) -> Result<(), MergePayloadError> { + let object = payload + .as_object() + .ok_or_else(|| invalid("merge payload is not an object"))?; + let source_id = required_string( + object, + "source_id", + "merge payload missing source entity id", + )?; + let target_id = required_string( + object, + "target_id", + "merge payload missing target entity id", + )?; + contained_path(journal, &format!("entities/{source_id}")).map_err(MergePayloadError::Path)?; + contained_path(journal, &format!("entities/{target_id}")).map_err(MergePayloadError::Path)?; + let source_state = required_object( + object, + "source_state", + "merge payload missing source_state", + "merge payload source_state is not an object", + )?; + let snapshots = required_array( + source_state, + "snapshots", + "merge payload missing snapshots", + "merge payload snapshots is not a list", + )?; + for snapshot in snapshots { + let snapshot = snapshot + .as_object() + .ok_or_else(|| invalid("merge payload snapshot is not an object"))?; + let rel = optional_string(snapshot, "rel", "manifest snapshot missing relative path")?; + contained_path(journal, rel).map_err(MergePayloadError::Path)?; + if let Some(files) = snapshot.get("files") { + for file in files + .as_array() + .ok_or_else(|| invalid("manifest snapshot files is not a list"))? + { + let file = file + .as_object() + .ok_or_else(|| invalid("manifest snapshot file is not an object"))?; + let item_rel = + optional_string(file, "rel", "manifest snapshot file missing relative path")?; + contained_path(journal, &format!("{rel}/{item_rel}")) + .map_err(MergePayloadError::Path)?; + } + } + } + let manifest = required_object( + object, + "manifest", + "merge payload missing manifest", + "merge payload manifest is not an object", + )?; + let identity = required_object( + manifest, + "identity", + "merge payload missing identity manifest", + "merge payload identity manifest is not an object", + )?; + for field in ["aka_support", "email_support", "scalar_support"] { + required_array( + identity, + field, + &format!("merge payload identity missing {field}"), + &format!("merge payload identity {field} is not a list"), + )?; + } + required_object( + identity, + "target_before", + "merge payload identity missing target_before", + "merge payload identity target_before is not an object", + )?; + let voiceprints = required_object( + manifest, + "voiceprints", + "merge payload missing voiceprints manifest", + "merge payload voiceprints manifest is not an object", + )?; + required_array( + voiceprints, + "support", + "merge payload voiceprints missing support", + "merge payload voiceprints support is not a list", + )?; + let facets = required_object( + manifest, + "facets", + "merge payload missing facets manifest", + "merge payload facets manifest is not an object", + )?; + let facets_entries = required_array( + facets, + "entries", + "merge payload facets missing entries", + "merge payload facets entries is not a list", + )?; + for entry in facets_entries { + let entry = entry + .as_object() + .ok_or_else(|| invalid("merge payload facet entry is not an object"))?; + let facet = required_string(entry, "facet", "manifest facet entry missing facet name")?; + contained_path(journal, &format!("facets/{facet}/entities/{target_id}")) + .map_err(MergePayloadError::Path)?; + } + for section in ["segments", "activities", "observation_relations"] { + let section_value = required_object( + manifest, + section, + &format!("merge payload missing {section} manifest"), + &format!("merge payload {section} manifest is not an object"), + )?; + let entries = required_array( + section_value, + "entries", + &format!("merge payload {section} missing entries"), + &format!("merge payload {section} entries is not a list"), + )?; + for entry in entries { + let entry = entry.as_object().ok_or_else(|| { + invalid(&format!("merge payload {section} entry is not an object")) + })?; + let path = optional_string( + entry, + "path", + &format!("manifest {section} entry missing path"), + )?; + contained_path(journal, path).map_err(MergePayloadError::Path)?; + } + } + required_array( + manifest, + "rebased_merge_ids", + "merge payload missing rebased_merge_ids", + "merge payload rebased_merge_ids is not a list", + )?; + Ok(()) +} +fn required_string<'a>( + object: &'a serde_json::Map, + key: &str, + missing: &str, +) -> Result<&'a str, MergePayloadError> { + object + .get(key) + .and_then(Value::as_str) + .filter(|value| !value.is_empty()) + .ok_or_else(|| invalid(missing)) +} +fn optional_string<'a>( + object: &'a serde_json::Map, + key: &str, + missing: &str, +) -> Result<&'a str, MergePayloadError> { + object + .get(key) + .and_then(Value::as_str) + .ok_or_else(|| invalid(missing)) +} +fn required_object<'a>( + object: &'a serde_json::Map, + key: &str, + missing: &str, + invalid_message: &str, +) -> Result<&'a serde_json::Map, MergePayloadError> { + object + .get(key) + .ok_or_else(|| invalid(missing))? + .as_object() + .ok_or_else(|| invalid(invalid_message)) +} +fn required_array<'a>( + object: &'a serde_json::Map, + key: &str, + missing: &str, + invalid_message: &str, +) -> Result<&'a Vec, MergePayloadError> { + object + .get(key) + .ok_or_else(|| invalid(missing))? + .as_array() + .ok_or_else(|| invalid(invalid_message)) +} +fn invalid(message: &str) -> MergePayloadError { + MergePayloadError::Invalid(message.to_owned()) +} diff --git a/core/crates/solstone-core-entity/src/store/merge_rollback.rs b/core/crates/solstone-core-entity/src/store/merge_rollback.rs new file mode 100644 index 000000000..8991c927f --- /dev/null +++ b/core/crates/solstone-core-entity/src/store/merge_rollback.rs @@ -0,0 +1,28 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::path::Path; + +use solstone_core_journal_io::JournalSnapshot; +use solstone_core_journal_io::SnapshotError; +use solstone_core_journal_io::capture_snapshot; +use solstone_core_journal_io::restore_snapshot; + +#[derive(Debug, Default)] +pub(crate) struct MergeRollback { + snapshots: Vec, +} + +impl MergeRollback { + pub(super) fn capture(&mut self, journal: &Path, path: &str) -> Result<(), SnapshotError> { + self.snapshots.push(capture_snapshot(journal, path)?); + Ok(()) + } + + pub(super) fn restore(&self, journal: &Path) -> Result<(), SnapshotError> { + for snapshot in self.snapshots.iter().rev() { + restore_snapshot(journal, snapshot)?; + } + Ok(()) + } +} -- 2.51.2