diff --git a/crates/logex-storage/src/native/segment.rs b/crates/logex-storage/src/native/segment.rs index ab8553da..a4986a42 100644 --- a/crates/logex-storage/src/native/segment.rs +++ b/crates/logex-storage/src/native/segment.rs @@ -1,6 +1,6 @@ use std::fs; use std::fs::File; -use std::io::{BufWriter, Write}; +use std::io::{BufWriter, Read, Seek, SeekFrom, Write}; use std::ops::Range; use std::path::{Path, PathBuf}; use std::thread; @@ -8,7 +8,7 @@ use std::thread; use alloy_primitives::{Address, B256}; use logex_types::LogRow; -use crate::column::{ColumnFile, NullBitmap}; +use crate::column::{ColumnFile, ColumnFileHeader, NullBitmap}; use crate::page::{ PageIndexEntry, encode_fixed_width_page, encode_u8_page, encode_u32_page, encode_u64_page, encode_var_bytes_page, write_page_index, @@ -23,6 +23,30 @@ use super::catalog::{ const DEFAULT_PAGE_ROWS: u32 = 16_384; const RECOMPACTED_COLUMNS_DIR: &str = "columns_profile_v2"; +const RAW_FIXED_COLUMNS: &[(&str, u64)] = &[ + ("address.col", 20), + ("block_number.col", 8), + ("block_hash.col", 32), + ("timestamp.col", 8), + ("tx_hash.col", 32), + ("tx_index.col", 4), + ("log_index.col", 4), + ("data_len.col", 4), + ("source.col", 1), + ("topic0.col", 32), + ("topic1.col", 32), + ("topic2.col", 32), + ("topic3.col", 32), +]; +const RAW_NULL_BITMAP_COLUMNS: &[&str] = + &["topic0.null", "topic1.null", "topic2.null", "topic3.null"]; +const RAW_BITMAP_COLUMNS: &[&str] = &[ + "topic0.null", + "topic1.null", + "topic2.null", + "topic3.null", + "canonical.bitmap", +]; pub(crate) fn append_rows( segment_dir: &Path, @@ -306,6 +330,7 @@ pub(crate) fn compact_segment( let segment_dir = paths.segment_dir(descriptor.id); fs::create_dir_all(segment_dir.join("columns"))?; + verify_raw_segment_files_complete(descriptor, &segment_dir)?; let columns = vec![ compact_address_column(&segment_dir)?, @@ -853,26 +878,12 @@ where } fn remove_raw_hot_files(segment_dir: &Path) -> std::io::Result<()> { - for name in [ - "address.col", - "block_number.col", - "block_hash.col", - "timestamp.col", - "tx_hash.col", - "tx_index.col", - "log_index.col", - "data.col", - "data_len.col", - "source.col", - "topic0.col", - "topic1.col", - "topic2.col", - "topic3.col", - "topic0.null", - "topic1.null", - "topic2.null", - "topic3.null", - ] { + for name in RAW_FIXED_COLUMNS + .iter() + .map(|(name, _width)| *name) + .chain(std::iter::once("data.col")) + .chain(RAW_NULL_BITMAP_COLUMNS.iter().copied()) + { let path = segment_dir.join(name); if path.exists() { fs::remove_file(path)?; @@ -881,6 +892,254 @@ fn remove_raw_hot_files(segment_dir: &Path) -> std::io::Result<()> { Ok(()) } +pub(crate) fn verify_raw_segment_files_complete( + descriptor: &SegmentDescriptor, + segment_dir: &Path, +) -> std::io::Result<()> { + for (name, width) in RAW_FIXED_COLUMNS { + verify_fixed_raw_column_file(descriptor, segment_dir, name, *width)?; + } + for name in RAW_BITMAP_COLUMNS { + verify_raw_bitmap_file(descriptor, segment_dir, name)?; + } + verify_raw_data_column_file(descriptor, segment_dir)?; + Ok(()) +} + +fn verify_fixed_raw_column_file( + descriptor: &SegmentDescriptor, + segment_dir: &Path, + name: &str, + width: u64, +) -> std::io::Result<()> { + let path = segment_dir.join(name); + let header = read_raw_column_header(descriptor, &path, name)?; + if header.row_count != descriptor.row_count { + return Err(raw_segment_error( + descriptor, + name, + format!( + "row-count mismatch: manifest={} header={}", + descriptor.row_count, header.row_count + ), + )); + } + + let expected_len = + (ColumnFileHeader::SIZE as u64) + .checked_add(header.row_count.checked_mul(width).ok_or_else(|| { + raw_segment_error(descriptor, name, "expected byte length overflow") + })?) + .ok_or_else(|| raw_segment_error(descriptor, name, "expected byte length overflow"))?; + let actual_len = raw_file_len(descriptor, &path, name)?; + if actual_len != expected_len { + return Err(raw_segment_error( + descriptor, + name, + format!("length mismatch: expected {expected_len} bytes, got {actual_len} bytes"), + )); + } + + Ok(()) +} + +fn verify_raw_bitmap_file( + descriptor: &SegmentDescriptor, + segment_dir: &Path, + name: &str, +) -> std::io::Result<()> { + let path = segment_dir.join(name); + let mut file = open_raw_file(descriptor, &path, name)?; + let mut len_bytes = [0u8; 8]; + file.read_exact(&mut len_bytes).map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to read bitmap length: {error}"), + ) + })?; + let bitmap_rows = u64::from_le_bytes(len_bytes); + if bitmap_rows != descriptor.row_count { + return Err(raw_segment_error( + descriptor, + name, + format!( + "row-count mismatch: manifest={} bitmap={bitmap_rows}", + descriptor.row_count + ), + )); + } + + let expected_len = 8u64 + .checked_add(bitmap_rows.div_ceil(8)) + .ok_or_else(|| raw_segment_error(descriptor, name, "expected byte length overflow"))?; + let actual_len = file + .metadata() + .map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to read file metadata: {error}"), + ) + })? + .len(); + if actual_len != expected_len { + return Err(raw_segment_error( + descriptor, + name, + format!("length mismatch: expected {expected_len} bytes, got {actual_len} bytes"), + )); + } + + Ok(()) +} + +fn verify_raw_data_column_file( + descriptor: &SegmentDescriptor, + segment_dir: &Path, +) -> std::io::Result<()> { + let name = "data.col"; + let path = segment_dir.join(name); + let mut file = open_raw_file(descriptor, &path, name)?; + let header = read_raw_column_header_from_file(descriptor, &mut file, name)?; + if header.row_count != descriptor.row_count { + return Err(raw_segment_error( + descriptor, + name, + format!( + "row-count mismatch: manifest={} header={}", + descriptor.row_count, header.row_count + ), + )); + } + + let offsets_size = header + .row_count + .checked_add(1) + .and_then(|count| count.checked_mul(8)) + .ok_or_else(|| raw_segment_error(descriptor, name, "offset table size overflow"))?; + let blob_start = (ColumnFileHeader::SIZE as u64) + .checked_add(offsets_size) + .ok_or_else(|| raw_segment_error(descriptor, name, "blob offset overflow"))?; + let actual_len = file + .metadata() + .map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to read file metadata: {error}"), + ) + })? + .len(); + if actual_len < blob_start { + return Err(raw_segment_error( + descriptor, + name, + format!( + "offset table is truncated: expected at least {blob_start} bytes, got {actual_len} bytes" + ), + )); + } + + let last_offset_position = + (ColumnFileHeader::SIZE as u64) + .checked_add(header.row_count.checked_mul(8).ok_or_else(|| { + raw_segment_error(descriptor, name, "last offset position overflow") + })?) + .ok_or_else(|| raw_segment_error(descriptor, name, "last offset position overflow"))?; + file.seek(SeekFrom::Start(last_offset_position)) + .map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to seek final data offset: {error}"), + ) + })?; + let mut last_offset_bytes = [0u8; 8]; + file.read_exact(&mut last_offset_bytes).map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to read final data offset: {error}"), + ) + })?; + let last_offset = u64::from_le_bytes(last_offset_bytes); + let blob_len = actual_len - blob_start; + if last_offset != blob_len { + return Err(raw_segment_error( + descriptor, + name, + format!( + "blob length mismatch: final offset={last_offset} bytes, blob={blob_len} bytes" + ), + )); + } + + Ok(()) +} + +fn read_raw_column_header( + descriptor: &SegmentDescriptor, + path: &Path, + name: &str, +) -> std::io::Result { + let mut file = open_raw_file(descriptor, path, name)?; + read_raw_column_header_from_file(descriptor, &mut file, name) +} + +fn open_raw_file(descriptor: &SegmentDescriptor, path: &Path, name: &str) -> std::io::Result { + File::open(path).map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to open {}: {error}", path.display()), + ) + }) +} + +fn raw_file_len(descriptor: &SegmentDescriptor, path: &Path, name: &str) -> std::io::Result { + fs::metadata(path) + .map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to read metadata for {}: {error}", path.display()), + ) + }) + .map(|metadata| metadata.len()) +} + +fn read_raw_column_header_from_file( + descriptor: &SegmentDescriptor, + file: &mut File, + name: &str, +) -> std::io::Result { + let mut header_buf = [0u8; ColumnFileHeader::SIZE]; + file.read_exact(&mut header_buf).map_err(|error| { + raw_segment_error( + descriptor, + name, + format!("failed to read column header: {error}"), + ) + })?; + ColumnFileHeader::read_from(&header_buf) + .ok_or_else(|| raw_segment_error(descriptor, name, "corrupt column header")) +} + +fn raw_segment_error( + descriptor: &SegmentDescriptor, + name: &str, + detail: impl std::fmt::Display, +) -> std::io::Error { + std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!( + "segment {} raw column {name} is incomplete: {detail}", + descriptor.id + ), + ) +} + fn remove_superseded_column_dirs(segment_dir: &Path, active_dir: &str) -> std::io::Result<()> { for entry in fs::read_dir(segment_dir)? { let entry = entry?; diff --git a/crates/logex-storage/src/native/storage.rs b/crates/logex-storage/src/native/storage.rs index 20dfb1f6..196bdacc 100644 --- a/crates/logex-storage/src/native/storage.rs +++ b/crates/logex-storage/src/native/storage.rs @@ -20,7 +20,8 @@ use super::catalog::{ use super::segment::{ append_rows, apply_ordered_rows_to_descriptor, apply_rows_to_descriptor, compact_segment, persist_segment_manifest, persist_segment_manifest_with_columns, - segment_uses_current_compaction_profile, write_compacted_rows, + segment_uses_current_compaction_profile, verify_raw_segment_files_complete, + write_compacted_rows, }; const STORAGE_STATE_FILE: &str = "storage_state.json"; @@ -63,7 +64,22 @@ impl SegmentCompactionTask { } pub fn compact(&self) -> std::io::Result<()> { - compact_segment(&self.paths, &self.descriptor) + compact_segment(&self.paths, &self.descriptor).map_err(|error| { + io::Error::new( + error.kind(), + format!( + "failed to compact storage segment {} rows={} blocks=[{}, {}]: {error}", + self.descriptor.id, + self.descriptor.row_count, + self.descriptor + .min_block + .map_or_else(|| "?".to_owned(), |block| block.to_string()), + self.descriptor + .max_block + .map_or_else(|| "?".to_owned(), |block| block.to_string()) + ), + ) + }) } } @@ -1354,6 +1370,7 @@ fn verify_segment_integrity( } if dir.join("address.col").exists() { + verify_raw_segment_files_complete(descriptor, &dir)?; for (name, physical_rows) in hot_segment_physical_row_counts(&dir)? { if physical_rows != row_count { return Err(io::Error::new( @@ -2366,6 +2383,84 @@ mod tests { assert_eq!(err.kind(), io::ErrorKind::InvalidData); } + #[test] + fn startup_integrity_rejects_truncated_raw_column_body() { + let tmp = TempDir::new().unwrap(); + let config = NativeStorageConfig { + data_dir: tmp.path().to_path_buf(), + hot_target_rows: 5, + compaction_safety_margin_blocks: 2_048, + }; + let sealed_id; + + { + let mut storage = NativeStorage::open(config.clone()).unwrap(); + storage.write_batch(&make_rows(6, 100)).unwrap(); + let sealed = storage + .segments() + .iter() + .find(|segment| segment.kind == SegmentKind::Sealed) + .cloned() + .unwrap(); + sealed_id = sealed.id; + + let path = storage.segment_path(sealed.id).join("topic2.col"); + let mut data = fs::read(&path).unwrap(); + data.truncate(data.len() - 32); + fs::write(path, data).unwrap(); + } + + let err = NativeStorage::open(config) + .err() + .expect("truncated raw column should fail startup integrity"); + assert_eq!(err.kind(), io::ErrorKind::InvalidData); + let message = err.to_string(); + assert!(message.contains(&format!("segment {sealed_id} raw column topic2.col"))); + assert!(message.contains("length mismatch")); + } + + #[test] + fn compaction_error_identifies_truncated_raw_column() { + let tmp = TempDir::new().unwrap(); + let mut storage = NativeStorage::open(NativeStorageConfig { + data_dir: tmp.path().to_path_buf(), + hot_target_rows: 5, + compaction_safety_margin_blocks: 100, + }) + .unwrap(); + + storage.write_batch(&make_rows(6, 100)).unwrap(); + let sealed = storage + .segments() + .iter() + .find(|segment| segment.kind == SegmentKind::Sealed) + .cloned() + .unwrap(); + storage + .record_sync_head( + sealed.max_block.unwrap() + 200, + B256::repeat_byte(0xAA), + 999, + ) + .unwrap(); + + let path = storage.segment_path(sealed.id).join("topic2.col"); + let mut data = fs::read(&path).unwrap(); + data.truncate(data.len() - 32); + fs::write(path, data).unwrap(); + + let err = storage + .segment_compaction_plan(1) + .unwrap() + .compact() + .unwrap_err(); + assert_eq!(err.kind(), io::ErrorKind::InvalidData); + let message = err.to_string(); + assert!(message.contains(&format!("failed to compact storage segment {}", sealed.id))); + assert!(message.contains("raw column topic2.col")); + assert!(message.contains("length mismatch")); + } + #[test] fn historical_floor_only_moves_toward_older_blocks() { let tmp = TempDir::new().unwrap();