diff options
| author | Jan Tuomi <jan@jantuomi.fi> | 2024-11-21 21:09:52 +0200 |
|---|---|---|
| committer | Jan Tuomi <jan@jantuomi.fi> | 2024-11-21 21:10:07 +0200 |
| commit | 987f0c0374bf8456fe38a1368103133789f74e62 (patch) | |
| tree | 0660049426bf9e71fdb332b46f414a3639467ab9 /log_db/src/lib.rs | |
| parent | 34730bbe58b77c7ee3d5e6344846a8a8e6595f17 (diff) | |
Refactor do_maintenance_tasks
Diffstat (limited to 'log_db/src/lib.rs')
| -rw-r--r-- | log_db/src/lib.rs | 204 |
1 files changed, 106 insertions, 98 deletions
diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index f049295..664e110 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -365,7 +365,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let segment_num = log_key.segment_num(); let segment_index = log_key.index(); - let metadata_path = &self.data_dir.join(format!("metadata.{}", segment_num)); + let metadata_path = &self.data_dir.join(metadata_filename(segment_num)); let mut metadata_file = READ_MODE.open(&metadata_path)?; request_shared_lock(&self.data_dir, &mut metadata_file)?; @@ -412,7 +412,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let mut found_record: Option<Record> = None; for segment_num in (1..=greatest).rev() { - let segment_path = &self.data_dir.join(format!("metadata.{}", segment_num)); + let segment_path = &self.data_dir.join(metadata_filename(segment_num)); debug!( "Opening segment {} in read mode and acquiring shared lock...", @@ -525,7 +525,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let mut found_records = vec![]; for segment_num in (1..=greatest).rev() { - let segment_path = &self.data_dir.join(format!("metadata.{}", segment_num)); + let segment_path = &self.data_dir.join(metadata_filename(segment_num)); debug!( "Opening segment {} in read mode and acquiring shared lock...", @@ -610,7 +610,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { /// This function should be called periodically to ensure that the database remains in an optimal state. /// Note that this function is synchronous and may block for a relatively long time. /// You may call this function in a separate thread or process to avoid blocking the main thread. - /// However, the database will be exclusively locked, so all writes will be blocked during the tasks. + /// However, the database will be exclusively locked, so all writes and reads will be blocked during the tasks. pub fn do_maintenance_tasks(&mut self) -> Result<(), io::Error> { request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?; @@ -618,121 +618,129 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let metadata_size = self.active_metadata_file.seek(SeekFrom::End(0))?; if metadata_size >= self.config.segment_size as u64 { - debug!("Active log size exceeds threshold, starting rotation and compaction..."); + self.rotate_and_compact()?; + } - self.active_data_file.lock_shared()?; - let original_data_len = self.active_data_file.seek(SeekFrom::End(0))?; + self.active_metadata_file.unlock()?; - debug!("Reading segment data into a BTreeMap"); - let mut map = BTreeMap::new(); - let forward_log_reader = ForwardLogReader::new( - self.active_metadata_file.try_clone()?, - self.active_data_file.try_clone()?, - ); + Ok(()) + } - self.active_data_file.unlock()?; + fn rotate_and_compact(&mut self) -> Result<(), io::Error> { + debug!("Active log size exceeds threshold, starting rotation and compaction..."); - let mut read_n = 0; - for (original_index, record) in forward_log_reader.enumerate() { - let primary_key = record - .at(self.primary_key_index) - .as_indexable() - .expect("Primary key was not indexable"); - map.insert(primary_key, (original_index, record)); - read_n += 1; - } + self.active_data_file.lock_shared()?; + let original_data_len = self.active_data_file.seek(SeekFrom::End(0))?; + let metadata_size = self.active_metadata_file.seek(SeekFrom::End(0))?; - debug!( - "Read {} records, out of which {} were unique", - read_n, - map.len() - ); - debug!("Opening temporary files for writing compacted data"); - // Create a new log data file - let (new_data_uuid, new_data_path) = create_segment_data_file(&self.data_dir)?; - let mut new_data_file = APPEND_MODE.open(&new_data_path)?; + debug!("Reading segment data into a BTreeMap"); + let mut map = BTreeMap::new(); + let forward_log_reader = ForwardLogReader::new( + self.active_metadata_file.try_clone()?, + self.active_data_file.try_clone()?, + ); - let temp_metadata_file = tempfile::NamedTempFile::new()?; - let temp_metadata_path = temp_metadata_file.as_ref(); - let mut temp_metadata_file = WRITE_MODE.open(temp_metadata_path)?; + self.active_data_file.unlock()?; - debug!("Writing compacted data to temporary files"); - let metadata_header = MetadataHeader { - version: 1, - uuid: new_data_uuid, - }; + let mut read_n = 0; + for (original_index, record) in forward_log_reader.enumerate() { + let primary_key = record + .at(self.primary_key_index) + .as_indexable() + .expect("Primary key was not indexable"); + map.insert(primary_key, (original_index, record)); + read_n += 1; + } - temp_metadata_file.write_all(&metadata_header.serialize())?; + debug!( + "Read {} records, out of which {} were unique", + read_n, + map.len() + ); + debug!("Opening temporary files for writing compacted data"); + // Create a new log data file + let (new_data_uuid, new_data_path) = create_segment_data_file(&self.data_dir)?; + let mut new_data_file = APPEND_MODE.open(&new_data_path)?; - let mut metadata_rows_buf = vec![0; metadata_size as usize - METADATA_FILE_HEADER_SIZE]; + let temp_metadata_file = tempfile::NamedTempFile::new()?; + let temp_metadata_path = temp_metadata_file.as_ref(); + let mut temp_metadata_file = WRITE_MODE.open(temp_metadata_path)?; - let mut offset = 0u64; - for (original_index, record) in map.values() { - let serialized = record.serialize(); - let len = serialized.len() as u64; - new_data_file.write_all(&serialized)?; + debug!("Writing compacted data to temporary files"); + let metadata_header = MetadataHeader { + version: 1, + uuid: new_data_uuid, + }; - let metadata_offset = original_index * 16; + temp_metadata_file.write_all(&metadata_header.serialize())?; - for (i, byte) in offset.to_be_bytes().iter().enumerate() { - metadata_rows_buf[metadata_offset + i] = *byte; - } - for (i, byte) in len.to_be_bytes().iter().enumerate() { - metadata_rows_buf[metadata_offset + 8 + i] = *byte; - } + let mut metadata_rows_buf = vec![0; metadata_size as usize - METADATA_FILE_HEADER_SIZE]; - offset += len; + let mut offset = 0u64; + for (original_index, record) in map.values() { + let serialized = record.serialize(); + let len = serialized.len() as u64; + new_data_file.write_all(&serialized)?; + + let metadata_offset = original_index * 16; + + for (i, byte) in offset.to_be_bytes().iter().enumerate() { + metadata_rows_buf[metadata_offset + i] = *byte; + } + for (i, byte) in len.to_be_bytes().iter().enumerate() { + metadata_rows_buf[metadata_offset + 8 + i] = *byte; } - temp_metadata_file.write_all(&metadata_rows_buf)?; + offset += len; + } - // Sync the temporary files to disk - // This is fine to do without consulting WriteDurability because this is a one-off - // operation that is not part of the normal write path. - temp_metadata_file.flush()?; - new_data_file.flush()?; + temp_metadata_file.write_all(&metadata_rows_buf)?; - let final_len = new_data_file.seek(io::SeekFrom::End(0))?; + // Sync the temporary files to disk + // This is fine to do without consulting WriteDurability because this is a one-off + // operation that is not part of the normal write path. + temp_metadata_file.flush()?; + new_data_file.flush()?; - let active_num = greatest_segment_number(&self.data_dir)?; + let final_len = new_data_file.seek(io::SeekFrom::End(0))?; - debug!("Moving temporary files to their final locations"); - let new_data_path = &self.data_dir.join(new_data_uuid.to_string()); - let active_metadata_path = &self.data_dir.join(format!("metadata.{}", active_num)); // overwrite active + let active_num = greatest_segment_number(&self.data_dir)?; - fs::rename(&temp_metadata_path, &active_metadata_path)?; + debug!("Moving temporary files to their final locations"); + let new_data_path = &self.data_dir.join(new_data_uuid.to_string()); + let active_metadata_path = &self.data_dir.join(metadata_filename(active_num)); // overwrite active - debug!( - "Compaction complete, reduced data size: {} -> {}", - original_data_len, final_len - ); + fs::rename(&temp_metadata_path, &active_metadata_path)?; - let new_segment_num = active_num + 1; - let new_metadata_path = self.data_dir.join(format!("metadata.{}", new_segment_num)); - let mut new_metadata_file = - APPEND_MODE.clone().create(true).open(&new_metadata_path)?; + debug!( + "Compaction complete, reduced data size: {} -> {}", + original_data_len, final_len + ); - let new_metadata_header = MetadataHeader { - version: 1, - uuid: new_data_uuid, - }; + let new_segment_num = active_num + 1; + let new_metadata_path = self.data_dir.join(metadata_filename(new_segment_num)); + let mut new_metadata_file = APPEND_MODE.clone().create(true).open(&new_metadata_path)?; + + let new_metadata_header = MetadataHeader { + version: 1, + uuid: new_data_uuid, + }; - new_metadata_file.write_all(&new_metadata_header.serialize())?; + new_metadata_file.write_all(&new_metadata_header.serialize())?; - set_active_segment(&self.data_dir, new_segment_num)?; - self.active_metadata_file.unlock()?; + set_active_segment(&self.data_dir, new_segment_num)?; - self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?; - self.active_data_file = APPEND_MODE.open(&new_data_path)?; + // Old active metadata file should lose lock by RAII, or by + // the manual unlock call in the do_maintenance_tasks method. - // The new active log file is not locked by this client so it cannot be touched. - debug!( - "Active log file rotated and compacted, new segment: {}", - new_segment_num - ); - } else { - self.active_metadata_file.unlock()?; - } + self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?; + self.active_data_file = APPEND_MODE.open(&new_data_path)?; + + // The new active log file is not locked by this client so it cannot be touched. + debug!( + "Active log file rotated and compacted, new segment: {}", + new_segment_num + ); Ok(()) } @@ -769,7 +777,7 @@ mod tests { db.upsert(&record).expect("Failed to insert record"); } - let mut segment1_file = READ_MODE.open(data_dir.join("metadata.1")).unwrap(); + let mut segment1_file = READ_MODE.open(data_dir.join(metadata_filename(1))).unwrap(); let segment1_metadata_size_original = segment1_file.seek(io::SeekFrom::End(0)).unwrap(); let segment1_header = read_metadata_header(&mut segment1_file).unwrap(); @@ -787,14 +795,14 @@ mod tests { db.upsert(&record).expect("Failed to insert record"); // Check that rotation resulted in 2 segments - assert!(fs::exists(data_dir.join("metadata.1")).unwrap()); - assert!(fs::exists(data_dir.join("metadata.2")).unwrap()); + assert!(fs::exists(data_dir.join(metadata_filename(1))).unwrap()); + assert!(fs::exists(data_dir.join(metadata_filename(2))).unwrap()); // Note negation here - assert!(!fs::exists(data_dir.join("metadata.3")).unwrap()); + assert!(!fs::exists(data_dir.join(metadata_filename(3))).unwrap()); // Check that the compacted metadata file has the same size let mut segment1_metadata_file_compacted = - READ_MODE.open(data_dir.join("metadata.1")).unwrap(); + READ_MODE.open(data_dir.join(metadata_filename(1))).unwrap(); let segment1_metadata_size_compacted = segment1_metadata_file_compacted .seek(io::SeekFrom::End(0)) .unwrap(); @@ -862,7 +870,7 @@ mod tests { } // Open the segment file and write garbage to it to simulate corruption - let segment_metadata_path = data_dir.join("metadata.1"); + let segment_metadata_path = data_dir.join(metadata_filename(1)); let mut file = APPEND_MODE .open(&segment_metadata_path) .expect("Failed to open file"); |
