use hashtree::mst::{value_hash, Mst, MstDiff};
use hashtree::Hash;
use crate::datastore::noxu::{NoxuDatastore, NoxuDatastoreError};
const COMPOSITE_SEP: u8 = 0x00;
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum MstReconcileError {
#[error("mst reconcile datastore: {0}")]
Datastore(#[from] NoxuDatastoreError),
#[error("mst reconcile: malformed composite key")]
MalformedKey,
}
#[must_use]
pub fn composite_key(bucket: &[u8], key: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(bucket.len() + 1 + key.len());
out.extend_from_slice(bucket);
out.push(COMPOSITE_SEP);
out.extend_from_slice(key);
out
}
pub fn split_composite(composite: &[u8]) -> Result<(&[u8], &[u8]), MstReconcileError> {
let idx = composite
.iter()
.position(|b| *b == COMPOSITE_SEP)
.ok_or(MstReconcileError::MalformedKey)?;
Ok((&composite[..idx], &composite[idx + 1..]))
}
#[must_use]
pub fn value_digest(value: &[u8]) -> Hash {
value_hash(value)
}
pub fn build_mst(db: &NoxuDatastore) -> Result<Mst, MstReconcileError> {
let mut pairs: Vec<(Vec<u8>, Hash)> = Vec::new();
db.fold_primary(|bucket, key, value| {
pairs.push((composite_key(bucket, key), value_digest(value)));
Ok(())
})?;
Ok(Mst::from_pairs(pairs))
}
#[derive(Debug, Clone, Eq, PartialEq)]
pub struct ReconcileOutcome {
pub applied: usize,
pub to_push: usize,
pub comparisons: usize,
}
pub trait ObjectSource {
fn fetch(&self, bucket: &[u8], key: &[u8]) -> Result<Option<Vec<u8>>, MstReconcileError>;
}
pub struct DatastoreSource<'a> {
db: &'a NoxuDatastore,
}
impl<'a> DatastoreSource<'a> {
#[must_use]
pub fn new(db: &'a NoxuDatastore) -> Self {
Self { db }
}
}
impl ObjectSource for DatastoreSource<'_> {
fn fetch(&self, bucket: &[u8], key: &[u8]) -> Result<Option<Vec<u8>>, MstReconcileError> {
Ok(self.db.get_object(bucket, key)?)
}
}
pub fn reconcile_pull<S: ObjectSource>(
local_db: &NoxuDatastore,
local: &Mst,
peer: &Mst,
source: &S,
) -> Result<ReconcileOutcome, MstReconcileError> {
let diff: MstDiff = local.diff(peer);
let mut applied = 0usize;
for composite in diff.only_there() {
let (bucket, key) = split_composite(composite)?;
if let Some(value) = source.fetch(bucket, key)? {
local_db.put_object(bucket, key, &value, &[])?;
applied += 1;
}
}
Ok(ReconcileOutcome {
applied,
to_push: diff.only_here().len(),
comparisons: diff.comparisons(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn open_ds() -> (TempDir, NoxuDatastore) {
let dir = TempDir::new().expect("tempdir");
let ds = NoxuDatastore::open_in(dir.path()).expect("open");
(dir, ds)
}
#[test]
fn composite_round_trips() {
let c = composite_key(b"users", b"alice");
let (b, k) = split_composite(&c).expect("split");
assert_eq!(b, b"users");
assert_eq!(k, b"alice");
}
#[test]
fn identical_stores_have_equal_mst_roots() {
let (_da, a) = open_ds();
let (_db, b) = open_ds();
for i in 0..1000u32 {
let k = format!("k{i:06}");
let v = format!("v{i}");
a.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put a");
b.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put b");
}
let ma = build_mst(&a).expect("build a");
let mb = build_mst(&b).expect("build b");
assert_eq!(ma.root(), mb.root());
assert_eq!(ma.diff(&mb).diff_len(), 0);
}
#[test]
fn reconcile_pulls_missing_objects_and_converges() {
let (_da, a) = open_ds();
let (_db, b) = open_ds();
for i in 0..900u32 {
let k = format!("k{i:06}");
let v = format!("v{i}");
a.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put a");
b.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put b");
}
for i in 900..1000u32 {
let k = format!("k{i:06}");
let v = format!("v{i}");
b.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put b extra");
}
let ma = build_mst(&a).expect("build a");
let mb = build_mst(&b).expect("build b");
let source = DatastoreSource::new(&b);
let outcome = reconcile_pull(&a, &ma, &mb, &source).expect("reconcile");
assert_eq!(outcome.applied, 100, "a should pull 100 missing objects");
let ma2 = build_mst(&a).expect("rebuild a");
assert_eq!(ma2.root(), mb.root(), "reconcile must converge roots");
assert_eq!(ma2.diff(&mb).diff_len(), 0);
}
#[test]
fn reconcile_is_divergence_proportional() {
let (_da, a) = open_ds();
let (_db, b) = open_ds();
for i in 0..10000u32 {
let k = format!("k{i:06}");
let v = format!("v{i}");
a.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put a");
b.put_object(b"users", k.as_bytes(), v.as_bytes(), &[])
.expect("put b");
}
for i in 0..20u32 {
let k = format!("k{i:06}");
b.put_object(b"users", k.as_bytes(), b"CHANGED", &[])
.expect("update b");
}
let ma = build_mst(&a).expect("build a");
let mb = build_mst(&b).expect("build b");
let source = DatastoreSource::new(&b);
let outcome = reconcile_pull(&a, &ma, &mb, &source).expect("reconcile");
assert_eq!(outcome.applied, 20);
assert!(
outcome.comparisons < 2000,
"reconcile not divergence-proportional: {} comparisons for 20 diffs / 10000 keys",
outcome.comparisons
);
let ma2 = build_mst(&a).expect("rebuild a");
assert_eq!(ma2.root(), mb.root());
}
}