use crate::db::Db;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::BTreeSet;
use std::sync::atomic::Ordering;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum ChangeKind {
Added,
Removed,
Modified,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct FieldDelta {
pub added: Vec<String>,
pub removed: Vec<String>,
pub changed: Vec<String>,
}
impl FieldDelta {
pub fn is_empty(&self) -> bool {
self.added.is_empty() && self.removed.is_empty() && self.changed.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TemporalDelta {
pub valid_from_before: Option<String>,
pub valid_from_after: Option<String>,
pub valid_to_before: Option<String>,
pub valid_to_after: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DocChange {
pub coll: String,
pub id: String,
pub kind: ChangeKind,
pub before: Option<Value>,
pub after: Option<Value>,
pub fields: Option<FieldDelta>,
pub temporal: Option<TemporalDelta>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CollChange {
pub name: String,
pub kind: ChangeKind,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StateDiff {
pub from_seq: u64,
pub to_seq: u64,
pub collections: Vec<CollChange>,
pub documents: Vec<DocChange>,
}
impl StateDiff {
pub fn is_empty(&self) -> bool {
self.collections.is_empty() && self.documents.is_empty()
}
}
pub const HISTORY_PRUNED: &str = "HISTORY_PRUNED";
pub const REVERSED_RANGE: &str = "REVERSED_RANGE";
pub fn diff(db: &Db, from_seq: u64, to_seq: u64) -> Result<StateDiff, String> {
if from_seq > to_seq {
return Err(format!(
"{}: from_seq {} is after to_seq {}; a diff's argument order fixes the \
sign of every change in it. Ask for diff({}, {}) if you want the \
reverse.",
REVERSED_RANGE, from_seq, to_seq, to_seq, from_seq
));
}
let floor = db.history_floor();
if from_seq < floor || to_seq < floor {
let which = if from_seq < floor { "from_seq" } else { "to_seq" };
let bad = if from_seq < floor { from_seq } else { to_seq };
return Err(format!(
"{}: {} {} is below the history floor {}. compact() discarded the \
superseded versions needed to reconstruct that state, so a diff there \
could not distinguish an unchanged document from an unreadable one. \
The oldest diffable sequence is {}.",
HISTORY_PRUNED, which, bad, floor, floor
));
}
if from_seq == to_seq {
return Ok(StateDiff { from_seq, to_seq, collections: vec![], documents: vec![] });
}
let before_colls: BTreeSet<String> = live_collections(db, from_seq);
let after_colls: BTreeSet<String> = live_collections(db, to_seq);
let mut collections = Vec::new();
for name in after_colls.difference(&before_colls) {
collections.push(CollChange { name: name.clone(), kind: ChangeKind::Added });
}
for name in before_colls.difference(&after_colls) {
collections.push(CollChange { name: name.clone(), kind: ChangeKind::Removed });
}
collections.sort_by(|a, b| a.name.cmp(&b.name));
let mut documents = Vec::new();
for coll in before_colls.union(&after_colls) {
let live_before = before_colls.contains(coll);
let live_after = after_colls.contains(coll);
for id in db.list_ids_including_deleted(coll) {
let before = live_before.then(|| db.get_as_of(coll, &id, from_seq)).flatten();
let after = live_after.then(|| db.get_as_of(coll, &id, to_seq)).flatten();
if let Some(change) = classify(coll, &id, before, after) {
documents.push(change);
}
}
}
documents.sort_by(|a, b| a.coll.cmp(&b.coll).then_with(|| a.id.cmp(&b.id)));
Ok(StateDiff { from_seq, to_seq, collections, documents })
}
fn live_collections(db: &Db, seq: u64) -> BTreeSet<String> {
db.collections_as_of(seq)
.into_iter()
.filter(|c| !crate::namespace::is_reserved(c))
.collect()
}
fn classify(
coll: &str,
id: &str,
before: Option<crate::store::Node>,
after: Option<crate::store::Node>,
) -> Option<DocChange> {
match (before, after) {
(None, None) => None,
(None, Some(a)) => Some(DocChange {
coll: coll.to_string(),
id: id.to_string(),
kind: ChangeKind::Added,
before: None,
after: Some(a.data),
fields: None,
temporal: None,
}),
(Some(b), None) => Some(DocChange {
coll: coll.to_string(),
id: id.to_string(),
kind: ChangeKind::Removed,
before: Some(b.data),
after: None,
fields: None,
temporal: None,
}),
(Some(b), Some(a)) => {
let temporal_moved =
b.valid_from != a.valid_from || b.valid_to != a.valid_to;
if b.data == a.data && !temporal_moved {
return None;
}
let fields = field_delta(&b.data, &a.data);
Some(DocChange {
coll: coll.to_string(),
id: id.to_string(),
kind: ChangeKind::Modified,
before: Some(b.data),
after: Some(a.data),
fields,
temporal: temporal_moved.then(|| TemporalDelta {
valid_from_before: b.valid_from,
valid_from_after: a.valid_from,
valid_to_before: b.valid_to,
valid_to_after: a.valid_to,
}),
})
}
}
}
fn field_delta(before: &Value, after: &Value) -> Option<FieldDelta> {
let (b, a) = match (before.as_object(), after.as_object()) {
(Some(b), Some(a)) => (b, a),
_ => return None,
};
let mut delta = FieldDelta::default();
let keys: BTreeSet<&String> = b.keys().chain(a.keys()).collect();
for k in keys {
match (b.get(k), a.get(k)) {
(None, Some(_)) => delta.added.push(k.clone()),
(Some(_), None) => delta.removed.push(k.clone()),
(Some(bv), Some(av)) if bv != av => delta.changed.push(k.clone()),
_ => {}
}
}
Some(delta)
}
pub fn tip(db: &Db) -> u64 {
db.seq.load(Ordering::SeqCst).saturating_sub(1)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn db() -> Db {
Db::in_memory()
}
fn put(db: &Db, coll: &str, id: &str, data: Value) -> u64 {
db.put(coll, id, data, vec![], None, None)
.expect("put should succeed")
.seq
}
fn find<'a>(d: &'a StateDiff, coll: &str, id: &str) -> Option<&'a DocChange> {
d.documents.iter().find(|c| c.coll == coll && c.id == id)
}
#[test]
fn added_document_is_detected() {
let db = db();
let from = put(&db, "users", "a", json!({"n": 1}));
put(&db, "users", "b", json!({"n": 2}));
let d = diff(&db, from, tip(&db)).unwrap();
let c = find(&d, "users", "b").expect("b should appear");
assert_eq!(c.kind, ChangeKind::Added);
assert_eq!(c.before, None);
assert_eq!(c.after, Some(json!({"n": 2})));
assert!(c.fields.is_none(), "an add has no field delta; `after` is the whole story");
}
#[test]
fn removed_document_is_detected() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
db.delete("users", "a").unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
let c = find(&d, "users", "a").expect("a should appear");
assert_eq!(c.kind, ChangeKind::Removed);
assert_eq!(c.before, Some(json!({"n": 1})));
assert_eq!(c.after, None);
}
#[test]
fn modified_document_is_detected() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
put(&db, "users", "a", json!({"n": 2}));
let d = diff(&db, from, tip(&db)).unwrap();
let c = find(&d, "users", "a").expect("a should appear");
assert_eq!(c.kind, ChangeKind::Modified);
assert_eq!(c.before, Some(json!({"n": 1})));
assert_eq!(c.after, Some(json!({"n": 2})));
}
#[test]
fn document_created_and_deleted_inside_the_range_does_not_appear() {
let db = db();
put(&db, "users", "keep", json!({"n": 0}));
let from = tip(&db);
put(&db, "users", "ghost", json!({"n": 1}));
db.delete("users", "ghost").unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
assert!(
find(&d, "users", "ghost").is_none(),
"absent on both sides is not a state change: {:?}",
d.documents
);
assert!(d.documents.is_empty(), "{:?}", d.documents);
}
#[test]
fn unchanged_document_does_not_appear() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
put(&db, "users", "b", json!({"n": 2}));
let d = diff(&db, from, tip(&db)).unwrap();
assert!(find(&d, "users", "a").is_none());
assert_eq!(d.documents.len(), 1);
}
#[test]
fn rewriting_identical_content_is_not_a_change() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
for _ in 0..5 {
put(&db, "users", "a", json!({"n": 1}));
}
let d = diff(&db, from, tip(&db)).unwrap();
assert!(d.is_empty(), "{:?}", d);
}
#[test]
fn value_restored_to_its_original_is_not_a_change() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
put(&db, "users", "a", json!({"n": 99}));
put(&db, "users", "a", json!({"n": 1}));
let d = diff(&db, from, tip(&db)).unwrap();
assert!(d.is_empty(), "net state is identical: {:?}", d);
}
#[test]
fn created_collection_appears_as_added() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
put(&db, "orders", "o1", json!({"total": 10}));
let d = diff(&db, from, tip(&db)).unwrap();
assert_eq!(
d.collections,
vec![CollChange { name: "orders".into(), kind: ChangeKind::Added }]
);
assert_eq!(find(&d, "orders", "o1").unwrap().kind, ChangeKind::Added);
}
#[test]
fn emptied_but_live_collection_is_not_a_removed_collection() {
let db = db();
put(&db, "orders", "o1", json!({"total": 10}));
let from = tip(&db);
db.delete("orders", "o1").unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
assert!(
d.collections.is_empty(),
"an empty collection still exists; only a drop removes it: {:?}",
d.collections
);
assert_eq!(find(&d, "orders", "o1").unwrap().kind, ChangeKind::Removed);
}
#[test]
fn dropped_collection_appears_as_removed_with_its_rows() {
let db = db();
put(&db, "orders", "o1", json!({"total": 10}));
put(&db, "orders", "o2", json!({"total": 20}));
let from = tip(&db);
assert!(db.drop_collection("orders").unwrap());
let d = diff(&db, from, tip(&db)).unwrap();
assert_eq!(
d.collections,
vec![CollChange { name: "orders".into(), kind: ChangeKind::Removed }]
);
assert_eq!(d.documents.len(), 2);
assert!(d.documents.iter().all(|c| c.kind == ChangeKind::Removed));
}
#[test]
fn reserved_collections_never_appear() {
let db = db();
put(&db, "seed", "s", json!({}));
let from = tip(&db);
put(&db, "users", "a", json!({"n": 1}));
db.create_root().unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
assert!(
d.collections.iter().all(|c| !c.name.starts_with("_nedb")),
"{:?}",
d.collections
);
assert!(
d.documents.iter().all(|c| !c.coll.starts_with("_nedb")),
"{:?}",
d.documents
);
assert_eq!(
d.collections,
vec![CollChange { name: "users".into(), kind: ChangeKind::Added }]
);
}
#[test]
fn below_the_history_floor_is_refused_not_silently_empty() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
put(&db, "users", "a", json!({"n": 2}));
db.set_history_floor(tip(&db)).unwrap();
let floor = db.history_floor();
assert!(floor > 0, "precondition: the database is in the pruned state");
let err = diff(&db, 0, tip(&db)).unwrap_err();
assert!(err.starts_with(HISTORY_PRUNED), "{}", err);
assert!(err.contains(&floor.to_string()), "the reason must name the floor: {}", err);
assert!(diff(&db, floor, tip(&db)).is_ok());
}
#[test]
fn history_floor_refusal_also_covers_to_seq() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
put(&db, "users", "a", json!({"n": 2}));
db.set_history_floor(tip(&db)).unwrap();
let floor = db.history_floor();
let err = diff(&db, 0, 0).unwrap_err();
assert!(err.starts_with(HISTORY_PRUNED), "{}", err);
assert!(err.contains("from_seq"), "{}", err);
let err = diff(&db, floor, floor - 1).unwrap_err();
assert!(err.starts_with(REVERSED_RANGE), "{}", err);
}
#[test]
fn reversed_range_is_refused() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let t = tip(&db);
let err = diff(&db, t, 0).unwrap_err();
assert!(err.starts_with(REVERSED_RANGE), "{}", err);
assert!(err.contains(&format!("diff({}, {})", 0, t)), "must name the fix: {}", err);
}
#[test]
fn field_delta_names_added_removed_and_changed_keys() {
let db = db();
put(&db, "users", "a", json!({"keep": 1, "drop": 2, "move": 3}));
let from = tip(&db);
put(&db, "users", "a", json!({"keep": 1, "move": 4, "new": 5}));
let d = diff(&db, from, tip(&db)).unwrap();
let f = find(&d, "users", "a").unwrap().fields.as_ref().expect("objects get a delta");
assert_eq!(f.added, vec!["new".to_string()]);
assert_eq!(f.removed, vec!["drop".to_string()]);
assert_eq!(f.changed, vec!["move".to_string()]);
}
#[test]
fn no_field_delta_when_either_side_is_not_an_object() {
let db = db();
put(&db, "vals", "scalar", json!(1));
put(&db, "vals", "arr", json!([1, 2]));
let from = tip(&db);
put(&db, "vals", "scalar", json!({"n": 1}));
put(&db, "vals", "arr", json!([1, 2, 3]));
let d = diff(&db, from, tip(&db)).unwrap();
assert!(find(&d, "vals", "scalar").unwrap().fields.is_none());
assert!(find(&d, "vals", "arr").unwrap().fields.is_none());
}
#[test]
fn valid_from_change_alone_counts_as_modified() {
let db = db();
db.put("users", "a", json!({"n": 1}), vec![], Some("2020-01-01".into()), None)
.unwrap();
let from = tip(&db);
db.put("users", "a", json!({"n": 1}), vec![], Some("2021-01-01".into()), None)
.unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
let c = find(&d, "users", "a").expect("a temporal move is a state change");
assert_eq!(c.kind, ChangeKind::Modified);
assert_eq!(c.before, c.after, "payload identical; only the window moved");
let t = c.temporal.as_ref().expect("the reason must be reported");
assert_eq!(t.valid_from_before.as_deref(), Some("2020-01-01"));
assert_eq!(t.valid_from_after.as_deref(), Some("2021-01-01"));
assert!(
c.fields.as_ref().unwrap().is_empty(),
"validity is not a field, so no key moved"
);
}
#[test]
fn valid_to_change_alone_counts_as_modified() {
let db = db();
db.put("users", "a", json!({"n": 1}), vec![], None, None).unwrap();
let from = tip(&db);
db.put("users", "a", json!({"n": 1}), vec![], None, Some("2030-01-01".into()))
.unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
let c = find(&d, "users", "a").unwrap();
assert_eq!(c.kind, ChangeKind::Modified);
let t = c.temporal.as_ref().unwrap();
assert_eq!(t.valid_to_before, None);
assert_eq!(t.valid_to_after.as_deref(), Some("2030-01-01"));
}
#[test]
fn no_temporal_delta_when_the_window_did_not_move() {
let db = db();
db.put("users", "a", json!({"n": 1}), vec![], Some("2020-01-01".into()), None)
.unwrap();
let from = tip(&db);
db.put("users", "a", json!({"n": 2}), vec![], Some("2020-01-01".into()), None)
.unwrap();
let d = diff(&db, from, tip(&db)).unwrap();
assert!(find(&d, "users", "a").unwrap().temporal.is_none());
}
#[test]
fn diff_of_a_point_with_itself_is_empty() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
put(&db, "orders", "o1", json!({"n": 2}));
let t = tip(&db);
let d = diff(&db, t, t).unwrap();
assert!(d.is_empty(), "{:?}", d);
assert_eq!((d.from_seq, d.to_seq), (t, t));
assert!(diff(&db, 1, 1).unwrap().is_empty());
}
#[test]
fn output_ordering_is_deterministic() {
let db = db();
put(&db, "seed", "s", json!({}));
let from = tip(&db);
for (coll, id) in [
("zeta", "m"), ("alpha", "z"), ("zeta", "a"),
("alpha", "b"), ("mid", "q"), ("alpha", "a"),
] {
put(&db, coll, id, json!({"v": id}));
}
let to = tip(&db);
let first = diff(&db, from, to).unwrap();
let second = diff(&db, from, to).unwrap();
assert_eq!(first, second, "two runs over the same range must agree");
let colls: Vec<&str> = first.collections.iter().map(|c| c.name.as_str()).collect();
assert_eq!(colls, vec!["alpha", "mid", "zeta"]);
let docs: Vec<(&str, &str)> = first
.documents
.iter()
.map(|c| (c.coll.as_str(), c.id.as_str()))
.collect();
assert_eq!(
docs,
vec![
("alpha", "a"), ("alpha", "b"), ("alpha", "z"),
("mid", "q"),
("zeta", "a"), ("zeta", "m"),
]
);
}
#[test]
fn recreated_document_is_modified_not_added() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
db.delete("users", "a").unwrap();
put(&db, "users", "a", json!({"n": 2}));
let d = diff(&db, from, tip(&db)).unwrap();
let c = find(&d, "users", "a").unwrap();
assert_eq!(c.kind, ChangeKind::Modified);
assert_eq!(c.before, Some(json!({"n": 1})));
assert_eq!(c.after, Some(json!({"n": 2})));
}
#[test]
fn an_empty_diff_implies_equal_state_roots() {
let db = db();
put(&db, "users", "a", json!({"n": 1}));
let from = tip(&db);
put(&db, "users", "a", json!({"n": 99}));
put(&db, "users", "a", json!({"n": 1}));
let to = tip(&db);
assert!(diff(&db, from, to).unwrap().is_empty());
assert_eq!(
db.state_root_as_of(from).unwrap().state_root,
db.state_root_as_of(to).unwrap().state_root
);
}
}