use crate::stores::database::DatabaseStore;
use bytes::{BufMut, Bytes, BytesMut};
use flexbuffers::Reader;
use ordinary_config::DatabaseModelConfig;
use ordinary_types::{Kind, TimeUnit};
use rustc_hash::FxHashSet;
use saferlmdb::{WriteTransaction, put};
use thiserror::Error;
use tracing::instrument;
use uuid::Uuid;
#[derive(Error, Debug)]
pub enum InsertError {
#[error("no model for idx {0}")]
ModelNotFound(u8),
#[error("value has already been indexed")]
IndexCollision,
#[error("invalid kind for index: {0}")]
InvalidIndexKind(Kind),
#[error("invalid kind for query: {0}")]
InvalidQueryKind(Kind),
#[error("item size exceeds limits")]
ExceedsSizeLimits,
#[error("(FlexBuffers) {0}")]
FlexBuffersReaderError(#[from] flexbuffers::ReaderError),
#[error("(LMDB) {0}")]
LmdbError(#[from] saferlmdb::Error),
#[error(transparent)]
Anyhow(#[from] anyhow::Error),
#[error(transparent)]
TryFromIntError(#[from] std::num::TryFromIntError),
}
#[inline]
#[allow(clippy::too_many_lines)]
#[instrument(skip_all, fields(i = model_config.idx, nm = tracing::field::display(&model_config.name), uuid = tracing::field::display(uuid)), err)]
pub(crate) fn insert(
database: &DatabaseStore,
model_config: &DatabaseModelConfig,
uuid: Uuid,
payload: &[u8],
txn: Option<WriteTransaction<'_>>,
) -> anyhow::Result<[u8; 16], InsertError> {
let uuid_bytes = *uuid.as_bytes();
let mut size: isize = 0;
let mut inserts = 0;
let root_vec = Reader::get_root(payload)?.as_vector();
let mut builder = flexbuffers::Builder::new(&flexbuffers::BuilderOptions::SHARE_NONE);
let mut item = builder.start_vector();
let mut model_prefixed_uuid = BytesMut::new();
model_prefixed_uuid.put_u8(model_config.idx);
model_prefixed_uuid.put_u8(0);
model_prefixed_uuid.put(uuid.as_ref());
let curr = if let Some(auto_int) = database.auto_ints.get(model_config.idx as usize) {
let mut auto_int = auto_int.lock();
*auto_int += 1;
*auto_int
} else {
return Err(InsertError::ModelNotFound(model_config.idx));
};
let mut item_key = BytesMut::new();
item_key.put_u8(model_config.idx);
match curr {
0..=255 => {
item_key.put_u8(u8::try_from(curr)?);
}
256..=65_535 => {
item_key.put(u16::try_from(curr)?.to_be_bytes().as_slice());
}
65_536..=4_294_967_295 => {
item_key.put(u32::try_from(curr)?.to_be_bytes().as_slice());
}
_ => {
item_key.put((curr as u64).to_be_bytes().as_slice());
}
}
item.push(flexbuffers::Blob(&uuid_bytes[..]));
let txn = match txn {
Some(txn) => txn,
None => WriteTransaction::new(database.env.clone())?,
};
{
let mut access = txn.access();
access.put(
&database.index_db,
model_prefixed_uuid.as_ref(),
item_key.as_ref(),
&put::Flags::empty(),
)?;
inserts += 1;
size += isize::try_from(model_prefixed_uuid.len())? + isize::try_from(item_key.len())?;
for field in &model_config.fields {
let reader = root_vec.idx(field.idx as usize - 1);
if let Kind::Ref {
model: _,
field: _,
many,
} = &field.kind
{
if let Some((ref_model_idx, ref_field_idx, ref_model_fields, ref_is_many)) =
database.cross_refs.get(&(model_config.idx, field.idx))
{
if many == &Some(true) && *ref_is_many {
let mut refs = item.start_vector();
let mut seen_uuids = FxHashSet::default();
for uuid in &reader.as_vector() {
let uuid = uuid.as_blob();
let uuid_bytes = uuid.0;
let mut lookup = BytesMut::new();
lookup.put_u8(*ref_model_idx);
lookup.put_u8(0);
lookup.put(uuid_bytes);
if seen_uuids.contains(uuid_bytes) {
continue;
}
seen_uuids.insert(Bytes::copy_from_slice(uuid_bytes));
let ref_item_id =
access.get::<[u8], [u8]>(&database.index_db, lookup.as_ref())?;
let ref_item_id = Bytes::copy_from_slice(ref_item_id);
refs.push(flexbuffers::Blob(ref_item_id.as_ref()));
let ref_item = access
.get::<[u8], [u8]>(&database.item_db, ref_item_id.as_ref())?;
let root_vec = Reader::get_root(ref_item)?.as_vector();
let mut builder =
flexbuffers::Builder::new(&flexbuffers::BuilderOptions::SHARE_NONE);
let mut builder_vec = builder.start_vector();
builder_vec.push(uuid);
for ref_field in ref_model_fields {
if &ref_field.idx == ref_field_idx {
let mut curr_refs = builder_vec.start_vector();
curr_refs.push(flexbuffers::Blob(item_key.as_ref()));
for curr_ref in
&root_vec.idx(ref_field.idx as usize).as_vector()
{
let curr_ref = curr_ref.as_blob();
let curr_ref_bytes = curr_ref.0;
if curr_ref_bytes != item_key {
curr_refs.push(curr_ref);
}
}
curr_refs.end_vector();
} else {
ref_field.kind.copy_to(
&root_vec.idx(ref_field.idx as usize),
&mut builder_vec,
None,
)?;
}
}
builder_vec.end_vector();
let updated_ref_item = builder.view();
let ref_item_len = isize::try_from(ref_item.len())?;
access.put(
&database.item_db,
ref_item_id.as_ref(),
updated_ref_item,
&put::Flags::empty(),
)?;
inserts += 1;
size += isize::try_from(updated_ref_item.len())? - ref_item_len;
}
refs.end_vector();
} else if *ref_is_many {
let uuid = reader.as_blob();
let uuid_bytes = uuid.0;
let mut lookup = BytesMut::new();
lookup.put_u8(*ref_model_idx);
lookup.put_u8(0);
lookup.put(uuid_bytes);
let ref_item_id =
access.get::<[u8], [u8]>(&database.index_db, lookup.as_ref())?;
let ref_item_id = Bytes::copy_from_slice(ref_item_id);
item.push(flexbuffers::Blob(ref_item_id.as_ref()));
let ref_item =
access.get::<[u8], [u8]>(&database.item_db, ref_item_id.as_ref())?;
let root_vec = Reader::get_root(ref_item)?.as_vector();
let mut builder =
flexbuffers::Builder::new(&flexbuffers::BuilderOptions::SHARE_NONE);
let mut builder_vec = builder.start_vector();
builder_vec.push(uuid);
for ref_field in ref_model_fields {
if &ref_field.idx == ref_field_idx {
let mut curr_refs = builder_vec.start_vector();
curr_refs.push(flexbuffers::Blob(item_key.as_ref()));
for curr_ref in &root_vec.idx(ref_field.idx as usize).as_vector() {
let curr_ref = curr_ref.as_blob();
let curr_ref_bytes = curr_ref.0;
if curr_ref_bytes != item_key {
curr_refs.push(curr_ref);
}
}
curr_refs.end_vector();
} else {
ref_field.kind.copy_to(
&root_vec.idx(ref_field.idx as usize),
&mut builder_vec,
None,
)?;
}
}
builder_vec.end_vector();
let updated_ref_item = builder.view();
let ref_item_len = isize::try_from(ref_item.len())?;
access.put(
&database.item_db,
ref_item_id.as_ref(),
updated_ref_item,
&put::Flags::empty(),
)?;
inserts += 1;
size += isize::try_from(updated_ref_item.len())? - ref_item_len;
} else {
let empty = item.start_vector();
empty.end_vector();
}
}
} else {
if field.indexed == Some(true) {
let mut key = BytesMut::new();
key.put_u8(model_config.idx);
key.put_u8(field.idx);
match field.kind {
Kind::String | Kind::Url => {
key.put(reader.as_str().as_bytes());
}
Kind::Uuid => {
key.put(reader.as_blob().0);
}
_ => {
return Err(InsertError::InvalidIndexKind(field.kind.clone()));
}
}
if access
.get::<[u8], [u8]>(&database.index_db, key.as_ref())
.is_ok()
{
return Err(InsertError::IndexCollision);
}
access.put(
&database.index_db,
key.as_ref(),
item_key.as_ref(),
&put::Flags::empty(),
)?;
inserts += 1;
size += isize::try_from(key.len())? + isize::try_from(item_key.len())?;
}
if field.queryable == Some(true) {
let mut key = BytesMut::new();
key.put_u8(model_config.idx);
key.put_u8(field.idx);
match &field.kind {
Kind::Bool => {
if reader.as_bool() {
key.put_u8(1);
} else {
key.put_u8(0);
}
}
Kind::String | Kind::Url => key.put(reader.as_str().as_bytes()),
Kind::Uuid => key.put(reader.as_blob().0),
Kind::Timestamp { unit, .. } => match &unit {
TimeUnit::Seconds => key.put_i64(reader.as_i64()),
},
Kind::U8 => key.put_u8(reader.as_u8()),
Kind::U16 => key.put_u16(reader.as_u16()),
Kind::U32 => key.put_u32(reader.as_u32()),
Kind::U64 => key.put_u64(reader.as_u64()),
Kind::I8 => key.put_i8(reader.as_i8()),
Kind::I16 => key.put_i16(reader.as_i16()),
Kind::I32 => key.put_i32(reader.as_i32()),
Kind::I64 => key.put_i64(reader.as_i64()),
Kind::F32 => key.put_f32(reader.as_f32()),
Kind::F64 => key.put_f64(reader.as_f64()),
_ => {
return Err(InsertError::InvalidQueryKind(field.kind.clone()));
}
}
let mut builder =
flexbuffers::Builder::new(&flexbuffers::BuilderOptions::SHARE_NONE);
let mut builder_vec = builder.start_vector();
builder_vec.push(flexbuffers::Blob(item_key.as_ref()));
let list_len = if let Ok(list) =
access.get::<[u8], [u8]>(&database.query_db, key.as_ref())
{
let root = Reader::get_root(list)?;
for val in &root.as_vector() {
builder_vec.push(val.as_blob());
}
list.len()
} else {
0
};
builder_vec.end_vector();
let updated_list = builder.view();
access.put(
&database.query_db,
key.as_ref(),
updated_list,
&put::Flags::empty(),
)?;
inserts += 1;
size += isize::try_from(updated_list.len())? - isize::try_from(list_len)?;
}
database.encrypt_compress_copy_to(field, &reader, &mut item)?;
}
}
item.end_vector();
let item = builder.view();
if item.len() as u64 > database.limits.max_item_size {
return Err(InsertError::ExceedsSizeLimits);
}
access.put(
&database.item_db,
item_key.as_ref(),
item,
&put::Flags::empty(),
)?;
inserts += 1;
size += isize::try_from(item_key.len() + item.len())?;
}
txn.commit()?;
let tracing_size = database.log_sizes.then_some(if size > 0 {
tracing::field::display(format!(
"{}",
bytesize::ByteSize(size as u64).display().si_short()
))
} else {
tracing::field::display(format!(
"-{}",
bytesize::ByteSize(-size as u64).display().si_short()
))
});
tracing::info!(inserts, size = tracing_size);
Ok(uuid_bytes)
}