diff --git a/readme.md b/readme.md index 5bcde68..cff439a 100644 --- a/readme.md +++ b/readme.md @@ -74,6 +74,34 @@ fn main() { } ``` +### Application-supplied merge policies + +`Repo::with_merge_policy` supplies an implementation of `MergePolicy` for a +named custom policy. A locally created history records its selected policy in +the genesis metadata. Subsequent automatic merges, including `merge_heads`, +select the implementation named by that genesis; installing a different custom +policy does not override an existing history's choice. Histories using `lww` +continue to use the built-in LWW implementation. An unavailable custom policy +causes a merge error rather than a fallback to a different rule. + +The application owns the custom policy's semantics and must deploy compatible, +convergent implementations to its replicas. Policy names are identifiers, not +proof that two implementations behave identically: use a new name when changing +semantics. This API does not migrate existing histories to a new policy, nor does +it prove that an imported merge payload was computed correctly. + +Parent-dependent policies require all immediate-parent payloads before a merge +can be saved. If a parent has not synced, import it and retry. Policies independent +of parents can opt out with `MergePolicy::requires_parent_payloads`. + +Replication must preserve `Operation::node_metadata` as well as the node +timestamp. The metadata is the exact serialized value, not just a policy name; +replacing it with local defaults can change the node's CID. Legacy operations +without this field remain readable. Newly written LevelDB operations use a +versioned storage format: this version reads old records, but older binaries +cannot read the new records. Do not downgrade a database after writing with this +version. Older replicas are not compatible with new custom-policy histories. + ## 🖥️ CLI Tool CRSL includes a command-line interface for easy content management. diff --git a/src/convergence/policies/lww.rs b/src/convergence/policies/lww.rs index dc2892e..e77218e 100644 --- a/src/convergence/policies/lww.rs +++ b/src/convergence/policies/lww.rs @@ -7,6 +7,10 @@ use crate::convergence::policy::{MergePolicy, ResolveInput}; pub struct LwwMergePolicy; impl MergePolicy

for LwwMergePolicy { + fn requires_parent_payloads(&self) -> bool { + false + } + fn resolve(&self, nodes: &[ResolveInput

]) -> P { let winner = nodes .iter() diff --git a/src/convergence/policy.rs b/src/convergence/policy.rs index 827d105..7559fe7 100644 --- a/src/convergence/policy.rs +++ b/src/convergence/policy.rs @@ -6,6 +6,20 @@ pub struct ResolveInput

{ pub cid: Cid, pub payload: P, pub timestamp: u64, + /// Payloads of this head's parents, in the node's parent order. + /// + /// A policy that treats the payload as more than an opaque value needs to + /// know what a head *changed*, not just what it holds: a node that copied + /// its parent's value for one field while updating another should not win + /// a last-writer race on the field it left alone. The library cannot make + /// that distinction — it does not know the payload's shape — so it hands + /// the parents over and lets the policy compare. + /// + /// The resolver supplies all immediate parents when the policy requires + /// them, or returns an error if any are missing. Empty for a genesis head, + /// policies that opt out of parent payloads (such as `LwwMergePolicy`), + /// and inputs constructed by callers that have no DAG at hand. + pub parent_payloads: Vec

, } impl

ResolveInput

{ @@ -14,15 +28,49 @@ impl

ResolveInput

{ cid, payload, timestamp, + parent_payloads: Vec::new(), + } + } + + pub fn with_parents(cid: Cid, payload: P, timestamp: u64, parent_payloads: Vec

) -> Self { + Self { + cid, + payload, + timestamp, + parent_payloads, } } } /// A merge strategy that produces a converged payload from candidate nodes. +/// +/// The library ships [`LwwMergePolicy`](crate::convergence::policies::lww::LwwMergePolicy) +/// and selects it by the `policy_type` recorded in the genesis metadata. An +/// application whose payload is a composite — fields with different +/// convergence rules — supplies its own implementation through +/// [`Repo::with_merge_policy`](crate::repo::Repo::with_merge_policy); the +/// library calls it only when its name matches the genesis metadata. The +/// built-in `lww` name is reserved and always selects the library's LWW rule. pub trait MergePolicy

: Send + Sync { + /// Whether resolution needs the complete payloads of each head's immediate parents. + /// + /// Defaults to true: the resolver returns a missing-node error before + /// invoking `resolve` if any parent has not synced yet. The caller can + /// retry after importing the missing parents; no merge is persisted. + /// Return false only when resolution is independent of parent payloads. + /// In that case the resolver does not load them and supplies empty vectors. + fn requires_parent_payloads(&self) -> bool { + true + } + /// Resolve competing nodes into a single payload. fn resolve(&self, nodes: &[ResolveInput

]) -> P; - /// Return a descriptive name of the policy (e.g. "lww"). + /// Return the stable policy identifier recorded in genesis metadata. + /// + /// All implementations sharing a name must have identical deterministic + /// merge semantics across replicas. Use a new name when those semantics + /// change; installing it does not migrate existing content. `lww` is + /// reserved for the built-in rule, not an application override. fn name(&self) -> &str; } diff --git a/src/convergence/resolver.rs b/src/convergence/resolver.rs index 2fadac2..7594f27 100644 --- a/src/convergence/resolver.rs +++ b/src/convergence/resolver.rs @@ -52,7 +52,7 @@ where )); } - let inputs = self.collect_inputs(heads, dag)?; + let inputs = self.collect_inputs(heads, dag, policy.requires_parent_payloads())?; let merged_payload = policy.resolve(&inputs); let metadata = self.merge_metadata(heads, dag)?; Ok(Node::new_child( @@ -68,6 +68,7 @@ where &self, heads: &[Cid], dag: &DagGraph, + requires_parent_payloads: bool, ) -> CrdtResult>> where S: NodeStorage, @@ -80,10 +81,24 @@ where .get_node(&cid) .map_err(CrdtError::Graph)? .ok_or_else(|| CrdtError::Internal(format!("Head node not found: {cid}")))?; - inputs.push(ResolveInput::new( + let mut parent_payloads = Vec::new(); + if requires_parent_payloads { + parent_payloads.reserve(node.parents().len()); + for parent in node.parents() { + let parent_node = + dag.get_node(parent) + .map_err(CrdtError::Graph)? + .ok_or(CrdtError::Graph( + crate::graph::error::GraphError::NodeNotFound(*parent), + ))?; + parent_payloads.push(parent_node.payload().clone()); + } + } + inputs.push(ResolveInput::with_parents( cid, node.payload().clone(), node.timestamp(), + parent_payloads, )); } Ok(inputs) @@ -193,6 +208,63 @@ mod tests { Cid::new_v1(0x55, digest) } + /// Each `ResolveInput` handed to the policy carries its head's parent + /// payloads, so a policy can tell a head that changed the value from one + /// that only re-committed its parent's. + #[test] + fn create_merge_node_attaches_parent_payloads() { + struct Capture(std::sync::Mutex>>); + impl MergePolicy for Capture { + fn resolve(&self, nodes: &[ResolveInput]) -> String { + let mut seen = self.0.lock().unwrap(); + for n in nodes { + seen.push(n.parent_payloads.clone()); + } + nodes[0].payload.clone() + } + fn name(&self) -> &str { + "capture" + } + } + + let storage = MemoryNodeStorage::::default(); + let dag = DagGraph::new(storage.clone()); + let metadata = ContentMetadata::with_policy("capture"); + let genesis_node = Node::new_genesis("genesis".to_string(), 1, metadata.clone()); + let genesis_cid = genesis_node.content_id().unwrap(); + dag.storage.put(&genesis_node).unwrap(); + let head_a = Node::new_child( + "payload-a".to_string(), + vec![genesis_cid], + genesis_cid, + 10, + metadata.clone(), + ); + let head_a_cid = head_a.content_id().unwrap(); + dag.storage.put(&head_a).unwrap(); + let head_b = Node::new_child( + "payload-b".to_string(), + vec![genesis_cid], + genesis_cid, + 11, + metadata, + ); + let head_b_cid = head_b.content_id().unwrap(); + dag.storage.put(&head_b).unwrap(); + + let policy = Capture(std::sync::Mutex::new(Vec::new())); + let resolver = ConflictResolver::::new(); + resolver + .create_merge_node(&[head_a_cid, head_b_cid], &dag, genesis_cid, 20, &policy) + .unwrap(); + + let seen = policy.0.into_inner().unwrap(); + assert_eq!( + seen, + vec![vec!["genesis".to_string()], vec!["genesis".to_string()]] + ); + } + #[test] fn create_merge_node_merges_heads() { let storage = MemoryNodeStorage::::default(); diff --git a/src/crdt/operation.rs b/src/crdt/operation.rs index 846bc66..20ebf3d 100644 --- a/src/crdt/operation.rs +++ b/src/crdt/operation.rs @@ -2,6 +2,7 @@ use serde::{Deserialize, Serialize}; use std::fmt::Debug; use ulid::Ulid; +use crate::convergence::metadata::ContentMetadata; use crate::crdt::timestamp::next_monotonic_timestamp; /// Unique identifier for operations (based on Ulid) @@ -66,6 +67,17 @@ pub struct Operation { /// This ensures CID consistency across replicas. #[serde(default)] pub node_timestamp: Option, + /// Exact DAG metadata for replication. Missing on legacy operations. + /// A local Create may explicitly select metadata here; otherwise the Repo's + /// installed policy is recorded. Imported Creates never infer it from the + /// receiver's policy: absence means the historical default metadata. + /// Repo commits populate this for every node. Updates and Merges carrying + /// their payload can therefore reconstruct metadata even before ancestry + /// arrives; Deletes still require existing payload history. + /// Local Updates/Deletes inherit metadata instead of using this field to + /// change policy. Keep the exact metadata representation, not just its name. + #[serde(default)] + pub node_metadata: Option, } impl Operation @@ -95,6 +107,7 @@ where author, parents: Vec::new(), node_timestamp: None, + node_metadata: None, } } diff --git a/src/crdt/reducer.rs b/src/crdt/reducer.rs index a9444b8..80dc6cd 100644 --- a/src/crdt/reducer.rs +++ b/src/crdt/reducer.rs @@ -53,6 +53,7 @@ mod tests { author: "test".into(), parents: Vec::new(), node_timestamp: None, + node_metadata: None, } } @@ -71,6 +72,7 @@ mod tests { author: "test".into(), parents: Vec::new(), node_timestamp: None, + node_metadata: None, } } diff --git a/src/crdt/storage.rs b/src/crdt/storage.rs index 9cb2480..62fbcbc 100644 --- a/src/crdt/storage.rs +++ b/src/crdt/storage.rs @@ -1,5 +1,5 @@ use crate::crdt::error::{CrdtError, Result}; -use crate::crdt::operation::Operation; +use crate::crdt::operation::{Operation, OperationType}; use crate::storage::{BatchError, LeveldbBatchGuard, SharedLeveldb, SharedLeveldbAccess}; use bincode; use rusty_leveldb::LdbIterator; @@ -8,6 +8,20 @@ use std::path::Path; use std::sync::Arc; use ulid::Ulid; +// Legacy records start with the bincode string length (26) of a ULID. +// A distinct versioned envelope prevents truncated new metadata from ever +// being accepted as a legacy operation with default policy metadata. +const OPERATION_V1: &[u8] = b"CRSLop\x01"; +type LegacyOperation = ( + Ulid, + ContentId, + OperationType, + u64, + String, + Vec, + Option, +); + /// Abstraction over the persistent storage used by `CrdtState`. pub trait OperationStorage: Send + Sync { fn save_operation(&self, op: &Operation) -> Result<()>; @@ -20,6 +34,10 @@ pub trait OperationStorage: Send + Sync { } /// LevelDB-backed implementation of [`OperationStorage`]. +/// +/// Reads legacy unversioned bincode operations and versioned records carrying +/// node metadata. New writes use the versioned format; older library versions +/// cannot read them. Invalid records are errors, never silently skipped. #[derive(Clone)] pub struct LeveldbStorage { shared: Arc, @@ -53,10 +71,56 @@ impl LeveldbStorage { ContentId: serde::Serialize, T: serde::Serialize, { - let value = bincode::serde::encode_to_vec(op, bincode::config::standard())?; + let mut value = OPERATION_V1.to_vec(); + value.extend(bincode::serde::encode_to_vec( + op, + bincode::config::standard(), + )?); Ok(value) } + fn decode_operation(raw: &[u8]) -> Result> + where + ContentId: for<'de> serde::Deserialize<'de>, + T: for<'de> serde::Deserialize<'de>, + { + let (op, consumed, expected) = if let Some(body) = raw.strip_prefix(OPERATION_V1) { + let (op, consumed) = + bincode::serde::decode_from_slice(body, bincode::config::standard())?; + (op, consumed, body.len()) + } else if raw.first() == Some(&26) { + let ((id, genesis, kind, timestamp, author, parents, node_timestamp), consumed) = + bincode::serde::decode_from_slice::, _>( + raw, + bincode::config::standard(), + )?; + ( + Operation { + id, + genesis, + kind, + timestamp, + author, + parents, + node_timestamp, + node_metadata: None, + }, + consumed, + raw.len(), + ) + } else { + return Err(CrdtError::Internal( + "Unknown operation storage format".into(), + )); + }; + if consumed != expected { + return Err(CrdtError::Internal( + "Trailing bytes in stored operation".into(), + )); + } + Ok(op) + } + /// Writes value bytes either to the active batch or directly to the DB. fn put_bytes(&self, key: &[u8], value: &[u8]) -> Result<()> { if self @@ -117,10 +181,8 @@ where let mut value = Vec::new(); while iter.valid() { iter.current(&mut key, &mut value); - if let Ok((op, _)) = bincode::serde::decode_from_slice::, _>( - &value, - bincode::config::standard(), - ) { + if key.first() == Some(&0x01) { + let op = Self::decode_operation(&value)?; if op.genesis == *genesis { result.push(op); } @@ -134,13 +196,7 @@ where fn get_operation(&self, op_id: &Ulid) -> Result>> { let key = Self::make_key(op_id); match self.shared.db().get(&key) { - Some(raw) => { - let (op, _) = bincode::serde::decode_from_slice::, _>( - &raw, - bincode::config::standard(), - )?; - Ok(Some(op)) - } + Some(raw) => Self::decode_operation(&raw).map(Some), None => Ok(None), } } diff --git a/src/repo.rs b/src/repo.rs index 57042a2..50010a7 100644 --- a/src/repo.rs +++ b/src/repo.rs @@ -36,6 +36,8 @@ where pub state: CrdtState, pub dag: DagGraph, resolver: ConflictResolver, + /// Application-supplied implementation selected by genesis policy name. + merge_policy: Option>>, } impl Repo @@ -52,14 +54,54 @@ where state, dag, resolver: ConflictResolver::new(), + merge_policy: None, } } + /// Install one named application-supplied [`MergePolicy`] for auto-merges. + /// + /// The library's only built-in policy is last-writer-wins over the whole + /// payload: whichever head has the newest timestamp is copied into the + /// merge node. That is right for a payload that is one value. It is wrong + /// for a payload that is several — a content body next to an access + /// policy, say — where a head that only advanced one field must not carry + /// the *other* field's stale value over a concurrent head that changed + /// it. The library cannot merge such a payload field by field, because + /// it does not know the fields; the application does, so it passes the + /// rule in here. + /// + /// Genesis metadata is authoritative: this implementation is used only + /// for its matching name. `lww` always selects the built-in rule; unknown + /// custom names error rather than fall back. A local Create records this + /// name unless `Operation::node_metadata` is explicitly supplied. Imports + /// use their own metadata, never the receiver's installed policy. + /// + /// Every replica merging a custom-policy content must install an + /// implementation with the same name and deterministic semantics. The + /// name travels with the data; executable code is process-local. Installing + /// another implementation does not migrate existing content. + /// + /// Policies require complete immediate-parent payloads by default. If a + /// parent has not synced, both lazy auto-merge and [`Self::merge_heads`] + /// return an error without committing a merge. Retry after importing the + /// missing parents. A policy independent of parent payloads can opt out + /// through [`MergePolicy::requires_parent_payloads`]. + pub fn with_merge_policy(mut self, policy: Box>) -> Self { + self.merge_policy = Some(policy); + self + } + /// Commits an operation to the repository. /// /// If `op.node_timestamp` is set, the operation is treated as an import from /// another replica, preserving the original timestamp for CID consistency. /// Otherwise, the current time is used for the DAG node timestamp. + /// Successful commits store the exact node timestamp and metadata on the + /// operation for subsequent export. Imported Create CIDs are checked; + /// imported Merge payloads are trusted, not recomputed by the local policy. + /// Child policy names are checked against genesis when it is available. + /// Out-of-order imports are checked again as heads before a merge; this is + /// not a validation of every ancestor or of the imported payload semantics. /// /// # Arguments /// @@ -90,6 +132,48 @@ where self.dag.calculate_latest(genesis_id).ok().flatten() } + /// The current heads of a content: every node no other node names as a + /// parent. One head means the history is linear at the tip; more means + /// concurrent versions are waiting to be merged. + pub fn heads(&self, genesis: &Cid) -> Result> { + self.find_heads(genesis) + } + + /// Merge concurrent heads now, and return the resulting single head. + /// + /// Auto-merge normally runs lazily, inside the next commit. That is too + /// late for a caller that wants to *read* the converged state first — + /// to copy the current payload into a new version, or to evaluate a + /// policy it carries — because [`latest`](Self::latest) alone picks one + /// of the concurrent heads by timestamp and says nothing about the + /// others. Calling this first makes the read see what the merge policy + /// decides, not what one branch happens to hold. + /// + /// Commits the merge node exactly as the lazy path would; returns + /// `Ok(None)` when the content has no versions, and the existing head + /// when there was nothing to merge. + pub fn merge_heads(&mut self, genesis: &Cid) -> Result> { + let shared = self.shared_leveldb()?; + let batch_guard = Self::begin_shared_batch(&shared)?; + let mut pending_nodes: Vec = Vec::new(); + + let merged = match self.check_and_merge(genesis, &mut pending_nodes) { + Ok(merged) => merged, + Err(err) => { + self.rollback_pending_nodes(&pending_nodes); + return Err(err); + } + }; + if let Err(status) = batch_guard.commit() { + self.rollback_pending_nodes(&pending_nodes); + return Err(CrdtError::Storage(status)); + } + match merged { + Some(cid) => Ok(Some(cid)), + None => self.dag.calculate_latest(genesis).map_err(CrdtError::Graph), + } + } + /// Convenience wrapper around `DagGraph::get_genesis` pub fn get_genesis(&self, cid: &Cid) -> Result { self.dag.get_genesis(cid).map_err(CrdtError::Graph) @@ -227,6 +311,9 @@ where } }; + // Persist the exact reconstruction inputs, not the receiver's defaults. + op.node_timestamp = Some(timestamp); + op.node_metadata = pending_nodes.last().map(|node| node.metadata.clone()); if let Err(err) = self.state.apply(op) { self.rollback_pending_nodes(&pending_nodes); return Err(err); @@ -306,9 +393,20 @@ where timestamp: u64, pending_nodes: &mut Vec, ) -> Result { - let (genesis_cid, node) = - self.dag - .prepare_genesis_node(payload, timestamp, ContentMetadata::default())?; + let metadata = if let Some(metadata) = &op.node_metadata { + metadata.clone() + } else if op.node_timestamp.is_some() { + ContentMetadata::default() + } else { + self.merge_policy + .as_ref() + .map_or_else(ContentMetadata::default, |policy| { + ContentMetadata::with_policy(policy.name()) + }) + }; + let (genesis_cid, node) = self + .dag + .prepare_genesis_node(payload, timestamp, metadata)?; if op.node_timestamp.is_some() { // Import: verify that the computed CID matches the expected genesis @@ -335,8 +433,12 @@ where pending_nodes: &mut Vec, ) -> Result { let lenient = op.node_timestamp.is_some(); - let metadata = - self.resolve_metadata(&op.genesis, &op.parents, pending_nodes.as_slice(), lenient)?; + let metadata = match &op.node_metadata { + Some(metadata) if lenient => metadata.clone(), + _ => { + self.resolve_metadata(&op.genesis, &op.parents, pending_nodes.as_slice(), lenient)? + } + }; let (cid, node) = self.dag.prepare_child_node( payload, op.parents.clone(), @@ -373,8 +475,12 @@ where })?; let lenient = op.node_timestamp.is_some(); - let metadata = - self.resolve_metadata(&op.genesis, &op.parents, pending_nodes.as_slice(), lenient)?; + let metadata = match &op.node_metadata { + Some(metadata) if lenient => metadata.clone(), + _ => { + self.resolve_metadata(&op.genesis, &op.parents, pending_nodes.as_slice(), lenient)? + } + }; let (cid, node) = self.dag.prepare_child_node( last_payload, op.parents.clone(), @@ -394,8 +500,12 @@ where pending_nodes: &mut Vec, ) -> Result { // Merge operations are always imports, so use lenient metadata resolution - let metadata = - self.resolve_metadata(&op.genesis, &op.parents, pending_nodes.as_slice(), true)?; + let metadata = match &op.node_metadata { + Some(metadata) => metadata.clone(), + None => { + self.resolve_metadata(&op.genesis, &op.parents, pending_nodes.as_slice(), true)? + } + }; let (cid, node) = self.dag.prepare_child_node( payload, op.parents.clone(), @@ -422,6 +532,11 @@ where cid: Cid, node: &Node, ) -> Result { + if let Some(genesis) = node.genesis { + if let Some(root) = self.dag.get_node(&genesis).map_err(CrdtError::Graph)? { + Self::validate_policy_binding(node.metadata(), root.metadata())?; + } + } self.dag.storage.put(node).map_err(CrdtError::Graph)?; self.dag .register_prepared_node(cid, node) @@ -462,8 +577,17 @@ where .get_node(genesis) .map_err(CrdtError::Graph)? .ok_or_else(|| CrdtError::Internal(format!("Genesis not found: {genesis}")))?; - let policy_type = genesis_node.metadata().policy_type(); - let policy = self.create_policy(policy_type)?; + // Out-of-order imports may have arrived before the genesis. Check + // every head now, before invoking a policy or persisting any merge. + for head in &heads { + let node = self + .dag + .get_node(head) + .map_err(CrdtError::Graph)? + .ok_or_else(|| CrdtError::Internal(format!("Head node not found: {head}")))?; + Self::validate_policy_binding(node.metadata(), genesis_node.metadata())?; + } + let policy = self.create_policy(genesis_node.metadata().policy_type())?; self.validate_parent_genesis(genesis, &heads)?; @@ -473,7 +597,7 @@ where &self.dag, *genesis, merge_timestamp, - policy.as_ref(), + policy, )?; let (merge_cid, node) = self @@ -494,6 +618,8 @@ where "auto-merge".to_string(), ); merge_op.parents = heads; + merge_op.node_timestamp = Some(merge_timestamp); + merge_op.node_metadata = Some(merge_node.metadata().clone()); if let Err(err) = self.state.apply(merge_op) { self.dag .rollback_pending_node(&pending.cid, &pending.parents); @@ -533,10 +659,25 @@ where .collect()) } - fn create_policy(&self, policy_type: &str) -> Result>> { - match policy_type { - "lww" => Ok(Box::new(LwwMergePolicy)), - other => Err(CrdtError::Internal(format!("Unknown policy type: {other}"))), + fn validate_policy_binding( + metadata: &ContentMetadata, + genesis: &ContentMetadata, + ) -> Result<()> { + if metadata.policy_type() != genesis.policy_type() { + return Err(CrdtError::Internal(format!( + "Policy mismatch: genesis requires {}, node declares {}", + genesis.policy_type(), + metadata.policy_type() + ))); + } + Ok(()) + } + + fn create_policy(&self, policy_type: &str) -> Result<&dyn MergePolicy> { + match (policy_type, self.merge_policy.as_deref()) { + ("lww", _) => Ok(&LwwMergePolicy), + (name, Some(policy)) if policy.name() == name => Ok(policy), + (other, _) => Err(CrdtError::Internal(format!("Unknown policy type: {other}"))), } } @@ -1382,6 +1523,177 @@ mod tests { assert!(!heads_after_merge.contains(&branch2_cid)); } + /// A policy installed with `with_merge_policy` is recorded in a local + /// genesis and decides that content's merge payload. It sees each + /// head's parent payloads, so it can tell what a head changed. + #[test] + fn installed_merge_policy_is_used_and_sees_parents() { + use crate::convergence::policy::{MergePolicy, ResolveInput}; + use std::sync::{Arc, Mutex}; + + /// Concatenates every head that differs from its parent, in CID order, + /// and records what it was given. Deliberately not LWW so the test can + /// tell the two apart. + type Seen = Arc)>>>; + struct Concat { + seen: Seen, + } + impl MergePolicy for Concat { + fn resolve(&self, nodes: &[ResolveInput]) -> TestPayload { + let mut seen = self.seen.lock().unwrap(); + for n in nodes { + seen.push(( + n.payload.0.clone(), + n.parent_payloads.iter().map(|p| p.0.clone()).collect(), + )); + } + let mut changed: Vec<&str> = nodes + .iter() + .filter(|n| n.parent_payloads.iter().all(|p| p != &n.payload)) + .map(|n| n.payload.0.as_str()) + .collect(); + changed.sort(); + TestPayload(changed.join("+")) + } + fn name(&self) -> &str { + "concat" + } + } + + let seen = Arc::new(Mutex::new(Vec::new())); + let (repo, _dir) = setup_test_repo(); + let mut repo = repo.with_merge_policy(Box::new(Concat { seen: seen.clone() })); + + let initial_genesis = Cid::new_v1( + 0x55, + multihash::Multihash::<64>::wrap(0x12, b"installedPolicy").unwrap(), + ); + let create = make_test_operation( + initial_genesis, + OperationType::Create(TestPayload("root".into())), + ); + let genesis = repo.commit_operation(create).unwrap(); + + // Two concurrent heads off the genesis. One repeats its parent's value + // (a "policy-only" head, as far as this payload can express it). + let mut a = make_test_operation(genesis, OperationType::Update(TestPayload("a".into()))); + a.parents.push(genesis); + repo.commit_operation(a).unwrap(); + sleep_for_ordering(); + let mut same = + make_test_operation(genesis, OperationType::Update(TestPayload("root".into()))); + same.parents.push(genesis); + repo.commit_operation(same).unwrap(); + sleep_for_ordering(); + + // Triggers the auto-merge. + let update = make_test_operation(genesis, OperationType::Update(TestPayload("z".into()))); + repo.commit_operation(update).unwrap(); + + let ops = repo.state.get_operations_by_genesis(&genesis).unwrap(); + let merge = ops + .iter() + .find_map(|op| match &op.kind { + OperationType::Merge(p) => Some(p.clone()), + _ => None, + }) + .expect("auto-merge must have run"); + // LWW would have picked "root" (the newer head); the installed policy + // kept only the head that changed something. + assert_eq!(merge, TestPayload("a".into())); + + // Both heads were handed over with their parent ("root") attached. + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 2); + for (_, parents) in seen.iter() { + assert_eq!(parents, &vec!["root".to_string()]); + } + } + + /// `merge_heads` converges concurrent heads on demand, so a caller can + /// read the merged payload before committing on top of it. Without it, + /// `latest` returns one branch's tip and the merge only happens inside + /// the next commit — after the caller has already copied the wrong value. + #[test] + fn merge_heads_converges_before_the_next_commit() { + let (mut repo, _dir) = setup_test_repo(); + let initial_genesis = Cid::new_v1( + 0x55, + multihash::Multihash::<64>::wrap(0x12, b"mergeHeads").unwrap(), + ); + let create = make_test_operation( + initial_genesis, + OperationType::Create(TestPayload("root".into())), + ); + let genesis = repo.commit_operation(create).unwrap(); + let mut a = make_test_operation(genesis, OperationType::Update(TestPayload("a".into()))); + a.parents.push(genesis); + let a_cid = repo.commit_operation(a).unwrap(); + sleep_for_ordering(); + let mut b = make_test_operation(genesis, OperationType::Update(TestPayload("b".into()))); + b.parents.push(genesis); + let b_cid = repo.commit_operation(b).unwrap(); + + assert_eq!(repo.heads(&genesis).unwrap().len(), 2); + // `latest` alone just picks the newer branch. + assert_eq!(repo.latest(&genesis), Some(b_cid)); + + let merged = repo.merge_heads(&genesis).unwrap().unwrap(); + assert_ne!(merged, a_cid); + assert_ne!(merged, b_cid); + assert_eq!(repo.heads(&genesis).unwrap(), vec![merged]); + assert_eq!(repo.latest(&genesis), Some(merged)); + let node = repo.dag.get_node(&merged).unwrap().unwrap(); + assert_eq!(node.parents().len(), 2); + + // Idempotent: nothing left to merge, same head comes back. + assert_eq!(repo.merge_heads(&genesis).unwrap(), Some(merged)); + // Unknown content: nothing to merge, no head. + let unknown = Cid::new_v1( + 0x55, + multihash::Multihash::<64>::wrap(0x12, b"nothing").unwrap(), + ); + assert_eq!(repo.merge_heads(&unknown).unwrap(), None); + } + + /// Without an installed policy, the genesis metadata's policy_type still + /// selects the built-in LWW — the extension point changes nothing for + /// callers that do not use it. + #[test] + fn named_policy_still_applies_when_none_installed() { + let (mut repo, _dir) = setup_test_repo(); + let initial_genesis = Cid::new_v1( + 0x55, + multihash::Multihash::<64>::wrap(0x12, b"namedPolicy").unwrap(), + ); + let create = make_test_operation( + initial_genesis, + OperationType::Create(TestPayload("root".into())), + ); + let genesis = repo.commit_operation(create).unwrap(); + let mut a = make_test_operation(genesis, OperationType::Update(TestPayload("a".into()))); + a.parents.push(genesis); + repo.commit_operation(a).unwrap(); + sleep_for_ordering(); + let mut b = make_test_operation(genesis, OperationType::Update(TestPayload("b".into()))); + b.parents.push(genesis); + repo.commit_operation(b).unwrap(); + sleep_for_ordering(); + let update = make_test_operation(genesis, OperationType::Update(TestPayload("z".into()))); + repo.commit_operation(update).unwrap(); + + let ops = repo.state.get_operations_by_genesis(&genesis).unwrap(); + let merge = ops + .iter() + .find_map(|op| match &op.kind { + OperationType::Merge(p) => Some(p.clone()), + _ => None, + }) + .expect("auto-merge must have run"); + // LWW: the later head wins outright. + assert_eq!(merge, TestPayload("b".into())); + } + #[test] fn test_auto_merge_from_intermediate_branch() { let (mut repo, _) = setup_test_repo(); diff --git a/tests/merge_heads_errors.rs b/tests/merge_heads_errors.rs new file mode 100644 index 0000000..0636379 --- /dev/null +++ b/tests/merge_heads_errors.rs @@ -0,0 +1,116 @@ +use cid::Cid; +use crsl_lib::{ + convergence::metadata::ContentMetadata, + crdt::{ + crdt_state::CrdtState, + operation::{Operation, OperationType}, + storage::LeveldbStorage, + }, + dasl::node::Node, + graph::{ + dag::DagGraph, + error::{GraphError, Result}, + storage::{LeveldbNodeStorage, NodeStorage}, + }, + repo::Repo, + storage::{SharedLeveldb, SharedLeveldbAccess}, +}; +use rusty_leveldb::{Status, StatusCode}; +use std::{ + collections::HashMap, + sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }, +}; + +struct ReadFault { + inner: LeveldbNodeStorage, + calls: AtomicUsize, + fail_at: AtomicUsize, +} +impl ReadFault { + fn arm(&self, fail_at: usize) { + self.calls.store(0, Ordering::SeqCst); + self.fail_at.store(fail_at, Ordering::SeqCst); + } +} +impl SharedLeveldbAccess for ReadFault { + fn shared_leveldb(&self) -> Option> { + self.inner.shared_leveldb() + } +} +impl NodeStorage for ReadFault { + fn get(&self, cid: &Cid) -> Result>> { + self.inner.get(cid) + } + fn put(&self, node: &Node) -> Result<()> { + self.inner.put(node) + } + fn delete(&self, cid: &Cid) -> Result<()> { + self.inner.delete(cid) + } + fn get_node_map(&self) -> Result>> { + let call = self.calls.fetch_add(1, Ordering::SeqCst) + 1; + if call == self.fail_at.load(Ordering::SeqCst) { + return Err(GraphError::Storage(Status::new( + StatusCode::IOError, + "review read failure", + ))); + } + let map = self.inner.get_node_map()?; + Ok(map) + } +} + +#[test] +fn merge_heads_propagates_read_errors() { + let dir = tempfile::tempdir().unwrap(); + let shared = SharedLeveldb::open(dir.path().join("store")).unwrap(); + let state = CrdtState::new(LeveldbStorage::new(shared.clone())); + let storage = ReadFault { + inner: LeveldbNodeStorage::new(shared), + calls: AtomicUsize::new(0), + fail_at: AtomicUsize::new(0), + }; + let mut repo = Repo::new(state, DagGraph::new(storage)); + let seed = Cid::new_v1( + 0x55, + multihash::Multihash::<64>::wrap(0x12, b"review").unwrap(), + ); + let genesis = repo + .commit_operation(Operation::new( + seed, + OperationType::Create("persisted".to_owned()), + "review".into(), + )) + .unwrap(); + assert_eq!(repo.heads(&genesis).unwrap(), vec![genesis]); + assert_eq!(repo.merge_heads(&genesis).unwrap(), Some(genesis)); + + repo.dag.storage.arm(1); + let first = repo.merge_heads(&genesis); + assert!(matches!( + first, + Err(crsl_lib::crdt::error::CrdtError::Graph( + GraphError::Storage(_) + )) + )); + + repo.dag.storage.arm(2); + let result = repo.merge_heads(&genesis); + assert_eq!(repo.dag.storage.calls.load(Ordering::SeqCst), 2); + assert!( + matches!( + result, + Err(crsl_lib::crdt::error::CrdtError::Graph( + GraphError::Storage(_) + )) + ), + "read failure must propagate: {result:?}" + ); + + repo.dag.storage.arm(0); + assert_eq!(repo.heads(&genesis).unwrap(), vec![genesis]); + assert_eq!(repo.merge_heads(&genesis).unwrap(), Some(genesis)); +} diff --git a/tests/merge_parent_completeness.rs b/tests/merge_parent_completeness.rs new file mode 100644 index 0000000..3d54883 --- /dev/null +++ b/tests/merge_parent_completeness.rs @@ -0,0 +1,211 @@ +use cid::Cid; +use crsl_lib::{ + convergence::{ + metadata::ContentMetadata, + policy::{MergePolicy, ResolveInput}, + }, + crdt::{ + crdt_state::CrdtState, + operation::{Operation, OperationType}, + storage::LeveldbStorage, + }, + graph::{dag::DagGraph, storage::LeveldbNodeStorage}, + repo::Repo, + storage::SharedLeveldb, +}; +use std::sync::{Arc, Mutex}; + +type TestRepo = + Repo, LeveldbNodeStorage, String>; +type Seen = Arc)>>>; +// Same parent-aware rule as the PR's installed_merge_policy_is_used_and_sees_parents. +struct Changed(Seen); +impl MergePolicy for Changed { + fn resolve(&self, nodes: &[ResolveInput]) -> String { + *self.0.lock().unwrap() = nodes + .iter() + .map(|n| (n.payload.clone(), n.parent_payloads.clone())) + .collect(); + let mut changed: Vec<_> = nodes + .iter() + .filter(|n| n.parent_payloads.iter().all(|p| p != &n.payload)) + .map(|n| n.payload.clone()) + .collect(); + changed.sort(); + changed.join("+") + } + fn name(&self) -> &str { + "changed" + } +} +fn repo() -> (TestRepo, tempfile::TempDir, Seen) { + let dir = tempfile::tempdir().unwrap(); + let db = SharedLeveldb::open(dir.path().join("db")).unwrap(); + let seen = Arc::new(Mutex::new(Vec::new())); + let repo = Repo::new( + CrdtState::new(LeveldbStorage::new(db.clone())), + DagGraph::new(LeveldbNodeStorage::new(db)), + ) + .with_merge_policy(Box::new(Changed(seen.clone()))); + (repo, dir, seen) +} +fn op(g: Cid, kind: OperationType, parents: Vec, t: u64) -> Operation { + let mut op = Operation::new(g, kind, "replica".into()); + op.parents = parents; + op.timestamp = t; + op.node_timestamp = Some(t); + op +} +fn assert_merge_waits_for_parents(partial_multi_parent: bool) { + let (mut complete, _d1, seen_complete) = repo(); + let (mut incomplete, _d2, seen_incomplete) = repo(); + let placeholder = crsl_lib::dasl::node::Node::new_genesis( + "root".to_string(), + 1, + ContentMetadata::with_policy("changed"), + ) + .content_id() + .unwrap(); + let mut create = op(placeholder, OperationType::Create("root".into()), vec![], 1); + create.node_metadata = Some(ContentMetadata::with_policy("changed")); + let g = complete.commit_operation(create.clone()).unwrap(); + assert_eq!(g, incomplete.commit_operation(create).unwrap()); + let parent = op(g, OperationType::Update("old".into()), vec![g], 2); + let p = complete.commit_operation(parent.clone()).unwrap(); + let parents = if partial_multi_parent { + vec![g, p] + } else { + vec![p] + }; + let copy = op(g, OperationType::Update("old".into()), parents, 4); + let edit = op(g, OperationType::Update("new".into()), vec![g], 3); + for operation in [copy, edit] { + assert_eq!( + complete.commit_operation(operation.clone()).unwrap(), + incomplete.commit_operation(operation).unwrap() + ); + } + let mut h1 = complete.heads(&g).unwrap(); + let mut h2 = incomplete.heads(&g).unwrap(); + h1.sort(); + h2.sort(); + assert_eq!(h1, h2); + assert_eq!(h1.len(), 2); + assert!(incomplete.dag.get_node(&p).unwrap().is_none()); + let good = complete.merge_heads(&g).unwrap().unwrap(); + let original_ops = incomplete.state.get_operations_by_genesis(&g).unwrap(); + let error = incomplete + .merge_heads(&g) + .expect_err("missing parent must defer merge"); + assert!(error.to_string().contains(&p.to_string())); + assert!( + seen_incomplete.lock().unwrap().is_empty(), + "policy must not see partial inputs" + ); + let mut remaining = incomplete.heads(&g).unwrap(); + remaining.sort(); + assert_eq!(remaining, h2); + assert_eq!( + incomplete.state.get_operations_by_genesis(&g).unwrap(), + original_ops + ); + + // A failed explicit merge must release its batch, and the lazy path must + // enforce the same rule without persisting its caller's update either. + let update = Operation::new( + g, + OperationType::Update("after".to_string()), + "local".into(), + ); + assert!(incomplete.commit_operation(update).is_err()); + assert_eq!( + incomplete.state.get_operations_by_genesis(&g).unwrap(), + original_ops + ); + assert!(seen_incomplete.lock().unwrap().is_empty()); + + assert_eq!(incomplete.commit_operation(parent).unwrap(), p); + let merged = incomplete.merge_heads(&g).unwrap().unwrap(); + let expected = complete + .dag + .get_node(&good) + .unwrap() + .unwrap() + .payload() + .clone(); + assert_eq!(expected, "new"); + assert_eq!( + incomplete.dag.get_node(&merged).unwrap().unwrap().payload(), + &expected + ); + assert_eq!(incomplete.heads(&g).unwrap(), vec![merged]); + assert_eq!(incomplete.merge_heads(&g).unwrap(), Some(merged)); + let mut complete_inputs = seen_complete.lock().unwrap().clone(); + let mut retried_inputs = seen_incomplete.lock().unwrap().clone(); + complete_inputs.sort(); + retried_inputs.sort(); + assert_eq!(complete_inputs, retried_inputs); +} +#[test] +fn merge_waits_for_missing_parent() { + assert_merge_waits_for_parents(false); +} +#[test] +fn merge_waits_for_one_of_two_parents() { + assert_merge_waits_for_parents(true); +} + +fn assert_lww_allows_missing_parents(installed: bool) { + use crsl_lib::convergence::policies::lww::LwwMergePolicy; + use crsl_lib::dasl::node::Node; + + let (original, _dir, _) = repo(); + let mut repo = Repo::new(original.state, original.dag); + if installed { + repo = repo.with_merge_policy(Box::new(LwwMergePolicy)); + } + let root = Node::new_genesis("root".to_string(), 1, ContentMetadata::default()); + let genesis = root.content_id().unwrap(); + repo.commit_operation(op(genesis, OperationType::Create("root".into()), vec![], 1)) + .unwrap(); + let missing = Node::new_child( + "missing".to_string(), + vec![genesis], + genesis, + 2, + ContentMetadata::default(), + ) + .content_id() + .unwrap(); + repo.commit_operation(op( + genesis, + OperationType::Update("older".into()), + vec![genesis], + 3, + )) + .unwrap(); + repo.commit_operation(op( + genesis, + OperationType::Update("newer".into()), + vec![missing], + 4, + )) + .unwrap(); + assert!(repo.dag.get_node(&missing).unwrap().is_none()); + let merged = repo.merge_heads(&genesis).unwrap().unwrap(); + assert_eq!( + repo.dag.get_node(&merged).unwrap().unwrap().payload(), + "newer" + ); + assert_eq!(repo.heads(&genesis).unwrap(), vec![merged]); +} + +#[test] +fn named_lww_allows_missing_parents() { + assert_lww_allows_missing_parents(false); +} + +#[test] +fn installed_lww_allows_missing_parents() { + assert_lww_allows_missing_parents(true); +} diff --git a/tests/operation_storage_compatibility.rs b/tests/operation_storage_compatibility.rs new file mode 100644 index 0000000..2207d9d --- /dev/null +++ b/tests/operation_storage_compatibility.rs @@ -0,0 +1,74 @@ +use crsl_lib::{ + convergence::metadata::ContentMetadata, + crdt::{ + operation::{Operation, OperationType}, + storage::{LeveldbStorage, OperationStorage}, + }, + storage::SharedLeveldb, +}; +use serde::Serialize; +use ulid::Ulid; + +// Exact pre-metadata persisted schema. Do not derive this from today's Operation. +#[derive(Serialize)] +struct LegacyOperation { + id: Ulid, + genesis: u64, + kind: OperationType, + timestamp: u64, + author: String, + parents: Vec, + node_timestamp: Option, +} +fn key(id: Ulid) -> Vec { + let mut key = vec![1]; + key.extend_from_slice(&id.to_bytes()); + key +} +#[test] +fn reads_real_legacy_bincode_records_alongside_new_records() { + let dir = tempfile::tempdir().unwrap(); + let db = SharedLeveldb::open(dir.path()).unwrap(); + let storage = LeveldbStorage::::new(db.clone()); + let id = Ulid::new(); + let old = LegacyOperation { + id, + genesis: 42, + kind: OperationType::Create("old".into()), + timestamp: 9, + author: "old-client".into(), + parents: vec![], + node_timestamp: Some(9), + }; + let bytes = bincode::serde::encode_to_vec(&old, bincode::config::standard()).unwrap(); + db.db().put(&key(id), &bytes).unwrap(); + let old_read = storage.get_operation(&id).unwrap().unwrap(); + assert_eq!(old_read.node_metadata, None); + assert_eq!(old_read.node_timestamp, Some(9)); + assert_eq!(old_read.kind, old.kind); + let mut new = Operation::new(42, OperationType::Update("new".into()), "new-client".into()); + new.node_metadata = Some(ContentMetadata::with_policy("custom-v1")); + storage.save_operation(&new).unwrap(); + assert_eq!(storage.get_operation(&new.id).unwrap(), Some(new.clone())); + let all = storage.load_operations(&42).unwrap(); + assert_eq!(all.len(), 2); + assert!(all.contains(&old_read)); + assert!(all.contains(&new)); +} +#[test] +fn corrupt_or_trailing_operation_bytes_are_errors_not_dropped_or_downgraded() { + let dir = tempfile::tempdir().unwrap(); + let db = SharedLeveldb::open(dir.path()).unwrap(); + let storage = LeveldbStorage::::new(db.clone()); + let mut op = Operation::new(42, OperationType::Create("new".into()), "local".into()); + op.node_metadata = Some(ContentMetadata::with_policy("custom-v1")); + storage.save_operation(&op).unwrap(); + let bytes = db.db().get(&key(op.id)).unwrap(); + let mut trailing = bytes.to_vec(); + trailing.push(0); + for corrupted in [bytes[..bytes.len() - 1].to_vec(), trailing, vec![255]] { + db.db().put(&key(op.id), &corrupted).unwrap(); + assert!(storage.get_operation(&op.id).is_err()); + assert!(storage.load_operations(&42).is_err()); + } +} diff --git a/tests/policy_binding.rs b/tests/policy_binding.rs new file mode 100644 index 0000000..c2ccce5 --- /dev/null +++ b/tests/policy_binding.rs @@ -0,0 +1,390 @@ +use cid::Cid; +use crsl_lib::{ + convergence::{ + metadata::ContentMetadata, + policy::{MergePolicy, ResolveInput}, + }, + crdt::{ + crdt_state::CrdtState, + operation::{Operation, OperationType}, + storage::LeveldbStorage, + }, + dasl::node::Node, + graph::{dag::DagGraph, storage::LeveldbNodeStorage}, + repo::Repo, + storage::SharedLeveldb, +}; +type TestRepo = + Repo, LeveldbNodeStorage, String>; +struct Custom(&'static str); +impl MergePolicy for Custom { + fn name(&self) -> &str { + self.0 + } + fn resolve(&self, _: &[ResolveInput]) -> String { + "custom-result".into() + } +} +fn repo(policy: Option<&'static str>) -> (TestRepo, tempfile::TempDir) { + let dir = tempfile::tempdir().unwrap(); + let db = SharedLeveldb::open(dir.path()).unwrap(); + let repo = Repo::new( + CrdtState::new(LeveldbStorage::new(db.clone())), + DagGraph::new(LeveldbNodeStorage::new(db)), + ); + ( + match policy { + Some(name) => repo.with_merge_policy(Box::new(Custom(name))), + None => repo, + }, + dir, + ) +} +fn seed() -> Cid { + Node::new_genesis("seed".to_string(), 1, ContentMetadata::default()) + .content_id() + .unwrap() +} +#[test] +fn local_create_records_selected_policy() { + let (mut repo, _dir) = repo(Some("custom-v1")); + let g = repo + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + assert_eq!( + repo.dag + .get_node(&g) + .unwrap() + .unwrap() + .metadata() + .policy_type(), + "custom-v1" + ); +} + +#[test] +fn exported_operations_reconstruct_exact_custom_nodes_without_receiver_policy() { + let (mut source, _dir) = repo(Some("custom-v1")); + let g = source + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + branches(&mut source, g); + source.merge_heads(&g).unwrap(); + let ops = source.get_operations_with_index(&g).unwrap(); + for installed in [None, Some("wrong-v1"), Some("custom-v1")] { + let (mut receiver, _dir) = repo(installed); + // Import children before genesis: metadata cannot depend on local ancestry. + for (_, op) in ops.iter().rev() { + let json = serde_json::to_vec(op).unwrap(); + let imported = serde_json::from_slice(&json).unwrap(); + let cid = receiver.commit_operation(imported).unwrap(); + assert_eq!( + receiver.dag.get_node(&cid).unwrap(), + source.dag.get_node(&cid).unwrap() + ); + assert!(source.dag.get_node(&cid).unwrap().is_some()); + } + assert_eq!(receiver.latest(&g), source.latest(&g)); + assert_eq!( + receiver.dag.get_node(&g).unwrap(), + source.dag.get_node(&g).unwrap() + ); + } +} + +#[test] +fn legacy_json_import_uses_historical_default_not_installed_policy() { + let expected = Node::new_genesis("legacy".to_string(), 7, ContentMetadata::default()); + let g = expected.content_id().unwrap(); + let mut op = Operation::new(g, OperationType::Create("legacy".to_string()), "old".into()); + op.node_timestamp = Some(7); + let mut json = serde_json::to_value(op).unwrap(); + json.as_object_mut().unwrap().remove("node_metadata"); + for installed in [None, Some("custom-v1")] { + let (mut receiver, _dir) = repo(installed); + let imported: Operation = serde_json::from_value(json.clone()).unwrap(); + assert_eq!(imported.node_metadata, None); + assert_eq!(receiver.commit_operation(imported).unwrap(), g); + assert_eq!(receiver.dag.get_node(&g).unwrap(), Some(expected.clone())); + } +} + +#[test] +fn explicit_metadata_preserves_representation_and_local_create_choice() { + for metadata in [ + ContentMetadata::default(), + ContentMetadata::with_policy("lww"), + ContentMetadata::with_policy("custom-v1"), + ] { + let expected = Node::new_genesis("prepared".to_string(), 7, metadata.clone()); + let g = expected.content_id().unwrap(); + let (mut receiver, _dir) = repo(Some("different")); + let mut imported = Operation::new( + g, + OperationType::Create("prepared".into()), + "prepared".into(), + ); + imported.node_metadata = Some(metadata.clone()); + imported.node_timestamp = Some(7); + assert_eq!(receiver.commit_operation(imported).unwrap(), g); + assert_eq!(receiver.dag.get_node(&g).unwrap(), Some(expected)); + let mut local = Operation::new( + seed(), + OperationType::Create("local".into()), + "local".into(), + ); + local.node_metadata = Some(metadata.clone()); + let local_g = receiver.commit_operation(local).unwrap(); + assert_eq!( + receiver.dag.get_node(&local_g).unwrap().unwrap().metadata(), + &metadata + ); + } +} + +#[test] +fn builtin_lww_name_is_reserved_even_for_an_installed_policy() { + let (mut repo, _dir) = repo(Some("lww")); + let g = repo + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + branches(&mut repo, g); + let merged = repo.merge_heads(&g).unwrap().unwrap(); + assert_eq!( + repo.dag.get_node(&merged).unwrap().unwrap().payload(), + "newer" + ); +} + +#[test] +fn imported_custom_genesis_cannot_lose_metadata_silently() { + let (mut source, _dir) = repo(Some("custom-v1")); + let g = source + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + let mut op = source.get_operations_with_index(&g).unwrap().remove(0).1; + op.node_metadata = None; + let (mut receiver, _dir) = repo(Some("custom-v1")); + assert!(receiver + .commit_operation(op) + .unwrap_err() + .to_string() + .contains("CID mismatch")); + assert!(receiver.dag.get_nodes_by_genesis(&g).unwrap().is_empty()); + assert!(receiver + .state + .get_operations_by_genesis(&g) + .unwrap() + .is_empty()); +} + +#[test] +fn imported_child_policy_must_match_known_genesis_without_writes() { + for kind in [ + OperationType::Update("bad".into()), + OperationType::Merge("bad".into()), + OperationType::Delete, + ] { + let (mut repo, _dir) = repo(Some("custom-v1")); + let g = repo + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + let before = repo.state.get_operations_by_genesis(&g).unwrap(); + let mut imported = Operation::new(g, kind, "remote".into()); + imported.parents = vec![g]; + imported.node_timestamp = Some(10); + imported.node_metadata = Some(ContentMetadata::default()); + assert!(repo + .commit_operation(imported) + .unwrap_err() + .to_string() + .contains("Policy mismatch")); + assert_eq!(repo.state.get_operations_by_genesis(&g).unwrap(), before); + assert_eq!(repo.dag.get_nodes_by_genesis(&g).unwrap(), vec![g]); + } +} + +#[test] +fn heads_imported_before_genesis_are_checked_before_merging() { + let (mut source, _dir) = repo(Some("custom-v1")); + let g = source + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + let create = source.get_operations_with_index(&g).unwrap().remove(0).1; + let (mut receiver, _dir) = repo(Some("custom-v1")); + let mut bad_head = None; + for (t, metadata) in [ + (10, ContentMetadata::default()), + (11, ContentMetadata::with_policy("custom-v1")), + ] { + let mut imported = Operation::new( + g, + OperationType::Update(format!("head-{t}")), + "remote".into(), + ); + imported.parents = vec![g]; + imported.node_timestamp = Some(t); + imported.node_metadata = Some(metadata); + let cid = receiver.commit_operation(imported).unwrap(); + if t == 10 { + bad_head = Some(cid); + } + } + receiver.commit_operation(create).unwrap(); + let before = receiver.state.get_operations_by_genesis(&g).unwrap(); + assert!(receiver + .merge_heads(&g) + .unwrap_err() + .to_string() + .contains("Policy mismatch")); + let mut update = Operation::new(g, OperationType::Update("after".into()), "local".into()); + assert!(receiver + .commit_operation(update.clone()) + .unwrap_err() + .to_string() + .contains("Policy mismatch")); + update.parents = vec![bad_head.unwrap()]; + assert!(receiver + .commit_operation(update) + .unwrap_err() + .to_string() + .contains("Policy mismatch")); + assert_eq!( + receiver.state.get_operations_by_genesis(&g).unwrap(), + before + ); + assert_eq!(receiver.heads(&g).unwrap().len(), 2); +} + +fn branches(repo: &mut TestRepo, g: Cid) { + for value in ["older", "newer"] { + let mut op = Operation::new(g, OperationType::Update(value.into()), "local".into()); + op.parents = vec![g]; + repo.commit_operation(op).unwrap(); + } +} + +#[test] +fn existing_lww_is_not_overridden() { + for lazy in [false, true] { + let (mut original, _dir) = repo(None); + let g = original + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + branches(&mut original, g); + let mut installed = original.with_merge_policy(Box::new(Custom("different"))); + if lazy { + installed + .commit_operation(Operation::new( + g, + OperationType::Update("after".into()), + "local".into(), + )) + .unwrap(); + } else { + installed.merge_heads(&g).unwrap(); + } + let ops = installed.state.get_operations_by_genesis(&g).unwrap(); + let merged = ops + .iter() + .find_map(|op| match &op.kind { + OperationType::Merge(p) => Some(p.as_str()), + _ => None, + }) + .unwrap(); + assert_eq!(merged, "newer"); + } +} + +#[test] +fn unavailable_custom_policy_rejects_merges_without_writes() { + for installed in [None, Some("wrong-v1")] { + let (mut original, _dir) = repo(Some("custom-v1")); + let g = original + .commit_operation(Operation::new( + seed(), + OperationType::Create("root".into()), + "local".into(), + )) + .unwrap(); + branches(&mut original, g); + let mut receiver = Repo::new(original.state, original.dag); + if let Some(name) = installed { + receiver = receiver.with_merge_policy(Box::new(Custom(name))); + } + let ops = receiver.state.get_operations_by_genesis(&g).unwrap(); + let nodes: std::collections::HashSet<_> = receiver + .dag + .get_nodes_by_genesis(&g) + .unwrap() + .into_iter() + .collect(); + let heads: std::collections::HashSet<_> = receiver.heads(&g).unwrap().into_iter().collect(); + assert!(receiver + .merge_heads(&g) + .unwrap_err() + .to_string() + .contains("custom-v1")); + assert!(receiver + .commit_operation(Operation::new( + g, + OperationType::Update("after".into()), + "local".into() + )) + .unwrap_err() + .to_string() + .contains("custom-v1")); + assert_eq!(receiver.state.get_operations_by_genesis(&g).unwrap(), ops); + assert_eq!( + receiver + .dag + .get_nodes_by_genesis(&g) + .unwrap() + .into_iter() + .collect::>(), + nodes + ); + assert_eq!( + receiver + .heads(&g) + .unwrap() + .into_iter() + .collect::>(), + heads + ); + let mut receiver = receiver.with_merge_policy(Box::new(Custom("custom-v1"))); + let merged = receiver.merge_heads(&g).unwrap().unwrap(); + assert_eq!( + receiver.dag.get_node(&merged).unwrap().unwrap().payload(), + "custom-result" + ); + } +}