kcode-kweb-db 0.1.0

A convergent signed-DAG store for Kweb nodes and objects
Documentation
use chrono::Utc;
use kcode_kweb_db::{
    Config, Gossip, KwebDb, NodeData, NoopGossip, Owner, Provenance, TransactionPackage, WriterId,
};
use std::sync::{Arc, Mutex};

#[derive(Clone, Default)]
struct Recorder(Arc<Mutex<Vec<TransactionPackage>>>);

impl Gossip for Recorder {
    fn announce(&self, package: TransactionPackage) {
        self.0.lock().unwrap().push(package);
    }
}

impl Recorder {
    fn take(&self) -> Vec<TransactionPackage> {
        std::mem::take(&mut *self.0.lock().unwrap())
    }
}

fn provenance(label: &str) -> Provenance {
    Provenance {
        author: label.into(),
        source: format!("replication/{label}"),
        source_created_at: Utc::now(),
        data: String::new(),
    }
}

fn data(name: &str) -> NodeData {
    NodeData {
        short_name: name.into(),
        short_description: String::new(),
        long_description: String::new(),
        owner: Owner::SelfNode,
        fixed_connections: [None; 3],
        objects: Vec::new(),
    }
}

fn config(key: [u8; 32], writers: Vec<WriterId>, gossip: Arc<dyn Gossip>) -> Config {
    Config {
        signing_key: key,
        writers_by_priority: writers,
        gossip,
    }
}

#[test]
fn missing_parent_activates_with_descendants() {
    let source_root = tempfile::tempdir().unwrap();
    let destination_root = tempfile::tempdir().unwrap();
    let source_key = [11; 32];
    let destination_key = [12; 32];
    let writers = vec![
        WriterId::from_signing_key(&source_key),
        WriterId::from_signing_key(&destination_key),
    ];
    let recorder = Recorder::default();
    let source = KwebDb::open(
        source_root.path(),
        config(source_key, writers.clone(), Arc::new(recorder.clone())),
    )
    .unwrap();
    let destination = KwebDb::open(
        destination_root.path(),
        config(destination_key, writers, Arc::new(NoopGossip)),
    )
    .unwrap();

    let mut first = source.start_transaction(provenance("first")).unwrap();
    let node = first.create_node(data("Initial node")).unwrap();
    first.finalize().unwrap();
    let mut second = source.start_transaction(provenance("second")).unwrap();
    second.update_node(node, data("Updated node")).unwrap();
    second.finalize().unwrap();
    let packages = recorder.take();

    assert!(destination.accept_transaction(packages[1].clone()).unwrap());
    assert!(destination.get_node(node).is_err());
    assert!(!destination.get_node_history(node).unwrap().entries[0].active);
    assert!(destination.accept_transaction(packages[0].clone()).unwrap());
    assert_eq!(
        destination.get_node(node).unwrap().data.short_name,
        "Updated node"
    );
}

#[test]
fn three_replicas_partition_merge_and_rejoin() {
    let roots = [
        tempfile::tempdir().unwrap(),
        tempfile::tempdir().unwrap(),
        tempfile::tempdir().unwrap(),
    ];
    let keys = [[21; 32], [22; 32], [23; 32]];
    let writers = keys
        .iter()
        .map(WriterId::from_signing_key)
        .collect::<Vec<_>>();
    let recorders = [
        Recorder::default(),
        Recorder::default(),
        Recorder::default(),
    ];
    let databases = [
        KwebDb::open(
            roots[0].path(),
            config(keys[0], writers.clone(), Arc::new(recorders[0].clone())),
        )
        .unwrap(),
        KwebDb::open(
            roots[1].path(),
            config(keys[1], writers.clone(), Arc::new(recorders[1].clone())),
        )
        .unwrap(),
        KwebDb::open(
            roots[2].path(),
            config(keys[2], writers, Arc::new(recorders[2].clone())),
        )
        .unwrap(),
    ];

    let mut initial = databases[0]
        .start_transaction(provenance("initial"))
        .unwrap();
    let object = initial.create_object(b"shared object".to_vec()).unwrap();
    let mut initial_data = data("Initial node");
    initial_data.objects.push(object);
    let node = initial.create_node(initial_data).unwrap();
    initial.finalize().unwrap();
    let initial_package = recorders[0].take().pop().unwrap();
    databases[1]
        .accept_transaction(initial_package.clone())
        .unwrap();
    databases[2].accept_transaction(initial_package).unwrap();
    recorders[1].take();
    recorders[2].take();

    let mut left = databases[0].start_transaction(provenance("left")).unwrap();
    let mut left_data = data("Left branch");
    left_data.objects.push(object);
    left.update_node(node, left_data).unwrap();
    let left_id = left.finalize().unwrap();
    let left_package = recorders[0].take().pop().unwrap();

    let mut right = databases[1].start_transaction(provenance("right")).unwrap();
    let mut right_data = data("Right branch");
    right_data.objects.push(object);
    right.update_node(node, right_data).unwrap();
    let right_id = right.finalize().unwrap();
    let right_package = recorders[1].take().pop().unwrap();

    databases[0]
        .accept_transaction(right_package.clone())
        .unwrap();
    databases[1]
        .accept_transaction(left_package.clone())
        .unwrap();
    databases[2].accept_transaction(right_package).unwrap();
    databases[2].accept_transaction(left_package).unwrap();
    assert_eq!(
        databases[0].get_node_history(node).unwrap().frontier.len(),
        2
    );
    assert_eq!(
        databases[0].get_node(node).unwrap().data.short_name,
        "Left branch"
    );

    recorders[2].take();
    let mut merge = databases[2].start_transaction(provenance("merge")).unwrap();
    merge.merge(left_id, right_id).unwrap();
    let mut merged_data = data("Merged node");
    merged_data.objects.push(object);
    merge.update_node(node, merged_data).unwrap();
    merge.finalize().unwrap();
    let merge_package = recorders[2].take().pop().unwrap();
    databases[0]
        .accept_transaction(merge_package.clone())
        .unwrap();
    databases[1].accept_transaction(merge_package).unwrap();

    for database in &databases {
        assert_eq!(
            database.get_node(node).unwrap().data.short_name,
            "Merged node"
        );
        assert_eq!(database.get_node_history(node).unwrap().frontier.len(), 1);
        assert_eq!(database.get_object(object).unwrap(), b"shared object");
    }
}