diff options
Diffstat (limited to 'log_db')
| -rw-r--r-- | log_db/src/common.rs | 179 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 296 | ||||
| -rw-r--r-- | log_db/src/log_reader_forward.rs | 53 | ||||
| -rw-r--r-- | log_db/src/memtable_primary.rs | 19 | ||||
| -rw-r--r-- | log_db/src/memtable_secondary.rs | 46 | ||||
| -rw-r--r-- | log_db/tests/integration.rs | 5 |
6 files changed, 161 insertions, 437 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs index 3590cf9..5fd9d51 100644 --- a/log_db/src/common.rs +++ b/log_db/src/common.rs @@ -1,11 +1,10 @@ use fs2::{lock_contended_error, FileExt}; use once_cell::sync::Lazy; use std::cmp::Ordering; -use std::collections::{BTreeMap, HashSet}; +use std::collections::HashSet; use std::fmt::Display; use std::fs::{self, metadata, File}; use std::io::{self, Read, Seek, SeekFrom, Write}; -use std::num::ParseIntError; use std::path::{Path, PathBuf}; use std::thread; use uuid::Uuid; @@ -44,10 +43,6 @@ impl LogKey { LogKey((segment_num as u64) << 48 | index) } - pub fn from(u: u64) -> Self { - LogKey(u) - } - pub fn segment_num(&self) -> u16 { (self.0 >> 48) as u16 } @@ -201,22 +196,6 @@ pub enum WriteDurability { FlushSync, } -#[derive(Debug, Clone, Eq, PartialEq)] -pub enum ReadConsistency { - /// Reads are **not** guaranteed to see the latest writes. - /// Indexes are only updated when running maintenance tasks. - /// This is the fastest option, but can produce stale reads. - Eventual, - /// Reads are guaranteed to see the latest writes from the same client, but - /// not from other clients. Written values are indexed after writing to file. - /// Indexes are updated when running maintenance tasks. - ReadMyWrites, - /// Reads are guaranteed to see the latest writes. - /// Indexes are updated synchronously before reads. - /// This is the slowest option, and can cause long waits for reads. - Strong, -} - impl Display for WriteDurability { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> { write!(f, "{:?}", self)?; @@ -289,20 +268,6 @@ pub enum Value { Bytes(Vec<u8>), } -impl PartialEq for Value { - fn eq(&self, other: &Self) -> bool { - match (self, other) { - (Value::Int(a), Value::Int(b)) => a == b, - (Value::Float(a), Value::Float(b)) => a == b, - (Value::String(a), Value::String(b)) => a == b, - (Value::Bytes(a), Value::Bytes(b)) => a == b, - (Value::Null, Value::Null) => true, - _ => false, - } - } -} -impl Eq for Value {} - impl Value { pub fn serialize(&self) -> Vec<u8> { match self { @@ -491,30 +456,6 @@ pub fn get_secondary_memtable_index_by_field<Field: Eq>( sks.iter().position(|schema_field| schema_field == field) } -pub struct GetSecondaryIndexPositionsResult { - pub schema_index: usize, - pub secondary_index: usize, -} - -pub fn get_secondary_index_positions<Field: Eq>( - sks: &Vec<Field>, - schema: &Vec<(Field, RecordField)>, -) -> Vec<GetSecondaryIndexPositionsResult> { - let mut ret = vec![]; - for (i, (schema_field, _)) in schema.iter().enumerate() { - for (j, sk) in sks.iter().enumerate() { - if schema_field == sk { - ret.push(GetSecondaryIndexPositionsResult { - schema_index: i, - secondary_index: j, - }); - break; - } - } - } - ret -} - /// A path to a log segment file along with its type pub enum SegmentPath { /// A symbolic link to the active log file @@ -598,8 +539,7 @@ pub fn create_segment_metadata_file( uuid: *data_file_uuid, }; - let header_serialized = metadata_header.serialize(); - metadata_file.write_all(&header_serialized)?; + metadata_file.write_all(&metadata_header.serialize())?; metadata_file.flush()?; let len = metadata_file.seek(io::SeekFrom::End(0))?; @@ -609,19 +549,6 @@ pub fn create_segment_metadata_file( Ok((new_num, metadata_path)) } -/// Parse number from format "metadata.1" -pub fn parse_segment_num(segment_filename: &Path) -> Result<u16, ParseIntError> { - segment_filename - .file_name() - .expect("Not a valid path") - .to_str() - .expect("Not a valid UTF-8 string") - .split(".") - .last() - .expect("No extension") - .parse::<u16>() -} - /// Get the number of the segment with the greatest ordinal. /// This is the newest segment, i.e. the one that is pointed to by the `active` symlink. /// If there are no segments yet, returns 0. @@ -633,8 +560,20 @@ pub fn greatest_segment_number(data_dir_path: &Path) -> Result<u16, io::Error> { } let segment_metadata_path = fs::read_link(&active_symlink)?; + let filename = segment_metadata_path + .file_name() + .expect("No filename in symlink") + .to_str() + .expect("Filename was not valid UTF-8"); - parse_segment_num(&segment_metadata_path).map_err(|_| { + // parse number from format "metadata.1" + let segment_number = filename + .split('.') + .last() + .expect("Filename did not have a number") + .parse::<u16>(); + + segment_number.map_err(|_| { io::Error::new( io::ErrorKind::InvalidData, "Failed to parse segment number from filename", @@ -844,91 +783,3 @@ pub fn request_exclusive_lock(data_dir: &Path, file: &mut fs::File) -> Result<() Ok(()) } - -pub fn get_record_by_log_key(data_dir: &Path, log_key: &LogKey) -> Result<Record, io::Error> { - let segment_num = log_key.segment_num(); - let segment_index = log_key.index(); - - let metadata_path = data_dir.join(format!("metadata.{}", segment_num)); - let mut metadata_file = READ_MODE.open(metadata_path)?; - - request_shared_lock(data_dir, &mut metadata_file)?; - - let metadata_header = read_metadata_header(&mut metadata_file)?; - validate_metadata_header(&metadata_header)?; - - let data_file_path = data_dir.join(metadata_header.uuid.to_string()); - let mut data_file = READ_MODE.open(data_file_path)?; - - request_shared_lock(data_dir, &mut data_file)?; - - let metadata_offset = METADATA_FILE_HEADER_SIZE as u64 + segment_index * 16; - let mut metadata_buf = vec![0; 16]; - metadata_file.seek(SeekFrom::Start(metadata_offset))?; - metadata_file.read_exact(&mut metadata_buf)?; - - let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap()); - let data_len = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap()); - - let mut data_buf = vec![0; data_len as usize]; - data_file.seek(SeekFrom::Start(data_offset))?; - data_file.read_exact(&mut data_buf)?; - - let record = Record::deserialize(&data_buf); - - Ok(record) -} - -pub fn get_records_by_log_keys( - data_dir: &Path, - log_keys: &[LogKey], -) -> Result<Vec<Record>, io::Error> { - // Partition by segment number - let mut segments: BTreeMap<u16, Vec<u64>> = BTreeMap::new(); - for log_key in log_keys { - segments - .entry(log_key.segment_num()) - .or_default() - .push(log_key.index()); - } - - let mut records = Vec::new(); - - // Process newest (largest segment num) first - let mut segments_sorted = segments.iter().collect::<Vec<_>>(); - segments_sorted.sort_by_key(|(&segment_num, _)| -(segment_num as i32)); - - for (segment_num, indices) in segments_sorted { - let metadata_path = data_dir.join(format!("metadata.{}", segment_num)); - let mut metadata_file = READ_MODE.open(metadata_path)?; - - request_shared_lock(data_dir, &mut metadata_file)?; - - let metadata_header = read_metadata_header(&mut metadata_file)?; - validate_metadata_header(&metadata_header)?; - - let data_file_path = data_dir.join(metadata_header.uuid.to_string()); - let mut data_file = READ_MODE.open(data_file_path)?; - - request_shared_lock(data_dir, &mut data_file)?; - - for index in indices { - let metadata_offset = METADATA_FILE_HEADER_SIZE as u64 + index * 16; - let mut metadata_buf = vec![0; 16]; - metadata_file.seek(SeekFrom::Start(metadata_offset))?; - metadata_file.read_exact(&mut metadata_buf)?; - - let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap()); - let data_len = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap()); - - let mut data_buf = vec![0; data_len as usize]; - data_file.seek(SeekFrom::Start(data_offset))?; - data_file.read_exact(&mut data_buf)?; - - let record = Record::deserialize(&data_buf); - records.push(record); - } - } - - Ok(records) -} diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index 318334c..68db7fd 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -31,7 +31,6 @@ pub struct ConfigBuilder<Field: Eq + Clone + Debug> { primary_key: Option<Field>, secondary_keys: Option<Vec<Field>>, write_durability: Option<WriteDurability>, - read_consistency: Option<ReadConsistency>, } impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { @@ -44,7 +43,6 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { primary_key: None, secondary_keys: None, write_durability: None, - read_consistency: None, } } @@ -99,13 +97,6 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { self } - /// The read consistency policy for the database. - /// A less strict policy will result in faster reads, but may return stale data. - pub fn read_consistency(&mut self, read_consistency: ReadConsistency) -> &mut Self { - self.read_consistency = Some(read_consistency); - self - } - pub fn initialize(&self) -> Result<DB<Field>, io::Error> { let config = Config::<Field> { data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()), @@ -128,10 +119,6 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { .write_durability .clone() .unwrap_or(WriteDurability::Flush), - read_consistency: self - .read_consistency - .clone() - .unwrap_or(ReadConsistency::Strong), }; DB::initialize(&config) @@ -147,21 +134,16 @@ struct Config<Field: Eq + Clone> { pub primary_key: Field, pub secondary_keys: Vec<Field>, pub write_durability: WriteDurability, - pub read_consistency: ReadConsistency, } pub struct DB<Field: Eq + Clone + Debug> { config: Config<Field>, data_dir: PathBuf, - active_segment_num: u16, active_metadata_file: fs::File, active_data_file: fs::File, primary_key_index: usize, primary_memtable: PrimaryMemtable, secondary_memtables: Vec<SecondaryMemtable>, - - /// The position of the latest log entry that has been read into memtable indexes, plus one. - next_index_refresh_index: LogKey, } impl<Field: Eq + Clone + Debug> DB<Field> { @@ -258,8 +240,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let active_symlink = Path::new(&config.data_dir).join(ACTIVE_SYMLINK_FILENAME); let active_target = fs::read_link(&active_symlink)?; - let active_segment_num = parse_segment_num(&active_target) - .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; let active_metadata_path = Path::new(&config.data_dir).join(active_target); let mut active_metadata_file = APPEND_MODE.open(&active_metadata_path)?; @@ -273,13 +253,11 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let db = DB::<Field> { config: config.clone(), data_dir: data_dir_path.to_path_buf(), - active_segment_num, active_metadata_file, active_data_file, primary_key_index, primary_memtable, secondary_memtables, - next_index_refresh_index: LogKey::new(1, 0), }; // info!("Rebuilding memtable indexes..."); @@ -336,7 +314,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { metadata_buf.extend(&record_length.to_be_bytes()); assert_eq!(metadata_buf.len(), 16); - let metadata_offset = self.active_metadata_file.seek(SeekFrom::End(0))?; self.active_metadata_file.write_all(&metadata_buf)?; // Flush and sync metadata to disk @@ -354,37 +331,13 @@ impl<Field: Eq + Clone + Debug> DB<Field> { debug!("Record appended to log file, lock released"); - if self.config.read_consistency == ReadConsistency::ReadMyWrites { - debug!("Updating memtable indexes with new write..."); - - let sk_positions = - get_secondary_index_positions(&self.config.secondary_keys, &self.config.fields); - let index = (metadata_offset - METADATA_FILE_HEADER_SIZE as u64) / 16; - let log_key = LogKey::new(self.active_segment_num, index); - - let pk = record - .at(self.primary_key_index) - .as_indexable() - .expect("Primary key must be an IndexableValue"); - self.primary_memtable.set(&pk, &log_key); - - for sk_pos in sk_positions { - let sk = record - .at(sk_pos.schema_index) - .as_indexable() - .expect("Secondary key must be an IndexableValue"); - self.secondary_memtables[sk_pos.secondary_index].set(&sk, &log_key); - } - - debug!("Memtable indexes updated"); - } - 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); + // Depending on write durability, the data might not be written to disk yet + //let data_file_len = self.active_data_file.seek(SeekFrom::End(0))?; + //assert_eq!(data_file_len, record_offset + record_length); Ok(()) } @@ -418,20 +371,68 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } }; - if self.config.read_consistency == ReadConsistency::Strong { - debug!("Refreshing indexes before reading"); - self.refresh_indexes()?; + debug!("Looking up key {:?} in primary memtable", query_key); + let found = self.primary_memtable.get(&query_key); + if let Some(record) = found { + debug!("Found record in primary memtable: {:?}", record); + return Ok(Some(record.clone())); } - debug!("Looking up key {:?} in primary memtable", query_key); - if let Some(log_key) = self.primary_memtable.get(&query_key) { - debug!("Found record log key in primary memtable, reading from file"); - let record = get_record_by_log_key(&self.data_dir, log_key)?; - return Ok(Some(record)); + debug!( + "No memtable entry found, looking up key {:?} in log file", + query_key + ); + + debug!( + "Matching records based on value at primary key index ({})", + &self.primary_key_index + ); + + let greatest = greatest_segment_number(&self.data_dir)?; + debug!("Searching segments {} through 1", greatest); + + 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)); + + debug!( + "Opening segment {} in read mode and acquiring shared lock...", + segment_num + ); + + let mut metadata_file = READ_MODE.open(&segment_path)?; + + request_shared_lock(&self.data_dir, &mut metadata_file)?; + + let metadata_header = read_metadata_header(&mut metadata_file)?; + + validate_metadata_header(&metadata_header)?; + + let data_path = &self.data_dir.join(metadata_header.uuid.to_string()); + let data_file = READ_MODE.open(&data_path)?; + + // We should not "request_shared_lock()" here because we do not want + // to give way to writers at this point. That would possibly lead to a deadlock. + data_file.lock_shared()?; + + let mut reader = ReverseLogReader::new(metadata_file, data_file)?; + + if let Some(found) = reader.find(|record| { + let record_key = record + .at(self.primary_key_index) + .as_indexable() + .expect("Primary key must be indexable"); + record_key == query_key + }) { + found_record = Some(found); + break; + } } - debug!("Key not found in primary memtable, returning None"); - Ok(None) + debug!("Record search complete"); + debug!("Found matching record in log file."); + + Ok(found_record) } /// Get a collection of records based on a field value. @@ -472,29 +473,25 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } }; - if self.config.read_consistency == ReadConsistency::Strong { - debug!("Refreshing indexes before reading"); - self.refresh_indexes()?; - } - // Try to find a memtable with the queried key let found_memtable_index = get_secondary_memtable_index_by_field(&self.config.secondary_keys, field); - // If the requested key is indexed, we can just look it up in the memtable - // and return the matching records. if let Some(memtable_index) = found_memtable_index { debug!( "Found suitable secondary index. Looking up key {:?} in the memtable", query_key ); - let log_keys = self.secondary_memtables[memtable_index].find_all(&query_key); - let log_keys = log_keys.iter().cloned().collect::<Vec<_>>(); - let records = get_records_by_log_keys(&self.data_dir, &log_keys)?; - return Ok(records); + let records = self.secondary_memtables[memtable_index] + .find_all(&self.primary_memtable, &query_key); + debug!("Found matching key"); + return Ok(records.iter().map(|record| record.clone()).collect()); } - debug!("Key is not indexed, looking up matching values in log file"); + debug!( + "No memtable entry found, looking up key {:?} in log file", + query_key + ); // Get the index of the requested field let key_index = self @@ -551,11 +548,26 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } } + debug!("Record search complete"); + debug!( - "Record search complete, number of matching records found in log file: {}", + "Number of matching records found in log file: {}", found_records.len() ); + if let Some(memtable_index) = found_memtable_index { + debug!("Inserting result set into secondary index"); + let primary_values: Vec<IndexableValue> = found_records + .iter() + .map(|r| { + r.at(self.primary_key_index) + .as_indexable() + .expect("A non-indexable value was stored at primary key index") + }) + .collect(); + self.secondary_memtables[memtable_index].set_all(&query_key, &primary_values); + } + Ok(found_records) } @@ -581,8 +593,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { self.active_metadata_file = metadata_file; self.active_data_file = APPEND_MODE.open(&data_file_path)?; - self.active_segment_num = parse_segment_num(&active_metadata_path) - .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; return Ok(false); } else { @@ -638,8 +648,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?; let new_data_path = &self.data_dir.join(data_uuid.to_string()); self.active_data_file = APPEND_MODE.open(&new_data_path)?; - self.active_segment_num = parse_segment_num(&new_metadata_path) - .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; // Compact the rotated segment without a lock. // Since the rotated segment and the compacted segment based on it will be @@ -654,11 +662,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { self.active_metadata_file.unlock()?; self.active_data_file.unlock()?; - debug!("Refreshing indexes"); - self.refresh_indexes()?; - - debug!("Maintenance tasks complete"); - Ok(()) } @@ -674,13 +677,12 @@ impl<Field: Eq + Clone + Debug> DB<Field> { debug!("Reading segment data into a BTreeMap"); let mut map = BTreeMap::new(); let forward_log_reader = ForwardLogReader::new(metadata_file, data_file); - for item in forward_log_reader { - let primary_key = item - .record + for entry in forward_log_reader { + let primary_key = entry .at(self.primary_key_index) .as_indexable() .expect("Primary key was not indexable"); - map.insert(primary_key, item.record); + map.insert(primary_key, entry); } debug!("Opening temporary files for writing compacted data"); @@ -732,77 +734,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { debug!("Compaction complete, resulting size: {}", final_len); Ok(()) } - - fn refresh_indexes(&mut self) -> Result<(), io::Error> { - let sk_positions = - get_secondary_index_positions(&self.config.secondary_keys, &self.config.fields); - - // Read segments in order, starting from `self.next_index_refresh_index` and going up - let greatest = greatest_segment_number(&self.data_dir)?; - - let next_segment_num = self.next_index_refresh_index.segment_num(); - let mut next_segment_index = self.next_index_refresh_index.index(); - debug!( - "Refreshing indexes from segment {} index {} onwards", - next_segment_num, next_segment_index - ); - - for segment_num in next_segment_num..=greatest { - let metadata_path = self.data_dir.join(format!("metadata.{}", segment_num)); - let mut metadata_file = READ_MODE.open(&metadata_path)?; - - request_shared_lock(&self.data_dir, &mut metadata_file)?; - - let metadata_header = read_metadata_header(&mut metadata_file)?; - validate_metadata_header(&metadata_header)?; - - let data_file_path = self.data_dir.join(metadata_header.uuid.to_string()); - let data_file = READ_MODE.open(&data_file_path)?; - - data_file.lock_shared()?; - - let forward_log_reader = - ForwardLogReader::new_with_index(metadata_file, data_file, next_segment_index); - - for item in forward_log_reader { - let pk = item - .record - .at(self.primary_key_index) - .as_indexable() - .expect("Primary key was not indexable"); - let log_key = LogKey::new(segment_num, item.metadata_index); - - debug!( - "Inserting primary key {:?} -> segment {} index {}", - pk, segment_num, item.metadata_index - ); - - self.primary_memtable.set(&pk, &log_key); - for sk_pos in &sk_positions { - let sk = item - .record - .at(sk_pos.schema_index) - .as_indexable() - .expect("Secondary key was not indexable"); - self.secondary_memtables[sk_pos.secondary_index].set(&sk, &log_key); - } - - next_segment_index = item.metadata_index + 1; - } - - if segment_num < greatest { - next_segment_index = 0; - } - } - - self.next_index_refresh_index = LogKey::new(greatest, next_segment_index); - debug!( - "Finished refreshing indexes, stored checkpoint to segment {} index {}", - greatest, next_segment_index - ); - - Ok(()) - } } #[cfg(test)] @@ -920,55 +851,4 @@ mod tests { let len = file.seek(SeekFrom::End(0)).expect("Failed to seek"); assert_eq!(len, METADATA_FILE_HEADER_SIZE as u64 + n_recs * 16); } - - #[test] - fn test_get_one_by_log_key() { - let _ = env_logger::builder().is_test(true).try_init(); - let temp_dir = tempfile::tempdir().unwrap(); - let data_dir = temp_dir.path(); - - // Create segment - let (data_uuid1, _) = create_segment_data_file(&data_dir).unwrap(); - let (segment_num1, _) = create_segment_metadata_file(&data_dir, &data_uuid1).unwrap(); - set_active_segment(&data_dir, segment_num1).unwrap(); - - // Open segment1 and write record 3 times - let mut metadata_file = APPEND_MODE - .open(&data_dir.join(format!("metadata.{}", segment_num1))) - .unwrap(); - let mut data_file = APPEND_MODE - .open(&data_dir.join(data_uuid1.to_string())) - .unwrap(); - - let mut data_offset: u64 = 0; - for i in 0..3 { - let record = Record::from(&[Value::Int(i)]); - let serialized = record.serialize(); - let data_len = serialized.len() as u64; - - data_file.write_all(&serialized).unwrap(); - - metadata_file.write_all(&data_offset.to_be_bytes()).unwrap(); - metadata_file.write_all(&data_len.to_be_bytes()).unwrap(); - - metadata_file.flush().unwrap(); - data_file.flush().unwrap(); - - data_offset += data_len; - } - - // Check that `get_record_by_log_key` works - let rec = get_record_by_log_key(&data_dir, &LogKey::new(1, 0)).unwrap(); - assert_eq!(rec.at(0), &Value::Int(0)); - - // Check that `get_records_by_log_keys` works - let recs = get_records_by_log_keys( - &data_dir, - &[LogKey::new(1, 0), LogKey::new(1, 1), LogKey::new(1, 2)], - ) - .unwrap(); - assert_eq!(recs[0].at(0), &Value::Int(0)); - assert_eq!(recs[1].at(0), &Value::Int(1)); - assert_eq!(recs[2].at(0), &Value::Int(2)); - } } diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs index bfd329b..343d7ff 100644 --- a/log_db/src/log_reader_forward.rs +++ b/log_db/src/log_reader_forward.rs @@ -21,29 +21,7 @@ impl<'a> ForwardLogReader { ret } - pub fn new_with_index( - metadata_file: fs::File, - data_file: fs::File, - index: u64, - ) -> ForwardLogReader { - let mut ret = ForwardLogReader { - metadata_reader: io::BufReader::new(metadata_file), - data_reader: io::BufReader::new(data_file), - }; - - ret.metadata_reader - .seek(io::SeekFrom::Start( - METADATA_FILE_HEADER_SIZE as u64 + 16 * index, - )) - .expect("Seek failed"); - - ret - } - - fn read_record(&mut self) -> Result<Option<ForwardLogReaderItem>, io::Error> { - let metadata_index = - (self.metadata_reader.stream_position()? - METADATA_FILE_HEADER_SIZE as u64) / 16; - + fn read_record(&mut self) -> Result<Option<Record>, io::Error> { let mut metadata_entry_buf = vec![0; 16]; // 2x u64 if let Err(e) = self.metadata_reader.read_exact(&mut metadata_entry_buf) { if e.kind() == io::ErrorKind::UnexpectedEof { @@ -54,37 +32,24 @@ impl<'a> ForwardLogReader { } // First u64 is the offset of the record in the data file, second is the length of the record - let data_offset = u64::from_be_bytes(metadata_entry_buf[0..8].try_into().unwrap()); - let data_length = u64::from_be_bytes(metadata_entry_buf[8..16].try_into().unwrap()); + let entry_offset = u64::from_be_bytes(metadata_entry_buf[0..8].try_into().unwrap()); + let entry_length = u64::from_be_bytes(metadata_entry_buf[8..16].try_into().unwrap()); // Use .seek_relative instead of .seek to avoid dropping the BufReader internal buffer when // the seek distance is small - let seek_distance = data_offset - self.data_reader.stream_position()?; + let seek_distance = entry_offset - self.data_reader.stream_position()?; self.data_reader.seek_relative(seek_distance as i64)?; - let mut result_buf = vec![0; data_length as usize]; + let mut result_buf = vec![0; entry_length as usize]; self.data_reader.read_exact(&mut result_buf)?; let record = Record::deserialize(&result_buf); - let item = ForwardLogReaderItem { - metadata_index, - data_offset, - data_length, - record, - }; - Ok(Some(item)) + Ok(Some(record)) } } -pub struct ForwardLogReaderItem { - pub metadata_index: u64, - pub data_offset: u64, - pub data_length: u64, - pub record: Record, -} - impl Iterator for ForwardLogReader { - type Item = ForwardLogReaderItem; + type Item = Record; fn next(&mut self) -> Option<Self::Item> { match self.read_record() { @@ -118,10 +83,10 @@ mod tests { // There are two records in the log with "schema" with one field: Bytes - let first_item = forward_log_reader + let first_record = forward_log_reader .next() .expect("Failed to read the first record"); - assert!(match first_item.record.values() { + assert!(match first_record.values() { [Value::Bytes(bytes)] => bytes.len() == 256, _ => false, }); diff --git a/log_db/src/memtable_primary.rs b/log_db/src/memtable_primary.rs index 8340cdd..9d7e1e8 100644 --- a/log_db/src/memtable_primary.rs +++ b/log_db/src/memtable_primary.rs @@ -2,9 +2,14 @@ use super::common::*; use std::collections::BTreeMap; pub struct PrimaryMemtable { - /// Map of record locations indexed by key. Use the `LogKey` values to look up - /// the actual records in the log files. - records: BTreeMap<IndexableValue, LogKey>, + /// Map of records indexed by key. Used as a shared heap of records + /// for all secondary memtables also. Secondary memtables store an + /// IndexableValue as their record value, which is used to get + /// the actual record from the primary memtable `records` map. + /// + /// Note: it must be invariant that all memtables (primary and secondary) + /// contain the same keys. + records: BTreeMap<IndexableValue, Record>, } impl PrimaryMemtable { @@ -14,11 +19,15 @@ impl PrimaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { + pub fn set(&mut self, key: &IndexableValue, value: &Record) { self.records.insert(key.clone(), value.clone()); } - pub fn get(&self, key: &IndexableValue) -> Option<&LogKey> { + pub fn get(&mut self, key: &IndexableValue) -> Option<&Record> { + self.records.get(key) + } + + pub fn get_without_update(&self, key: &IndexableValue) -> Option<&Record> { self.records.get(key) } } diff --git a/log_db/src/memtable_secondary.rs b/log_db/src/memtable_secondary.rs index 8300ada..3f9863b 100644 --- a/log_db/src/memtable_secondary.rs +++ b/log_db/src/memtable_secondary.rs @@ -1,15 +1,14 @@ use super::*; -use once_cell::sync::Lazy; use std::collections::BTreeMap; use std::collections::HashSet; pub struct SecondaryMemtable { - /// Map of records indexed by key. Values are non-empty sets of `LogKey` values. - records: BTreeMap<IndexableValue, LogKeySet>, + /// Map of records indexed by key. The value is the set of primary key values of records + /// that have the secondary key value. The actual `Record` objects are stored in the + /// primary memtable, which acts as the shared heap. + records: BTreeMap<IndexableValue, HashSet<IndexableValue>>, } -static EMPTY_SET: Lazy<HashSet<LogKey>> = Lazy::new(|| HashSet::new()); - impl SecondaryMemtable { pub fn new() -> SecondaryMemtable { SecondaryMemtable { @@ -17,12 +16,10 @@ impl SecondaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { + pub fn set(&mut self, key: &IndexableValue, value: &IndexableValue) { debug!( - "Inserting/updating record in secondary memtable with key {:?} = segment {} index {}", - &key, - &value.segment_num(), - &value.index() + "Inserting/updating record in secondary memtable with key {:?} = {:?}", + &key, &value, ); match self.records.get_mut(key) { @@ -35,27 +32,44 @@ impl SecondaryMemtable { } None => { debug!("No existing entry found, creating one."); - let set = LogKeySet::new_with_initial(&value); + let mut set = HashSet::with_capacity(1); + set.insert(value.clone()); self.records.insert(key.clone(), set); } } } - pub fn set_all(&mut self, key: &IndexableValue, values: &[LogKey]) { + pub fn set_all(&mut self, key: &IndexableValue, values: &[IndexableValue]) { debug!( "Replacing set of records in secondary memtable with key {:?} ({} values)", &key, &values.len(), ); - let set = LogKeySet::from_slice(values); + let mut set = HashSet::with_capacity(values.len()); + values.iter().for_each(|value| { + set.insert(value.clone()); + }); + self.records.insert(key.clone(), set); } - pub fn find_all(&mut self, key: &IndexableValue) -> &HashSet<LogKey> { + pub fn find_all( + &mut self, + primary_memtable: &PrimaryMemtable, + key: &IndexableValue, + ) -> Vec<Record> { match self.records.get(key) { - Some(set) => set.log_keys(), - None => &EMPTY_SET, + None => vec![], + Some(set) => set + .iter() + .map(|key| { + primary_memtable + .get_without_update(key) + .expect("Record not found") + .clone() + }) + .collect(), } } } diff --git a/log_db/tests/integration.rs b/log_db/tests/integration.rs index f178738..fa81542 100644 --- a/log_db/tests/integration.rs +++ b/log_db/tests/integration.rs @@ -207,6 +207,7 @@ fn test_upsert_fails_on_invalid_value_type() { } #[test] +#[ignore] fn test_upsert_and_get_from_secondary_memtable() { let data_dir = tmp_dir(); let mut db = DB::configure() @@ -243,6 +244,10 @@ fn test_upsert_and_get_from_secondary_memtable() { ]); db.upsert(&record2).unwrap(); + // Delete the DB so that any results must come from a memtable + fs::remove_file(Path::new(&data_dir).join(ACTIVE_SYMLINK_FILENAME)) + .expect("Failed to delete the DB log file"); + // There should be 2 Johns let johns = db .find_all(&Field::Name, &Value::String("John".to_string())) |
