mod worktable_impls;
use proc_macro2::{Literal, TokenStream};
use quote::quote;
use crate::common::name_generator::WorktableNameGenerator;
use crate::persist_table::generator::Generator;
pub const WT_INDEX_EXTENSION: &str = ".wt.idx";
pub const WT_DATA_EXTENSION: &str = ".wt.data";
impl Generator {
pub fn gen_space_file_def(&self) -> TokenStream {
let type_ = self.gen_space_file_type();
let impls = self.gen_space_file_impls();
let worktable_impl = self.gen_space_file_worktable_impl();
quote! {
#type_
#impls
#worktable_impl
}
}
fn gen_space_file_type(&self) -> TokenStream {
let name_generator = WorktableNameGenerator::from_struct_ident(&self.struct_def.ident);
let index_persisted_ident = name_generator.get_persisted_index_ident();
let inner_const_name = name_generator.get_page_inner_size_const_ident();
let pk_type = name_generator.get_primary_key_type_ident();
let space_file_ident = name_generator.get_space_file_ident();
let primary_index = if self.attributes.pk_unsized {
quote! {
pub primary_index: (Vec<GeneralPage<TableOfContentsPage<(#pk_type, Link)>>>, Vec<GeneralPage<UnsizedIndexPage<#pk_type, {#inner_const_name as u32}>>>),
}
} else {
quote! {
pub primary_index: (Vec<GeneralPage<TableOfContentsPage<(#pk_type, Link)>>>, Vec<GeneralPage<IndexPage<#pk_type>>>),
}
};
quote! {
#[derive(Debug)]
pub struct #space_file_ident {
#primary_index
pub indexes: #index_persisted_ident,
pub data: Vec<GeneralPage<DataPage<#inner_const_name>>>,
pub data_info: GeneralPage<SpaceInfoPage<<<#pk_type as TablePrimaryKey>::Generator as PrimaryKeyGeneratorState>::State>>,
}
}
}
fn gen_space_file_get_primary_index_info_fn(&self) -> TokenStream {
let name_generator = WorktableNameGenerator::from_struct_ident(&self.struct_def.ident);
let literal_name = name_generator.get_work_table_literal_name();
let version_const = name_generator.get_version_const_ident();
quote! {
fn get_primary_index_info(&self) -> eyre::Result<GeneralPage<SpaceInfoPage<()>>> {
let mut info = {
let inner = SpaceInfoPage {
id: 0.into(),
version: #version_const,
page_count: 0,
name: #literal_name.to_string(),
pk_gen_state: (),
empty_links_list: vec![],
primary_key_fields: vec![],
row_schema: vec![],
secondary_index_types: vec![],
};
let header = GeneralHeader {
data_version: DATA_VERSION,
page_id: 0.into(),
previous_id: 0.into(),
next_id: 0.into(),
page_type: PageType::SpaceInfo,
space_id: 0.into(),
data_length: 0,
};
GeneralPage {
header,
inner
}
};
info.inner.page_count = self.primary_index.0.len() as u32 + self.primary_index.1.len() as u32;
Ok(info)
}
}
}
pub fn gen_space_file_impls(&self) -> TokenStream {
let name_generator = WorktableNameGenerator::from_struct_ident(&self.struct_def.ident);
let space_ident = name_generator.get_space_file_ident();
let into_worktable_fn = self.gen_space_file_into_worktable_fn();
let parse_file_fn = self.gen_space_file_parse_file_fn();
let get_primary_index_info_fn = self.gen_space_file_get_primary_index_info_fn();
quote! {
impl #space_ident {
#into_worktable_fn
#parse_file_fn
#get_primary_index_info_fn
}
}
}
fn gen_space_file_into_worktable_fn(&self) -> TokenStream {
let wt_ident = &self.struct_def.ident;
let name_generator = WorktableNameGenerator::from_struct_ident(&self.struct_def.ident);
let index_ident = name_generator.get_index_type_ident();
let task_ident = name_generator.get_persistence_task_ident();
let const_name = name_generator.get_page_inner_size_const_ident();
let pk_type = name_generator.get_primary_key_type_ident();
let lock_type = name_generator.get_lock_type_ident();
let table_name = name_generator.get_work_table_literal_name();
let secondary_index_events = name_generator.get_space_secondary_index_events_ident();
let avt_index_ident = name_generator.get_available_indexes_ident();
let primary_index_init = if self.attributes.pk_unsized {
let pk_ident = &self.pk_ident;
quote! {
let pk_map = IndexMap::<#pk_ident, OffsetEqLink<#const_name>, UnsizedNode<_>>::with_maximum_node_size(#const_name);
for page in self.primary_index.1 {
let node = page
.inner
.get_node()
.into_iter()
.map(|p| IndexPair {
key: p.key,
value: p.value.into(),
})
.collect();
pk_map.attach_node(UnsizedNode::from_inner(node, #const_name));
}
let mut reverse_pk_map = IndexMap::<OffsetEqLink<#const_name>, #pk_ident>::new();
for entry in pk_map.iter() {
let (pk, link) = entry;
reverse_pk_map.insert(*link, pk.clone());
}
let primary_index = PrimaryIndex { pk_map, reverse_pk_map };
}
} else {
quote! {
let size = get_index_page_size_from_data_length::<#pk_type>(#const_name);
let pk_map = IndexMap::<_, OffsetEqLink<#const_name>>::with_maximum_node_size(size);
for page in self.primary_index.1 {
let node = page
.inner
.get_node()
.into_iter()
.map(|p| IndexPair {
key: p.key,
value: p.value.into(),
})
.collect();
pk_map.attach_node(node);
}
let mut reverse_pk_map = IndexMap::<OffsetEqLink<#const_name>, #pk_type>::new();
for entry in pk_map.iter() {
let (pk, link) = entry;
reverse_pk_map.insert(*link, pk.clone());
}
let primary_index = PrimaryIndex { pk_map, reverse_pk_map };
}
};
if self.attributes.read_only {
quote! {
pub fn into_worktable(self) -> #wt_ident {
let mut page_id = 1;
let data = self.data.into_iter().map(|p| {
let mut data = Data::from_data_page(p);
data.set_page_id(page_id.into());
page_id += 1;
std::sync::Arc::new(data)
})
.collect();
let data = DataPages::from_data(data)
.with_empty_links(self.data_info.inner.empty_links_list);
let indexes = #index_ident::from_persisted(self.indexes);
#primary_index_init
let table = WorkTable {
data: std::sync::Arc::new(data),
primary_index: std::sync::Arc::new(primary_index),
indexes: std::sync::Arc::new(indexes),
pk_gen: PrimaryKeyGeneratorState::from_state(self.data_info.inner.pk_gen_state),
lock_manager: std::sync::Arc::new(LockMap::<#lock_type, #pk_type>::default()),
update_state: IndexMap::default(),
table_name: #table_name,
pk_phantom: std::marker::PhantomData,
};
#wt_ident(table)
}
}
} else {
quote! {
pub async fn into_worktable<E, C>(self, engine: E) -> #wt_ident
where
E: PersistenceEngine<
<<#pk_type as TablePrimaryKey>::Generator as PrimaryKeyGeneratorState>::State,
#pk_type,
#secondary_index_events,
#avt_index_ident,
Config=C
> + Send
+ 'static,
C: Clone + PersistenceConfig,
{
let mut page_id = 1;
let data = self.data.into_iter().map(|p| {
let mut data = Data::from_data_page(p);
data.set_page_id(page_id.into());
page_id += 1;
std::sync::Arc::new(data)
})
.collect();
let data = DataPages::from_data(data)
.with_empty_links(self.data_info.inner.empty_links_list);
let indexes = #index_ident::from_persisted(self.indexes);
#primary_index_init
let table = WorkTable {
data: std::sync::Arc::new(data),
primary_index: std::sync::Arc::new(primary_index),
indexes: std::sync::Arc::new(indexes),
pk_gen: PrimaryKeyGeneratorState::from_state(self.data_info.inner.pk_gen_state),
lock_manager: std::sync::Arc::new(LockMap::<#lock_type, #pk_type>::default()),
update_state: IndexMap::default(),
table_name: #table_name,
pk_phantom: std::marker::PhantomData,
};
#wt_ident(
table,
#task_ident::run_engine(engine)
)
}
}
}
}
fn gen_space_file_parse_file_fn(&self) -> TokenStream {
let name_generator = WorktableNameGenerator::from_struct_ident(&self.struct_def.ident);
let pk_type = name_generator.get_primary_key_type_ident();
let page_const_name = name_generator.get_page_size_const_ident();
let inner_const_name = name_generator.get_page_inner_size_const_ident();
let persisted_index_name = name_generator.get_persisted_index_ident();
let index_extension = Literal::string(WT_INDEX_EXTENSION);
let data_extension = Literal::string(WT_DATA_EXTENSION);
let parse_pk_page = if self.attributes.pk_unsized {
quote! {
let index = parse_page::<UnsizedIndexPage<#pk_type, {#inner_const_name as u32}>, { #page_const_name as u32 }>(&mut primary_file, (*page_id).into()).await?;
}
} else {
quote! {
let index = parse_page::<IndexPage<#pk_type>, { #page_const_name as u32 }>(&mut primary_file, (*page_id).into()).await?;
}
};
quote! {
pub async fn parse_file(path: &str) -> eyre::Result<Self> {
let mut primary_index = {
let mut primary_index = vec![];
let mut primary_file = tokio::fs::File::open(format!("{}/primary{}", path, #index_extension)).await?;
let info = parse_page::<SpaceInfoPage<()>, { #page_const_name as u32 }>(&mut primary_file, 0).await?;
let file_length = primary_file.metadata().await?.len();
let count = file_length / (#page_const_name as u64 + GENERAL_HEADER_SIZE as u64);
let next_page_id = std::sync::Arc::new(std::sync::atomic::AtomicU32::new(count as u32));
let toc = IndexTableOfContents::<_, { #page_const_name as u32 }>::parse_from_file(&mut primary_file, 0.into(), next_page_id.clone()).await?;
for page_id in toc.iter().map(|(_, page_id)| page_id) {
#parse_pk_page
primary_index.push(index);
}
(toc.pages, primary_index)
};
let indexes = #persisted_index_name::parse_from_file(path).await?;
let (data, data_info) = {
let mut data = vec![];
let mut data_file = tokio::fs::File::open(format!("{}/{}", path, #data_extension)).await?;
let info = parse_page::<SpaceInfoPage<<<#pk_type as TablePrimaryKey>::Generator as PrimaryKeyGeneratorState>::State>, { #page_const_name as u32 }>(&mut data_file, 0).await?;
let file_length = data_file.metadata().await?.len();
let count = file_length / (#inner_const_name as u64 + GENERAL_HEADER_SIZE as u64);
for page_id in 1..=count {
let index = parse_data_page::<{ #page_const_name as u32}, { #inner_const_name as usize }>(&mut data_file, page_id as u32).await?;
data.push(index);
}
(data, info)
};
Ok(Self {
primary_index,
indexes,
data,
data_info
})
}
}
}
}