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