kcode-kweb-db 1.0.1

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

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

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

#[derive(Default)]
struct RequestRecorder {
    requests: Mutex<Vec<Vec<kcode_kweb_db::TransactionId>>>,
}

impl TransactionSource for RequestRecorder {
    fn request_transactions(&self, transactions: Vec<kcode_kweb_db::TransactionId>) {
        self.requests.lock().unwrap().push(transactions);
    }
}

fn copy_package(package: &TransactionPackage) -> TransactionPackage {
    TransactionPackage {
        transaction: package.transaction.clone(),
        objects: package
            .objects
            .iter()
            .map(|object| kcode_kweb_db::ObjectPayload {
                id: object.id,
                bytes: object.bytes.clone(),
            })
            .collect(),
    }
}

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: Vec::new(),
        recent_connections: Vec::new(),
        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 whole_node_conflict_merge() {
    let roots = [
        tempfile::tempdir().unwrap(),
        tempfile::tempdir().unwrap(),
        tempfile::tempdir().unwrap(),
    ];
    let keys = [[11; 32], [12; 32], [13; 32]];
    let writers = keys
        .iter()
        .map(WriterId::from_signing_key)
        .collect::<Vec<_>>();
    let recorders = [
        Arc::new(Recorder::default()),
        Arc::new(Recorder::default()),
        Arc::new(Recorder::default()),
    ];
    let databases = [
        KwebDb::open(
            roots[0].path(),
            config(keys[0], writers.clone(), recorders[0].clone()),
        )
        .unwrap(),
        KwebDb::open(
            roots[1].path(),
            config(keys[1], writers.clone(), recorders[1].clone()),
        )
        .unwrap(),
        KwebDb::open(
            roots[2].path(),
            config(keys[2], writers, recorders[2].clone()),
        )
        .unwrap(),
    ];

    let mut initial = databases[0]
        .start_transaction(provenance("initial"))
        .unwrap();
    let node = initial.create_node(data("Initial node")).unwrap();
    initial.finalize().unwrap();
    let initial_package = recorders[0].0.lock().unwrap().remove(0);
    databases[1]
        .accept_transaction(copy_package(&initial_package))
        .unwrap();
    databases[2].accept_transaction(initial_package).unwrap();
    recorders[1].0.lock().unwrap().clear();
    recorders[2].0.lock().unwrap().clear();

    let mut left = databases[0].start_transaction(provenance("left")).unwrap();
    let mut left_data = data("Left branch");
    left_data.recent_connections = vec![node];
    left.update_node(node, left_data).unwrap();
    let left_id = left.finalize().unwrap();
    let left_package = recorders[0].0.lock().unwrap().remove(0);

    let mut right = databases[1].start_transaction(provenance("right")).unwrap();
    let mut right_data = data("Right branch");
    right_data.fixed_connections = vec![node];
    right.update_node(node, right_data).unwrap();
    let right_id = right.finalize().unwrap();
    let right_package = recorders[1].0.lock().unwrap().remove(0);

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

    databases[0].accept_transaction(right_package).unwrap();
    databases[1].accept_transaction(left_package).unwrap();
    recorders[2].0.lock().unwrap().clear();
    let mut merge = databases[2].start_transaction(provenance("merge")).unwrap();
    merge.merge(left_id, right_id).unwrap();
    merge.update_node(node, data("Merged node")).unwrap();
    merge.finalize().unwrap();
    let merge_package = recorders[2].0.lock().unwrap().remove(0);
    databases[0]
        .accept_transaction(copy_package(&merge_package))
        .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);
    }
}

