use core_api::{BatchOp, FsyncPolicy, GraphDb, MutationEvent, SharedDb};
use core_storage::fs::{FileId, Fs, FsIntrospect};
use std::collections::HashMap;
use std::sync::{Arc, Barrier, Mutex};
use std::thread;
fn tmp(name: &str) -> std::path::PathBuf {
let d = std::env::temp_dir().join(format!(
"graphdb-gc-{}-{}-{}",
name,
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.subsec_nanos()
));
let _ = std::fs::remove_dir_all(&d);
d
}
#[derive(Default)]
struct CountingFs {
files: HashMap<FileId, Vec<u8>>,
syncs: usize,
}
impl Fs for CountingFs {
fn append(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()> {
self.files.entry(file).or_default().extend_from_slice(data);
Ok(())
}
fn sync(&mut self, _file: FileId) -> std::io::Result<()> {
self.syncs += 1;
Ok(())
}
fn read(&self, file: FileId) -> std::io::Result<Vec<u8>> {
Ok(self.files.get(&file).cloned().unwrap_or_default())
}
fn write_atomic(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()> {
self.files.insert(file, data.to_vec());
Ok(())
}
}
impl FsIntrospect for CountingFs {
fn total_appended(&self) -> usize {
0
}
fn sync_count(&self) -> usize {
self.syncs
}
}
fn counting_db() -> GraphDb<CountingFs> {
GraphDb::open_with(CountingFs::default()).unwrap()
}
#[test]
fn group_atomicity_each_submission_is_a_separate_wal_frame() {
use core_storage::wal::decode_all;
use sim_harness::SimFs;
let fs = SimFs::new();
let mut db = GraphDb::open_with(fs).unwrap();
let g = vec![
vec![BatchOp::InsertNode {
label: "A".into(),
key: "sub1".into(),
props: vec![],
}],
vec![BatchOp::InsertNode {
label: "A".into(),
key: "sub2".into(),
props: vec![],
}],
];
let (results, sync_err) = db.commit_group(g);
assert!(results.iter().all(|r| r.is_ok()), "both submissions ok");
assert!(sync_err.is_none(), "no sync error");
assert!(db.has_node("sub1"));
assert!(db.has_node("sub2"));
let fs = db.into_fs();
let wal = fs.read(FileId::Wal).unwrap();
let (records, _) = decode_all(&wal);
let batch_count = records
.iter()
.filter(|r| matches!(r, core_storage::wal::WalRecord::Batch(_)))
.count();
assert_eq!(
batch_count, 2,
"two submissions must produce two WAL Batch frames"
);
}
#[test]
fn one_fsync_per_group_strict_policy() {
let mut db = counting_db();
let groups: Vec<Vec<BatchOp>> = (0..8)
.map(|i| {
vec![BatchOp::InsertNode {
label: "A".into(),
key: format!("n{i}"),
props: vec![],
}]
})
.collect();
let (results, sync_err) = db.commit_group(groups);
assert!(results.iter().all(|r| r.is_ok()), "all submissions ok");
assert!(sync_err.is_none(), "sync succeeded");
assert_eq!(
db.fs_sync_count(),
1,
"group of 8 submissions must use exactly ONE fsync"
);
assert_eq!(db.node_count(), 8);
}
#[test]
fn two_groups_produce_two_fsyncs() {
let mut db = counting_db();
let g1 = vec![vec![BatchOp::InsertNode {
label: "A".into(),
key: "a".into(),
props: vec![],
}]];
let g2 = vec![vec![BatchOp::InsertNode {
label: "A".into(),
key: "b".into(),
props: vec![],
}]];
db.commit_group(g1);
db.commit_group(g2);
assert_eq!(
db.fs_sync_count(),
2,
"two separate commit_group calls = two fsyncs"
);
}
#[test]
fn relaxed_policy_group_skips_fsync() {
let mut db = counting_db();
db.set_fsync_policy(FsyncPolicy::Relaxed);
let groups: Vec<Vec<BatchOp>> = (0..4)
.map(|i| {
vec![BatchOp::InsertNode {
label: "A".into(),
key: format!("r{i}"),
props: vec![],
}]
})
.collect();
let (results, sync_err) = db.commit_group(groups);
assert!(results.iter().all(|r| r.is_ok()));
assert!(sync_err.is_none());
assert_eq!(
db.fs_sync_count(),
0,
"Relaxed policy must skip all fsyncs even in a group"
);
}
#[test]
fn fifo_ordering_within_single_caller() {
let dir = tmp("fifo");
let db = SharedDb::open(&dir).unwrap();
const N: usize = 20;
let mut prev_nodes = 0usize;
for i in 0..N {
let ops = vec![BatchOp::InsertNode {
label: "A".into(),
key: format!("seq{i}"),
props: vec![],
}];
db.submit_batch(ops).unwrap();
let n = db.read().node_count();
assert!(
n >= prev_nodes,
"node count must not decrease: was {prev_nodes}, now {n}"
);
prev_nodes = n;
}
assert_eq!(db.read().node_count(), N);
}
#[test]
fn concurrent_submitters_all_commit() {
let dir = tmp("conc");
let db = SharedDb::open(&dir).unwrap();
const WRITERS: usize = 8;
const OPS_PER_WRITER: usize = 25;
let start = Arc::new(Barrier::new(WRITERS));
let handles: Vec<_> = (0..WRITERS)
.map(|w| {
let db = db.clone();
let start = Arc::clone(&start);
thread::spawn(move || {
start.wait(); for i in 0..OPS_PER_WRITER {
let key = format!("w{w}_n{i}");
db.submit_batch(vec![BatchOp::InsertNode {
label: "N".into(),
key,
props: vec![],
}])
.expect("submit_batch must succeed");
}
})
})
.collect();
for h in handles {
h.join().expect("writer thread panicked");
}
let expected = WRITERS * OPS_PER_WRITER;
let actual = db.read().node_count();
assert_eq!(
actual, expected,
"all {expected} concurrent submissions must commit"
);
}
#[test]
fn crash_before_group_fsync_loses_unsynced_group() {
use core_storage::fs::FileId;
use sim_harness::SimFs;
let fs = SimFs::new();
let mut db = GraphDb::open_with(fs).unwrap();
let g1 = vec![vec![BatchOp::InsertNode {
label: "A".into(),
key: "g1".into(),
props: vec![],
}]];
db.commit_group(g1);
let after_g1_bytes = db.fs_total_appended();
let g2 = vec![
vec![BatchOp::InsertNode {
label: "A".into(),
key: "g2a".into(),
props: vec![],
}],
vec![BatchOp::InsertNode {
label: "A".into(),
key: "g2b".into(),
props: vec![],
}],
];
db.commit_group_nosync(g2);
let fs = db.into_fs();
let wal = fs.read(FileId::Wal).unwrap();
assert!(
wal.len() > after_g1_bytes,
"WAL must contain g2 bytes before crash"
);
let mut survivor = SimFs::new();
let snap = fs.read(FileId::Snapshot).unwrap();
if !snap.is_empty() {
survivor.write_atomic(FileId::Snapshot, &snap).unwrap();
}
survivor
.write_atomic(FileId::Wal, &wal[..after_g1_bytes])
.unwrap();
let db2 = GraphDb::open_with(survivor).unwrap();
assert!(db2.has_node("g1"), "g1 (synced group) must survive");
assert!(
!db2.has_node("g2a"),
"g2a (unsynced group) must be lost on crash"
);
assert!(
!db2.has_node("g2b"),
"g2b (unsynced group) must be lost on crash"
);
}
#[test]
fn intra_group_prefix_survives_crash() {
use sim_harness::SimFs;
let probe_fs = SimFs::new();
let mut probe = GraphDb::open_with(probe_fs).unwrap();
probe
.commit_group_nosync(vec![vec![BatchOp::InsertNode {
label: "A".into(),
key: "s1".into(),
props: vec![],
}]])
.into_iter()
.for_each(|r| {
r.unwrap();
});
let frame1_bytes = probe.fs_total_appended(); drop(probe);
let crash_at = frame1_bytes + 3;
let fs = SimFs::with_crash_after(crash_at);
let mut db = GraphDb::open_with(fs).unwrap();
let results = db.commit_group_nosync(vec![
vec![BatchOp::InsertNode {
label: "A".into(),
key: "s1".into(),
props: vec![],
}],
vec![BatchOp::InsertNode {
label: "A".into(),
key: "s2".into(),
props: vec![],
}],
]);
let _ = results;
let fs = db.into_fs();
let survivor = fs.surviving_state();
let db2 = GraphDb::open_with(survivor).unwrap();
assert!(
db2.has_node("s1"),
"first submission frame (before crash point) must survive"
);
assert!(
!db2.has_node("s2"),
"torn second-frame submission must be dropped on recovery"
);
assert_eq!(
db2.node_count(),
1,
"only the complete first frame survives"
);
}
#[test]
fn direct_apis_unchanged_alongside_queue() {
let dir = tmp("direct");
let db = SharedDb::open(&dir).unwrap();
db.write()
.insert_node("N", "direct1", vec![])
.expect("direct insert_node must work");
db.write()
.insert_node("N", "direct2", vec![])
.expect("direct insert_node must work");
db.submit_batch(vec![BatchOp::InsertNode {
label: "N".into(),
key: "queued1".into(),
props: vec![],
}])
.expect("submit_batch must work");
db.write()
.write_batch(|b| {
b.insert_node("N", "batch1", vec![]);
b.insert_node("N", "batch2", vec![]);
})
.expect("write_batch must work");
let n = db.read().node_count();
assert_eq!(n, 5, "direct + queued + batch all committed");
}
#[test]
fn deferred_events_fire_after_flush() {
use sim_harness::SimFs;
let received: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
let received2 = Arc::clone(&received);
let fs = SimFs::new();
let mut db = GraphDb::open_with(fs).unwrap();
db.set_event_sink(Box::new(move |ev| {
if let MutationEvent::NodeInserted { key, .. } = ev {
received2.lock().unwrap().push(key);
}
}));
db.set_deferred_events_mode(true);
db.commit_group_nosync(vec![
vec![BatchOp::InsertNode {
label: "A".into(),
key: "ev1".into(),
props: vec![],
}],
vec![BatchOp::InsertNode {
label: "A".into(),
key: "ev2".into(),
props: vec![],
}],
]);
assert!(
received.lock().unwrap().is_empty(),
"events must not fire before flush"
);
db.flush_deferred_events();
db.set_deferred_events_mode(false);
let keys = received.lock().unwrap().clone();
assert!(keys.contains(&"ev1".to_string()), "ev1 must be delivered");
assert!(keys.contains(&"ev2".to_string()), "ev2 must be delivered");
}
#[test]
fn deferred_events_discarded_on_failure() {
use sim_harness::SimFs;
let received: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
let received2 = Arc::clone(&received);
let fs = SimFs::new();
let mut db = GraphDb::open_with(fs).unwrap();
db.set_event_sink(Box::new(move |ev| {
if let MutationEvent::NodeInserted { key, .. } = ev {
received2.lock().unwrap().push(key);
}
}));
db.set_deferred_events_mode(true);
db.commit_group_nosync(vec![vec![BatchOp::InsertNode {
label: "A".into(),
key: "lost".into(),
props: vec![],
}]]);
assert!(
received.lock().unwrap().is_empty(),
"events must not fire before discard"
);
db.discard_deferred_events();
db.set_deferred_events_mode(false);
assert!(
received.lock().unwrap().is_empty(),
"discarded events must never be delivered to subscribers"
);
}
#[test]
#[ignore]
fn group_commit_throughput_bench() {
use std::time::Instant;
const WRITERS: usize = 8;
const OPS_PER_WRITER: usize = 200;
const TOTAL_OPS: usize = WRITERS * OPS_PER_WRITER;
let dir_serial = tmp("bench-serial");
let db_serial = SharedDb::open(&dir_serial).unwrap();
db_serial.write().insert_node("W", "warm", vec![]).unwrap();
let t0 = Instant::now();
for i in 0..TOTAL_OPS {
db_serial
.write()
.insert_node("W", &format!("s{i}"), vec![])
.unwrap();
}
let serial_elapsed = t0.elapsed();
let serial_ops_per_s = TOTAL_OPS as f64 / serial_elapsed.as_secs_f64();
let dir_conc = tmp("bench-conc");
let db_conc = SharedDb::open(&dir_conc).unwrap();
{
let db = db_conc.clone();
db.submit_batch(vec![BatchOp::InsertNode {
label: "W".into(),
key: "warm".into(),
props: vec![],
}])
.unwrap();
}
let start = Arc::new(Barrier::new(WRITERS));
let t1 = Instant::now();
let handles: Vec<_> = (0..WRITERS)
.map(|w| {
let db = db_conc.clone();
let start = Arc::clone(&start);
thread::spawn(move || {
start.wait();
for i in 0..OPS_PER_WRITER {
db.submit_batch(vec![BatchOp::InsertNode {
label: "W".into(),
key: format!("w{w}n{i}"),
props: vec![],
}])
.unwrap();
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let conc_elapsed = t1.elapsed();
let conc_ops_per_s = TOTAL_OPS as f64 / conc_elapsed.as_secs_f64();
let ratio = conc_ops_per_s / serial_ops_per_s;
let dir_reader = tmp("bench-reader");
let db_reader = SharedDb::open(&dir_reader).unwrap();
for i in 0..10 {
db_reader
.write()
.insert_node("R", &format!("pre{i}"), vec![])
.unwrap();
}
let start2 = Arc::new(Barrier::new(WRITERS + 1));
let read_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
let read_latencies: Arc<Mutex<Vec<u128>>> = Arc::new(Mutex::new(Vec::new()));
let writer_handles: Vec<_> = (0..WRITERS)
.map(|w| {
let db = db_reader.clone();
let start2 = Arc::clone(&start2);
let read_done = Arc::clone(&read_done);
thread::spawn(move || {
start2.wait();
let mut i = 0usize;
while !read_done.load(std::sync::atomic::Ordering::Relaxed) {
let _ = db.submit_batch(vec![BatchOp::InsertNode {
label: "W".into(),
key: format!("bw{w}_{i}"),
props: vec![],
}]);
i += 1;
}
})
})
.collect();
let lat_db = db_reader.clone();
let lat_lats = Arc::clone(&read_latencies);
let lat_start = Arc::clone(&start2);
let lat_read_done = Arc::clone(&read_done);
let reader_handle = thread::spawn(move || {
lat_start.wait();
let deadline = Instant::now() + std::time::Duration::from_millis(500);
while Instant::now() < deadline {
let t = Instant::now();
let _ = lat_db.reader().query(
"MATCH (n:R) RETURN n.id",
&std::collections::BTreeMap::new(),
);
let elapsed_us = t.elapsed().as_micros();
lat_lats.lock().unwrap().push(elapsed_us);
}
lat_read_done.store(true, std::sync::atomic::Ordering::Relaxed);
});
reader_handle.join().unwrap();
for h in writer_handles {
let _ = h.join();
}
let mut lats = read_latencies.lock().unwrap().clone();
let reader_p95_us = if lats.is_empty() {
0u128
} else {
lats.sort_unstable();
lats[lats.len() * 95 / 100]
};
println!(
"{}",
serde_json::json!({
"serialized_writer_ops_per_s": serial_ops_per_s as u64,
"eight_writer_ops_per_s": conc_ops_per_s as u64,
"ratio": format!("{ratio:.2}"),
"reader_under_burst_p95_us": reader_p95_us,
"gate_pass": ratio >= 3.0,
})
);
}
#[test]
#[ignore]
fn group_commit_simfs_amortization_bench() {
use sim_harness::SimFs;
use std::time::Instant;
const FSYNC_DELAY_US: u64 = 5_000; const WRITERS: usize = 8;
const OPS_PER_WRITER: usize = 50; const TOTAL_OPS: usize = WRITERS * OPS_PER_WRITER;
let fs_serial = SimFs::with_sync_delay_us(FSYNC_DELAY_US);
let mut db_serial = GraphDb::open_with(fs_serial).unwrap();
let t0 = Instant::now();
for i in 0..TOTAL_OPS {
db_serial
.insert_node("W", &format!("s{i}"), vec![])
.unwrap();
}
let serial_elapsed = t0.elapsed();
let serial_ops_per_s = TOTAL_OPS as f64 / serial_elapsed.as_secs_f64();
let fs_group = SimFs::with_sync_delay_us(FSYNC_DELAY_US);
let mut db_group = GraphDb::open_with(fs_group).unwrap();
let t1 = Instant::now();
for g in 0..(TOTAL_OPS / WRITERS) {
let batches: Vec<Vec<BatchOp>> = (0..WRITERS)
.map(|i| {
vec![BatchOp::InsertNode {
label: "W".into(),
key: format!("g{g}n{i}"),
props: vec![],
}]
})
.collect();
let (results, sync_err) = db_group.commit_group(batches);
assert!(sync_err.is_none(), "simfs sync must not fail");
for r in results {
r.unwrap();
}
}
let group_elapsed = t1.elapsed();
let group_ops_per_s = TOTAL_OPS as f64 / group_elapsed.as_secs_f64();
let ratio = group_ops_per_s / serial_ops_per_s;
println!(
"{}",
serde_json::json!({
"bench": "simfs_amortization",
"fsync_delay_us": FSYNC_DELAY_US,
"writers": WRITERS,
"total_ops": TOTAL_OPS,
"serial_ops_per_s": serial_ops_per_s as u64,
"group_ops_per_s": group_ops_per_s as u64,
"ratio": format!("{ratio:.2}"),
"gate_pass": ratio >= 3.0,
})
);
assert!(
ratio >= 3.0,
"group-commit ({group_ops_per_s:.0} ops/s) must be >= 3x serial \
({serial_ops_per_s:.0} ops/s) under {FSYNC_DELAY_US}µs injected fsync \
latency; ratio = {ratio:.2}"
);
}
#[test]
fn shared_db_fsync_failure_degrades_and_truncates_wal() {
use std::sync::atomic::{AtomicBool, Ordering};
let dir = tmp("f1c-single");
let fail = Arc::new(AtomicBool::new(false));
let fail2 = Arc::clone(&fail);
let db = SharedDb::open_with_test_sync(&dir, move |path| {
if fail2.load(Ordering::Acquire) {
Err(std::io::Error::other("injected fsync failure"))
} else {
core_storage::sync_wal_at(path)
}
})
.unwrap();
db.submit_batch(vec![BatchOp::InsertNode {
label: "A".into(),
key: "pre".into(),
props: vec![],
}])
.unwrap();
let pre_group_wal_len = std::fs::metadata(dir.join("wal.bin")).unwrap().len();
fail.store(true, Ordering::Release);
let result = db.submit_batch(vec![BatchOp::InsertNode {
label: "A".into(),
key: "fail-group".into(),
props: vec![],
}]);
assert!(result.is_err(), "group with failing fsync must return Err");
let post_wal_len = std::fs::metadata(dir.join("wal.bin")).unwrap().len();
assert_eq!(
post_wal_len, pre_group_wal_len,
"WAL must be truncated to pre-group length after fsync failure"
);
let result2 = db.submit_batch(vec![BatchOp::InsertNode {
label: "A".into(),
key: "post-fail".into(),
props: vec![],
}]);
assert!(
result2.is_err(),
"subsequent submit_batch must return Err after degradation"
);
drop(db);
let db2 = SharedDb::open(&dir).unwrap();
assert!(
db2.read().has_node("pre"),
"pre-group node must survive replay"
);
assert!(
!db2.read().has_node("fail-group"),
"failed-group node must be absent on replay"
);
}
#[test]
fn direct_write_before_group_survives_group_fsync_failure() {
use std::sync::atomic::{AtomicBool, Ordering};
let dir = tmp("f1c-concurrent");
let fail = Arc::new(AtomicBool::new(false));
let fail2 = Arc::clone(&fail);
let db = SharedDb::open_with_test_sync(&dir, move |path| {
if fail2.load(Ordering::Acquire) {
Err(std::io::Error::other("injected fsync failure"))
} else {
core_storage::sync_wal_at(path)
}
})
.unwrap();
db.write().insert_node("A", "direct-ok", vec![]).unwrap();
let pre_group_wal_len = std::fs::metadata(dir.join("wal.bin")).unwrap().len();
fail.store(true, Ordering::Release);
let result = db.submit_batch(vec![BatchOp::InsertNode {
label: "A".into(),
key: "group-fail".into(),
props: vec![],
}]);
assert!(result.is_err(), "group fsync failure must return Err");
let post_wal_len = std::fs::metadata(dir.join("wal.bin")).unwrap().len();
assert_eq!(
post_wal_len, pre_group_wal_len,
"WAL truncated to pre-group boundary; direct write frames are preserved"
);
drop(db);
let db2 = SharedDb::open(&dir).unwrap();
assert!(
db2.read().has_node("direct-ok"),
"directly-acknowledged write must survive replay"
);
assert!(
!db2.read().has_node("group-fail"),
"failed group node must be absent on replay"
);
}
#[test]
fn write_batch_strict_always_fsyncs() {
let mut db = counting_db();
assert_eq!(db.fsync_policy(), FsyncPolicy::Strict);
for i in 0..5usize {
db.write_batch(|b| {
b.insert_node("X", &format!("single{i}"), vec![]);
})
.expect("write_batch must succeed");
}
assert_eq!(
db.fs_sync_count(),
5,
"5 single-op write_batch calls under Strict must produce exactly 5 fsyncs"
);
assert_eq!(db.node_count(), 5);
db.write_batch(|b| {
for j in 0..5usize {
b.insert_node("X", &format!("multi{j}"), vec![]);
}
})
.expect("5-op write_batch must succeed");
assert_eq!(
db.fs_sync_count(),
6,
"one 5-op write_batch under Strict must produce exactly 1 additional fsync (total 6)"
);
assert_eq!(db.node_count(), 10);
}
#[test]
fn concurrent_overlapping_key_inserts_land_exactly_once() {
let dir = tmp("conc-overlap");
let db = SharedDb::open(&dir).unwrap();
const WRITERS: usize = 16;
let start = Arc::new(Barrier::new(WRITERS));
let ok_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let handles: Vec<_> = (0..WRITERS)
.map(|_| {
let db = db.clone();
let start = Arc::clone(&start);
let ok_count = Arc::clone(&ok_count);
thread::spawn(move || {
start.wait();
let r = db.submit_batch(vec![BatchOp::InsertNode {
label: "N".into(),
key: "shared".into(),
props: vec![],
}]);
if r.is_ok() {
ok_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
})
})
.collect();
for h in handles {
h.join().expect("writer thread panicked");
}
assert_eq!(
ok_count.load(std::sync::atomic::Ordering::Relaxed),
1,
"exactly one insert of the shared key may succeed"
);
assert!(db.read().has_node("shared"), "the node must exist");
assert_eq!(db.read().node_count(), 1, "exactly one node total");
}
#[test]
fn concurrent_writes_keep_rules_and_index_consistent() {
use core_api::{Predicate, RuleDef};
use std::collections::BTreeMap;
let dir = tmp("conc-rules-index");
let db = SharedDb::open(&dir).unwrap();
db.write()
.create_rule(RuleDef {
name: "same_city".into(),
src_label: "Talent".into(),
dst_label: "Company".into(),
predicate: Predicate::FieldEqual {
field: "city".into(),
},
edge_type: "IN_CITY".into(),
weight_prop: None,
max_edges: None,
approximate: false,
via_label: None,
via_edge: None,
via_dir: None,
namespace: None,
})
.unwrap();
db.write().enable_index("Talent", "city").unwrap();
const N_TALENT: usize = 30;
const N_COMPANY: usize = 10;
let start = Arc::new(Barrier::new(2));
let db_t = db.clone();
let start_t = Arc::clone(&start);
let t_thread = thread::spawn(move || {
start_t.wait();
for i in 0..N_TALENT {
db_t.submit_batch(vec![BatchOp::InsertNode {
label: "Talent".into(),
key: format!("t{i}"),
props: vec![("city".into(), core_api::Value::Str("austin".into()))],
}])
.unwrap();
}
});
let db_c = db.clone();
let start_c = Arc::clone(&start);
let c_thread = thread::spawn(move || {
start_c.wait();
for i in 0..N_COMPANY {
db_c.submit_batch(vec![BatchOp::InsertNode {
label: "Company".into(),
key: format!("c{i}"),
props: vec![("city".into(), core_api::Value::Str("austin".into()))],
}])
.unwrap();
}
});
t_thread.join().unwrap();
c_thread.join().unwrap();
let edges = db
.read()
.query(
"MATCH (t:Talent)-[r:IN_CITY]->(c:Company) RETURN t",
&BTreeMap::new(),
)
.unwrap();
assert_eq!(
edges.len(),
N_TALENT * N_COMPANY,
"all derived edges must be present after concurrent rule-fires"
);
let indexed = db
.read()
.query(
"MATCH (t:Talent {city: 'austin'}) RETURN t",
&BTreeMap::new(),
)
.unwrap();
assert_eq!(
indexed.len(),
N_TALENT,
"the equality index must be consistent under concurrent writes"
);
assert_eq!(db.read().node_count(), N_TALENT + N_COMPANY);
}