Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 22 additions & 2 deletions readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,10 @@ CRSL is a Rust library for content versioning and CRDT (Conflict-free Replicated

- **Content Versioning**: Content creation, update, deletion, and history management
- **CRDT Support**: Conflict resolution through Last-Write-Wins (LWW) reducer
- **Auto-Merge**: Automatic conflict resolution when multiple heads exist
- **DAG (Directed Acyclic Graph)**: Efficient version history management
- **LevelDB Storage**: High-performance persistent storage
- **Thread-Safe**: Safe to use in async/await environments with `Mutex`-based storage
- **CID (Content Identifier)**: IPFS-compatible content identifiers

## 🛠️ Usage
Expand Down Expand Up @@ -105,11 +107,17 @@ crsl-lib/
│ │ ├── crdt_state.rs # CRDT state management
│ │ ├── operation.rs # Operation definitions
│ │ ├── reducer.rs # LWW reducer
│ │ ├── storage.rs # Operation storage
│ │ ├── storage.rs # Operation storage (thread-safe)
│ │ └── error.rs # Error definitions
│ ├── convergence/ # Conflict resolution
│ │ ├── resolver.rs # Merge orchestration
│ │ ├── policy.rs # MergePolicy trait
│ │ ├── policies/ # Policy implementations
│ │ │ └── lww.rs # Last-Write-Wins policy
│ │ └── metadata.rs # Content metadata
│ ├── graph/ # DAG graph implementation
│ │ ├── dag.rs # DAG graph management
│ │ ├── storage.rs # Node storage
│ │ ├── storage.rs # Node storage (thread-safe)
│ │ └── error.rs # Graph errors
│ ├── dasl/ # DASL (Distributed Application Storage Layer)
│ ├── masl/ # MASL (Multi-Agent Storage Layer)
Expand Down Expand Up @@ -183,6 +191,11 @@ cargo doc --open
- Conflict resolution through LWW reducer
- Integration with operation storage

### Convergence (`src/convergence/`)
- **MergePolicy trait**: Customizable merge strategies
- **LwwMergePolicy**: Last-Write-Wins merge implementation
- **ConflictResolver**: Automatic merge node creation

### DAG Graph (`src/graph/dag.rs`)
- DAG management for version history
- Node addition, retrieval, and history tracking
Expand All @@ -191,12 +204,19 @@ cargo doc --open
### Repository (`src/repo.rs`)
- Integration of CRDT State and DAG Graph
- Operation commit and history management
- Auto-merge when multiple heads exist
- High-level API provision

### Operations (`src/crdt/operation.rs`)
- Create: New content creation
- Update: Content updates
- Delete: Content deletion
- Merge: Automatic merge operations

### Thread Safety
- `LeveldbStorage` and `LeveldbNodeStorage` use `Mutex` internally
- `OperationStorage` and `NodeStorage` traits require `Send + Sync`
- Safe to use with `Arc<Mutex<Repo>>` in async/await environments

## 📄 License

Expand Down
17 changes: 8 additions & 9 deletions src/convergence/resolver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,13 +117,12 @@ mod tests {
use crate::graph::storage::NodeStorage;
use multihash::Multihash;
use serde::{Deserialize, Serialize};
use std::cell::RefCell;
use std::collections::HashMap;
use std::rc::Rc;
use std::sync::{Arc, Mutex};

#[derive(Clone, Default)]
struct MemoryNodeStorage<P, M> {
nodes: Rc<RefCell<HashMap<Cid, Node<P, M>>>>,
nodes: Arc<Mutex<HashMap<Cid, Node<P, M>>>>,
}

impl<P, M> MemoryNodeStorage<P, M>
Expand All @@ -135,32 +134,32 @@ mod tests {
let cid = node
.content_id()
.map_err(|e| GraphError::NodeOperation(e.to_string()))?;
self.nodes.borrow_mut().insert(cid, node.clone());
self.nodes.lock().unwrap().insert(cid, node.clone());
Ok(())
}
}

impl<P, M> NodeStorage<P, M> for MemoryNodeStorage<P, M>
where
P: Clone + Serialize + for<'de> Deserialize<'de>,
M: Clone + Serialize + for<'de> Deserialize<'de>,
P: Clone + Serialize + for<'de> Deserialize<'de> + Send + Sync,
M: Clone + Serialize + for<'de> Deserialize<'de> + Send + Sync,
{
fn get(&self, content_id: &Cid) -> GraphResult<Option<Node<P, M>>> {
Ok(self.nodes.borrow().get(content_id).cloned())
Ok(self.nodes.lock().unwrap().get(content_id).cloned())
}

fn put(&self, node: &Node<P, M>) -> GraphResult<()> {
self.insert(node)
}

fn delete(&self, content_id: &Cid) -> GraphResult<()> {
self.nodes.borrow_mut().remove(content_id);
self.nodes.lock().unwrap().remove(content_id);
Ok(())
}

fn get_node_map(&self) -> GraphResult<HashMap<Cid, Vec<Cid>>> {
let mut map = HashMap::new();
for (cid, node) in self.nodes.borrow().iter() {
for (cid, node) in self.nodes.lock().unwrap().iter() {
map.insert(*cid, node.parents().to_vec());
}
Ok(map)
Expand Down
32 changes: 16 additions & 16 deletions src/crdt/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,11 @@ use bincode;
use rusty_leveldb::LdbIterator;
use std::marker::PhantomData;
use std::path::Path;
use std::rc::Rc;
use std::sync::Arc;
use ulid::Ulid;

/// Abstraction over the persistent storage used by `CrdtState`.
pub trait OperationStorage<ContentId, T> {
pub trait OperationStorage<ContentId, T>: Send + Sync {
fn save_operation(&self, op: &Operation<ContentId, T>) -> Result<()>;
fn load_operations(&self, genesis: &ContentId) -> Result<Vec<Operation<ContentId, T>>>;
fn get_operation(&self, op_id: &Ulid) -> Result<Option<Operation<ContentId, T>>>;
Expand All @@ -22,7 +22,7 @@ pub trait OperationStorage<ContentId, T> {
/// LevelDB-backed implementation of [`OperationStorage`].
#[derive(Clone)]
pub struct LeveldbStorage<ContentId, T> {
shared: Rc<SharedLeveldb>,
shared: Arc<SharedLeveldb>,
_marker: PhantomData<(ContentId, T)>,
}

Expand All @@ -32,7 +32,7 @@ impl<ContentId, T> LeveldbStorage<ContentId, T> {
Ok(Self::new(shared))
}

pub fn new(shared: Rc<SharedLeveldb>) -> Self {
pub fn new(shared: Arc<SharedLeveldb>) -> Self {
Self {
shared,
_marker: PhantomData,
Expand Down Expand Up @@ -64,7 +64,7 @@ impl<ContentId, T> LeveldbStorage<ContentId, T> {
.with_active_batch(|batch| batch.put(key, value))
.is_none()
{
self.shared.db().borrow_mut().put(key, value)?;
self.shared.db().put(key, value)?;
}
Ok(())
}
Expand All @@ -76,22 +76,27 @@ impl<ContentId, T> LeveldbStorage<ContentId, T> {
.with_active_batch(|batch| batch.delete(key))
.is_none()
{
self.shared.db().borrow_mut().delete(key)?;
self.shared.db().delete(key)?;
}
Ok(())
}
}

impl<ContentId, T> SharedLeveldbAccess for LeveldbStorage<ContentId, T> {
fn shared_leveldb(&self) -> Option<Rc<SharedLeveldb>> {
fn shared_leveldb(&self) -> Option<Arc<SharedLeveldb>> {
Some(self.shared.clone())
}
}

impl<ContentId, T> OperationStorage<ContentId, T> for LeveldbStorage<ContentId, T>
where
ContentId: serde::Serialize + for<'de> serde::Deserialize<'de> + PartialEq + std::fmt::Debug,
T: serde::Serialize + for<'de> serde::Deserialize<'de> + std::fmt::Debug,
ContentId: serde::Serialize
+ for<'de> serde::Deserialize<'de>
+ PartialEq
+ std::fmt::Debug
+ Send
+ Sync,
T: serde::Serialize + for<'de> serde::Deserialize<'de> + std::fmt::Debug + Send + Sync,
{
fn begin_batch(&self) -> std::result::Result<LeveldbBatchGuard<'_>, BatchError> {
self.shared.begin_batch()
Expand All @@ -105,12 +110,7 @@ where

fn load_operations(&self, genesis: &ContentId) -> Result<Vec<Operation<ContentId, T>>> {
let mut result = Vec::new();
let mut iter = self
.shared
.db()
.borrow_mut()
.new_iter()
.map_err(CrdtError::Storage)?;
let mut iter = self.shared.db().new_iter().map_err(CrdtError::Storage)?;
iter.seek_to_first();

let mut key = Vec::new();
Expand All @@ -133,7 +133,7 @@ where

fn get_operation(&self, op_id: &Ulid) -> Result<Option<Operation<ContentId, T>>> {
let key = Self::make_key(op_id);
match self.shared.db().borrow_mut().get(&key) {
match self.shared.db().get(&key) {
Some(raw) => {
let (op, _) = bincode::serde::decode_from_slice::<Operation<ContentId, T>, _>(
&raw,
Expand Down
63 changes: 44 additions & 19 deletions src/graph/dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -442,28 +442,28 @@ where
mod tests {
use super::*;
use crate::graph::storage::LeveldbNodeStorage;
use std::cell::RefCell;
use std::collections::BTreeMap;
use std::sync::Mutex;
use tempfile::tempdir;

type TestDag = DagGraph<MockStorage, String, BTreeMap<String, String>>;

#[derive(Debug)]
struct MockStorage {
edges: std::cell::RefCell<HashMap<Cid, Vec<Cid>>>,
timestamps: std::cell::RefCell<HashMap<Cid, u64>>,
edges: Mutex<HashMap<Cid, Vec<Cid>>>,
timestamps: Mutex<HashMap<Cid, u64>>,
}
impl MockStorage {
fn new() -> Self {
Self {
edges: RefCell::new(HashMap::new()),
timestamps: RefCell::new(HashMap::new()),
edges: Mutex::new(HashMap::new()),
timestamps: Mutex::new(HashMap::new()),
}
}

fn setup_graph(&mut self, structure: &[(Cid, Cid)]) {
let mut edges = self.edges.borrow_mut();
let mut timestamps = self.timestamps.borrow_mut();
let mut edges = self.edges.lock().unwrap();
let mut timestamps = self.timestamps.lock().unwrap();
let mut ts = 1;

for (parent, child) in structure {
Expand All @@ -480,19 +480,24 @@ mod tests {

impl<P, M> NodeStorage<P, M> for MockStorage
where
P: Default + serde::Serialize + serde::de::DeserializeOwned,
M: Default + serde::Serialize + serde::de::DeserializeOwned,
P: Default + serde::Serialize + serde::de::DeserializeOwned + Send + Sync,
M: Default + serde::Serialize + serde::de::DeserializeOwned + Send + Sync,
{
fn get(&self, content_id: &Cid) -> Result<Option<Node<P, M>>> {
let edges = self.edges.borrow();
let edges = self.edges.lock().unwrap();
let parents = edges.get(content_id).cloned();

let parents = match parents {
Some(p) => p,
None => return Ok(None),
};

let ts = *self.timestamps.borrow().get(content_id).unwrap_or(&0);
let ts = *self
.timestamps
.lock()
.unwrap()
.get(content_id)
.unwrap_or(&0);

fn find_genesis(edges: &HashMap<Cid, Vec<Cid>>, cid: &Cid) -> Cid {
let mut current = *cid;
Expand All @@ -504,7 +509,7 @@ mod tests {
}
current
}
let genesis_cid = find_genesis(&self.edges.borrow(), content_id);
let genesis_cid = find_genesis(&edges, content_id);

let node = if parents.is_empty() {
Node::new_genesis(P::default(), ts, M::default())
Expand All @@ -520,8 +525,8 @@ mod tests {
let parents = node.parents().to_vec();
let ts = node.timestamp();

self.edges.borrow_mut().insert(cid, parents);
self.timestamps.borrow_mut().insert(cid, ts);
self.edges.lock().unwrap().insert(cid, parents);
self.timestamps.lock().unwrap().insert(cid, ts);

Ok(())
}
Expand All @@ -531,7 +536,7 @@ mod tests {
}

fn get_node_map(&self) -> Result<HashMap<Cid, Vec<Cid>>> {
Ok(self.edges.borrow().clone())
Ok(self.edges.lock().unwrap().clone())
}
}

Expand Down Expand Up @@ -863,7 +868,12 @@ mod tests {
fn test_get_genesis_from_genesis_node() {
let storage = MockStorage::new();
let genesis_cid = create_test_content_id(b"genesis");
storage.edges.borrow_mut().entry(genesis_cid).or_default();
storage
.edges
.lock()
.unwrap()
.entry(genesis_cid)
.or_default();
let dag = DagGraph::<MockStorage, String, BTreeMap<String, String>>::new(storage);

let result = dag.get_genesis(&genesis_cid);
Expand Down Expand Up @@ -898,7 +908,12 @@ mod tests {
fn test_calculate_latest_genesis_only() {
let storage = MockStorage::new();
let genesis_cid = create_test_content_id(b"genesis");
storage.edges.borrow_mut().entry(genesis_cid).or_default();
storage
.edges
.lock()
.unwrap()
.entry(genesis_cid)
.or_default();
let dag = DagGraph::<MockStorage, String, BTreeMap<String, String>>::new(storage);
let result = dag.calculate_latest(&genesis_cid).unwrap();
assert_eq!(result, Some(genesis_cid));
Expand Down Expand Up @@ -943,7 +958,12 @@ mod tests {
fn test_get_nodes_by_genesis_genesis_only() {
let storage = MockStorage::new();
let genesis_cid = create_test_content_id(b"genesis");
storage.edges.borrow_mut().entry(genesis_cid).or_default();
storage
.edges
.lock()
.unwrap()
.entry(genesis_cid)
.or_default();
let dag = DagGraph::<MockStorage, String, BTreeMap<String, String>>::new(storage);
let result = dag.get_nodes_by_genesis(&genesis_cid).unwrap();
assert_eq!(result, vec![genesis_cid]);
Expand Down Expand Up @@ -971,7 +991,12 @@ mod tests {
let v1_cid = create_test_content_id(b"v1");
let unrelated_cid = create_test_content_id(b"unrelated");
storage.setup_graph(&[(genesis1_cid, v1_cid)]);
storage.edges.borrow_mut().entry(unrelated_cid).or_default();
storage
.edges
.lock()
.unwrap()
.entry(unrelated_cid)
.or_default();
let dag = DagGraph::<MockStorage, String, BTreeMap<String, String>>::new(storage);
let mut result = dag.get_nodes_by_genesis(&genesis1_cid).unwrap();
result.sort();
Expand Down
Loading