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
);
}