aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
-rw-r--r--log_db/src/common.rs29
-rw-r--r--log_db/src/lib.rs63
-rw-r--r--log_db/src/log_reader_forward.rs37
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,
});