pub mod checkpoint;
pub mod mmap;
pub mod wal;
pub use checkpoint::{Checkpoint, CheckpointConfig, CheckpointManager, CheckpointStats};
#[cfg(feature = "wal")]
pub use wal::{SyncMode, WalConfig, WalEntry, WalOperation, WriteAheadLog};
use crate::error::{GraphError, Result};
use crate::graph::{Id, Node, Relationship};
use crate::index::IndexManager;
use serde_json::Value;
use std::collections::HashMap;
use std::path::Path;
#[cfg(feature = "wal")]
use parking_lot::Mutex;
#[cfg(feature = "wal")]
use std::path::PathBuf;
#[cfg(feature = "wal")]
use std::sync::Arc;
pub enum StorageBackend {
InMemory(InMemoryStorage),
MmapFile(mmap::MmapStorage),
}
pub struct Storage {
backend: StorageBackend,
index_manager: IndexManager,
#[cfg(feature = "wal")]
wal: Option<Arc<Mutex<WriteAheadLog>>>,
#[cfg(feature = "wal")]
#[allow(dead_code)] storage_path: Option<PathBuf>,
#[cfg(feature = "wal")]
checkpoint_manager: Option<CheckpointManager>,
}
pub struct InMemoryStorage {
nodes: HashMap<Id, Node>,
relationships: HashMap<Id, Relationship>,
node_relationships: HashMap<Id, Vec<Id>>, }
impl Storage {
pub fn new() -> Result<Self> {
Ok(Storage {
backend: StorageBackend::InMemory(InMemoryStorage {
nodes: HashMap::new(),
relationships: HashMap::new(),
node_relationships: HashMap::new(),
}),
index_manager: IndexManager::new(),
#[cfg(feature = "wal")]
wal: None,
#[cfg(feature = "wal")]
storage_path: None,
#[cfg(feature = "wal")]
checkpoint_manager: None,
})
}
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
let path_ref = path.as_ref();
let mmap_storage = mmap::MmapStorage::open(path_ref)?;
#[cfg(feature = "wal")]
let wal_path = path_ref.with_extension("wal");
let mut storage = Storage {
backend: StorageBackend::MmapFile(mmap_storage),
index_manager: IndexManager::new(),
#[cfg(feature = "wal")]
wal: None,
#[cfg(feature = "wal")]
storage_path: Some(path_ref.to_path_buf()),
#[cfg(feature = "wal")]
checkpoint_manager: None,
};
#[cfg(feature = "wal")]
{
let wal = WriteAheadLog::open(&wal_path)?;
storage.wal = Some(Arc::new(Mutex::new(wal)));
storage.recover_from_wal()?;
}
storage.rebuild_indexes()?;
Ok(storage)
}
#[cfg(feature = "wal")]
pub fn open_with_wal<P: AsRef<Path>>(path: P, wal_config: WalConfig) -> Result<Self> {
let path_ref = path.as_ref();
let mmap_storage = mmap::MmapStorage::open(path_ref)?;
let wal_path = path_ref.with_extension("wal");
let wal = WriteAheadLog::open_with_config(&wal_path, wal_config)?;
let checkpoint_dir = path_ref.parent().unwrap_or(Path::new("."));
let checkpoint_manager =
CheckpointManager::new(checkpoint_dir, CheckpointConfig::default())?;
let mut storage = Storage {
backend: StorageBackend::MmapFile(mmap_storage),
index_manager: IndexManager::new(),
wal: Some(Arc::new(Mutex::new(wal))),
storage_path: Some(path_ref.to_path_buf()),
checkpoint_manager: Some(checkpoint_manager),
};
storage.recover_from_wal()?;
storage.rebuild_indexes()?;
Ok(storage)
}
#[cfg(feature = "wal")]
fn recover_from_wal(&mut self) -> Result<()> {
let entries = if let Some(wal_arc) = &self.wal {
let wal = wal_arc.lock();
wal.recover()?
} else {
return Ok(());
};
for entry in entries {
match entry.operation {
WalOperation::CreateNode { id: _, data } => {
if let Ok(node) = bincode::deserialize::<Node>(&data) {
self.apply_node_internal(node)?;
}
}
WalOperation::CreateRelationship {
id: _,
from_id: _,
to_id: _,
rel_type: _,
data,
} => {
if let Ok(rel) = bincode::deserialize::<Relationship>(&data) {
self.apply_relationship_internal(rel)?;
}
}
WalOperation::DeleteNode { id: _ } => {
}
WalOperation::CommitTransaction { tx_id: _ } => {
}
WalOperation::Checkpoint { .. } => {
}
_ => {
}
}
}
Ok(())
}
#[cfg(feature = "wal")]
fn apply_node_internal(&mut self, node: Node) -> Result<()> {
match &mut self.backend {
StorageBackend::InMemory(storage) => storage.store_node(node),
StorageBackend::MmapFile(storage) => storage.store_node(node),
}
}
#[cfg(feature = "wal")]
fn apply_relationship_internal(&mut self, rel: Relationship) -> Result<()> {
match &mut self.backend {
StorageBackend::InMemory(storage) => storage.store_relationship(rel),
StorageBackend::MmapFile(storage) => storage.store_relationship(rel),
}
}
pub fn store_node(&mut self, node: Node) -> Result<()> {
#[cfg(feature = "wal")]
if let Some(wal_arc) = &self.wal {
let data =
bincode::serialize(&node).map_err(|e| GraphError::Serialization(e.to_string()))?;
let mut wal = wal_arc.lock();
wal.append(WalOperation::CreateNode { id: node.id, data })?;
}
self.index_manager.add_node(&node);
match &mut self.backend {
StorageBackend::InMemory(storage) => storage.store_node(node),
StorageBackend::MmapFile(storage) => storage.store_node(node),
}
}
pub fn get_node(&self, id: Id) -> Result<Option<Node>> {
match &self.backend {
StorageBackend::InMemory(storage) => storage.get_node(id),
StorageBackend::MmapFile(storage) => storage.get_node(id),
}
}
pub fn store_relationship(&mut self, relationship: Relationship) -> Result<()> {
#[cfg(feature = "wal")]
if let Some(wal_arc) = &self.wal {
let data = bincode::serialize(&relationship)
.map_err(|e| GraphError::Serialization(e.to_string()))?;
let mut wal = wal_arc.lock();
wal.append(WalOperation::CreateRelationship {
id: relationship.id,
from_id: relationship.from_id,
to_id: relationship.to_id,
rel_type: relationship.rel_type.clone(),
data,
})?;
}
self.index_manager.add_relationship(&relationship);
match &mut self.backend {
StorageBackend::InMemory(storage) => storage.store_relationship(relationship),
StorageBackend::MmapFile(storage) => storage.store_relationship(relationship),
}
}
pub fn get_relationship(&self, id: Id) -> Result<Option<Relationship>> {
match &self.backend {
StorageBackend::InMemory(storage) => storage.get_relationship(id),
StorageBackend::MmapFile(storage) => storage.get_relationship(id),
}
}
pub fn get_relationships_for_node(&self, node_id: Id) -> Result<Vec<Relationship>> {
match &self.backend {
StorageBackend::InMemory(storage) => storage.get_relationships_for_node(node_id),
StorageBackend::MmapFile(storage) => storage.get_relationships_for_node(node_id),
}
}
pub fn delete_node(&mut self, id: Id) -> Result<Option<Node>> {
match &mut self.backend {
StorageBackend::InMemory(storage) => storage.delete_node(id),
StorageBackend::MmapFile(storage) => storage.delete_node(id),
}
}
pub fn delete_relationship(&mut self, id: Id) -> Result<Option<Relationship>> {
match &mut self.backend {
StorageBackend::InMemory(storage) => storage.delete_relationship(id),
StorageBackend::MmapFile(storage) => storage.delete_relationship(id),
}
}
pub fn node_has_relationships(&self, node_id: Id) -> bool {
match &self.backend {
StorageBackend::InMemory(storage) => storage.node_has_relationships(node_id),
StorageBackend::MmapFile(storage) => storage.node_has_relationships(node_id),
}
}
pub fn node_count(&self) -> usize {
match &self.backend {
StorageBackend::InMemory(storage) => storage.node_count(),
StorageBackend::MmapFile(storage) => storage.node_count(),
}
}
pub fn relationship_count(&self) -> usize {
match &self.backend {
StorageBackend::InMemory(storage) => storage.relationship_count(),
StorageBackend::MmapFile(storage) => storage.relationship_count(),
}
}
pub fn flush(&self) -> Result<()> {
#[cfg(feature = "wal")]
if let Some(wal_arc) = &self.wal {
let mut wal = wal_arc.lock();
wal.sync()?;
}
match &self.backend {
StorageBackend::InMemory(_) => Ok(()), StorageBackend::MmapFile(storage) => storage.flush(),
}
}
#[cfg(feature = "wal")]
pub fn sync_wal(&self) -> Result<()> {
if let Some(wal_arc) = &self.wal {
let mut wal = wal_arc.lock();
wal.sync()?;
}
Ok(())
}
#[cfg(feature = "wal")]
pub fn checkpoint(&self) -> Result<u64> {
if let Some(wal_arc) = &self.wal {
let mut wal = wal_arc.lock();
wal.checkpoint()
} else {
Err(GraphError::Storage("WAL not enabled".to_string()))
}
}
#[cfg(feature = "wal")]
pub fn has_wal(&self) -> bool {
self.wal.is_some()
}
#[cfg(not(feature = "wal"))]
pub fn has_wal(&self) -> bool {
false
}
#[cfg(feature = "wal")]
pub fn has_checkpoint_manager(&self) -> bool {
self.checkpoint_manager.is_some()
}
#[cfg(not(feature = "wal"))]
pub fn has_checkpoint_manager(&self) -> bool {
false
}
#[cfg(feature = "wal")]
pub fn checkpoint_manager(&self) -> Option<&CheckpointManager> {
self.checkpoint_manager.as_ref()
}
#[cfg(feature = "wal")]
pub fn checkpoint_manager_mut(&mut self) -> Option<&mut CheckpointManager> {
self.checkpoint_manager.as_mut()
}
pub fn create_property_index(&mut self, property_key: String) {
self.index_manager.create_property_index(property_key);
}
pub fn create_range_index(&mut self, property_key: String) {
self.index_manager.create_range_index(property_key);
}
pub fn create_composite_index(&mut self, property_keys: Vec<String>) {
self.index_manager.create_composite_index(property_keys);
}
pub fn find_by_property(&self, property_key: &str, property_value: &Value) -> Vec<Id> {
self.index_manager
.find_by_property(property_key, property_value)
}
pub fn find_in_range(
&self,
property_key: &str,
min_value: &Value,
max_value: &Value,
) -> Vec<Id> {
self.index_manager
.find_in_range(property_key, min_value, max_value)
}
pub fn find_relationships_by_type(&self, rel_type: &str) -> Vec<Id> {
self.index_manager.find_relationships_by_type(rel_type)
}
pub fn find_outgoing_relationships(&self, from_id: Id, rel_type: &str) -> Vec<Id> {
self.index_manager
.find_outgoing_relationships(from_id, rel_type)
}
pub fn find_incoming_relationships(&self, to_id: Id, rel_type: &str) -> Vec<Id> {
self.index_manager
.find_incoming_relationships(to_id, rel_type)
}
pub fn get_index_stats(&self) -> crate::index::IndexStats {
self.index_manager.get_index_stats()
}
fn rebuild_indexes(&mut self) -> Result<()> {
match &self.backend {
StorageBackend::InMemory(storage) => {
for node in storage.nodes.values() {
self.index_manager.add_node(node);
}
for relationship in storage.relationships.values() {
self.index_manager.add_relationship(relationship);
}
}
StorageBackend::MmapFile(_) => {
}
}
Ok(())
}
pub fn index_manager(&self) -> &IndexManager {
&self.index_manager
}
pub fn index_manager_mut(&mut self) -> &mut IndexManager {
&mut self.index_manager
}
}
impl InMemoryStorage {
pub fn store_node(&mut self, node: Node) -> Result<()> {
let id = node.id;
self.nodes.insert(id, node);
self.node_relationships.entry(id).or_default();
Ok(())
}
pub fn get_node(&self, id: Id) -> Result<Option<Node>> {
Ok(self.nodes.get(&id).cloned())
}
pub fn store_relationship(&mut self, relationship: Relationship) -> Result<()> {
let id = relationship.id;
let from_id = relationship.from_id;
let to_id = relationship.to_id;
if !self.nodes.contains_key(&from_id) {
return Err(GraphError::NotFound(format!("Node {from_id} not found")));
}
if !self.nodes.contains_key(&to_id) {
return Err(GraphError::NotFound(format!("Node {to_id} not found")));
}
self.relationships.insert(id, relationship);
self.node_relationships.entry(from_id).or_default().push(id);
if from_id != to_id {
self.node_relationships.entry(to_id).or_default().push(id);
}
Ok(())
}
pub fn get_relationship(&self, id: Id) -> Result<Option<Relationship>> {
Ok(self.relationships.get(&id).cloned())
}
pub fn get_relationships_for_node(&self, node_id: Id) -> Result<Vec<Relationship>> {
let empty_vec = Vec::new();
let relationship_ids = self.node_relationships.get(&node_id).unwrap_or(&empty_vec);
let mut relationships = Vec::new();
for &rel_id in relationship_ids {
if let Some(rel) = self.relationships.get(&rel_id) {
relationships.push(rel.clone());
}
}
Ok(relationships)
}
pub fn node_count(&self) -> usize {
self.nodes.len()
}
pub fn relationship_count(&self) -> usize {
self.relationships.len()
}
pub fn delete_node(&mut self, id: Id) -> Result<Option<Node>> {
self.node_relationships.remove(&id);
Ok(self.nodes.remove(&id))
}
pub fn delete_relationship(&mut self, id: Id) -> Result<Option<Relationship>> {
if let Some(rel) = self.relationships.remove(&id) {
if let Some(rels) = self.node_relationships.get_mut(&rel.from_id) {
rels.retain(|&r| r != id);
}
if rel.from_id != rel.to_id {
if let Some(rels) = self.node_relationships.get_mut(&rel.to_id) {
rels.retain(|&r| r != id);
}
}
Ok(Some(rel))
} else {
Ok(None)
}
}
pub fn node_has_relationships(&self, node_id: Id) -> bool {
self.node_relationships
.get(&node_id)
.map(|rels| !rels.is_empty())
.unwrap_or(false)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
#[test]
fn test_storage_node_operations() {
let mut storage = Storage::new().unwrap();
let node = Node::new(1, HashMap::new());
storage.store_node(node.clone()).unwrap();
let retrieved = storage.get_node(1).unwrap();
assert_eq!(retrieved, Some(node));
assert_eq!(storage.node_count(), 1);
}
#[test]
fn test_storage_relationship_operations() {
let mut storage = Storage::new().unwrap();
let node1 = Node::new(1, HashMap::new());
let node2 = Node::new(2, HashMap::new());
storage.store_node(node1).unwrap();
storage.store_node(node2).unwrap();
let rel = Relationship::new(1, 1, 2, "KNOWS".to_string(), HashMap::new());
storage.store_relationship(rel.clone()).unwrap();
let retrieved = storage.get_relationship(1).unwrap();
assert_eq!(retrieved, Some(rel));
assert_eq!(storage.relationship_count(), 1);
}
#[test]
fn test_node_relationships() {
let mut storage = Storage::new().unwrap();
let node1 = Node::new(1, HashMap::new());
let node2 = Node::new(2, HashMap::new());
storage.store_node(node1).unwrap();
storage.store_node(node2).unwrap();
let rel = Relationship::new(1, 1, 2, "KNOWS".to_string(), HashMap::new());
storage.store_relationship(rel.clone()).unwrap();
let rels = storage.get_relationships_for_node(1).unwrap();
assert_eq!(rels.len(), 1);
assert_eq!(rels[0], rel);
let rels = storage.get_relationships_for_node(2).unwrap();
assert_eq!(rels.len(), 1);
assert_eq!(rels[0], rel);
let knows_rels = storage.find_relationships_by_type("KNOWS");
assert_eq!(knows_rels.len(), 1);
assert!(knows_rels.contains(&1));
let outgoing = storage.find_outgoing_relationships(1, "KNOWS");
assert_eq!(outgoing.len(), 1);
assert!(outgoing.contains(&1));
}
#[test]
fn test_index_stats() {
let mut storage = Storage::new().unwrap();
storage.create_property_index("name".to_string());
storage.create_range_index("age".to_string());
let stats = storage.get_index_stats();
assert_eq!(stats.property_index_count, 1);
assert_eq!(stats.range_index_count, 1);
assert_eq!(stats.composite_index_count, 0);
}
}