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, ObjectPayload, Owner, Provenance,
    TransactionPackage, WriterId,
};
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
    time::Duration,
};

#[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: "Kennedy".into(),
        source: label.into(),
        source_created_at: Utc::now(),
        data: String::new(),
    }
}

fn data(name: &str, objects: Vec<kcode_kweb_db::ObjectId>) -> NodeData {
    NodeData {
        short_name: name.into(),
        short_description: "A test node".into(),
        long_description: "Complete node data for an integration test.".into(),
        owner: Owner::SelfNode,
        fixed_connections: [None; 3],
        objects,
    }
}

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

#[test]
fn builder_drop_reads_history_objects_and_callbacks() {
    let root = tempfile::tempdir().unwrap();
    let key = [1; 32];
    let writer = WriterId::from_signing_key(&key);
    let recorder = Recorder::default();
    let database = KwebDb::open(
        root.path(),
        config(key, vec![writer], Arc::new(recorder.clone())),
    )
    .unwrap();

    let mut transaction = database.start_transaction(provenance("create")).unwrap();
    let object = transaction.create_object(b"exact bytes".to_vec()).unwrap();
    let node = transaction
        .create_node(data("First node", vec![object]))
        .unwrap();
    assert!(database.get_node(node).is_err());
    let transaction_id = transaction.finalize().unwrap();

    assert_eq!(
        database.get_node(node).unwrap().data.short_name,
        "First node"
    );
    assert_eq!(database.get_object(object).unwrap(), b"exact bytes");
    let history = database.get_node_history(node).unwrap();
    assert_eq!(history.frontier, vec![transaction_id]);
    assert_eq!(history.visible, Some(transaction_id));
    assert_eq!(history.entries.len(), 1);
    assert_eq!(recorder.take().len(), 1);

    let mut dropped = database.start_transaction(provenance("drop")).unwrap();
    let dropped_node = dropped
        .create_node(data("Dropped node", Vec::new()))
        .unwrap();
    drop(dropped);
    assert!(database.get_node(dropped_node).is_err());
}

#[test]
fn ingress_duplicates_and_exclusive_builders() {
    let source_root = tempfile::tempdir().unwrap();
    let destination_root = tempfile::tempdir().unwrap();
    let source_key = [2; 32];
    let destination_key = [3; 32];
    let source_writer = WriterId::from_signing_key(&source_key);
    let destination_writer = WriterId::from_signing_key(&destination_key);
    let writers = vec![source_writer, destination_writer];
    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 transaction = source.start_transaction(provenance("source")).unwrap();
    let node = transaction
        .create_node(data("Replicated", Vec::new()))
        .unwrap();
    transaction.finalize().unwrap();
    let package = recorder.take().pop().unwrap();
    assert!(destination.accept_transaction(package.clone()).unwrap());
    assert!(!destination.accept_transaction(package).unwrap());
    assert_eq!(
        destination.get_node(node).unwrap().data.short_name,
        "Replicated"
    );

    let held = destination.start_transaction(provenance("held")).unwrap();
    let clone = destination.clone();
    let (sender, receiver) = mpsc::channel();
    let handle = thread::spawn(move || {
        let _second = clone.start_transaction(provenance("second")).unwrap();
        sender.send(()).unwrap();
    });
    assert!(receiver.recv_timeout(Duration::from_millis(100)).is_err());
    drop(held);
    receiver.recv_timeout(Duration::from_secs(2)).unwrap();
    handle.join().unwrap();
}

#[test]
fn connections_are_additive_and_recency_ordered() {
    let root = tempfile::tempdir().unwrap();
    let key = [4; 32];
    let writer = WriterId::from_signing_key(&key);
    let database =
        KwebDb::open(root.path(), config(key, vec![writer], Arc::new(NoopGossip))).unwrap();
    let mut create = database.start_transaction(provenance("nodes")).unwrap();
    let first = create.create_node(data("First node", Vec::new())).unwrap();
    let second = create.create_node(data("Second node", Vec::new())).unwrap();
    let third = create.create_node(data("Third node", Vec::new())).unwrap();
    create.finalize().unwrap();

    let mut older = database.start_transaction(provenance("older")).unwrap();
    older.connect_node(first, second).unwrap();
    older.finalize().unwrap();
    let mut newer = database.start_transaction(provenance("newer")).unwrap();
    newer.connect_node(first, third).unwrap();
    newer.finalize().unwrap();
    assert_eq!(
        database.get_node(first).unwrap().connections,
        vec![third, second]
    );
}

#[test]
fn invalid_config_references_and_id_text_fail_closed() {
    let root = tempfile::tempdir().unwrap();
    let key = [5; 32];
    assert!(KwebDb::open(root.path(), config(key, Vec::new(), Arc::new(NoopGossip))).is_err());

    let root = tempfile::tempdir().unwrap();
    let writer = WriterId::from_signing_key(&key);
    let database =
        KwebDb::open(root.path(), config(key, vec![writer], Arc::new(NoopGossip))).unwrap();
    let mut transaction = database.start_transaction(provenance("bad-ref")).unwrap();
    transaction
        .create_node(data("Bad reference", vec![kcode_kweb_db::ObjectId([9; 6])]))
        .unwrap();
    assert!(transaction.finalize().is_err());
    assert!("ABCDEFABCDEF".parse::<kcode_kweb_db::NodeId>().is_err());
}

#[test]
fn package_serde_preserves_exact_bytes() {
    let package = TransactionPackage {
        transaction: vec![0, 1, 255],
        objects: vec![ObjectPayload {
            id: kcode_kweb_db::ObjectId([1; 6]),
            bytes: vec![2, 3, 4],
        }],
    };
    let json = serde_json::to_string(&package).unwrap();
    assert_eq!(
        serde_json::from_str::<TransactionPackage>(&json).unwrap(),
        package
    );
}