use super::action::{DagAction, DagActionHash, DagPayload};
use super::store::json_to_graph_value;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TimeTravelSnapshot {
pub target_hash: DagActionHash,
pub target_timestamp: DateTime<Utc>,
pub actions_replayed: usize,
pub triple_count: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DagDiff {
pub from: DagActionHash,
pub to: DagActionHash,
pub actions: Vec<DagAction>,
}
pub(crate) fn replay_payload(db: &crate::GraphDB, payload: &DagPayload) -> crate::Result<()> {
match payload {
DagPayload::TripleInsert { triples } => {
for t in triples {
let triple = crate::Triple::new(
crate::NodeId::named(&t.subject),
crate::Predicate::named(&t.predicate),
json_to_graph_value(&t.object),
);
let _ = db.insert(triple);
}
}
DagPayload::TripleDelete { triple_ids, .. } => {
for tid_bytes in triple_ids {
let tid = crate::TripleId::new(*tid_bytes);
let _ = db.delete(&tid);
}
}
DagPayload::Batch { ops } => {
for op in ops {
replay_payload(db, op)?;
}
}
_ => {}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::dag::{DagStore, TripleInsertPayload};
use crate::{GraphDB, NodeId, Predicate, Triple, Value};
use chrono::Utc;
fn insert_action(
store: &DagStore,
seq: u64,
subject: &str,
object: &str,
parents: Vec<DagActionHash>,
) -> DagActionHash {
let action = DagAction {
parents,
author: NodeId::named("node:1"),
seq,
timestamp: Utc::now(),
payload: DagPayload::TripleInsert {
triples: vec![TripleInsertPayload {
subject: subject.into(),
predicate: "knows".into(),
object: serde_json::json!(object),
provenance: None,
}],
},
signature: None,
};
store.put(&action).unwrap()
}
#[test]
fn test_replay_triple_insert() {
let db = GraphDB::memory().unwrap();
let payload = DagPayload::TripleInsert {
triples: vec![TripleInsertPayload {
subject: "alice".into(),
predicate: "knows".into(),
object: serde_json::json!("bob"),
provenance: None,
}],
};
replay_payload(&db, &payload).unwrap();
assert_eq!(db.count(), 1);
}
#[test]
fn test_replay_triple_delete() {
let db = GraphDB::memory().unwrap();
let triple = Triple::new(
NodeId::named("alice"),
Predicate::named("knows"),
Value::String("bob".into()),
);
let tid = db.insert(triple).unwrap();
assert_eq!(db.count(), 1);
let payload = DagPayload::TripleDelete {
triple_ids: vec![*tid.as_bytes()],
subjects: vec![],
};
replay_payload(&db, &payload).unwrap();
assert_eq!(db.count(), 0);
}
#[test]
fn test_replay_batch() {
let db = GraphDB::memory().unwrap();
let payload = DagPayload::Batch {
ops: vec![
DagPayload::TripleInsert {
triples: vec![TripleInsertPayload {
subject: "alice".into(),
predicate: "knows".into(),
object: serde_json::json!("bob"),
provenance: None,
}],
},
DagPayload::TripleInsert {
triples: vec![TripleInsertPayload {
subject: "bob".into(),
predicate: "knows".into(),
object: serde_json::json!("charlie"),
provenance: None,
}],
},
],
};
replay_payload(&db, &payload).unwrap();
assert_eq!(db.count(), 2);
}
#[test]
fn test_replay_noop_and_genesis_are_no_ops() {
let db = GraphDB::memory().unwrap();
replay_payload(&db, &DagPayload::Noop).unwrap();
replay_payload(
&db,
&DagPayload::Genesis {
triple_count: 0,
description: "test".into(),
},
)
.unwrap();
assert_eq!(db.count(), 0);
}
#[test]
fn test_dag_at_linear_chain() {
let db = GraphDB::memory_with_dag().unwrap();
let store = db.dag_store().unwrap();
let h1 = insert_action(store, 1, "alice", "bob", vec![]);
let h2 = insert_action(store, 2, "bob", "charlie", vec![h1]);
let h3 = insert_action(store, 3, "charlie", "dave", vec![h2]);
let (snap1, info1) = db.dag_at(&h1).unwrap();
assert_eq!(info1.triple_count, 1);
assert_eq!(info1.actions_replayed, 1);
assert_eq!(snap1.count(), 1);
let (_snap2, info2) = db.dag_at(&h2).unwrap();
assert_eq!(info2.triple_count, 2);
assert_eq!(info2.actions_replayed, 2);
let (snap3, info3) = db.dag_at(&h3).unwrap();
assert_eq!(info3.triple_count, 3);
assert_eq!(info3.actions_replayed, 3);
assert_eq!(snap3.count(), 3);
}
#[test]
fn test_dag_at_branching() {
let db = GraphDB::memory_with_dag().unwrap();
let store = db.dag_store().unwrap();
let h0 = insert_action(store, 0, "root", "x", vec![]);
let ha = insert_action(store, 1, "alice", "bob", vec![h0]);
let hb = insert_action(store, 2, "charlie", "dave", vec![h0]);
let (snap_a, _) = db.dag_at(&ha).unwrap();
assert_eq!(snap_a.count(), 2);
let (snap_b, _) = db.dag_at(&hb).unwrap();
assert_eq!(snap_b.count(), 2);
}
#[test]
fn test_dag_diff() {
let db = GraphDB::memory_with_dag().unwrap();
let store = db.dag_store().unwrap();
let h1 = insert_action(store, 1, "alice", "bob", vec![]);
let h2 = insert_action(store, 2, "bob", "charlie", vec![h1]);
let h3 = insert_action(store, 3, "charlie", "dave", vec![h2]);
let diff = db.dag_diff(&h1, &h3).unwrap();
assert_eq!(diff.actions.len(), 2);
assert_eq!(diff.actions[0].seq, 2); assert_eq!(diff.actions[1].seq, 3); }
#[test]
fn test_dag_at_timestamp() {
let db = GraphDB::memory_with_dag().unwrap();
let store = db.dag_store().unwrap();
let before = Utc::now();
let h1 = insert_action(store, 1, "alice", "bob", vec![]);
let _h2 = insert_action(store, 2, "bob", "charlie", vec![h1]);
let result = db.dag_at_timestamp(&(before - chrono::Duration::seconds(10)));
assert!(
result.is_err() || {
true
}
);
let (snap, info) = db.dag_at_timestamp(&Utc::now()).unwrap();
assert_eq!(info.triple_count, 2);
assert_eq!(snap.count(), 2);
}
}