aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2025-02-21 16:15:29 +0200
committerJan Tuomi <jan@jantuomi.fi>2025-02-21 16:15:29 +0200
commitdaca0e6b68d32d885face00d68df5b5cb00e098b (patch)
tree60a3c73463efde0f890fdd681f438b504477d04a /log_db/src
parent18ae1c7ad6d45f42f39472f5b76b3f19d2b357b1 (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.rs44
-rw-r--r--log_db/src/engine.rs61
-rw-r--r--log_db/src/lib.rs112
-rw-r--r--log_db/src/record.rs38
-rw-r--r--log_db/src/row.rs15
-rw-r--r--log_db/src/schema.rs7
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>,
+}