#[test]
fn exact_duplicate_is_harmless() {
    let source_root = tempfile::tempdir().unwrap();
    let destination_root = tempfile::tempdir().unwrap();
    let source_key = [21; 32];
    let destination_key = [22; 32];
    let writers = vec![
        WriterId::from_signing_key(&source_key),
        WriterId::from_signing_key(&destination_key),
    ];
    let recorder = Arc::new(Recorder::default());
    let source = KwebDb::open(
        source_root.path(),
        config(source_key, writers.clone(), recorder.clone()),
    )
    .unwrap();
    let destination = KwebDb::open(
        destination_root.path(),
        config(destination_key, writers, Arc::new(NoopGossip)),
    )
    .unwrap();
    let mut transaction = source.start_transaction(provenance("source")).unwrap();
    transaction.create_node(data("Replicated node")).unwrap();
    transaction.finalize().unwrap();
    let package = recorder.0.lock().unwrap().remove(0);
    let duplicate = copy_package(&package);
    assert!(destination.accept_transaction(package).unwrap());
    assert!(!destination.accept_transaction(duplicate).unwrap());
}

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

    let mut parent = source.start_transaction(provenance("parent")).unwrap();
    let node = parent.create_node(data("Parent value")).unwrap();
    let parent_id = parent.finalize().unwrap();

    let mut child = source.start_transaction(provenance("child")).unwrap();
    child.update_node(node, data("Middle value")).unwrap();
    let child_id = child.finalize().unwrap();

    let mut grandchild = source.start_transaction(provenance("grandchild")).unwrap();
    let object = grandchild.create_object(b"spooled bytes".to_vec()).unwrap();
    let mut final_data = data("Child value");
    final_data.objects.push(object);
    grandchild.update_node(node, final_data).unwrap();
    grandchild.finalize().unwrap();

    let mut packages = recorder.0.lock().unwrap();
    let parent_package = packages.remove(0);
    let child_package = packages.remove(0);
    let grandchild_package = packages.remove(0);
    drop(packages);

    let strict_grandchild = copy_package(&grandchild_package);
    let duplicate_grandchild = copy_package(&grandchild_package);
    let duplicate_child = copy_package(&child_package);
    assert!(matches!(
        destination.accept_transaction(strict_grandchild),
        Err(Error::InvalidTransaction(message))
            if message.contains("only gossip admission may wait for dependencies")
    ));
    assert_eq!(
        fs::read_dir(destination_root.path().join("incoming"))
            .unwrap()
            .count(),
        0
    );
    assert!(
        destination
            .accept_gossip_transaction(grandchild_package, &request_recorder)
            .unwrap()
    );
    assert_eq!(
        request_recorder.requests.lock().unwrap().as_slice(),
        &[vec![child_id]]
    );
    assert!(
        !destination
            .accept_gossip_transaction(duplicate_grandchild, &request_recorder)
            .unwrap()
    );
    assert_eq!(
        fs::metadata(destination_root.path().join("transactions.kwl"))
            .unwrap()
            .len(),
        0
    );
    assert!(destination.get_node(node).is_err());
    assert!(destination.get_object(object).is_err());

    assert!(
        destination
            .accept_gossip_transaction(child_package, &request_recorder)
            .unwrap()
    );
    assert_eq!(
        request_recorder.requests.lock().unwrap().as_slice(),
        &[vec![child_id], vec![parent_id]]
    );
    assert_eq!(
        fs::metadata(destination_root.path().join("transactions.kwl"))
            .unwrap()
            .len(),
        0
    );

    assert!(destination.accept_transaction(parent_package).unwrap());
    assert_eq!(
        destination.get_node(node).unwrap().data.short_name,
        "Child value"
    );
    assert_eq!(destination.get_object(object).unwrap(), b"spooled bytes");
    let history = destination.get_node_history(node).unwrap();
    assert_eq!(history.entries.len(), 3);
    assert!(history.entries.iter().all(|entry| entry.active));

    let committed_log_length = fs::metadata(destination_root.path().join("transactions.kwl"))
        .unwrap()
        .len();
    assert!(!destination.accept_transaction(duplicate_child).unwrap());
    assert_eq!(
        fs::metadata(destination_root.path().join("transactions.kwl"))
            .unwrap()
            .len(),
        committed_log_length
    );
}