use std::{collections::HashSet, sync::Arc};
use rocksdb::{Direction as ScanDir, IteratorMode, OptimisticTransactionDB, ReadOptions};
use std::collections::HashMap;
use super::{CF_EDGES_IN, CF_EDGES_OUT, CF_SCHEMA, CF_VERTEX_DEGREE, CF_VERTICES};
use crate::{
store::traits::GraphTransaction,
types::{
kv_codec::{
build_lazy_edge, build_lazy_vertex, decode_edge_key, decode_vertex_key, edge_scan_prefix, encode_edge_key,
encode_schema_key, encode_vertex_key, prefix_upper_bound, EdgeValue, VertexDegree, VertexValue,
EDGE_KEY_SIZE,
},
prop_codec::encode_props,
AdjacentEdgeCursor, AdjacentEdgesOptions, CanonicalEdgeKey, Direction, Edge, EdgeKey, LabelId, Primitive, Rank,
StoreError, Vertex, VertexKey,
},
};
type OwnedRocksTxn = rocksdb::Transaction<'static, OptimisticTransactionDB>;
type OwnedRocksTxnSnap = rocksdb::SnapshotWithThreadMode<'static, OwnedRocksTxn>;
fn begin_txn(db: &Arc<OptimisticTransactionDB>) -> OwnedRocksTxn {
let mut tx_opts = rocksdb::OptimisticTransactionOptions::default();
tx_opts.set_snapshot(true);
let txn = db.transaction_opt(&rocksdb::WriteOptions::default(), &tx_opts);
unsafe { std::mem::transmute(txn) }
}
pub struct Transaction {
db_txn_snap: Option<OwnedRocksTxnSnap>,
db_txn: Option<OwnedRocksTxn>,
db: Arc<OptimisticTransactionDB>,
}
impl Drop for Transaction {
fn drop(&mut self) {
self.db_txn_snap.take();
if let Some(txn) = self.db_txn.take() {
let _ = txn.rollback();
}
}
}
impl Transaction {
#[allow(clippy::missing_transmute_annotations)]
pub fn new(db: Arc<OptimisticTransactionDB>) -> Self {
let db_txn = begin_txn(&db);
let snap = db_txn.snapshot();
let db_txn_snap = unsafe { std::mem::transmute(snap) };
Self { db_txn_snap: Some(db_txn_snap), db_txn: Some(db_txn), db }
}
#[inline]
fn read_opts(&self) -> ReadOptions {
let mut opts = ReadOptions::default();
if let Some(ref snap) = self.db_txn_snap {
opts.set_snapshot(snap);
}
opts
}
}
impl GraphTransaction for Transaction {
fn get_vertex(&mut self, key: VertexKey) -> Result<Option<Vertex>, StoreError> {
let cf_vertices = self.db.cf_handle(CF_VERTICES).ok_or(StoreError::MissingColumnFamily("vertices"))?;
let vv_raw = self
.db_txn
.as_ref()
.expect("no active transaction")
.get_for_update_cf_opt(&cf_vertices, encode_vertex_key(key), true, &self.read_opts())
.map_err(StoreError::RocksDb)?;
match vv_raw {
Some(vv_bytes) => {
let vv = VertexValue::decode(&vv_bytes).ok_or(StoreError::CorruptData("vertex value"))?;
Ok(Some(build_lazy_vertex(key, &vv)))
}
_ => Ok(None),
}
}
fn get_vertex_degree(&mut self, key: VertexKey) -> Result<Option<(u32, u32, LabelId)>, StoreError> {
let cf_degree = self.db.cf_handle(CF_VERTEX_DEGREE).ok_or(StoreError::MissingColumnFamily("vertex_degree"))?;
let vd_raw = self
.db_txn
.as_ref()
.expect("no active transaction")
.get_for_update_cf_opt(&cf_degree, encode_vertex_key(key), true, &self.read_opts())
.map_err(StoreError::RocksDb)?;
match vd_raw {
Some(vd_bytes) => {
let vd = VertexDegree::decode(&vd_bytes).ok_or(StoreError::CorruptData("vertex degree"))?;
Ok(Some((vd.out_e_cnt, vd.in_e_cnt, vd.vertex_label_id)))
}
_ => Ok(None),
}
}
fn get_edge(&mut self, key: &EdgeKey) -> Result<Option<Edge>, StoreError> {
let cf_name = match key.direction {
Direction::OUT => CF_EDGES_OUT,
Direction::IN => CF_EDGES_IN,
};
let key_bytes = encode_edge_key(key);
let cf = self.db.cf_handle(cf_name).ok_or(StoreError::MissingColumnFamily(cf_name))?;
let raw_opt = self
.db_txn
.as_ref()
.expect("no active transaction")
.get_for_update_cf_opt(&cf, key_bytes, false, &self.read_opts())
.map_err(StoreError::RocksDb)?;
match raw_opt {
None => Ok(None),
Some(raw) => {
let ev = EdgeValue::decode(&raw).ok_or(StoreError::CorruptData("edge value"))?;
Ok(Some(build_lazy_edge(key, &ev)))
}
}
}
fn get_vertices(&mut self, keys: &[VertexKey]) -> Result<Vec<Vertex>, StoreError> {
let cf = self.db.cf_handle(CF_VERTICES).ok_or(StoreError::MissingColumnFamily("vertices"))?;
let txn = self.db_txn.as_ref().expect("no active transaction");
let mut out = Vec::with_capacity(keys.len());
for &k in keys {
let vv_raw = txn
.get_for_update_cf_opt(&cf, encode_vertex_key(k), true, &self.read_opts())
.map_err(StoreError::RocksDb)?;
if let Some(bytes) = vv_raw {
let vv = VertexValue::decode(&bytes).ok_or(StoreError::CorruptData("vertex value"))?;
out.push(build_lazy_vertex(k, &vv));
}
}
Ok(out)
}
fn get_edges(&mut self, keys: &[EdgeKey]) -> Result<Vec<Edge>, StoreError> {
let cf_out = self.db.cf_handle(CF_EDGES_OUT).ok_or(StoreError::MissingColumnFamily(CF_EDGES_OUT))?;
let cf_in = self.db.cf_handle(CF_EDGES_IN).ok_or(StoreError::MissingColumnFamily(CF_EDGES_IN))?;
let txn = self.db_txn.as_ref().expect("no active transaction");
let mut out = Vec::with_capacity(keys.len());
for key in keys {
let cf = match key.direction {
Direction::OUT => &cf_out,
Direction::IN => &cf_in,
};
let raw = txn
.get_for_update_cf_opt(cf, encode_edge_key(key), false, &self.read_opts())
.map_err(StoreError::RocksDb)?;
if let Some(bytes) = raw {
let ev = EdgeValue::decode(&bytes).ok_or(StoreError::CorruptData("edge value"))?;
out.push(build_lazy_edge(key, &ev));
}
}
Ok(out)
}
fn get_adjacent_edges(
&mut self,
vertex: VertexKey,
direction: Direction,
opts: AdjacentEdgesOptions<'_>,
limit: Option<u32>,
) -> Result<(Vec<Edge>, Option<AdjacentEdgeCursor>), StoreError> {
let cf_name = match direction {
Direction::OUT => CF_EDGES_OUT,
Direction::IN => CF_EDGES_IN,
};
let prefix = edge_scan_prefix(vertex, opts.label);
let mut read_opts = self.read_opts();
read_opts.set_prefix_same_as_start(true);
if let Some(upper) = prefix_upper_bound(&prefix) {
read_opts.set_iterate_upper_bound(upper.to_vec());
}
let seek_key = if let Some(cursor) = opts.start_from {
let mut key = Vec::with_capacity(EDGE_KEY_SIZE);
key.extend_from_slice(&encode_vertex_key(vertex));
key.extend_from_slice(&cursor.label_id.to_be_bytes());
key.extend_from_slice(&encode_vertex_key(cursor.secondary_id));
key.extend_from_slice(&cursor.rank.to_be_bytes());
key
} else {
prefix.clone().into_vec()
};
let dst_set: Option<HashSet<VertexKey>> = opts.dst.map(|k| k.iter().copied().collect());
let rank_set: Option<HashSet<Rank>> = opts.rank.map(|r| r.iter().copied().collect());
let cf = self.db.cf_handle(cf_name).ok_or(StoreError::MissingColumnFamily(cf_name))?;
let txn = self.db_txn.as_ref().expect("no active transaction");
let iter = txn.iterator_cf_opt(&cf, read_opts, IteratorMode::From(&seek_key, ScanDir::Forward));
let mut result = Vec::new();
let mut first = true;
for item in iter {
let (key_bytes, val_bytes) = item.map_err(StoreError::RocksDb)?;
if !key_bytes.starts_with(&prefix) {
break;
}
let ek = decode_edge_key(&key_bytes, direction).ok_or(StoreError::CorruptData("edge key"))?;
let current_cursor =
AdjacentEdgeCursor { label_id: ek.label_id, secondary_id: ek.secondary_id, rank: ek.rank };
if first && opts.start_from.is_some() {
first = false;
if Some(current_cursor) == opts.start_from {
continue;
}
}
if let Some(ref set) = dst_set {
if !set.contains(&ek.secondary_id) {
continue;
}
}
if let Some(ref set) = rank_set {
if !set.contains(&ek.rank) {
continue;
}
}
let ev = EdgeValue::decode(&val_bytes).ok_or(StoreError::CorruptData("edge value"))?;
result.push(build_lazy_edge(&ek, &ev));
if let Some(max) = limit {
if result.len() >= max as usize {
break;
}
}
}
let next_cursor = if let Some(last_edge) = result.last() {
if limit.map(|l| result.len() >= l as usize).unwrap_or(false) {
Some(AdjacentEdgeCursor::from_edge(last_edge, direction))
} else {
None
}
} else {
None
};
Ok((result, next_cursor))
}
fn scan_vertices(
&mut self,
label: Option<LabelId>,
start_from: Option<VertexKey>,
limit: u32,
) -> Result<(Vec<Vertex>, Option<VertexKey>), StoreError> {
let cf = self.db.cf_handle(CF_VERTICES).ok_or(StoreError::MissingColumnFamily("vertices"))?;
let mut read_opts = self.read_opts();
read_opts.set_total_order_seek(true);
let seek_key = if let Some(vk) = start_from { encode_vertex_key(vk).to_vec() } else { Vec::new() };
let txn = self.db_txn.as_ref().expect("no active transaction");
let iter = txn.iterator_cf_opt(&cf, read_opts, IteratorMode::From(&seek_key, ScanDir::Forward));
let mut result = Vec::new();
let mut first = true;
for item in iter {
let (key_bytes, val_bytes) = item.map_err(StoreError::RocksDb)?;
let key = decode_vertex_key(&key_bytes).ok_or(StoreError::CorruptData("vertex key"))?;
if first && start_from.is_some() {
first = false;
if Some(key) == start_from {
continue;
}
}
let vv = VertexValue::decode(&val_bytes).ok_or(StoreError::CorruptData("vertex value"))?;
if let Some(lbl) = label {
if vv.label_id != lbl {
continue;
}
}
result.push(build_lazy_vertex(key, &vv));
if result.len() >= limit as usize {
break;
}
}
let next_cursor = if result.len() >= limit as usize { result.last().map(|v| v.id) } else { None };
Ok((result, next_cursor))
}
fn scan_edges(
&mut self,
label: Option<LabelId>,
start_from: Option<CanonicalEdgeKey>,
limit: u32,
) -> Result<(Vec<Edge>, Option<CanonicalEdgeKey>), StoreError> {
let cf = self.db.cf_handle(CF_EDGES_OUT).ok_or(StoreError::MissingColumnFamily("edges_out"))?;
let mut read_opts = self.read_opts();
read_opts.set_total_order_seek(true);
let seek_key = if let Some(cek) = start_from { encode_edge_key(&cek.out_key()).to_vec() } else { Vec::new() };
let txn = self.db_txn.as_ref().expect("no active transaction");
let iter = txn.iterator_cf_opt(&cf, read_opts, IteratorMode::From(&seek_key, ScanDir::Forward));
let mut result = Vec::new();
let mut first = true;
for item in iter {
let (key_bytes, val_bytes) = item.map_err(StoreError::RocksDb)?;
let ek = decode_edge_key(&key_bytes, Direction::OUT).ok_or(StoreError::CorruptData("edge key"))?;
let current_cek = ek.canonical_edge_key();
if first && start_from.is_some() {
first = false;
if Some(current_cek) == start_from {
continue;
}
}
if let Some(lbl) = label {
if current_cek.label_id != lbl {
continue;
}
}
let ev = EdgeValue::decode(&val_bytes).ok_or(StoreError::CorruptData("edge value"))?;
result.push(build_lazy_edge(&ek, &ev));
if result.len() >= limit as usize {
break;
}
}
let next_cursor = if result.len() >= limit as usize { result.last().map(|e| e.canonical_key()) } else { None };
Ok((result, next_cursor))
}
fn put_vertex(
&mut self,
key: VertexKey,
label_id: LabelId,
props: &HashMap<u16, Primitive>,
) -> Result<(), StoreError> {
let txn = self.db_txn.as_ref().expect("no active transaction");
let cf_vertices = self.db.cf_handle(CF_VERTICES).ok_or(StoreError::MissingColumnFamily("vertices"))?;
let vv = VertexValue { label_id, property_blob: encode_props(props) };
txn.put_cf(&cf_vertices, encode_vertex_key(key), vv.encode()).map_err(StoreError::RocksDb)
}
fn put_vertex_degree(
&mut self,
key: VertexKey,
out_e_cnt: u32,
in_e_cnt: u32,
vertex_label_id: LabelId,
) -> Result<(), StoreError> {
let txn = self.db_txn.as_ref().expect("no active transaction");
let cf_degree = self.db.cf_handle(CF_VERTEX_DEGREE).ok_or(StoreError::MissingColumnFamily("vertex_degree"))?;
let vd = VertexDegree { vertex_label_id, out_e_cnt, in_e_cnt };
txn.put_cf(&cf_degree, encode_vertex_key(key), vd.encode()).map_err(StoreError::RocksDb)
}
fn put_edge(
&mut self,
key: &EdgeKey,
end_vertex_label: LabelId,
props: &HashMap<u16, Primitive>,
) -> Result<(), StoreError> {
let txn = self.db_txn.as_ref().expect("no active transaction");
let cf_name = match key.direction {
Direction::OUT => CF_EDGES_OUT,
Direction::IN => CF_EDGES_IN,
};
let key_bytes = encode_edge_key(key);
let cf = self.db.cf_handle(cf_name).ok_or(StoreError::MissingColumnFamily(cf_name))?;
let ev_bytes = EdgeValue { end_vertex_label, property_blob: encode_props(props) }.encode();
txn.put_cf(&cf, key_bytes, &ev_bytes).map_err(StoreError::RocksDb)
}
fn put_schema_entry(&mut self, kind: u8, name: &str, value: &[u8]) -> Result<(), StoreError> {
let txn = self.db_txn.as_ref().expect("no active transaction");
let cf_schema = self.db.cf_handle(CF_SCHEMA).ok_or(StoreError::MissingColumnFamily(CF_SCHEMA))?;
let key = encode_schema_key(kind, name);
txn.put_cf(&cf_schema, key, value).map_err(StoreError::RocksDb)
}
fn delete_vertex(&mut self, key: VertexKey) -> Result<(), StoreError> {
let cf_vertices = self.db.cf_handle(CF_VERTICES).ok_or(StoreError::MissingColumnFamily("vertices"))?;
let txn = self.db_txn.as_ref().expect("no active transaction");
txn.delete_cf(&cf_vertices, encode_vertex_key(key)).map_err(StoreError::RocksDb)
}
fn delete_vertex_degree(&mut self, key: VertexKey) -> Result<(), StoreError> {
let cf_degree = self.db.cf_handle(CF_VERTEX_DEGREE).ok_or(StoreError::MissingColumnFamily("vertex_degree"))?;
let txn = self.db_txn.as_ref().expect("no active transaction");
txn.delete_cf(&cf_degree, encode_vertex_key(key)).map_err(StoreError::RocksDb)
}
fn delete_edge(&mut self, key: &EdgeKey) -> Result<(), StoreError> {
let cf_name = match key.direction {
Direction::OUT => CF_EDGES_OUT,
Direction::IN => CF_EDGES_IN,
};
let key_bytes = encode_edge_key(key);
let cf = self.db.cf_handle(cf_name).ok_or(StoreError::MissingColumnFamily(cf_name))?;
let txn = self.db_txn.as_ref().expect("no active transaction");
txn.delete_cf(&cf, key_bytes).map_err(StoreError::RocksDb)
}
#[allow(clippy::missing_transmute_annotations)]
fn commit(&mut self) -> Result<(), StoreError> {
self.db_txn_snap.take();
let txn = self.db_txn.take().expect("no active transaction");
let result = txn.commit().map_err(|e| {
if e.to_string().contains("Resource busy") {
StoreError::Conflict
} else {
StoreError::RocksDb(e)
}
});
let new_txn = begin_txn(&self.db);
self.db_txn = Some(new_txn);
let snap = self.db_txn.as_ref().unwrap().snapshot();
self.db_txn_snap = Some(unsafe { std::mem::transmute(snap) });
result
}
#[allow(clippy::missing_transmute_annotations)]
fn abort(&mut self) {
self.db_txn_snap.take();
if let Some(txn) = self.db_txn.take() {
let _ = txn.rollback();
}
let new_txn = begin_txn(&self.db);
self.db_txn = Some(new_txn);
let snap = self.db_txn.as_ref().unwrap().snapshot();
self.db_txn_snap = Some(unsafe { std::mem::transmute(snap) });
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use rocksdb::{DBCommon, OptimisticTransactionDB, Options, SingleThreaded, DB};
use smol_str::SmolStr;
use crate::{
store::{
traits::{GraphStore, GraphTransaction},
RocksStorage,
},
types::{AdjacentEdgesOptions, Direction, Edge, EdgeKey, LabelId, Primitive, Vertex, VertexKey},
};
#[test]
fn test_read_write_conflict() {
let dir = tempfile::tempdir().unwrap();
let db: DBCommon<SingleThreaded, _> = OptimisticTransactionDB::open_default(dir.path()).unwrap();
let txn = db.transaction();
txn.put(b"Key_A", b"initial_A").unwrap();
txn.put(b"Key_B", b"initial_B").unwrap();
txn.commit().unwrap();
let snapshot = false;
let txn1 = db.transaction();
println!("[Time 0] txn1 started.");
let txn2 = db.transaction();
println!("[Time 1] txn2 started.");
let _ = txn1.get_for_update(b"Key_A", snapshot).unwrap();
let _ = txn1.get_for_update(b"Key_B", snapshot).unwrap();
txn1.put(b"Key_B", b"new_value_1").unwrap();
println!("[Time 2] txn1 executed GetForUpdate(A, B) and Put(B).");
let _ = txn2.get_for_update(b"Key_A", snapshot).unwrap();
txn2.put(b"Key_A", b"new_value_2").unwrap();
println!("[Time 3] txn2 executed GetForUpdate(A) and Put(A).");
println!("\n--- Entering Commit Phase (Scenario B) ---");
assert!(txn2.commit().is_ok(), "[Result] txn2 failed to commit! (Unexpected)");
print!("[Result] txn2 committed successfully. ");
assert!(txn1.commit().is_err(), "[Result] txn1 successfully committed! (Unexpected)");
print!("[Result] txn1 failed to commit as expected due to conflict.");
let _ = DB::destroy(&Options::default(), dir.path());
}
fn open_temp_store() -> (RocksStorage, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let store = RocksStorage::open(dir.path(), &Default::default()).unwrap();
(store, dir)
}
fn ctx(store: &RocksStorage) -> super::Transaction {
store.begin()
}
fn get_adjacent_edges_test(
txn: &mut super::Transaction,
vertex: VertexKey,
direction: Direction,
label: Option<LabelId>,
dst: Option<&[VertexKey]>,
limit: Option<u32>,
) -> Vec<Edge> {
txn.get_adjacent_edges(
vertex,
direction,
AdjacentEdgesOptions { label, dst, rank: None, start_from: None },
limit,
)
.unwrap()
.0
}
fn create_test_vertex(id: i64, label_id: LabelId) -> Vertex {
Vertex::new(id, label_id)
}
fn create_test_edge(src: i64, label: LabelId, dst: i64, _dir: Direction) -> Edge {
Edge::new(src, label, dst, 0, None, None)
}
#[test]
fn test_put_and_get_vertex() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let mut v_pos = create_test_vertex(1, 100);
txn.put_vertex(v_pos.id, v_pos.label_id, v_pos.props()).unwrap();
let fetched_v_pos = txn.get_vertex(v_pos.id).unwrap().unwrap();
assert_eq!(fetched_v_pos.id, v_pos.id);
assert_eq!(fetched_v_pos.label_id, v_pos.label_id);
let mut v_neg = create_test_vertex(-2, 200);
txn.put_vertex(v_neg.id, v_neg.label_id, v_neg.props()).unwrap();
let fetched_v_neg = txn.get_vertex(v_neg.id).unwrap().unwrap();
assert_eq!(fetched_v_neg.id, v_neg.id);
assert_eq!(fetched_v_neg.label_id, v_neg.label_id);
assert!(txn.get_vertex(999).unwrap().is_none());
}
#[test]
fn test_put_and_get_vertex_degree() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
txn.put_vertex_degree(1, 5, 10, 0).unwrap();
let (out_pos, in_pos, _label) = txn.get_vertex_degree(1).unwrap().unwrap();
assert_eq!(out_pos, 5);
assert_eq!(in_pos, 10);
txn.put_vertex_degree(-2, 15, 20, 0).unwrap();
let (out_neg, in_neg, _label) = txn.get_vertex_degree(-2).unwrap().unwrap();
assert_eq!(out_neg, 15);
assert_eq!(in_neg, 20);
assert!(txn.get_vertex_degree(999).unwrap().is_none());
}
#[test]
fn test_put_schema_entry() {
use super::super::CF_SCHEMA;
use crate::types::kv_codec::{encode_schema_key, SCHEMA_KIND_VERTEX_LABEL};
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
txn.put_schema_entry(SCHEMA_KIND_VERTEX_LABEL, "person", &7u16.to_be_bytes()).unwrap();
txn.commit().unwrap();
let cf = store.db.cf_handle(CF_SCHEMA).unwrap();
let key = encode_schema_key(SCHEMA_KIND_VERTEX_LABEL, "person");
let value = store.db.get_cf(&cf, key).unwrap().unwrap();
assert_eq!(value, 7u16.to_be_bytes());
}
#[test]
fn test_put_and_get_edge() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let ek_pos = EdgeKey::out_e(1, 100, 2, 0);
let mut e_pos = create_test_edge(1, 100, 2, Direction::OUT);
txn.put_edge(&ek_pos, 0, e_pos.props()).unwrap();
let fetched_e_pos = txn.get_edge(&ek_pos).unwrap().unwrap();
assert_eq!(fetched_e_pos.src_id, ek_pos.primary_id);
assert_eq!(fetched_e_pos.dst_id, ek_pos.secondary_id);
let ek_neg = EdgeKey::in_e(-3, 200, -4, 0);
let mut e_neg = create_test_edge(-3, 200, -4, Direction::IN);
txn.put_edge(&ek_neg, 0, e_neg.props()).unwrap();
let fetched_e_neg = txn.get_edge(&ek_neg).unwrap().unwrap();
assert_eq!(fetched_e_neg.src_id, ek_neg.secondary_id); assert_eq!(fetched_e_neg.dst_id, ek_neg.primary_id);
let non_existent_ek = EdgeKey::out_e(999, 1, 1000, 0);
assert!(txn.get_edge(&non_existent_ek).unwrap().is_none());
}
#[test]
fn test_get_edges() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
txn.put_vertex(1, 1, &HashMap::new()).unwrap();
txn.put_vertex(2, 1, &HashMap::new()).unwrap();
txn.put_vertex(3, 1, &HashMap::new()).unwrap();
txn.put_vertex(-1, 1, &HashMap::new()).unwrap();
txn.put_vertex(-2, 1, &HashMap::new()).unwrap();
txn.put_edge(&EdgeKey::out_e(1, 10, 2, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::out_e(1, 10, 3, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::out_e(1, 20, 2, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::in_e(1, 10, 2, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::in_e(1, 20, 2, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::out_e(-1, 30, -2, 0), 0, &HashMap::new()).unwrap();
let edges = get_adjacent_edges_test(&mut txn, 1, Direction::OUT, None, None, None);
assert_eq!(edges.len(), 3);
let edges_label_10 = get_adjacent_edges_test(&mut txn, 1, Direction::OUT, Some(10), None, None);
assert_eq!(edges_label_10.len(), 2);
assert!(edges_label_10.iter().all(|e| e.label_id == 10));
let edges_dst_2 = get_adjacent_edges_test(&mut txn, 1, Direction::OUT, None, Some(&[2]), None);
assert_eq!(edges_dst_2.len(), 2);
let edges_limit_1 = get_adjacent_edges_test(&mut txn, 1, Direction::OUT, None, None, Some(1));
assert_eq!(edges_limit_1.len(), 1);
let edges_neg = get_adjacent_edges_test(&mut txn, -1, Direction::OUT, None, None, None);
assert_eq!(edges_neg.len(), 1);
assert_eq!(edges_neg[0].src_id, -1);
assert_eq!(edges_neg[0].dst_id, -2);
let edges_in_2 = get_adjacent_edges_test(&mut txn, 2, Direction::IN, None, None, None);
assert_eq!(edges_in_2.len(), 2); }
#[test]
fn test_delete_vertex() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let mut v_pos = create_test_vertex(1, 100);
txn.put_vertex(v_pos.id, v_pos.label_id, v_pos.props()).unwrap();
txn.put_vertex_degree(v_pos.id, 0, 0, 0).unwrap();
assert!(txn.get_vertex(v_pos.id).unwrap().is_some());
txn.delete_vertex(v_pos.id).unwrap();
txn.delete_vertex_degree(v_pos.id).unwrap();
assert!(txn.get_vertex(v_pos.id).unwrap().is_none());
assert!(txn.get_vertex_degree(v_pos.id).unwrap().is_none());
let mut v_neg = create_test_vertex(-2, 200);
txn.put_vertex(v_neg.id, v_neg.label_id, v_neg.props()).unwrap();
txn.put_vertex_degree(v_neg.id, 0, 0, 0).unwrap();
assert!(txn.get_vertex(v_neg.id).unwrap().is_some());
txn.delete_vertex(v_neg.id).unwrap();
txn.delete_vertex_degree(v_neg.id).unwrap();
assert!(txn.get_vertex(v_neg.id).unwrap().is_none());
assert!(txn.get_vertex_degree(v_neg.id).unwrap().is_none());
}
#[test]
fn test_delete_edge() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let ek_pos = EdgeKey::out_e(1, 100, 2, 0);
let mut e_pos = create_test_edge(1, 100, 2, Direction::OUT);
txn.put_edge(&ek_pos, 0, e_pos.props()).unwrap();
assert!(txn.get_edge(&ek_pos).unwrap().is_some());
txn.delete_edge(&ek_pos).unwrap();
assert!(txn.get_edge(&ek_pos).unwrap().is_none());
let ek_neg = EdgeKey::in_e(-3, 200, -4, 0);
let mut e_neg = create_test_edge(-3, 200, -4, Direction::IN);
txn.put_edge(&ek_neg, 0, e_neg.props()).unwrap();
assert!(txn.get_edge(&ek_neg).unwrap().is_some());
txn.delete_edge(&ek_neg).unwrap();
assert!(txn.get_edge(&ek_neg).unwrap().is_none());
}
#[test]
fn test_commit_and_abort() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let mut v1 = create_test_vertex(1, 1);
txn.put_vertex(v1.id, v1.label_id, v1.props()).unwrap();
txn.commit().unwrap();
let mut new_txn = ctx(&store);
assert!(new_txn.get_vertex(v1.id).unwrap().is_some());
let mut v2 = create_test_vertex(2, 2);
txn.put_vertex(v2.id, v2.label_id, v2.props()).unwrap();
txn.abort();
let mut new_txn_after_abort = ctx(&store);
assert!(new_txn_after_abort.get_vertex(v2.id).unwrap().is_none());
assert!(new_txn_after_abort.get_vertex(v1.id).unwrap().is_some());
}
#[test]
fn test_put_and_get_vertex_with_properties() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let v1_id = 1;
let v1_label = 10;
let props: HashMap<u16, Primitive> =
[(1u16, Primitive::String(SmolStr::new("Alice"))), (2u16, Primitive::Int32(30))].into();
txn.put_vertex(v1_id, v1_label, &props).unwrap();
let mut fetched_v = txn.get_vertex(v1_id).unwrap().unwrap();
assert_eq!(fetched_v.id, v1_id);
assert_eq!(fetched_v.label_id, v1_label);
let fetched_props = fetched_v.props();
assert_eq!(fetched_props.len(), 2);
assert_eq!(fetched_props.get(&1u16), Some(&Primitive::String(SmolStr::new("Alice"))));
assert_eq!(fetched_props.get(&2u16), Some(&Primitive::Int32(30)));
}
#[test]
fn test_get_edges_in_direction_filters() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
txn.put_edge(&EdgeKey::in_e(1, 10, 5, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::in_e(2, 10, 5, 0), 0, &HashMap::new()).unwrap(); txn.put_edge(&EdgeKey::in_e(3, 20, 5, 0), 0, &HashMap::new()).unwrap();
let by_label = get_adjacent_edges_test(&mut txn, 5, Direction::IN, Some(10), None, None);
assert_eq!(by_label.len(), 2);
assert!(by_label.iter().all(|e| e.label_id == 10));
let limited = get_adjacent_edges_test(&mut txn, 5, Direction::IN, None, None, Some(2));
assert_eq!(limited.len(), 2);
let by_src = get_adjacent_edges_test(&mut txn, 5, Direction::IN, None, Some(&[2, 3]), None);
assert_eq!(by_src.len(), 2);
assert!(by_src.iter().all(|e| e.src_id == 2 || e.src_id == 3));
let combined = get_adjacent_edges_test(&mut txn, 5, Direction::IN, Some(10), Some(&[2]), None);
assert_eq!(combined.len(), 1);
assert_eq!(combined[0].src_id, 2);
assert_eq!(combined[0].label_id, 10);
}
#[test]
fn test_commit_returns_conflict() {
let (store, _dir) = open_temp_store();
let mut seed = ctx(&store);
seed.put_vertex(42, 1, &HashMap::new()).unwrap();
seed.commit().unwrap();
let mut txn1 = ctx(&store);
txn1.get_vertex(42).unwrap();
let mut txn2 = ctx(&store);
txn2.put_vertex(42, 2, &HashMap::new()).unwrap();
txn2.commit().unwrap();
txn1.put_vertex(42, 3, &HashMap::new()).unwrap();
assert!(matches!(txn1.commit(), Err(crate::types::StoreError::Conflict)));
}
#[test]
fn test_put_vertex_overwrite() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
txn.put_vertex(7, 1, &HashMap::new()).unwrap();
let first = txn.get_vertex(7).unwrap().unwrap();
assert_eq!(first.label_id, 1);
txn.put_vertex(7, 99, &HashMap::new()).unwrap();
let second = txn.get_vertex(7).unwrap().unwrap();
assert_eq!(second.label_id, 99);
}
#[test]
fn g25_put_empty_props_get_value_returns_none() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
txn.put_vertex(1, 5, &HashMap::new()).unwrap();
txn.commit().unwrap();
let mut txn2 = ctx(&store);
let mut v = txn2.get_vertex(1).unwrap().unwrap();
assert_eq!(v.get_value(10), None);
assert_eq!(v.label_id, 5);
assert!(v.props().is_empty());
}
#[test]
fn test_put_edge_overwrite() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let ek = EdgeKey::out_e(1, 10, 2, 0);
txn.put_edge(&ek, 0, &HashMap::new()).unwrap();
let mut first = txn.get_edge(&ek).unwrap().unwrap();
assert_eq!(first.props().len(), 0);
let props: HashMap<u16, Primitive> = [(1u16, Primitive::Int32(7))].into();
txn.put_edge(&ek, 0, &props).unwrap();
let mut second = txn.get_edge(&ek).unwrap().unwrap();
let second_props = second.props();
assert_eq!(second_props.len(), 1);
assert_eq!(second_props.get(&1u16), Some(&Primitive::Int32(7)));
}
#[test]
fn test_edges_with_nonzero_rank() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let ek0 = EdgeKey::out_e(1, 10, 2, 0);
let ek1 = EdgeKey::out_e(1, 10, 2, 1);
txn.put_edge(&ek0, 0, &HashMap::new()).unwrap();
txn.put_edge(&ek1, 0, &HashMap::new()).unwrap();
assert!(txn.get_edge(&ek0).unwrap().is_some());
assert!(txn.get_edge(&ek1).unwrap().is_some());
let edges = get_adjacent_edges_test(&mut txn, 1, Direction::OUT, Some(10), None, None);
assert_eq!(edges.len(), 2);
txn.delete_edge(&ek0).unwrap();
assert!(txn.get_edge(&ek0).unwrap().is_none());
assert!(txn.get_edge(&ek1).unwrap().is_some());
}
#[test]
fn test_put_and_get_edge_with_properties() {
let (store, _dir) = open_temp_store();
let mut txn = ctx(&store);
let ek = EdgeKey::out_e(1, 100, 2, 0);
let props: HashMap<u16, Primitive> = [(1u16, Primitive::Float64(0.5)), (2u16, Primitive::Int64(12345))].into();
txn.put_edge(&ek, 0, &props).unwrap();
let mut fetched_e = txn.get_edge(&ek).unwrap().unwrap();
assert_eq!(fetched_e.src_id, ek.primary_id);
assert_eq!(fetched_e.dst_id, ek.secondary_id);
assert_eq!(fetched_e.label_id, ek.label_id);
let fetched_props = fetched_e.props();
assert_eq!(fetched_props.len(), 2);
assert_eq!(fetched_props.get(&1u16), Some(&Primitive::Float64(0.5)));
assert_eq!(fetched_props.get(&2u16), Some(&Primitive::Int64(12345)));
}
#[test]
fn test_transaction_repeatable_reads_all_scenarios() {
let (store, _dir) = open_temp_store();
let mut seed = ctx(&store);
seed.put_vertex(1, 1, &HashMap::new()).unwrap();
seed.put_vertex(2, 1, &HashMap::new()).unwrap();
seed.put_vertex_degree(1, 1, 0, 0).unwrap();
let ek_seed = EdgeKey::out_e(1, 10, 2, 0);
seed.put_edge(&ek_seed, 0, &HashMap::new()).unwrap();
seed.commit().unwrap();
let mut txn1 = ctx(&store);
let mut txn2 = ctx(&store);
txn2.put_vertex(3, 100, &HashMap::new()).unwrap();
txn2.put_vertex(1, 99, &HashMap::new()).unwrap();
txn2.put_vertex_degree(1, 1, 1, 0).unwrap();
let ek_new = EdgeKey::out_e(1, 20, 3, 0);
txn2.put_edge(&ek_new, 0, &HashMap::new()).unwrap();
txn2.commit().unwrap();
let v1 = txn1.get_vertex(1).unwrap().unwrap();
assert_eq!(v1.label_id, 1); let v3_opt = txn1.get_vertex(3).unwrap();
assert!(v3_opt.is_none());
let batch_v = txn1.get_vertices(&[1, 3]).unwrap();
assert_eq!(batch_v.len(), 1);
assert_eq!(batch_v[0].id, 1);
assert_eq!(batch_v[0].label_id, 1);
let (deg_out, deg_in, _label) = txn1.get_vertex_degree(1).unwrap().unwrap();
assert_eq!(deg_out, 1);
assert_eq!(deg_in, 0);
let e_seed = txn1.get_edge(&ek_seed).unwrap();
assert!(e_seed.is_some());
let e_new = txn1.get_edge(&ek_new).unwrap();
assert!(e_new.is_none());
let batch_e = txn1.get_edges(&[ek_seed, ek_new]).unwrap();
assert_eq!(batch_e.len(), 1);
assert_eq!(batch_e[0].src_id, 1);
assert_eq!(batch_e[0].dst_id, 2);
let (adj_edges, _) = txn1
.get_adjacent_edges(
1,
Direction::OUT,
AdjacentEdgesOptions { label: None, dst: None, rank: None, start_from: None },
None,
)
.unwrap();
assert_eq!(adj_edges.len(), 1);
assert_eq!(adj_edges[0].dst_id, 2);
let (vertices_scan, _) = txn1.scan_vertices(None, None, 10).unwrap();
let vertex_ids: Vec<_> = vertices_scan.iter().map(|v| v.id).collect();
assert!(vertex_ids.contains(&1));
assert!(vertex_ids.contains(&2));
assert!(!vertex_ids.contains(&3));
let (edges_scan, _) = txn1.scan_edges(None, None, 10).unwrap();
assert_eq!(edges_scan.len(), 1);
assert_eq!(edges_scan[0].dst_id, 2); }
#[test]
fn test_scan_edges_paginates_across_src_id_prefixes() {
let (store, _dir) = open_temp_store();
let cf = store.db.cf_handle(super::CF_EDGES_OUT).unwrap();
for src in [1i64, 2, 3, 4, 5] {
let mut txn = ctx(&store);
txn.put_edge(&EdgeKey::out_e(src, 10, 100, 0), 0, &HashMap::new()).unwrap();
txn.commit().unwrap();
store.db.flush_cf(&cf).unwrap();
}
let mut txn = ctx(&store);
let mut seen = Vec::new();
let mut cursor = None;
loop {
let (page, next) = txn.scan_edges(None, cursor, 2).unwrap();
if page.is_empty() {
break;
}
seen.extend(page.iter().map(|e| e.src_id));
if next.is_none() {
break;
}
cursor = next;
}
seen.sort_unstable();
assert_eq!(seen, vec![1, 2, 3, 4, 5]);
}
}