use std::collections::HashMap;
use convert_case::{Case, Casing};
use proc_macro2::{Ident, Span, TokenStream};
use quote::quote;
use crate::common::model::Index;
use crate::common::model::Operation;
use crate::common::name_generator::{WorktableNameGenerator, is_float};
use crate::generators::persist::PersistGenerator;
impl PersistGenerator {
pub fn gen_query_delete_impl(&mut self) -> syn::Result<TokenStream> {
let name_generator = WorktableNameGenerator::from_table_name(self.name.to_string());
let table_ident = name_generator.get_work_table_ident();
let custom_deletes = if let Some(q) = &self.queries {
let custom_deletes = self.gen_custom_deletes(q.deletes.clone());
quote! {
#custom_deletes
}
} else {
quote! {}
};
let full_row_delete = self.gen_full_row_delete();
let full_row_delete_without_lock = self.gen_full_row_delete_without_lock();
Ok(quote! {
impl #table_ident {
#full_row_delete
#full_row_delete_without_lock
#custom_deletes
}
})
}
fn gen_full_row_delete(&mut self) -> TokenStream {
let name_generator = WorktableNameGenerator::from_table_name(self.name.to_string());
let pk_ident = name_generator.get_primary_key_type_ident();
let delete_logic = self.gen_delete_logic(true);
let full_row_lock = self.gen_full_lock_for_update();
quote! {
pub async fn delete<Pk>(&self, pk: Pk) -> core::result::Result<(), WorkTableError>
where #pk_ident: From<Pk>
{
let pk: #pk_ident = pk.into();
let op_lock = { #full_row_lock };
let _guard = LockGuard::new(
op_lock,
self.0.lock_manager.clone(),
pk.clone(),
);
#delete_logic
core::result::Result::Ok(())
}
}
}
fn gen_full_row_delete_without_lock(&mut self) -> TokenStream {
let name_generator = WorktableNameGenerator::from_table_name(self.name.to_string());
let pk_ident = name_generator.get_primary_key_type_ident();
let delete_logic = self.gen_delete_logic(false);
quote! {
pub async fn delete_without_lock<Pk>(&self, pk: Pk) -> core::result::Result<(), WorkTableError>
where #pk_ident: From<Pk>
{
let pk: #pk_ident = pk.into();
#delete_logic
core::result::Result::Ok(())
}
}
}
fn gen_delete_logic(&self, is_locked: bool) -> TokenStream {
let name_generator = WorktableNameGenerator::from_table_name(self.name.to_string());
let pk_ident = name_generator.get_primary_key_type_ident();
let secondary_events_ident = name_generator.get_space_secondary_index_events_ident();
let process = quote! {
let (secondary_keys_events, res) = self.0.indexes.delete_row_cdc(row, link);
res?;
let (_, primary_key_events) = self.0.primary_index.remove_cdc(pk.clone(), link);
self.0.data.delete(link).map_err(WorkTableError::PagesError)?;
let mut op: Operation<
<<#pk_ident as TablePrimaryKey>::Generator as PrimaryKeyGeneratorState>::State,
#pk_ident,
#secondary_events_ident
> = Operation::Delete(DeleteOperation {
id: uuid::Uuid::now_v7().into(),
secondary_keys_events,
primary_key_events,
link,
});
self.1.apply_operation(op);
};
if is_locked {
quote! {
let link = match self.0
.primary_index
.pk_map
.get(&pk)
.map(|v| v.get().value.into())
.ok_or(WorkTableError::NotFound) {
Ok(l) => l,
Err(e) => {
return Err(e);
}
};
let row = self.0.select(pk.clone()).unwrap();
#process
}
} else {
quote! {
let link = self.0
.primary_index
.pk_map
.get(&pk)
.map(|v| v.get().value.into())
.ok_or(WorkTableError::NotFound)?;
let row = self.0.select(pk.clone()).unwrap();
#process
}
}
}
fn gen_custom_deletes(&mut self, deleted: HashMap<Ident, Operation>) -> TokenStream {
let defs = deleted
.iter()
.map(|(name, op)| {
let snake_case_name = name.to_string().from_case(Case::Pascal).to_case(Case::Snake);
let method_ident = Ident::new(format!("delete_{snake_case_name}").as_str(), Span::mixed_site());
let index = self.columns.indexes.values().find(|idx| idx.field == op.by);
let type_ = self.columns.columns_map.get(&op.by).unwrap();
if let Some(index) = index {
let index_name = &index.name;
if index.is_unique {
Self::gen_unique_delete(type_, &method_ident, index_name)
} else {
Self::gen_non_unique_delete(type_, &method_ident, index)
}
} else {
Self::gen_brute_force_delete_field(&op.by, type_, &method_ident)
}
})
.collect::<Vec<_>>();
quote! {
#(#defs)*
}
}
fn gen_brute_force_delete_field(field: &Ident, type_: &TokenStream, name: &Ident) -> TokenStream {
quote! {
pub async fn #name(&self, by: #type_) -> core::result::Result<(), WorkTableError> {
self.iter_with_async(|row| {
if row.#field == by {
futures::future::Either::Left(async move {
self.delete::<_>(row.get_primary_key()).await
})
} else {
futures::future::Either::Right(async {
Ok(())
})
}
}).await?;
core::result::Result::Ok(())
}
}
}
fn gen_non_unique_delete(type_: &TokenStream, name: &Ident, index: &Index) -> TokenStream {
let by_field = &index.field;
let index = &index.name;
let by = if is_float(type_.to_string().as_str()) {
quote! {
&OrderedFloat(by)
}
} else {
quote! {
&by
}
};
quote! {
pub async fn #name(&self, by: #type_) -> core::result::Result<(), WorkTableError> {
let mut pks: Vec<_> = Vec::new();
for link in self.0.indexes.#index.get(#by).map(|kv| kv.1.0) {
match self.0.data.select_non_ghosted(link) {
core::result::Result::Ok(r) => {
if r.#by_field == by {
pks.push(r.get_primary_key());
}
}
core::result::Result::Err(e) if e.is_row_absent() => {}
core::result::Result::Err(e) => {
return core::result::Result::Err(WorkTableError::PagesError(e));
}
}
}
pks.sort_unstable();
pks.dedup();
for pk in pks {
match self.delete(pk).await {
core::result::Result::Ok(()) => {}
core::result::Result::Err(WorkTableError::NotFound) => {}
core::result::Result::Err(e) => return core::result::Result::Err(e),
}
}
core::result::Result::Ok(())
}
}
}
fn gen_unique_delete(type_: &TokenStream, name: &Ident, index: &Ident) -> TokenStream {
let by = if is_float(type_.to_string().as_str()) {
quote! {
&OrderedFloat(by)
}
} else {
quote! {
&by
}
};
quote! {
pub async fn #name(&self, by: #type_) -> core::result::Result<(), WorkTableError> {
let row_to_update = self.0.indexes.#index.get(#by).map(|v| v.get().value.into());
if let Some(link) = row_to_update {
let row = self.0.data.select_non_ghosted(link).map_err(WorkTableError::PagesError)?;
self.delete(row.get_primary_key()).await?;
}
core::result::Result::Ok(())
}
}
}
}