From 147c589ce617c1d500c3b698321eeb71bacd74a3 Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Fri, 21 Feb 2025 14:25:15 +0200 Subject: Rename Record -> Row --- log_db/src/engine.rs | 48 +++++++++++------------- log_db/src/lib.rs | 18 ++++----- log_db/src/log_reader_forward.rs | 10 ++--- log_db/src/record.rs | 81 ---------------------------------------- log_db/src/row.rs | 81 ++++++++++++++++++++++++++++++++++++++++ 5 files changed, 117 insertions(+), 121 deletions(-) delete mode 100644 log_db/src/record.rs create mode 100644 log_db/src/row.rs (limited to 'log_db') 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 Engine { 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 Engine { 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 Engine { 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 Engine { } } - 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 Engine { field: &str, values: impl Iterator, params: &QueryParams, - ) -> DBResult> { + ) -> DBResult> { let indexables = values .map(|value| { value.as_indexable().ok_or(DBError::ValidationError( @@ -321,7 +321,7 @@ impl Engine { fn read_tagged_log_keys<'a>( &self, log_keys: impl Iterator, - ) -> DBResult> { + ) -> DBResult> { let mut records = vec![]; let mut log_keys_map = BTreeMap::new(); @@ -366,7 +366,7 @@ impl Engine { 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 Engine { field: &str, range: B, params: &QueryParams, - ) -> DBResult> { + ) -> DBResult> { fn range_bound_to_indexable(bound: Bound<&Value>) -> DBResult> { match bound { Bound::Included(value) => value @@ -458,8 +458,8 @@ impl Engine { } } - pub fn delete_by_field(&mut self, field: &str, value: &Value) -> DBResult> { - let recs: Vec = self + pub fn delete_by_field(&mut self, field: &str, value: &Value) -> DBResult> { + let recs: Vec = self .batch_find_by_records(field, std::iter::once(value), &DEFAULT_QUERY_PARAMS)? .into_iter() .map(|(_, mut rec)| { @@ -573,18 +573,18 @@ impl Engine { 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 { engine: Engine, @@ -55,11 +55,11 @@ impl 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, 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 DB { /// Get a record by its primary index value. /// E.g. `db.get(Value::Int(10))`. pub fn get(&mut self, value: &Value) -> DBResult> { - 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 DB { ) })?; - 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/record.rs deleted file mode 100644 index 8108fb5..0000000 --- a/log_db/src/record.rs +++ /dev/null @@ -1,81 +0,0 @@ -use super::*; - -#[derive(Debug, Clone)] -pub struct Record { - pub values: Vec, - pub tombstone: bool, -} - -impl Record { - pub fn serialize(&self) -> Vec { - let mut bytes = Vec::new(); - - if self.tombstone { - bytes.extend(&[B_TOMBSTONE]); - } else { - bytes.extend(&[B_LIVE]); - } - - for value in &self.values { - bytes.extend(value.serialize()); - } - bytes - } - - pub fn deserialize(bytes: &[u8]) -> Record { - assert!(bytes.len() > 0); - - let mut values = Vec::new(); - - let tombstone = bytes[0] == B_TOMBSTONE; - - let mut start = 1; - while start < bytes.len() { - let (rv, consumed) = Value::deserialize(&bytes[start..]); - values.push(rv); - start += consumed; - } - Record { values, tombstone } - } - - pub fn from(values: &[Value]) -> Record { - Record { - values: values.to_vec(), - tombstone: false, - } - } - - pub fn at(&self, index: usize) -> &Value { - &self.values[index] - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_record_serialize_deserialize() { - let record = Record { - values: vec![ - Value::Int(1), - Value::String("hello".to_string()), - Value::Bytes(vec![0, 1, 2, 3]), - ], - tombstone: true, - }; - - let serialized = record.serialize(); - let deserialized = Record::deserialize(&serialized); - let reserialized = deserialized.serialize(); - - assert_eq!(serialized.len(), reserialized.len()); - assert_eq!(record.values, deserialized.values); - } -} - -#[derive(Clone, Debug)] -pub enum TxEntry { - Upsert { record: Record }, - Delete { record: Record }, -} diff --git a/log_db/src/row.rs b/log_db/src/row.rs new file mode 100644 index 0000000..5ebb069 --- /dev/null +++ b/log_db/src/row.rs @@ -0,0 +1,81 @@ +use super::*; + +#[derive(Debug, Clone)] +pub struct Row { + pub values: Vec, + pub tombstone: bool, +} + +impl Row { + pub fn serialize(&self) -> Vec { + let mut bytes = Vec::new(); + + if self.tombstone { + bytes.extend(&[B_TOMBSTONE]); + } else { + bytes.extend(&[B_LIVE]); + } + + for value in &self.values { + bytes.extend(value.serialize()); + } + bytes + } + + pub fn deserialize(bytes: &[u8]) -> Row { + assert!(bytes.len() > 0); + + let mut values = Vec::new(); + + let tombstone = bytes[0] == B_TOMBSTONE; + + let mut start = 1; + while start < bytes.len() { + let (rv, consumed) = Value::deserialize(&bytes[start..]); + values.push(rv); + start += consumed; + } + Row { values, tombstone } + } + + pub fn from(values: &[Value]) -> Row { + Row { + values: values.to_vec(), + tombstone: false, + } + } + + pub fn at(&self, index: usize) -> &Value { + &self.values[index] + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_record_serialize_deserialize() { + let record = Row { + values: vec![ + Value::Int(1), + Value::String("hello".to_string()), + Value::Bytes(vec![0, 1, 2, 3]), + ], + tombstone: true, + }; + + let serialized = record.serialize(); + let deserialized = Row::deserialize(&serialized); + let reserialized = deserialized.serialize(); + + assert_eq!(serialized.len(), reserialized.len()); + assert_eq!(record.values, deserialized.values); + } +} + +#[derive(Clone, Debug)] +pub enum TxEntry { + Upsert { record: Row }, + Delete { record: Row }, +} -- cgit v1.3