use super::{BlobRef, BloomFilter, CompressionAlgorithm, Key, LSMConfig, Value, ValueData};
use crate::{Result, StorageError};
use std::fs::{File, OpenOptions};
use std::io::{BufReader, BufWriter, Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use memmap2::Mmap;
#[derive(Clone)]
pub struct ValueBytes {
pub(crate) block: Arc<Vec<u8>>,
pub(crate) start: usize,
pub(crate) len: usize,
}
impl ValueBytes {
pub fn as_slice(&self) -> &[u8] {
&self.block[self.start..self.start + self.len]
}
}
const SSTABLE_MAGIC: u32 = 0x4C534D54;
const SSTABLE_VERSION: u32 = 1;
pub struct SSTable {
path: PathBuf,
mmap: Option<Arc<Mmap>>,
#[allow(dead_code)]
file: File,
index: BlockIndex,
bloom: BloomFilter,
footer: Footer,
}
#[derive(Clone, Debug)]
pub struct BlockIndexEntry {
pub first_key: Key,
pub offset: u64,
pub size: u32,
pub last_key: Key,
}
#[derive(Clone, Debug)]
pub struct BlockIndex {
entries: Vec<BlockIndexEntry>,
}
#[derive(Clone, Debug)]
struct Footer {
magic: u32,
version: u32,
index_offset: u64,
index_size: u32,
bloom_offset: u64,
bloom_size: u32,
num_entries: u64,
min_timestamp: u64,
max_timestamp: u64,
max_key: u64,
}
impl SSTable {
pub fn read_metadata_with_keys<P: AsRef<Path>>(path: P) -> Result<(u64, u64, u64, Key, Key)> {
let path = path.as_ref();
let mut file = OpenOptions::new().read(true).open(path)?;
let footer = Self::read_footer(&mut file)?;
let file_size = file.metadata()?.len();
let max_key = footer.max_key;
let min_key = if footer.index_size > 0 {
file.seek(SeekFrom::Start(footer.index_offset))?;
let mut index_buf = vec![0u8; footer.index_size as usize];
file.read_exact(&mut index_buf)?;
match BlockIndex::deserialize(&index_buf) {
Ok(idx) if !idx.entries.is_empty() => idx.entries[0].first_key,
_ => 0u64,
}
} else {
0u64
};
Ok((
footer.num_entries,
footer.min_timestamp,
file_size,
min_key,
max_key,
))
}
pub fn read_metadata<P: AsRef<Path>>(path: P) -> Result<(u64, u64, u64)> {
let path = path.as_ref();
let mut file = OpenOptions::new().read(true).open(path)?;
let footer = Self::read_footer(&mut file)?;
let file_size = file.metadata()?.len();
Ok((footer.num_entries, footer.min_timestamp, file_size))
}
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
let path = path.as_ref().to_path_buf();
let mut file = OpenOptions::new().read(true).open(&path)?;
let footer = Self::read_footer(&mut file)?;
let mmap = unsafe { Mmap::map(&file).ok() }.map(Arc::new);
let index = if let Some(ref mmap_data) = mmap {
let start = footer.index_offset as usize;
let end = start + footer.index_size as usize;
if end <= mmap_data.len() {
BlockIndex::deserialize(&mmap_data[start..end])?
} else {
file.seek(SeekFrom::Start(footer.index_offset))?;
let mut index_buf = vec![0u8; footer.index_size as usize];
file.read_exact(&mut index_buf)?;
BlockIndex::deserialize(&index_buf)?
}
} else {
file.seek(SeekFrom::Start(footer.index_offset))?;
let mut index_buf = vec![0u8; footer.index_size as usize];
file.read_exact(&mut index_buf)?;
BlockIndex::deserialize(&index_buf)?
};
let bloom = if let Some(ref mmap_data) = mmap {
let start = footer.bloom_offset as usize;
let end = start + footer.bloom_size as usize;
if end <= mmap_data.len() {
BloomFilter::from_bytes_full(&mmap_data[start..end])
.ok_or_else(|| StorageError::InvalidData("Invalid Bloom filter".into()))?
} else {
file.seek(SeekFrom::Start(footer.bloom_offset))?;
let mut bloom_buf = vec![0u8; footer.bloom_size as usize];
file.read_exact(&mut bloom_buf)?;
BloomFilter::from_bytes_full(&bloom_buf)
.ok_or_else(|| StorageError::InvalidData("Invalid Bloom filter".into()))?
}
} else {
file.seek(SeekFrom::Start(footer.bloom_offset))?;
let mut bloom_buf = vec![0u8; footer.bloom_size as usize];
file.read_exact(&mut bloom_buf)?;
BloomFilter::from_bytes_full(&bloom_buf)
.ok_or_else(|| StorageError::InvalidData("Invalid Bloom filter".into()))?
};
Ok(Self {
path,
mmap,
file,
index,
bloom,
footer,
})
}
pub fn bloom_filter(&self) -> &BloomFilter {
&self.bloom
}
pub fn shared_mmap(&self) -> Option<Arc<Mmap>> {
self.mmap.clone()
}
pub fn shared_index_entries(&self) -> Arc<Vec<BlockIndexEntry>> {
self.index.shared_entries()
}
pub fn read_block_zero_copy(&self, offset: u64, size: u32) -> Result<Vec<u8>> {
self.read_block(offset, size)
}
fn read_block(&self, offset: u64, size: u32) -> Result<Vec<u8>> {
if size < 4 {
return Err(crate::StorageError::InvalidData(format!(
"Block too small at offset {}: {} bytes",
offset, size
)));
}
let buf: &[u8] = if let Some(ref mmap) = self.mmap {
let end = offset as usize + size as usize;
if end > mmap.len() {
return Err(crate::StorageError::InvalidData(format!(
"Block extends beyond mmap: offset {} + size {} > {}",
offset,
size,
mmap.len()
)));
}
&mmap[offset as usize..end]
} else {
return Self::read_block_fallback(&self.path, offset, size);
};
let data_len = buf.len() - 4;
let data = &buf[..data_len];
let stored_crc = u32::from_le_bytes([
buf[data_len],
buf[data_len + 1],
buf[data_len + 2],
buf[data_len + 3],
]);
let computed_crc = crc32fast::hash(data);
if stored_crc != computed_crc {
return Err(crate::StorageError::InvalidData(format!(
"CRC32 mismatch at offset {}: expected {:08x}, got {:08x}. Data may be corrupted!",
offset, stored_crc, computed_crc
)));
}
Ok(data.to_vec())
}
fn read_block_fallback(path: &Path, offset: u64, size: u32) -> Result<Vec<u8>> {
let mut file = OpenOptions::new().read(true).open(path)?;
file.seek(SeekFrom::Start(offset))?;
let mut buf = vec![0u8; size as usize];
file.read_exact(&mut buf)?;
let data_len = buf.len() - 4;
let data = &buf[..data_len];
let stored_crc = u32::from_le_bytes([
buf[data_len],
buf[data_len + 1],
buf[data_len + 2],
buf[data_len + 3],
]);
let computed_crc = crc32fast::hash(data);
if stored_crc != computed_crc {
return Err(crate::StorageError::InvalidData(format!(
"CRC32 mismatch at offset {}: expected {:08x}, got {:08x}. Data may be corrupted!",
offset, stored_crc, computed_crc
)));
}
Ok(data.to_vec())
}
pub fn get(&self, key: Key) -> Result<Option<Value>> {
let key_bytes = key.to_be_bytes();
if !self.bloom.may_contain(&key_bytes) {
return Ok(None);
}
let block_entry = match self.index.find_block(&key_bytes) {
Some(entry) => entry,
None => return Ok(None),
};
let block_buf = self.read_block(block_entry.offset, block_entry.size)?;
Self::get_from_block_data(&block_buf, &key_bytes)
}
fn get_from_block_data(data: &[u8], key_bytes: &[u8]) -> Result<Option<Value>> {
if data.is_empty() {
return Ok(None);
}
let uncompressed: Vec<u8>;
let buf: &[u8] = match data[0] {
0 => &data[1..],
1 => {
let mut decoder = snap::raw::Decoder::new();
uncompressed = decoder.decompress_vec(&data[1..]).map_err(|e| {
crate::StorageError::Io(std::io::Error::other(format!(
"Snappy decompression failed: {}",
e
)))
})?;
&uncompressed
}
2 => {
uncompressed =
zstd::bulk::decompress(&data[1..], 4 * 1024 * 1024).map_err(|e| {
crate::StorageError::Io(std::io::Error::other(format!(
"Zstd decompression failed: {}",
e
)))
})?;
&uncompressed
}
_ => &data[1..],
};
if buf.len() < 4 {
return Ok(None);
}
let num_entries = u32::from_le_bytes([buf[0], buf[1], buf[2], buf[3]]) as usize;
let target_key = u64::from_be_bytes([
key_bytes[0],
key_bytes[1],
key_bytes[2],
key_bytes[3],
key_bytes[4],
key_bytes[5],
key_bytes[6],
key_bytes[7],
]);
let max_entries = buf.len() / 8;
let num_entries = num_entries.min(max_entries);
let mut key_offsets: Vec<(u64, usize)> = Vec::with_capacity(num_entries);
let mut off = 4usize;
for _ in 0..num_entries {
if off + 8 > buf.len() {
break;
}
let k = u64::from_be_bytes([
buf[off],
buf[off + 1],
buf[off + 2],
buf[off + 3],
buf[off + 4],
buf[off + 5],
buf[off + 6],
buf[off + 7],
]);
key_offsets.push((k, off));
off += 8;
off += 10;
if off > buf.len() {
break;
}
let value_type = buf[off - 1];
match value_type {
0 => {
if off + 4 > buf.len() {
break;
}
let vlen =
u32::from_le_bytes([buf[off], buf[off + 1], buf[off + 2], buf[off + 3]])
as usize;
off += 4 + vlen;
}
1 => {
off += 16;
}
_ => {
break;
}
}
}
let found = key_offsets.binary_search_by_key(&target_key, |(k, _)| *k);
match found {
Ok(idx) => {
let entry_off = key_offsets[idx].1;
let mut pos = entry_off + 8;
if pos + 10 > buf.len() {
return Ok(None);
}
let timestamp = u64::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
let deleted = buf[pos] != 0;
pos += 1;
let value_type = buf[pos];
pos += 1;
let value_data = match value_type {
0 => {
if pos + 4 > buf.len() {
return Ok(None);
}
let vlen = u32::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
]) as usize;
pos += 4;
if pos + vlen > buf.len() {
return Ok(None);
}
ValueData::Inline(std::sync::Arc::new(buf[pos..pos + vlen].to_vec()))
}
1 => {
if pos + 16 > buf.len() {
return Ok(None);
}
let file_id = u32::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
]);
pos += 4;
let blob_offset = u64::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
let size = u32::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
]);
ValueData::Blob(BlobRef {
file_id,
offset: blob_offset,
size,
})
}
_ => return Ok(None),
};
Ok(Some(Value {
data: value_data,
timestamp,
deleted,
}))
}
Err(_) => Ok(None),
}
}
pub fn batch_get(&self, keys: &[Key]) -> Result<Vec<Option<Value>>> {
let mut results = vec![None; keys.len()];
let key_bytes: Vec<[u8; 8]> = keys.iter().map(|k| k.to_be_bytes()).collect();
let key_refs: Vec<&[u8]> = key_bytes.iter().map(|b| b.as_slice()).collect();
let bloom_results = self.bloom.may_contain_batch(&key_refs);
let mut candidates: Vec<(usize, Key)> = Vec::new();
for (i, &may_exist) in bloom_results.iter().enumerate() {
if may_exist {
candidates.push((i, keys[i]));
}
}
for (idx, key) in candidates {
let key_bytes = key.to_be_bytes();
let block_entry = match self.index.find_block(&key_bytes) {
Some(entry) => entry,
None => continue,
};
let block_buf = self.read_block(block_entry.offset, block_entry.size)?;
let block = DataBlock::deserialize(&block_buf)?;
results[idx] = block.get(&key_bytes);
}
Ok(results)
}
pub fn scan(&self, start: Key, end: Key) -> Result<Vec<(Key, Value)>> {
let estimated_size = ((end - start) as usize).min(1000);
let mut results = Vec::with_capacity(estimated_size);
let start_bytes = start.to_be_bytes();
let start_idx = self.index.find_block_index(&start_bytes);
for i in start_idx..self.index.entries.len() {
let entry = &self.index.entries[i];
if entry.last_key < start {
continue;
}
if entry.first_key >= end {
break;
}
let block_buf = self.read_block(entry.offset, entry.size)?;
let block = DataBlock::deserialize(&block_buf)?;
for (k, v) in block.entries.iter() {
if k >= &start && k < &end {
results.push((*k, v.clone()));
}
if k >= &end {
return Ok(results);
}
}
}
Ok(results)
}
pub fn scan_all(&mut self) -> Result<Vec<(Key, Value)>> {
let estimated_size = (self.footer.num_entries as usize).min(10000);
let mut results = Vec::with_capacity(estimated_size);
let block_entries: Vec<_> = self
.index
.entries
.iter()
.map(|e| (e.offset, e.size))
.collect();
for (offset, size) in block_entries {
let block_buf = self.read_block(offset, size)?;
let block = DataBlock::deserialize(&block_buf)?;
for (k, v) in block.entries.iter() {
results.push((*k, v.clone()));
}
}
Ok(results)
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn iter(&mut self) -> Result<SSTableIterator> {
SSTableIterator::new(self)
}
pub fn stats(&self) -> SSTableStats {
SSTableStats {
num_entries: self.footer.num_entries,
file_size: std::fs::metadata(&self.path).map(|m| m.len()).unwrap_or(0),
num_blocks: self.index.entries.len(),
min_timestamp: self.footer.min_timestamp,
max_timestamp: self.footer.max_timestamp,
}
}
fn read_footer(file: &mut File) -> Result<Footer> {
let file_size = file.metadata()?.len();
if file_size < 64 {
return Err(StorageError::InvalidData("SSTable file too small".into()));
}
file.seek(SeekFrom::End(-64))?;
let mut buf = [0u8; 64];
file.read_exact(&mut buf)?;
let footer = Footer::deserialize(&buf)?;
let data_end = file_size - 64; let index_end = footer.index_offset + footer.index_size as u64;
let bloom_end = footer.bloom_offset + footer.bloom_size as u64;
if index_end > data_end || bloom_end > data_end {
return Err(StorageError::InvalidData(format!(
"SSTable footer points beyond file: index_end={}, bloom_end={}, data_end={}",
index_end, bloom_end, data_end
)));
}
Ok(footer)
}
}
pub struct SSTableBuilder {
writer: BufWriter<File>,
path: PathBuf,
current_block: DataBlock,
index: BlockIndex,
bloom: BloomFilter,
config: LSMConfig,
num_entries: u64,
min_timestamp: u64,
max_timestamp: u64,
min_key: Option<Key>,
max_key: Option<Key>,
offset: u64,
}
impl SSTableBuilder {
pub fn new<P: AsRef<Path>>(path: P, config: LSMConfig, estimated_keys: usize) -> Result<Self> {
let final_path = path.as_ref().to_path_buf();
let tmp_path = final_path.with_extension("sst.tmp");
if let Some(parent) = tmp_path.parent() {
std::fs::create_dir_all(parent)?;
}
let file = OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(&tmp_path)?;
Ok(Self {
writer: BufWriter::with_capacity(64 * 1024, file),
path: final_path,
current_block: DataBlock::new(),
index: BlockIndex::new(),
bloom: BloomFilter::new(estimated_keys, config.bloom_bits_per_key),
config,
num_entries: 0,
min_timestamp: u64::MAX,
max_timestamp: 0,
min_key: None,
max_key: None,
offset: 0,
})
}
pub fn add(&mut self, key: Key, value: Value) -> Result<()> {
self.bloom.insert(&key.to_be_bytes());
self.num_entries += 1;
self.min_timestamp = self.min_timestamp.min(value.timestamp);
self.max_timestamp = self.max_timestamp.max(value.timestamp);
if self.min_key.is_none() {
self.min_key = Some(key);
}
self.max_key = Some(key);
self.current_block.add(key, value)?;
if self.current_block.size() >= self.config.block_size {
self.flush_block()?;
}
Ok(())
}
pub fn finish(mut self) -> Result<super::compaction::SSTableMeta> {
if !self.current_block.is_empty() {
self.flush_block()?;
}
let min_key = self.min_key.unwrap_or_default();
let max_key = self.max_key.unwrap_or_default();
let index_offset = self.offset;
let index_data = self.index.serialize()?;
let index_size = index_data.len() as u32;
self.writer.write_all(&index_data)?;
self.offset += index_size as u64;
let bloom_offset = self.offset;
let bloom_data = self.bloom.to_bytes();
let bloom_size = bloom_data.len() as u32;
self.writer.write_all(&bloom_data)?;
self.offset += bloom_size as u64;
let footer = Footer {
magic: SSTABLE_MAGIC,
version: SSTABLE_VERSION,
index_offset,
index_size,
bloom_offset,
bloom_size,
num_entries: self.num_entries,
min_timestamp: if self.min_timestamp == u64::MAX {
0
} else {
self.min_timestamp
},
max_timestamp: self.max_timestamp,
max_key: self.max_key.unwrap_or(u64::MAX),
};
let footer_data = footer.serialize()?;
self.writer.write_all(&footer_data)?;
self.writer.flush()?;
self.writer.get_mut().sync_data()?;
let tmp_path = self.path.with_extension("sst.tmp");
if let Err(e) = std::fs::rename(&tmp_path, &self.path) {
return Err(crate::StorageError::Io(e));
}
let file_size = self.offset + footer_data.len() as u64;
Ok(super::compaction::SSTableMeta {
path: self.path,
size: file_size,
num_entries: self.num_entries,
min_key,
max_key,
min_timestamp: if self.min_timestamp == u64::MAX {
0
} else {
self.min_timestamp
},
max_timestamp: self.max_timestamp,
bloom_filter: Some(Arc::new(self.bloom)),
})
}
fn flush_block(&mut self) -> Result<()> {
if self.current_block.is_empty() {
return Ok(());
}
let first_key = self
.current_block
.entries
.first()
.map(|(k, _)| *k) .ok_or_else(|| StorageError::InvalidData("Empty block".into()))?;
let last_key = self
.current_block
.entries
.last()
.map(|(k, _)| *k)
.ok_or_else(|| StorageError::InvalidData("Empty block".into()))?;
let block_data = self.current_block.serialize_compressed(
self.config.enable_compression,
self.config.compression_algorithm,
)?;
let block_size = block_data.len() as u32;
self.index.entries.push(BlockIndexEntry {
first_key,
offset: self.offset,
size: block_size + 4,
last_key,
});
self.writer.write_all(&block_data)?;
let crc = crc32fast::hash(&block_data);
self.writer.write_all(&crc.to_le_bytes())?;
self.offset += block_size as u64 + 4;
self.current_block = DataBlock::new();
Ok(())
}
}
struct DataBlock {
entries: Vec<(Key, Value)>,
}
impl DataBlock {
fn new() -> Self {
Self {
entries: Vec::new(),
}
}
fn add(&mut self, key: Key, value: Value) -> Result<()> {
self.entries.push((key, value));
Ok(())
}
fn get(&self, key_bytes: &[u8]) -> Option<Value> {
if key_bytes.len() != 8 {
return None;
}
let key = u64::from_be_bytes([
key_bytes[0],
key_bytes[1],
key_bytes[2],
key_bytes[3],
key_bytes[4],
key_bytes[5],
key_bytes[6],
key_bytes[7],
]);
self.entries
.binary_search_by_key(&key, |(k, _)| *k)
.ok()
.map(|idx| self.entries[idx].1.clone())
}
fn size(&self) -> usize {
self.entries
.iter()
.map(|(_, v)| 8 + v.data.len() + 24) .sum()
}
fn is_empty(&self) -> bool {
self.entries.is_empty()
}
fn serialize(&self) -> Result<Vec<u8>> {
let mut buf = Vec::new();
buf.extend_from_slice(&(self.entries.len() as u32).to_le_bytes());
for (key, value) in &self.entries {
buf.extend_from_slice(&key.to_be_bytes());
buf.extend_from_slice(&value.timestamp.to_le_bytes());
buf.extend_from_slice(&[if value.deleted { 1 } else { 0 }]);
match &value.data {
ValueData::Inline(data) => {
buf.push(0);
buf.extend_from_slice(&(data.len() as u32).to_le_bytes());
buf.extend_from_slice(data);
}
ValueData::Blob(blob_ref) => {
buf.push(1);
buf.extend_from_slice(&blob_ref.file_id.to_le_bytes());
buf.extend_from_slice(&blob_ref.offset.to_le_bytes());
buf.extend_from_slice(&blob_ref.size.to_le_bytes());
}
}
}
Ok(buf)
}
fn serialize_compressed(
&self,
enable_compression: bool,
algorithm: CompressionAlgorithm,
) -> Result<Vec<u8>> {
let uncompressed = self.serialize()?;
if !enable_compression || uncompressed.len() < 1024 {
let mut result = vec![0u8]; result.extend_from_slice(&uncompressed);
return Ok(result);
}
match algorithm {
CompressionAlgorithm::Zstd => {
let level = 1; let compressed = zstd::bulk::compress(&uncompressed, level).map_err(|e| {
StorageError::Io(std::io::Error::other(format!(
"Zstd compression failed: {}",
e
)))
})?;
if compressed.len() < uncompressed.len() {
let mut result = vec![2u8]; result.extend_from_slice(&compressed);
Ok(result)
} else {
let mut result = vec![0u8];
result.extend_from_slice(&uncompressed);
Ok(result)
}
}
CompressionAlgorithm::Snappy => {
let mut encoder = snap::raw::Encoder::new();
let compressed = encoder.compress_vec(&uncompressed).map_err(|e| {
StorageError::Io(std::io::Error::other(format!(
"Snappy compression failed: {}",
e
)))
})?;
if compressed.len() < uncompressed.len() {
let mut result = vec![1u8]; result.extend_from_slice(&compressed);
Ok(result)
} else {
let mut result = vec![0u8];
result.extend_from_slice(&uncompressed);
Ok(result)
}
}
CompressionAlgorithm::None => {
let mut result = vec![0u8];
result.extend_from_slice(&uncompressed);
Ok(result)
}
}
}
fn deserialize(data: &[u8]) -> Result<Self> {
if data.is_empty() {
return Err(StorageError::InvalidData("Empty block data".into()));
}
let compression_flag = data[0];
let actual_data = &data[1..];
let uncompressed = match compression_flag {
0 => actual_data.to_vec(),
1 => {
let mut decoder = snap::raw::Decoder::new();
decoder.decompress_vec(actual_data).map_err(|e| {
StorageError::Io(std::io::Error::other(format!(
"Snappy decompression failed: {}",
e
)))
})?
}
2 => {
zstd::bulk::decompress(actual_data, 4 * 1024 * 1024).map_err(|e| {
StorageError::Io(std::io::Error::other(format!(
"Zstd decompression failed: {}",
e
)))
})?
}
_ => {
return Err(StorageError::InvalidData(format!(
"Unknown compression flag: {}",
compression_flag
)));
}
};
Self::deserialize_raw(&uncompressed)
}
fn deserialize_raw(data: &[u8]) -> Result<Self> {
let mut offset = 0;
if data.len() < 4 {
return Err(StorageError::InvalidData("Block too small".into()));
}
let num_entries = u32::from_le_bytes([data[0], data[1], data[2], data[3]]) as usize;
offset += 4;
let mut entries = Vec::with_capacity(num_entries);
for _ in 0..num_entries {
if offset + 8 > data.len() {
return Err(StorageError::InvalidData(
"Insufficient data for key".into(),
));
}
let key = u64::from_be_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
if offset + 10 > data.len() {
return Err(StorageError::InvalidData(
"Insufficient data for value metadata".into(),
));
}
let timestamp = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let deleted = data[offset] != 0;
offset += 1;
let value_type = data[offset];
offset += 1;
let value_data = match value_type {
0 => {
if offset + 4 > data.len() {
return Err(StorageError::InvalidData(
"Insufficient data for value length".into(),
));
}
let value_len = u32::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]) as usize;
offset += 4;
if offset + value_len > data.len() {
return Err(StorageError::InvalidData(format!(
"Value data exceeds block: need {} bytes, have {}",
value_len,
data.len() - offset
)));
}
let inline_data = data[offset..offset + value_len].to_vec();
offset += value_len;
ValueData::Inline(std::sync::Arc::new(inline_data))
}
1 => {
if offset + 16 > data.len() {
return Err(StorageError::InvalidData(
"Insufficient data for blob reference".into(),
));
}
let file_id = u32::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]);
offset += 4;
let blob_offset = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let size = u32::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]);
offset += 4;
ValueData::Blob(BlobRef {
file_id,
offset: blob_offset,
size,
})
}
_ => {
return Err(StorageError::InvalidData(format!(
"Unknown value type: {}",
value_type
)));
}
};
entries.push((
key,
Value {
data: value_data,
timestamp,
deleted,
},
));
}
Ok(Self { entries })
}
}
fn decompress_block(data: &[u8]) -> Result<Vec<u8>> {
if data.is_empty() {
return Err(StorageError::InvalidData("Empty block data".into()));
}
let compression_flag = data[0];
let actual_data = &data[1..];
match compression_flag {
0 => Ok(actual_data.to_vec()),
1 => {
let mut decoder = snap::raw::Decoder::new();
decoder.decompress_vec(actual_data).map_err(|e| {
StorageError::Io(std::io::Error::other(format!(
"Snappy decompression failed: {}",
e
)))
})
}
2 => zstd::bulk::decompress(actual_data, 4 * 1024 * 1024).map_err(|e| {
StorageError::Io(std::io::Error::other(format!(
"Zstd decompression failed: {}",
e
)))
}),
_ => Err(StorageError::InvalidData(format!(
"Unknown compression flag: {}",
compression_flag
))),
}
}
struct LazyEntryCursor {
data: Arc<Vec<u8>>,
num_entries: u32,
pos: usize,
entries_consumed: u32,
}
impl LazyEntryCursor {
fn new(data: Arc<Vec<u8>>) -> Result<Self> {
if data.len() < 4 {
return Err(StorageError::InvalidData(
"Block too small for header".into(),
));
}
let num_entries = u32::from_le_bytes([data[0], data[1], data[2], data[3]]);
let max_entries = (data.len() / 22) as u32;
let num_entries = num_entries.min(max_entries);
Ok(Self {
data,
num_entries,
pos: 4,
entries_consumed: 0,
})
}
fn peek_key(&self) -> Option<Key> {
if self.entries_consumed >= self.num_entries {
return None;
}
if self.pos + 8 > self.data.len() {
return None;
}
let key = u64::from_be_bytes([
self.data[self.pos],
self.data[self.pos + 1],
self.data[self.pos + 2],
self.data[self.pos + 3],
self.data[self.pos + 4],
self.data[self.pos + 5],
self.data[self.pos + 6],
self.data[self.pos + 7],
]);
Some(key)
}
fn skip_entry(&mut self) -> Result<()> {
if self.entries_consumed >= self.num_entries {
return Ok(());
}
let buf = &self.data;
let mut pos = self.pos;
pos += 8;
if pos + 10 > buf.len() {
return Err(StorageError::InvalidData("Truncated entry in skip".into()));
}
pos += 8;
pos += 1;
let value_type = buf[pos];
pos += 1;
match value_type {
0 => {
if pos + 4 > buf.len() {
return Err(StorageError::InvalidData(
"Truncated inline length in skip".into(),
));
}
let vlen = u32::from_le_bytes([buf[pos], buf[pos + 1], buf[pos + 2], buf[pos + 3]])
as usize;
pos += 4 + vlen;
}
1 => {
pos += 16;
}
_ => {
return Err(StorageError::InvalidData(format!(
"Unknown value type in skip: {}",
value_type
)));
}
}
self.pos = pos;
self.entries_consumed += 1;
Ok(())
}
fn next_entry(&mut self) -> Result<Option<(Key, Value)>> {
if self.entries_consumed >= self.num_entries {
return Ok(None);
}
let buf = &self.data;
let mut pos = self.pos;
if pos + 8 > buf.len() {
return Err(StorageError::InvalidData("Truncated key".into()));
}
let key = u64::from_be_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
if pos + 10 > buf.len() {
return Err(StorageError::InvalidData("Truncated value metadata".into()));
}
let timestamp = u64::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
let deleted = buf[pos] != 0;
pos += 1;
let value_type = buf[pos];
pos += 1;
let value_data = match value_type {
0 => {
if pos + 4 > buf.len() {
return Err(StorageError::InvalidData("Truncated inline length".into()));
}
let vlen = u32::from_le_bytes([buf[pos], buf[pos + 1], buf[pos + 2], buf[pos + 3]])
as usize;
pos += 4;
if pos + vlen > buf.len() {
return Err(StorageError::InvalidData(format!(
"Inline data exceeds block: need {} bytes at pos {}",
vlen, pos
)));
}
let inline_data = buf[pos..pos + vlen].to_vec();
pos += vlen;
ValueData::Inline(std::sync::Arc::new(inline_data))
}
1 => {
if pos + 16 > buf.len() {
return Err(StorageError::InvalidData("Truncated blob reference".into()));
}
let file_id =
u32::from_le_bytes([buf[pos], buf[pos + 1], buf[pos + 2], buf[pos + 3]]);
pos += 4;
let blob_offset = u64::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
let size = u32::from_le_bytes([buf[pos], buf[pos + 1], buf[pos + 2], buf[pos + 3]]);
pos += 4;
ValueData::Blob(BlobRef {
file_id,
offset: blob_offset,
size,
})
}
_ => {
return Err(StorageError::InvalidData(format!(
"Unknown value type: {}",
value_type
)));
}
};
self.pos = pos;
self.entries_consumed += 1;
Ok(Some((
key,
Value {
data: value_data,
timestamp,
deleted,
},
)))
}
#[allow(dead_code)]
fn next_entry_raw(&mut self) -> Result<Option<(Key, u64, bool, &[u8])>> {
if self.entries_consumed >= self.num_entries {
return Ok(None);
}
let buf = &self.data;
let mut pos = self.pos;
if pos + 8 > buf.len() {
return Err(StorageError::InvalidData("Truncated key".into()));
}
let key = u64::from_be_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
if pos + 10 > buf.len() {
return Err(StorageError::InvalidData("Truncated value metadata".into()));
}
let timestamp = u64::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
let deleted = buf[pos] != 0;
pos += 1;
let value_type = buf[pos];
pos += 1;
let value_start = match value_type {
0 => {
if pos + 4 > buf.len() {
return Err(StorageError::InvalidData("Truncated inline len".into()));
}
let vlen = u32::from_le_bytes([buf[pos], buf[pos + 1], buf[pos + 2], buf[pos + 3]])
as usize;
pos += 4;
if pos + vlen > buf.len() {
return Err(StorageError::InvalidData(format!(
"Inline data exceeds block: need {} bytes at pos {}",
vlen, pos
)));
}
let start = pos;
pos += vlen;
start
}
1 => {
pos += 16;
pos
}
_ => {
return Err(StorageError::InvalidData(format!(
"Unknown value type: {}",
value_type
)))
}
};
let value_len = pos - value_start;
self.pos = pos;
self.entries_consumed += 1;
let slice = &self.data[value_start..value_start + value_len];
Ok(Some((key, timestamp, deleted, slice)))
}
fn next_entry_arc(&mut self) -> Result<Option<(Key, u64, bool, Arc<Vec<u8>>, usize, usize)>> {
if self.entries_consumed >= self.num_entries {
return Ok(None);
}
let buf = &self.data;
let mut pos = self.pos;
if pos + 8 > buf.len() {
return Err(StorageError::InvalidData("Truncated key".into()));
}
let key = u64::from_be_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
if pos + 10 > buf.len() {
return Err(StorageError::InvalidData("Truncated value metadata".into()));
}
let timestamp = u64::from_le_bytes([
buf[pos],
buf[pos + 1],
buf[pos + 2],
buf[pos + 3],
buf[pos + 4],
buf[pos + 5],
buf[pos + 6],
buf[pos + 7],
]);
pos += 8;
let deleted = buf[pos] != 0;
pos += 1;
let value_type = buf[pos];
pos += 1;
let value_start = match value_type {
0 => {
if pos + 4 > buf.len() {
return Err(StorageError::InvalidData("Truncated inline length".into()));
}
let vlen = u32::from_le_bytes([buf[pos], buf[pos + 1], buf[pos + 2], buf[pos + 3]])
as usize;
pos += 4;
if pos + vlen > buf.len() {
return Err(StorageError::InvalidData(format!(
"Inline data exceeds block: need {} bytes at pos {}",
vlen, pos
)));
}
let start = pos;
pos += vlen;
start
}
1 => {
pos += 16;
pos }
_ => {
return Err(StorageError::InvalidData(format!(
"Unknown value type: {}",
value_type
)));
}
};
let value_len = pos - value_start;
self.pos = pos;
self.entries_consumed += 1;
let block = Arc::clone(&self.data);
Ok(Some((
key,
timestamp,
deleted,
block,
value_start,
value_len,
)))
}
}
impl BlockIndex {
fn new() -> Self {
Self {
entries: Vec::new(),
}
}
fn find_block(&self, key_bytes: &[u8]) -> Option<&BlockIndexEntry> {
if key_bytes.len() != 8 {
return None;
}
let key = u64::from_be_bytes([
key_bytes[0],
key_bytes[1],
key_bytes[2],
key_bytes[3],
key_bytes[4],
key_bytes[5],
key_bytes[6],
key_bytes[7],
]);
match self.entries.binary_search_by(|e| e.first_key.cmp(&key)) {
Ok(idx) => Some(&self.entries[idx]),
Err(idx) => {
if idx == 0 {
None
} else {
Some(&self.entries[idx - 1])
}
}
}
}
fn find_block_index(&self, key_bytes: &[u8]) -> usize {
if key_bytes.len() != 8 {
return 0;
}
let key = u64::from_be_bytes([
key_bytes[0],
key_bytes[1],
key_bytes[2],
key_bytes[3],
key_bytes[4],
key_bytes[5],
key_bytes[6],
key_bytes[7],
]);
match self.entries.binary_search_by(|e| e.first_key.cmp(&key)) {
Ok(idx) => idx,
Err(idx) => {
if idx == 0 {
0
} else {
idx - 1
}
}
}
}
fn serialize(&self) -> Result<Vec<u8>> {
let mut buf = Vec::new();
buf.extend_from_slice(&(self.entries.len() as u32).to_le_bytes());
for entry in &self.entries {
buf.extend_from_slice(&entry.first_key.to_be_bytes());
buf.extend_from_slice(&entry.offset.to_le_bytes());
buf.extend_from_slice(&entry.size.to_le_bytes());
buf.extend_from_slice(&entry.last_key.to_be_bytes());
}
Ok(buf)
}
fn deserialize(data: &[u8]) -> Result<Self> {
if data.len() < 4 {
return Err(StorageError::InvalidData("BlockIndex too small".into()));
}
let num_entries = u32::from_le_bytes([data[0], data[1], data[2], data[3]]) as usize;
let new_size = 4 + num_entries * 28;
let old_size = 4 + num_entries * 20;
if data.len() >= new_size {
Self::deserialize_v2(data, num_entries)
} else if data.len() >= old_size {
Self::deserialize_v1(data, num_entries)
} else {
Err(StorageError::InvalidData(format!(
"BlockIndex truncated: expected {} (or {} legacy) bytes, got {}",
new_size,
old_size,
data.len()
)))
}
}
fn deserialize_v2(data: &[u8], num_entries: usize) -> Result<Self> {
let mut off = 4usize;
let mut entries = Vec::with_capacity(num_entries);
for _ in 0..num_entries {
let first_key = u64::from_be_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
data[off + 4],
data[off + 5],
data[off + 6],
data[off + 7],
]);
off += 8;
let block_offset = u64::from_le_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
data[off + 4],
data[off + 5],
data[off + 6],
data[off + 7],
]);
off += 8;
let block_size =
u32::from_le_bytes([data[off], data[off + 1], data[off + 2], data[off + 3]]);
off += 4;
let last_key = u64::from_be_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
data[off + 4],
data[off + 5],
data[off + 6],
data[off + 7],
]);
off += 8;
entries.push(BlockIndexEntry {
first_key,
offset: block_offset,
size: block_size,
last_key,
});
}
Ok(Self { entries })
}
fn deserialize_v1(data: &[u8], num_entries: usize) -> Result<Self> {
let mut off = 4usize;
let mut entries = Vec::with_capacity(num_entries);
for i in 0..num_entries {
let first_key = u64::from_be_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
data[off + 4],
data[off + 5],
data[off + 6],
data[off + 7],
]);
off += 8;
let block_offset = u64::from_le_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
data[off + 4],
data[off + 5],
data[off + 6],
data[off + 7],
]);
off += 8;
let block_size =
u32::from_le_bytes([data[off], data[off + 1], data[off + 2], data[off + 3]]);
off += 4;
let last_key = if i + 1 < num_entries {
let next_first = u64::from_be_bytes([
data[off],
data[off + 1],
data[off + 2],
data[off + 3],
data[off + 4],
data[off + 5],
data[off + 6],
data[off + 7],
]);
next_first.saturating_sub(1)
} else {
u64::MAX };
entries.push(BlockIndexEntry {
first_key,
offset: block_offset,
size: block_size,
last_key,
});
}
Ok(Self { entries })
}
fn shared_entries(&self) -> Arc<Vec<BlockIndexEntry>> {
Arc::new(self.entries.clone())
}
}
impl Footer {
fn serialize(&self) -> Result<Vec<u8>> {
let mut buf = vec![0u8; 64];
let mut offset = 0;
buf[offset..offset + 4].copy_from_slice(&self.magic.to_le_bytes());
offset += 4;
buf[offset..offset + 4].copy_from_slice(&self.version.to_le_bytes());
offset += 4;
buf[offset..offset + 8].copy_from_slice(&self.index_offset.to_le_bytes());
offset += 8;
buf[offset..offset + 4].copy_from_slice(&self.index_size.to_le_bytes());
offset += 4;
buf[offset..offset + 8].copy_from_slice(&self.bloom_offset.to_le_bytes());
offset += 8;
buf[offset..offset + 4].copy_from_slice(&self.bloom_size.to_le_bytes());
offset += 4;
buf[offset..offset + 8].copy_from_slice(&self.num_entries.to_le_bytes());
offset += 8;
buf[offset..offset + 8].copy_from_slice(&self.min_timestamp.to_le_bytes());
offset += 8;
buf[offset..offset + 8].copy_from_slice(&self.max_timestamp.to_le_bytes());
offset += 8;
buf[offset..offset + 8].copy_from_slice(&self.max_key.to_le_bytes());
Ok(buf)
}
fn deserialize(data: &[u8]) -> Result<Self> {
if data.len() < 64 {
return Err(StorageError::InvalidData("Footer too small".into()));
}
let mut offset = 0;
let magic = u32::from_le_bytes([data[0], data[1], data[2], data[3]]);
offset += 4;
if magic != SSTABLE_MAGIC {
return Err(StorageError::InvalidData("Invalid SSTable magic".into()));
}
let version = u32::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]);
offset += 4;
let index_offset = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let index_size = u32::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]);
offset += 4;
let bloom_offset = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let bloom_size = u32::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
]);
offset += 4;
let num_entries = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let min_timestamp = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let max_timestamp = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
offset += 8;
let max_key = u64::from_le_bytes([
data[offset],
data[offset + 1],
data[offset + 2],
data[offset + 3],
data[offset + 4],
data[offset + 5],
data[offset + 6],
data[offset + 7],
]);
let max_key = if max_key == 0 && num_entries > 0 {
u64::MAX
} else {
max_key
};
Ok(Self {
magic,
version,
index_offset,
index_size,
bloom_offset,
bloom_size,
num_entries,
min_timestamp,
max_timestamp,
max_key,
})
}
}
#[derive(Debug, Clone)]
pub struct SSTableStats {
pub num_entries: u64,
pub file_size: u64,
pub num_blocks: usize,
pub min_timestamp: u64,
pub max_timestamp: u64,
}
pub struct SSTableIterator {
mmap: Option<Arc<Mmap>>,
index_entries: Arc<Vec<BlockIndexEntry>>,
file: Option<BufReader<File>>,
#[allow(dead_code)]
path: PathBuf,
current_block_idx: usize,
current_cursor: Option<LazyEntryCursor>,
start_key: Option<Key>,
end_key: Option<Key>,
verify_crc: bool,
}
impl SSTableIterator {
fn new(sstable: &SSTable) -> Result<Self> {
Self::with_range(sstable, None, None)
}
pub fn with_range(
sstable: &SSTable,
start_key: Option<Key>,
end_key: Option<Key>,
) -> Result<Self> {
let mmap = sstable.shared_mmap();
let index_entries = sstable.shared_index_entries();
let start_block_idx = if let Some(start) = start_key {
let start_bytes = start.to_be_bytes();
sstable.index.find_block_index(&start_bytes)
} else {
0
};
let (file, path) = if mmap.is_none() {
let file = BufReader::new(File::open(&sstable.path).map_err(StorageError::Io)?);
(Some(file), sstable.path.clone())
} else {
(None, sstable.path.clone())
};
Ok(Self {
mmap,
index_entries,
file,
path,
current_block_idx: start_block_idx,
current_cursor: None,
start_key,
end_key,
verify_crc: true, })
}
fn load_next_block(&mut self) -> Result<bool> {
loop {
if self.current_block_idx >= self.index_entries.len() {
return Ok(false); }
let entry = &self.index_entries[self.current_block_idx];
if let Some(start) = self.start_key {
if entry.last_key < start {
self.current_block_idx += 1;
continue; }
}
if let Some(end) = self.end_key {
if entry.first_key >= end {
return Ok(false); }
}
break; }
let offset = self.index_entries[self.current_block_idx].offset;
let size = self.index_entries[self.current_block_idx].size;
if size < 4 {
return Err(crate::StorageError::InvalidData(
"Block too small for CRC".into(),
));
}
let block_bytes: Vec<u8> = if let Some(ref mmap) = self.mmap {
let start = offset as usize;
let end = start + size as usize;
if end > mmap.len() {
return Err(crate::StorageError::InvalidData(format!(
"Block extends beyond mmap: offset {} + size {} > {}",
offset,
size,
mmap.len()
)));
}
let data_len = size as usize - 4;
if self.verify_crc {
let stored_crc = u32::from_le_bytes([
mmap[start + data_len],
mmap[start + data_len + 1],
mmap[start + data_len + 2],
mmap[start + data_len + 3],
]);
let computed_crc = crc32fast::hash(&mmap[start..start + data_len]);
if stored_crc != computed_crc {
return Err(crate::StorageError::InvalidData(
format!("CRC32 mismatch in iterator block at offset {}: expected {:08x}, got {:08x}", offset, stored_crc, computed_crc)
));
}
}
decompress_block(&mmap[start..start + data_len])?
} else {
let file = self.file.as_mut().unwrap();
file.seek(SeekFrom::Start(offset))?;
let mut buf = vec![0u8; size as usize];
file.read_exact(&mut buf)?;
let data_len = buf.len() - 4;
if self.verify_crc {
let stored_crc = u32::from_le_bytes([
buf[data_len],
buf[data_len + 1],
buf[data_len + 2],
buf[data_len + 3],
]);
let computed_crc = crc32fast::hash(&buf[..data_len]);
if stored_crc != computed_crc {
return Err(crate::StorageError::InvalidData(
format!("CRC32 mismatch in iterator block at offset {}: expected {:08x}, got {:08x}", offset, stored_crc, computed_crc)
));
}
}
decompress_block(&buf[..data_len])?
};
self.current_cursor = Some(LazyEntryCursor::new(Arc::new(block_bytes))?);
self.current_block_idx += 1;
Ok(true)
}
pub fn set_verify_crc(&mut self, verify: bool) {
self.verify_crc = verify;
}
pub fn next_raw(&mut self) -> Option<(Key, u64, bool, ValueBytes)> {
loop {
if let Some(ref mut cursor) = self.current_cursor {
if let Some(key) = cursor.peek_key() {
if let Some(start) = self.start_key {
if key < start {
if cursor.skip_entry().is_err() {
return None;
}
continue;
}
}
if let Some(end) = self.end_key {
if key >= end {
self.current_cursor = None;
return None;
}
}
}
match cursor.next_entry_arc() {
Ok(Some((key, ts, deleted, block, start, len))) => {
return Some((key, ts, deleted, ValueBytes { block, start, len }));
}
Ok(None) => {
self.current_cursor = None;
}
Err(e) => {
eprintln!("[MoteDB] SSTableIterator::next_raw: parse error: {}", e);
return None;
}
}
}
match self.load_next_block() {
Ok(true) => continue,
Ok(false) => return None,
Err(e) => {
eprintln!(
"[MoteDB] SSTableIterator::next_raw: block load error: {}",
e
);
return None;
}
}
}
}
}
impl Iterator for SSTableIterator {
type Item = (Key, Value);
fn next(&mut self) -> Option<Self::Item> {
loop {
if let Some(ref mut cursor) = self.current_cursor {
if let Some(key) = cursor.peek_key() {
if let Some(start) = self.start_key {
if key < start {
if cursor.skip_entry().is_err() {
return None;
}
continue;
}
}
if let Some(end) = self.end_key {
if key >= end {
self.current_cursor = None;
return None;
}
}
}
match cursor.next_entry() {
Ok(Some((key, value))) => return Some((key, value)),
Ok(None) => {
self.current_cursor = None; }
Err(e) => {
eprintln!("[MoteDB] SSTableIterator: failed to parse entry: {}", e);
return None;
}
}
}
match self.load_next_block() {
Ok(true) => continue,
Ok(false) => return None,
Err(e) => {
eprintln!("[MoteDB] SSTableIterator: failed to load block: {}", e);
return None;
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[test]
fn test_sstable_basic() {
let temp_dir = TempDir::new().unwrap();
let path = temp_dir.path().join("test.sst");
{
let mut builder = SSTableBuilder::new(&path, LSMConfig::default(), 100).unwrap();
for i in 0..100 {
let key = i as u64; let value = Value::new(format!("value_{}", i).into_bytes(), i as u64);
builder.add(key, value).unwrap();
}
builder.finish().unwrap();
}
{
let sst = SSTable::open(&path).unwrap();
let key = 50u64; let value = sst.get(key).unwrap().unwrap();
assert_eq!(
value.data,
ValueData::Inline(std::sync::Arc::new(b"value_50".to_vec()))
);
assert_eq!(value.timestamp, 50);
let result = sst.get(999u64).unwrap();
assert!(result.is_none());
}
}
}