Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -3,23 +3,23 @@ members = ["codegen", "tools/create-data-file", "tools/dump-data-file"]

[package]
name = "data_bucket"
version = "0.4.1"
version = "0.5.0"
edition = "2021"
authors = ["Handy-caT"]
license = "MIT"
repository = "https://github.com/pathscale/DataBucket"
description = "DataBucket is container for WorkTable's data"

[dependencies]
data_bucket_derive = { path = "codegen", version = "=0.3.16" }
data_bucket_derive = { path = "codegen", version = "=0.3.17" }

eyre = "0.6.12"
derive_more = { version = "1.0.0", features = ["from", "error", "display", "into"] }
rkyv = { version = "0.8.9", features = ["uuid-1"] }
rkyv = { version = "0.8.17", features = ["uuid-1"] }
uuid = { version = "1.11.0", features = ["v4"] }
psc-nanoid = { version = "3.1.1", features = ["rkyv", "packed"] }
ordered-float = "5.0.0"
indexset = { package = "WorkTablesIndex", version = "=0.0.1", features = ["concurrent", "cdc", "multimap"] }
indexset = { package = "WorkTablesIndex", version = "=0.0.3", features = ["concurrent", "cdc", "multimap"] }
# indexset = { package = "wt-indexset", path = "../indexset", version = "0.12.10", features = ["concurrent", "cdc", "multimap"] }
# indexset = { package = "wt-indexset", version = "0.12.12", features = ["concurrent", "cdc", "multimap"] }
tokio = { version = "1", features = ["full"] }
Expand Down
10 changes: 1 addition & 9 deletions codegen/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "data_bucket_derive"
version = "0.3.16"
version = "0.3.17"
edition = "2021"
authors = ["Handy-caT"]
license = "MIT"
Expand All @@ -13,14 +13,6 @@ path = "src/lib.rs"
proc-macro = true

[dependencies]
rkyv = { version = "0.7.45" }
syn = { version = "2.0.74", features = ["full"] }
quote = "1.0.36"
proc-macro2 = "1.0.86"
regex = "1.10.6"
convert_case = "0.6.0"

