aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src/lib.rs
diff options
context:
space:
mode:
Diffstat (limited to 'log_db/src/lib.rs')
-rw-r--r--log_db/src/lib.rs106
1 files changed, 90 insertions, 16 deletions
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<Field: Eq + Clone + Debug> DB<Field> {
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::<Field> {
@@ -251,8 +251,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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<Field: Eq + Clone + Debug> DB<Field> {
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<Field: Eq + Clone + Debug> DB<Field> {
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<Field: Eq + Clone + Debug> DB<Field> {
}
}
- 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<Field: Eq + Clone + Debug> DB<Field> {
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);
+ }
+ }
+}