From c9d2a58c94e30ea615376fbd10d8c6299470e6ea Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Sat, 5 Oct 2024 18:45:50 +0200 Subject: Handle Null values properly, refactor schema creation to support nullable fields --- benches/benchmark.rs | 36 +++++++++++----------- src/common.rs | 42 ++++++++++++++++++++++++++ src/lib.rs | 62 +++++++++++++++++++++++++------------- tests/integration.rs | 85 +++++++++++++++++++++++++++++++--------------------- 4 files changed, 153 insertions(+), 72 deletions(-) diff --git a/benches/benchmark.rs b/benches/benchmark.rs index 43d53f7..d993a1c 100644 --- a/benches/benchmark.rs +++ b/benches/benchmark.rs @@ -26,9 +26,9 @@ pub fn upsert_various_initial_sizes(c: &mut Criterion) { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -61,9 +61,9 @@ pub fn upsert_write_durability(c: &mut Criterion) { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .write_durability(mode.clone()) .primary_key(Field::Id) @@ -91,9 +91,9 @@ pub fn get_from_disk_various_initial_sizes(c: &mut Criterion) { .data_dir(&data_dir) .memtable_capacity(0) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -123,9 +123,9 @@ pub fn get_various_memtable_capacities(c: &mut Criterion) { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -141,9 +141,9 @@ pub fn get_various_memtable_capacities(c: &mut Criterion) { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .memtable_capacity(size) .primary_key(Field::Id) @@ -176,9 +176,9 @@ fn reverse_read_file_with_various_buffer_sizes(c: &mut Criterion) { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() diff --git a/src/common.rs b/src/common.rs index 2f0f942..c2af766 100644 --- a/src/common.rs +++ b/src/common.rs @@ -106,6 +106,48 @@ pub enum RecordFieldType { Bytes, } +#[derive(Debug, Clone)] +pub struct RecordField { + pub field_type: RecordFieldType, + pub nullable: bool, +} + +impl RecordField { + pub fn int() -> Self { + RecordField { + field_type: RecordFieldType::Int, + nullable: false, + } + } + + pub fn float() -> Self { + RecordField { + field_type: RecordFieldType::Float, + nullable: false, + } + } + + pub fn string() -> Self { + RecordField { + field_type: RecordFieldType::String, + nullable: false, + } + } + + pub fn bytes() -> Self { + RecordField { + field_type: RecordFieldType::Bytes, + nullable: false, + } + } + + pub fn nullable(&mut self) -> Self { + let mut new = self.clone(); + new.nullable = true; + new + } +} + #[derive(Debug, Clone)] pub enum RecordValue { Null, diff --git a/src/lib.rs b/src/lib.rs index 0ce40bb..1050cb1 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -24,7 +24,7 @@ pub struct ConfigBuilder<'a, Field: Eq + Clone + Debug> { data_dir: Option, segment_size: Option, memtable_capacity: Option, - fields: Option<&'a Vec<(Field, RecordFieldType)>>, + fields: Option<&'a Vec<(Field, RecordField)>>, primary_key: Option, secondary_keys: Option>, memtable_evict_policy: Option, @@ -67,7 +67,7 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<'a, Field> { } /// The field schema of the database. - pub fn fields(&mut self, fields: &'a Vec<(Field, RecordFieldType)>) -> &mut Self { + pub fn fields(&mut self, fields: &'a Vec<(Field, RecordField)>) -> &mut Self { self.fields = Some(fields); self } @@ -142,7 +142,7 @@ struct Config { pub data_dir: String, pub segment_size: usize, pub memtable_capacity: usize, - pub fields: Vec<(Field, RecordFieldType)>, + pub fields: Vec<(Field, RecordField)>, pub primary_key: Field, pub secondary_keys: Vec, pub memtable_evict_policy: MemtableEvictPolicy, @@ -203,15 +203,14 @@ impl DB { // If any of the keys is not in the schema or // is not an IndexableValue, return an error for &key in &all_keys { - let (_, field_type) = - config - .fields - .iter() - .find(|(field, _)| field == key) - .ok_or(io::Error::new( - io::ErrorKind::InvalidInput, - "Secondary key must be present in the field schema", - ))?; + let (_, RecordField { field_type, .. }) = config + .fields + .iter() + .find(|(field, _)| field == key) + .ok_or(io::Error::new( + io::ErrorKind::InvalidInput, + "Secondary key must be present in the field schema", + ))?; match field_type { RecordFieldType::Int | RecordFieldType::String => {} @@ -280,19 +279,42 @@ impl DB { } // Validate that record fields match schema types - // TODO: handle Null - for (i, (_, field_type)) in self.config.fields.iter().enumerate() { - match (&record.values[i], field_type) { - (RecordValue::Int(_), RecordFieldType::Int) => {} - (RecordValue::Float(_), RecordFieldType::Float) => {} - (RecordValue::String(_), RecordFieldType::String) => {} - (RecordValue::Bytes(_), RecordFieldType::Bytes) => {} + for (i, (_, field)) in self.config.fields.iter().enumerate() { + match (&record.values[i], field) { + ( + RecordValue::Null, + RecordField { + nullable: true, + field_type: _, + }, + ) => {} + ( + RecordValue::Int(_), + RecordField { + field_type: RecordFieldType::Int, + .. + }, + ) => {} + ( + RecordValue::String(_), + RecordField { + field_type: RecordFieldType::String, + .. + }, + ) => {} + ( + RecordValue::Bytes(_), + RecordField { + field_type: RecordFieldType::Bytes, + .. + }, + ) => {} _ => { return Err(io::Error::new( io::ErrorKind::InvalidInput, format!( "Record field {} has incorrect type: {:?}, expected {:?}", - &i, &record.values[i], &field_type + &i, &record.values[i], &field.field_type ), )) } diff --git a/tests/integration.rs b/tests/integration.rs index fc1d0cc..cb3af56 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -4,8 +4,8 @@ extern crate tempfile; use ctor::ctor; use env_logger; use log::debug; -use log_db; -use log_db::{Record, RecordFieldType, RecordValue, DB, TEST_RESOURCES_DIR}; +use log_db::{self, RecordField}; +use log_db::{Record, RecordValue, DB, TEST_RESOURCES_DIR}; use std::fs; use std::path::Path; use std::thread; @@ -41,9 +41,9 @@ fn test_initialize() { let _db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -56,9 +56,9 @@ fn test_upsert_and_get_with_primary_memtable() { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -89,9 +89,9 @@ fn test_upsert_and_get_without_memtable() { .data_dir(&data_dir) .memtable_capacity(0) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string().nullable()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -101,7 +101,7 @@ fn test_upsert_and_get_without_memtable() { let record0 = Record { values: vec![ RecordValue::Int(0), - RecordValue::String("John".to_string()), + RecordValue::Null, RecordValue::Bytes(vec![3, 4, 5]), ], }; @@ -143,7 +143,7 @@ fn test_upsert_and_get_without_memtable() { _ => false, }); assert!(match (&result.values[1], &record0.values[1]) { - (RecordValue::String(a), RecordValue::String(b)) => a == b, + (RecordValue::Null, RecordValue::Null) => true, _ => false, }); @@ -161,15 +161,32 @@ fn test_upsert_and_get_without_memtable() { }); } +#[test] +fn test_upsert_fails_on_null_in_non_nullable_field() { + let data_dir = tmp_dir(); + let mut db = DB::configure() + .data_dir(&data_dir) + .fields(&vec![(Field::Id, RecordField::int())]) + .primary_key(Field::Id) + .initialize() + .expect("Failed to initialize DB instance"); + + let record = Record { + // Null value + values: vec![RecordValue::Null], + }; + assert!(db.upsert(&record).is_err()); +} + #[test] fn test_upsert_fails_on_invalid_number_of_values() { let data_dir = tmp_dir(); let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -191,9 +208,9 @@ fn test_upsert_fails_on_invalid_value_type() { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -215,9 +232,9 @@ fn test_upsert_and_get_from_secondary_memtable() { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .secondary_keys(vec![Field::Name]) @@ -277,9 +294,9 @@ fn test_initialize_and_read_from_primary_memtable_fixture_db2() { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() @@ -312,9 +329,9 @@ fn test_initialize_without_memtables_fixture_db3() { let mut db = DB::configure() .data_dir(&data_dir) .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Name, RecordFieldType::String), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Name, RecordField::string()), + (Field::Data, RecordField::bytes()), ]) .memtable_capacity(0) .primary_key(Field::Id) @@ -342,7 +359,7 @@ fn test_multiple_writing_threads() { threads.push(thread::spawn(move || { let mut db = DB::configure() .data_dir(&data_dir) - .fields(&vec![(Field::Id, RecordFieldType::Int)]) + .fields(&vec![(Field::Id, RecordField::int())]) .primary_key(Field::Id) .initialize() .expect("Failed to initialize DB instance"); @@ -361,7 +378,7 @@ fn test_multiple_writing_threads() { // Read the records let mut db = DB::configure() .data_dir(&data_dir) - .fields(&vec![(Field::Id, RecordFieldType::Int)]) + .fields(&vec![(Field::Id, RecordField::int())]) .primary_key(Field::Id) .initialize() .expect("Failed to initialize DB instance"); @@ -391,7 +408,7 @@ fn test_one_writer_and_multiple_reading_threads() { threads.push(thread::spawn(move || { let mut db = DB::configure() .data_dir(&data_dir) - .fields(&vec![(Field::Id, RecordFieldType::Int)]) + .fields(&vec![(Field::Id, RecordField::int())]) .primary_key(Field::Id) .initialize() .expect("Failed to initialize DB instance"); @@ -422,7 +439,7 @@ fn test_one_writer_and_multiple_reading_threads() { threads.push(thread::spawn(move || { let mut db = DB::configure() .data_dir(&data_dir) - .fields(&vec![(Field::Id, RecordFieldType::Int)]) + .fields(&vec![(Field::Id, RecordField::int())]) .primary_key(Field::Id) .initialize() .expect("Failed to initialize DB instance"); @@ -449,8 +466,8 @@ fn test_literal_escape_is_escaped() { .data_dir(&data_dir) .memtable_capacity(0) // disable memtables .fields(&vec![ - (Field::Id, RecordFieldType::Int), - (Field::Data, RecordFieldType::Bytes), + (Field::Id, RecordField::int()), + (Field::Data, RecordField::bytes()), ]) .primary_key(Field::Id) .initialize() -- cgit v1.3