From 08cd6fb2a9356646238682f495f2a2b7d417affe Mon Sep 17 00:00:00 2001 From: somasekimoto <0421.soma@gmail.com> Date: Sun, 13 Sep 2026 22:26:02 +0900 Subject: [PATCH 1/3] Let the application supply the auto-merge policy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Auto-merge resolves concurrent heads with the policy named in the genesis metadata, and the only one that exists is last-writer-wins over the whole payload: the newest head's payload is copied into the merge node. That is right when the payload is one value. It is wrong when the payload is several values with different convergence rules — a content body next to an access policy, say. A head that only advanced the policy still carries the body it inherited; a concurrent head that only changed the body still carries the policy it inherited. Whichever is newer wins both fields, and the other head's change is lost even though nothing competed with it. The library cannot merge such a payload field by field: it does not know the fields. So it now lets the application say how. - `Repo::with_merge_policy(Box>)` installs a policy that is used for every auto-merge in place of the named one. Nothing changes for callers that do not install one. - `ResolveInput` gains `parent_payloads`: the payloads of the head's parents. A field-aware policy needs to know what a head changed, not just what it holds, and only the parent can tell it. `LwwMergePolicy` ignores the field; a parent missing from storage (ancestry not yet synced) is skipped rather than an error. The policy is process-local and never travels with the data, so every replica must install the same one — replicas merging the same heads under different rules produce different merge nodes and keep re-merging. That is documented on the method. Tests: the installed policy is used over the genesis's "lww" and sees each head's parents (repo); the resolver attaches parent payloads (resolver); the named policy still applies when nothing is installed (repo). Co-Authored-By: Claude Code --- src/convergence/policy.rs | 29 +++++++ src/convergence/resolver.rs | 69 ++++++++++++++- src/repo.rs | 163 +++++++++++++++++++++++++++++++++++- 3 files changed, 257 insertions(+), 4 deletions(-) diff --git a/src/convergence/policy.rs b/src/convergence/policy.rs index 827d105..8c3e9d7 100644 --- a/src/convergence/policy.rs +++ b/src/convergence/policy.rs @@ -6,6 +6,18 @@ 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. + /// + /// Empty for a genesis head (no parents) and for inputs constructed by + /// callers that have no DAG at hand; `LwwMergePolicy` ignores it. + pub parent_payloads: Vec

, } impl

ResolveInput

