From ee0a81cf0277a7ff8683f94cc95a23f29996d980 Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Thu, 10 Oct 2024 23:41:10 +0300 Subject: Add invariant that memtables should always have the same set of primary key references --- log_db/src/lib.rs | 106 +++++++++++++++++++++++++++++++++++++++++++++--------- 1 file changed, 90 insertions(+), 16 deletions(-) (limited to 'log_db/src/lib.rs') diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index d704c7d..bcdaefe 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -234,7 +234,7 @@ impl DB { let secondary_memtables = config .secondary_keys .iter() - .map(|key| SecondaryMemtable::new(key)) + .map(|key| SecondaryMemtable::new(&config.fields, key, primary_key_index)) .collect(); let mut db = DB:: { @@ -251,8 +251,7 @@ impl DB { let forward_log_reader = ForwardLogReader::new(&mut file); for record in forward_log_reader { - db.update_primary_index(&record); - db.update_secondary_indexes(&record); + db.insert_to_memtables(&record); } info!("Database ready."); @@ -351,11 +350,8 @@ impl DB { debug!("Record appended to log file, lock released"); - debug!("Updating primary memtable"); - self.update_primary_index(record); - - debug!("Updating secondary memtables"); - self.update_secondary_indexes(record); + debug!("Updating memtables"); + self.insert_to_memtables(record); Ok(()) } @@ -456,11 +452,8 @@ impl DB { debug!("Found matching record in log file."); - debug!("Updating primary memtable"); - self.update_primary_index(&result_value); - - debug!("Updating secondary memtables"); - self.update_secondary_indexes(&result_value); + debug!("Updating memtables"); + self.insert_to_memtables(&result_value); Ok(result) } @@ -593,14 +586,30 @@ impl DB { } } - fn update_primary_index(&mut self, record: &Record) { + fn insert_to_memtables(&mut self, record: &Record) { let key = record.values[self.primary_key_index] .as_indexable() .expect("A non-indexable value was stored at key index"); + + if self.primary_memtable.capacity == 0 { + return; + } + + debug!( + "Inserting/updating record in primary memtable with key {:?} = {:?}", + &key, &record, + ); + + if let Some(evicted) = self.primary_memtable.evict_if_necessary() { + self.secondary_memtables + .iter_mut() + .for_each(|secondary_memtable| { + secondary_memtable.remove(&evicted); + }); + } + self.primary_memtable.set(&key, record); - } - fn update_secondary_indexes(&mut self, record: &Record) { self.secondary_memtables .iter_mut() .for_each(|secondary_memtable| { @@ -840,3 +849,68 @@ impl DB { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use rand::distributions::Alphanumeric; + use rand::Rng; + use std::collections::HashSet; + use tempfile::tempdir; + + #[derive(Eq, PartialEq, Clone, Debug)] + enum Field { + Id, + Name, + Data, + } + + fn tmp_dir() -> String { + let dir = tempdir() + .expect("Failed to create temporary directory") + .path() + .to_str() + .expect("Failed to convert temporary directory path to string") + .to_string(); + fs::create_dir_all(&dir).expect("Failed to create temporary directory"); + dir + } + + #[test] + fn memtables_always_have_the_same_primary_keys() { + let data_dir = tmp_dir(); + + let mut db = DB::configure() + .data_dir(&data_dir) + .fields(vec![ + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + ]) + .primary_key(Field::Id) + .secondary_keys(vec![Field::Name]) + .initialize() + .expect("Failed to initialize DB instance"); + + let mut rng = rand::thread_rng(); + for _ in 0..100 { + let id = rng.gen_range(0..100); + let name = (0..5).map(|_| rng.sample(Alphanumeric) as char).collect(); + + let record = Record { + values: vec![RecordValue::Int(id), RecordValue::String(name)], + }; + db.upsert(&record).expect("Failed to upsert record"); + + let p_set: HashSet<&IndexableValue> = db.primary_memtable.records.keys().collect(); + let mut s_set: HashSet<&IndexableValue> = HashSet::new(); + + for table in db.secondary_memtables.iter() { + table.records.values().for_each(|r| { + s_set.extend(r); + }); + } + + assert_eq!(p_set, s_set); + } + } +} -- cgit v1.3