#[path = "common/harness.rs"]
mod harness;
use std::time::Duration;
use harness::TestHarness;
use macrame::graph::{EdgeAssertion, TraversalBuilder};
use macrame::integrity::audit_current;
use macrame::{BranchId, ConceptUpsert, Database};
const EPOCH: &str = "1970-01-01T00:00:00.000000Z";
const OPEN: &str = "9999-12-31T23:59:59.999999Z";
const T1: &str = "1970-01-02T00:00:00.000000Z";
const T2: &str = "1970-01-03T00:00:00.000000Z";
const LATE: &str = "2999-01-01T00:00:00.000000Z";
const STEP: Duration = Duration::from_secs(3_600);
async fn seed(h: &TestHarness) -> Database {
let db = h.db_with_fake_clock().await;
db.write_concepts(
["a", "b", "c"]
.iter()
.map(|n| ConceptUpsert::new(*n, "n").valid_from(EPOCH))
.collect(),
)
.await
.unwrap();
db.bulk_import(vec![
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN),
EdgeAssertion::new("b", "c", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN),
])
.await
.unwrap();
h.advance(STEP);
db
}
fn id(name: &str) -> BranchId {
BranchId::new(name).unwrap()
}
async fn reached_at(db: &Database, branch: Option<&str>, ts: &str) -> Vec<String> {
let mut b = TraversalBuilder::new("a");
if let Some(n) = branch {
b = b.on_branch(id(n));
}
let mut v = b.execute_ids(db.read_conn(), ts).await.unwrap();
v.sort();
v
}
async fn reached(db: &Database, branch: Option<&str>) -> Vec<String> {
reached_at(db, branch, EPOCH).await
}
async fn rows_at(db: &Database, source: &str, target: &str) -> Vec<String> {
let mut rows = db
.read_conn()
.query(
"SELECT branch_id || ' ' || valid_to FROM links \
WHERE source_id = ?1 AND target_id = ?2 ORDER BY branch_id, valid_to",
libsql::params![source, target],
)
.await
.unwrap();
let mut out = Vec::new();
while let Some(r) = rows.next().await.unwrap() {
out.push(r.get::<String>(0).unwrap());
}
out
}
async fn counts(db: &Database) -> (i64, i64) {
let conn = db.read_conn();
let one = |sql: &'static str| {
let conn = conn.clone();
async move {
conn.query(sql, ())
.await
.unwrap()
.next()
.await
.unwrap()
.unwrap()
.get::<i64>(0)
.unwrap()
}
};
(
one("SELECT COUNT(*) FROM links").await,
one("SELECT COUNT(*) FROM transaction_log WHERE table_name = 'links'").await,
)
}
#[tokio::test]
async fn a_branch_writing_at_the_trunks_key_does_not_archive_the_trunks_belief() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(2.0)
.on_branch(alt.id.clone()),
)
.await
.unwrap();
h.advance(STEP);
assert_eq!(reached(&db, None).await, ["a", "b", "c"]);
db.archive(LATE).await.unwrap();
assert_eq!(
reached(&db, None).await,
["a", "b", "c"],
"the trunk lost an edge it still believed, because a branch disagreed with it"
);
assert_eq!(reached(&db, Some("alt")).await, ["a", "b", "c"]);
db.close().await.unwrap();
}
#[tokio::test]
async fn the_trunk_writing_after_a_fork_does_not_archive_the_branchs_belief() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(2.0)
.on_branch(alt.id.clone()),
)
.await
.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(3.0),
)
.await
.unwrap();
h.advance(STEP);
db.archive(LATE).await.unwrap();
assert_eq!(
reached(&db, Some("alt")).await,
["a", "b", "c"],
"the branch lost an edge because the trunk wrote at the same key after it forked"
);
assert_eq!(reached(&db, None).await, ["a", "b", "c"]);
db.close().await.unwrap();
}
#[tokio::test]
async fn archiving_a_shadow_row_does_not_resurrect_what_the_branch_retired() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.retire_edge_on("b", "c", "LEADSTO", EPOCH, T1, alt.id.clone())
.await
.unwrap();
h.advance(STEP);
assert_eq!(reached_at(&db, Some("alt"), T2).await, ["a", "b"]);
assert_eq!(reached_at(&db, None, T2).await, ["a", "b", "c"]);
db.archive(LATE).await.unwrap();
assert_eq!(
reached_at(&db, Some("alt"), T2).await,
["a", "b"],
"the archive un-retired an edge -- a belief resurrected by an operation \
that asserts nothing"
);
assert_eq!(
reached_at(&db, None, T2).await,
["a", "b", "c"],
"and the ancestor's own row was never the branch's to affect"
);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_key_two_lineages_hold_keeps_its_closed_intervals_hot() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.retire_edge_on("b", "c", "LEADSTO", EPOCH, T1, alt.id.clone())
.await
.unwrap();
h.advance(STEP);
db.retire_edge("b", "c", "LEADSTO", EPOCH, T1)
.await
.unwrap();
h.advance(STEP);
let report = db.archive(LATE).await.unwrap();
assert_eq!(report.links_archived, 0);
assert_eq!(
rows_at(&db, "b", "c").await,
[
format!("alt {T1}"),
format!("main {T1}"),
format!("main {OPEN}")
],
"both closed rows stay hot: two lineages hold the key, so the \
closed-interval arm stands down for both. The trunk's open row is the \
third and is held for the other reason ([D-269]): `alt` forked after it \
and reads it still"
);
assert_eq!(reached_at(&db, None, T2).await, ["a", "b"]);
assert_eq!(reached_at(&db, Some("alt"), T2).await, ["a", "b"]);
assert_eq!(reached_at(&db, Some("alt"), EPOCH).await, ["a", "b", "c"]);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_key_one_lineage_holds_still_archives_its_closed_interval() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "c", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN),
)
.await
.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("b", "c", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(2.0)
.on_branch(alt.id.clone()),
)
.await
.unwrap();
h.advance(STEP);
db.retire_edge("a", "c", "LEADSTO", EPOCH, T1)
.await
.unwrap();
h.advance(STEP);
let report = db.archive(LATE).await.unwrap();
assert_eq!(report.links_archived, 2);
assert_eq!(
rows_at(&db, "a", "c").await,
[] as [String; 0],
"the trunk's closed `a -> c` is nobody's shadow and belongs in the cold file"
);
assert_eq!(
rows_at(&db, "b", "c").await.len(),
2,
"and the key two lineages do hold is untouched by this session"
);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_fork_keeps_believing_the_row_the_trunk_restated_over() {
let h = TestHarness::new();
let db = seed(&h).await;
let before = reached_at(&db, None, T2).await;
db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN),
)
.await
.unwrap();
h.advance(STEP);
assert_eq!(reached_at(&db, Some("alt"), T2).await, before);
let report = db.archive(LATE).await.unwrap();
assert_eq!(
reached_at(&db, Some("alt"), T2).await,
before,
"the branch believes what it believed; an archive is a move, not a \
retirement (report: {report:?})"
);
assert_eq!(
reached_at(&db, None, T2).await,
before,
"and the trunk, which was never at risk, is the control"
);
db.close().await.unwrap();
}
#[tokio::test]
async fn the_row_a_fork_pins_does_not_resurrect_the_retirement_over_it() {
let h = TestHarness::new();
let db = seed(&h).await;
db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.retire_edge("a", "b", "LEADSTO", EPOCH, T1)
.await
.unwrap();
h.advance(STEP);
let before = reached_at(&db, None, T2).await;
assert_eq!(before, ["a"], "the trunk retired its way out of the graph");
let report = db.archive(LATE).await.unwrap();
assert_eq!(
reached_at(&db, None, T2).await,
before,
"the retirement stands: archiving the closed row would let the open row the fork pinned win again (report: {report:?})"
);
assert_eq!(
reached_at(&db, Some("alt"), T2).await,
["a", "b", "c"],
"and `alt`, which forked before the retirement, never stopped believing it"
);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_lineage_still_supersedes_itself() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
for (weight, branch) in [(2.0, None), (3.0, Some(alt.id.clone()))] {
let mut e = EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(weight);
if let Some(b) = branch {
e = e.on_branch(b);
}
db.assert_edge(e).await.unwrap();
h.advance(STEP);
}
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(4.0),
)
.await
.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(5.0)
.on_branch(alt.id.clone()),
)
.await
.unwrap();
h.advance(STEP);
let (links_before, _) = counts(&db).await;
let report = db.archive(LATE).await.unwrap();
let (links_after, _) = counts(&db).await;
assert!(
report.links_archived >= 2,
"each lineage superseded itself once and both should have gone cold, got {}",
report.links_archived
);
assert_eq!(links_before - links_after, report.links_archived as i64);
assert_eq!(
reached(&db, None).await,
["a", "b", "c"],
"and current belief is untouched on both"
);
assert_eq!(reached(&db, Some("alt")).await, ["a", "b", "c"]);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_ledger_that_never_forked_archives_what_it_always_did() {
let h = TestHarness::new();
let db = seed(&h).await;
for weight in [2.0, 3.0] {
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(weight),
)
.await
.unwrap();
h.advance(STEP);
}
let report = db.archive(LATE).await.unwrap();
assert_eq!(
report.links_archived, 2,
"two superseded generations, and nothing about lineage in the way"
);
assert_eq!(reached(&db, None).await, ["a", "b", "c"]);
assert_eq!(audit_current(db.read_conn()).await.unwrap(), 0);
db.close().await.unwrap();
}
#[tokio::test]
async fn the_drift_audit_cannot_see_a_wrongly_pruned_ledger() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(2.0)
.on_branch(alt.id.clone()),
)
.await
.unwrap();
h.advance(STEP);
assert_eq!(audit_current(db.read_conn()).await.unwrap(), 0);
db.archive(LATE).await.unwrap();
assert_eq!(
audit_current(db.read_conn()).await.unwrap(),
0,
"zero before and zero after — the same answer the defect gave"
);
db.close().await.unwrap();
}
#[tokio::test]
async fn every_lineage_keeps_the_newest_log_entry_for_its_own_key() {
let h = TestHarness::new();
let db = seed(&h).await;
let alt = db.fork(id("alt"), BranchId::main()).await.unwrap();
h.advance(STEP);
db.assert_edge(
EdgeAssertion::new("a", "b", "LEADSTO")
.valid_from(EPOCH)
.valid_to(OPEN)
.weight(2.0)
.on_branch(alt.id.clone()),
)
.await
.unwrap();
h.advance(STEP);
db.archive(LATE).await.unwrap();
let mut rows = db
.read_conn()
.query(
"SELECT branch_id, COUNT(*) FROM transaction_log \
WHERE table_name = 'links' AND entity_id = ?1 GROUP BY branch_id ORDER BY branch_id",
libsql::params![format!("a|b|LEADSTO|{EPOCH}")],
)
.await
.unwrap();
let mut per_lineage = Vec::new();
while let Some(r) = rows.next().await.unwrap() {
per_lineage.push((r.get::<String>(0).unwrap(), r.get::<i64>(1).unwrap()));
}
assert_eq!(
per_lineage,
[("alt".to_string(), 1), ("main".to_string(), 1)],
"one surviving entry per lineage, which is one per fold partition"
);
db.close().await.unwrap();
}