{ @@ -14,11 +26,28 @@ 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 then calls it for every auto-merge instead of the named policy. pub trait MergePolicy

: Send + Sync { /// Resolve competing nodes into a single payload. fn resolve(&self, nodes: &[ResolveInput

]) -> P; diff --git a/src/convergence/resolver.rs b/src/convergence/resolver.rs index 2fadac2..941424b 100644 --- a/src/convergence/resolver.rs +++ b/src/convergence/resolver.rs @@ -80,10 +80,20 @@ where .get_node(&cid) .map_err(CrdtError::Graph)? .ok_or_else(|| CrdtError::Internal(format!("Head node not found: {cid}")))?; - inputs.push(ResolveInput::new( + // A parent that is not in storage is not an error here: a replica + // may hold a head whose ancestry has not fully synced yet. The + // policy sees fewer parents and must cope (LWW never looks). + let mut parent_payloads = Vec::with_capacity(node.parents().len()); + for parent in node.parents() { + if let Some(parent_node) = dag.get_node(parent).map_err(CrdtError::Graph)? { + parent_payloads.push(parent_node.payload().clone()); + } + } + inputs.push(ResolveInput::with_parents( cid, node.payload().clone(), node.timestamp(), + parent_payloads, )); } Ok(inputs) @@ -193,6 +203,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/repo.rs b/src/repo.rs index 57042a2..7c060b6 100644 --- a/src/repo.rs +++ b/src/repo.rs @@ -36,6 +36,9 @@ where pub state: CrdtState, pub dag: DagGraph, resolver: ConflictResolver, + /// Application-supplied merge policy. When set it is used for every + /// auto-merge, regardless of the `policy_type` in the genesis metadata. + merge_policy: Option>>, } impl Repo @@ -52,9 +55,31 @@ where state, dag, resolver: ConflictResolver::new(), + merge_policy: None, } } + /// Use an 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. + /// + /// Every replica must install the same policy: replicas merging the same + /// heads under different rules produce different merge nodes and keep + /// re-merging. The policy is process-local and never travels with the + /// data. + 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 @@ -462,8 +487,15 @@ 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)?; + let named_policy; + let policy: &dyn MergePolicy = match &self.merge_policy { + Some(installed) => installed.as_ref(), + None => { + let policy_type = genesis_node.metadata().policy_type(); + named_policy = self.create_policy(policy_type)?; + named_policy.as_ref() + } + }; self.validate_parent_genesis(genesis, &heads)?; @@ -473,7 +505,7 @@ where &self.dag, *genesis, merge_timestamp, - policy.as_ref(), + policy, )?; let (merge_cid, node) = self @@ -1382,6 +1414,131 @@ mod tests { assert!(!heads_after_merge.contains(&branch2_cid)); } + /// A policy installed with `with_merge_policy` decides the merge node's + /// payload — not the "lww" named in the genesis metadata — and 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()]); + } + } + + /// 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(); From 0fbabbf0c054444dc284f7b0f90fc8ae56a52798 Mon Sep 17 00:00:00 2001 From: somasekimoto <0421.soma@gmail.com> Date: Sun, 13 Sep 2026 22:36:12 +0900 Subject: [PATCH 2/3] Add Repo::merge_heads so readers can converge concurrent heads first MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Auto-merge runs lazily, inside the next commit. A caller that wants to *read* the converged state before committing — to copy the current payload into a new version, or to evaluate a policy it carries — gets nothing from `latest()` but one of the concurrent heads picked by timestamp. Any value it copies from there is one branch's, and the merge that follows inside its commit resolves the *old* heads, not the value it just wrote on top of one of them. With a field-wise merge policy that is exactly the case that matters: the reader sees the branch that lacks the other branch's field. `merge_heads(genesis)` runs the same merge the lazy path would, commits it in its own batch, and returns the single head (or the existing one when there was nothing to merge). `heads(genesis)` exposes the current head set for callers that want to know whether a merge is pending. Both are idempotent. Co-Authored-By: Claude Code --- src/repo.rs | 85 +++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 85 insertions(+) diff --git a/src/repo.rs b/src/repo.rs index 7c060b6..2a212e0 100644 --- a/src/repo.rs +++ b/src/repo.rs @@ -115,6 +115,45 @@ 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)); + } + Ok(merged.or_else(|| self.latest(genesis))) + } + /// Convenience wrapper around `DagGraph::get_genesis` pub fn get_genesis(&self, cid: &Cid) -> Result { self.dag.get_genesis(cid).map_err(CrdtError::Graph) @@ -1501,6 +1540,52 @@ mod tests { } } + /// `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. From 2528d7557c52f3d9972e738494d54747fe2d44f2 Mon Sep 17 00:00:00 2001 From: somasekimoto <0421.soma@gmail.com> Date: Tue, 22 Sep 2026 18:40:11 +0900 Subject: [PATCH 3/3] fix: bind merge policy selection to genesis metadata --- readme.md | 28 ++ src/convergence/policies/lww.rs | 4 + src/convergence/policy.rs | 27 +- src/convergence/resolver.rs | 19 +- src/crdt/operation.rs | 13 + src/crdt/reducer.rs | 2 + src/crdt/storage.rs | 82 ++++- src/repo.rs | 134 ++++++-- tests/merge_heads_errors.rs | 116 +++++++ tests/merge_parent_completeness.rs | 211 ++++++++++++ tests/operation_storage_compatibility.rs | 74 +++++ tests/policy_binding.rs | 390 +++++++++++++++++++++++ 12 files changed, 1044 insertions(+), 56 deletions(-) create mode 100644 tests/merge_heads_errors.rs create mode 100644 tests/merge_parent_completeness.rs create mode 100644 tests/operation_storage_compatibility.rs create mode 100644 tests/policy_binding.rs 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 8c3e9d7..7559fe7 100644 --- a/src/convergence/policy.rs +++ b/src/convergence/policy.rs @@ -15,8 +15,10 @@ pub struct ResolveInput

{ /// that distinction — it does not know the payload's shape — so it hands /// the parents over and lets the policy compare. /// - /// Empty for a genesis head (no parents) and for inputs constructed by - /// callers that have no DAG at hand; `LwwMergePolicy` ignores it. + /// 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

, } @@ -47,11 +49,28 @@ impl

ResolveInput

{ /// 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 then calls it for every auto-merge instead of the named policy. +/// 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 941424b..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,12 +81,16 @@ where .get_node(&cid) .map_err(CrdtError::Graph)? .ok_or_else(|| CrdtError::Internal(format!("Head node not found: {cid}")))?; - // A parent that is not in storage is not an error here: a replica - // may hold a head whose ancestry has not fully synced yet. The - // policy sees fewer parents and must cope (LWW never looks). - let mut parent_payloads = Vec::with_capacity(node.parents().len()); - for parent in node.parents() { - if let Some(parent_node) = dag.get_node(parent).map_err(CrdtError::Graph)? { + 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()); } } 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 2a212e0..50010a7 100644 --- a/src/repo.rs +++ b/src/repo.rs @@ -36,8 +36,7 @@ where pub state: CrdtState, pub dag: DagGraph, resolver: ConflictResolver, - /// Application-supplied merge policy. When set it is used for every - /// auto-merge, regardless of the `policy_type` in the genesis metadata. + /// Application-supplied implementation selected by genesis policy name. merge_policy: Option>>, } @@ -59,7 +58,7 @@ where } } - /// Use an application-supplied [`MergePolicy`] for auto-merges. + /// 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 @@ -71,10 +70,22 @@ where /// it does not know the fields; the application does, so it passes the /// rule in here. /// - /// Every replica must install the same policy: replicas merging the same - /// heads under different rules produce different merge nodes and keep - /// re-merging. The policy is process-local and never travels with the - /// data. + /// 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 @@ -85,6 +96,12 @@ where /// 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 /// @@ -151,7 +168,10 @@ where self.rollback_pending_nodes(&pending_nodes); return Err(CrdtError::Storage(status)); } - Ok(merged.or_else(|| self.latest(genesis))) + match merged { + Some(cid) => Ok(Some(cid)), + None => self.dag.calculate_latest(genesis).map_err(CrdtError::Graph), + } } /// Convenience wrapper around `DagGraph::get_genesis` @@ -291,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); @@ -370,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 @@ -399,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(), @@ -437,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(), @@ -458,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(), @@ -486,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) @@ -526,15 +577,17 @@ where .get_node(genesis) .map_err(CrdtError::Graph)? .ok_or_else(|| CrdtError::Internal(format!("Genesis not found: {genesis}")))?; - let named_policy; - let policy: &dyn MergePolicy = match &self.merge_policy { - Some(installed) => installed.as_ref(), - None => { - let policy_type = genesis_node.metadata().policy_type(); - named_policy = self.create_policy(policy_type)?; - named_policy.as_ref() - } - }; + // 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)?; @@ -565,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); @@ -604,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}"))), } } @@ -1453,8 +1523,8 @@ mod tests { assert!(!heads_after_merge.contains(&branch2_cid)); } - /// A policy installed with `with_merge_policy` decides the merge node's - /// payload — not the "lww" named in the genesis metadata — and sees each + /// 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() { 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" + ); + } +}