use std::sync::Arc;
use topodb::*;
struct Fx {
db: Db,
_dir: tempfile::TempDir,
}
fn fx() -> Fx {
let dir = tempfile::tempdir().unwrap();
let db = Db::open(dir.path().join("t.redb")).unwrap();
Fx { db, _dir: dir }
}
fn create_node(id: NodeId, scope: Scope, label: &str) -> Op {
Op::CreateNode {
id,
scope,
label: label.into(),
props: Default::default(),
}
}
#[test]
fn duplicate_id_within_the_same_batch_is_rejected() {
let f = fx();
let scope = Scope::Shared;
let id = NodeId::new();
let other = NodeId::new();
let err =
f.db.submit(vec![
create_node(other, scope, "Entity"),
create_node(id, scope, "Memory"),
create_node(id, scope, "Memory"),
])
.unwrap_err();
assert!(
matches!(err, TopoError::Rejected(_)),
"expected Rejected, got {err:?}"
);
let scopes = ScopeSet::default().with_shared();
assert!(
f.db.node(&scopes, id).is_none(),
"the duplicated id must not exist"
);
assert!(
f.db.node(&scopes, other).is_none(),
"a rejected batch must leave storage untouched, including its OTHER ops"
);
}
#[test]
fn duplicate_id_against_an_earlier_committed_batch_is_rejected() {
let f = fx();
let scope = Scope::Shared;
let id = NodeId::new();
f.db.submit(vec![create_node(id, scope, "Memory")]).unwrap();
let err =
f.db.submit(vec![create_node(id, scope, "Entity")])
.unwrap_err();
assert!(
matches!(err, TopoError::Rejected(_)),
"expected Rejected, got {err:?}"
);
let scopes = ScopeSet::default().with_shared();
let rec =
f.db.node(&scopes, id)
.expect("the original node must still exist");
assert_eq!(
rec.label, "Memory",
"the original node's label must survive unchanged — no silent upsert"
);
}
#[test]
fn duplicate_id_raced_by_two_threads_commits_exactly_once() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Barrier;
let f = fx();
let db = Arc::new(f.db);
let scope = Scope::Shared;
const ROUNDS: usize = 64;
let mut a_ok = 0usize;
let mut b_ok = 0usize;
let rejected = AtomicUsize::new(0);
for _ in 0..ROUNDS {
let id = NodeId::new();
let barrier = Arc::new(Barrier::new(2));
let a_db = db.clone();
let a_barrier = barrier.clone();
let a = std::thread::spawn(move || {
a_barrier.wait();
a_db.submit(vec![create_node(id, scope, "WinnerA")])
});
let b_db = db.clone();
let b_barrier = barrier.clone();
let b = std::thread::spawn(move || {
b_barrier.wait();
b_db.submit(vec![create_node(id, scope, "WinnerB")])
});
let a_res = a.join().unwrap();
let b_res = b.join().unwrap();
let outcomes = [&a_res, &b_res];
let successes = outcomes.iter().filter(|r| r.is_ok()).count();
assert_eq!(
successes, 1,
"exactly one of the two racing same-id CreateNodes must commit \
per round — got {a_res:?} / {b_res:?}"
);
for r in outcomes {
match r {
Ok(_) => {}
Err(TopoError::Rejected(_)) => {
rejected.fetch_add(1, Ordering::Relaxed);
}
Err(other) => panic!("expected the loser to be Rejected, got {other:?}"),
}
}
if a_res.is_ok() {
a_ok += 1;
} else {
b_ok += 1;
}
}
assert_eq!(
a_ok + b_ok,
ROUNDS,
"every round must have exactly one winner"
);
assert_eq!(rejected.load(Ordering::Relaxed), ROUNDS);
let scopes = ScopeSet::default().with_shared();
let a_hits = db.nodes_by_label(&scopes, "WinnerA");
let b_hits = db.nodes_by_label(&scopes, "WinnerB");
assert_eq!(a_hits.len(), a_ok);
assert_eq!(b_hits.len(), b_ok);
}