diff options
Diffstat (limited to 'log_db/src/lib.rs')
| -rw-r--r-- | log_db/src/lib.rs | 887 |
1 files changed, 49 insertions, 838 deletions
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, ); } |
