diff options
Diffstat (limited to 'log_db/src')
| -rw-r--r-- | log_db/src/engine.rs | 48 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 18 | ||||
| -rw-r--r-- | log_db/src/log_reader_forward.rs | 10 | ||||
| -rw-r--r-- | log_db/src/row.rs (renamed from log_db/src/record.rs) | 20 |
4 files changed, 46 insertions, 50 deletions
diff --git a/log_db/src/engine.rs b/log_db/src/engine.rs index b06bfd2..b5ece69 100644 --- a/log_db/src/engine.rs +++ b/log_db/src/engine.rs @@ -155,15 +155,15 @@ impl<T> Engine<T> { let data_path = self.data_dir_path.join(metadata_header.uuid.to_string()); let data_file = READ_MODE.open(data_path)?; - for ForwardLogReaderItem { record, index } in + for ForwardLogReaderItem { row, 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); + if row.tombstone { + self.remove_record_from_memtables(&row); } else { - self.insert_record_to_memtables(log_key, record); + self.insert_record_to_memtables(log_key, row); } // Update from_index in case this is the last iteration: we need to know the next @@ -183,7 +183,7 @@ impl<T> Engine<T> { Ok(()) } - fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Record) { + fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Row) { let pk = record.at(self.primary_key_index).as_indexable().unwrap(); for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() { @@ -203,7 +203,7 @@ impl<T> Engine<T> { self.primary_memtable.set(pk, log_key); } - fn remove_record_from_memtables(&mut self, record: &Record) { + fn remove_record_from_memtables(&mut self, record: &Row) { let pk = record.at(self.primary_key_index).as_indexable().unwrap(); if let Some(_) = self.primary_memtable.remove(&pk) { @@ -222,7 +222,7 @@ impl<T> Engine<T> { } } - pub fn upsert_record(&mut self, record: Record) -> DBResult<()> { + pub fn upsert_record(&mut self, record: Row) -> DBResult<()> { debug!("Opening file in append mode..."); if !self.ensure_metadata_file_is_active()? @@ -250,7 +250,7 @@ impl<T> Engine<T> { field: &str, values: impl Iterator<Item = &'a Value>, params: &QueryParams, - ) -> DBResult<Vec<(usize, Record)>> { + ) -> DBResult<Vec<(usize, Row)>> { let indexables = values .map(|value| { value.as_indexable().ok_or(DBError::ValidationError( @@ -321,7 +321,7 @@ impl<T> Engine<T> { fn read_tagged_log_keys<'a>( &self, log_keys: impl Iterator<Item = &'a (usize, &'a LogKey)>, - ) -> DBResult<Vec<(usize, Record)>> { + ) -> DBResult<Vec<(usize, Row)>> { let mut records = vec![]; let mut log_keys_map = BTreeMap::new(); @@ -366,7 +366,7 @@ impl<T> Engine<T> { let mut data_buf = vec![0; data_length as usize]; data_file.read_exact(&mut data_buf)?; - let record = Record::deserialize(&data_buf); + let record = Row::deserialize(&data_buf); records.push((*tag, record)); current_metadata_offset = new_metadata_offset + row_length; @@ -381,7 +381,7 @@ impl<T> Engine<T> { field: &str, range: B, params: &QueryParams, - ) -> DBResult<Vec<Record>> { + ) -> DBResult<Vec<Row>> { fn range_bound_to_indexable(bound: Bound<&Value>) -> DBResult<Bound<IndexableValue>> { match bound { Bound::Included(value) => value @@ -458,8 +458,8 @@ impl<T> Engine<T> { } } - pub fn delete_by_field(&mut self, field: &str, value: &Value) -> DBResult<Vec<Record>> { - let recs: Vec<Record> = self + pub fn delete_by_field(&mut self, field: &str, value: &Value) -> DBResult<Vec<Row>> { + let recs: Vec<Row> = self .batch_find_by_records(field, std::iter::once(value), &DEFAULT_QUERY_PARAMS)? .into_iter() .map(|(_, mut rec)| { @@ -573,18 +573,18 @@ impl<T> Engine<T> { 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( + let mut pk_to_item_map: BTreeMap<&IndexableValue, &Row> = BTreeMap::new(); + let forward_read_items: Vec<(IndexableValue, Row)> = ForwardLogReader::new( self.active_metadata_file.try_clone()?, self.active_data_file.try_clone()?, ) .map(|item| { ( - item.record + item.row .at(self.primary_key_index) .as_indexable() .expect("Primary key was not indexable"), - item.record, + item.row, ) }) .collect(); @@ -806,10 +806,8 @@ mod tests { .len(), 0 ); - engine.insert_record_to_memtables( - LogKey::new(1, 0), - Record::from(&inst.clone().into_record()), - ); + engine + .insert_record_to_memtables(LogKey::new(1, 0), Row::from(&inst.clone().into_record())); assert_eq!(engine.primary_memtable.get(&id), Some(&LogKey::new(1, 0))); assert_eq!( engine.secondary_memtables[0] @@ -818,10 +816,8 @@ mod tests { 1 ); - engine.insert_record_to_memtables( - LogKey::new(1, 1), - Record::from(&inst.clone().into_record()), - ); + engine + .insert_record_to_memtables(LogKey::new(1, 1), Row::from(&inst.clone().into_record())); assert_eq!(engine.primary_memtable.get(&id), Some(&LogKey::new(1, 1))); assert_eq!( engine.secondary_memtables[0] @@ -830,7 +826,7 @@ mod tests { 1 ); - engine.remove_record_from_memtables(&Record::from(&inst.into_record())); + engine.remove_record_from_memtables(&Row::from(&inst.into_record())); assert_eq!(engine.primary_memtable.get(&id), None); assert_eq!( engine.secondary_memtables[0] diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index fddb344..b7229fe 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -23,7 +23,7 @@ mod lock; mod log_reader_forward; mod memtable_primary; mod memtable_secondary; -mod record; +mod row; pub use common::{DBError, DBResult, OwnedBounds, QueryParams, Value, DEFAULT_QUERY_PARAMS}; pub use config::{ReadConsistency, Schema, WriteDurability}; @@ -35,7 +35,7 @@ use lock::*; use log_reader_forward::*; use memtable_primary::PrimaryMemtable; use memtable_secondary::SecondaryMemtable; -use record::*; +use row::*; pub struct DB<T> { engine: Engine<T>, @@ -55,11 +55,11 @@ impl<T> DB<T> { /// 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, recordable: T) -> DBResult<()> { - let record = Record::from(&(self.engine.config.into_record)(recordable)); - debug!("Upserting record: {:?}", record); + let row = Row::from(&(self.engine.config.into_record)(recordable)); + debug!("Upserting record: {:?}", row); self.engine - .with_exclusive_lock(move |engine| engine.upsert_record(record))?; + .with_exclusive_lock(move |engine| engine.upsert_record(row))?; Ok(()) } @@ -67,7 +67,7 @@ impl<T> DB<T> { /// Get a record by its primary index value. /// E.g. `db.get(Value::Int(10))`. pub fn get(&mut self, value: &Value) -> DBResult<Option<T>> { - let recs = self.engine.with_shared_lock(|engine| { + let tagged_rows = self.engine.with_shared_lock(|engine| { engine.batch_find_by_records( // TODO: This clone is only here to appease the borrow checker &engine.config.primary_key.clone(), @@ -76,12 +76,12 @@ impl<T> DB<T> { ) })?; - assert!(recs.len() <= 1); + assert!(tagged_rows.len() <= 1); - Ok(recs + Ok(tagged_rows .into_iter() .next() - .map(|(_, rec)| (self.engine.config.from_record)(rec.values))) + .map(|(_, row)| (self.engine.config.from_record)(row.values))) } /// Get a collection of records based on an indexed field value. diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs index 66b0b47..f3fc16c 100644 --- a/log_db/src/log_reader_forward.rs +++ b/log_db/src/log_reader_forward.rs @@ -6,7 +6,7 @@ pub struct ForwardLogReader { } pub struct ForwardLogReaderItem { - pub record: Record, + pub row: Row, pub index: u64, } @@ -74,8 +74,8 @@ impl ForwardLogReader { let mut result_buf = vec![0; entry_length as usize]; self.data_reader.read_exact(&mut result_buf)?; - let record = Record::deserialize(&result_buf); - return Ok(Some(ForwardLogReaderItem { record, index })); + let row = Row::deserialize(&result_buf); + return Ok(Some(ForwardLogReaderItem { row, index })); } } } @@ -121,10 +121,10 @@ mod tests { // There are two records in the log with "schema" with one field: Bytes - let first_record = forward_log_reader + let ForwardLogReaderItem { row, index: _ } = forward_log_reader .next() .expect("Failed to read the first record"); - assert!(match &first_record.record.values[..] { + assert!(match &row.values[..] { [Value::Bytes(bytes)] => bytes.len() == 256, _ => false, }); diff --git a/log_db/src/record.rs b/log_db/src/row.rs index 8108fb5..5ebb069 100644 --- a/log_db/src/record.rs +++ b/log_db/src/row.rs @@ -1,12 +1,12 @@ use super::*; #[derive(Debug, Clone)] -pub struct Record { +pub struct Row { pub values: Vec<Value>, pub tombstone: bool, } -impl Record { +impl Row { pub fn serialize(&self) -> Vec<u8> { let mut bytes = Vec::new(); @@ -22,7 +22,7 @@ impl Record { bytes } - pub fn deserialize(bytes: &[u8]) -> Record { + pub fn deserialize(bytes: &[u8]) -> Row { assert!(bytes.len() > 0); let mut values = Vec::new(); @@ -35,11 +35,11 @@ impl Record { values.push(rv); start += consumed; } - Record { values, tombstone } + Row { values, tombstone } } - pub fn from(values: &[Value]) -> Record { - Record { + pub fn from(values: &[Value]) -> Row { + Row { values: values.to_vec(), tombstone: false, } @@ -56,7 +56,7 @@ mod tests { #[test] fn test_record_serialize_deserialize() { - let record = Record { + let record = Row { values: vec![ Value::Int(1), Value::String("hello".to_string()), @@ -66,7 +66,7 @@ mod tests { }; let serialized = record.serialize(); - let deserialized = Record::deserialize(&serialized); + let deserialized = Row::deserialize(&serialized); let reserialized = deserialized.serialize(); assert_eq!(serialized.len(), reserialized.len()); @@ -76,6 +76,6 @@ mod tests { #[derive(Clone, Debug)] pub enum TxEntry { - Upsert { record: Record }, - Delete { record: Record }, + Upsert { record: Row }, + Delete { record: Row }, } |