[dev-dependencies]
derive_more = { version = "1.0.0", features = ["from", "error", "display", "into"] }
rkyv = { version = "0.7.45", features = ["uuid"] }
scc = "2.1.16"
10 changes: 5 additions & 5 deletions codegen/src/persistable/generator/obj_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,7 @@ impl Generator {
);
quote! {
pub fn #fn_ident() -> usize {
<#ty as Default>::default().aligned_size()
<#ty as SizeMeasurable>::default_aligned_size()
}
}
}
Expand Down Expand Up @@ -218,7 +218,7 @@ impl Generator {
);
let len_in_vec = if is_primitive(&inner_ty_str) {
quote! {
align(length * <#inner_ty as Default>::default().aligned_size()) + 8
align(length * <#inner_ty as SizeMeasurable>::default_aligned_size()) + 8
}
} else {
quote! {
Expand All @@ -227,14 +227,14 @@ impl Generator {
};
let len_value = if is_primitive(&inner_ty_str) {
quote! {
<#inner_ty as Default>::default().aligned_size()
<#inner_ty as SizeMeasurable>::default_aligned_size()
}
} else {
quote! {
if <#inner_ty as SizeMeasurable>::align() == Some(8) {
align8(<#inner_ty as Default>::default().aligned_size())
align8(<#inner_ty as SizeMeasurable>::default_aligned_size())
} else {
<#inner_ty as Default>::default().aligned_size()
<#inner_ty as SizeMeasurable>::default_aligned_size()
}
}
};
Expand Down
6 changes: 3 additions & 3 deletions codegen/src/persistable/generator/persistable_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,7 @@ impl Generator {
let size_type = &f.ty;
let size_ident = f.ident.as_ref().unwrap();
quote! {
let size_length = <#size_type as Default>::default().aligned_size();
let size_length = <#size_type as SizeMeasurable>::default_aligned_size();
let archived =
data_bucket::access_archived::<<#size_type as Archive>::Archived>(&bytes[offset..offset + size_length]).expect("torn or corrupt page part: a size field fails validation");
let #size_ident =
Expand All @@ -237,7 +237,7 @@ impl Generator {

fn gen_from_bytes_for_primitive(&self, ty: &Type, ident: &Ident) -> TokenStream {
quote! {
let length = <#ty as Default>::default().aligned_size();
let length = <#ty as SizeMeasurable>::default_aligned_size();
let mut v = rkyv::util::AlignedVec::<4>::new();
v.extend_from_slice(&bytes[offset..offset + length]);
let archived = data_bucket::access_archived::<<#ty as Archive>::Archived>(&v[..]).expect("torn or corrupt page part: a field fails validation");
Expand Down Expand Up @@ -271,7 +271,7 @@ impl Generator {
let value_fn_ident = Ident::new(format!("{ident}_value_size").as_str(), Span::call_site());
let len = if is_primitive(&inner_ty_str) {
quote! {
let values_len = align(#size_ident as usize * <#inner_ty as Default>::default().aligned_size()) + 8;
let values_len = align(#size_ident as usize * <#inner_ty as SizeMeasurable>::default_aligned_size()) + 8;
}
} else {
quote! {
Expand Down
62 changes: 62 additions & 0 deletions codegen/src/size_measure/enum_generator.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
use proc_macro2::TokenStream;
use quote::quote;
use syn::{Fields, ItemEnum};

pub struct EnumGenerator {
pub enum_def: ItemEnum,
}

impl EnumGenerator {
pub fn gen_impl(&self) -> syn::Result<TokenStream> {
if !self.enum_def.generics.params.is_empty() {
return Err(syn::Error::new_spanned(
&self.enum_def.generics,
"SizeMeasure does not yet support generic enums",
));
}

if let Some(variant) = self
.enum_def
.variants
.iter()
.find(|variant| !matches!(variant.fields, Fields::Unit))
{
return Err(syn::Error::new_spanned(
variant,
"SizeMeasure supports only fieldless enums; payload sizes may be data-dependent",
));
}

let enum_ident = &self.enum_def.ident;
Ok(quote! {
impl SizeMeasurable for #enum_ident
where
#enum_ident: rkyv::Archive,
<#enum_ident as rkyv::Archive>::Archived: Sized,
{
fn aligned_size(&self) -> usize {
std::mem::size_of::<<#enum_ident as rkyv::Archive>::Archived>()
}
}
})
}
}

#[cfg(test)]
mod tests {
use super::EnumGenerator;
use syn::parse_quote;

#[test]
fn rejects_payload_bearing_enums() {
let enum_def = parse_quote! {
enum Payload {
Empty,
Value(u64),
}
};

let error = EnumGenerator { enum_def }.gen_impl().unwrap_err();
assert!(error.to_string().contains("fieldless enums"));
}
}
8 changes: 4 additions & 4 deletions codegen/src/size_measure/generator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ pub struct Generator {
}

impl Generator {
pub fn gen_impl(&self) -> TokenStream {
pub fn gen_impl(&self) -> syn::Result<TokenStream> {
let struct_ident = &self.struct_def.ident;

let mut num = 0;
Expand Down Expand Up @@ -37,14 +37,14 @@ impl Generator {
.map(|f| {
let t = &f.ty;
quote! {
if #t::align() == Some(8) {
if <#t as SizeMeasurable>::align() == Some(8) {
return Some(8)
}
}
})
.collect::<Vec<_>>();

quote! {
Ok(quote! {
impl SizeMeasurable for #struct_ident {
fn aligned_size(&self) -> usize {
let len = #(#sum+)* 0;
Expand All @@ -55,6 +55,6 @@ impl Generator {
None
}
}
}
})
}
}
12 changes: 6 additions & 6 deletions codegen/src/size_measure/mod.rs
Original file line number Diff line number Diff line change
@@ -1,20 +1,20 @@
mod enum_generator;
mod generator;
mod parser;

use proc_macro2::TokenStream;
use quote::quote;

use crate::size_measure::enum_generator::EnumGenerator;
use crate::size_measure::generator::Generator;
use crate::size_measure::parser::Parser;
use crate::size_measure::parser::{ParsedItem, Parser};

pub fn expand(input: &TokenStream) -> syn::Result<TokenStream> {
let input_fn = Parser::parse_struct(input)?;
let gen = Generator {
struct_def: input_fn,
let impl_def = match Parser::parse(input)? {
ParsedItem::Struct(struct_def) => Generator { struct_def }.gen_impl()?,
ParsedItem::Enum(enum_def) => EnumGenerator { enum_def }.gen_impl()?,
};

let impl_def = gen.gen_impl();

Ok(quote! {
#impl_def
})
Expand Down
18 changes: 14 additions & 4 deletions codegen/src/size_measure/parser.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,23 @@
use proc_macro2::TokenStream;
use syn::spanned::Spanned;
use syn::ItemStruct;
use syn::{Item, ItemEnum, ItemStruct};

pub enum ParsedItem {
Struct(ItemStruct),
Enum(ItemEnum),
}

pub struct Parser;

impl Parser {
pub fn parse_struct(input: &TokenStream) -> syn::Result<ItemStruct> {
match syn::parse2::<ItemStruct>(input.clone()) {
Ok(data) => Ok(data),
pub fn parse(input: &TokenStream) -> syn::Result<ParsedItem> {
match syn::parse2::<Item>(input.clone()) {
Ok(Item::Struct(data)) => Ok(ParsedItem::Struct(data)),
Ok(Item::Enum(data)) => Ok(ParsedItem::Enum(data)),
Ok(item) => Err(syn::Error::new_spanned(
item,
"SizeMeasure supports structs and fieldless enums",
)),
Err(err) => Err(syn::Error::new(input.span(), err.to_string())),
}
}
Expand Down
2 changes: 1 addition & 1 deletion src/page/index/table_of_contents_page.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ where
Self {
records: BTreeMap::new(),
empty_pages: vec![],
estimated_size: usize::default().aligned_size() + 12,
estimated_size: <usize as SizeMeasurable>::default_aligned_size() + 12,
}
}
}
Expand Down
6 changes: 3 additions & 3 deletions src/page/iterators.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use rkyv::{de::Pool, rancor::Strategy, Archive, DeserializeUnsized};

use crate::{
page::util::parse_general_header,
persistence::data::{rkyv_data::parse_archived_row, DataTypeValue},
persistence::data::{rkyv_data::parse_archived_row, DataDecodeError, DataTypeValue},
IndexData, Link,
};

Expand Down Expand Up @@ -136,7 +136,7 @@ impl DataIterator<'_> {
}

impl Iterator for DataIterator<'_> {
type Item = Vec<DataTypeValue>;
type Item = Result<Vec<DataTypeValue>, DataDecodeError>;

fn next(&mut self) -> Option<Self::Item> {
if self.link_index >= self.links.len() {
Expand Down Expand Up @@ -210,7 +210,7 @@ mod test {
let data_iterator: DataIterator<'_> =
DataIterator::new(&mut file, space_info.row_schema, links);
assert_eq!(
data_iterator.collect::<Vec<_>>(),
data_iterator.collect::<Result<Vec<_>, _>>().unwrap(),
vec![
vec![
DataTypeValue::I32(1),
Expand Down
56 changes: 53 additions & 3 deletions src/persistence/data/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,60 @@ mod util;

pub use types::DataTypeValue;

use std::fmt;

#[derive(Clone, Debug, Eq, PartialEq)]
pub enum DataDecodeError {
UnsupportedDataType {
data_type: String,
},
BufferTooShort {
required: usize,
actual: usize,
},
FieldOutOfBounds {
field_index: usize,
field_end: usize,
actual: usize,
},
InvalidArchive {
data_type: &'static str,
message: String,
},
}

impl fmt::Display for DataDecodeError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::UnsupportedDataType { data_type } => {
write!(f, "unsupported archived row data type `{data_type}`")
}
Self::BufferTooShort { required, actual } => write!(
f,
"archived row buffer is too short: schema requires {required} bytes, buffer has {actual}"
),
Self::FieldOutOfBounds {
field_index,
field_end,
actual,
} => write!(
f,
"archived row field {field_index} ends at byte {field_end}, beyond buffer length {actual}"
),
Self::InvalidArchive { data_type, message } => {
write!(f, "invalid archived `{data_type}` field: {message}")
}
}
}
}

impl std::error::Error for DataDecodeError {}

pub trait DataType {
/// Advances an offset past this type, including its required padding.
fn advance_accum(&self, accum: &mut usize);

/// Validates and decodes a value rooted at the end of `bytes`.
#[allow(clippy::wrong_self_convention)]
fn from_pointer(&self, pointer: *const u8, start_pointer: *const u8) -> DataTypeValue;
fn advance_pointer_for_padding(&self, pointer: &mut *const u8, start_pointer: *const u8);
fn advance_pointer(&self, pointer: &mut *const u8);
fn from_archived_bytes(&self, bytes: &[u8]) -> Result<DataTypeValue, DataDecodeError>;
}
Loading
Loading