From b5abd6dfea8b26b73cc9b489695ffdd6de6fe1cd Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Mon, 6 Jan 2025 15:51:45 +0200 Subject: Implement remove --- log_db/src/lib.rs | 80 ++++++++++++++++++++++++++++++++++++++++++++----------- 1 file changed, 65 insertions(+), 15 deletions(-) (limited to 'log_db/src/lib.rs') diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index 194b6d6..3c3f7bb 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -355,19 +355,20 @@ impl DB { fn remove_record_from_memtables(&mut self, record: &Record) { let pk = record.at(self.primary_key_index).as_indexable().unwrap(); - self.primary_memtable.remove(&pk); - for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() { - let secondary_memtable = &mut self.secondary_memtables[sk_index]; - let sk_field_index = self - .config - .fields - .iter() - .position(|(f, _)| sk_field == f) - .unwrap(); - let sk = record.at(sk_field_index).as_indexable().unwrap(); - - secondary_memtable.remove(&sk); + if let Some(plk) = self.primary_memtable.remove(&pk) { + for (sk_index, sk_field) in self.config.secondary_keys.iter_mut().enumerate() { + let secondary_memtable = &mut self.secondary_memtables[sk_index]; + let sk_field_index = self + .config + .fields + .iter() + .position(|(f, _)| sk_field == f) + .unwrap(); + let sk = record.at(sk_field_index).as_indexable().unwrap(); + + secondary_memtable.remove(&sk, &plk); + } } } @@ -651,9 +652,58 @@ impl DB { } } - /// Delete records by a field value. - pub fn delete(&mut self, field: &Field, value: &Value) -> Result { - todo!() + /// Delete record by primary key. + pub fn delete(&mut self, pk: &Value) -> Result, DBError> { + let record = match self.get(pk)? { + Some(record) => record, + None => return Ok(None), + }; + + // The Record interface does not allow manually setting the tombstone flag, + // so we have to serialize the record and manually set the first byte to B_TOMBSTONE. + let mut record_serialized = vec![B_TOMBSTONE]; + record_serialized.extend(&record.serialize()[1..]); + + request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?; + self.active_data_file.lock_exclusive()?; + + let offset = self.active_data_file.seek(SeekFrom::End(0))?; + let length = record_serialized.len() as u64; + + self.active_data_file.write_all(&record_serialized)?; + + // Flush and sync data to disk + if self.config.write_durability == WriteDurability::Flush { + self.active_data_file.flush()?; + } + if self.config.write_durability == WriteDurability::FlushSync { + self.active_data_file.flush()?; + self.active_data_file.sync_all()?; + } + + let mut metadata_entry = vec![]; + metadata_entry.extend(offset.to_be_bytes().iter()); + metadata_entry.extend(length.to_be_bytes().iter()); + + self.active_metadata_file.write_all(&metadata_entry)?; + + // Flush and sync metadata to disk + if self.config.write_durability == WriteDurability::Flush { + self.active_metadata_file.flush()?; + } + if self.config.write_durability == WriteDurability::FlushSync { + self.active_metadata_file.flush()?; + self.active_metadata_file.sync_all()?; + } + + self.remove_record_from_memtables(&record); + + self.active_metadata_file.unlock()?; + self.active_data_file.unlock()?; + + debug!("Record deleted, returning from delete"); + + Ok(Some(record)) } /// Check if there are any pending tasks and do them. Tasks include: -- cgit v1.3