diff options
Diffstat (limited to 'log_db/src/lib.rs')
| -rw-r--r-- | log_db/src/lib.rs | 794 |
1 files changed, 794 insertions, 0 deletions
diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs new file mode 100644 index 0000000..6787cf1 --- /dev/null +++ b/log_db/src/lib.rs @@ -0,0 +1,794 @@ +#[macro_use] +extern crate log; + +mod common; +mod forward_log_reader; +mod primary_memtable; +mod reverse_log_reader; +mod secondary_memtable; + +pub use common::*; +pub use forward_log_reader::ForwardLogReader; +use fs2::lock_contended_error; +use fs2::FileExt; +use primary_memtable::PrimaryMemtable; +pub use reverse_log_reader::ReverseLogReader; +use secondary_memtable::SecondaryMemtable; +use std::fmt::Debug; +use std::fs::{self}; +use std::io::{self, Write}; +use std::os::unix::fs::MetadataExt; +use std::path::{Path, PathBuf}; +use std::thread; + +pub struct ConfigBuilder<'a, Field: Eq + Clone + Debug> { + data_dir: Option<String>, + segment_size: Option<usize>, + memtable_capacity: Option<usize>, + fields: Option<&'a Vec<(Field, RecordField)>>, + primary_key: Option<Field>, + secondary_keys: Option<Vec<Field>>, + memtable_evict_policy: Option<MemtableEvictPolicy>, + write_durability: Option<WriteDurability>, +} + +impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<'a, Field> { + pub fn new() -> ConfigBuilder<'a, Field> { + ConfigBuilder::<Field> { + data_dir: None, + segment_size: None, + memtable_capacity: None, + fields: None, + primary_key: None, + secondary_keys: None, + memtable_evict_policy: None, + write_durability: None, + } + } + + /// 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 is closed and a new one is created. + /// Closed segment files can be compacted. + pub fn segment_size(&mut self, segment_size: usize) -> &mut Self { + self.segment_size = Some(segment_size); + self + } + + /// The maximum size of a single memtable in terms of records. + /// Note that each secondary index will have its own memtable. + pub fn memtable_capacity(&mut self, memtable_capacity: usize) -> &mut Self { + self.memtable_capacity = Some(memtable_capacity); + self + } + + /// The field schema of the database. + pub fn fields(&mut self, fields: &'a Vec<(Field, RecordField)>) -> &mut Self { + self.fields = Some(fields); + self + } + + /// The primary key of the database, used to construct + /// the primary memtable index. This should be the field + /// that is most frequently queried. + pub fn primary_key(&mut self, primary_key: Field) -> &mut Self { + self.primary_key = Some(primary_key); + self + } + + /// The secondary keys of the database, used to construct + /// the secondary memtable indexes. + pub fn secondary_keys(&mut self, secondary_keys: Vec<Field>) -> &mut Self { + self.secondary_keys = Some(secondary_keys); + self + } + + /// The eviction policy for the memtables. Determines which + /// record will be dropped from a memtable when it reaches + /// capacity. + pub fn memtable_evict_policy( + &mut self, + memtable_evict_policy: MemtableEvictPolicy, + ) -> &mut Self { + self.memtable_evict_policy = Some(memtable_evict_policy); + 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 + } + + pub fn initialize(&self) -> Result<DB<Field>, io::Error> { + let config = Config::<Field> { + data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()), + segment_size: self.segment_size.unwrap_or(4 * 1024 * 1024), // 4MB + memtable_capacity: self.memtable_capacity.unwrap_or(1_000_000), + fields: self + .fields + .ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Required config value \"fields\" is not set", + ))? + .clone(), + primary_key: self.primary_key.clone().ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Required config value \"primary_key\" is not set", + ))?, + secondary_keys: self.secondary_keys.clone().unwrap_or(Vec::new()), + memtable_evict_policy: self + .memtable_evict_policy + .clone() + .unwrap_or(MemtableEvictPolicy::LeastReadOrWritten), + write_durability: self + .write_durability + .clone() + .unwrap_or(WriteDurability::Flush), + }; + + DB::initialize(&config) + } +} + +#[derive(Clone)] +struct Config<Field: Eq + Clone> { + pub data_dir: String, + pub segment_size: usize, + pub memtable_capacity: usize, + pub fields: Vec<(Field, RecordField)>, + pub primary_key: Field, + pub secondary_keys: Vec<Field>, + pub memtable_evict_policy: MemtableEvictPolicy, + pub write_durability: WriteDurability, +} + +pub struct DB<Field: Eq + Clone + Debug> { + config: Config<Field>, + log_path: PathBuf, + log_file: fs::File, + primary_key_index: usize, + primary_memtable: PrimaryMemtable, + secondary_memtables: Vec<SecondaryMemtable<Field>>, +} + +impl<Field: Eq + Clone + Debug> DB<Field> { + /// Create a new database configuration builder. + pub fn configure() -> ConfigBuilder<'static, Field> { + ConfigBuilder::new() + } + + fn initialize(config: &Config<Field>) -> Result<DB<Field>, io::Error> { + info!("Initializing DB..."); + // If data_dir does not exist, create it + if !fs::exists(&config.data_dir)? { + fs::create_dir_all(&config.data_dir)?; + } + + let log_path = Path::new(&config.data_dir).join(ACTIVE_LOG_FILENAME); + + // Create the log file if it does not exist + let log_file_file = fs::OpenOptions::new() + .create(true) + .read(true) + .append(true) + .open(&log_path)?; + + // Create the exclusive lock request file if it does not exist + fs::OpenOptions::new() + .create(true) + .write(true) + .open(&Path::new(&config.data_dir).join(EXCL_LOCK_REQUEST_FILENAME))?; + + // 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 (_, RecordField { field_type, .. }) = config + .fields + .iter() + .find(|(field, _)| field == key) + .ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Secondary key must be present in the field schema", + ))?; + + match field_type { + RecordFieldType::Int | RecordFieldType::String => {} + _ => { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "Secondary key must be an IndexableValue", + )) + } + } + } + let primary_memtable = PrimaryMemtable::new( + config.memtable_capacity, + config.memtable_evict_policy.clone(), + ); + let secondary_memtables = config + .secondary_keys + .iter() + .map(|key| { + SecondaryMemtable::new( + key, + primary_key_index, + config.memtable_capacity, + config.memtable_evict_policy.clone(), + ) + }) + .collect(); + + let mut db = DB::<Field> { + config: config.clone(), + log_path: log_path.clone(), + log_file: log_file_file, + primary_key_index, + primary_memtable, + secondary_memtables, + }; + + info!("Rebuilding memtable indexes..."); + let mut file = fs::OpenOptions::new().read(true).open(&log_path)?; + + let forward_log_reader = ForwardLogReader::new(&mut file); + for record in forward_log_reader { + db.update_primary_index(&record); + db.update_secondary_indexes(&record); + } + + info!("Database ready."); + + Ok(db) + } + + /// Insert a record into the database. If the primary key value already exists, + /// the existing record will be replaced by the supplied one. + pub fn upsert(&mut self, record: &Record) -> Result<(), io::Error> { + debug!("Upserting record: {:?}", record); + // Validate the record length + if record.values.len() != self.config.fields.len() { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!( + "Record has an incorrect number of fields: {}, expected {}", + record.values.len(), + self.config.fields.len() + ), + )); + } + + // Validate that record fields match schema types + for (i, (_, field)) in self.config.fields.iter().enumerate() { + match (&record.values[i], field) { + ( + RecordValue::Null, + RecordField { + nullable: true, + field_type: _, + }, + ) => {} + ( + RecordValue::Int(_), + RecordField { + field_type: RecordFieldType::Int, + .. + }, + ) => {} + ( + RecordValue::String(_), + RecordField { + field_type: RecordFieldType::String, + .. + }, + ) => {} + ( + RecordValue::Bytes(_), + RecordField { + field_type: RecordFieldType::Bytes, + .. + }, + ) => {} + _ => { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!( + "Record field {} has incorrect type: {:?}, expected {:?}", + &i, &record.values[i], &field.field_type + ), + )) + } + } + } + + debug!("Record is valid"); + debug!("Opening file in append mode and acquiring exclusive lock..."); + + // Acquire an exclusive lock for writing + self.request_exclusive_lock()?; + + if self.ensure_correct_file_is_open()? { + // The log file has been rotated, so we must try again + return self.upsert(record); + } + + debug!("Lock acquired, appending to log file"); + + // Write the record to the log + // Each serialized row is suffixed with the field separator character sequence + let mut serialized_record = record.serialize(); + serialized_record.extend(SEQ_RECORD_SEP); + self.log_file.write_all(&serialized_record)?; + + // Flush and sync to disk + if self.config.write_durability == WriteDurability::Flush { + self.log_file.flush()?; + } + if self.config.write_durability == WriteDurability::FlushSync { + self.log_file.flush()?; + self.log_file.sync_all()?; + } + + self.log_file.unlock()?; + + debug!("Record appended to log file, lock released"); + + debug!("Updating primary memtable"); + self.update_primary_index(record); + + debug!("Updating secondary memtables"); + self.update_secondary_indexes(record); + + Ok(()) + } + + /// Get a record by its primary index value. + /// E.g. `db.get(RecordValue::Int(10))`. + pub fn get(&mut self, query_key: &RecordValue) -> Result<Option<Record>, io::Error> { + let query_key_original = query_key; + debug!( + "Getting record with field {:?} = {:?}", + &self.config.primary_key, query_key + ); + let query_key = query_key_original.as_indexable().ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Queried value must be indexable", + ))?; + + debug!("Looking up key {:?} in primary memtable", query_key); + let found = self.primary_memtable.get(&query_key); + if let Some(record) = found { + debug!("Found record in primary memtable: {:?}", record); + return Ok(Some(record.clone())); + } + + debug!( + "No memtable entry found, looking up key {:?} in log file", + query_key + ); + + debug!( + "Matching records based on value at primary key index ({})", + &self.primary_key_index + ); + debug!("Opening file in read mode and acquiring shared lock..."); + + // Open the file and acquire a shared lock for reading + let mut file = fs::OpenOptions::new().read(true).open(&self.log_path)?; + + self.request_shared_lock(&mut file)?; + + if !is_file_same_as_path(&file, &self.log_path)? { + // The log file has been rotated, so we must try again + debug!("Lock acquired, but the log file has been rotated. Retrying get..."); + file.unlock()?; + drop(file); + return self.get(query_key_original); + } + + debug!("Lock acquired, searching log files for record"); + + let segment_numbers = self.segment_numbers()?; + let mut result: Option<Record> = None; + for &n in &segment_numbers { + if n == 0 { + debug!("Searching the active log file..."); + result = ReverseLogReader::new(&mut file)?.find(|record| { + let record_key = record.values[self.primary_key_index] + .as_indexable() + .expect("A non-indexable value was stored at key index"); + record_key == query_key + }); + + debug!("Active log file searched, releasing shared lock..."); + file.unlock()?; + } else { + debug!("Locking and searching rotated log segment file {}...", n); + let path = Path::new(&self.config.data_dir) + .join(ACTIVE_LOG_FILENAME) + .with_extension(n.to_string()); + + let mut segm_file = fs::OpenOptions::new().read(true).open(&path)?; + self.request_shared_lock(&mut segm_file)?; + result = ReverseLogReader::new(&mut segm_file)?.find(|record| { + let record_key = record.values[self.primary_key_index] + .as_indexable() + .expect("A non-indexable value was stored at key index"); + record_key == query_key + }); + + debug!("Segment file searched, releasing shared lock..."); + segm_file.unlock()?; + }; + + if result.is_some() { + break; + } + } + + debug!("Record search complete"); + + let result_value = match &result { + Some(record) => record, + None => { + debug!("No record found for key {:?}", query_key); + return Ok(None); + } + }; + + debug!("Found matching record in log file."); + + debug!("Updating primary memtable"); + self.update_primary_index(&result_value); + + debug!("Updating secondary memtables"); + self.update_secondary_indexes(&result_value); + + Ok(result) + } + + /// Get a collection of records based on a field value. + /// Indexes will be used if they contain the requested key. + pub fn find_all( + &mut self, + field: &Field, + query_key: &RecordValue, + ) -> Result<Vec<Record>, io::Error> { + // If querying by primary key, return the result of `get` wrapped in a vec. + if field == &self.config.primary_key { + return match self.get(query_key)? { + Some(record) => Ok(vec![record.clone()]), + None => Ok(vec![]), + }; + } + + // Otherwise, continue with querying secondary indexes. + let query_key_original = query_key; + debug!( + "Finding all records with field {:?} = {:?}", + field, query_key + ); + let query_key = query_key_original.as_indexable().ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Queried value must be indexable", + ))?; + + // Try to find a memtable with the queried key + let found_memtable_index = self + .secondary_memtables + .iter_mut() + .position(|mt| &mt.field == field); + + if let Some(memtable_index) = found_memtable_index { + debug!( + "Found suitable secondary index. Looking up key {:?} in the memtable", + query_key + ); + let records = self.secondary_memtables[memtable_index].find_all(&query_key); + debug!("Found matching key"); + return Ok(records.iter().map(|&record| record.clone()).collect()); + } + + debug!( + "No memtable entry found, looking up key {:?} in log file", + query_key + ); + + // Get the index of the requested field + let key_index = self + .config + .fields + .iter() + .position(|(schema_field, _)| schema_field == field) + .ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Key not found in schema after initialize", + ))?; + + debug!("Matching key index {}", key_index); + debug!("Acquiring shared lock..."); + + // Acquire a shared lock for reading + self.log_file.lock_shared()?; + + if self.ensure_correct_file_is_open()? { + // The log file has been rotated, so we must try again + return self.find_all(field, query_key_original); + } + + debug!("Lock acquired, searching log file for record"); + + let result = ReverseLogReader::new(&mut self.log_file)? + .filter(|record| { + let record_key = record.values[key_index] + .as_indexable() + .expect("A non-indexable value was stored at key index"); + record_key == query_key + }) + .collect::<Vec<Record>>(); + + self.log_file.unlock()?; + debug!("Record search complete, lock released"); + + debug!( + "Number of matching records found in log file: {}", + result.len() + ); + + if let Some(memtable_index) = found_memtable_index { + debug!("Inserting result set into secondary index"); + self.secondary_memtables[memtable_index].set_all(&query_key, &result); + } + + Ok(result) + } + + /// Ensures that the `self.log_file` handle is still pointing to the correct file. + /// If the file has been rotated, the handle will be closed and reopened. + /// Returns `true` if the file has been rotated and the handle has been reopened. + fn ensure_correct_file_is_open(&mut self) -> Result<bool, io::Error> { + if !is_file_same_as_path(&self.log_file, &self.log_path)? { + // The log file has been rotated, so we must try again + debug!( + "Lock acquired, but the log file has been rotated. Reopening file and retrying..." + ); + self.log_file.unlock()?; + + self.log_file = fs::OpenOptions::new() + .create(true) + .read(true) + .append(true) + .open(&self.log_path)?; + + Ok(true) + } else { + Ok(false) + } + } + + fn update_primary_index(&mut self, record: &Record) { + let key = record.values[self.primary_key_index] + .as_indexable() + .expect("A non-indexable value was stored at key index"); + self.primary_memtable.set(&key, record); + } + + fn update_secondary_indexes(&mut self, record: &Record) { + self.secondary_memtables + .iter_mut() + .for_each(|secondary_memtable| { + debug!( + "Updating memtable for index on {:?}", + &secondary_memtable.field + ); + for (index, (schema_field, _)) in self.config.fields.iter().enumerate() { + if schema_field == &secondary_memtable.field { + let key = record.values[index] + .as_indexable() + .expect("Secondary index key was not indexable"); + secondary_memtable.set(&key, record); + } + } + }); + } + + fn request_exclusive_lock(&mut self) -> Result<(), io::Error> { + // Create a lock on the exclusive lock request file to signal to readers that they should wait + let lock_request_path = Path::new(&self.config.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 + .open(&lock_request_path)?; + + // Attempt to acquire an exclusive lock on the lock request file + // This will block until the lock is acquired + lock_request_file.lock_exclusive()?; + + // Check that the exclusive lock request file is still the same as the one we opened + // NOTE: this isn't strictly necessary, but it's a good sanity check. Disabled for now. + // if !is_file_same_as_path(&lock_request_file, &lock_request_path)? { + // // The lock request file has been removed + // return Err(io::Error::new( + // io::ErrorKind::Other, + // "Lock request file was removed unexpectedly", + // )); + // } + + // Acquire an exclusive lock on the log file + self.log_file.lock_exclusive()?; + + // Unlock the request file + lock_request_file.unlock()?; + + Ok(()) + } + + fn is_exclusive_lock_requested(&self, data_dir: &str) -> Result<bool, io::Error> { + let lock_request_path = Path::new(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 + .open(&lock_request_path)?; + + // Attempt to acquire a shared lock on the lock request file + // If the file is already locked, return false + match lock_request_file.try_lock_shared() { + Err(e) => { + if e.kind() == lock_contended_error().kind() { + return Ok(true); + } + return Err(e); + } + Ok(_) => { + // Check that the exclusive lock request file is still the same as the one we opened + if !is_file_same_as_path(&lock_request_file, &lock_request_path)? { + // The lock request file has been removed + return Err(io::Error::new( + io::ErrorKind::Other, + "Lock request file was removed unexpectedly", + )); + } + + lock_request_file.unlock()?; + return Ok(false); + } + } + } + + fn request_shared_lock(&self, file: &mut fs::File) -> Result<(), io::Error> { + const SHARED_LOCK_WAIT_MAX_MS: u64 = 100; + let mut timeout = 5; + loop { + if self.is_exclusive_lock_requested(&self.config.data_dir)? { + debug!("Exclusive lock requested, waiting for {}ms before requesting a shared lock again", timeout); + thread::sleep(std::time::Duration::from_millis(timeout)); + timeout = std::cmp::min(timeout * 2, SHARED_LOCK_WAIT_MAX_MS); + } else { + file.lock_shared()?; + return Ok(()); + } + } + } + + /// 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. + /// + /// This function should be called periodically to ensure that the database remains in an optimal state. + /// Note that this function is synchronous and may block for a relatively long time. + /// 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 will be blocked during the tasks. + pub fn do_maintenance_tasks(&mut self) -> Result<(), io::Error> { + let active_log_path = Path::new(&self.config.data_dir).join(ACTIVE_LOG_FILENAME); + let active_log_md = fs::metadata(&active_log_path)?; + + if active_log_md.size() >= self.config.segment_size as u64 { + // Rotate the active log file + + debug!("Starting rotation, requesting exclusive lock..."); + self.request_exclusive_lock()?; + + debug!("Exclusive lock acquired, rotating active log file..."); + let next_segment_number = self.next_segment_number()?; + let next_segment_path = + &active_log_path.with_extension(next_segment_number.to_string()); + + debug!("Renaming active log file to {:?}", &next_segment_path); + fs::rename(&active_log_path, &next_segment_path)?; + + // Create a new active log file + self.log_file = fs::OpenOptions::new() + .create(true) + .write(true) + .append(true) + .open(&active_log_path)?; + + // The new active log file is not locked by this client so it cannot be touched. + // Compact the rotated segment. + + debug!("Active log file rotated"); + + // TODO compaction of rotated segments + } + + Ok(()) + } + + fn next_segment_number(&self) -> Result<u64, io::Error> { + match self.segment_numbers()?.iter().max() { + Some(greatest) => Ok(greatest + 1), + None => Ok(1), + } + } + + /// Query the filesystem to get the numbers of existing segments + /// in the intended reading order: first the active log (signaled with 0), + /// then the segments from the greatest ordinal (newest) to the least (oldest). + /// E.g. `vec![0, 4, 3, 2, 1]`. + fn segment_numbers(&self) -> Result<Vec<u64>, io::Error> { + // TODO: optimize the vecs out of here + let files = fs::read_dir(&self.config.data_dir)?; + let mut nums: Vec<u64> = files + .filter_map(|f| { + let f_path = match f { + Ok(f) => f.path(), + Err(_) => return None, + }; + + if !f_path.is_file() { + return None; + } + + let name = &f_path + .with_extension("") + .file_name() + .expect("File did not have a name?") + .to_str() + .expect("Failed to convert file name to string") + .to_string(); + + if name != ACTIVE_LOG_FILENAME { + return None; + } + + let ext = match f_path.extension() { + Some(ext) => ext, + None => return None, + }; + + let ext_num = ext + .to_str() + .expect("Extension was not a valid UTF-8 string") + .parse::<u64>() + .expect("Extension was not a valid number"); + + Some(ext_num) + }) + .collect(); + + nums.sort(); + nums.push(0); + nums.reverse(); + Ok(nums) + } +} |
