diff options
| -rw-r--r-- | log_db/src/common.rs | 29 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 63 | ||||
| -rw-r--r-- | log_db/src/log_reader_forward.rs | 37 |
3 files changed, 106 insertions, 23 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs index 6fd4c1f..e384a8f 100644 --- a/log_db/src/common.rs +++ b/log_db/src/common.rs @@ -564,18 +564,9 @@ pub fn create_segment_metadata_file( Ok((new_num, metadata_path)) } -/// 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. -pub fn greatest_segment_number(data_dir_path: &Path) -> Result<u16, io::Error> { - let active_symlink = data_dir_path.join(ACTIVE_SYMLINK_FILENAME); - - if !fs::exists(&active_symlink)? { - return Ok(0); - } - - let segment_metadata_path = fs::read_link(&active_symlink)?; - let filename = segment_metadata_path +/// Parse the segment number from a metadata file path +pub fn parse_segment_number(metadata_path: &Path) -> Result<u16, io::Error> { + let filename = metadata_path .file_name() .expect("No filename in symlink") .to_str() @@ -596,6 +587,20 @@ pub fn greatest_segment_number(data_dir_path: &Path) -> Result<u16, io::Error> { }) } +/// 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. +pub fn greatest_segment_number(data_dir_path: &Path) -> Result<u16, io::Error> { + let active_symlink = data_dir_path.join(ACTIVE_SYMLINK_FILENAME); + + if !fs::exists(&active_symlink)? { + return Ok(0); + } + + let segment_metadata_path = fs::read_link(&active_symlink)?; + parse_segment_number(&segment_metadata_path) +} + /// Create a new segment data file and return its UUID. /// A data file contains the segment data, tightly packed without separators. /// An accompanying metadata file is required to interpret the data. diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index 664e110..3d8ddb9 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -10,6 +10,7 @@ mod memtable_secondary; pub use common::*; use fs2::FileExt; pub use log_reader_forward::ForwardLogReader; +use log_reader_forward::ForwardLogReaderItem; pub use log_reader_reverse::ReverseLogReader; use memtable_primary::PrimaryMemtable; use memtable_secondary::SecondaryMemtable; @@ -131,6 +132,7 @@ pub struct DB<Field: Eq + Clone + Debug> { primary_key_index: usize, primary_memtable: PrimaryMemtable, secondary_memtables: Vec<SecondaryMemtable>, + refresh_next_logkey: LogKey, } impl<Field: Eq + Clone + Debug> DB<Field> { @@ -237,7 +239,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { Path::new(&config.data_dir).join(active_metadata_header.uuid.to_string()); let active_data_file = APPEND_MODE.open(&active_data_path)?; - let db = DB::<Field> { + let mut db = DB::<Field> { config: config.clone(), data_dir: data_dir_path.to_path_buf(), active_metadata_file, @@ -245,16 +247,62 @@ impl<Field: Eq + Clone + Debug> DB<Field> { primary_key_index, primary_memtable, secondary_memtables, + refresh_next_logkey: LogKey::new(1, 0), }; - // info!("Rebuilding memtable indexes..."); - // TODO FIXME build memtable indexes + info!("Rebuilding memtable indexes..."); + db.refresh_indexes()?; info!("Database ready."); Ok(db) } + fn refresh_indexes(&mut self) -> Result<(), io::Error> { + let active_symlink_path = self.data_dir.join(ACTIVE_SYMLINK_FILENAME); + let active_target = fs::read_link(active_symlink_path)?; + let active_metadata_path = self.data_dir.join(active_target); + + let to_segnum = parse_segment_number(&active_metadata_path)?; + let from_segnum = self.refresh_next_logkey.segment_num(); + let mut from_index = self.refresh_next_logkey.index(); + + for segnum in from_segnum..=to_segnum { + let metadata_path = self.data_dir.join(metadata_filename(segnum)); + 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_path = self.data_dir.join(metadata_header.uuid.to_string()); + let data_file = READ_MODE.open(data_path)?; + + for ForwardLogReaderItem { record, index } in + ForwardLogReader::new_with_index(metadata_file, data_file, from_index) + { + let pk = record.at(self.primary_key_index).as_indexable().unwrap(); + let log_key = LogKey::new(segnum, index); + self.primary_memtable.set(&pk, &log_key); + + // Update from_index in case this is the last iteration: we need to know the next + // index that should be read on later invocations of refresh_indexes. + from_index = index + 1 + } + + // If there are still segments to read, set from_index to zero to read them + // from beginning. Otherwise we leave from_index as the index of the next record to read. + if segnum != to_segnum { + from_index = 0 + } + } + + self.refresh_next_logkey = LogKey::new(to_segnum, from_index); + + Ok(()) + } + /// Insert a record into the database. If the primary key value already exists, /// the existing record will be replaced by the supplied one. pub fn upsert(&mut self, record: &Record) -> Result<(), io::Error> { @@ -623,6 +671,8 @@ impl<Field: Eq + Clone + Debug> DB<Field> { self.active_metadata_file.unlock()?; + self.refresh_indexes()?; + Ok(()) } @@ -643,12 +693,13 @@ impl<Field: Eq + Clone + Debug> DB<Field> { self.active_data_file.unlock()?; let mut read_n = 0; - for (original_index, record) in forward_log_reader.enumerate() { - let primary_key = record + for (original_index, item) in forward_log_reader.enumerate() { + let primary_key = item + .record .at(self.primary_key_index) .as_indexable() .expect("Primary key was not indexable"); - map.insert(primary_key, (original_index, record)); + map.insert(primary_key, (original_index, item.record)); read_n += 1; } diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs index f9f1a88..5f1d054 100644 --- a/log_db/src/log_reader_forward.rs +++ b/log_db/src/log_reader_forward.rs @@ -7,7 +7,12 @@ pub struct ForwardLogReader { data_reader: io::BufReader<fs::File>, } -impl<'a> ForwardLogReader { +pub struct ForwardLogReaderItem { + pub record: Record, + pub index: u64, +} + +impl ForwardLogReader { pub fn new(metadata_file: fs::File, data_file: fs::File) -> ForwardLogReader { let mut ret = ForwardLogReader { metadata_reader: io::BufReader::new(metadata_file), @@ -21,8 +26,30 @@ impl<'a> ForwardLogReader { ret } - fn read_record(&mut self) -> Result<Option<Record>, io::Error> { + 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 + METADATA_ROW_LENGTH as u64 * index, + )) + .expect("Seek failed"); + + ret + } + + fn read_record(&mut self) -> Result<Option<ForwardLogReaderItem>, io::Error> { loop { + let pos = self.metadata_reader.stream_position()?; + let index = (pos - METADATA_FILE_HEADER_SIZE as u64) / METADATA_ROW_LENGTH as u64; + 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 { @@ -50,13 +77,13 @@ impl<'a> ForwardLogReader { self.data_reader.read_exact(&mut result_buf)?; let record = Record::deserialize(&result_buf); - return Ok(Some(record)); + return Ok(Some(ForwardLogReaderItem { record, index })); } } } impl Iterator for ForwardLogReader { - type Item = Record; + type Item = ForwardLogReaderItem; fn next(&mut self) -> Option<Self::Item> { match self.read_record() { @@ -93,7 +120,7 @@ mod tests { let first_record = forward_log_reader .next() .expect("Failed to read the first record"); - assert!(match first_record.values() { + assert!(match first_record.record.values() { [Value::Bytes(bytes)] => bytes.len() == 256, _ => false, }); |
