From 641e3cbdc4c82796d6a48c247fb679c81f4e21bd Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Wed, 6 Nov 2024 14:08:06 +0200 Subject: Implement auto-repair for active metadata file --- log_db/src/lib.rs | 213 +++++++++++++++++++++++++++++++++++++-- log_db/src/log_reader_reverse.rs | 5 - log_db/tests/integration.rs | 28 ----- 3 files changed, 203 insertions(+), 43 deletions(-) diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index d6e8de1..a926ea8 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -147,6 +147,12 @@ pub struct DB { secondary_memtables: Vec, } +enum IsActiveMetadataValidResult { + Ok, + ReplaceFile, + TruncateToSize(u64), +} + impl DB { /// Create a new database configuration builder. pub fn configure() -> ConfigBuilder { @@ -345,12 +351,14 @@ impl DB { // Acquire an exclusive lock for writing self.request_exclusive_lock_on_active()?; - if self.ensure_active_file_is_open()? { + if !self.ensure_active_file_is_open()? || !self.ensure_active_metadata_is_valid()? { // The log file has been rotated, so we must try again + self.active_metadata_file.unlock()?; + self.active_data_file.unlock()?; return self.upsert(record); } - debug!("Lock acquired, appending to log file"); + debug!("Exclusive lock acquired, appending to log file"); // Write the record to the log let serialized = &record.serialize(); @@ -371,6 +379,8 @@ impl DB { let mut metadata_buf = vec![]; metadata_buf.extend(&record_offset.to_be_bytes()); metadata_buf.extend(&record_length.to_be_bytes()); + + assert_eq!(metadata_buf.len(), 16); self.active_metadata_file.write_all(&metadata_buf)?; // Flush and sync metadata to disk @@ -390,6 +400,10 @@ impl DB { debug!("Updating memtables"); self.insert_to_memtables(record); + let len = self.active_metadata_file.seek(SeekFrom::End(0))?; + assert!(len >= METADATA_FILE_HEADER_SIZE as u64); + assert_eq!((len - METADATA_FILE_HEADER_SIZE as u64) % 16, 0); + Ok(()) } @@ -406,6 +420,19 @@ impl DB { "Queried value must be indexable", ))?; + match self.is_active_metadata_valid()? { + IsActiveMetadataValidResult::Ok => {} + _ => { + debug!("Active metadata file is invalid, acquiring exclusive lock to start autorepair..."); + self.request_exclusive_lock_on_active()?; + + self.ensure_active_metadata_is_valid()?; + + self.active_metadata_file.unlock()?; + self.active_data_file.unlock()?; + } + }; + debug!("Looking up key {:?} in primary memtable", query_key); let found = self.primary_memtable.get(&query_key); if let Some(record) = found { @@ -513,6 +540,19 @@ impl DB { "Queried value must be indexable", ))?; + match self.is_active_metadata_valid()? { + IsActiveMetadataValidResult::Ok => {} + _ => { + debug!("Active metadata file is invalid, acquiring exclusive lock to start autorepair..."); + self.request_exclusive_lock_on_active()?; + + self.ensure_active_metadata_is_valid()?; + + self.active_metadata_file.unlock()?; + self.active_data_file.unlock()?; + } + }; + // Try to find a memtable with the queried key let found_memtable_index = self.get_secondary_memtable_index_by_field(field); @@ -619,7 +659,7 @@ impl DB { /// Ensures that the `self.metadata_file` and `self.data_file` handles are still pointing to the correct files. /// If the segment has been rotated, the handle will be closed and reopened. - /// Returns `true` if the file has been rotated and the handle has been reopened. + /// Returns `false` if the file has been rotated and the handle has been reopened, `true` otherwise. fn ensure_active_file_is_open(&mut self) -> Result { let data_dir_path = Path::new(&self.config.data_dir); let active_target = fs::read_link(data_dir_path.join("active"))?; @@ -650,9 +690,9 @@ impl DB { self.active_metadata_file = metadata_file; self.active_data_file = fs::OpenOptions::new().append(true).open(&data_file_path)?; - return Ok(true); - } else { return Ok(false); + } else { + return Ok(true); } } @@ -781,11 +821,26 @@ impl DB { let active_metadata_path = data_dir_path.join(active_target); let active_log_md = fs::metadata(&active_log_path)?; + let mut already_locked = false; + match self.is_active_metadata_valid()? { + IsActiveMetadataValidResult::Ok => {} + _ => { + debug!("Active metadata is invalid, acquiring exclusive lock..."); + self.request_exclusive_lock_on_active()?; + already_locked = true; + + self.ensure_active_metadata_is_valid()?; + } + }; + if active_log_md.size() >= self.config.segment_size as u64 { // Rotate the active log file - debug!("Starting rotation, requesting exclusive lock..."); - self.request_exclusive_lock_on_active()?; + debug!("Starting rotation"); + if !already_locked { + debug!("Requesting exclusive lock on active log file..."); + self.request_exclusive_lock_on_active()?; + } debug!("Exclusive lock acquired, rotating active log file..."); @@ -806,6 +861,9 @@ impl DB { debug!("Segment compacted"); } + self.active_metadata_file.unlock()?; + self.active_data_file.unlock()?; + Ok(()) } @@ -875,12 +933,14 @@ impl DB { offset += len; } + let final_len = temp_metadata_file.seek(io::SeekFrom::End(0))?; + debug!("Moving temporary files to their final locations"); let target_data_file_path = data_dir_path.join(new_data_uuid.to_string()); fs::rename(&temp_data_path, &target_data_file_path)?; fs::rename(&temp_metadata_path, metadata_path)?; - debug!("Compaction complete"); + debug!("Compaction complete, resulting size: {}", final_len); Ok(()) } @@ -957,6 +1017,10 @@ impl DB { metadata_file.write_all(&metadata_header.serialize())?; + let len = metadata_file.seek(io::SeekFrom::End(0))?; + assert!(len >= METADATA_FILE_HEADER_SIZE as u64); + assert_eq!((len - METADATA_FILE_HEADER_SIZE as u64) % 16, 0); + Ok((new_num, metadata_path)) } @@ -986,6 +1050,88 @@ impl DB { let header = MetadataHeader::deserialize(&buf); Ok(header) } + + fn is_active_metadata_valid(&mut self) -> Result { + let size = self.active_metadata_file.seek(SeekFrom::End(0))? as usize; + + if size < METADATA_FILE_HEADER_SIZE { + return Ok(IsActiveMetadataValidResult::ReplaceFile); + } + + // The data section must be a multiple of 16 bytes. + // Otherwise, the non-aligned part of the file is dropped. + let data_section_len = size - METADATA_FILE_HEADER_SIZE; + let remainder = data_section_len % 16; + if remainder != 0 { + return Ok(IsActiveMetadataValidResult::TruncateToSize( + (size - remainder) as u64, + )); + } + + Ok(IsActiveMetadataValidResult::Ok) + } + + /// Check that the active metadata file is well-formed and repair it if necessary. + /// The metadata file is considered well-formed if its size is, in pseudocode, `header_size + n * record_size`. + /// If the file is not well-formed, it is truncated to the last well-formed record using + /// a temporary file and an atomic move operation. + /// + /// `self.active_metadata_file` must be a locked file handle opened with read permissions. + /// The function leaves the seek head in an unspecified position. + /// + /// Returns `false` if the file was repaired and rotated, `true` if no action was taken. + fn ensure_active_metadata_is_valid(&mut self) -> Result { + let current_len = self.active_metadata_file.seek(SeekFrom::End(0))? as usize; + + match self.is_active_metadata_valid()? { + IsActiveMetadataValidResult::Ok => return Ok(true), + IsActiveMetadataValidResult::ReplaceFile => { + let data_dir_path = Path::new(&self.config.data_dir); + let active_target = fs::read_link(data_dir_path.join(ACTIVE_SYMLINK_FILENAME))?; + let active_path = data_dir_path.join(&active_target); + warn!( + "Metadata file \"{}\" is malformed ({} bytes), replacing it with an empty file", + active_target.display(), + current_len, + ); + let mut tmp_file = tempfile::NamedTempFile::new()?; + + let header = MetadataHeader { + version: 1, + uuid: Uuid::new_v4(), + }; + + tmp_file.write_all(&header.serialize())?; + fs::rename(tmp_file.path(), active_path)?; + + debug!("Replaced metadata file"); + return Ok(false); + } + IsActiveMetadataValidResult::TruncateToSize(new_size) => { + let data_dir_path = Path::new(&self.config.data_dir); + let active_target = fs::read_link(data_dir_path.join(ACTIVE_SYMLINK_FILENAME))?; + let active_path = data_dir_path.join(&active_target); + warn!( + "Metadata file \"{}\" is malformed ({} bytes), truncating it to {} bytes", + active_target.display(), + current_len, + new_size + ); + + let mut tmp_file = tempfile::NamedTempFile::new()?; + + let mut buf = vec![0; new_size as usize]; + self.active_metadata_file.seek(SeekFrom::Start(0))?; + self.active_metadata_file.read_exact(&mut buf)?; + + tmp_file.write_all(&buf)?; + fs::rename(tmp_file.path(), active_path)?; + + debug!("Truncated metadata file"); + return Ok(false); + } + } + } } #[cfg(test)] @@ -995,8 +1141,6 @@ mod tests { #[derive(Eq, PartialEq, Clone, Debug)] enum Field { Id, - Name, - Data, } #[test] @@ -1064,4 +1208,53 @@ mod tests { _ => false, }); } + + #[test] + fn test_repair() { + let _ = env_logger::builder().is_test(true).try_init(); + let temp_dir = tempfile::tempdir().unwrap(); + let data_dir = temp_dir.path(); + + let mut db = DB::configure() + .data_dir(data_dir.to_str().unwrap()) + .memtable_capacity(0) + .fields(vec![(Field::Id, RecordField::int())]) + .primary_key(Field::Id) + .initialize() + .expect("Failed to create DB"); + + // Insert records + let n_recs = 100; + for i in 0..n_recs { + let record = Record { + values: vec![RecordValue::Int(i as i64)], + }; + db.upsert(&record).expect("Failed to insert record"); + } + + // Open the segment file and write garbage to it to simulate corruption + let segment_metadata_path = data_dir.join("metadata.1"); + let mut file = fs::OpenOptions::new() + .read(true) + .append(true) + .open(&segment_metadata_path) + .expect("Failed to open file"); + + file.write_all(&[0, 1, 2, 3]) + .expect("Failed to write garbage"); + + let len = file.seek(SeekFrom::End(0)).expect("Failed to seek"); + assert_ne!(len, METADATA_FILE_HEADER_SIZE as u64 + n_recs * 16); + + // Try to read from the file, triggering autorepair + db.get(&RecordValue::Int(0)).expect("Failed to get record"); + + // Reopen file and check that it has the correct size + let mut file = fs::OpenOptions::new() + .read(true) + .open(&segment_metadata_path) + .expect("Failed to open file"); + let len = file.seek(SeekFrom::End(0)).expect("Failed to seek"); + assert_eq!(len, METADATA_FILE_HEADER_SIZE as u64 + n_recs * 16); + } } diff --git a/log_db/src/log_reader_reverse.rs b/log_db/src/log_reader_reverse.rs index ce17e59..1023524 100644 --- a/log_db/src/log_reader_reverse.rs +++ b/log_db/src/log_reader_reverse.rs @@ -50,11 +50,6 @@ impl ReverseLogReader { let entry_offset = u64::from_be_bytes(self.metadata_buf[i..i + 8].try_into().unwrap()); let entry_length = u64::from_be_bytes(self.metadata_buf[i + 8..i + 16].try_into().unwrap()); - debug!( - "Read offset {} and length {} from metadata file at position {}", - entry_offset, entry_length, i - ); - self.data_reader.seek(io::SeekFrom::Start(entry_offset))?; let mut result_buf = vec![0; entry_length as usize]; self.data_reader.read_exact(&mut result_buf)?; diff --git a/log_db/tests/integration.rs b/log_db/tests/integration.rs index c431636..5542a3c 100644 --- a/log_db/tests/integration.rs +++ b/log_db/tests/integration.rs @@ -427,32 +427,4 @@ fn test_log_is_rotated_when_capacity_reached() { // 3rd segment should not exist (note negation) assert!(!data_dir_path.join("metadata").with_extension("3").exists()); - - // // Check that the active file only contains five rows - // let mut file = fs::OpenOptions::new() - // .read(true) - // .open(data_dir_path.join(ACTIVE_SYMLINK_FILENAME)) - // .expect("File could not be opened"); - - // let records_in_active_log = ForwardLogReader::new(&mut file).count(); - // assert_eq!(records_in_active_log, 5); - - // // Check that each rotated file contains only 1 record - // // because of compaction - // for i in &[1, 2] { - // let mut file = OpenOptions::new() - // .read(true) - // .open( - // Path::new(&data_dir) - // .join(ACTIVE_SYMLINK_FILENAME) - // .with_extension(i.to_string()), - // ) - // .expect("File could not be opened"); - // let records_in_rotated_log = ForwardLogReader::new(&mut file).count(); - // assert_eq!(records_in_rotated_log, 1); - // } - - // // Look for nonexistant record to scan all segment files - // let found = db.get(&RecordValue::Int(2)).expect("Failed to get record"); - // assert!(found.is_none()); } -- cgit v1.3