aboutsummaryrefslogtreecommitdiffstats
path: root/src/lib.rs
diff options
context:
space:
mode:
Diffstat (limited to 'src/lib.rs')
-rw-r--r--src/lib.rs794
1 files changed, 0 insertions, 794 deletions
diff --git a/src/lib.rs b/src/lib.rs
deleted file mode 100644
index 6787cf1..0000000
--- a/src/lib.rs
+++ /dev/null
@@ -1,794 +0,0 @@
-#[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)
- }
-}