#[path = "common/harness.rs"]
mod harness;
use harness::TestHarness;
use macrame::graph::EdgeAssertion;
use macrame::metrics::CommandKind;
use macrame::{ConceptUpsert, Database};
const T0: &str = "2026-01-01T00:00:00.000000Z";
const OPEN: &str = "9999-12-31T23:59:59.999999Z";
fn turns_for(snap: ¯ame::metrics::MetricsSnapshot, kind: CommandKind) -> u64 {
snap.kinds.iter().find(|k| k.kind == kind).unwrap().turns
}
#[tokio::test]
async fn every_write_method_is_attributed_to_its_own_command_kind() {
let harness = TestHarness::new();
let db = Database::open(&harness.db_path).await.unwrap();
db.upsert_concept(ConceptUpsert::new("a", "A").valid_from(T0))
.await
.unwrap();
db.upsert_concept(ConceptUpsert::new("b", "B").valid_from(T0))
.await
.unwrap();
db.assert_edge(
EdgeAssertion::new("a", "b", "KNOWS")
.valid_from(T0)
.valid_to(OPEN),
)
.await
.unwrap();
db.retire_edge("a", "b", "KNOWS", T0, "2026-06-01T00:00:00.000000Z")
.await
.unwrap();
db.rebuild_current().await.unwrap();
db.bulk_import(vec![EdgeAssertion::new("a", "b", "CITES")
.valid_from(T0)
.valid_to(OPEN)])
.await
.unwrap();
let snap = db.metrics();
assert_eq!(turns_for(&snap, CommandKind::UpsertConcept), 2);
assert_eq!(turns_for(&snap, CommandKind::AssertEdge), 1);
assert_eq!(turns_for(&snap, CommandKind::RetireEdge), 1);
assert_eq!(turns_for(&snap, CommandKind::RebuildCurrent), 1);
assert_eq!(turns_for(&snap, CommandKind::BulkImportChunk), 1);
assert_eq!(turns_for(&snap, CommandKind::Archive), 0);
let per_kind: u64 = snap.kinds.iter().map(|k| k.turns).sum();
assert_eq!(
per_kind, snap.turns,
"the loop counted {} turns but the kinds account for {per_kind}",
snap.turns
);
db.close().await.unwrap();
}
#[tokio::test]
async fn the_longest_hold_is_a_real_duration_and_names_its_command() {
let harness = TestHarness::new();
let db = Database::open(&harness.db_path).await.unwrap();
db.upsert_concept(ConceptUpsert::new("a", "A").valid_from(T0))
.await
.unwrap();
let leaves: Vec<_> = (0..500).map(|i| format!("c{i}")).collect();
db.write_concepts(
leaves
.iter()
.map(|id| ConceptUpsert::new(id, id).valid_from(T0))
.collect(),
)
.await
.unwrap();
let edges: Vec<_> = leaves
.iter()
.map(|id| {
EdgeAssertion::new("a", id, "KNOWS")
.valid_from(T0)
.valid_to(OPEN)
})
.collect();
db.bulk_import(edges).await.unwrap();
db.rebuild_current().await.unwrap();
let snap = db.metrics();
let (kind, held) = snap.longest.expect("some turn took at least a microsecond");
assert!(
held > std::time::Duration::ZERO,
"the timer is not running: longest hold is {held:?}"
);
let rebuild = snap
.kinds
.iter()
.find(|k| k.kind == CommandKind::RebuildCurrent)
.unwrap();
assert!(
held >= rebuild.mean,
"the high-water mark ({held:?}) is below a mean it should dominate \
({:?}) — {kind} was recorded as the longest",
rebuild.mean
);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_windowed_archive_takes_one_actor_turn_per_session() {
use std::time::Duration;
let harness = TestHarness::new();
let db = harness.db_with_fake_clock().await;
let ids: Vec<String> = (0..9).map(|i| format!("c{i:03}")).collect();
db.write_concepts(
ids.iter()
.map(|id| ConceptUpsert::new(id, "n").valid_from(T0))
.collect(),
)
.await
.unwrap();
for generation in 0..4 {
let batch: Vec<_> = (0..8)
.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();
harness.advance(Duration::from_secs(3_600));
}
let cutoff = harness.clock.peek();
let reports = db
.archive_windowed(&cutoff, Duration::from_secs(3_600))
.await
.unwrap();
assert!(reports.len() > 1, "the fixture produced one window");
let snap = db.metrics();
assert_eq!(
turns_for(&snap, CommandKind::Archive),
reports.len() as u64,
"{} sessions were reported but the actor spent a different number of \
turns on them — the loop is inside the actor, and windowing buys \
nothing",
reports.len()
);
db.close().await.unwrap();
}
#[tokio::test]
async fn a_backlog_shows_up_in_the_queue_depth() {
let harness = TestHarness::new();
let db = std::sync::Arc::new(Database::open(&harness.db_path).await.unwrap());
db.upsert_concept(ConceptUpsert::new("a", "A").valid_from(T0))
.await
.unwrap();
db.upsert_concept(ConceptUpsert::new("b", "B").valid_from(T0))
.await
.unwrap();
let mut tasks = Vec::new();
for i in 0..64 {
let db = std::sync::Arc::clone(&db);
tasks.push(tokio::spawn(async move {
db.assert_edge(
EdgeAssertion::new("a", "b", "KNOWS")
.valid_from(format!("2026-02-{:02}T00:00:00.000000Z", (i % 28) + 1))
.valid_to(format!("2026-03-{:02}T00:00:00.000000Z", (i % 28) + 1)),
)
.await
}));
}
for t in tasks {
let _ = t.await.unwrap();
}
let snap = db.metrics();
assert!(
snap.high_depth_max > 0,
"64 concurrent assertions produced no observed backlog at all — the \
depth is being sampled after the queue drains, not before the turn"
);
std::sync::Arc::into_inner(db)
.unwrap()
.close()
.await
.unwrap();
}