aboutsummaryrefslogtreecommitdiffstats
path: root/src
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2024-10-02 12:26:20 +0300
committerJan Tuomi <jan@jantuomi.fi>2024-10-02 12:26:20 +0300
commit7b84f3676b37a85e84b1364e5252283b16f4cb8b (patch)
treee26163435330dceef4778413a9bba9cbb3234c10 /src
parentf58051f117056f5449cf97243042a2222d1cd318 (diff)
Refactor modules
Diffstat (limited to 'src')
-rw-r--r--src/common.rs7
-rw-r--r--src/lib.rs129
-rw-r--r--src/log_reader.rs108
-rw-r--r--src/primary_memtable.rs13
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,
diff --git a/src/lib.rs b/src/lib.rs
index 3aae5fb..4afeb4c 100644
--- a/src/lib.rs
+++ b/src/lib.rs
@@ -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(),