Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ api = []

[dependencies]
bloomfilter = "3.0"
crc32fast = "1.4"
bincode = "1.3"
lz4_flex = "0.11"
serde = { version = "1.0", features = ["derive"] }
Expand Down
153 changes: 129 additions & 24 deletions src/storage/block.rs

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟑 current_size() does not account for the new CRC32 checksum bytes, causing add() to overfill blocks

Block::current_size() (src/storage/block.rs:40-42) computes data.len() + metadata_size(offsets.len()) but does not include the 4-byte CRC32 checksum that encode() now appends. This means add() at line 47-50 compares against block_size without accounting for the CRC32 overhead, allowing blocks to be filled 4 bytes beyond the intended block_size limit when encoded. For small block sizes (e.g., block_size = 128 used in tests), this is a ~3% overshoot. It also causes encode() at line 70 to under-allocate Vec capacity by 4 bytes, triggering an unnecessary reallocation.

(Refers to lines 40-42)

Prompt for agents
The `current_size()` method needs to account for the 4-byte CRC32 checksum that `encode()` now appends. This affects two call sites: (1) `add()` uses it to check if the block is full β€” without accounting for CRC32, blocks can be slightly overfilled relative to `block_size`. (2) `encode()` uses it for `Vec::with_capacity` β€” causing an unnecessary reallocation. The fix should add `U32_SIZE` (for CRC32) to `metadata_size()` or `current_size()`. Note that `metadata_size` is called with `self.offsets.len()` (current number of offsets), and in `add()` the check also adds `new_offset_size` separately, so changing `metadata_size` should be safe. Alternatively, add a constant for the CRC32 footer size.
Open in Devin Review

Was this helpful? React with πŸ‘ or πŸ‘Ž to provide feedback.

Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
use crate::infra::config::StorageConfig;
use crate::infra::{config::StorageConfig, error::LsmError};
use crc32fast::Hasher;
use std::mem::size_of;

pub const BLOCK_SIZE: usize = 4096;
Expand Down Expand Up @@ -76,48 +77,76 @@ impl Block {
let num_elements = self.offsets.len() as u32;
encoded.extend_from_slice(&num_elements.to_le_bytes());

// Calculate and append CRC32 checksum (Little Endian)
let mut hasher = Hasher::new();
hasher.update(&encoded);
let checksum = hasher.finalize();
encoded.extend_from_slice(&checksum.to_le_bytes());

encoded
}

