diff options
| -rw-r--r-- | src/common.rs | 7 | ||||
| -rw-r--r-- | src/lib.rs | 129 | ||||
| -rw-r--r-- | src/log_reader.rs | 108 | ||||
| -rw-r--r-- | src/primary_memtable.rs | 13 |
4 files changed, 123 insertions, 134 deletions
diff --git a/src/common.rs b/src/common.rs index b47a68b..f7aa755 100644 --- a/src/common.rs +++ b/src/common.rs @@ -16,6 +16,13 @@ pub const SEQ_RECORD_SEP: &[u8] = &[FIELD_SEPARATOR, FIELD_SEPARATOR, ESCAPE_CHA pub const SEQ_LIT_ESCAPE: &[u8] = &[ESCAPE_CHARACTER, ESCAPE_CHARACTER, ESCAPE_CHARACTER]; pub const SEQ_LIT_FIELD_SEP: &[u8] = &[ESCAPE_CHARACTER, FIELD_SEPARATOR, ESCAPE_CHARACTER]; +#[derive(Debug, Eq, PartialEq)] +pub enum SpecialSequence { + RecordSeparator, + LiteralFieldSeparator, + LiteralEscape, +} + #[derive(Debug, Clone, Eq, PartialEq)] pub enum MemtableEvictPolicy { LeastWritten, @@ -3,17 +3,18 @@ extern crate log; extern crate rev_buf_reader; mod common; +mod log_reader; mod primary_memtable; mod secondary_memtable; pub use common::*; use fs2::FileExt; +pub use log_reader::LogReader; use primary_memtable::PrimaryMemtable; -use rev_buf_reader::RevBufReader; use secondary_memtable::SecondaryMemtable; use std::fmt::Debug; use std::fs::{self}; -use std::io::{self, BufRead, Read, Seek, Write}; +use std::io::{self, Write}; use std::path::{Path, PathBuf}; pub struct ConfigBuilder<'a, Field: Eq + Clone + Debug> { @@ -130,18 +131,11 @@ struct Config<Field: Eq + Clone> { pub memtable_evict_policy: MemtableEvictPolicy, } -trait Memtable<Field: Eq + Clone + Debug, T: Debug> { - fn field(&self) -> &Field; - fn new(field: &Field, capacity: usize, evict_policy: MemtableEvictPolicy) -> Self; - fn set(&mut self, key: &IndexableValue, value: &T); - fn get(&mut self, key: &IndexableValue) -> Option<&T>; -} - pub struct DB<Field: Eq + Clone + Debug> { config: Config<Field>, log_path: PathBuf, primary_key_index: usize, - primary_memtable: PrimaryMemtable<Field>, + primary_memtable: PrimaryMemtable, secondary_memtables: Vec<SecondaryMemtable<Field>>, } @@ -204,7 +198,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } } let primary_memtable = PrimaryMemtable::new( - &config.primary_key, config.memtable_capacity, config.memtable_evict_policy.clone(), ); @@ -473,8 +466,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> { debug!("Lock acquired, searching log file for record"); - let mut log_reader = LogReader::new(&mut file)?; - let result = log_reader + let result = LogReader::new(&mut file)? .filter(|record| { let record_key = record.values[key_index] .as_indexable() @@ -499,114 +491,3 @@ impl<Field: Eq + Clone + Debug> DB<Field> { Ok(result) } } - -#[derive(Debug, Eq, PartialEq)] -enum SpecialSequence { - RecordSeparator, - LiteralFieldSeparator, - LiteralEscape, -} - -pub struct LogReader<'a> { - rev_reader: RevBufReader<&'a mut fs::File>, -} - -impl<'a> LogReader<'a> { - pub fn new(file: &mut fs::File) -> Result<LogReader, io::Error> { - let rev_reader = RevBufReader::new(file); - Ok(LogReader { rev_reader }) - } - - fn read_record(&mut self) -> Result<Option<Record>, io::Error> { - if self.rev_reader.stream_position()? == 0 { - return Ok(None); - } - - // Check that the record starts with the record separator - if self.read_special_sequence()? != SpecialSequence::RecordSeparator { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Record candidate does not end with record separator", - )); - } - - // The buffer that stores all the bytes of the record read so far in reverse order. - let mut result_buf: Vec<u8> = Vec::new(); - // The buffer that stores the bytes read from the file. - let mut read_buf: Vec<u8> = Vec::new(); - - loop { - read_buf.clear(); - self.rev_reader - .read_until(ESCAPE_CHARACTER, &mut read_buf)?; - - result_buf.extend(read_buf.iter().rev()); - - if self.rev_reader.stream_position()? == 0 { - // If we've reached the beginning of the file, we've read the entire record. - break; - } - - // Otherwise, we must have encountered an escape character. - match self.read_special_sequence()? { - SpecialSequence::RecordSeparator => { - // The record is complete, so we can break out of the loop. - // Move the cursor back to the beginning of the special sequence. - self.rev_reader.seek_relative(3)?; - break; - } - SpecialSequence::LiteralFieldSeparator => { - // The field separator is escaped, so we need to add it to the result buffer. - result_buf.push(FIELD_SEPARATOR); - } - SpecialSequence::LiteralEscape => { - // The escape character is escaped, so we need to add it to the result buffer. - result_buf.push(ESCAPE_CHARACTER); - } - } - } - - result_buf.reverse(); - let record = Record::deserialize(&result_buf); - Ok(Some(record)) - } - - fn read_special_sequence(&mut self) -> Result<SpecialSequence, io::Error> { - let mut special_buf: Vec<u8> = vec![0; 3]; - self.rev_reader.read_exact(&mut special_buf)?; - - match validate_special(&special_buf.as_slice()) { - Some(special) => Ok(special), - None => Err(io::Error::new( - io::ErrorKind::InvalidData, - "Not a special sequence", - )), - } - } -} - -impl Iterator for LogReader<'_> { - type Item = Record; - - fn next(&mut self) -> Option<Self::Item> { - match self.read_record() { - Ok(Some(record)) => Some(record), - Ok(None) => None, - Err(err) => panic!("Error reading record: {:?}", err), - } - } -} - -/// There are three special characters that need to be handled: -/// Here: SC = escape char, FS = field separator. -/// - FS FS SC -> actual record separator -/// - SC FS SC -> literal FS -/// - SC SC SC -> literal SC -fn validate_special(buf: &[u8]) -> Option<SpecialSequence> { - match buf { - SEQ_RECORD_SEP => Some(SpecialSequence::RecordSeparator), - SEQ_LIT_FIELD_SEP => Some(SpecialSequence::LiteralFieldSeparator), - SEQ_LIT_ESCAPE => Some(SpecialSequence::LiteralEscape), - _ => None, - } -} diff --git a/src/log_reader.rs b/src/log_reader.rs new file mode 100644 index 0000000..0ee0a81 --- /dev/null +++ b/src/log_reader.rs @@ -0,0 +1,108 @@ +use super::common::*; +use rev_buf_reader::RevBufReader; +use std::fs::{self}; +use std::io::{self, BufRead, Read, Seek}; + +pub struct LogReader<'a> { + rev_reader: RevBufReader<&'a mut fs::File>, +} + +impl<'a> LogReader<'a> { + pub fn new(file: &mut fs::File) -> Result<LogReader, io::Error> { + let rev_reader = RevBufReader::new(file); + Ok(LogReader { rev_reader }) + } + + fn read_record(&mut self) -> Result<Option<Record>, io::Error> { + if self.rev_reader.stream_position()? == 0 { + return Ok(None); + } + + // Check that the record starts with the record separator + if self.read_special_sequence()? != SpecialSequence::RecordSeparator { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "Record candidate does not end with record separator", + )); + } + + // The buffer that stores all the bytes of the record read so far in reverse order. + let mut result_buf: Vec<u8> = Vec::new(); + // The buffer that stores the bytes read from the file. + let mut read_buf: Vec<u8> = Vec::new(); + + loop { + read_buf.clear(); + self.rev_reader + .read_until(ESCAPE_CHARACTER, &mut read_buf)?; + + result_buf.extend(read_buf.iter().rev()); + + if self.rev_reader.stream_position()? == 0 { + // If we've reached the beginning of the file, we've read the entire record. + break; + } + + // Otherwise, we must have encountered an escape character. + match self.read_special_sequence()? { + SpecialSequence::RecordSeparator => { + // The record is complete, so we can break out of the loop. + // Move the cursor back to the beginning of the special sequence. + self.rev_reader.seek_relative(3)?; + break; + } + SpecialSequence::LiteralFieldSeparator => { + // The field separator is escaped, so we need to add it to the result buffer. + result_buf.push(FIELD_SEPARATOR); + } + SpecialSequence::LiteralEscape => { + // The escape character is escaped, so we need to add it to the result buffer. + result_buf.push(ESCAPE_CHARACTER); + } + } + } + + result_buf.reverse(); + let record = Record::deserialize(&result_buf); + Ok(Some(record)) + } + + fn read_special_sequence(&mut self) -> Result<SpecialSequence, io::Error> { + let mut special_buf: Vec<u8> = vec![0; 3]; + self.rev_reader.read_exact(&mut special_buf)?; + + match validate_special(&special_buf.as_slice()) { + Some(special) => Ok(special), + None => Err(io::Error::new( + io::ErrorKind::InvalidData, + "Not a special sequence", + )), + } + } +} + +impl Iterator for LogReader<'_> { + type Item = Record; + + fn next(&mut self) -> Option<Self::Item> { + match self.read_record() { + Ok(Some(record)) => Some(record), + Ok(None) => None, + Err(err) => panic!("Error reading record: {:?}", err), + } + } +} + +/// There are three special characters that need to be handled: +/// Here: SC = escape char, FS = field separator. +/// - FS FS SC -> actual record separator +/// - SC FS SC -> literal FS +/// - SC SC SC -> literal SC +fn validate_special(buf: &[u8]) -> Option<SpecialSequence> { + match buf { + SEQ_RECORD_SEP => Some(SpecialSequence::RecordSeparator), + SEQ_LIT_FIELD_SEP => Some(SpecialSequence::LiteralFieldSeparator), + SEQ_LIT_ESCAPE => Some(SpecialSequence::LiteralEscape), + _ => None, + } +} diff --git a/src/primary_memtable.rs b/src/primary_memtable.rs index ddca084..0c94ceb 100644 --- a/src/primary_memtable.rs +++ b/src/primary_memtable.rs @@ -1,10 +1,8 @@ use super::common::*; use priority_queue::PriorityQueue; use std::collections::BTreeMap; -use std::fmt::Debug; -pub struct PrimaryMemtable<Field: Eq + Clone + Debug> { - pub field: Field, +pub struct PrimaryMemtable { capacity: usize, /// Running counter of memtable operations, used as priority /// in evict_queue. @@ -16,14 +14,9 @@ pub struct PrimaryMemtable<Field: Eq + Clone + Debug> { evict_policy: MemtableEvictPolicy, } -impl<Field: Eq + Clone + Debug> PrimaryMemtable<Field> { - pub fn new( - field: &Field, - capacity: usize, - evict_policy: MemtableEvictPolicy, - ) -> PrimaryMemtable<Field> { +impl PrimaryMemtable { + pub fn new(capacity: usize, evict_policy: MemtableEvictPolicy) -> PrimaryMemtable { PrimaryMemtable { - field: field.clone(), capacity, n_operations: 0, records: BTreeMap::new(), |
