worktable_macros 0.9.0

Proc-macro companion crate for worktable: the worktable! macro and its derives. Formerly published as worktable_codegen.
Documentation
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));
                }
                // Reconstruct reverse_pk_map by iterating over pk_map
                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);
                }
                // Reconstruct reverse_pk_map by iterating over pk_map
                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
                })
            }
        }
    }
}