#[path = "common/harness.rs"]
mod harness;
use harness::TestHarness;
use macrame::graph::EdgeAssertion;
use macrame::integrity::audit_current;
use macrame::{ConceptUpsert, Database, DbError};
const T0: &str = "1970-01-01T00:00:00.000000Z";
const OPEN: &str = "9999-12-31T23:59:59.999999Z";
async fn seed(db: &Database, keys: usize, generations: usize) {
let ids: Vec<String> = (0..=keys).map(|i| format!("c{i:05}")).collect();
for chunk in ids.chunks(2_000) {
db.write_concepts(
chunk
.iter()
.map(|id| ConceptUpsert::new(id, "n").valid_from(T0))
.collect(),
)
.await
.unwrap();
}
for generation in 0..generations {
let batch: Vec<_> = (0..keys)
.map(|k| {
EdgeAssertion::new(&ids[k], &ids[k + 1], "LINKS")
.valid_from(T0)
.valid_to(OPEN)
.weight(1.0 + generation as f64)
})
.collect();
db.bulk_import(batch).await.unwrap();
}
}
async fn snapshot(db: &Database) -> Vec<(String, String, String, f64, String)> {
let mut out = Vec::new();
let mut rows = db
.read_conn()
.query(
"SELECT source_id, target_id, edge_type, weight, recorded_at \
FROM links_current ORDER BY source_id, target_id, edge_type, valid_from",
(),
)
.await
.unwrap();
while let Some(r) = rows.next().await.unwrap() {
out.push((
r.get(0).unwrap(),
r.get(1).unwrap(),
r.get(2).unwrap(),
r.get(3).unwrap(),
r.get(4).unwrap(),
));
}
out
}
#[tokio::test]
async fn the_chunked_rebuild_produces_what_the_atomic_one_does() {
let atomic = {
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
seed(&db, 900, 3).await;
db.rebuild_current().await.unwrap();
let s = snapshot(&db).await;
db.close().await.unwrap();
s
};
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
seed(&db, 900, 3).await;
let report = db.rebuild_current_chunked().await.unwrap();
assert_eq!(report.rows_rebuilt, 900);
let chunked = snapshot(&db).await;
assert_eq!(
chunked, atomic,
"the two repairs disagree about current belief"
);
assert_eq!(
audit_current(db.read_conn()).await.unwrap(),
0,
"the chunked rebuild left drift"
);
db.close().await.unwrap();
}
#[tokio::test]
async fn writing_after_a_swap_still_maintains_the_projection() {
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
seed(&db, 40, 2).await;
db.rebuild_current_chunked().await.unwrap();
let before = snapshot(&db).await.len();
db.upsert_concept(ConceptUpsert::new("fresh", "f").valid_from(T0))
.await
.unwrap();
db.assert_edge(
EdgeAssertion::new("c00000", "fresh", "CITES")
.valid_from(T0)
.valid_to(OPEN),
)
.await
.unwrap();
assert_eq!(
snapshot(&db).await.len(),
before + 1,
"trg_links_current_sync did not survive the swap"
);
db.assert_edge(
EdgeAssertion::new("c00000", "c00001", "LINKS")
.valid_from(T0)
.valid_to("2026-01-01T00:00:00.000000Z")
.weight(99.0),
)
.await
.unwrap();
assert_eq!(
snapshot(&db).await.len(),
before + 1,
"a re-assertion added a row: the swapped table has no primary key"
);
let err = db
.assert_edge(
EdgeAssertion::new("c00000", "fresh", "CITES")
.valid_from("2027-01-01T00:00:00.000000Z")
.valid_to(OPEN),
)
.await
.expect_err("trg_links_single_open did not survive the swap");
assert!(
matches!(err, DbError::SingleOpenViolation { .. }),
"got {err:?}"
);
assert_eq!(audit_current(db.read_conn()).await.unwrap(), 0);
db.close().await.unwrap();
}
#[tokio::test]
async fn the_indexes_come_back_with_their_declared_names() {
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
seed(&db, 60, 2).await;
db.rebuild_current_chunked().await.unwrap();
let conn = db.read_conn();
let mut names = Vec::new();
let mut rows = conn
.query(
"SELECT name FROM sqlite_master WHERE type = 'index' \
AND tbl_name = 'links_current' AND name NOT LIKE 'sqlite_%' \
ORDER BY name",
(),
)
.await
.unwrap();
while let Some(r) = rows.next().await.unwrap() {
names.push(r.get::<String>(0).unwrap());
}
assert_eq!(
names,
vec![
"idx_lc_open_interval".to_string(),
"idx_lc_tgt_active".to_string(),
"idx_lc_traversal_cover".to_string(),
],
"the swap left links_current indexed under different names than \
CREATE_INDICES declares, so the next migration would create a second \
copy of each"
);
let leftover: i64 = conn
.query(
"SELECT COUNT(*) FROM sqlite_master WHERE name LIKE '%shadow%'",
(),
)
.await
.unwrap()
.next()
.await
.unwrap()
.unwrap()
.get(0)
.unwrap();
assert_eq!(leftover, 0, "the shadow table outlived the swap");
db.close().await.unwrap();
}
#[tokio::test]
async fn a_write_during_the_build_survives_the_swap() {
let harness = TestHarness::new();
let db = std::sync::Arc::new(harness.db_with_fake_clock().await);
seed(&db, 900, 2).await;
let writer = {
let db = std::sync::Arc::clone(&db);
tokio::spawn(async move {
db.assert_edge(
EdgeAssertion::new("c00000", "c00001", "LINKS")
.valid_from(T0)
.valid_to("2030-01-01T00:00:00.000000Z")
.weight(1234.0),
)
.await
})
};
db.rebuild_current_chunked().await.unwrap();
writer.await.unwrap().unwrap();
assert_eq!(
audit_current(db.read_conn()).await.unwrap(),
0,
"a write during the shadow build left links_current diverged"
);
std::sync::Arc::into_inner(db)
.unwrap()
.close()
.await
.unwrap();
}
#[tokio::test]
async fn a_leftover_shadow_table_is_discarded_not_reused() {
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
seed(&db, 30, 2).await;
let raw = db.raw().connect().unwrap();
raw.execute(
"CREATE TABLE links_current_shadow (source_id TEXT, target_id TEXT)",
(),
)
.await
.unwrap();
raw.execute(
"INSERT INTO links_current_shadow VALUES ('junk', 'junk')",
(),
)
.await
.unwrap();
db.rebuild_current_chunked().await.unwrap();
assert_eq!(audit_current(db.read_conn()).await.unwrap(), 0);
assert_eq!(snapshot(&db).await.len(), 30);
db.close().await.unwrap();
}
#[tokio::test]
async fn an_archive_during_the_build_abandons_the_rebuild() {
use macrame::integrity::{ShadowOutcome, ShadowStep};
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
seed(&db, 30, 3).await;
let ShadowOutcome::Started { build_start, epoch } =
db.shadow_step(ShadowStep::Begin).await.unwrap()
else {
panic!("Begin returned the wrong outcome")
};
db.shadow_step(ShadowStep::Fill { after: None })
.await
.unwrap();
harness.advance(std::time::Duration::from_secs(3_600));
db.archive(&harness.clock.peek()).await.unwrap();
let err = db
.shadow_step(ShadowStep::Swap { build_start, epoch })
.await
.expect_err("an archive committed during the build");
assert!(
matches!(err, DbError::RebuildInterrupted { .. }),
"got {err:?}"
);
assert_eq!(audit_current(db.read_conn()).await.unwrap(), 0);
db.rebuild_current_chunked().await.unwrap();
assert_eq!(audit_current(db.read_conn()).await.unwrap(), 0);
db.close().await.unwrap();
}