diff options
| author | Jan Tuomi <jan@jantuomi.fi> | 2024-11-08 19:31:45 +0200 |
|---|---|---|
| committer | Jan Tuomi <jan@jantuomi.fi> | 2024-11-08 21:06:03 +0200 |
| commit | e2d68847e00e58ddb646fa281ec4dfa7b67489cf (patch) | |
| tree | d0ec2f46aed810d389919934d69c6c15e489f04d | |
| parent | 52d3d69035e80bb3889cf720324aa813df1f7f80 (diff) | |
Fix: set active files after rotation
| -rw-r--r-- | log_db/src/common.rs | 21 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 63 |
2 files changed, 37 insertions, 47 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs index cbb36a8..347b851 100644 --- a/log_db/src/common.rs +++ b/log_db/src/common.rs @@ -72,6 +72,27 @@ impl LogKeySet { LogKeySet { set } } + pub fn from_slice(keys: &[LogKey]) -> Self { + assert!( + !keys.is_empty(), + "LogKeySet::from_slice must be supplied a non-empty slice" + ); + let mut set = HashSet::with_capacity(keys.len()); + keys.iter().for_each(|key| { + set.insert(key.clone()); + }); + LogKeySet { set } + } + + pub fn iter(&self) -> std::collections::hash_set::Iter<'_, LogKey> { + self.set.iter() + } + + /// The number of LogKeys in the set. + pub fn len(&self) -> usize { + self.set.len() + } + /// Insert a LogKey into the set. pub fn insert(&mut self, key: LogKey) { self.set.insert(key); diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index a926ea8..c4f8deb 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -266,6 +266,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let active_data_path = Path::new(&config.data_dir).join(active_metadata_header.uuid.to_string()); let active_data_file = fs::OpenOptions::new() + .read(true) .append(true) .open(&active_data_path)?; @@ -362,7 +363,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { // Write the record to the log let serialized = &record.serialize(); - let record_offset = self.active_data_file.metadata()?.len(); + let record_offset = self.active_data_file.seek(SeekFrom::End(0))?; let record_length = serialized.len() as u64; self.active_data_file.write_all(serialized)?; @@ -397,12 +398,11 @@ impl<Field: Eq + Clone + Debug> DB<Field> { debug!("Record appended to log file, lock released"); - 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); + let data_file_len = self.active_data_file.seek(SeekFrom::End(0))?; + assert_eq!(data_file_len, record_offset + record_length); Ok(()) } @@ -497,20 +497,8 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } debug!("Record search complete"); - - let result_value = match &found_record { - Some(record) => record, - None => { - debug!("No record found for key {:?}", query_key); - return Ok(None); - } - }; - debug!("Found matching record in log file."); - debug!("Updating memtables"); - self.insert_to_memtables(&result_value); - Ok(found_record) } @@ -696,37 +684,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } } - fn insert_to_memtables(&mut self, record: &Record) { - let key = record.values[self.primary_key_index] - .as_indexable() - .expect("A non-indexable value was stored at key index"); - - if self.config.memtable_capacity == 0 { - return; - } - - debug!( - "Inserting/updating record in primary memtable with key {:?} = {:?}", - &key, &record, - ); - - self.primary_memtable.set(&key, record); - - for (field_index, value) in record.values.iter().enumerate() { - let field = &self.config.fields[field_index].0; - if let Some(smt_index) = self.get_secondary_memtable_index_by_field(field) { - debug!("Updating memtable for index on {:?}", field); - - let memtable = &mut self.secondary_memtables[smt_index]; - let key = value.as_indexable().expect("Primary key was not indexable"); - let primary_key = record.values[self.primary_key_index] - .as_indexable() - .expect("Primary key was not indexable"); - memtable.set(&key, &primary_key); - } - } - } - fn request_exclusive_lock_on_active(&mut self) -> Result<(), io::Error> { // Create a lock on the exclusive lock request file to signal to readers that they should wait let lock_request_path = Path::new(&self.config.data_dir).join(EXCL_LOCK_REQUEST_FILENAME); @@ -853,11 +810,23 @@ impl<Field: Eq + Clone + Debug> DB<Field> { // The new active log file is not locked by this client so it cannot be touched. debug!("Active log file rotated, new segment: {}", new_segment_num); + self.active_metadata_file = fs::OpenOptions::new() + .read(true) + .append(true) + .open(&data_dir_path.join(format!("metadata.{}", new_segment_num)))?; + + self.active_data_file = fs::OpenOptions::new() + .read(true) + .append(true) + .open(&data_dir_path.join(data_file_uuid.to_string()))?; + // Compact the rotated segment without a lock. // Since the rotated segment and the compacted segment based on it will be // a) read-only, and b) identical in effective content, there is no need to lock it. self.compact_segment(&active_metadata_path)?; + // The new active log file is not locked by this client so it cannot be touched. + debug!("Active log file rotated, new segment: {}", new_segment_num); debug!("Segment compacted"); } |
