use derive_more::{Deref, DerefMut, From, Into};
use serde::{Deserialize, Serialize};
use std::path::PathBuf;
use tokio::sync::RwLock;
use std::sync::Arc;
pub mod segment;
pub mod wal;
pub mod event;
pub mod error;
pub mod memory;
pub mod ultra_performance;
pub mod gpu_acceleration;
pub mod security;
pub use error::{StorageError, StorageResult};
pub use event::{Event, EventId, EventData};
pub use segment::{Segment, SegmentId, SegmentManager};
pub use wal::{WriteAheadLog, WalConfig};
pub use memory::{MemoryPool, PooledBuffer};
pub use ultra_performance::{
UltraPerformanceEngine,
LockFreeRingBuffer,
SIMDBatchProcessor,
PerformanceMetrics,
NUMATopology,
};
pub use gpu_acceleration::{
GPUAcceleratedEngine,
GPUAccelerationType,
GPUContext,
GPUMetrics,
HybridConfig,
DistributionStrategy,
};
pub use security::{
SecurityManager,
SecurityConfig,
SecurityMetrics,
SecurityError,
RateLimiter,
ResourceMonitor,
GPUTimeoutManager,
};
#[derive(Debug)]
pub struct StorageEngine {
wal: RwLock<WriteAheadLog>,
segment_manager: RwLock<SegmentManager>,
memory_pool: Arc<MemoryPool>,
config: StorageConfig,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StorageConfig {
pub data_dir: PathBuf,
pub segment_size: u64,
pub sync_interval_ms: u64,
pub enable_compression: bool,
pub memory_pool_size: usize,
pub batch_size: usize,
pub worker_count: Option<usize>,
}
impl Default for StorageConfig {
fn default() -> Self {
Self {
data_dir: PathBuf::from("./data"),
segment_size: 1024 * 1024 * 1024, sync_interval_ms: 1000,
enable_compression: false,
memory_pool_size: 128 * 1024 * 1024, batch_size: 1000,
worker_count: None, }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize, Deref, DerefMut, From, Into)]
pub struct Offset(pub u64);
impl Offset {
pub fn new(value: u64) -> Self {
Self(value)
}
pub fn next(self) -> Self {
Self(self.0 + 1)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, Deref, DerefMut, From, Into)]
pub struct Partition(pub u32);
impl Partition {
pub fn new(value: u32) -> Self {
Self(value)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, Deref, DerefMut, From, Into)]
pub struct Topic(pub String);
impl Topic {
pub fn new(name: impl Into<String>) -> Self {
Self(name.into())
}
}
impl StorageEngine {
pub async fn new(config: StorageConfig) -> StorageResult<Self> {
tokio::fs::create_dir_all(&config.data_dir).await?;
let memory_pool = Arc::new(MemoryPool::new(
config.memory_pool_size,
config.batch_size,
)?);
let wal_config = WalConfig {
data_dir: config.data_dir.clone(),
segment_size: config.segment_size,
sync_interval_ms: config.sync_interval_ms,
};
let wal = WriteAheadLog::new(wal_config).await?;
let segment_manager = SegmentManager::new(config.data_dir.clone()).await?;
Ok(Self {
wal: RwLock::new(wal),
segment_manager: RwLock::new(segment_manager),
memory_pool,
config,
})
}
pub async fn append_event(
&self,
topic: &Topic,
partition: Partition,
event_data: EventData,
) -> StorageResult<(EventId, Offset)> {
let event_id = EventId::new();
let event = Event::new(event_id, topic.clone(), partition, event_data);
let mut wal = self.wal.write().await;
let offset = wal.append(&event).await?;
Ok((event_id, offset))
}
pub async fn read_events(
&self,
topic: &Topic,
partition: Partition,
start_offset: Offset,
max_events: usize,
) -> StorageResult<Vec<Event>> {
let wal = self.wal.read().await;
wal.read_events(topic, partition, start_offset, max_events).await
}
pub async fn get_latest_offset(
&self,
topic: &Topic,
partition: Partition,
) -> StorageResult<Option<Offset>> {
let wal = self.wal.read().await;
wal.get_latest_offset(topic, partition).await
}
pub async fn append_events_batch(
&self,
events: Vec<(Topic, Partition, EventData)>,
) -> StorageResult<Vec<(EventId, Offset)>> {
if events.is_empty() {
return Ok(Vec::new());
}
let mut pooled_buffer = self.memory_pool.get_buffer().await?;
let mut batch_events = Vec::with_capacity(events.len());
let mut result_ids = Vec::with_capacity(events.len());
for (topic, partition, event_data) in events {
let event_id = EventId::new();
let event = Event::new(event_id, topic, partition, event_data);
batch_events.push(event);
result_ids.push(event_id);
}
let mut wal = self.wal.write().await;
let offsets = wal.append_batch(&batch_events, &mut pooled_buffer).await?;
let results: Vec<(EventId, Offset)> = result_ids.into_iter()
.zip(offsets.into_iter())
.collect();
Ok(results)
}
pub fn memory_pool_stats(&self) -> memory::PoolStats {
self.memory_pool.stats()
}
pub async fn shutdown(&self) -> StorageResult<()> {
let wal = self.wal.write().await;
wal.sync().await?;
tracing::info!("Storage engine shutdown completed");
Ok(())
}
}