diff options
| author | Jan Tuomi <jan@jantuomi.fi> | 2025-02-21 16:15:29 +0200 |
|---|---|---|
| committer | Jan Tuomi <jan@jantuomi.fi> | 2025-02-21 16:15:29 +0200 |
| commit | daca0e6b68d32d885face00d68df5b5cb00e098b (patch) | |
| tree | 60a3c73463efde0f890fdd681f438b504477d04a /log_db/src | |
| parent | 18ae1c7ad6d45f42f39472f5b76b3f19d2b357b1 (diff) | |
Remove <T> polymorphism from DB, replace with From<T> based approach
Diffstat (limited to 'log_db/src')
| -rw-r--r-- | log_db/src/config.rs | 44 | ||||
| -rw-r--r-- | log_db/src/engine.rs | 61 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 112 | ||||
| -rw-r--r-- | log_db/src/record.rs | 38 | ||||
| -rw-r--r-- | log_db/src/row.rs | 15 | ||||
| -rw-r--r-- | log_db/src/schema.rs | 7 |
6 files changed, 134 insertions, 143 deletions
diff --git a/log_db/src/config.rs b/log_db/src/config.rs index ec74164..e8c6fa1 100644 --- a/log_db/src/config.rs +++ b/log_db/src/config.rs @@ -1,12 +1,6 @@ use super::*; -pub struct Schema<F> { - pub fields: Vec<F>, - pub primary_key: F, - pub secondary_keys: Vec<F>, -} - -pub struct ConfigBuilder<T> { +pub struct ConfigBuilder { data_dir: Option<String>, segment_size: Option<usize>, write_durability: Option<WriteDurability>, @@ -15,14 +9,10 @@ pub struct ConfigBuilder<T> { fields: Option<Vec<String>>, primary_key: Option<String>, secondary_keys: Option<Vec<String>>, - from_record: Option<fn(Vec<Value>) -> T>, - into_record: Option<fn(T) -> Vec<Value>>, - - _marker: PhantomData<T>, } -impl<T> ConfigBuilder<T> { - pub fn new() -> ConfigBuilder<T> { +impl ConfigBuilder { + pub fn new() -> ConfigBuilder { ConfigBuilder { data_dir: None, segment_size: None, @@ -32,10 +22,6 @@ impl<T> ConfigBuilder<T> { fields: None, primary_key: None, secondary_keys: None, - from_record: None, - into_record: None, - - _marker: PhantomData, } } @@ -86,36 +72,18 @@ impl<T> ConfigBuilder<T> { self } - pub fn from_record(mut self, from_record: fn(Vec<Value>) -> T) -> Self { - self.from_record = Some(from_record); - self - } - - pub fn into_record(mut self, into_record: fn(T) -> Vec<Value>) -> Self { - self.into_record = Some(into_record); - self - } - - pub fn initialize(self) -> DBResult<DB<T>> { + pub fn initialize(self) -> DBResult<DB> { let schema = self .fields .ok_or_else(|| DBError::ValidationError("Schema not set".to_string()))?; let primary_key = self .primary_key .ok_or_else(|| DBError::ValidationError("Primary key not set".to_string()))?; - let from_record = self - .from_record - .ok_or_else(|| DBError::ValidationError("Callback from_record not set".to_string()))?; - let into_record = self - .into_record - .ok_or_else(|| DBError::ValidationError("Callback into_record not set".to_string()))?; let config = Config { schema, primary_key, secondary_keys: self.secondary_keys.unwrap_or_default(), - from_record, - into_record, data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()), segment_size: self.segment_size.unwrap_or(4 * 1024 * 1024), // 4MB @@ -134,12 +102,10 @@ impl<T> ConfigBuilder<T> { } #[derive(Clone)] -pub struct Config<T> { +pub struct Config { pub schema: Vec<String>, pub primary_key: String, pub secondary_keys: Vec<String>, - pub from_record: fn(Vec<Value>) -> T, - pub into_record: fn(T) -> Vec<Value>, pub data_dir: String, pub segment_size: usize, pub write_durability: WriteDurability, diff --git a/log_db/src/engine.rs b/log_db/src/engine.rs index b5ece69..0701191 100644 --- a/log_db/src/engine.rs +++ b/log_db/src/engine.rs @@ -1,7 +1,7 @@ use super::*; -pub struct Engine<T> { - pub config: Config<T>, +pub struct Engine { + pub config: Config, pub lock_manager: LockManager, data_dir_path: PathBuf, @@ -19,8 +19,8 @@ pub struct Engine<T> { pub secondary_memtables: Vec<SecondaryMemtable>, } -impl<T> Engine<T> { - pub fn initialize(config: Config<T>) -> DBResult<Engine<T>> { +impl Engine { + pub fn initialize(config: Config) -> DBResult<Engine> { info!("Initializing DB..."); // If data_dir does not exist or is empty, create it and any necessary files. // After creation, the directory should always be in a complete state without missing files. @@ -104,7 +104,7 @@ impl<T> Engine<T> { Path::new(&config.data_dir).join(active_metadata_header.uuid.to_string()); let active_data_file = APPEND_MODE.open(&active_data_path)?; - let mut engine = Engine::<T> { + let mut engine = Engine { config, lock_manager, data_dir_path, @@ -161,9 +161,9 @@ impl<T> Engine<T> { let log_key = LogKey::new(segnum, index); if row.tombstone { - self.remove_record_from_memtables(&row); + self.remove_row_from_memtables(&row.values); } else { - self.insert_record_to_memtables(log_key, row); + self.insert_row_to_memtables(log_key, row.values); } // Update from_index in case this is the last iteration: we need to know the next @@ -183,8 +183,8 @@ impl<T> Engine<T> { Ok(()) } - fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Row) { - let pk = record.at(self.primary_key_index).as_indexable().unwrap(); + fn insert_row_to_memtables(&mut self, log_key: LogKey, row_values: Vec<Value>) { + let pk = row_values[self.primary_key_index].as_indexable().unwrap(); for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() { let secondary_memtable = &mut self.secondary_memtables[sk_index]; @@ -194,7 +194,7 @@ impl<T> Engine<T> { .iter() .position(|f| sk_field == f) .unwrap(); - let sk = record.at(sk_field_index).as_indexable().unwrap(); + let sk = row_values[sk_field_index].as_indexable().unwrap(); secondary_memtable.set(pk.clone(), sk, log_key.clone()); } @@ -203,8 +203,8 @@ impl<T> Engine<T> { self.primary_memtable.set(pk, log_key); } - fn remove_record_from_memtables(&mut self, record: &Row) { - let pk = record.at(self.primary_key_index).as_indexable().unwrap(); + fn remove_row_from_memtables(&mut self, row_values: &Vec<Value>) { + let pk = row_values[self.primary_key_index].as_indexable().unwrap(); if let Some(_) = self.primary_memtable.remove(&pk) { for (sk_index, sk_field) in self.config.secondary_keys.iter_mut().enumerate() { @@ -215,7 +215,7 @@ impl<T> Engine<T> { .iter() .position(|f| sk_field == f) .unwrap(); - let sk = record.at(sk_field_index).as_indexable().unwrap(); + let sk = row_values[sk_field_index].as_indexable().unwrap(); secondary_memtable.remove(&pk, &sk); } @@ -235,7 +235,7 @@ impl<T> Engine<T> { return self.upsert_record(record); } - self.tx_log.push(TxEntry::Upsert { record }); + self.tx_log.push(TxEntry::Upsert { row: record }); if !self.tx_active { self.commit_transaction()?; @@ -471,7 +471,7 @@ impl<T> Engine<T> { // TODO: refactor the clone out of here for record in &recs { self.tx_log.push(TxEntry::Delete { - record: record.clone(), + row: record.clone(), }); } @@ -501,8 +501,8 @@ impl<T> Engine<T> { debug!("Serializing tx_log to byte arrays"); for tx_entry in &self.tx_log { let record = match tx_entry { - TxEntry::Upsert { record } => record, - TxEntry::Delete { record } => record, + TxEntry::Upsert { row: record } => record, + TxEntry::Delete { row: record } => record, }; let serialized = record.serialize(); @@ -544,8 +544,8 @@ impl<T> Engine<T> { debug!("Updating memtables"); for (log_key, tx_entry) in pending_memtable_ops { match tx_entry { - TxEntry::Upsert { record } => self.insert_record_to_memtables(log_key, record), - TxEntry::Delete { record } => self.remove_record_from_memtables(&record), + TxEntry::Upsert { row } => self.insert_row_to_memtables(log_key, row.values), + TxEntry::Delete { row } => self.remove_row_from_memtables(&row.values), } } debug!("Commit done"); @@ -580,8 +580,7 @@ impl<T> Engine<T> { ) .map(|item| { ( - item.row - .at(self.primary_key_index) + item.row.values[self.primary_key_index] .as_indexable() .expect("Primary key was not indexable"), item.row, @@ -752,12 +751,14 @@ mod tests { name: String, } - impl TestInst2 { - fn into_record(self) -> Vec<Value> { - vec![Value::Int(self.id), Value::String(self.name)] + impl From<TestInst2> for Vec<Value> { + fn from(inst: TestInst2) -> Self { + vec![Value::Int(inst.id), Value::String(inst.name)] } + } - fn from_record(record: Vec<Value>) -> Self { + impl From<Vec<Value>> for TestInst2 { + fn from(record: Vec<Value>) -> Self { let mut it = record.into_iter(); TestInst2 { id: match it.next().unwrap() { @@ -785,8 +786,6 @@ mod tests { .fields(vec![Field::Id, Field::Name]) .primary_key(Field::Id) .secondary_keys(vec![Field::Name]) - .from_record(TestInst2::from_record) - .into_record(TestInst2::into_record) .segment_size(segment_size) .initialize() .expect("Failed to create DB"); @@ -806,8 +805,7 @@ mod tests { .len(), 0 ); - engine - .insert_record_to_memtables(LogKey::new(1, 0), Row::from(&inst.clone().into_record())); + engine.insert_row_to_memtables(LogKey::new(1, 0), inst.clone().into()); assert_eq!(engine.primary_memtable.get(&id), Some(&LogKey::new(1, 0))); assert_eq!( engine.secondary_memtables[0] @@ -816,8 +814,7 @@ mod tests { 1 ); - engine - .insert_record_to_memtables(LogKey::new(1, 1), Row::from(&inst.clone().into_record())); + engine.insert_row_to_memtables(LogKey::new(1, 1), inst.clone().into()); assert_eq!(engine.primary_memtable.get(&id), Some(&LogKey::new(1, 1))); assert_eq!( engine.secondary_memtables[0] @@ -826,7 +823,7 @@ mod tests { 1 ); - engine.remove_record_from_memtables(&Row::from(&inst.into_record())); + engine.remove_row_from_memtables(&inst.into()); 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 b7229fe..006056c 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -8,7 +8,6 @@ use std::fmt::Debug; use std::fmt::Display; use std::fs::{self, metadata, File}; use std::io::{self, Read, Seek, SeekFrom, Write}; -use std::marker::PhantomData; use std::ops::*; use std::path::{Path, PathBuf}; use std::thread; @@ -23,10 +22,14 @@ mod lock; mod log_reader_forward; mod memtable_primary; mod memtable_secondary; +mod record; mod row; +mod schema; pub use common::{DBError, DBResult, OwnedBounds, QueryParams, Value, DEFAULT_QUERY_PARAMS}; -pub use config::{ReadConsistency, Schema, WriteDurability}; +pub use config::{ReadConsistency, WriteDurability}; +pub use record::Record; +pub use schema::Schema; use common::*; use config::*; @@ -37,25 +40,28 @@ use memtable_primary::PrimaryMemtable; use memtable_secondary::SecondaryMemtable; use row::*; -pub struct DB<T> { - engine: Engine<T>, +pub struct DB { + engine: Engine, } -impl<T> DB<T> { +impl DB { /// Create a new database configuration builder. - pub fn configure() -> ConfigBuilder<T> { + pub fn configure() -> ConfigBuilder { ConfigBuilder::new() } - fn initialize(config: Config<T>) -> DBResult<DB<T>> { + fn initialize(config: Config) -> DBResult<DB> { let engine = Engine::initialize(config)?; Ok(DB { engine }) } /// 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 row = Row::from(&(self.engine.config.into_record)(recordable)); + pub fn upsert(&mut self, record: impl Into<Record>) -> DBResult<()> { + let row = Row { + values: record.into().into(), + tombstone: false, + }; debug!("Upserting record: {:?}", row); self.engine @@ -66,7 +72,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>> { + pub fn get(&mut self, value: &Value) -> DBResult<Option<Record>> { 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 @@ -81,12 +87,12 @@ impl<T> DB<T> { Ok(tagged_rows .into_iter() .next() - .map(|(_, row)| (self.engine.config.from_record)(row.values))) + .map(|(_, row)| Record::from(row))) } /// Get a collection of records based on an indexed field value. - pub fn find_by(&mut self, field: impl AsRef<str>, value: &Value) -> DBResult<Vec<T>> { - let recs = self.engine.with_shared_lock(|engine| { + pub fn find_by(&mut self, field: impl AsRef<str>, value: &Value) -> DBResult<Vec<Record>> { + let tagged_rows = self.engine.with_shared_lock(|engine| { engine.batch_find_by_records( field.as_ref(), std::iter::once(value), @@ -94,9 +100,9 @@ impl<T> DB<T> { ) })?; - Ok(recs + Ok(tagged_rows .into_iter() - .map(|(_, rec)| (self.engine.config.from_record)(rec.values)) + .map(|(_, row)| Record::from(row)) .collect()) } @@ -106,15 +112,12 @@ impl<T> DB<T> { field: impl AsRef<str>, value: &Value, params: &QueryParams, - ) -> DBResult<Vec<T>> { + ) -> DBResult<Vec<Record>> { let recs = self.engine.with_shared_lock(|engine| { engine.batch_find_by_records(field.as_ref(), std::iter::once(value), params) })?; - Ok(recs - .into_iter() - .map(|(_, rec)| (self.engine.config.from_record)(rec.values)) - .collect()) + Ok(recs.into_iter().map(|(_, row)| Record::from(row)).collect()) } /// Get a collection of records based on a sequence of indexed field values. @@ -124,14 +127,14 @@ impl<T> DB<T> { &mut self, field: impl Into<String>, values: &[Value], - ) -> DBResult<Vec<(usize, T)>> { + ) -> DBResult<Vec<(usize, Record)>> { let recs = self.engine.with_shared_lock(|engine| { engine.batch_find_by_records(&field.into(), values.iter(), &DEFAULT_QUERY_PARAMS) })?; Ok(recs .into_iter() - .map(|(tag, rec)| (tag, (self.engine.config.from_record)(rec.values))) + .map(|(tag, row)| (tag, Record::from(row))) .collect()) } @@ -143,14 +146,14 @@ impl<T> DB<T> { field: impl AsRef<str>, values: &[Value], params: &QueryParams, - ) -> DBResult<Vec<(usize, T)>> { + ) -> DBResult<Vec<(usize, Record)>> { let recs = self.engine.with_shared_lock(|engine| { engine.batch_find_by_records(field.as_ref(), values.iter(), params) })?; Ok(recs .into_iter() - .map(|(tag, rec)| (tag, (self.engine.config.from_record)(rec.values))) + .map(|(tag, row)| (tag, Record::from(row))) .collect()) } @@ -161,15 +164,12 @@ impl<T> DB<T> { &mut self, field: impl AsRef<str>, range: B, - ) -> DBResult<Vec<T>> { + ) -> DBResult<Vec<Record>> { let recs = self.engine.with_shared_lock(|engine| { engine.range_by_records(field.as_ref(), range, &DEFAULT_QUERY_PARAMS) })?; - Ok(recs - .into_iter() - .map(|rec| (self.engine.config.from_record)(rec.values)) - .collect()) + Ok(recs.into_iter().map(|row| Record::from(row)).collect()) } /// Get a collection of records based on a range of indexed field values, with additional parameters. @@ -180,15 +180,12 @@ impl<T> DB<T> { field: impl AsRef<str>, range: B, params: &QueryParams, - ) -> DBResult<Vec<T>> { + ) -> DBResult<Vec<Record>> { let recs = self .engine .with_shared_lock(|engine| engine.range_by_records(field.as_ref(), range, params))?; - Ok(recs - .into_iter() - .map(|rec| (self.engine.config.from_record)(rec.values)) - .collect()) + Ok(recs.into_iter().map(|row| Record::from(row)).collect()) } /// Delete records by a field value. @@ -197,19 +194,19 @@ impl<T> DB<T> { /// /// Deletion is done by marking the record as a tombstone. The record will still be present in the log file, /// but will be ignored by reads. Upon compaction, tombstoned records will be removed. - pub fn delete_by(&mut self, field: impl AsRef<str>, value: &Value) -> DBResult<Vec<T>> { + pub fn delete_by(&mut self, field: impl AsRef<str>, value: &Value) -> DBResult<Vec<Record>> { let recs = self .engine .with_exclusive_lock(|engine| engine.delete_by_field(field.as_ref(), value))?; Ok(recs .into_iter() - .map(|rec| (self.engine.config.from_record)(rec.values)) + .map(|row| Record::from(row.values)) .collect()) } /// Delete record by primary key. - pub fn delete(&mut self, pk: &Value) -> DBResult<Option<T>> { + pub fn delete(&mut self, pk: &Value) -> DBResult<Option<Record>> { let recs = self.engine.with_exclusive_lock(|engine| { engine // TODO: This clone is only here to appease the borrow checker @@ -218,10 +215,7 @@ impl<T> DB<T> { assert!(recs.len() <= 1); - Ok(recs - .into_iter() - .next() - .map(|rec| (self.engine.config.from_record)(rec.values))) + Ok(recs.into_iter().next().map(|row| Record::from(row.values))) } /// Check if there are any pending tasks and do them. Tasks include: @@ -320,12 +314,14 @@ mod tests { id: i64, } - impl TestInst1 { - fn into_record(self) -> Vec<Value> { - vec![Value::Int(self.id)] + impl From<TestInst1> for Record { + fn from(inst: TestInst1) -> Self { + vec![Value::Int(inst.id)].into() } + } - fn from_record(record: Vec<Value>) -> Self { + impl From<Record> for TestInst1 { + fn from(record: Record) -> Self { let mut it = record.into_iter(); TestInst1 { id: match it.next().unwrap() { @@ -341,12 +337,14 @@ mod tests { name: String, } - impl TestInst2 { - fn into_record(self) -> Vec<Value> { - vec![Value::Int(self.id), Value::String(self.name)] + impl From<TestInst2> for Record { + fn from(inst: TestInst2) -> Self { + vec![Value::Int(inst.id), Value::String(inst.name)].into() } + } - fn from_record(record: Vec<Value>) -> Self { + impl From<Record> for TestInst2 { + fn from(record: Record) -> Self { let mut it = record.into_iter(); TestInst2 { id: match it.next().unwrap() { @@ -373,8 +371,6 @@ mod tests { .data_dir(data_dir.to_str().unwrap()) .fields(vec![Field::Id]) .primary_key(Field::Id) - .from_record(TestInst1::from_record) - .into_record(TestInst1::into_record) .segment_size(segment_size) .initialize() .expect("Failed to create DB"); @@ -436,17 +432,19 @@ mod tests { ); // Check that the records can be read - let inst0 = db + let inst0: TestInst1 = db .get(&Value::Int(0 as i64)) .expect("Failed to get record") - .expect("Record not found"); + .expect("Record not found") + .into(); assert!(inst0.id == 0); - let inst1 = db + let inst1: TestInst1 = db .get(&Value::Int(1 as i64)) .expect("Failed to get record") - .expect("Record not found"); + .expect("Record not found") + .into(); assert!(inst1.id == 1); } @@ -460,8 +458,6 @@ mod tests { .data_dir(data_dir.to_str().unwrap()) .fields(vec![Field::Id]) .primary_key(Field::Id) - .from_record(TestInst1::from_record) - .into_record(TestInst1::into_record) .initialize() .expect("Failed to create DB"); @@ -515,8 +511,6 @@ mod tests { .fields(vec![Field::Id, Field::Name]) .primary_key(Field::Id) .secondary_keys(vec![Field::Name]) - .from_record(TestInst2::from_record) - .into_record(TestInst2::into_record) .initialize() .expect("Failed to create DB"); diff --git a/log_db/src/record.rs b/log_db/src/record.rs new file mode 100644 index 0000000..b349853 --- /dev/null +++ b/log_db/src/record.rs @@ -0,0 +1,38 @@ +use super::*; + +pub struct Record { + values: Vec<Value>, +} + +impl Record { + pub fn values(&self) -> &[Value] { + &self.values + } +} + +impl IntoIterator for Record { + type Item = Value; + type IntoIter = std::vec::IntoIter<Self::Item>; + + fn into_iter(self) -> Self::IntoIter { + self.values.into_iter() + } +} + +impl From<Vec<Value>> for Record { + fn from(values: Vec<Value>) -> Self { + Record { values } + } +} + +impl From<Record> for Vec<Value> { + fn from(record: Record) -> Self { + record.values + } +} + +impl From<Row> for Record { + fn from(row: Row) -> Self { + Record { values: row.values } + } +} diff --git a/log_db/src/row.rs b/log_db/src/row.rs index 5ebb069..d8140d2 100644 --- a/log_db/src/row.rs +++ b/log_db/src/row.rs @@ -37,17 +37,6 @@ impl Row { } 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)] @@ -76,6 +65,6 @@ mod tests { #[derive(Clone, Debug)] pub enum TxEntry { - Upsert { record: Row }, - Delete { record: Row }, + Upsert { row: Row }, + Delete { row: Row }, } diff --git a/log_db/src/schema.rs b/log_db/src/schema.rs new file mode 100644 index 0000000..98d0df7 --- /dev/null +++ b/log_db/src/schema.rs @@ -0,0 +1,7 @@ +use super::*; + +pub struct Schema { + pub fields: Vec<String>, + pub primary_key: String, + pub secondary_keys: Vec<String>, +} |
