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