aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2025-01-17 15:25:46 +0200
committerJan Tuomi <jan@jantuomi.fi>2025-01-17 15:25:46 +0200
commit98430c35eeb27b1bb3c7ac7a9e30a6918514358a (patch)
tree5a319c629d508d91ad4dcce325aa7b1956afdfad /log_db/src
parent953f6d0b2a8c85f08ade5cc5526d4c44f2abcddd (diff)
Refactor DB into DB and Engine
Diffstat (limited to 'log_db/src')
-rw-r--r--log_db/src/common.rs30
-rw-r--r--log_db/src/config.rs111
-rw-r--r--log_db/src/engine.rs747
-rw-r--r--log_db/src/lib.rs887
4 files changed, 908 insertions, 867 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs
index 5539689..09de2d7 100644
--- a/log_db/src/common.rs
+++ b/log_db/src/common.rs
@@ -3,7 +3,6 @@ use once_cell::sync::Lazy;
use rust_decimal::Decimal;
use std::cmp::Ordering;
use std::collections::HashSet;
-use std::fmt::Display;
use std::fs::{self, metadata, File};
use std::io::{self, Read, Seek, SeekFrom, Write};
use std::ops::{Bound, RangeBounds};
@@ -201,33 +200,6 @@ impl MetadataHeader {
}
}
-#[derive(Debug, Clone, Eq, PartialEq)]
-pub enum ReadConsistency {
- /// Reads by client A are guaranteed to see writes by themselves and any writes by other clients B
- /// that were done before last index refresh. You must call `refresh_indexes()` manually to refresh indexes.
- Eventual,
- /// Reads by client A are guaranteed to see all writes. This is slower: all reads must first
- /// refresh indexes.
- Strong,
-}
-
-#[derive(Debug, Clone, Eq, PartialEq)]
-pub enum WriteDurability {
- /// Changes are written to the OS write buffer but not immediately synced to disk.
- /// This is generally recommended. Most OSes will sync the write buffer to disk within a few seconds.
- Flush,
- /// Changes are written to the OS write buffer and synced to disk immediately.
- /// Offers the best durability guarantees but is a lot slower.
- FlushSync,
-}
-
-impl Display for WriteDurability {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> {
- write!(f, "{:?}", self)?;
- Ok(())
- }
-}
-
#[derive(Debug, Clone, Ord, PartialOrd, Eq, PartialEq, Hash)]
pub enum IndexableValue {
Null,
@@ -680,8 +652,8 @@ pub fn ensure_active_metadata_is_valid(
const LOCK_WAIT_MAX_MS: u64 = 1000;
- let lock_request_path = data_dir.join(EXCL_LOCK_REQUEST_FILENAME);
pub fn is_exclusive_lock_requested(data_dir: &Path) -> DBResult<bool> {
+ let lock_request_path = data_dir.join(EXCL_LOCK_REQUEST_FILENAME);
let lock_request_file = fs::OpenOptions::new()
.create(true)
.write(true) // When requesting a lock, we need to have either read or write permissions
diff --git a/log_db/src/config.rs b/log_db/src/config.rs
new file mode 100644
index 0000000..daf0983
--- /dev/null
+++ b/log_db/src/config.rs
@@ -0,0 +1,111 @@
+use super::*;
+
+pub struct ConfigBuilder<R: Recordable> {
+ data_dir: Option<String>,
+ segment_size: Option<usize>,
+ write_durability: Option<WriteDurability>,
+ read_consistency: Option<ReadConsistency>,
+ _marker: PhantomData<R>,
+}
+
+impl<R: Recordable> ConfigBuilder<R> {
+ pub fn new() -> ConfigBuilder<R> {
+ ConfigBuilder {
+ data_dir: None,
+ segment_size: None,
+ write_durability: None,
+ read_consistency: None,
+ _marker: PhantomData,
+ }
+ }
+
+ /// The directory where the database will store its data.
+ pub fn data_dir(&mut self, data_dir: &str) -> &mut Self {
+ self.data_dir = Some(data_dir.to_string());
+ self
+ }
+
+ /// The maximum size of a segment file in bytes.
+ /// Once a segment file reaches this size, it can be closed, rotated and compacted.
+ /// Note that this is not a hard limit: if `db.do_maintenance_tasks()` is not called,
+ /// the segment file may continue to grow.
+ pub fn segment_size(&mut self, segment_size: usize) -> &mut Self {
+ self.segment_size = Some(segment_size);
+ self
+ }
+
+ /// The write durability policy for the database.
+ /// This determines how writes are persisted to disk.
+ /// The default is WriteDurability::Flush.
+ pub fn write_durability(&mut self, write_durability: WriteDurability) -> &mut Self {
+ self.write_durability = Some(write_durability);
+ self
+ }
+
+ /// The read consistency policy for the database.
+ /// This determines how recent writes are visible when reading.
+ /// See individual `ReadConsistency` enum values for more information.
+ /// The default is ReadConsistency::Strong.
+ pub fn read_consistency(&mut self, read_consistency: ReadConsistency) -> &mut Self {
+ self.read_consistency = Some(read_consistency);
+ self
+ }
+
+ pub fn initialize(&self) -> DBResult<DB<R>> {
+ let config = Config {
+ fields: R::schema(),
+ primary_key: R::primary_key(),
+ secondary_keys: R::secondary_keys(),
+ data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()),
+ segment_size: self.segment_size.unwrap_or(4 * 1024 * 1024), // 4MB
+ write_durability: self
+ .write_durability
+ .clone()
+ .unwrap_or(WriteDurability::Flush),
+ read_consistency: self
+ .read_consistency
+ .clone()
+ .unwrap_or(ReadConsistency::Strong),
+ };
+
+ DB::initialize(config)
+ }
+}
+
+#[derive(Clone)]
+pub struct Config<R: Recordable> {
+ pub fields: Vec<(R::Field, ValueType)>,
+ pub primary_key: R::Field,
+ pub secondary_keys: Vec<R::Field>,
+ pub data_dir: String,
+ pub segment_size: usize,
+ pub write_durability: WriteDurability,
+ pub read_consistency: ReadConsistency,
+}
+
+#[derive(Debug, Clone, Eq, PartialEq)]
+pub enum ReadConsistency {
+ /// Reads by client A are guaranteed to see writes by themselves and any writes by other clients B
+ /// that were done before last index refresh. You must call `refresh_indexes()` manually to refresh indexes.
+ Eventual,
+ /// Reads by client A are guaranteed to see all writes. This is slower: all reads must first
+ /// refresh indexes.
+ Strong,
+}
+
+#[derive(Debug, Clone, Eq, PartialEq)]
+pub enum WriteDurability {
+ /// Changes are written to the OS write buffer but not immediately synced to disk.
+ /// This is generally recommended. Most OSes will sync the write buffer to disk within a few seconds.
+ Flush,
+ /// Changes are written to the OS write buffer and synced to disk immediately.
+ /// Offers the best durability guarantees but is a lot slower.
+ FlushSync,
+}
+
+impl Display for WriteDurability {
+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> {
+ write!(f, "{:?}", self)?;
+ Ok(())
+ }
+}
diff --git a/log_db/src/engine.rs b/log_db/src/engine.rs
new file mode 100644
index 0000000..ef73db2
--- /dev/null
+++ b/log_db/src/engine.rs
@@ -0,0 +1,747 @@
+use super::*;
+
+pub struct Engine<R: Recordable> {
+ pub config: Config<R>,
+
+ data_dir: PathBuf,
+ active_metadata_file: fs::File,
+ active_data_file: fs::File,
+ primary_key_index: usize,
+ refresh_next_logkey: LogKey,
+
+ // TODO: these could be made private. Currently they are public for testing in lib.rs.
+ pub primary_memtable: PrimaryMemtable,
+ pub secondary_memtables: Vec<SecondaryMemtable>,
+}
+
+impl<R: Recordable> Engine<R> {
+ pub fn initialize(config: Config<R>) -> DBResult<Engine<R>> {
+ info!("Initializing DB...");
+ // If data_dir does not exist or is empty, create it and any necessary files
+ // After creation, the directory should always be in a complete state
+ // without missing files.
+ // A tempdir-move strategy is used to achieve one-phase commit.
+
+ // Ensure the data directory exists
+ let data_dir = Path::new(&config.data_dir).to_path_buf();
+ match fs::create_dir(&data_dir) {
+ Ok(_) => {}
+ Err(e) => {
+ if e.kind() != io::ErrorKind::AlreadyExists {
+ return Err(DBError::IOError(e));
+ }
+ }
+ }
+
+ // Create an initialize lock file to prevent multiple concurrent initializations
+ let init_lock_file = fs::OpenOptions::new()
+ .create(true)
+ .write(true)
+ .open(&data_dir.join(INIT_LOCK_FILENAME))?;
+
+ init_lock_file.lock_exclusive()?;
+
+ // We have acquired the lock, check if the data directory is in a complete state
+ // If not, initialize it, otherwise skip.
+ if !fs::exists(data_dir.join(ACTIVE_SYMLINK_FILENAME))? {
+ let (segment_uuid, _) = create_segment_data_file(&data_dir)?;
+ let (segment_num, _) = create_segment_metadata_file(&data_dir, &segment_uuid)?;
+ set_active_segment(&data_dir, segment_num)?;
+
+ // Create the exclusive lock request file
+ fs::OpenOptions::new()
+ .create(true)
+ .write(true)
+ .open(data_dir.join(EXCL_LOCK_REQUEST_FILENAME))?;
+ }
+
+ init_lock_file.unlock()?;
+
+ // Calculate the index of the primary value in a record
+ let primary_key_index = config
+ .fields
+ .iter()
+ .position(|(field, _)| field == &config.primary_key)
+ .ok_or(io::Error::new(
+ io::ErrorKind::InvalidInput,
+ "Primary key not found in schema after initialize",
+ ))?;
+
+ // Join primary key and secondary keys vec into a single vec
+ let mut all_keys = vec![&config.primary_key];
+ all_keys.extend(&config.secondary_keys);
+
+ // If any of the keys is not in the schema or
+ // is not an IndexableValue, return an error
+ for &key in &all_keys {
+ let (_, value_type) = config.fields.iter().find(|(field, _)| field == key).ok_or(
+ DBError::ValidationError("Key must be present in the field schema".to_owned()),
+ )?;
+
+ match value_type.prim_value_type {
+ PrimValueType::Int | PrimValueType::String => {}
+ _ => return Err(DBError::ValidationError("Key must be indexable".to_owned())),
+ }
+ }
+ let primary_memtable = PrimaryMemtable::new();
+ let secondary_memtables = config
+ .secondary_keys
+ .iter()
+ .map(|_| SecondaryMemtable::new())
+ .collect();
+
+ let active_symlink = Path::new(&config.data_dir).join(ACTIVE_SYMLINK_FILENAME);
+
+ let active_target = fs::read_link(&active_symlink)?;
+ let active_metadata_path = Path::new(&config.data_dir).join(active_target);
+ let mut active_metadata_file = APPEND_MODE.open(&active_metadata_path)?;
+
+ let active_metadata_header = read_metadata_header(&mut active_metadata_file)?;
+ validate_metadata_header(&active_metadata_header)?;
+
+ let active_data_path =
+ Path::new(&config.data_dir).join(active_metadata_header.uuid.to_string());
+ let active_data_file = APPEND_MODE.open(&active_data_path)?;
+
+ let mut engine = Engine::<R> {
+ config,
+ data_dir,
+ active_metadata_file,
+ active_data_file,
+ primary_key_index,
+ primary_memtable,
+ secondary_memtables,
+ refresh_next_logkey: LogKey::new(1, 0),
+ };
+
+ info!("Rebuilding memtable indexes...");
+ engine.refresh_indexes()?;
+
+ info!("Database ready.");
+
+ Ok(engine)
+ }
+
+ pub fn refresh_indexes(&mut self) -> DBResult<()> {
+ 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)?;
+
+ let metadata_len = metadata_file.seek(SeekFrom::End(0))?;
+ if (metadata_len - METADATA_FILE_HEADER_SIZE as u64) % METADATA_ROW_LENGTH as u64 != 0 {
+ return Err(DBError::ConsistencyError(format!(
+ "Metadata file {} has invalid size: {}",
+ metadata_path.display(),
+ metadata_len
+ )));
+ }
+
+ 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 log_key = LogKey::new(segnum, index);
+
+ if record.tombstone {
+ self.remove_record_from_memtables(&record);
+ } else {
+ self.insert_record_to_memtables(log_key, record);
+ }
+
+ // 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(())
+ }
+
+ fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Record) {
+ for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() {
+ let secondary_memtable = &mut self.secondary_memtables[sk_index];
+ let sk_field_index = self
+ .config
+ .fields
+ .iter()
+ .position(|(f, _)| sk_field == f)
+ .unwrap();
+ let sk = record.at(sk_field_index).as_indexable().unwrap();
+
+ secondary_memtable.set(sk, log_key.clone());
+ }
+
+ // Doing this last because this moves log_key
+ let pk = record.at(self.primary_key_index).as_indexable().unwrap();
+ self.primary_memtable.set(pk, log_key);
+ }
+
+ fn remove_record_from_memtables(&mut self, record: &Record) {
+ let pk = record.at(self.primary_key_index).as_indexable().unwrap();
+
+ if let Some(plk) = self.primary_memtable.remove(&pk) {
+ for (sk_index, sk_field) in self.config.secondary_keys.iter_mut().enumerate() {
+ let secondary_memtable = &mut self.secondary_memtables[sk_index];
+ let sk_field_index = self
+ .config
+ .fields
+ .iter()
+ .position(|(f, _)| sk_field == f)
+ .unwrap();
+ let sk = record.at(sk_field_index).as_indexable().unwrap();
+
+ secondary_memtable.remove(&sk, &plk);
+ }
+ }
+ }
+
+ pub fn batch_upsert_records(&mut self, records: impl Iterator<Item = Record>) -> DBResult<()> {
+ debug!("Opening file in append mode and acquiring exclusive lock...");
+
+ // Acquire an exclusive lock for writing
+ request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
+
+ if !self.ensure_metadata_file_is_active()?
+ || !ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?
+ {
+ // The log file has been rotated, so we must try again
+ self.active_metadata_file.unlock()?;
+ return self.batch_upsert_records(records);
+ }
+
+ self.active_data_file.lock_exclusive()?;
+
+ let active_symlink_path = self.data_dir.join(ACTIVE_SYMLINK_FILENAME);
+ let active_target = fs::read_link(active_symlink_path)?;
+ let segment_num = parse_segment_number(&active_target)?;
+
+ debug!("Exclusive lock acquired, appending to log file");
+
+ let mut serialized_data: Vec<u8> = vec![];
+ let mut serialized_metadata: Vec<u8> = vec![];
+ let mut pending_memtable_insertions: Vec<(LogKey, Record)> = vec![];
+ for record in records {
+ // Write the record to the log
+ let serialized = &record.serialize();
+ let record_offset = self.active_data_file.seek(SeekFrom::End(0))?;
+ let record_length = serialized.len() as u64;
+ assert!(record_length > 0);
+
+ serialized_data.extend(serialized);
+
+ let metadata_pos = self.active_metadata_file.seek(SeekFrom::End(0))?;
+ let metadata_index =
+ (metadata_pos - METADATA_FILE_HEADER_SIZE as u64) / METADATA_ROW_LENGTH as u64;
+
+ // Write the record metadata to the metadata file
+ let mut metadata_buf = vec![];
+ metadata_buf.extend(record_offset.to_be_bytes().into_iter());
+ metadata_buf.extend(record_length.to_be_bytes().into_iter());
+
+ assert_eq!(metadata_buf.len(), 16);
+
+ serialized_metadata.extend(metadata_buf);
+
+ let log_key = LogKey::new(segment_num, metadata_index);
+
+ pending_memtable_insertions.push((log_key, record));
+ }
+
+ self.active_data_file.write_all(&serialized_data)?;
+ self.active_metadata_file.write_all(&serialized_metadata)?;
+
+ // Flush and sync data and metadata to disk
+ if self.config.write_durability == WriteDurability::Flush {
+ self.active_data_file.flush()?;
+ self.active_metadata_file.flush()?;
+ } else if self.config.write_durability == WriteDurability::FlushSync {
+ self.active_data_file.flush()?;
+ self.active_data_file.sync_all()?;
+ self.active_metadata_file.flush()?;
+ self.active_metadata_file.sync_all()?;
+ }
+
+ debug!("Records appended to log file, releasing locks");
+
+ // Manually release the locks because the file handles are left open
+ self.active_data_file.unlock()?;
+ self.active_metadata_file.unlock()?;
+
+ for (log_key, record) in pending_memtable_insertions {
+ self.insert_record_to_memtables(log_key, record);
+ }
+
+ Ok(())
+ }
+
+ pub fn batch_find_by_records<'a>(
+ &mut self,
+ field: &R::Field,
+ values: impl Iterator<Item = &'a Value>,
+ ) -> DBResult<Vec<(usize, Record)>> {
+ let field_type = self.get_field_type(field).ok_or(DBError::ValidationError(
+ "Field not found in schema".to_owned(),
+ ))?;
+
+ let indexables = values
+ .map(|value| {
+ if type_check(&value, &field_type) {
+ value.as_indexable().ok_or(DBError::ValidationError(
+ "Queried value must be indexable".to_owned(),
+ ))
+ } else {
+ Err(DBError::ValidationError(format!(
+ "Queried value {:?} does not match key type: {:?}",
+ value, field_type
+ )))
+ }
+ })
+ .collect::<DBResult<Vec<IndexableValue>>>()?;
+
+ // Otherwise, continue with querying secondary indexes.
+ debug!(
+ "Finding all records with fields {:?} = {:?}",
+ field, indexables
+ );
+
+ if self.config.read_consistency == ReadConsistency::Strong {
+ self.refresh_indexes()?;
+ }
+
+ let log_key_batches = indexables
+ .into_iter()
+ .map(|query_key| {
+ if field == &self.config.primary_key {
+ let opt = self.primary_memtable.get(&query_key);
+ let log_keys = match opt {
+ Some(log_key) => vec![log_key],
+ None => vec![],
+ };
+ Ok(log_keys)
+ } else {
+ let smemtable_index = match get_secondary_memtable_index_by_field(
+ &self.config.secondary_keys,
+ field,
+ ) {
+ Some(index) => index,
+ None => {
+ return Err(DBError::ValidationError(
+ "Cannot find_by by non-indexed key".to_owned(),
+ ))
+ }
+ };
+
+ let log_keys = self.secondary_memtables[smemtable_index]
+ .find_by(&query_key)
+ .into_iter()
+ .collect();
+ Ok(log_keys)
+ }
+ })
+ .collect::<DBResult<Vec<Vec<&LogKey>>>>()?;
+
+ debug!("Found log keys in memtable: {:?}", log_key_batches);
+
+ let mut tagged = vec![];
+ let mut tag: usize = 0;
+ for batch in log_key_batches {
+ let mapped = batch.into_iter().map(|log_key| (tag, log_key));
+ tagged.extend(mapped);
+ tag += 1;
+ }
+
+ let tagged_records = self.read_tagged_log_keys(tagged.into_iter())?;
+
+ debug!("Read {} records", tagged_records.len());
+
+ Ok(tagged_records)
+ }
+
+ /// Read records from segment files based on log keys.
+ /// The log keys are accompanied by an integer tag that can be used to identify and group them later.
+ fn read_tagged_log_keys<'a>(
+ &self,
+ log_keys: impl Iterator<Item = (usize, &'a LogKey)>,
+ ) -> DBResult<Vec<(usize, Record)>> {
+ let mut records = vec![];
+ let mut log_keys_map = BTreeMap::new();
+
+ for (tag, log_key) in log_keys {
+ if !log_keys_map.contains_key(&log_key.segment_num()) {
+ log_keys_map.insert(log_key.segment_num(), vec![(tag, log_key.index())]);
+ } else {
+ log_keys_map
+ .get_mut(&log_key.segment_num())
+ .unwrap()
+ .push((tag, log_key.index()));
+ }
+ }
+
+ for (segment_num, mut segment_indexes) in log_keys_map {
+ segment_indexes.sort_unstable();
+
+ let metadata_path = &self.data_dir.join(metadata_filename(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)?;
+
+ let data_path = &self.data_dir.join(metadata_header.uuid.to_string());
+ let mut data_file = READ_MODE.open(&data_path)?;
+
+ data_file.lock_shared()?;
+
+ let header_size = METADATA_FILE_HEADER_SIZE as i64;
+ let row_length = METADATA_ROW_LENGTH as i64;
+ let mut current_metadata_offset = header_size;
+ for (tag, segment_index) in segment_indexes {
+ let new_metadata_offset = header_size + segment_index as i64 * row_length;
+ metadata_file.seek_relative(new_metadata_offset - current_metadata_offset)?;
+
+ let mut metadata_buf = [0; METADATA_ROW_LENGTH];
+ metadata_file.read_exact(&mut metadata_buf)?;
+
+ let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap());
+ let data_length = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap());
+ assert!(data_length > 0);
+
+ data_file.seek(SeekFrom::Start(data_offset))?;
+
+ let mut data_buf = vec![0; data_length as usize];
+ data_file.read_exact(&mut data_buf)?;
+
+ let record = Record::deserialize(&data_buf);
+ records.push((tag, record));
+
+ current_metadata_offset = new_metadata_offset + row_length;
+ }
+
+ metadata_file.unlock()?;
+ data_file.unlock()?;
+ }
+
+ Ok(records)
+ }
+
+ pub fn range_by_records<B: RangeBounds<Value>>(
+ &mut self,
+ field: &R::Field,
+ range: B,
+ ) -> DBResult<Vec<Record>> {
+ fn range_bound_to_indexable(
+ bound: Bound<&Value>,
+ field_type: &ValueType,
+ ) -> DBResult<Bound<IndexableValue>> {
+ fn convert(value: &Value, field_type: &ValueType) -> DBResult<IndexableValue> {
+ if !type_check(&value, field_type) {
+ return Err(DBError::ValidationError(format!(
+ "Queried value does not match type: {:?}",
+ field_type
+ )));
+ }
+ value.as_indexable().ok_or(DBError::ValidationError(
+ "Queried value must be indexable".to_owned(),
+ ))
+ }
+
+ match bound {
+ Bound::Included(value) => convert(value, field_type).map(Bound::Included),
+ Bound::Excluded(value) => convert(value, field_type).map(Bound::Excluded),
+ Bound::Unbounded => Ok(Bound::Unbounded),
+ }
+ }
+
+ let field_type = self.get_field_type(field).ok_or(DBError::ValidationError(
+ "Field not found in schema".to_owned(),
+ ))?;
+
+ let start_indexable = range_bound_to_indexable(range.start_bound(), field_type)?;
+ let end_indexable = range_bound_to_indexable(range.end_bound(), field_type)?;
+
+ let indexable_bounds = OwnedBounds::new(start_indexable, end_indexable);
+
+ if self.config.read_consistency == ReadConsistency::Strong {
+ self.refresh_indexes()?;
+ }
+
+ let log_keys = if field == &self.config.primary_key {
+ self.primary_memtable.range(indexable_bounds)
+ } else {
+ let index = get_secondary_memtable_index_by_field(&self.config.secondary_keys, field)
+ .ok_or_else(|| {
+ DBError::ValidationError("Cannot range_by by non-indexed key".to_owned())
+ })?;
+
+ self.secondary_memtables[index].range(indexable_bounds)
+ };
+
+ let log_key_batches = log_keys.into_iter().map(|log_key| (0, log_key));
+
+ let tagged_records = self.read_tagged_log_keys(log_key_batches);
+
+ Ok(tagged_records?.into_iter().map(|(_, rec)| rec).collect())
+ }
+
+ /// Ensures that the `self.metadata_file` and `self.data_file` handles are still pointing to the correct files.
+ /// If the segment has been rotated, the handle will be closed and reopened.
+ /// Returns `false` if the file has been rotated and the handle has been reopened, `true` otherwise.
+ fn ensure_metadata_file_is_active(&mut self) -> DBResult<bool> {
+ let active_target = fs::read_link(&self.data_dir.join(ACTIVE_SYMLINK_FILENAME))?;
+ let active_metadata_path = &self.data_dir.join(active_target);
+
+ let correct = is_file_same_as_path(&self.active_metadata_file, &active_metadata_path)?;
+ if !correct {
+ debug!("Metadata file has been rotated. Reopening...");
+ let mut metadata_file = APPEND_MODE.open(&active_metadata_path)?;
+
+ request_shared_lock(&self.data_dir, &mut metadata_file)?;
+
+ let metadata_header = read_metadata_header(&mut self.active_metadata_file)?;
+
+ validate_metadata_header(&metadata_header)?;
+
+ let data_file_path = &self.data_dir.join(metadata_header.uuid.to_string());
+
+ self.active_metadata_file = metadata_file;
+ self.active_data_file = APPEND_MODE.open(&data_file_path)?;
+
+ return Ok(false);
+ } else {
+ return Ok(true);
+ }
+ }
+
+ pub fn delete_by_field(&mut self, field: &R::Field, value: &Value) -> DBResult<Vec<Record>> {
+ let value_batch = std::iter::once(value);
+ let recs: Vec<Record> = self
+ .batch_find_by_records(field, value_batch)?
+ .into_iter()
+ .map(|(_, mut rec)| {
+ rec.tombstone = true;
+ rec
+ })
+ .collect();
+
+ request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
+ self.active_data_file.lock_exclusive()?;
+
+ for record in &recs {
+ let record_serialized = record.serialize();
+
+ let offset = self.active_data_file.seek(SeekFrom::End(0))?;
+ let length = record_serialized.len() as u64;
+
+ self.active_data_file.write_all(&record_serialized)?;
+
+ // Flush and sync data to disk
+ if self.config.write_durability == WriteDurability::Flush {
+ self.active_data_file.flush()?;
+ }
+ if self.config.write_durability == WriteDurability::FlushSync {
+ self.active_data_file.flush()?;
+ self.active_data_file.sync_all()?;
+ }
+
+ let mut metadata_entry = vec![];
+ metadata_entry.extend(offset.to_be_bytes().into_iter());
+ metadata_entry.extend(length.to_be_bytes().into_iter());
+
+ self.active_metadata_file.write_all(&metadata_entry)?;
+
+ // Flush and sync metadata to disk
+ if self.config.write_durability == WriteDurability::Flush {
+ self.active_metadata_file.flush()?;
+ }
+ if self.config.write_durability == WriteDurability::FlushSync {
+ self.active_metadata_file.flush()?;
+ self.active_metadata_file.sync_all()?;
+ }
+
+ self.remove_record_from_memtables(&record);
+ }
+
+ self.active_metadata_file.unlock()?;
+ self.active_data_file.unlock()?;
+
+ debug!("Records deleted");
+
+ Ok(recs)
+ }
+
+ pub fn do_maintenance_tasks(&mut self) -> DBResult<()> {
+ request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
+
+ ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?;
+
+ let metadata_size = self.active_metadata_file.seek(SeekFrom::End(0))?;
+ if metadata_size >= self.config.segment_size as u64 {
+ self.rotate_and_compact()?;
+ }
+
+ self.active_metadata_file.unlock()?;
+
+ Ok(())
+ }
+
+ fn rotate_and_compact(&mut self) -> DBResult<()> {
+ debug!("Active log size exceeds threshold, starting rotation and compaction...");
+
+ self.active_data_file.lock_shared()?;
+ let original_data_len = self.active_data_file.seek(SeekFrom::End(0))?;
+
+ let active_target = fs::read_link(&self.data_dir.join(ACTIVE_SYMLINK_FILENAME))?;
+ let active_num = parse_segment_number(&active_target)?;
+
+ debug!("Reading segment data into a BTreeMap");
+ let mut pk_to_item_map: BTreeMap<&IndexableValue, &Record> = BTreeMap::new();
+ let forward_read_items: Vec<(IndexableValue, Record)> = ForwardLogReader::new(
+ self.active_metadata_file.try_clone()?,
+ self.active_data_file.try_clone()?,
+ )
+ .map(|item| {
+ (
+ item.record
+ .at(self.primary_key_index)
+ .as_indexable()
+ .expect("Primary key was not indexable"),
+ item.record,
+ )
+ })
+ .collect();
+
+ self.active_data_file.unlock()?;
+
+ for (pk, record) in forward_read_items.iter() {
+ pk_to_item_map.insert(pk, record);
+ }
+
+ debug!(
+ "Read {} records, out of which {} were unique",
+ forward_read_items.len(),
+ pk_to_item_map.len()
+ );
+
+ // Create a new log data file and write it
+ debug!("Opening new data file and writing compacted data");
+ let (new_data_uuid, new_data_path) = create_segment_data_file(&self.data_dir)?;
+ let mut new_data_file = APPEND_MODE.open(&new_data_path)?;
+
+ let mut pk_to_data_map = BTreeMap::new();
+ let mut offset = 0u64;
+ for (pk, record) in pk_to_item_map.into_iter() {
+ let serialized = record.serialize();
+ let len = serialized.len() as u64;
+ new_data_file.write_all(&serialized)?;
+
+ pk_to_data_map.insert(pk, (offset, len));
+ offset += len;
+ }
+
+ // Sync the data file to disk.
+ // This is fine to do without consulting WriteDurability because this is a one-off
+ // operation that is not part of the normal write path.
+ new_data_file.flush()?;
+ new_data_file.sync_all()?;
+
+ let final_data_len = new_data_file.seek(io::SeekFrom::End(0))?;
+ debug!(
+ "Wrote compacted data, reduced data size: {} -> {}",
+ original_data_len, final_data_len
+ );
+
+ // Create a new log metadata file and write it
+ debug!("Opening temp metadata file and writing pointers to compacted data file");
+ let temp_metadata_file = tempfile::NamedTempFile::new()?;
+ let temp_metadata_path = temp_metadata_file.as_ref();
+ let mut temp_metadata_file = WRITE_MODE.open(temp_metadata_path)?;
+
+ let metadata_header = MetadataHeader {
+ version: 1,
+ uuid: new_data_uuid,
+ };
+
+ temp_metadata_file.write_all(&metadata_header.serialize())?;
+
+ for (pk, _) in forward_read_items.iter() {
+ let (offset, len) = pk_to_data_map.get(&pk).unwrap();
+
+ let mut metadata_buf = vec![];
+ metadata_buf.extend(offset.to_be_bytes().into_iter());
+ metadata_buf.extend(len.to_be_bytes().into_iter());
+
+ temp_metadata_file.write_all(&metadata_buf)?;
+ }
+
+ // Sync the metadata file to disk, see comment above about sync.
+ temp_metadata_file.flush()?;
+ temp_metadata_file.sync_all()?;
+
+ debug!("Moving temporary files to their final locations");
+ let new_data_path = &self.data_dir.join(new_data_uuid.to_string());
+ let active_metadata_path = &self.data_dir.join(metadata_filename(active_num)); // overwrite active
+
+ fs::rename(&temp_metadata_path, &active_metadata_path)?;
+
+ debug!("Compaction complete, creating new segment");
+
+ let new_segment_num = active_num + 1;
+ let new_metadata_path = self.data_dir.join(metadata_filename(new_segment_num));
+ let mut new_metadata_file = APPEND_MODE.clone().create(true).open(&new_metadata_path)?;
+
+ let new_metadata_header = MetadataHeader {
+ version: 1,
+ uuid: new_data_uuid,
+ };
+
+ new_metadata_file.write_all(&new_metadata_header.serialize())?;
+
+ set_active_segment(&self.data_dir, new_segment_num)?;
+
+ // Old active metadata file should lose lock by RAII, or by
+ // the manual unlock call in the do_maintenance_tasks method.
+
+ self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?;
+ self.active_data_file = APPEND_MODE.open(&new_data_path)?;
+
+ // The new active log file is not locked by this client so it cannot be touched.
+ debug!(
+ "Active log file {} rotated and compacted, new segment: {}",
+ active_num, new_segment_num
+ );
+
+ Ok(())
+ }
+
+ #[inline]
+ fn get_field_type(&self, field: &R::Field) -> Option<&ValueType> {
+ self.config
+ .fields
+ .iter()
+ .find(|(f, _)| f == field)
+ .map(|(_, t)| t)
+ }
+}
diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs
index 97ea28f..9f2cac9 100644
--- a/log_db/src/lib.rs
+++ b/log_db/src/lib.rs
@@ -1,122 +1,38 @@
#[macro_use]
extern crate log;
-#[macro_use]
-mod common;
-mod log_reader_forward;
-mod memtable_primary;
-mod memtable_secondary;
-mod record;
-
-pub use common::*;
use fs2::FileExt;
-pub use log_reader_forward::ForwardLogReader;
-use log_reader_forward::ForwardLogReaderItem;
-use memtable_primary::PrimaryMemtable;
-use memtable_secondary::SecondaryMemtable;
-pub use record::Recordable;
-use record::*;
use std::collections::BTreeMap;
use std::fmt::Debug;
+use std::fmt::Display;
use std::fs::{self};
use std::io::{self, Read, Seek, SeekFrom, Write};
use std::marker::PhantomData;
use std::ops::*;
use std::path::{Path, PathBuf};
-pub struct ConfigBuilder<R: Recordable> {
- data_dir: Option<String>,
- segment_size: Option<usize>,
- write_durability: Option<WriteDurability>,
- read_consistency: Option<ReadConsistency>,
- _marker: PhantomData<R>,
-}
-
-impl<R: Recordable> ConfigBuilder<R> {
- pub fn new() -> ConfigBuilder<R> {
- ConfigBuilder {
- data_dir: None,
- segment_size: None,
- write_durability: None,
- read_consistency: None,
- _marker: PhantomData,
- }
- }
-
- /// The directory where the database will store its data.
- pub fn data_dir(&mut self, data_dir: &str) -> &mut Self {
- self.data_dir = Some(data_dir.to_string());
- self
- }
-
- /// The maximum size of a segment file in bytes.
- /// Once a segment file reaches this size, it can be closed, rotated and compacted.
- /// Note that this is not a hard limit: if `db.do_maintenance_tasks()` is not called,
- /// the segment file may continue to grow.
- pub fn segment_size(&mut self, segment_size: usize) -> &mut Self {
- self.segment_size = Some(segment_size);
- self
- }
-
- /// The write durability policy for the database.
- /// This determines how writes are persisted to disk.
- /// The default is WriteDurability::Flush.
- pub fn write_durability(&mut self, write_durability: WriteDurability) -> &mut Self {
- self.write_durability = Some(write_durability);
- self
- }
-
- /// The read consistency policy for the database.
- /// This determines how recent writes are visible when reading.
- /// See individual `ReadConsistency` enum values for more information.
- /// The default is ReadConsistency::Strong.
- pub fn read_consistency(&mut self, read_consistency: ReadConsistency) -> &mut Self {
- self.read_consistency = Some(read_consistency);
- self
- }
-
- pub fn initialize(&self) -> DBResult<DB<R>> {
- let config = Config {
- fields: R::schema(),
- primary_key: R::primary_key(),
- secondary_keys: R::secondary_keys(),
- data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()),
- segment_size: self.segment_size.unwrap_or(4 * 1024 * 1024), // 4MB
- write_durability: self
- .write_durability
- .clone()
- .unwrap_or(WriteDurability::Flush),
- read_consistency: self
- .read_consistency
- .clone()
- .unwrap_or(ReadConsistency::Strong),
- };
+#[macro_use]
+mod common;
+mod config;
+mod engine;
+mod log_reader_forward;
+mod memtable_primary;
+mod memtable_secondary;
+mod record;
- DB::initialize(config)
- }
-}
+pub use common::*;
+pub use config::{ReadConsistency, WriteDurability};
+pub use record::Recordable;
-#[derive(Clone)]
-struct Config<R: Recordable> {
- pub fields: Vec<(R::Field, ValueType)>,
- pub primary_key: R::Field,
- pub secondary_keys: Vec<R::Field>,
- pub data_dir: String,
- pub segment_size: usize,
- pub write_durability: WriteDurability,
- pub read_consistency: ReadConsistency,
-}
+use config::*;
+use engine::*;
+use log_reader_forward::*;
+use memtable_primary::PrimaryMemtable;
+use memtable_secondary::SecondaryMemtable;
+use record::*;
pub struct DB<R: Recordable> {
- config: Config<R>,
-
- data_dir: PathBuf,
- active_metadata_file: fs::File,
- active_data_file: fs::File,
- primary_key_index: usize,
- primary_memtable: PrimaryMemtable,
- secondary_memtables: Vec<SecondaryMemtable>,
- refresh_next_logkey: LogKey,
+ engine: Engine<R>,
}
impl<R: Recordable> DB<R> {
@@ -126,208 +42,8 @@ impl<R: Recordable> DB<R> {
}
fn initialize(config: Config<R>) -> DBResult<DB<R>> {
- info!("Initializing DB...");
- // If data_dir does not exist or is empty, create it and any necessary files
- // After creation, the directory should always be in a complete state
- // without missing files.
- // A tempdir-move strategy is used to achieve one-phase commit.
-
- // Ensure the data directory exists
- let data_dir = Path::new(&config.data_dir).to_path_buf();
- match fs::create_dir(&data_dir) {
- Ok(_) => {}
- Err(e) => {
- if e.kind() != io::ErrorKind::AlreadyExists {
- return Err(DBError::IOError(e));
- }
- }
- }
-
- // Create an initialize lock file to prevent multiple concurrent initializations
- let init_lock_file = fs::OpenOptions::new()
- .create(true)
- .write(true)
- .open(&data_dir.join(INIT_LOCK_FILENAME))?;
-
- init_lock_file.lock_exclusive()?;
-
- // We have acquired the lock, check if the data directory is in a complete state
- // If not, initialize it, otherwise skip.
- if !fs::exists(data_dir.join(ACTIVE_SYMLINK_FILENAME))? {
- let (segment_uuid, _) = create_segment_data_file(&data_dir)?;
- let (segment_num, _) = create_segment_metadata_file(&data_dir, &segment_uuid)?;
- set_active_segment(&data_dir, segment_num)?;
-
- // Create the exclusive lock request file
- fs::OpenOptions::new()
- .create(true)
- .write(true)
- .open(data_dir.join(EXCL_LOCK_REQUEST_FILENAME))?;
- }
-
- init_lock_file.unlock()?;
-
- // Calculate the index of the primary value in a record
- let primary_key_index = config
- .fields
- .iter()
- .position(|(field, _)| field == &config.primary_key)
- .ok_or(io::Error::new(
- io::ErrorKind::InvalidInput,
- "Primary key not found in schema after initialize",
- ))?;
-
- // Join primary key and secondary keys vec into a single vec
- let mut all_keys = vec![&config.primary_key];
- all_keys.extend(&config.secondary_keys);
-
- // If any of the keys is not in the schema or
- // is not an IndexableValue, return an error
- for &key in &all_keys {
- let (_, value_type) = config.fields.iter().find(|(field, _)| field == key).ok_or(
- DBError::ValidationError("Key must be present in the field schema".to_owned()),
- )?;
-
- match value_type.prim_value_type {
- PrimValueType::Int | PrimValueType::String => {}
- _ => return Err(DBError::ValidationError("Key must be indexable".to_owned())),
- }
- }
- let primary_memtable = PrimaryMemtable::new();
- let secondary_memtables = config
- .secondary_keys
- .iter()
- .map(|_| SecondaryMemtable::new())
- .collect();
-
- let active_symlink = Path::new(&config.data_dir).join(ACTIVE_SYMLINK_FILENAME);
-
- let active_target = fs::read_link(&active_symlink)?;
- let active_metadata_path = Path::new(&config.data_dir).join(active_target);
- let mut active_metadata_file = APPEND_MODE.open(&active_metadata_path)?;
-
- let active_metadata_header = read_metadata_header(&mut active_metadata_file)?;
- validate_metadata_header(&active_metadata_header)?;
-
- let active_data_path =
- Path::new(&config.data_dir).join(active_metadata_header.uuid.to_string());
- let active_data_file = APPEND_MODE.open(&active_data_path)?;
-
- let mut db = DB::<R> {
- config,
- data_dir,
- active_metadata_file,
- active_data_file,
- primary_key_index,
- primary_memtable,
- secondary_memtables,
- refresh_next_logkey: LogKey::new(1, 0),
- };
-
- info!("Rebuilding memtable indexes...");
- db.refresh_indexes()?;
-
- info!("Database ready.");
-
- Ok(db)
- }
-
- /// Refresh the in-memory indexes from the log files.
- /// This needs to only be called if the read consistency is set to `ReadConsistency::Eventual`.
- pub fn refresh_indexes(&mut self) -> DBResult<()> {
- 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)?;
-
- let metadata_len = metadata_file.seek(SeekFrom::End(0))?;
- if (metadata_len - METADATA_FILE_HEADER_SIZE as u64) % METADATA_ROW_LENGTH as u64 != 0 {
- return Err(DBError::ConsistencyError(format!(
- "Metadata file {} has invalid size: {}",
- metadata_path.display(),
- metadata_len
- )));
- }
-
- 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 log_key = LogKey::new(segnum, index);
-
- if record.tombstone {
- self.remove_record_from_memtables(&record);
- } else {
- self.insert_record_to_memtables(log_key, record);
- }
-
- // 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(())
- }
-
- fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Record) {
- for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() {
- let secondary_memtable = &mut self.secondary_memtables[sk_index];
- let sk_field_index = self
- .config
- .fields
- .iter()
- .position(|(f, _)| sk_field == f)
- .unwrap();
- let sk = record.at(sk_field_index).as_indexable().unwrap();
-
- secondary_memtable.set(sk, log_key.clone());
- }
-
- // Doing this last because this moves log_key
- let pk = record.at(self.primary_key_index).as_indexable().unwrap();
- self.primary_memtable.set(pk, log_key);
- }
-
- fn remove_record_from_memtables(&mut self, record: &Record) {
- let pk = record.at(self.primary_key_index).as_indexable().unwrap();
-
- if let Some(plk) = self.primary_memtable.remove(&pk) {
- for (sk_index, sk_field) in self.config.secondary_keys.iter_mut().enumerate() {
- let secondary_memtable = &mut self.secondary_memtables[sk_index];
- let sk_field_index = self
- .config
- .fields
- .iter()
- .position(|(f, _)| sk_field == f)
- .unwrap();
- let sk = record.at(sk_field_index).as_indexable().unwrap();
-
- secondary_memtable.remove(&sk, &plk);
- }
- }
+ let engine = Engine::initialize(config)?;
+ Ok(DB { engine })
}
/// Insert a record into the database. If the primary key value already exists,
@@ -336,10 +52,10 @@ impl<R: Recordable> DB<R> {
let record = Record::from(&recordable.into_record());
debug!("Upserting record: {:?}", record);
- record.validate(&self.config.fields)?;
+ record.validate(&self.engine.config.fields)?;
debug!("Record is valid");
- self.batch_upsert_records(std::iter::once(record))
+ self.engine.batch_upsert_records(std::iter::once(record))
}
/// Insert a batch of records into the database. If the primary key value for a record already exists,
@@ -352,97 +68,21 @@ impl<R: Recordable> DB<R> {
debug!("Batch upserting {} records", records.len());
for record in &records {
- record.validate(&self.config.fields)?;
+ record.validate(&self.engine.config.fields)?;
}
debug!("Records are valid");
- self.batch_upsert_records(records.into_iter())
- }
-
- fn batch_upsert_records(&mut self, records: impl Iterator<Item = Record>) -> DBResult<()> {
- debug!("Opening file in append mode and acquiring exclusive lock...");
-
- // Acquire an exclusive lock for writing
- request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
-
- if !self.ensure_metadata_file_is_active()?
- || !ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?
- {
- // The log file has been rotated, so we must try again
- self.active_metadata_file.unlock()?;
- return self.batch_upsert_records(records);
- }
-
- self.active_data_file.lock_exclusive()?;
-
- let active_symlink_path = self.data_dir.join(ACTIVE_SYMLINK_FILENAME);
- let active_target = fs::read_link(active_symlink_path)?;
- let segment_num = parse_segment_number(&active_target)?;
-
- debug!("Exclusive lock acquired, appending to log file");
-
- let mut serialized_data: Vec<u8> = vec![];
- let mut serialized_metadata: Vec<u8> = vec![];
- let mut pending_memtable_insertions: Vec<(LogKey, Record)> = vec![];
- for record in records {
- // Write the record to the log
- let serialized = &record.serialize();
- let record_offset = self.active_data_file.seek(SeekFrom::End(0))?;
- let record_length = serialized.len() as u64;
- assert!(record_length > 0);
-
- serialized_data.extend(serialized);
-
- let metadata_pos = self.active_metadata_file.seek(SeekFrom::End(0))?;
- let metadata_index =
- (metadata_pos - METADATA_FILE_HEADER_SIZE as u64) / METADATA_ROW_LENGTH as u64;
-
- // Write the record metadata to the metadata file
- let mut metadata_buf = vec![];
- metadata_buf.extend(record_offset.to_be_bytes().into_iter());
- metadata_buf.extend(record_length.to_be_bytes().into_iter());
-
- assert_eq!(metadata_buf.len(), 16);
-
- serialized_metadata.extend(metadata_buf);
-
- let log_key = LogKey::new(segment_num, metadata_index);
-
- pending_memtable_insertions.push((log_key, record));
- }
-
- self.active_data_file.write_all(&serialized_data)?;
- self.active_metadata_file.write_all(&serialized_metadata)?;
-
- // Flush and sync data and metadata to disk
- if self.config.write_durability == WriteDurability::Flush {
- self.active_data_file.flush()?;
- self.active_metadata_file.flush()?;
- } else if self.config.write_durability == WriteDurability::FlushSync {
- self.active_data_file.flush()?;
- self.active_data_file.sync_all()?;
- self.active_metadata_file.flush()?;
- self.active_metadata_file.sync_all()?;
- }
-
- debug!("Records appended to log file, releasing locks");
-
- // Manually release the locks because the file handles are left open
- self.active_data_file.unlock()?;
- self.active_metadata_file.unlock()?;
-
- for (log_key, record) in pending_memtable_insertions {
- self.insert_record_to_memtables(log_key, record);
- }
-
- Ok(())
+ self.engine.batch_upsert_records(records.into_iter())
}
/// Get a record by its primary index value.
/// E.g. `db.get(Value::Int(10))`.
pub fn get(&mut self, value: &Value) -> DBResult<Option<R>> {
let value_batch = std::iter::once(value);
- let records = self.batch_find_by_records(&self.config.primary_key.clone(), value_batch)?;
+ let records = self
+ .engine
+ // TODO: This clone is only here to appease the borrow checker
+ .batch_find_by_records(&self.engine.config.primary_key.clone(), value_batch)?;
assert!(records.len() <= 1);
Ok(records
@@ -456,6 +96,7 @@ impl<R: Recordable> DB<R> {
pub fn find_by(&mut self, field: &R::Field, value: &Value) -> DBResult<Vec<R>> {
let value_batch = std::iter::once(value);
Ok(self
+ .engine
.batch_find_by_records(field, value_batch)?
.into_iter()
.map(|(_, rec)| R::from_record(rec.values))
@@ -472,262 +113,26 @@ impl<R: Recordable> DB<R> {
values: &[Value],
) -> DBResult<Vec<(usize, R)>> {
Ok(self
+ .engine
.batch_find_by_records(field, values.iter())?
.into_iter()
.map(|(tag, rec)| (tag, R::from_record(rec.values)))
.collect())
}
- fn batch_find_by_records<'a>(
- &mut self,
- field: &R::Field,
- values: impl Iterator<Item = &'a Value>,
- ) -> DBResult<Vec<(usize, Record)>> {
- let field_type = self.get_field_type(field).ok_or(DBError::ValidationError(
- "Field not found in schema".to_owned(),
- ))?;
-
- let indexables = values
- .map(|value| {
- if type_check(&value, &field_type) {
- value.as_indexable().ok_or(DBError::ValidationError(
- "Queried value must be indexable".to_owned(),
- ))
- } else {
- Err(DBError::ValidationError(format!(
- "Queried value {:?} does not match key type: {:?}",
- value, field_type
- )))
- }
- })
- .collect::<DBResult<Vec<IndexableValue>>>()?;
-
- // Otherwise, continue with querying secondary indexes.
- debug!(
- "Finding all records with fields {:?} = {:?}",
- field, indexables
- );
-
- if self.config.read_consistency == ReadConsistency::Strong {
- self.refresh_indexes()?;
- }
-
- let log_key_batches = indexables
- .into_iter()
- .map(|query_key| {
- if field == &self.config.primary_key {
- let opt = self.primary_memtable.get(&query_key);
- let log_keys = match opt {
- Some(log_key) => vec![log_key],
- None => vec![],
- };
- Ok(log_keys)
- } else {
- let smemtable_index = match get_secondary_memtable_index_by_field(
- &self.config.secondary_keys,
- field,
- ) {
- Some(index) => index,
- None => {
- return Err(DBError::ValidationError(
- "Cannot find_by by non-indexed key".to_owned(),
- ))
- }
- };
-
- let log_keys = self.secondary_memtables[smemtable_index]
- .find_by(&query_key)
- .into_iter()
- .collect();
- Ok(log_keys)
- }
- })
- .collect::<DBResult<Vec<Vec<&LogKey>>>>()?;
-
- debug!("Found log keys in memtable: {:?}", log_key_batches);
-
- let mut tagged = vec![];
- let mut tag: usize = 0;
- for batch in log_key_batches {
- let mapped = batch.into_iter().map(|log_key| (tag, log_key));
- tagged.extend(mapped);
- tag += 1;
- }
-
- let tagged_records = self.read_tagged_log_keys(tagged.into_iter())?;
-
- debug!("Read {} records", tagged_records.len());
-
- Ok(tagged_records)
- }
-
- /// Read records from segment files based on log keys.
- /// The log keys are accompanied by an integer tag that can be used to identify and group them later.
- fn read_tagged_log_keys<'a>(
- &self,
- log_keys: impl Iterator<Item = (usize, &'a LogKey)>,
- ) -> DBResult<Vec<(usize, Record)>> {
- let mut records = vec![];
- let mut log_keys_map = BTreeMap::new();
-
- for (tag, log_key) in log_keys {
- if !log_keys_map.contains_key(&log_key.segment_num()) {
- log_keys_map.insert(log_key.segment_num(), vec![(tag, log_key.index())]);
- } else {
- log_keys_map
- .get_mut(&log_key.segment_num())
- .unwrap()
- .push((tag, log_key.index()));
- }
- }
-
- for (segment_num, mut segment_indexes) in log_keys_map {
- segment_indexes.sort_unstable();
-
- let metadata_path = &self.data_dir.join(metadata_filename(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)?;
-
- let data_path = &self.data_dir.join(metadata_header.uuid.to_string());
- let mut data_file = READ_MODE.open(&data_path)?;
-
- data_file.lock_shared()?;
-
- let header_size = METADATA_FILE_HEADER_SIZE as i64;
- let row_length = METADATA_ROW_LENGTH as i64;
- let mut current_metadata_offset = header_size;
- for (tag, segment_index) in segment_indexes {
- let new_metadata_offset = header_size + segment_index as i64 * row_length;
- metadata_file.seek_relative(new_metadata_offset - current_metadata_offset)?;
-
- let mut metadata_buf = [0; METADATA_ROW_LENGTH];
- metadata_file.read_exact(&mut metadata_buf)?;
-
- let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap());
- let data_length = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap());
- assert!(data_length > 0);
-
- data_file.seek(SeekFrom::Start(data_offset))?;
-
- let mut data_buf = vec![0; data_length as usize];
- data_file.read_exact(&mut data_buf)?;
-
- let record = Record::deserialize(&data_buf);
- records.push((tag, record));
-
- current_metadata_offset = new_metadata_offset + row_length;
- }
-
- metadata_file.unlock()?;
- data_file.unlock()?;
- }
-
- Ok(records)
- }
-
pub fn range_by<B: RangeBounds<Value>>(
&mut self,
field: &R::Field,
range: B,
) -> DBResult<Vec<R>> {
Ok(self
+ .engine
.range_by_records(field, range)?
.into_iter()
.map(|rec| R::from_record(rec.values))
.collect())
}
- fn range_by_records<B: RangeBounds<Value>>(
- &mut self,
- field: &R::Field,
- range: B,
- ) -> DBResult<Vec<Record>> {
- fn range_bound_to_indexable(
- bound: Bound<&Value>,
- field_type: &ValueType,
- ) -> DBResult<Bound<IndexableValue>> {
- fn convert(value: &Value, field_type: &ValueType) -> DBResult<IndexableValue> {
- if !type_check(&value, field_type) {
- return Err(DBError::ValidationError(format!(
- "Queried value does not match type: {:?}",
- field_type
- )));
- }
- value.as_indexable().ok_or(DBError::ValidationError(
- "Queried value must be indexable".to_owned(),
- ))
- }
-
- match bound {
- Bound::Included(value) => convert(value, field_type).map(Bound::Included),
- Bound::Excluded(value) => convert(value, field_type).map(Bound::Excluded),
- Bound::Unbounded => Ok(Bound::Unbounded),
- }
- }
-
- let field_type = self.get_field_type(field).ok_or(DBError::ValidationError(
- "Field not found in schema".to_owned(),
- ))?;
-
- let start_indexable = range_bound_to_indexable(range.start_bound(), field_type)?;
- let end_indexable = range_bound_to_indexable(range.end_bound(), field_type)?;
-
- let indexable_bounds = OwnedBounds::new(start_indexable, end_indexable);
-
- if self.config.read_consistency == ReadConsistency::Strong {
- self.refresh_indexes()?;
- }
-
- let log_keys = if field == &self.config.primary_key {
- self.primary_memtable.range(indexable_bounds)
- } else {
- let index = get_secondary_memtable_index_by_field(&self.config.secondary_keys, field)
- .ok_or_else(|| {
- DBError::ValidationError("Cannot range_by by non-indexed key".to_owned())
- })?;
-
- self.secondary_memtables[index].range(indexable_bounds)
- };
-
- let log_key_batches = log_keys.into_iter().map(|log_key| (0, log_key));
-
- let tagged_records = self.read_tagged_log_keys(log_key_batches);
-
- Ok(tagged_records?.into_iter().map(|(_, rec)| rec).collect())
- }
-
- /// Ensures that the `self.metadata_file` and `self.data_file` handles are still pointing to the correct files.
- /// If the segment has been rotated, the handle will be closed and reopened.
- /// Returns `false` if the file has been rotated and the handle has been reopened, `true` otherwise.
- fn ensure_metadata_file_is_active(&mut self) -> DBResult<bool> {
- let active_target = fs::read_link(&self.data_dir.join(ACTIVE_SYMLINK_FILENAME))?;
- let active_metadata_path = &self.data_dir.join(active_target);
-
- let correct = is_file_same_as_path(&self.active_metadata_file, &active_metadata_path)?;
- if !correct {
- debug!("Metadata file has been rotated. Reopening...");
- let mut metadata_file = APPEND_MODE.open(&active_metadata_path)?;
-
- request_shared_lock(&self.data_dir, &mut metadata_file)?;
-
- let metadata_header = read_metadata_header(&mut self.active_metadata_file)?;
-
- validate_metadata_header(&metadata_header)?;
-
- let data_file_path = &self.data_dir.join(metadata_header.uuid.to_string());
-
- self.active_metadata_file = metadata_file;
- self.active_data_file = APPEND_MODE.open(&data_file_path)?;
-
- return Ok(false);
- } else {
- return Ok(true);
- }
- }
-
/// Delete records by a field value.
/// E.g. `db.delete_by(Field::Name, "John")`, assuming `Field` is the DB field type and `Field::Name` is secondary indexed.
/// Returns a vector of deleted records. If no records were deleted, the vector will be empty.
@@ -735,7 +140,7 @@ impl<R: Recordable> DB<R> {
/// Deletion is done by marking the record as a tombstone. The record will still be present in the log file,
/// but will be ignored by reads. Upon compaction, tombstoned records will be removed.
pub fn delete_by(&mut self, field: &R::Field, value: &Value) -> DBResult<Vec<R>> {
- let recs = self.delete_by_field(field, value)?;
+ let recs = self.engine.delete_by_field(field, value)?;
Ok(recs
.into_iter()
@@ -745,7 +150,10 @@ impl<R: Recordable> DB<R> {
/// Delete record by primary key.
pub fn delete(&mut self, pk: &Value) -> DBResult<Option<R>> {
- let recs = self.delete_by_field(&self.config.primary_key.clone(), pk)?;
+ let recs = self
+ .engine
+ // TODO: This clone is only here to appease the borrow checker
+ .delete_by_field(&self.engine.config.primary_key.clone(), pk)?;
assert!(recs.len() <= 1);
Ok(recs
@@ -754,63 +162,6 @@ impl<R: Recordable> DB<R> {
.map(|rec| R::from_record(rec.values)))
}
- fn delete_by_field(&mut self, field: &R::Field, value: &Value) -> DBResult<Vec<Record>> {
- let value_batch = std::iter::once(value);
- let recs: Vec<Record> = self
- .batch_find_by_records(field, value_batch)?
- .into_iter()
- .map(|(_, mut rec)| {
- rec.tombstone = true;
- rec
- })
- .collect();
-
- request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
- self.active_data_file.lock_exclusive()?;
-
- for record in &recs {
- let record_serialized = record.serialize();
-
- let offset = self.active_data_file.seek(SeekFrom::End(0))?;
- let length = record_serialized.len() as u64;
-
- self.active_data_file.write_all(&record_serialized)?;
-
- // Flush and sync data to disk
- if self.config.write_durability == WriteDurability::Flush {
- self.active_data_file.flush()?;
- }
- if self.config.write_durability == WriteDurability::FlushSync {
- self.active_data_file.flush()?;
- self.active_data_file.sync_all()?;
- }
-
- let mut metadata_entry = vec![];
- metadata_entry.extend(offset.to_be_bytes().into_iter());
- metadata_entry.extend(length.to_be_bytes().into_iter());
-
- self.active_metadata_file.write_all(&metadata_entry)?;
-
- // Flush and sync metadata to disk
- if self.config.write_durability == WriteDurability::Flush {
- self.active_metadata_file.flush()?;
- }
- if self.config.write_durability == WriteDurability::FlushSync {
- self.active_metadata_file.flush()?;
- self.active_metadata_file.sync_all()?;
- }
-
- self.remove_record_from_memtables(&record);
- }
-
- self.active_metadata_file.unlock()?;
- self.active_data_file.unlock()?;
-
- debug!("Records deleted");
-
- Ok(recs)
- }
-
/// Check if there are any pending tasks and do them. Tasks include:
/// - Rotating the active log file if it has reached capacity and compacting it.
///
@@ -819,156 +170,13 @@ impl<R: Recordable> DB<R> {
/// You may call this function in a separate thread or process to avoid blocking the main thread.
/// However, the database will be exclusively locked, so all writes and reads will be blocked during the tasks.
pub fn do_maintenance_tasks(&mut self) -> DBResult<()> {
- request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
-
- ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?;
-
- let metadata_size = self.active_metadata_file.seek(SeekFrom::End(0))?;
- if metadata_size >= self.config.segment_size as u64 {
- self.rotate_and_compact()?;
- }
-
- self.active_metadata_file.unlock()?;
-
- Ok(())
- }
-
- fn rotate_and_compact(&mut self) -> DBResult<()> {
- debug!("Active log size exceeds threshold, starting rotation and compaction...");
-
- self.active_data_file.lock_shared()?;
- let original_data_len = self.active_data_file.seek(SeekFrom::End(0))?;
-
- let active_target = fs::read_link(&self.data_dir.join(ACTIVE_SYMLINK_FILENAME))?;
- let active_num = parse_segment_number(&active_target)?;
-
- debug!("Reading segment data into a BTreeMap");
- let mut pk_to_item_map: BTreeMap<&IndexableValue, &Record> = BTreeMap::new();
- let forward_read_items: Vec<(IndexableValue, Record)> = ForwardLogReader::new(
- self.active_metadata_file.try_clone()?,
- self.active_data_file.try_clone()?,
- )
- .map(|item| {
- (
- item.record
- .at(self.primary_key_index)
- .as_indexable()
- .expect("Primary key was not indexable"),
- item.record,
- )
- })
- .collect();
-
- self.active_data_file.unlock()?;
-
- for (pk, record) in forward_read_items.iter() {
- pk_to_item_map.insert(pk, record);
- }
-
- debug!(
- "Read {} records, out of which {} were unique",
- forward_read_items.len(),
- pk_to_item_map.len()
- );
-
- // Create a new log data file and write it
- debug!("Opening new data file and writing compacted data");
- let (new_data_uuid, new_data_path) = create_segment_data_file(&self.data_dir)?;
- let mut new_data_file = APPEND_MODE.open(&new_data_path)?;
-
- let mut pk_to_data_map = BTreeMap::new();
- let mut offset = 0u64;
- for (pk, record) in pk_to_item_map.into_iter() {
- let serialized = record.serialize();
- let len = serialized.len() as u64;
- new_data_file.write_all(&serialized)?;
-
- pk_to_data_map.insert(pk, (offset, len));
- offset += len;
- }
-
- // Sync the data file to disk.
- // This is fine to do without consulting WriteDurability because this is a one-off
- // operation that is not part of the normal write path.
- new_data_file.flush()?;
- new_data_file.sync_all()?;
-
- let final_data_len = new_data_file.seek(io::SeekFrom::End(0))?;
- debug!(
- "Wrote compacted data, reduced data size: {} -> {}",
- original_data_len, final_data_len
- );
-
- // Create a new log metadata file and write it
- debug!("Opening temp metadata file and writing pointers to compacted data file");
- let temp_metadata_file = tempfile::NamedTempFile::new()?;
- let temp_metadata_path = temp_metadata_file.as_ref();
- let mut temp_metadata_file = WRITE_MODE.open(temp_metadata_path)?;
-
- let metadata_header = MetadataHeader {
- version: 1,
- uuid: new_data_uuid,
- };
-
- temp_metadata_file.write_all(&metadata_header.serialize())?;
-
- for (pk, _) in forward_read_items.iter() {
- let (offset, len) = pk_to_data_map.get(&pk).unwrap();
-
- let mut metadata_buf = vec![];
- metadata_buf.extend(offset.to_be_bytes().into_iter());
- metadata_buf.extend(len.to_be_bytes().into_iter());
-
- temp_metadata_file.write_all(&metadata_buf)?;
- }
-
- // Sync the metadata file to disk, see comment above about sync.
- temp_metadata_file.flush()?;
- temp_metadata_file.sync_all()?;
-
- debug!("Moving temporary files to their final locations");
- let new_data_path = &self.data_dir.join(new_data_uuid.to_string());
- let active_metadata_path = &self.data_dir.join(metadata_filename(active_num)); // overwrite active
-
- fs::rename(&temp_metadata_path, &active_metadata_path)?;
-
- debug!("Compaction complete, creating new segment");
-
- let new_segment_num = active_num + 1;
- let new_metadata_path = self.data_dir.join(metadata_filename(new_segment_num));
- let mut new_metadata_file = APPEND_MODE.clone().create(true).open(&new_metadata_path)?;
-
- let new_metadata_header = MetadataHeader {
- version: 1,
- uuid: new_data_uuid,
- };
-
- new_metadata_file.write_all(&new_metadata_header.serialize())?;
-
- set_active_segment(&self.data_dir, new_segment_num)?;
-
- // Old active metadata file should lose lock by RAII, or by
- // the manual unlock call in the do_maintenance_tasks method.
-
- self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?;
- self.active_data_file = APPEND_MODE.open(&new_data_path)?;
-
- // The new active log file is not locked by this client so it cannot be touched.
- debug!(
- "Active log file {} rotated and compacted, new segment: {}",
- active_num, new_segment_num
- );
-
- Ok(())
+ self.engine.do_maintenance_tasks()
}
- #[inline]
- fn get_field_type(&self, field: &R::Field) -> Option<&ValueType> {
- self.config
- .fields
- .iter()
- .find(|(f, _)| f == field)
- .map(|(_, t)| t)
+ /// Refresh the in-memory indexes from the log files.
+ /// This needs to only be called if the read consistency is set to `ReadConsistency::Eventual`.
+ pub fn refresh_indexes(&mut self) -> DBResult<()> {
+ self.engine.refresh_indexes()
}
}
@@ -1201,9 +409,12 @@ mod tests {
.expect("Failed to create DB");
// Check that the key is not indexed before write
- assert_eq!(db.primary_memtable.get(&IndexableValue::Int(0)), None);
assert_eq!(
- db.secondary_memtables[0].find_by(&IndexableValue::String("John".to_string())),
+ db.engine.primary_memtable.get(&IndexableValue::Int(0)),
+ None
+ );
+ assert_eq!(
+ db.engine.secondary_memtables[0].find_by(&IndexableValue::String("John".to_string())),
&HashSet::new()
);
@@ -1217,13 +428,13 @@ mod tests {
// Check that the key is now indexed
let expected_log_key = LogKey::new(1, 0);
assert_eq!(
- db.primary_memtable.get(&IndexableValue::Int(0)),
+ db.engine.primary_memtable.get(&IndexableValue::Int(0)),
Some(&expected_log_key)
);
let mut expected_set: HashSet<LogKey> = HashSet::new();
expected_set.insert(expected_log_key);
assert_eq!(
- db.secondary_memtables[0].find_by(&IndexableValue::String("John".to_string())),
+ db.engine.secondary_memtables[0].find_by(&IndexableValue::String("John".to_string())),
&expected_set,
);
}