aboutsummaryrefslogtreecommitdiffstats
path: root/log_db
diff options
context:
space:
mode:
Diffstat (limited to 'log_db')
-rw-r--r--log_db/src/engine.rs48
-rw-r--r--log_db/src/lib.rs18
-rw-r--r--log_db/src/log_reader_forward.rs10
-rw-r--r--log_db/src/row.rs (renamed from log_db/src/record.rs)20
4 files changed, 46 insertions, 50 deletions
diff --git a/log_db/src/engine.rs b/log_db/src/engine.rs
index b06bfd2..b5ece69 100644
--- a/log_db/src/engine.rs
+++ b/log_db/src/engine.rs
@@ -155,15 +155,15 @@ impl<T> Engine<T> {
let data_path = self.data_dir_path.join(metadata_header.uuid.to_string());
let data_file = READ_MODE.open(data_path)?;
- for ForwardLogReaderItem { record, index } in
+ for ForwardLogReaderItem { row, index } in
ForwardLogReader::new_with_index(metadata_file, data_file, from_index)
{
let log_key = LogKey::new(segnum, index);
- if record.tombstone {
- self.remove_record_from_memtables(&record);
+ if row.tombstone {
+ self.remove_record_from_memtables(&row);
} else {
- self.insert_record_to_memtables(log_key, record);
+ self.insert_record_to_memtables(log_key, row);
}
// Update from_index in case this is the last iteration: we need to know the next
@@ -183,7 +183,7 @@ impl<T> Engine<T> {
Ok(())
}
- fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Record) {
+ fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Row) {
let pk = record.at(self.primary_key_index).as_indexable().unwrap();
for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() {
@@ -203,7 +203,7 @@ impl<T> Engine<T> {
self.primary_memtable.set(pk, log_key);
}
- fn remove_record_from_memtables(&mut self, record: &Record) {
+ fn remove_record_from_memtables(&mut self, record: &Row) {
let pk = record.at(self.primary_key_index).as_indexable().unwrap();
if let Some(_) = self.primary_memtable.remove(&pk) {
@@ -222,7 +222,7 @@ impl<T> Engine<T> {
}
}
- pub fn upsert_record(&mut self, record: Record) -> DBResult<()> {
+ pub fn upsert_record(&mut self, record: Row) -> DBResult<()> {
debug!("Opening file in append mode...");
if !self.ensure_metadata_file_is_active()?
@@ -250,7 +250,7 @@ impl<T> Engine<T> {
field: &str,
values: impl Iterator<Item = &'a Value>,
params: &QueryParams,
- ) -> DBResult<Vec<(usize, Record)>> {
+ ) -> DBResult<Vec<(usize, Row)>> {
let indexables = values
.map(|value| {
value.as_indexable().ok_or(DBError::ValidationError(
@@ -321,7 +321,7 @@ impl<T> Engine<T> {
fn read_tagged_log_keys<'a>(
&self,
log_keys: impl Iterator<Item = &'a (usize, &'a LogKey)>,
- ) -> DBResult<Vec<(usize, Record)>> {
+ ) -> DBResult<Vec<(usize, Row)>> {
let mut records = vec![];
let mut log_keys_map = BTreeMap::new();
@@ -366,7 +366,7 @@ impl<T> Engine<T> {
let mut data_buf = vec![0; data_length as usize];
data_file.read_exact(&mut data_buf)?;
- let record = Record::deserialize(&data_buf);
+ let record = Row::deserialize(&data_buf);
records.push((*tag, record));
current_metadata_offset = new_metadata_offset + row_length;
@@ -381,7 +381,7 @@ impl<T> Engine<T> {
field: &str,
range: B,
params: &QueryParams,
- ) -> DBResult<Vec<Record>> {
+ ) -> DBResult<Vec<Row>> {
fn range_bound_to_indexable(bound: Bound<&Value>) -> DBResult<Bound<IndexableValue>> {
match bound {
Bound::Included(value) => value
@@ -458,8 +458,8 @@ impl<T> Engine<T> {
}
}
- pub fn delete_by_field(&mut self, field: &str, value: &Value) -> DBResult<Vec<Record>> {
- let recs: Vec<Record> = self
+ pub fn delete_by_field(&mut self, field: &str, value: &Value) -> DBResult<Vec<Row>> {
+ let recs: Vec<Row> = self
.batch_find_by_records(field, std::iter::once(value), &DEFAULT_QUERY_PARAMS)?
.into_iter()
.map(|(_, mut rec)| {
@@ -573,18 +573,18 @@ impl<T> Engine<T> {
let active_num = parse_segment_number(&active_target)?;
debug!("Reading segment data into a BTreeMap");
- let mut pk_to_item_map: BTreeMap<&IndexableValue, &Record> = BTreeMap::new();
- let forward_read_items: Vec<(IndexableValue, Record)> = ForwardLogReader::new(
+ let mut pk_to_item_map: BTreeMap<&IndexableValue, &Row> = BTreeMap::new();
+ let forward_read_items: Vec<(IndexableValue, Row)> = ForwardLogReader::new(
self.active_metadata_file.try_clone()?,
self.active_data_file.try_clone()?,
)
.map(|item| {
(
- item.record
+ item.row
.at(self.primary_key_index)
.as_indexable()
.expect("Primary key was not indexable"),
- item.record,
+ item.row,
)
})
.collect();
@@ -806,10 +806,8 @@ mod tests {
.len(),
0
);
- engine.insert_record_to_memtables(
- LogKey::new(1, 0),
- Record::from(&inst.clone().into_record()),
- );
+ engine
+ .insert_record_to_memtables(LogKey::new(1, 0), Row::from(&inst.clone().into_record()));
assert_eq!(engine.primary_memtable.get(&id), Some(&LogKey::new(1, 0)));
assert_eq!(
engine.secondary_memtables[0]
@@ -818,10 +816,8 @@ mod tests {
1
);
- engine.insert_record_to_memtables(
- LogKey::new(1, 1),
- Record::from(&inst.clone().into_record()),
- );
+ engine
+ .insert_record_to_memtables(LogKey::new(1, 1), Row::from(&inst.clone().into_record()));
assert_eq!(engine.primary_memtable.get(&id), Some(&LogKey::new(1, 1)));
assert_eq!(
engine.secondary_memtables[0]
@@ -830,7 +826,7 @@ mod tests {
1
);
- engine.remove_record_from_memtables(&Record::from(&inst.into_record()));
+ engine.remove_record_from_memtables(&Row::from(&inst.into_record()));
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 fddb344..b7229fe 100644
--- a/log_db/src/lib.rs
+++ b/log_db/src/lib.rs
@@ -23,7 +23,7 @@ mod lock;
mod log_reader_forward;
mod memtable_primary;
mod memtable_secondary;
-mod record;
+mod row;
pub use common::{DBError, DBResult, OwnedBounds, QueryParams, Value, DEFAULT_QUERY_PARAMS};
pub use config::{ReadConsistency, Schema, WriteDurability};
@@ -35,7 +35,7 @@ use lock::*;
use log_reader_forward::*;
use memtable_primary::PrimaryMemtable;
use memtable_secondary::SecondaryMemtable;
-use record::*;
+use row::*;
pub struct DB<T> {
engine: Engine<T>,
@@ -55,11 +55,11 @@ impl<T> DB<T> {
/// 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 record = Record::from(&(self.engine.config.into_record)(recordable));
- debug!("Upserting record: {:?}", record);
+ let row = Row::from(&(self.engine.config.into_record)(recordable));
+ debug!("Upserting record: {:?}", row);
self.engine
- .with_exclusive_lock(move |engine| engine.upsert_record(record))?;
+ .with_exclusive_lock(move |engine| engine.upsert_record(row))?;
Ok(())
}
@@ -67,7 +67,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>> {
- let recs = self.engine.with_shared_lock(|engine| {
+ 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
&engine.config.primary_key.clone(),
@@ -76,12 +76,12 @@ impl<T> DB<T> {
)
})?;
- assert!(recs.len() <= 1);
+ assert!(tagged_rows.len() <= 1);
- Ok(recs
+ Ok(tagged_rows
.into_iter()
.next()
- .map(|(_, rec)| (self.engine.config.from_record)(rec.values)))
+ .map(|(_, row)| (self.engine.config.from_record)(row.values)))
}
/// Get a collection of records based on an indexed field value.
diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs
index 66b0b47..f3fc16c 100644
--- a/log_db/src/log_reader_forward.rs
+++ b/log_db/src/log_reader_forward.rs
@@ -6,7 +6,7 @@ pub struct ForwardLogReader {
}
pub struct ForwardLogReaderItem {
- pub record: Record,
+ pub row: Row,
pub index: u64,
}
@@ -74,8 +74,8 @@ impl ForwardLogReader {
let mut result_buf = vec![0; entry_length as usize];
self.data_reader.read_exact(&mut result_buf)?;
- let record = Record::deserialize(&result_buf);
- return Ok(Some(ForwardLogReaderItem { record, index }));
+ let row = Row::deserialize(&result_buf);
+ return Ok(Some(ForwardLogReaderItem { row, index }));
}
}
}
@@ -121,10 +121,10 @@ mod tests {
// There are two records in the log with "schema" with one field: Bytes
- let first_record = forward_log_reader
+ let ForwardLogReaderItem { row, index: _ } = forward_log_reader
.next()
.expect("Failed to read the first record");
- assert!(match &first_record.record.values[..] {
+ assert!(match &row.values[..] {
[Value::Bytes(bytes)] => bytes.len() == 256,
_ => false,
});
diff --git a/log_db/src/record.rs b/log_db/src/row.rs
index 8108fb5..5ebb069 100644
--- a/log_db/src/record.rs
+++ b/log_db/src/row.rs
@@ -1,12 +1,12 @@
use super::*;
#[derive(Debug, Clone)]
-pub struct Record {
+pub struct Row {
pub values: Vec<Value>,
pub tombstone: bool,
}
-impl Record {
+impl Row {
pub fn serialize(&self) -> Vec<u8> {
let mut bytes = Vec::new();
@@ -22,7 +22,7 @@ impl Record {
bytes
}
- pub fn deserialize(bytes: &[u8]) -> Record {
+ pub fn deserialize(bytes: &[u8]) -> Row {
assert!(bytes.len() > 0);
let mut values = Vec::new();
@@ -35,11 +35,11 @@ impl Record {
values.push(rv);
start += consumed;
}
- Record { values, tombstone }
+ Row { values, tombstone }
}
- pub fn from(values: &[Value]) -> Record {
- Record {
+ pub fn from(values: &[Value]) -> Row {
+ Row {
values: values.to_vec(),
tombstone: false,
}
@@ -56,7 +56,7 @@ mod tests {
#[test]
fn test_record_serialize_deserialize() {
- let record = Record {
+ let record = Row {
values: vec![
Value::Int(1),
Value::String("hello".to_string()),
@@ -66,7 +66,7 @@ mod tests {
};
let serialized = record.serialize();
- let deserialized = Record::deserialize(&serialized);
+ let deserialized = Row::deserialize(&serialized);
let reserialized = deserialized.serialize();
assert_eq!(serialized.len(), reserialized.len());
@@ -76,6 +76,6 @@ mod tests {
#[derive(Clone, Debug)]
pub enum TxEntry {
- Upsert { record: Record },
- Delete { record: Record },
+ Upsert { record: Row },
+ Delete { record: Row },
}