pub mod allocator_safe;
pub mod batch;
pub mod cache;
pub mod interning;
pub mod profiler;
pub use batch::{BatchContext, BatchManager, BatchMemoryTracker, BatchStats};
pub use profiler::{
AllocationCategory, AllocationEvent, MemoryProfiler, MemorySnapshot, SharedMemoryProfiler,
};
#[cfg(feature = "unsafe-allocators")]
pub mod allocator;
#[cfg(feature = "unsafe-allocators")]
pub mod pool;
use crate::error::{GraphError, Result};
use crate::graph::Id;
use bumpalo::Bump;
use lru::LruCache;
use parking_lot::RwLock;
use std::num::NonZeroUsize;
use std::sync::Arc;
pub struct ArenaHandle<T> {
value: T,
}
impl<T> ArenaHandle<T> {
fn new(value: T) -> Self {
Self { value }
}
pub fn get(&self) -> &T {
&self.value
}
pub fn get_mut(&mut self) -> &mut T {
&mut self.value
}
}
#[derive(Debug, Clone)]
pub struct MemoryConfig {
pub max_memory_bytes: usize,
pub arena_size_bytes: usize,
pub node_cache_capacity: usize,
pub relationship_cache_capacity: usize,
pub property_cache_capacity: usize,
pub string_table_capacity: usize,
}
impl Default for MemoryConfig {
fn default() -> Self {
Self {
max_memory_bytes: 1024 * 1024 * 1024, arena_size_bytes: 64 * 1024 * 1024, node_cache_capacity: 100_000,
relationship_cache_capacity: 400_000,
property_cache_capacity: 200_000,
string_table_capacity: 50_000,
}
}
}
#[allow(clippy::arc_with_non_send_sync)]
pub struct MemoryManager {
config: MemoryConfig,
arena: Arc<RwLock<Bump>>,
node_cache: Arc<RwLock<LruCache<Id, Vec<u8>>>>,
relationship_cache: Arc<RwLock<LruCache<Id, Vec<u8>>>>,
string_interner: Arc<RwLock<interning::StringInterner>>,
memory_stats: Arc<RwLock<MemoryStats>>,
batch_manager: Arc<RwLock<batch::BatchManager>>,
}
#[derive(Debug, Default)]
pub struct MemoryStats {
pub total_allocated: usize,
pub node_cache_bytes: usize,
pub relationship_cache_bytes: usize,
pub string_table_bytes: usize,
pub arena_bytes: usize,
pub cache_hit_rate: f64,
pub cache_hits: u64,
pub cache_misses: u64,
}
impl MemoryManager {
#[allow(clippy::arc_with_non_send_sync)]
pub fn new(config: MemoryConfig) -> Result<Self> {
let node_cache_cap = NonZeroUsize::new(config.node_cache_capacity)
.ok_or_else(|| GraphError::Memory("Invalid node cache capacity".to_string()))?;
let rel_cache_cap = NonZeroUsize::new(config.relationship_cache_capacity)
.ok_or_else(|| GraphError::Memory("Invalid relationship cache capacity".to_string()))?;
Ok(Self {
arena: Arc::new(RwLock::new(Bump::with_capacity(config.arena_size_bytes))),
node_cache: Arc::new(RwLock::new(LruCache::new(node_cache_cap))),
relationship_cache: Arc::new(RwLock::new(LruCache::new(rel_cache_cap))),
string_interner: Arc::new(RwLock::new(interning::StringInterner::new(
config.string_table_capacity,
))),
memory_stats: Arc::new(RwLock::new(MemoryStats::default())),
batch_manager: Arc::new(RwLock::new(batch::BatchManager::default())),
config,
})
}
pub fn with_defaults() -> Result<Self> {
Self::new(MemoryConfig::default())
}
pub fn arena(&self) -> &Arc<RwLock<Bump>> {
&self.arena
}
pub fn reset_arena(&self) {
let mut arena = self.arena.write();
arena.reset();
let mut stats = self.memory_stats.write();
stats.arena_bytes = 0;
}
pub fn arena_alloc<T>(&self, value: T) -> ArenaHandle<T> {
ArenaHandle::new(value)
}
pub fn get_cached_node(&self, id: Id) -> Option<Vec<u8>> {
let mut cache = self.node_cache.write();
let result = cache.get(&id).cloned();
let mut stats = self.memory_stats.write();
if result.is_some() {
stats.cache_hits += 1;
} else {
stats.cache_misses += 1;
}
stats.cache_hit_rate =
stats.cache_hits as f64 / (stats.cache_hits + stats.cache_misses) as f64;
result
}
pub fn cache_node(&self, id: Id, data: Vec<u8>) {
let data_size = data.len();
let mut cache = self.node_cache.write();
cache.put(id, data);
let mut stats = self.memory_stats.write();
stats.node_cache_bytes += data_size;
}
pub fn get_cached_relationship(&self, id: Id) -> Option<Vec<u8>> {
let mut cache = self.relationship_cache.write();
let result = cache.get(&id).cloned();
let mut stats = self.memory_stats.write();
if result.is_some() {
stats.cache_hits += 1;
} else {
stats.cache_misses += 1;
}
stats.cache_hit_rate =
stats.cache_hits as f64 / (stats.cache_hits + stats.cache_misses) as f64;
result
}
pub fn cache_relationship(&self, id: Id, data: Vec<u8>) {
let data_size = data.len();
let mut cache = self.relationship_cache.write();
cache.put(id, data);
let mut stats = self.memory_stats.write();
stats.relationship_cache_bytes += data_size;
}
pub fn intern_string(&self, s: &str) -> u32 {
let interner = self.string_interner.write();
interner.intern(s)
}
pub fn get_interned_string(&self, id: u32) -> Option<String> {
let interner = self.string_interner.read();
interner.get(id)
}
pub fn memory_stats(&self) -> MemoryStats {
let stats = self.memory_stats.read();
MemoryStats {
total_allocated: stats.total_allocated,
node_cache_bytes: stats.node_cache_bytes,
relationship_cache_bytes: stats.relationship_cache_bytes,
string_table_bytes: stats.string_table_bytes,
arena_bytes: stats.arena_bytes,
cache_hit_rate: stats.cache_hit_rate,
cache_hits: stats.cache_hits,
cache_misses: stats.cache_misses,
}
}
pub fn check_memory_pressure(&self) -> Result<()> {
let stats = self.memory_stats.read();
let total_used = stats.total_allocated;
if total_used > (self.config.max_memory_bytes * 9) / 10 {
drop(stats);
self.cleanup_memory()?;
}
Ok(())
}
pub fn cleanup_memory(&self) -> Result<()> {
self.reset_arena();
{
let mut node_cache = self.node_cache.write();
let current_size = node_cache.len();
for _ in 0..(current_size / 2) {
node_cache.pop_lru();
}
}
{
let mut rel_cache = self.relationship_cache.write();
let current_size = rel_cache.len();
for _ in 0..(current_size / 2) {
rel_cache.pop_lru();
}
}
let mut stats = self.memory_stats.write();
stats.node_cache_bytes /= 2;
stats.relationship_cache_bytes /= 2;
stats.total_allocated =
stats.node_cache_bytes + stats.relationship_cache_bytes + stats.string_table_bytes;
Ok(())
}
pub fn config(&self) -> &MemoryConfig {
&self.config
}
pub fn batch_manager(&self) -> &Arc<RwLock<batch::BatchManager>> {
&self.batch_manager
}
pub fn start_batch(&self) -> batch::BatchContext {
self.batch_manager.read().start_batch()
}
pub fn batch_stats(&self) -> batch::BatchStats {
self.batch_manager.read().get_stats()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_memory_manager_creation() {
let manager = MemoryManager::with_defaults().unwrap();
let stats = manager.memory_stats();
assert_eq!(stats.cache_hits, 0);
assert_eq!(stats.cache_misses, 0);
}
#[test]
fn test_arena_allocation() {
let manager = MemoryManager::with_defaults().unwrap();
let handle = manager.arena_alloc(42u32);
assert_eq!(*handle.get(), 42);
manager.reset_arena();
let stats = manager.memory_stats();
assert_eq!(stats.arena_bytes, 0);
}
#[test]
fn test_node_caching() {
let manager = MemoryManager::with_defaults().unwrap();
let node_id = 1;
let node_data = vec![1, 2, 3, 4];
assert!(manager.get_cached_node(node_id).is_none());
manager.cache_node(node_id, node_data.clone());
assert_eq!(manager.get_cached_node(node_id), Some(node_data));
let stats = manager.memory_stats();
assert_eq!(stats.cache_hits, 1);
assert_eq!(stats.cache_misses, 1);
assert!(stats.cache_hit_rate > 0.0);
}
#[test]
fn test_string_interning() {
let manager = MemoryManager::with_defaults().unwrap();
let id1 = manager.intern_string("hello");
let id2 = manager.intern_string("world");
let id3 = manager.intern_string("hello");
assert_eq!(id1, id3); assert_ne!(id1, id2);
assert_eq!(manager.get_interned_string(id1), Some("hello".to_string()));
assert_eq!(manager.get_interned_string(id2), Some("world".to_string()));
}
}