use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use super::version_store::VersionedTriple;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum ChangeOperation {
Insert,
Delete,
}
impl std::fmt::Display for ChangeOperation {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ChangeOperation::Insert => write!(f, "INSERT"),
ChangeOperation::Delete => write!(f, "DELETE"),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ChangeEntry {
pub sequence: u64,
pub operation: ChangeOperation,
pub triple: VersionedTriple,
pub transaction_id: u64,
pub timestamp: DateTime<Utc>,
}
impl ChangeEntry {
pub fn describe(&self) -> String {
format!(
"[{}] seq={} txn={} ({}, {}, {})",
self.operation,
self.sequence,
self.transaction_id,
self.triple.subject,
self.triple.predicate,
self.triple.object,
)
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct ChangeLogStats {
pub total_entries: usize,
pub inserts: usize,
pub deletes: usize,
pub distinct_transactions: usize,
pub earliest: Option<DateTime<Utc>>,
pub latest: Option<DateTime<Utc>>,
}
#[derive(Debug, Default)]
pub struct ChangeLog {
entries: Vec<ChangeEntry>,
next_sequence: u64,
}
impl ChangeLog {
pub fn new() -> Self {
Self {
entries: Vec::new(),
next_sequence: 1,
}
}
pub fn record_insert(&mut self, triple: VersionedTriple, transaction_id: u64) {
let seq = self.next_sequence;
self.next_sequence += 1;
self.entries.push(ChangeEntry {
sequence: seq,
operation: ChangeOperation::Insert,
triple,
transaction_id,
timestamp: Utc::now(),
});
}
pub fn record_delete(&mut self, triple: VersionedTriple, transaction_id: u64) {
let seq = self.next_sequence;
self.next_sequence += 1;
self.entries.push(ChangeEntry {
sequence: seq,
operation: ChangeOperation::Delete,
triple,
transaction_id,
timestamp: Utc::now(),
});
}
pub fn all_entries(&self) -> &[ChangeEntry] {
&self.entries
}
pub fn entries_in_range(&self, from: DateTime<Utc>, to: DateTime<Utc>) -> Vec<&ChangeEntry> {
self.entries
.iter()
.filter(|e| e.timestamp >= from && e.timestamp < to)
.collect()
}
pub fn entries_for_transaction(&self, transaction_id: u64) -> Vec<&ChangeEntry> {
self.entries
.iter()
.filter(|e| e.transaction_id == transaction_id)
.collect()
}
pub fn inserts(&self) -> Vec<&ChangeEntry> {
self.entries
.iter()
.filter(|e| e.operation == ChangeOperation::Insert)
.collect()
}
pub fn deletes(&self) -> Vec<&ChangeEntry> {
self.entries
.iter()
.filter(|e| e.operation == ChangeOperation::Delete)
.collect()
}
pub fn stats(&self) -> ChangeLogStats {
self.stats_for(&self.entries)
}
pub fn stats_in_range(&self, from: DateTime<Utc>, to: DateTime<Utc>) -> ChangeLogStats {
let range: Vec<&ChangeEntry> = self.entries_in_range(from, to);
let cloned: Vec<ChangeEntry> = range.into_iter().cloned().collect();
self.stats_for(&cloned)
}
pub fn len(&self) -> usize {
self.entries.len()
}
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
fn stats_for(&self, entries: &[ChangeEntry]) -> ChangeLogStats {
let mut stats = ChangeLogStats {
total_entries: entries.len(),
..Default::default()
};
let mut txn_ids = std::collections::HashSet::new();
for e in entries {
match e.operation {
ChangeOperation::Insert => stats.inserts += 1,
ChangeOperation::Delete => stats.deletes += 1,
}
txn_ids.insert(e.transaction_id);
stats.earliest = Some(match stats.earliest {
None => e.timestamp,
Some(prev) => prev.min(e.timestamp),
});
stats.latest = Some(match stats.latest {
None => e.timestamp,
Some(prev) => prev.max(e.timestamp),
});
}
stats.distinct_transactions = txn_ids.len();
stats
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;
fn ts(y: i32, m: u32, d: u32) -> DateTime<Utc> {
Utc.with_ymd_and_hms(y, m, d, 0, 0, 0).unwrap()
}
fn make_triple(s: &str, p: &str, o: &str) -> VersionedTriple {
VersionedTriple::new(s, p, o, Utc::now(), 1)
}
#[test]
fn test_record_and_retrieve() {
let mut log = ChangeLog::new();
log.record_insert(make_triple("s", "p", "o"), 1);
log.record_delete(make_triple("s", "p", "o"), 2);
assert_eq!(log.len(), 2);
assert_eq!(log.inserts().len(), 1);
assert_eq!(log.deletes().len(), 1);
}
#[test]
fn test_sequence_numbers() {
let mut log = ChangeLog::new();
log.record_insert(make_triple("a", "b", "c"), 10);
log.record_insert(make_triple("x", "y", "z"), 11);
let entries = log.all_entries();
assert_eq!(entries[0].sequence, 1);
assert_eq!(entries[1].sequence, 2);
}
#[test]
fn test_entries_for_transaction() {
let mut log = ChangeLog::new();
log.record_insert(make_triple("s1", "p", "o"), 42);
log.record_insert(make_triple("s2", "p", "o"), 43);
log.record_delete(make_triple("s1", "p", "o"), 42);
let txn42 = log.entries_for_transaction(42);
assert_eq!(txn42.len(), 2);
let txn43 = log.entries_for_transaction(43);
assert_eq!(txn43.len(), 1);
}
#[test]
fn test_stats() {
let mut log = ChangeLog::new();
log.record_insert(make_triple("s", "p", "o1"), 1);
log.record_insert(make_triple("s", "p", "o2"), 2);
log.record_delete(make_triple("s", "p", "o1"), 3);
let stats = log.stats();
assert_eq!(stats.total_entries, 3);
assert_eq!(stats.inserts, 2);
assert_eq!(stats.deletes, 1);
assert_eq!(stats.distinct_transactions, 3);
}
#[test]
fn test_describe() {
let mut log = ChangeLog::new();
log.record_insert(make_triple("Alice", "knows", "Bob"), 7);
let desc = log.all_entries()[0].describe();
assert!(desc.contains("INSERT"));
assert!(desc.contains("Alice"));
assert!(desc.contains("Bob"));
}
#[test]
fn test_empty_log() {
let log = ChangeLog::new();
assert!(log.is_empty());
let stats = log.stats();
assert_eq!(stats.total_entries, 0);
assert!(stats.earliest.is_none());
}
}