pub fn decode(data: &[u8]) -> Self {
pub fn decode(data: &[u8]) -> std::result::Result<Self, LsmError> {
if data.len() < U32_SIZE {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

πŸ”΄ Insufficient minimum length check in Block::decode causes panic on short input

Block::decode at src/storage/block.rs:90 checks data.len() < U32_SIZE (4 bytes), but the minimum valid encoded block requires 2 * U32_SIZE (8 bytes): 4 for the num_elements field + 4 for the CRC32 checksum. If data is 4–7 bytes and the CRC32 happens to match, line 120 computes data_without_checksum.len() - U32_SIZE which underflows usize, causing a panic.

This is deterministically triggerable: the input [0, 0, 0, 0] always panics because CRC32 of empty data is 0, so the checksum verification passes, then 0 - 4 underflows. This defeats the purpose of the CRC32 integrity check being addedβ€”corrupted or truncated block data that happens to be exactly 4–7 bytes crashes the process instead of returning a clean error.

Suggested change
if data.len() < U32_SIZE {
if data.len() < 2 * U32_SIZE {
Open in Devin Review

Was this helpful? React with πŸ‘ or πŸ‘Ž to provide feedback.

return Self {
data: Vec::new(),
offsets: Vec::new(),
block_size: BLOCK_SIZE,
};
return Err(LsmError::CorruptedData(
"Data too short to contain checksum".to_string(),
));
}

// Read stored checksum (last 4 bytes)
let checksum_start = data.len() - U32_SIZE;
let stored_checksum = u32::from_le_bytes([
data[checksum_start],
data[checksum_start + 1],
data[checksum_start + 2],
data[checksum_start + 3],
]);

// Extract data without checksum for verification
let data_without_checksum = &data[..checksum_start];

// Calculate actual checksum
let mut hasher = Hasher::new();
hasher.update(data_without_checksum);
let calculated_checksum = hasher.finalize();

// Verify checksum
if stored_checksum != calculated_checksum {
return Err(LsmError::CorruptedData(
"CRC32 checksum mismatch: data corruption detected".to_string(),
));
}

let num_elements_start = data.len() - U32_SIZE;
let num_elements_start = data_without_checksum.len() - U32_SIZE;
let num_elements = u32::from_le_bytes([
data[num_elements_start],
data[num_elements_start + 1],
data[num_elements_start + 2],
data[num_elements_start + 3],
data_without_checksum[num_elements_start],
data_without_checksum[num_elements_start + 1],
data_without_checksum[num_elements_start + 2],
data_without_checksum[num_elements_start + 3],
]) as usize;

let offsets_start = data.len() - U32_SIZE - (num_elements * U32_SIZE);
let records_data = data[..offsets_start].to_vec();
let offsets_start = data_without_checksum.len() - U32_SIZE - (num_elements * U32_SIZE);
let records_data = data_without_checksum[..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 = u32::from_le_bytes([
data[offset_pos],
data[offset_pos + 1],
data[offset_pos + 2],
data[offset_pos + 3],
data_without_checksum[offset_pos],
data_without_checksum[offset_pos + 1],
data_without_checksum[offset_pos + 2],
data_without_checksum[offset_pos + 3],
]);
offsets.push(offset);
offset_pos += U32_SIZE;
}

Self {
Ok(Self {
data: records_data,
offsets,
block_size: BLOCK_SIZE,
}
})
}

pub fn len(&self) -> usize {
Expand Down Expand Up @@ -221,7 +250,7 @@ mod tests {

// Verify integrity
let encoded = block.encode();
let decoded = Block::decode(&encoded);
let decoded = Block::decode(&encoded).unwrap();

assert_eq!(decoded.len(), block.len());
assert_eq!(decoded.offsets.len(), block.offsets.len());
Expand All @@ -245,7 +274,7 @@ mod tests {
fn test_encode_decode_empty_block() {
let block = Block::new(BLOCK_SIZE);
let encoded = block.encode();
let decoded = Block::decode(&encoded);
let decoded = Block::decode(&encoded).unwrap();
assert_eq!(decoded.len(), 0);
assert!(decoded.is_empty());
}
Expand All @@ -255,7 +284,7 @@ mod tests {
let mut block = Block::new(BLOCK_SIZE);
block.add(b"key1", b"value1");
let encoded = block.encode();
let decoded = Block::decode(&encoded);
let decoded = Block::decode(&encoded).unwrap();
assert_eq!(decoded.len(), 1);
assert_eq!(decoded.data_size(), block.data_size());
assert_eq!(decoded.data, block.data);
Expand All @@ -278,9 +307,85 @@ mod tests {
}

let encoded = block.encode();
let decoded = Block::decode(&encoded);
let decoded = Block::decode(&encoded).unwrap();
assert_eq!(decoded.len(), entries.len());
assert_eq!(decoded.data, block.data);
assert_eq!(decoded.offsets, block.offsets);
}

#[test]
fn test_crc32_corruption_detected() {
let mut block = Block::new(BLOCK_SIZE);
block.add(b"test_key", b"test_value");

let encoded = block.encode();
let mut corrupted = encoded.clone();

// Corrupt a byte in the data section (not the checksum)
corrupted[10] ^= 0xFF;

// Verify that decode returns a corruption error
let result = Block::decode(&corrupted);
assert!(result.is_err());

let err = result.unwrap_err();
assert!(matches!(err, LsmError::CorruptedData(_)));
assert!(err.to_string().contains("CRC32"));
}

#[test]
fn test_crc32_valid_checksum() {
let mut block = Block::new(BLOCK_SIZE);
block.add(b"key1", b"value1");
block.add(b"key2", b"value2");

let encoded = block.encode();
let decoded = Block::decode(&encoded).unwrap();

assert_eq!(decoded.len(), 2);
assert_eq!(decoded.data, block.data);
assert_eq!(decoded.offsets, block.offsets);
}

#[test]
fn test_crc32_checksum_mismatch_single_bit_flip() {
let mut block = Block::new(BLOCK_SIZE);
for i in 0..50 {
let key = format!("key_{:03}", i);
let value = format!("value_{:03}", i);
assert!(block.add(key.as_bytes(), value.as_bytes()));
}

let encoded = block.encode();
let corrupted = corrupt_byte(&encoded, 100);

let result = Block::decode(&corrupted);
assert!(result.is_err());

let err = result.unwrap_err();
assert!(matches!(err, LsmError::CorruptedData(_)));
assert!(err.to_string().contains("mismatch"));
}

#[test]
fn test_crc32_truncated_file_detected() {
let mut block = Block::new(BLOCK_SIZE);
block.add(b"short_key", b"short_value");

let encoded = block.encode();
// Truncate by removing the checksum bytes
let truncated = &encoded[..encoded.len() - U32_SIZE];

let result = Block::decode(truncated);
assert!(result.is_err());
}

/// Helper to corrupt a specific byte in the data
fn corrupt_byte(data: &[u8], pos: usize) -> Vec<u8> {
let mut corrupted = data.to_vec();
if pos < corrupted.len() {
corrupted[pos] ^= 0xFF;
}
corrupted
}
}
98 changes: 96 additions & 2 deletions src/storage/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,7 @@ impl SstableReader {
let block_data = self.read_block(&block_meta)?;

// Deserialize block (no lock needed)
let block = Block::decode(&block_data);
let block = Block::decode(&block_data)?;

// Linear scan within the block to find the key (no lock needed)
Self::search_in_block(&block, key.as_bytes())
Expand Down Expand Up @@ -188,7 +188,7 @@ impl SstableReader {

for block_meta in &blocks {
let block_data = self.read_block(block_meta)?;
let block = Block::decode(&block_data);
let block = Block::decode(&block_data)?;

// Access block data through pub(crate) fields
for &offset in &block.offsets {
Expand Down Expand Up @@ -861,4 +861,98 @@ mod tests {
handle.join().unwrap();
}
}

#[test]
fn test_sstable_data_corruption_detected() {
use std::fs::File;
use std::io::Read;
use std::io::Write;

let dir = tempdir().unwrap();
let path = dir.path().join("corruption_test.sst");
// Use a small block size to create many blocks in a large file
let config = StorageConfig {
block_size: 128, // Very small block size
..Default::default()
};
let _cache = create_test_cache(&config);

// Write an SSTable with enough data to create many blocks
let mut builder = SstableBuilder::new(path.clone(), config.clone(), 12345).unwrap();
for i in 0..100 {
let key = format!("key_{:03}", i);
let value = format!("value_{:03}", i);
builder
.add(key.as_bytes(), &create_test_record(&key, value.as_bytes()))
.unwrap();
}
builder.finish().unwrap();

// Open the SSTable to get block metadata
let reader = SstableReader::open(path.clone(), config.clone(), _cache);
let reader = match reader {
Ok(r) => r,
Err(e) => {
// If we can't even open the original SSTable, something is wrong
panic!("Failed to open original SSTable: {:?}", e);
}
};

let metadata = reader.metadata();

// Corrupt the last block's data
// Get the last block's offset and size from metadata
let last_block = metadata.blocks.last().unwrap();
let last_block_start = last_block.offset as usize;
let last_block_size = last_block.size as usize;

// Read the file
let mut file = File::open(&path).unwrap();
let mut original_data = Vec::new();
file.read_to_end(&mut original_data).unwrap();

// Corrupt a byte in the last compressed block
let corrupt_offset = last_block_start + last_block_size / 2;

if corrupt_offset < original_data.len() - 8 {
// Keep footer intact
original_data[corrupt_offset] ^= 0xFF;

// Write the corrupted data back
let mut file = File::create(&path).unwrap();
file.write_all(&original_data).unwrap();
drop(file);

// Re-open reader with a fresh cache
let fresh_cache = create_test_cache(&config);
let reader = SstableReader::open(path, config, fresh_cache).unwrap();

// Try to read the last block which should trigger the CRC32 check
// Use the key from the last block (we'll use a key we know exists)
// For simplicity, get the last key by reading metadata.max_key
let result = reader.get(&String::from_utf8_lossy(&metadata.max_key));

// The corruption should cause CRC32 verification to fail
assert!(
result.is_err(),
"Should fail to read from corrupted SSTable, got: {:?}",
result
);

let err = result.unwrap_err();
assert!(
matches!(err, LsmError::CorruptedData(_)),
"Expected CorruptedData error, got: {:?}",
err
);
assert!(
err.to_string().contains("CRC32"),
"Error should mention CRC32"
);
} else {
// Fallback: if corruption position is invalid, just verify the block tests still work
// This shouldn't happen in practice
panic!("Corruption position calculation failed");
}
}
}
2 changes: 1 addition & 1 deletion src/storage/sst_iterator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ impl SstableIterator {
}
let block_meta = meta.blocks[block_idx].clone();
let raw = self.reader.read_block(&block_meta)?;
self.current_block = Some(Block::decode(&raw));
self.current_block = Some(Block::decode(&raw)?);
self.block_index = block_idx;
self.offset_index = 0;
Ok(())
Expand Down
Loading