diff options
Diffstat (limited to 'log_db/src')
| -rw-r--r-- | log_db/src/common.rs | 98 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 338 | ||||
| -rw-r--r-- | log_db/src/memtable_secondary.rs | 29 |
3 files changed, 248 insertions, 217 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs index e384a8f..e379364 100644 --- a/log_db/src/common.rs +++ b/log_db/src/common.rs @@ -29,13 +29,6 @@ pub fn metadata_filename(num: u16) -> String { format!("metadata.{}", num) } -#[derive(Debug, Eq, PartialEq)] -pub enum SpecialSequence { - RecordSeparator, - LiteralFieldSeparator, - LiteralEscape, -} - /// LogKey is a packed struct that contains: /// - a log segment number (16 bits) /// - a log index within the segment (48 bits) @@ -192,6 +185,16 @@ 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. + 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. @@ -214,45 +217,47 @@ pub enum IndexableValue { String(String), } +/// A primitive type #[derive(Debug, Clone)] -pub enum RecordFieldType { +pub enum PrimValueType { Int, Float, String, Bytes, } +/// A primitive type + a nullability bit #[derive(Debug, Clone)] -pub struct RecordField { - pub field_type: RecordFieldType, +pub struct ValueType { + pub prim_value_type: PrimValueType, pub nullable: bool, } -impl RecordField { +impl ValueType { pub fn int() -> Self { - RecordField { - field_type: RecordFieldType::Int, + ValueType { + prim_value_type: PrimValueType::Int, nullable: false, } } pub fn float() -> Self { - RecordField { - field_type: RecordFieldType::Float, + ValueType { + prim_value_type: PrimValueType::Float, nullable: false, } } pub fn string() -> Self { - RecordField { - field_type: RecordFieldType::String, + ValueType { + prim_value_type: PrimValueType::String, nullable: false, } } pub fn bytes() -> Self { - RecordField { - field_type: RecordFieldType::Bytes, + ValueType { + prim_value_type: PrimValueType::Bytes, nullable: false, } } @@ -401,7 +406,7 @@ impl Record { &self.0[index] } - pub fn validate<Field: Eq>(&self, schema: &Vec<(Field, RecordField)>) -> Result<(), io::Error> { + pub fn validate<Field: Eq>(&self, schema: &Vec<(Field, ValueType)>) -> Result<(), io::Error> { // Validate the record length if self.0.len() != schema.len() { return Err(io::Error::new( @@ -419,29 +424,29 @@ impl Record { match (&self.0[i], field) { ( Value::Null, - RecordField { + ValueType { nullable: true, - field_type: _, + prim_value_type: _, }, ) => {} ( Value::Int(_), - RecordField { - field_type: RecordFieldType::Int, + ValueType { + prim_value_type: PrimValueType::Int, .. }, ) => {} ( Value::String(_), - RecordField { - field_type: RecordFieldType::String, + ValueType { + prim_value_type: PrimValueType::String, .. }, ) => {} ( Value::Bytes(_), - RecordField { - field_type: RecordFieldType::Bytes, + ValueType { + prim_value_type: PrimValueType::Bytes, .. }, ) => {} @@ -450,7 +455,7 @@ impl Record { io::ErrorKind::InvalidInput, format!( "Record field {} has incorrect type: {:?}, expected {:?}", - &i, &self.0[i], &field.field_type + &i, &self.0[i], &field.prim_value_type ), )) } @@ -460,6 +465,41 @@ impl Record { } } +pub fn type_check(value: &Value, value_type: &ValueType) -> bool { + match (value, value_type) { + ( + Value::Int(_), + ValueType { + prim_value_type: PrimValueType::Int, + .. + }, + ) => true, + ( + Value::Float(_), + ValueType { + prim_value_type: PrimValueType::Float, + .. + }, + ) => true, + ( + Value::Bytes(_), + ValueType { + prim_value_type: PrimValueType::Bytes, + .. + }, + ) => true, + ( + Value::String(_), + ValueType { + prim_value_type: PrimValueType::String, + .. + }, + ) => true, + (Value::Null, ValueType { nullable: true, .. }) => true, + _ => false, + } +} + /// A trait that describes how to convert a data structure into a database `Record` and vice versa. pub trait Recordable { /// Convert the data structure implementing the `Recordable` trait into a database `Record`. diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index 3d8ddb9..3dc9150 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -25,10 +25,11 @@ use uuid::Uuid; pub struct ConfigBuilder<Field: Eq + Clone + Debug> { data_dir: Option<String>, segment_size: Option<usize>, - fields: Option<Vec<(Field, RecordField)>>, + fields: Option<Vec<(Field, ValueType)>>, primary_key: Option<Field>, secondary_keys: Option<Vec<Field>>, write_durability: Option<WriteDurability>, + read_consistency: Option<ReadConsistency>, } impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { @@ -40,6 +41,7 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { primary_key: None, secondary_keys: None, write_durability: None, + read_consistency: None, } } @@ -59,7 +61,7 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { } /// The field schema of the database. - pub fn fields(&mut self, fields: &[(Field, RecordField)]) -> &mut Self { + pub fn fields(&mut self, fields: &[(Field, ValueType)]) -> &mut Self { self.fields = Some(fields.to_vec()); self } @@ -87,6 +89,15 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { 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) -> Result<DB<Field>, io::Error> { let config = Config::<Field> { data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()), @@ -108,6 +119,10 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { .write_durability .clone() .unwrap_or(WriteDurability::Flush), + read_consistency: self + .read_consistency + .clone() + .unwrap_or(ReadConsistency::Strong), }; DB::initialize(&config) @@ -118,10 +133,11 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> { struct Config<Field: Eq + Clone> { pub data_dir: String, pub segment_size: usize, - pub fields: Vec<(Field, RecordField)>, + pub fields: Vec<(Field, ValueType)>, pub primary_key: Field, pub secondary_keys: Vec<Field>, pub write_durability: WriteDurability, + pub read_consistency: ReadConsistency, } pub struct DB<Field: Eq + Clone + Debug> { @@ -200,7 +216,12 @@ impl<Field: Eq + Clone + Debug> DB<Field> { // 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 + let ( + _, + ValueType { + prim_value_type, .. + }, + ) = config .fields .iter() .find(|(field, _)| field == key) @@ -209,8 +230,8 @@ impl<Field: Eq + Clone + Debug> DB<Field> { "Secondary key must be present in the field schema", ))?; - match field_type { - RecordFieldType::Int | RecordFieldType::String => {} + match prim_value_type { + PrimValueType::Int | PrimValueType::String => {} _ => { return Err(io::Error::new( io::ErrorKind::InvalidInput, @@ -269,7 +290,19 @@ impl<Field: Eq + Clone + Debug> DB<Field> { 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 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(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "Metadata file {} has invalid size: {}", + metadata_path.display(), + metadata_len + ), + )); + } request_shared_lock(&self.data_dir, &mut metadata_file)?; @@ -286,6 +319,19 @@ impl<Field: Eq + Clone + Debug> DB<Field> { let log_key = LogKey::new(segnum, index); self.primary_memtable.set(&pk, &log_key); + 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); + } + // 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 @@ -380,126 +426,77 @@ impl<Field: Eq + Clone + Debug> DB<Field> { /// Get a record by its primary index value. /// E.g. `db.get(Value::Int(10))`. pub fn get(&mut self, query_key: &Value) -> Result<Option<Record>, io::Error> { - let query_key_original = query_key; + let pk_type = &self.config.fields[self.primary_key_index].1; + if !type_check(&query_key, &pk_type) { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!( + "Queried value does not match primary key type: {:?}", + pk_type + ), + )); + } + debug!( "Getting record with field {:?} = {:?}", &self.config.primary_key, query_key ); - let query_key = query_key_original.as_indexable().ok_or(io::Error::new( + let query_key = query_key.as_indexable().ok_or(io::Error::new( io::ErrorKind::InvalidInput, "Queried value must be indexable", ))?; - match is_metadata_file_valid(&mut self.active_metadata_file)? { - IsMetadatafileValidResult::Ok => {} - _ => { - debug!("Active metadata file is invalid, acquiring exclusive lock to start autorepair..."); - request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?; - - ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?; - self.ensure_metadata_file_is_active()?; - // The lock should be dropped by RAII, but just in case - self.active_metadata_file.unlock()?; + if self.config.read_consistency == ReadConsistency::Strong { + self.refresh_indexes()?; + } - debug!("Active metadata file is now valid, retrying get operation..."); - return self.get(query_key_original); + debug!("Looking up key {:?} in primary memtable", query_key); + let log_key = match self.primary_memtable.get(&query_key) { + Some(log_key) => log_key, + None => { + debug!("Not found in primary memtable, returning None"); + return Ok(None); } }; - debug!("Looking up key {:?} in primary memtable", query_key); - let found = self.primary_memtable.get(&query_key); - if let Some(log_key) = found { - debug!("Found log_key in primary memtable: {:?}", log_key); - let segment_num = log_key.segment_num(); - let segment_index = log_key.index(); + debug!("Found log_key in primary memtable: {:?}", log_key); + let segment_num = log_key.segment_num(); + let segment_index = log_key.index(); - let metadata_path = &self.data_dir.join(metadata_filename(segment_num)); - let mut metadata_file = READ_MODE.open(&metadata_path)?; + 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)?; + request_shared_lock(&self.data_dir, &mut metadata_file)?; - let metadata_header = read_metadata_header(&mut metadata_file)?; + let metadata_header = read_metadata_header(&mut metadata_file)?; - metadata_file.seek_relative(segment_index as i64 * 16)?; + metadata_file.seek_relative(segment_index as i64 * 16)?; - let mut metadata_buf = [0; 2 * 8]; - metadata_file.read_exact(&mut metadata_buf)?; + let mut metadata_buf = [0; 2 * 8]; + metadata_file.read_exact(&mut metadata_buf)?; - metadata_file.unlock()?; + metadata_file.unlock()?; - 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()); - - let data_path = &self.data_dir.join(metadata_header.uuid.to_string()); - let mut data_file = READ_MODE.open(&data_path)?; + 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()); - request_shared_lock(&self.data_dir, &mut data_file)?; + let data_path = &self.data_dir.join(metadata_header.uuid.to_string()); + let mut data_file = READ_MODE.open(&data_path)?; - data_file.seek(SeekFrom::Start(data_offset))?; + request_shared_lock(&self.data_dir, &mut data_file)?; - let mut data_buf = vec![0; data_length as usize]; - data_file.read_exact(&mut data_buf)?; + data_file.seek(SeekFrom::Start(data_offset))?; - let record = Record::deserialize(&data_buf); - - return Ok(Some(record)); - } - - debug!( - "No memtable entry found, looking up key {:?} in log file", - query_key - ); + let mut data_buf = vec![0; data_length as usize]; + data_file.read_exact(&mut data_buf)?; debug!( - "Matching records based on value at primary key index ({})", - &self.primary_key_index + "Read matching record with size {} from log file, deserializing and returning.", + data_buf.len() ); - let greatest = greatest_segment_number(&self.data_dir)?; - debug!("Searching segments {} through 1", greatest); - - let mut found_record: Option<Record> = None; - for segment_num in (1..=greatest).rev() { - let segment_path = &self.data_dir.join(metadata_filename(segment_num)); - - debug!( - "Opening segment {} in read mode and acquiring shared lock...", - segment_num - ); - - let mut metadata_file = READ_MODE.open(&segment_path)?; - - 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)?; - - // We should not "request_shared_lock()" here because we do not want - // to give way to writers at this point. That would possibly lead to a deadlock. - data_file.lock_shared()?; - - let mut reader = ReverseLogReader::new(metadata_file, data_file)?; - - if let Some(found) = reader.find(|record| { - let record_key = record - .at(self.primary_key_index) - .as_indexable() - .expect("Primary key must be indexable"); - record_key == query_key - }) { - found_record = Some(found); - break; - } - } - - debug!("Record search complete"); - debug!("Found matching record in log file."); - - Ok(found_record) + let record = Record::deserialize(&data_buf); + return Ok(Some(record)); } /// Get a collection of records based on a field value. @@ -514,113 +511,83 @@ impl<Field: Eq + Clone + Debug> DB<Field> { } // 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( + let query_key = query_key.as_indexable().ok_or(io::Error::new( io::ErrorKind::InvalidInput, "Queried value must be indexable", ))?; - match is_metadata_file_valid(&mut self.active_metadata_file)? { - IsMetadatafileValidResult::Ok => {} - _ => { - debug!("Active metadata file is invalid, acquiring exclusive lock to start autorepair..."); - request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?; - - ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?; - self.ensure_metadata_file_is_active()?; - // The lock should be dropped by RAII, but just in case - self.active_metadata_file.unlock()?; - - debug!("Active metadata file is now valid, retrying find_all operation..."); - return self.find_all(field, query_key_original); - } - }; - // Try to find a memtable with the queried key - let found_memtable_index = - get_secondary_memtable_index_by_field(&self.config.secondary_keys, field); - - if let Some(memtable_index) = found_memtable_index { - debug!( - "Found suitable secondary index. Looking up key {:?} in the memtable", - query_key - ); + let memtable_index = + match get_secondary_memtable_index_by_field(&self.config.secondary_keys, field) { + Some(index) => index, + None => { + return Err(io::Error::new( + io::ErrorKind::NotFound, + "Cannot find_all by non-secondary key", + )) + } + }; - // TODO: Implement secondary memtable search + if self.config.read_consistency == ReadConsistency::Strong { + self.refresh_indexes()?; } debug!( - "No memtable entry found, looking up key {:?} in log file", + "Found suitable secondary index. Looking up key {:?} in the memtable", query_key ); + let memtable = &self.secondary_memtables[memtable_index]; + let log_keys = memtable.find_all(&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", - ))?; - - let greatest = greatest_segment_number(&self.data_dir)?; + debug!("Found log keys in secondary memtable: {:?}", log_keys); - let mut found_records = vec![]; - for segment_num in (1..=greatest).rev() { - let segment_path = &self.data_dir.join(metadata_filename(segment_num)); - - debug!( - "Opening segment {} in read mode and acquiring shared lock...", - segment_num - ); + let mut records = vec![]; + for log_key in log_keys.iter() { + // TODO optimize this so that a given segment is only opened once per find_all, and not for every log key + let segment_num = log_key.segment_num(); + let segment_index = log_key.index(); - let mut metadata_file = READ_MODE.open(&segment_path)?; + 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)?; - if metadata_header.version != 1 { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "Unsupported segment version", - )); - } + metadata_file.seek_relative(segment_index as i64 * 16)?; + + let mut metadata_buf = [0; 2 * 8]; + metadata_file.read_exact(&mut metadata_buf)?; + + metadata_file.unlock()?; + + 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()); let data_path = &self.data_dir.join(metadata_header.uuid.to_string()); - let data_file = READ_MODE.open(&data_path)?; + let mut data_file = READ_MODE.open(&data_path)?; - // We should not "request_shared_lock()" here because we do not want - // to give way to writers at this point. That would possibly lead to a deadlock. - data_file.lock_shared()?; + request_shared_lock(&self.data_dir, &mut data_file)?; - let reader = ReverseLogReader::new(metadata_file, data_file)?; + data_file.seek(SeekFrom::Start(data_offset))?; - for record in reader { - let record_key = record - .at(key_index) - .as_indexable() - .expect("Secondary key must be indexable"); - if record_key == query_key { - found_records.push(record); - } - } - } + let mut data_buf = vec![0; data_length as usize]; + data_file.read_exact(&mut data_buf)?; - debug!("Record search complete"); + debug!( + "Read matching record with size {} from log file, deserializing and adding to result set.", + data_buf.len() + ); - debug!( - "Number of matching records found in log file: {}", - found_records.len() - ); + let record = Record::deserialize(&data_buf); + records.push(record); + } - Ok(found_records) + Ok(records) } /// Ensures that the `self.metadata_file` and `self.data_file` handles are still pointing to the correct files. @@ -817,7 +784,7 @@ mod tests { let mut db = DB::configure() .data_dir(data_dir.to_str().unwrap()) .segment_size(segment_size) - .fields(&[(Field::Id, RecordField::int())]) + .fields(&[(Field::Id, ValueType::int())]) .primary_key(Field::Id) .initialize() .expect("Failed to create DB"); @@ -908,7 +875,7 @@ mod tests { let mut db = DB::configure() .data_dir(data_dir.to_str().unwrap()) - .fields(&[(Field::Id, RecordField::int())]) + .fields(&[(Field::Id, ValueType::int())]) .primary_key(Field::Id) .initialize() .expect("Failed to create DB"); @@ -926,15 +893,24 @@ mod tests { .open(&segment_metadata_path) .expect("Failed to open file"); - file.write_all(&[0, 1, 2, 3]) + file.write_all(&[1, 0, 0, 0]) // A partially written integer value ([1] + some bytes) .expect("Failed to write garbage"); file.flush().unwrap(); let len = file.seek(SeekFrom::End(0)).expect("Failed to seek"); assert_ne!(len, METADATA_FILE_HEADER_SIZE as u64 + n_recs * 16); - // Try to read from the file, triggering autorepair - db.get(&Value::Int(0)).expect("Failed to get record"); + // Try to refresh indexes, reading the file from beginning to end: should lead to error + db.refresh_indexes() + .expect_err("refresh_indexes should fail because of partial write"); + + // Trigger autorepair + db.do_maintenance_tasks() + .expect("Failed to run maintenance tasks"); + + // Try to refresh indexes, reading the file from beginning to end: should work now + db.refresh_indexes() + .expect("refresh_indexes should succeed"); // Reopen file and check that it has the correct size let mut file = READ_MODE diff --git a/log_db/src/memtable_secondary.rs b/log_db/src/memtable_secondary.rs index 371cb84..5c6cd92 100644 --- a/log_db/src/memtable_secondary.rs +++ b/log_db/src/memtable_secondary.rs @@ -1,5 +1,7 @@ +use once_cell::sync::Lazy; + use super::*; -use std::collections::BTreeMap; +use std::collections::{BTreeMap, HashSet}; pub struct SecondaryMemtable { /// Map of records indexed by key. The value is the set of primary key values of records @@ -8,6 +10,8 @@ pub struct SecondaryMemtable { records: BTreeMap<IndexableValue, LogKeySet>, } +static EMPTY_SET: Lazy<HashSet<LogKey>> = Lazy::new(|| HashSet::new()); + impl SecondaryMemtable { pub fn new() -> SecondaryMemtable { SecondaryMemtable { @@ -15,15 +19,26 @@ impl SecondaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &IndexableValue) { - unimplemented!(); + pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { + match self.records.get_mut(key) { + Some(set) => { + set.insert(value.clone()); + } + None => { + self.records + .insert(key.clone(), LogKeySet::new_with_initial(&value)); + } + }; } - pub fn set_all(&mut self, key: &IndexableValue, values: &LogKeySet) { - unimplemented!(); + pub fn replace(&mut self, key: &IndexableValue, values: &LogKeySet) { + self.records.insert(key.clone(), values.clone()); } - pub fn find_all(&self, key: &IndexableValue) -> &LogKeySet { - unimplemented!(); + pub fn find_all(&self, key: &IndexableValue) -> &HashSet<LogKey> { + match self.records.get(key) { + Some(set) => set.log_keys(), + None => &EMPTY_SET, + } } } |
