diff --git a/src/core/engine.rs b/src/core/engine.rs index dcb35db..cf96581 100644 --- a/src/core/engine.rs +++ b/src/core/engine.rs @@ -1,9 +1,9 @@ use crate::core::log_record::LogRecord; use crate::core::memtable::MemTable; -use crate::infra::codec::decode; use crate::infra::config::LsmConfig; use crate::infra::error::{LsmError, Result}; -use crate::storage::sstable::SStable; +use crate::storage::builder::SstableBuilder; +use crate::storage::reader::SstableReader; use crate::storage::wal::WriteAheadLog; use std::collections::HashMap; @@ -29,7 +29,7 @@ pub struct LsmStats { pub struct LsmEngine { pub(crate) memtable: Mutex, pub(crate) wal: WriteAheadLog, - pub(crate) sstables: Mutex>, + pub(crate) sstables: Mutex>, pub(crate) dir_path: PathBuf, pub(crate) config: LsmConfig, } @@ -46,14 +46,15 @@ impl LsmEngine { let entry = entry?; let path = entry.path(); if path.extension().is_some_and(|ext| ext == "sst") { - match SStable::load(&path) { + match SstableReader::open(path.clone(), config.storage.clone()) { Ok(sst) => sstables.push(sst), Err(e) => warn!("Failed to load SSTable {}: {}", path.display(), e), } } } - sstables.sort_by(|a, b| b.metadata.timestamp.cmp(&a.metadata.timestamp)); + // Sort by timestamp descending (newest first) + sstables.sort_by(|a, b| b.metadata().timestamp.cmp(&a.metadata().timestamp)); let mut memtable = MemTable::new(config.core.memtable_max_size); for record in wal_records { @@ -81,7 +82,7 @@ impl LsmEngine { .map_err(|_| LsmError::LockPoisoned("memtable")) } - fn sstables_lock(&self) -> Result>> { + fn sstables_lock(&self) -> Result>> { self.sstables .lock() .map_err(|_| LsmError::LockPoisoned("sstables")) @@ -128,7 +129,7 @@ impl LsmEngine { } drop(memtable); - // 2. Verificar SSTables (da mais recente para a mais antiga) + // 2. Check SSTables (newest to oldest) let mut sstables = self.sstables_lock()?; for sst in sstables.iter_mut() { if let Some(record) = sst.get(key)? { @@ -189,11 +190,21 @@ impl LsmEngine { } let timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos(); + let filename = format!("{}.sst", timestamp); + let path = self.dir_path.join(filename); - let sst = SStable::create(&self.dir_path, timestamp, &self.config.storage, &records)?; + // Create new SSTable using Builder (V2) + let mut builder = SstableBuilder::new(path, self.config.storage.clone(), timestamp)?; + for (key, record) in records { + builder.add(key.as_bytes(), &record)?; + } + let sst_path = builder.finish()?; + + // Open the new SSTable as Reader (V2) + let reader = SstableReader::open(sst_path, self.config.storage.clone())?; let mut sstables = self.sstables_lock()?; - sstables.insert(0, sst); + sstables.insert(0, reader); let cleared = memtable.clear(); info!( @@ -222,11 +233,12 @@ impl LsmEngine { } drop(memtable); - let sstables = self.sstables_lock()?; - for sst in sstables.iter() { - let records = self.read_all_from_sstable(sst)?; - for record in records { - result_map.entry(record.key.clone()).or_insert(( + let mut sstables = self.sstables_lock()?; + for sst in sstables.iter_mut() { + let records = sst.scan()?; + for (key_bytes, record) in records { + let key = String::from_utf8(key_bytes).map_err(|e| LsmError::CorruptedData(e.to_string()))?; + result_map.entry(key).or_insert(( record.value, record.timestamp, record.is_deleted, @@ -250,35 +262,6 @@ impl LsmEngine { Ok(results) } - fn read_all_from_sstable(&self, sst: &SStable) -> Result> { - use std::fs::File; - use std::io::{BufReader, Read, Seek, SeekFrom}; - - let mut file = BufReader::new(File::open(&sst.path)?); - file.seek(SeekFrom::Current(8))?; - - let mut len_buf = [0u8; 4]; - - file.read_exact(&mut len_buf)?; - let bloom_len = u32::from_le_bytes(len_buf) as usize; - file.seek(SeekFrom::Current(bloom_len as i64))?; - - file.read_exact(&mut len_buf)?; - let meta_len = u32::from_le_bytes(len_buf) as usize; - file.seek(SeekFrom::Current(meta_len as i64))?; - - let mut records = Vec::new(); - for _ in 0..sst.metadata.record_count { - file.read_exact(&mut len_buf)?; - let record_len = u32::from_le_bytes(len_buf) as usize; - let mut record_data = vec![0u8; record_len]; - file.read_exact(&mut record_data)?; - let record: LogRecord = decode(&record_data)?; - records.push(record); - } - Ok(records) - } - pub fn keys(&self) -> Result> { let all_data = self.scan()?; Ok(all_data.into_iter().map(|(k, _)| k).collect()) @@ -313,12 +296,12 @@ impl LsmEngine { let mem_records = memtable.data.len(); let sst_records_total: u64 = sstables .iter() - .map(|s| s.metadata.record_count as u64) + .map(|s| s.metadata().record_count) .sum(); let sst_bytes_total: u64 = sstables .iter() - .map(|s| std::fs::metadata(&s.path).map(|m| m.len()).unwrap_or(0)) + .map(|s| std::fs::metadata(s.path()).map(|m| m.len()).unwrap_or(0)) .sum(); let wal_bytes: u64 = std::fs::metadata(&self.wal.path) diff --git a/src/storage/block.rs b/src/storage/block.rs index 6aa4cdf..8640563 100644 --- a/src/storage/block.rs +++ b/src/storage/block.rs @@ -2,12 +2,12 @@ use crate::infra::config::StorageConfig; use std::mem::size_of; pub const BLOCK_SIZE: usize = 4096; -const U16_SIZE: usize = size_of::(); +const U32_SIZE: usize = size_of::(); #[derive(Debug, Clone)] pub struct Block { pub(crate) data: Vec, - pub(crate) offsets: Vec, + pub(crate) offsets: Vec, block_size: usize, } @@ -25,11 +25,15 @@ impl Block { } fn entry_size(key: &[u8], value: &[u8]) -> usize { - U16_SIZE + key.len() + U16_SIZE + value.len() + // KeyLen(2) + Key + ValLen(2) + Value + // Note: Using u16 for key/value length storage within the block data + // to maintain compactness for individual entries, while allowing + // the overall block to be larger than 64KB via u32 offsets. + 2 + key.len() + 2 + value.len() } fn metadata_size(num_entries: usize) -> usize { - (num_entries * U16_SIZE) + U16_SIZE + (num_entries * U32_SIZE) + U32_SIZE } fn current_size(&self) -> usize { @@ -38,16 +42,18 @@ impl Block { pub fn add(&mut self, key: &[u8], value: &[u8]) -> bool { let entry_size = Self::entry_size(key, value); - let new_offset_size = U16_SIZE; + let new_offset_size = U32_SIZE; let total_needed = self.current_size() + entry_size + new_offset_size; if total_needed > self.block_size { return false; } - let offset = self.data.len() as u16; + let offset = self.data.len() as u32; self.offsets.push(offset); + // Cast to u16 is safe for key/value lengths as we assume + // individual entries don't exceed 64KB, even if the block does. let key_len = key.len() as u16; let val_len = value.len() as u16; @@ -67,14 +73,14 @@ impl Block { encoded.extend_from_slice(&offset.to_le_bytes()); } - let num_elements = self.offsets.len() as u16; + let num_elements = self.offsets.len() as u32; encoded.extend_from_slice(&num_elements.to_le_bytes()); encoded } pub fn decode(data: &[u8]) -> Self { - if data.len() < U16_SIZE { + if data.len() < U32_SIZE { return Self { data: Vec::new(), offsets: Vec::new(), @@ -82,20 +88,20 @@ impl Block { }; } - let num_elements_start = data.len() - U16_SIZE; + let num_elements_start = data.len() - U32_SIZE; let num_elements = - u16::from_le_bytes([data[num_elements_start], data[num_elements_start + 1]]) as usize; + u32::from_le_bytes([data[num_elements_start], data[num_elements_start + 1], data[num_elements_start + 2], data[num_elements_start + 3]]) as usize; - let offsets_start = data.len() - U16_SIZE - (num_elements * U16_SIZE); + let offsets_start = data.len() - U32_SIZE - (num_elements * U32_SIZE); let records_data = data[..offsets_start].to_vec(); let mut offsets = Vec::with_capacity(num_elements); let mut offset_pos = offsets_start; for _ in 0..num_elements { - let offset = u16::from_le_bytes([data[offset_pos], data[offset_pos + 1]]); + let offset = u32::from_le_bytes([data[offset_pos], data[offset_pos + 1], data[offset_pos + 2], data[offset_pos + 3]]); offsets.push(offset); - offset_pos += U16_SIZE; + offset_pos += U32_SIZE; } Self { @@ -178,6 +184,41 @@ mod tests { assert!(!result, "Should reject entry when block is full"); } + #[test] + fn test_block_overflow_u16() { + // Create a block larger than 64KB (u16::MAX is 65535) + let block_size = 70_000; + let mut block = Block::new(block_size); + + // Fill with enough data to exceed 64KB + // Each entry approx 1024 bytes + let val_size = 1000; + let large_value = vec![b'x'; val_size]; + let key_base = "key"; + + let mut count = 0; + while block.data_size() < 66000 { + let key = format!("{}{}", key_base, count); + if !block.add(key.as_bytes(), &large_value) { + break; + } + count += 1; + } + + assert!(block.data_size() > 65535, "Block data size should be > 64KB"); + + // Verify integrity + let encoded = block.encode(); + let decoded = Block::decode(&encoded); + + assert_eq!(decoded.len(), block.len()); + assert_eq!(decoded.offsets.len(), block.offsets.len()); + + // Verify last entry is correct + let last_offset = *decoded.offsets.last().unwrap(); + assert!(last_offset > 65535, "Last offset should exceed u16 limit"); + } + #[test] fn test_overflow_large_entry() { let mut block = Block::new(128); diff --git a/src/storage/builder.rs b/src/storage/builder.rs index e18d92f..1a948ac 100644 --- a/src/storage/builder.rs +++ b/src/storage/builder.rs @@ -10,7 +10,7 @@ use std::fs::File; use std::io::{BufWriter, Write}; use std::path::PathBuf; -const SST_MAGIC_V2: &[u8; 8] = b"LSMSST02"; +const SST_MAGIC_V2: &[u8; 8] = b"LSMSST03"; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct BlockMeta { diff --git a/src/storage/config.rs b/src/storage/config.rs index cea64fc..4ee1284 100644 --- a/src/storage/config.rs +++ b/src/storage/config.rs @@ -1,8 +1,18 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum CompactionStrategy { + SizeTiered, + Leveled, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] pub struct StorageConfig { pub block_size: usize, pub block_cache_size_mb: usize, pub sparse_index_interval: usize, pub compaction_strategy: CompactionStrategy, + pub bloom_false_positive_rate: f64, } impl Default for StorageConfig { @@ -12,13 +22,7 @@ impl Default for StorageConfig { block_cache_size_mb: 64, sparse_index_interval: 16, compaction_strategy: CompactionStrategy::SizeTiered, + bloom_false_positive_rate: 0.01, } } } - -// src/core/engine.rs -pub struct LsmConfig { - pub dir_path: PathBuf, - pub memtable_max_size: usize, - pub storage: StorageConfig, // ✅ Composição -} diff --git a/src/storage/mod.rs b/src/storage/mod.rs index b7cbadb..4107c64 100644 --- a/src/storage/mod.rs +++ b/src/storage/mod.rs @@ -1,5 +1,5 @@ pub mod block; pub mod builder; +pub mod config; pub mod reader; -pub mod sstable; pub mod wal; diff --git a/src/storage/reader.rs b/src/storage/reader.rs index 4da5fb2..0c83f62 100644 --- a/src/storage/reader.rs +++ b/src/storage/reader.rs @@ -12,7 +12,7 @@ use std::io::{Read, Seek, SeekFrom}; use std::num::NonZeroUsize; use std::path::PathBuf; -const SST_MAGIC_V2: &[u8; 8] = b"LSMSST02"; +const SST_MAGIC_V2: &[u8; 8] = b"LSMSST03"; const FOOTER_SIZE: u64 = 8; /// SSTable V2 Reader with sparse index, Bloom filter, and block caching diff --git a/src/storage/sstable/mod.rs b/src/storage/sstable/mod.rs deleted file mode 100644 index 00fe013..0000000 --- a/src/storage/sstable/mod.rs +++ /dev/null @@ -1,392 +0,0 @@ -use crate::core::log_record::LogRecord; -use crate::infra::codec::{decode, encode}; -use crate::infra::config::StorageConfig; -use crate::infra::error::{LsmError, Result}; -use bloomfilter::Bloom; -use crc32fast; -use serde::{Deserialize, Serialize}; -use std::fs::File; -use std::io::{BufWriter, Read, Seek, SeekFrom, Write}; -use std::path::{Path, PathBuf}; -use tracing::debug; - -const SST_MAGIC: &[u8; 8] = b"LSMSST01"; -const BLOCK_SIZE: usize = 4096; // 4KB blocks - -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct SstableMetadata { - pub timestamp: u128, - pub min_key: String, - pub max_key: String, - pub record_count: u32, - pub checksum: u32, -} - -/// Metadata for a single block in the sparse index -#[derive(Clone, Debug, Serialize, Deserialize)] -pub struct BlockMeta { - pub first_key: String, - pub offset: u64, - pub size: u32, -} - -#[derive(Debug)] -pub struct SStable { - pub(crate) metadata: SstableMetadata, - pub(crate) bloom_filter: Bloom<[u8]>, - pub(crate) index: Vec, - pub(crate) file: File, - pub(crate) path: PathBuf, -} - -fn read_and_check_magic(mut r: R) -> Result<()> { - let mut magic = [0u8; 8]; - r.read_exact(&mut magic)?; - if &magic != SST_MAGIC { - return Err(LsmError::InvalidSstable); - } - Ok(()) -} - -impl SStable { - pub fn create( - dir_path: &Path, - timestamp: u128, - config: &StorageConfig, - records: &[(String, LogRecord)], - ) -> Result { - if records.is_empty() { - return Err(LsmError::CompactionFailed( - "Cannot create SSTable with empty records".to_string(), - )); - } - - let path = dir_path.join(format!("{timestamp}.sst")); - let mut file = BufWriter::new(File::create(&path)?); - file.write_all(SST_MAGIC)?; - - let mut bloom = - Bloom::<[u8]>::new_for_fp_rate(records.len(), config.bloom_false_positive_rate) - .map_err(|e| LsmError::CompactionFailed(e.to_string()))?; - - for (key, _) in records.iter() { - bloom.set(key.as_bytes()); - } - - let bloom_bytes = bloom.into_bytes(); - file.write_all(&(bloom_bytes.len() as u32).to_le_bytes())?; - file.write_all(&bloom_bytes)?; - - // 2) Metadata - let checksum = crc32fast::hash(&encode(&records)?); // Checksum over serialized records - - let metadata = SstableMetadata { - timestamp, - min_key: records[0].0.clone(), - max_key: records[records.len() - 1].0.clone(), - record_count: records.len() as u32, - checksum, - }; - - let metadata_bytes = encode(&metadata)?; - file.write_all(&(metadata_bytes.len() as u32).to_le_bytes())?; - file.write_all(&metadata_bytes)?; - - // 3) Write blocks and build sparse index - let mut index = Vec::new(); - let mut current_block = Vec::new(); - let mut current_block_size = 0usize; - let blocks_start_offset = file.stream_position()?; - let mut current_offset = blocks_start_offset; - - for (key, record) in records.iter() { - let record_bytes = encode(record)?; - let entry_size = 4 + record_bytes.len(); // u32 length + data - - // Check if adding this record would exceed block size - if current_block_size + entry_size > BLOCK_SIZE && !current_block.is_empty() { - // Write current block - let block_data = serialize_block(¤t_block)?; - file.write_all(&block_data)?; - - // Add to index - index.push(BlockMeta { - first_key: current_block[0].0.clone(), - offset: current_offset, - size: block_data.len() as u32, - }); - - current_offset += block_data.len() as u64; - current_block.clear(); - current_block_size = 0; - } - - current_block.push((key.clone(), record.clone())); - current_block_size += entry_size; - } - - // Write last block - if !current_block.is_empty() { - let block_data = serialize_block(¤t_block)?; - file.write_all(&block_data)?; - - index.push(BlockMeta { - first_key: current_block[0].0.clone(), - offset: current_offset, - size: block_data.len() as u32, - }); - } - - // 4) Write sparse index - let index_offset = file.stream_position()?; - let index_bytes = encode(&index)?; - file.write_all(&(index_bytes.len() as u32).to_le_bytes())?; - file.write_all(&index_bytes)?; - - // 5) Write footer (index offset) - file.write_all(&index_offset.to_le_bytes())?; - - file.flush()?; - file.get_ref().sync_all()?; - - debug!( - "SSTable created: {}, records={}, blocks={}, checksum={}", - path.display(), - metadata.record_count, - index.len(), - metadata.checksum - ); - - // 6) Open file for reading and rebuild bloom from bytes - let read_file = File::open(&path)?; - let bloom_filter = Bloom::<[u8]>::from_bytes(bloom_bytes) - .map_err(|e| LsmError::CompactionFailed(e.to_string()))?; - - Ok(Self { - metadata, - bloom_filter, - index, - file: read_file, - path, - }) - } - - /// Open an existing SSTable file with lazy loading (only footer + index) - pub fn open(path: &Path) -> Result { - let mut file = File::open(path)?; - - // 1. Read footer (last 8 bytes = index offset) - file.seek(SeekFrom::End(-8))?; - let mut footer = [0u8; 8]; - file.read_exact(&mut footer)?; - let index_offset = u64::from_le_bytes(footer); - - // 2. Read sparse index - file.seek(SeekFrom::Start(index_offset))?; - let mut len_buf = [0u8; 4]; - file.read_exact(&mut len_buf)?; - let index_len = u32::from_le_bytes(len_buf) as usize; - let mut index_data = vec![0u8; index_len]; - file.read_exact(&mut index_data)?; - let index: Vec = decode(&index_data)?; - - if index.is_empty() { - return Err(LsmError::InvalidSstable); - } - - // 3. Read header, bloom filter, and metadata - file.seek(SeekFrom::Start(0))?; - read_and_check_magic(&mut file)?; - - // Bloom filter - file.read_exact(&mut len_buf)?; - let bloom_len = u32::from_le_bytes(len_buf) as usize; - let mut bloom_data = vec![0u8; bloom_len]; - file.read_exact(&mut bloom_data)?; - let bloom = Bloom::<[u8]>::from_bytes(bloom_data).map_err(|_| LsmError::InvalidSstable)?; - - // Metadata - file.read_exact(&mut len_buf)?; - let meta_len = u32::from_le_bytes(len_buf) as usize; - let mut meta_data = vec![0u8; meta_len]; - file.read_exact(&mut meta_data)?; - let metadata: SstableMetadata = decode(&meta_data)?; - - debug!( - "SSTable opened: {}, records={}, blocks={}", - path.display(), - metadata.record_count, - index.len() - ); - - Ok(Self { - metadata, - bloom_filter: bloom, - index, - file, - path: path.to_path_buf(), - }) - } - - /// Legacy load method for backward compatibility (delegates to open) - pub fn load(path: &Path) -> Result { - Self::open(path) - } - - /// Read a specific block from disk - fn read_block(&mut self, block_meta: &BlockMeta) -> Result> { - self.file.seek(SeekFrom::Start(block_meta.offset))?; - - let mut block_data = vec![0u8; block_meta.size as usize]; - self.file.read_exact(&mut block_data)?; - - deserialize_block(&block_data) - } - - pub fn get(&mut self, key: &str) -> Result> { - // 1. Check bloom filter - if !self.bloom_filter.check(key.as_bytes()) { - return Ok(None); - } - - // 2. Binary search on sparse index using partition_point - // Find the first block where first_key > search_key - let block_idx = self - .index - .partition_point(|block_meta| block_meta.first_key.as_str() <= key); - - // Edge case: key is smaller than the first key of the first block - if block_idx == 0 { - return Ok(None); - } - - // The candidate block is at index block_idx - 1 - let candidate_idx = block_idx - 1; - let block_meta = &self.index[candidate_idx].clone(); - - // 3. Load the block from disk - let records = self.read_block(block_meta)?; - - // 4. Linear search within the block - for record in records { - if record.key == key { - return Ok(Some(record)); - } - } - - Ok(None) - } -} - -/// Serialize a block of records into bytes -fn serialize_block(records: &[(String, LogRecord)]) -> Result> { - let mut block_data = Vec::new(); - for (_key, record) in records { - let record_bytes = encode(record)?; - let len = record_bytes.len() as u32; - block_data.extend_from_slice(&len.to_le_bytes()); - block_data.extend_from_slice(&record_bytes); - } - Ok(block_data) -} - -/// Deserialize a block of bytes into records -fn deserialize_block(block_data: &[u8]) -> Result> { - let mut cursor = std::io::Cursor::new(block_data); - let mut records = Vec::new(); - let mut len_buf = [0u8; 4]; - - while cursor.position() < block_data.len() as u64 { - cursor.read_exact(&mut len_buf)?; - let record_len = u32::from_le_bytes(len_buf) as usize; - - let mut record_data = vec![0u8; record_len]; - cursor.read_exact(&mut record_data)?; - - let record: LogRecord = decode(&record_data)?; - records.push(record); - } - - Ok(records) -} - -#[cfg(test)] -mod tests { - use super::*; - use tempfile::tempdir; - - // NOTE: Assuming StorageConfig implements the Default trait or has a public constructor. - // This change is necessary to satisfy the function signature of SStable::create. - - #[test] - fn test_sstable_create_and_open() { - let dir = tempdir().unwrap(); - let timestamp = 12345u128; - - // Create test records - let records: Vec<(String, LogRecord)> = (0..100) - .map(|i| { - let key = format!("key_{:03}", i); - let record = LogRecord::new(key.clone(), format!("value_{}", i).into_bytes()); - (key, record) - }) - .collect(); - - // Create SSTable - let config = StorageConfig::default(); // Assuming StorageConfig implements Default - let sstable = SStable::create(dir.path(), timestamp, &config, &records).unwrap(); - assert_eq!(sstable.metadata.record_count, 100); - assert!(sstable.index.len() > 0); - - // Close and reopen - drop(sstable); - let mut reopened = SStable::open(&dir.path().join(format!("{}.sst", timestamp))).unwrap(); - assert_eq!(reopened.metadata.record_count, 100); - assert!(reopened.index.len() > 0); - - // Test get operations - let result = reopened.get("key_050").unwrap(); - assert!(result.is_some()); - assert_eq!(result.unwrap().key, "key_050"); - - // Test non-existent key - let result = reopened.get("key_999").unwrap(); - assert!(result.is_none()); - } - - #[test] - fn test_sparse_index_edge_cases() { - let dir = tempdir().unwrap(); - let timestamp = 67890u128; - - let records: Vec<(String, LogRecord)> = vec![ - ( - "apple".to_string(), - LogRecord::new("apple".to_string(), b"a".to_vec()), - ), - ( - "banana".to_string(), - LogRecord::new("banana".to_string(), b"b".to_vec()), - ), - ( - "cherry".to_string(), - LogRecord::new("cherry".to_string(), b"c".to_vec()), - ), - ]; - - let config = StorageConfig::default(); // Assuming StorageConfig implements Default - let mut sstable = SStable::create(dir.path(), timestamp, &config, &records).unwrap(); - - // Key before first key - assert!(sstable.get("aardvark").unwrap().is_none()); - - // Exact first key - assert!(sstable.get("apple").unwrap().is_some()); - - // Key after last key - assert!(sstable.get("zebra").unwrap().is_none()); - - // Middle key - assert!(sstable.get("banana").unwrap().is_some()); - } -}