diff --git a/benchmarks/compress-bench/src/gpu/vortex.rs b/benchmarks/compress-bench/src/gpu/vortex.rs index bb48b461e87..8070a75ef22 100644 --- a/benchmarks/compress-bench/src/gpu/vortex.rs +++ b/benchmarks/compress-bench/src/gpu/vortex.rs @@ -22,7 +22,6 @@ use vortex::array::IntoArray; use vortex::array::VortexSessionExecute; use vortex::array::arrays::StructArray; use vortex::array::arrays::struct_::StructArrayExt; -use vortex::compressor::BtrBlocksCompressorBuilder; use vortex::error::VortexResult; use vortex::file::OpenOptionsSessionExt; use vortex::file::WriteOptionsSessionExt; @@ -36,8 +35,8 @@ use vortex_bench::compress::Compressed; use vortex_bench::compress::CompressedData; use vortex_bench::compress::Compressor; use vortex_bench::compress::Uncompressed; +use vortex_bench::compressor_builder_for_session; use vortex_bench::conversions::parquet_to_vortex_chunks_with_batch_size; -use vortex_bench::retain_edition_encodings; use vortex_cuda::CanonicalCudaExt; use vortex_cuda::CudaExecutionCtx; use vortex_cuda::CudaOpenOptionsExt; @@ -100,11 +99,9 @@ impl Compressor for GpuVortexCompressor { // partition rather than whatever the default strategy would regroup them into. let strategy = Arc::new(ChunkedLayoutStrategy::new(CompressingStrategy::new( CudaFlatLayoutStrategy::default(), - retain_edition_encodings( - &SESSION, - BtrBlocksCompressorBuilder::default().only_cuda_compatible(), - ) - .build(), + compressor_builder_for_session(&SESSION) + .only_cuda_compatible() + .build(), ))); let start = Instant::now(); SESSION diff --git a/vortex-bench/src/conversions.rs b/vortex-bench/src/conversions.rs index 6b4ed871f2e..596aca0cef8 100644 --- a/vortex-bench/src/conversions.rs +++ b/vortex-bench/src/conversions.rs @@ -37,7 +37,6 @@ use vortex::array::arrays::struct_::StructArrayExt; use vortex::array::builders::builder_with_capacity_in; use vortex::array::stream::ArrayStreamAdapter; use vortex::array::stream::ArrayStreamExt; -use vortex::compressor::BtrBlocksCompressorBuilder; use vortex::dtype::DType; use vortex::dtype::FieldPath; use vortex::dtype::StructFields; @@ -66,7 +65,7 @@ use wkb::writer::write_geometry; use crate::CompactionStrategy; use crate::Format; use crate::SESSION; -use crate::retain_edition_encodings; +use crate::compressor_builder_for_session; use crate::utils::file::idempotent_async; /// Memory budget per concurrent conversion stream in GB. This is somewhat arbitary. @@ -248,10 +247,8 @@ fn write_options_for( let mut builder = WriteStrategyBuilder::default(); if matches!(compaction, CompactionStrategy::Compact) { - builder = builder.with_btrblocks_builder(retain_edition_encodings( - &SESSION, - BtrBlocksCompressorBuilder::default().with_compact(), - )); + builder = + builder.with_btrblocks_builder(compressor_builder_for_session(&SESSION).with_compact()); } for name in binary_fields { builder = builder.with_field_writer(FieldPath::from_name(name), no_dict_layout()); @@ -263,7 +260,7 @@ fn write_options_for( fn no_dict_layout() -> Arc { Arc::new(CompressingStrategy::new( ChunkedLayoutStrategy::new(FlatLayoutStrategy::default()), - retain_edition_encodings(&SESSION, BtrBlocksCompressorBuilder::default()).build(), + compressor_builder_for_session(&SESSION).build(), )) } diff --git a/vortex-bench/src/lib.rs b/vortex-bench/src/lib.rs index 569dfc74a4d..7dcfb8a01f3 100644 --- a/vortex-bench/src/lib.rs +++ b/vortex-bench/src/lib.rs @@ -255,10 +255,7 @@ impl CompactionStrategy { match self { CompactionStrategy::Compact => options.with_strategy( WriteStrategyBuilder::default() - .with_btrblocks_builder(retain_edition_encodings( - &SESSION, - BtrBlocksCompressorBuilder::default().with_compact(), - )) + .with_btrblocks_builder(compressor_builder_for_session(&SESSION).with_compact()) .build(), ), CompactionStrategy::Default => options, @@ -266,19 +263,16 @@ impl CompactionStrategy { } } -/// Restrict `builder` to the encodings permitted by the session's enabled editions. +/// Create a compressor builder permitting the session's enabled array encodings. /// -/// The default writer applies this filter itself. An explicit strategy bypasses it, so a -/// benchmark that builds its own compressor applies it here to stay within editions. -pub fn retain_edition_encodings( - session: &VortexSession, - builder: BtrBlocksCompressorBuilder, -) -> BtrBlocksCompressorBuilder { +/// Benchmarks supplying an explicit strategy use the session's permissions, including opt-in +/// editions, instead of the compressor builder's default core edition. +pub fn compressor_builder_for_session(session: &VortexSession) -> BtrBlocksCompressorBuilder { let allowed = session .enabled_component_ids(ComponentKind::Array) .into_iter() .collect(); - builder.retain_allowed_encodings(&allowed) + BtrBlocksCompressorBuilder::new(allowed) } /// Verify that local data has already been prepared for the requested benchmark formats. diff --git a/vortex-btrblocks/Cargo.toml b/vortex-btrblocks/Cargo.toml index 24a03337768..a92c0a27fa9 100644 --- a/vortex-btrblocks/Cargo.toml +++ b/vortex-btrblocks/Cargo.toml @@ -26,6 +26,7 @@ vortex-buffer = { workspace = true } vortex-compressor = { workspace = true } vortex-datetime-parts = { workspace = true } vortex-decimal-byte-parts = { workspace = true } +vortex-edition = { workspace = true } vortex-error = { workspace = true } vortex-fastlanes = { workspace = true } vortex-fsst = { workspace = true } @@ -49,7 +50,6 @@ tpchgen = { workspace = true } tpchgen-arrow = { workspace = true } vortex-array = { workspace = true, features = ["_test-harness"] } vortex-arrow = { workspace = true } -vortex-edition = { workspace = true } vortex-mask = { workspace = true } vortex-session = { workspace = true } diff --git a/vortex-btrblocks/src/builder.rs b/vortex-btrblocks/src/builder.rs index 3bcda909227..61e56951874 100644 --- a/vortex-btrblocks/src/builder.rs +++ b/vortex-btrblocks/src/builder.rs @@ -4,8 +4,12 @@ //! Builder for configuring `BtrBlocksCompressor` instances. use vortex_array::ArrayId; +use vortex_decimal_byte_parts::decimal_byte_parts_v2_id; +use vortex_edition::DEFAULT_CORE_EDITION; +use vortex_edition::array_ids_for_edition; use vortex_utils::aliases::hash_set::HashSet; +use crate::AllowedSerializedIds; use crate::BtrBlocksCompressor; use crate::CascadingCompressor; use crate::Scheme; @@ -18,10 +22,14 @@ use crate::schemes::integer; use crate::schemes::string; use crate::schemes::temporal; -/// All available compression schemes. +/// All compression schemes. /// /// This list is order-sensitive: the builder preserves this order when constructing /// the final scheme list, so that tie-breaking is deterministic. +/// +/// If a scheme can be configured to support different editions like +/// [`DecimalScheme`](crate::schemes::decimal), put oldest version here to defer upgrading +/// to newest supported version when compressor is built. pub const ALL_SCHEMES: &[&dyn Scheme] = &[ //////////////////////////////////////////////////////////////////////////////////////////////// // Integer schemes. @@ -60,32 +68,41 @@ pub const ALL_SCHEMES: &[&dyn Scheme] = &[ &binary::BinaryDictScheme, &binary::VarBinScheme, // Decimal schemes. - &decimal::DecimalScheme, + // Use v1 by default and let builder upgrade to v2 if permitted by edition. + &decimal::DecimalScheme::v1(), // Temporal schemes. &temporal::TemporalScheme, ]; /// Delta, kept out of [`ALL_SCHEMES`] because it is slower to decompress than the schemes that /// would otherwise win. Callers that want it opt in with -/// [`with_new_scheme`](BtrBlocksCompressorBuilder::with_new_scheme). +/// [`with_new_scheme`](BtrBlocksCompressorBuilder::with_new_scheme) and permit `fastlanes.delta` +/// with [`BtrBlocksCompressorBuilder::allow_encodings`]. /// /// TODO(robert): Return it to [`ALL_SCHEMES`] once we have scheme filtering. pub static DELTA_SCHEME: integer::DeltaScheme = integer::DeltaScheme::new(1.25); /// Builder for creating configured [`BtrBlocksCompressor`] instances. /// -/// By default, all schemes in [`ALL_SCHEMES`] are enabled in a deterministic order. Feature-gated +/// By default, all schemes in [`ALL_SCHEMES`] are registered in a deterministic order. Feature-gated /// schemes (Pco, Zstd) are not in `ALL_SCHEMES` and must be added explicitly via /// [`with_new_scheme`](BtrBlocksCompressorBuilder::with_new_scheme) or `with_compact` when the /// `zstd` feature is enabled. /// +/// [`Self::new`] takes the initial permitted serialized IDs. [`Self::set_allowed_encodings`] +/// replaces these permissions, and [`Self::allow_encodings`] extends them. During [`Self::build`], +/// the final permissions allow scheme upgrades and filter all registered schemes. The default +/// builder permits the array IDs in [`DEFAULT_CORE_EDITION`]. Decimal defaults to v1 and upgrades +/// to v2 when both serialized IDs are permitted. +/// /// # Examples /// /// ```rust /// use vortex_btrblocks::{BtrBlocksCompressorBuilder, Scheme, SchemeExt}; /// use vortex_btrblocks::schemes::integer::IntDictScheme; /// -/// // Default compressor with all schemes in ALL_SCHEMES. +/// // Default compressor with all schemes in ALL_SCHEMES, restricted to the +/// // default core edition. /// let compressor = BtrBlocksCompressorBuilder::default().build(); /// /// // Remove specific schemes. @@ -96,30 +113,74 @@ pub static DELTA_SCHEME: integer::DeltaScheme = integer::DeltaScheme::new(1.25); #[derive(Debug, Clone)] pub struct BtrBlocksCompressorBuilder { schemes: Vec<&'static dyn Scheme>, + allowed_serialized_ids: AllowedSerializedIds, } impl Default for BtrBlocksCompressorBuilder { + /// Uses the default core edition's serialized array IDs. Use [`Self::allow_encodings`] to + /// permit additional encodings; otherwise, schemes requiring them are omitted at build. fn default() -> Self { + Self::new(array_ids_for_edition(&DEFAULT_CORE_EDITION).collect()) + } +} + +impl BtrBlocksCompressorBuilder { + /// Creates a builder with all default schemes and the supplied serialized ID permissions. + /// + /// An empty set permits no serialized IDs. Upgrades and filtering are deferred until + /// [`Self::build`], including for schemes registered later. Schemes are never downgraded. + pub fn new(allowed_serialized_ids: AllowedSerializedIds) -> Self { Self { schemes: ALL_SCHEMES.to_vec(), + allowed_serialized_ids, } } -} -impl BtrBlocksCompressorBuilder { - /// Creates a builder with no schemes registered. + /// Creates a builder with no registered schemes and no permitted serialized IDs. /// /// Useful when the caller wants explicit, scheme-by-scheme control over the compressor. + /// Register schemes with [`Self::with_new_scheme`] and permit their serialized IDs with + /// [`Self::allow_encodings`]. + /// + /// ```rust + /// use vortex_btrblocks::{BtrBlocksCompressorBuilder, Scheme}; + /// use vortex_btrblocks::schemes::integer::FoRScheme; + /// + /// let compressor = BtrBlocksCompressorBuilder::empty() + /// .allow_encodings(FoRScheme.produced_encodings()) + /// .with_new_scheme(&FoRScheme) + /// .build(); + /// ``` pub fn empty() -> Self { Self { schemes: Vec::new(), + allowed_serialized_ids: HashSet::new(), } } + /// Replaces the permitted serialized IDs, including any default edition permissions. + /// + /// An empty iterator permits no serialized IDs. The final permissions apply to every + /// registered scheme during [`Self::build`]. + pub fn set_allowed_encodings(mut self, ids: impl IntoIterator) -> Self { + self.allowed_serialized_ids = ids.into_iter().collect(); + self + } + + /// Adds permitted serialized IDs while preserving existing permissions. + /// + /// The final permissions apply to every registered scheme during [`Self::build`]. + pub fn allow_encodings(mut self, ids: impl IntoIterator) -> Self { + self.allowed_serialized_ids.extend(ids); + self + } + /// Adds an external compression scheme not in [`ALL_SCHEMES`]. /// /// This allows encoding crates outside of `vortex-btrblocks` to register their own schemes /// with the compressor. + /// Schemes with unpermitted outputs are silently omitted during [`Self::build`]; use + /// [`Self::allow_encodings`] to permit outputs outside the default core edition. /// /// # Panics /// @@ -163,8 +224,13 @@ impl BtrBlocksCompressorBuilder { /// /// Both the array-level and the buffer-level Zstd schemes are added. Buffer-level /// compression preserves binary arrays' buffer layout for zero-conversion GPU decompression, - /// but belongs to the opt-in `zstd` edition, so callers filter the two through - /// [`retain_allowed_encodings`](Self::retain_allowed_encodings). + /// but belongs to the opt-in `zstd` edition. The final permissions determine which schemes + /// survive. + /// The decimal v2 serialized ID is excluded because CUDA does not support lower decimal + /// parts. This prevents v1 schemes from upgrading and filters out explicitly registered v2 + /// schemes, including Decimal schemes registered after this call. + /// Apply this preset after changing permissions: [`Self::set_allowed_encodings`] and + /// [`Self::allow_encodings`] can re-enable decimal v2 if called afterwards. /// /// This preset is intended for files that will be decoded by CUDA kernels. It may choose a /// larger encoded representation than the default compressor. @@ -189,7 +255,10 @@ impl BtrBlocksCompressorBuilder { excluded.push(integer::DeltaScheme::default().id()); #[cfg(feature = "pco")] excluded.extend([integer::PcoScheme.id(), float::PcoScheme.id()]); - let builder = self.exclude_schemes(excluded); + let mut builder = self.exclude_schemes(excluded); + builder + .allowed_serialized_ids + .remove(&decimal_byte_parts_v2_id()); #[cfg(feature = "zstd")] let builder = builder @@ -206,33 +275,47 @@ impl BtrBlocksCompressorBuilder { self } - /// Retains only schemes whose produced serialized IDs all belong to `allowed`. - /// - /// `allowed` holds serialized IDs. The file writer passes the array IDs its enabled editions - /// permit. - pub fn retain_allowed_encodings(mut self, allowed: &HashSet) -> Self { - self.schemes - .retain(|s| s.produced_encodings().iter().all(|id| allowed.contains(id))); - self - } - /// Builds the configured [`BtrBlocksCompressor`]. pub fn build(self) -> BtrBlocksCompressor { - BtrBlocksCompressor(CascadingCompressor::new(self.schemes)) + BtrBlocksCompressor(CascadingCompressor::new(self.configured_schemes())) + } + + fn configured_schemes(self) -> Vec<&'static dyn Scheme> { + let mut final_schemes = Vec::with_capacity(self.schemes.len()); + let allowed = &self.allowed_serialized_ids; + for scheme in self.schemes { + let scheme = scheme.try_upgrade(allowed).unwrap_or(scheme); + if scheme + .produced_encodings() + .iter() + .all(|id| allowed.contains(id)) + { + final_schemes.push(scheme); + } + } + final_schemes } } #[cfg(test)] mod tests { use vortex_array::VTable; + use vortex_fastlanes::Delta; use vortex_fastlanes::FoR; use super::*; #[test] - fn empty_starts_with_no_schemes() { + fn empty_starts_with_no_schemes_or_permissions() { let builder = BtrBlocksCompressorBuilder::empty(); assert!(builder.schemes.is_empty()); + assert!(builder.allowed_serialized_ids.is_empty()); + + let builder = builder.with_new_scheme(&integer::FoRScheme); + assert!(builder.clone().configured_schemes().is_empty()); + let schemes = builder.allow_encodings([FoR.id()]).configured_schemes(); + assert_eq!(schemes.len(), 1); + assert_eq!(schemes[0].id(), integer::FoRScheme.id()); } #[test] @@ -242,24 +325,81 @@ mod tests { } #[test] - fn retain_allowed_encodings_filters_schemes() { - let allowed: HashSet = [FoR.id()].into_iter().collect(); - let builder = BtrBlocksCompressorBuilder::default().retain_allowed_encodings(&allowed); - assert_eq!(builder.schemes.len(), 1); - assert_eq!(builder.schemes[0].id(), integer::FoRScheme.id()); - - let none = BtrBlocksCompressorBuilder::default().retain_allowed_encodings(&HashSet::new()); - assert!(none.schemes.is_empty()); + fn allowed_encodings_filter_schemes() { + let schemes = + BtrBlocksCompressorBuilder::new(HashSet::from([FoR.id()])).configured_schemes(); + assert_eq!(schemes.len(), 1); + assert_eq!(schemes[0].id(), integer::FoRScheme.id()); + + let none = BtrBlocksCompressorBuilder::new(HashSet::new()).configured_schemes(); + assert!(none.is_empty()); } #[test] - fn retaining_all_declared_outputs_keeps_every_scheme() { - let allowed: HashSet = ALL_SCHEMES + fn set_allowed_encodings_replaces_permissions() { + let builder = BtrBlocksCompressorBuilder::default().set_allowed_encodings([FoR.id()]); + let schemes = builder.clone().configured_schemes(); + assert_eq!(schemes.len(), 1); + assert_eq!(schemes[0].id(), integer::FoRScheme.id()); + assert!( + builder + .set_allowed_encodings([]) + .configured_schemes() + .is_empty() + ); + } + + #[test] + fn allow_encodings_keeps_defaults_and_enables_registered_schemes() { + let builder = BtrBlocksCompressorBuilder::default().with_new_scheme(&DELTA_SCHEME); + assert_eq!(builder.clone().configured_schemes(), ALL_SCHEMES); + + let schemes = builder.allow_encodings([Delta.id()]).configured_schemes(); + assert_eq!(&schemes[..ALL_SCHEMES.len()], ALL_SCHEMES); + assert_eq!(schemes.len(), ALL_SCHEMES.len() + 1); + assert_eq!(schemes[ALL_SCHEMES.len()].id(), DELTA_SCHEME.id()); + } + + #[test] + fn default_configuration_preserves_scheme_order() { + let schemes = BtrBlocksCompressorBuilder::default().configured_schemes(); + assert_eq!(schemes, ALL_SCHEMES); + } + + #[test] + fn allowing_all_declared_outputs_keeps_every_scheme() { + let allowed = ALL_SCHEMES .iter() .flat_map(|scheme| scheme.produced_encodings()) .collect(); - let builder = BtrBlocksCompressorBuilder::default().retain_allowed_encodings(&allowed); - assert_eq!(builder.schemes.len(), ALL_SCHEMES.len()); + let schemes = BtrBlocksCompressorBuilder::new(allowed).configured_schemes(); + assert_eq!(schemes, ALL_SCHEMES); + } + + #[test] + fn every_declared_output_must_be_permitted() { + for scheme in ALL_SCHEMES { + for id in scheme.produced_encodings() { + let mut allowed: HashSet<_> = scheme.produced_encodings().into_iter().collect(); + allowed.remove(&id); + let builder = BtrBlocksCompressorBuilder::empty() + .allow_encodings(allowed) + .with_new_scheme(*scheme); + assert!( + builder.configured_schemes().is_empty(), + "{} requires {id}", + scheme.scheme_name() + ); + } + } + } + + #[test] + fn permissions_apply_to_later_registrations() { + let builder = BtrBlocksCompressorBuilder::new(HashSet::new()) + .exclude_schemes([integer::FoRScheme.id()]) + .with_new_scheme(&integer::FoRScheme); + assert!(builder.configured_schemes().is_empty()); } #[test] diff --git a/vortex-btrblocks/src/canonical_compressor.rs b/vortex-btrblocks/src/canonical_compressor.rs index d93be365550..f080de89ba9 100644 --- a/vortex-btrblocks/src/canonical_compressor.rs +++ b/vortex-btrblocks/src/canonical_compressor.rs @@ -292,9 +292,8 @@ mod tests { // The CUDA preset carries both Zstd schemes; the edition filter decides which one // survives. - let compressor = BtrBlocksCompressorBuilder::default() + let compressor = BtrBlocksCompressorBuilder::new(HashSet::from([allowed])) .only_cuda_compatible() - .retain_allowed_encodings(&HashSet::from([allowed])) .build(); let mut ctx = SESSION.create_execution_ctx(); let compressed = compressor.compress(&array.clone().into_array(), &mut ctx)?; diff --git a/vortex-btrblocks/src/lib.rs b/vortex-btrblocks/src/lib.rs index 2e8ae484f90..680bfd57676 100644 --- a/vortex-btrblocks/src/lib.rs +++ b/vortex-btrblocks/src/lib.rs @@ -82,6 +82,7 @@ pub use builder::DELTA_SCHEME; pub use canonical_compressor::BtrBlocksCompressor; pub use schemes::patches::compress_patches; pub use vortex_compressor::CascadingCompressor; +pub use vortex_compressor::scheme::AllowedSerializedIds; pub use vortex_compressor::scheme::CompressorContext; pub use vortex_compressor::scheme::MAX_CASCADE; pub use vortex_compressor::scheme::Scheme; diff --git a/vortex-btrblocks/src/schemes/decimal.rs b/vortex-btrblocks/src/schemes/decimal.rs index f77a77d8c50..e50ce4ee5ca 100644 --- a/vortex-btrblocks/src/schemes/decimal.rs +++ b/vortex-btrblocks/src/schemes/decimal.rs @@ -9,13 +9,16 @@ use vortex_array::Canonical; use vortex_array::ExecutionCtx; use vortex_array::IntoArray; use vortex_array::arrays::DecimalArray; -use vortex_array::arrays::PrimitiveArray; use vortex_array::arrays::decimal::narrowed_decimal; use vortex_array::dtype::DecimalType; +use vortex_compressor::scheme::AllowedSerializedIds; use vortex_compressor::scheme::CompressionEstimate; use vortex_compressor::scheme::EstimateVerdict; use vortex_decimal_byte_parts::DecimalByteParts; +use vortex_decimal_byte_parts::DecimalBytePartsSlots; use vortex_decimal_byte_parts::decimal_byte_parts_v1_id; +use vortex_decimal_byte_parts::decimal_byte_parts_v2_id; +use vortex_decimal_byte_parts::split_decimal; use vortex_error::VortexResult; use crate::ArrayAndStats; @@ -24,12 +27,53 @@ use crate::CompressorContext; use crate::Scheme; use crate::SchemeExt; +#[derive(Debug, Copy, Clone, PartialEq, Eq)] +enum DecimalSchemeMode { + V1, + V2, +} + +static DECIMAL_V2: DecimalScheme = DecimalScheme::v2(); + /// Compression scheme for decimal arrays via byte-part decomposition. /// -/// Narrows the decimal to the smallest integer type, compresses the underlying primitive, and wraps -/// the result in a `DecimalBytePartsArray`. +/// Narrows the decimal to the smallest integer type and compresses its byte parts independently. +/// The v1 mode leaves values wider than `i64` canonical; v2 splits them into a signed most +/// significant part and up to three unsigned lower parts. Single-part arrays serialize as v1 +/// in either mode, while arrays with lower parts serialize as v2. +/// +/// The default uses v1. Permitting both decimal IDs lets the builder upgrade v1 to v2. +/// A v2 scheme is filtered out if either ID is forbidden, including under the CUDA preset. #[derive(Debug, Copy, Clone, PartialEq, Eq)] -pub struct DecimalScheme; +pub struct DecimalScheme { + mode: DecimalSchemeMode, +} + +impl DecimalScheme { + /// Creates a decimal scheme configured for v1, disallowing splitting of wide decimals. + /// + /// Values that remain wider than `i64` after narrowing stay canonical. + /// The builder may upgrade to v2 if its permissions include both serialized IDs. + pub const fn v1() -> Self { + Self { + mode: DecimalSchemeMode::V1, + } + } + + /// Creates a decimal scheme configured for v2, allowing splitting of wide decimals. + /// The builder filters this scheme out if either decimal serialized ID is forbidden. + pub const fn v2() -> Self { + Self { + mode: DecimalSchemeMode::V2, + } + } +} + +impl Default for DecimalScheme { + fn default() -> Self { + Self::v1() + } +} impl Scheme for DecimalScheme { fn scheme_name(&self) -> &'static str { @@ -41,14 +85,27 @@ impl Scheme for DecimalScheme { } fn produced_encodings(&self) -> Vec { - // This scheme only builds single-part arrays, which serialize under the frozen v1 ID. - // The in-memory ID is the v2 wire ID, which no edition permits yet. - vec![decimal_byte_parts_v1_id()] + match self.mode { + DecimalSchemeMode::V1 => vec![decimal_byte_parts_v1_id()], + DecimalSchemeMode::V2 => { + vec![decimal_byte_parts_v1_id(), decimal_byte_parts_v2_id()] + } + } } - /// Children: primitive=0. + fn try_upgrade(&self, allowed_serialized_ids: &AllowedSerializedIds) -> Option<&dyn Scheme> { + (self.mode == DecimalSchemeMode::V1 + && allowed_serialized_ids.contains(&decimal_byte_parts_v1_id()) + && allowed_serialized_ids.contains(&decimal_byte_parts_v2_id())) + .then_some(&DECIMAL_V2 as &dyn Scheme) + } + + /// Children: msp=0, then up to three lower parts in v2 mode. fn num_children(&self) -> usize { - 1 + match self.mode { + DecimalSchemeMode::V1 => 1, + DecimalSchemeMode::V2 => 4, + } } fn expected_compression_ratio( @@ -68,22 +125,38 @@ impl Scheme for DecimalScheme { compress_ctx: CompressorContext, exec_ctx: &mut ExecutionCtx, ) -> VortexResult { - // TODO(joe): add support splitting i128/256 buffers into chunks of primitive values - // for compression. 2 for i128 and 4 for i256. let decimal = data.array().clone().execute::(exec_ctx)?; let decimal = narrowed_decimal(decimal); - let validity = decimal.validity()?; - let prim = match decimal.values_type() { - DecimalType::I8 => PrimitiveArray::new(decimal.buffer::(), validity), - DecimalType::I16 => PrimitiveArray::new(decimal.buffer::(), validity), - DecimalType::I32 => PrimitiveArray::new(decimal.buffer::(), validity), - DecimalType::I64 => PrimitiveArray::new(decimal.buffer::(), validity), - _ => return Ok(decimal.into_array()), - }; - - let compressed = - compressor.compress_child(&prim.into_array(), &compress_ctx, self.id(), 0, exec_ctx)?; - - DecimalByteParts::try_new(compressed, decimal.decimal_dtype()).map(|d| d.into_array()) + if self.mode == DecimalSchemeMode::V1 + && matches!(decimal.values_type(), DecimalType::I128 | DecimalType::I256) + { + return Ok(decimal.into_array()); + } + + let parts = split_decimal(&decimal, exec_ctx)?; + let msp = compressor.compress_child( + &parts.msp, + &compress_ctx, + self.id(), + DecimalBytePartsSlots::MSP, + exec_ctx, + )?; + let lower_parts = parts + .lower_parts + .iter() + .enumerate() + .map(|(idx, part)| { + compressor.compress_child( + part, + &compress_ctx, + self.id(), + DecimalBytePartsSlots::LOWER_PARTS_OFFSET + idx, + exec_ctx, + ) + }) + .collect::>>()?; + + DecimalByteParts::try_new_with_lower_parts(msp, lower_parts, decimal.decimal_dtype()) + .map(IntoArray::into_array) } } diff --git a/vortex-btrblocks/src/schemes/integer/scheme_selection_tests.rs b/vortex-btrblocks/src/schemes/integer/scheme_selection_tests.rs index b4726dab9b5..03888f6baf5 100644 --- a/vortex-btrblocks/src/schemes/integer/scheme_selection_tests.rs +++ b/vortex-btrblocks/src/schemes/integer/scheme_selection_tests.rs @@ -10,6 +10,7 @@ use rand::Rng; use rand::SeedableRng; use rand::rngs::StdRng; use vortex_array::IntoArray; +use vortex_array::VTable; use vortex_array::VortexSessionExecute; use vortex_array::arrays::Constant; use vortex_array::arrays::Dict; @@ -19,8 +20,11 @@ use vortex_array::expr::stats::Stat; use vortex_array::expr::stats::StatsProviderExt; use vortex_array::validity::Validity; use vortex_buffer::Buffer; +use vortex_edition::DEFAULT_CORE_EDITION; +use vortex_edition::array_ids_for_edition; use vortex_error::VortexResult; use vortex_fastlanes::BitPacked; +use vortex_fastlanes::Delta; use vortex_fastlanes::FoR; use vortex_runend::RunEnd; use vortex_sequence::Sequence; @@ -32,6 +36,16 @@ use crate::BtrBlocksCompressorBuilder; use crate::DELTA_SCHEME; static SESSION: LazyLock = LazyLock::new(vortex_array::array_session); +fn delta_compressor() -> BtrBlocksCompressor { + // Delta is outside the default core edition, so it needs an explicit permission too. + let allowed = array_ids_for_edition(&DEFAULT_CORE_EDITION) + .chain(iter::once(Delta.id())) + .collect(); + BtrBlocksCompressorBuilder::new(allowed) + .with_new_scheme(&DELTA_SCHEME) + .build() +} + #[test] fn test_constant_compressed() -> VortexResult<()> { let values: Vec = iter::repeat_n(42, 100).collect(); @@ -164,7 +178,6 @@ fn test_rle_compressed() -> VortexResult<()> { fn test_delta_compressed() -> VortexResult<()> { let mut ctx = SESSION.create_execution_ctx(); use vortex_array::assert_arrays_eq; - use vortex_fastlanes::Delta; let mut rng = StdRng::seed_from_u64(7u64); let mut value = 500_000i32; @@ -176,9 +189,7 @@ fn test_delta_compressed() -> VortexResult<()> { .collect(); let array = PrimitiveArray::new(Buffer::copy_from(&values), Validity::NonNullable); - let btr = BtrBlocksCompressorBuilder::default() - .with_new_scheme(&DELTA_SCHEME) - .build(); + let btr = delta_compressor(); let compressed = btr.compress( &array.clone().into_array(), &mut SESSION.create_execution_ctx(), @@ -204,7 +215,6 @@ fn test_delta_compressed() -> VortexResult<()> { fn test_delta_compressed_unaligned_length() -> VortexResult<()> { let mut ctx = SESSION.create_execution_ctx(); use vortex_array::assert_arrays_eq; - use vortex_fastlanes::Delta; let mut rng = StdRng::seed_from_u64(7u64); let mut value = 500_000i32; @@ -216,9 +226,7 @@ fn test_delta_compressed_unaligned_length() -> VortexResult<()> { .collect(); let array = PrimitiveArray::new(Buffer::copy_from(&values), Validity::NonNullable); - let btr = BtrBlocksCompressorBuilder::default() - .with_new_scheme(&DELTA_SCHEME) - .build(); + let btr = delta_compressor(); let compressed = btr.compress( &array.clone().into_array(), &mut SESSION.create_execution_ctx(), @@ -239,15 +247,12 @@ fn test_delta_compressed_unaligned_length() -> VortexResult<()> { fn test_delta_nullable_unaligned_sum() -> VortexResult<()> { use vortex_array::aggregate_fn::fns::sum::sum; use vortex_array::assert_arrays_eq; - use vortex_fastlanes::Delta; let mut ctx = SESSION.create_execution_ctx(); let array = PrimitiveArray::from_option_iter(iter::once(None).chain((1i32..=100_000).map(Some))); - let btr = BtrBlocksCompressorBuilder::default() - .with_new_scheme(&DELTA_SCHEME) - .build(); + let btr = delta_compressor(); let compressed = btr.compress(&array.clone().into_array(), &mut ctx)?; assert!( compressed.is::(), @@ -266,8 +271,6 @@ fn test_delta_nullable_unaligned_sum() -> VortexResult<()> { /// Returns true if any `Delta` array appears below an ancestor `Delta` in the tree. fn has_nested_delta(array: &vortex_array::ArrayRef, under_delta: bool) -> bool { - use vortex_fastlanes::Delta; - let is_delta = array.is::(); if is_delta && under_delta { return true; diff --git a/vortex-btrblocks/src/schemes/string/scheme_selection_tests.rs b/vortex-btrblocks/src/schemes/string/scheme_selection_tests.rs index aac0b4de4de..60a23d488a3 100644 --- a/vortex-btrblocks/src/schemes/string/scheme_selection_tests.rs +++ b/vortex-btrblocks/src/schemes/string/scheme_selection_tests.rs @@ -17,6 +17,11 @@ use vortex_fsst::FSST; use vortex_session::VortexSession; use crate::BtrBlocksCompressor; +use crate::BtrBlocksCompressorBuilder; +use crate::Scheme; +use crate::SchemeExt; +use crate::schemes::string::FSSTScheme; +use crate::schemes::string::onpair::OnPairScheme; static SESSION: LazyLock = LazyLock::new(vortex_array::array_session); @@ -48,9 +53,6 @@ fn test_dict_compressed() -> VortexResult<()> { #[test] fn test_all_schemes_includes_onpair() { - use crate::SchemeExt; - use crate::schemes::string::onpair::OnPairScheme; - let ids: Vec<_> = crate::ALL_SCHEMES.iter().map(|s| s.id()).collect(); assert!( ids.contains(&OnPairScheme.id()), @@ -85,10 +87,6 @@ fn test_default_btrblocks_compressor_selects_onpair() -> VortexResult<()> { /// still produces an FSST array. #[test] fn test_fsst_in_default_scheme_list() -> VortexResult<()> { - use crate::BtrBlocksCompressorBuilder; - use crate::SchemeExt; - use crate::schemes::string::FSSTScheme; - // FSST is registered by default. assert!( crate::ALL_SCHEMES.iter().any(|s| s.id() == FSSTScheme.id()), @@ -107,6 +105,7 @@ fn test_fsst_in_default_scheme_list() -> VortexResult<()> { let array_ref = array.into_array(); let compressor = BtrBlocksCompressorBuilder::empty() + .allow_encodings(FSSTScheme.produced_encodings()) .with_new_scheme(&FSSTScheme) .build(); let compressed = compressor.compress(&array_ref, &mut SESSION.create_execution_ctx())?; diff --git a/vortex-btrblocks/src/trace_tests.rs b/vortex-btrblocks/src/trace_tests.rs index e23e4ef0244..39058b3b769 100644 --- a/vortex-btrblocks/src/trace_tests.rs +++ b/vortex-btrblocks/src/trace_tests.rs @@ -11,6 +11,7 @@ //! rules and execute kernels fire for the scan operations TPC-H queries perform over those //! encodings. +use std::iter; use std::sync::LazyLock; use arrow_array::RecordBatch; @@ -22,6 +23,7 @@ use vortex_array::ArrayRef; use vortex_array::Canonical; use vortex_array::ExecutionCtx; use vortex_array::IntoArray; +use vortex_array::VTable; use vortex_array::arrays::ConstantArray; use vortex_array::arrays::DictArray; use vortex_array::arrays::FilterArray; @@ -49,7 +51,10 @@ use vortex_array::session::ArraySessionExt; use vortex_array::test_harness::trace::Traced; use vortex_array::test_harness::trace::trace_op; use vortex_arrow::ArrowSessionExt; +use vortex_edition::DEFAULT_CORE_EDITION; +use vortex_edition::array_ids_for_edition; use vortex_error::VortexResult; +use vortex_fastlanes::Delta; use vortex_mask::Mask; use vortex_session::VortexSession; @@ -125,9 +130,12 @@ fn lineitem() -> VortexResult { .from_arrow_record_batch(batch, &schema) } -/// Delta is opt-in, and these traces cover the delta-encoded FSST offsets, so enable it here. +/// These traces cover delta-encoded FSST offsets, so permit and register the opt-in Delta scheme. fn compressed_lineitem() -> VortexResult { - BtrBlocksCompressorBuilder::default() + let allowed = array_ids_for_edition(&DEFAULT_CORE_EDITION) + .chain(iter::once(Delta.id())) + .collect(); + BtrBlocksCompressorBuilder::new(allowed) .with_new_scheme(&DELTA_SCHEME) .build() .compress(&lineitem()?, &mut execution_ctx()) diff --git a/vortex-btrblocks/tests/decimal_config.rs b/vortex-btrblocks/tests/decimal_config.rs new file mode 100644 index 00000000000..e20b52aca65 --- /dev/null +++ b/vortex-btrblocks/tests/decimal_config.rs @@ -0,0 +1,271 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright the Vortex contributors + +//! Decimal mode selection, serialized permissions, and compression of wide decimal parts. + +#![cfg(test)] + +use std::sync::LazyLock; + +use rstest::rstest; +use vortex_array::ArrayContext; +use vortex_array::ArrayId; +use vortex_array::ArrayRef; +use vortex_array::IntoArray; +use vortex_array::VortexSessionExecute; +use vortex_array::arrays::Decimal; +use vortex_array::arrays::DecimalArray; +use vortex_array::assert_arrays_eq; +use vortex_array::dtype::DecimalDType; +use vortex_array::dtype::i256; +use vortex_array::serde::SerializeOptions; +use vortex_array::serde::SerializedArray; +use vortex_array::session::ArraySessionExt; +use vortex_array::validity::Validity; +use vortex_btrblocks::BtrBlocksCompressorBuilder; +use vortex_btrblocks::Scheme; +use vortex_btrblocks::SchemeExt; +use vortex_btrblocks::schemes::decimal::DecimalScheme; +use vortex_btrblocks::schemes::integer::BitPackingScheme; +use vortex_btrblocks::schemes::integer::FoRScheme; +use vortex_buffer::Buffer; +use vortex_buffer::ByteBufferMut; +use vortex_decimal_byte_parts::DecimalByteParts; +use vortex_decimal_byte_parts::DecimalBytePartsArraySlotsExt; +use vortex_decimal_byte_parts::decimal_byte_parts_v1_id; +use vortex_decimal_byte_parts::decimal_byte_parts_v2_id; +use vortex_error::VortexResult; +use vortex_error::vortex_err; +use vortex_session::VortexSession; +use vortex_session::registry::ReadContext; +use vortex_utils::aliases::hash_set::HashSet; + +static DECIMAL_V1: DecimalScheme = DecimalScheme::v1(); +static DECIMAL_V2: DecimalScheme = DecimalScheme::v2(); + +static SESSION: LazyLock = LazyLock::new(|| { + let session = vortex_array::array_session(); + vortex_decimal_byte_parts::initialize(&session); + vortex_fastlanes::initialize(&session); + session +}); + +fn decimal_array(wide: bool) -> ArrayRef { + let base = if wide { 1i128 << 70 } else { 0 }; + DecimalArray::new( + (0..128i128).map(|i| base + i).collect::>(), + DecimalDType::new(38, 2), + Validity::NonNullable, + ) + .into_array() +} + +fn assert_decimal_output( + builder: BtrBlocksCompressorBuilder, + wide: bool, + expected_id: Option, +) -> VortexResult<()> { + let array = decimal_array(wide); + let mut ctx = SESSION.create_execution_ctx(); + let compressed = builder.build().compress(&array, &mut ctx)?; + if let Some(expected_id) = expected_id { + assert!(compressed.is::()); + let serialized = SESSION + .array_serialize(&compressed)? + .ok_or_else(|| vortex_err!("expected serializable decimal byte parts"))?; + assert_eq!(serialized.serialized_id, expected_id); + } else { + assert!(compressed.is::()); + } + assert_arrays_eq!(array, compressed, &mut ctx); + Ok(()) +} + +#[rstest] +#[case::default_edition(None, Some(false))] +#[case::neither(Some(vec![]), None)] +#[case::v1(Some(vec![decimal_byte_parts_v1_id()]), Some(false))] +#[case::v2_only(Some(vec![decimal_byte_parts_v2_id()]), None)] +#[case::both(Some(vec![decimal_byte_parts_v1_id(), decimal_byte_parts_v2_id()]), Some(true))] +fn decimal_mode_follows_permissions( + #[case] ids: Option>, + #[case] v2: Option, + #[values(false, true)] wide: bool, +) -> VortexResult<()> { + let builder = ids + .map(|ids| BtrBlocksCompressorBuilder::new(ids.into_iter().collect())) + .unwrap_or_default(); + let expected = match v2 { + Some(true) if wide => Some(decimal_byte_parts_v2_id()), + Some(_) if !wide => Some(decimal_byte_parts_v1_id()), + _ => None, + }; + assert_decimal_output(builder, wide, expected) +} + +#[rstest] +fn decimal_mode_uses_final_permissions(#[values(false, true)] wide: bool) -> VortexResult<()> { + let builder = + BtrBlocksCompressorBuilder::default().allow_encodings([decimal_byte_parts_v2_id()]); + assert_decimal_output( + builder.clone(), + wide, + Some(if wide { + decimal_byte_parts_v2_id() + } else { + decimal_byte_parts_v1_id() + }), + )?; + + assert_decimal_output( + builder.set_allowed_encodings([decimal_byte_parts_v1_id()]), + wide, + (!wide).then(decimal_byte_parts_v1_id), + ) +} + +#[rstest] +fn permissions_can_reenable_decimal_v2_after_cuda( + #[values(false, true)] replace: bool, +) -> VortexResult<()> { + let builder = BtrBlocksCompressorBuilder::default().only_cuda_compatible(); + let builder = if replace { + builder.set_allowed_encodings([decimal_byte_parts_v1_id(), decimal_byte_parts_v2_id()]) + } else { + builder.allow_encodings([decimal_byte_parts_v2_id()]) + }; + assert_decimal_output(builder, true, Some(decimal_byte_parts_v2_id())) +} + +#[rstest] +fn explicit_decimal_modes_only_upgrade( + #[values(false, true)] allowed_v2: bool, + #[values(false, true)] initial_v2: bool, + #[values(false, true)] wide: bool, +) -> VortexResult<()> { + let scheme: &'static dyn Scheme = if initial_v2 { &DECIMAL_V2 } else { &DECIMAL_V1 }; + let mut allowed = HashSet::from([decimal_byte_parts_v1_id()]); + if allowed_v2 { + allowed.insert(decimal_byte_parts_v2_id()); + } + let builder = BtrBlocksCompressorBuilder::empty() + .allow_encodings(allowed) + .with_new_scheme(scheme); + let mode = (!initial_v2 || allowed_v2).then_some(allowed_v2); + let expected = match mode { + Some(true) if wide => Some(decimal_byte_parts_v2_id()), + Some(_) if !wide => Some(decimal_byte_parts_v1_id()), + _ => None, + }; + assert_decimal_output(builder, wide, expected) +} + +#[test] +fn upgrades_do_not_restore_excluded_decimal() -> VortexResult<()> { + let allowed = HashSet::from([decimal_byte_parts_v1_id(), decimal_byte_parts_v2_id()]); + let builder = + BtrBlocksCompressorBuilder::new(allowed).exclude_schemes([DecimalScheme::default().id()]); + assert_decimal_output(builder, false, None) +} + +#[rstest] +fn cuda_never_uses_decimal_v2( + #[values(false, true)] initial_v2: bool, + #[values(false, true)] register_later: bool, + #[values(false, true)] wide: bool, +) -> VortexResult<()> { + let allowed = HashSet::from([decimal_byte_parts_v1_id(), decimal_byte_parts_v2_id()]); + let scheme: &'static dyn Scheme = if initial_v2 { &DECIMAL_V2 } else { &DECIMAL_V1 }; + let mut builder = BtrBlocksCompressorBuilder::empty().allow_encodings(allowed); + if !register_later { + builder = builder.with_new_scheme(scheme); + } + builder = builder.only_cuda_compatible(); + if register_later { + builder = builder.with_new_scheme(scheme); + } + assert_decimal_output( + builder, + wide, + (!initial_v2 && !wide).then(decimal_byte_parts_v1_id), + ) +} + +#[rstest] +#[case::i128(false, 1)] +#[case::i256(true, 3)] +fn wide_decimal_parts_roundtrip( + #[case] use_i256: bool, + #[case] lower_part_count: usize, + #[values(false, true)] negative: bool, + #[values(false, true)] nullable: bool, + #[values(false, true)] compress_children: bool, +) -> VortexResult<()> { + let validity = if nullable { + Validity::from_iter((0..2048).map(|i| i % 7 != 0)) + } else { + Validity::NonNullable + }; + let array = if use_i256 { + let values = (1..=2048u32) + .map(|i| { + let msp = if negative { + -i128::from(i) + } else { + i128::from(i) + }; + i256::from_parts( + (u128::from(i) << 64) | u128::from(i * 131 + 17), + (msp << 64) | i128::from(i * 3 + 1), + ) + }) + .collect::>(); + DecimalArray::new(values, DecimalDType::new(76, 2), validity) + } else { + let values = (1..=2048i128) + .map(|i| { + let msp = if negative { -i } else { i }; + (msp << 70) + i * 131 + 17 + }) + .collect::>(); + DecimalArray::new(values, DecimalDType::new(38, 2), validity) + } + .into_array(); + let mut allowed = HashSet::from([decimal_byte_parts_v1_id(), decimal_byte_parts_v2_id()]); + if compress_children { + allowed.extend(FoRScheme.produced_encodings()); + allowed.extend(BitPackingScheme.produced_encodings()); + } + let compressor = BtrBlocksCompressorBuilder::empty() + .allow_encodings(allowed) + .with_new_scheme(&DECIMAL_V2) + .with_new_scheme(&FoRScheme) + .with_new_scheme(&BitPackingScheme) + .build(); + let mut ctx = SESSION.create_execution_ctx(); + let compressed = compressor.compress(&array, &mut ctx)?; + // Every part varies, so without child compression splitting alone cannot save space. + assert_eq!(compressed.is::(), compress_children); + if compress_children { + let parts = compressed + .as_opt::() + .ok_or_else(|| vortex_err!("expected decimal byte parts"))?; + assert_eq!(parts.lower_parts().len(), lower_part_count); + assert!(!parts.msp().is_canonical()); + assert!(parts.lower_parts().iter().all(|part| !part.is_canonical())); + } + + let array_ctx = ArrayContext::empty(); + let mut bytes = ByteBufferMut::empty(); + for buffer in compressed.serialize(&array_ctx, &SESSION, &SerializeOptions::default())? { + bytes.extend_from_slice(buffer.as_ref()); + } + let decoded = SerializedArray::try_from(bytes.freeze())?.decode( + array.dtype(), + array.len(), + &ReadContext::new(array_ctx.to_ids()), + &SESSION, + )?; + assert_arrays_eq!(array, decoded, &mut ctx); + Ok(()) +} diff --git a/vortex-btrblocks/tests/golden.rs b/vortex-btrblocks/tests/golden.rs index fc636252c27..c0539a7407f 100644 --- a/vortex-btrblocks/tests/golden.rs +++ b/vortex-btrblocks/tests/golden.rs @@ -414,35 +414,18 @@ fn edition_session(editions: &[EditionId]) -> VortexResult { Ok(session) } -fn compressor_for_session( - session: &VortexSession, - builder: BtrBlocksCompressorBuilder, -) -> BtrBlocksCompressor { +fn compressor_builder_for_session(session: &VortexSession) -> BtrBlocksCompressorBuilder { let allowed = session .enabled_component_ids(ComponentKind::Array) .into_iter() .collect(); - without_onpair(builder) - .retain_allowed_encodings(&allowed) - .build() -} - -/// Like [`compressor_for_session`] but keeps OnPair in the scheme pool. -fn compressor_with_onpair( - session: &VortexSession, - builder: BtrBlocksCompressorBuilder, -) -> BtrBlocksCompressor { - let allowed = session - .enabled_component_ids(ComponentKind::Array) - .into_iter() - .collect(); - builder.retain_allowed_encodings(&allowed).build() + BtrBlocksCompressorBuilder::new(allowed) } #[test] fn golden_regular() -> VortexResult<()> { let session = edition_session(&[CORE_2026_08_3])?; - let compressor = compressor_for_session(&session, BtrBlocksCompressorBuilder::default()); + let compressor = without_onpair(compressor_builder_for_session(&session)).build(); golden_corpus_snapshots("regular", &compressor) } @@ -450,7 +433,7 @@ fn golden_regular() -> VortexResult<()> { #[test] fn golden_onpair() -> VortexResult<()> { let session = edition_session(&[CORE_2026_08_3])?; - let compressor = compressor_with_onpair(&session, BtrBlocksCompressorBuilder::default()); + let compressor = compressor_builder_for_session(&session).build(); golden_snapshots( "onpair", &compressor, @@ -464,9 +447,7 @@ fn golden_compact() -> VortexResult<()> { let session = edition_session(&[CORE_2026_08_3])?; vortex_zstd::initialize(&session); session.enable_edition(vortex_zstd::editions::ZSTD_2026_02)?; - let compressor = compressor_for_session( - &session, - BtrBlocksCompressorBuilder::default().with_compact(), - ); + let compressor = + without_onpair(compressor_builder_for_session(&session).with_compact()).build(); golden_corpus_snapshots("compact", &compressor) } diff --git a/vortex-compressor/src/scheme/mod.rs b/vortex-compressor/src/scheme/mod.rs index 0ba1c90202a..6af68761f2e 100644 --- a/vortex-compressor/src/scheme/mod.rs +++ b/vortex-compressor/src/scheme/mod.rs @@ -28,11 +28,15 @@ use vortex_array::ArrayRef; use vortex_array::Canonical; use vortex_array::ExecutionCtx; use vortex_error::VortexResult; +use vortex_utils::aliases::hash_set::HashSet; use crate::CascadingCompressor; use crate::stats::ArrayAndStats; use crate::stats::GenerateStatsOptions; +/// The serialized IDs permitted for compression schemes. +pub type AllowedSerializedIds = HashSet; + /// Unique identifier for a compression scheme. /// /// The only way to obtain a [`SchemeId`] is through [`SchemeExt::id()`], which is auto-implemented @@ -135,6 +139,15 @@ pub trait Scheme: Debug + Send + Sync { /// formats declares the wire IDs the scheme writes, which may differ from its in-memory ID. fn produced_encodings(&self) -> Vec; + /// Returns a newer configuration supported by the permitted serialized IDs, if available. + /// + /// `None` keeps the original scheme. An upgrade must preserve the [`SchemeId`] and must not + /// downgrade the registered configuration. The caller checks the selected scheme's + /// [`produced_encodings`](Self::produced_encodings) before using it. + fn try_upgrade(&self, _allowed_serialized_ids: &AllowedSerializedIds) -> Option<&dyn Scheme> { + None + } + /// Returns the stats generation options this scheme requires. The compressor merges all /// eligible schemes' options before generating stats so that a single stats pass satisfies /// every scheme. diff --git a/vortex-cuda/src/layout.rs b/vortex-cuda/src/layout.rs index 7cc7d0322da..732527530be 100644 --- a/vortex-cuda/src/layout.rs +++ b/vortex-cuda/src/layout.rs @@ -557,9 +557,7 @@ pub fn cuda_write_strategy(session: &VortexSession, block_rows: usize) -> Arc impl Iterator + '_ { + EDITION_DECLARATIONS + .iter() + .filter(move |declaration| declaration.edition.id.is_at_or_before(edition)) + .flat_map(|declaration| declaration.added) + .filter(|member| member.kind == ComponentKind::Array) + .map(|member| member.component.component_id()) +} diff --git a/vortex-edition/src/lib.rs b/vortex-edition/src/lib.rs index ba2382ec02b..e882f4263eb 100644 --- a/vortex-edition/src/lib.rs +++ b/vortex-edition/src/lib.rs @@ -42,8 +42,10 @@ use std::fmt::Debug; use std::fmt::Display; use std::fmt::Formatter; +pub use declarations::DEFAULT_CORE_EDITION; pub use declarations::EDITION_DECLARATIONS; pub use declarations::EDITION_FAMILIES; +pub use declarations::array_ids_for_edition; pub use session::EditionSession; pub use session::EditionSessionExt; pub use session::EnabledEditions; diff --git a/vortex-edition/src/tests.rs b/vortex-edition/src/tests.rs index 63cf003d11f..58bd30472f0 100644 --- a/vortex-edition/src/tests.rs +++ b/vortex-edition/src/tests.rs @@ -4,6 +4,7 @@ use vortex_error::VortexResult; use vortex_error::vortex_err; use vortex_session::VortexSession; +use vortex_session::registry::Id; use crate::ComponentKind; use crate::Edition; @@ -15,6 +16,29 @@ use crate::EditionMember; use crate::EditionSession; use crate::EditionSessionExt; use crate::EnabledEditions; +use crate::array_ids_for_edition; +use crate::declarations::core::CORE_2025_05_0; +use crate::declarations::core::CORE_2025_06_0; + +#[test] +fn static_array_ids_inherit_only_within_the_edition_family() { + let first_ids: Vec<_> = array_ids_for_edition(&CORE_2025_05_0).collect(); + let second_ids: Vec<_> = array_ids_for_edition(&CORE_2025_06_0).collect(); + let first: Vec<_> = first_ids.iter().map(Id::as_str).collect(); + let second: Vec<_> = second_ids.iter().map(Id::as_str).collect(); + + assert!(first.contains(&"vortex.primitive")); + assert!(!first.contains(&"vortex.sequence")); + assert!(first.iter().all(|id| second.contains(id))); + assert!(second.contains(&"vortex.sequence")); + assert!(!second.contains(&"vortex.flat")); + assert!(!second.contains(&"vortex.date")); + assert!( + array_ids_for_edition(&EditionId::new("other", 2025, 6, 0)) + .next() + .is_none() + ); +} static TEST_FAMILY: EditionFamily = EditionFamily { name: "test", diff --git a/vortex-file/src/tests.rs b/vortex-file/src/tests.rs index 640de874d2b..0175bbd0d8d 100644 --- a/vortex-file/src/tests.rs +++ b/vortex-file/src/tests.rs @@ -38,6 +38,7 @@ use vortex_array::dtype::Nullability; use vortex_array::dtype::PType; use vortex_array::dtype::PType::I32; use vortex_array::dtype::StructFields; +use vortex_array::dtype::i256; use vortex_array::expr::BoundExpression; use vortex_array::expr::Expression; use vortex_array::expr::and; @@ -72,7 +73,12 @@ use vortex_buffer::Buffer; use vortex_buffer::ByteBuffer; use vortex_buffer::ByteBufferMut; use vortex_buffer::buffer; +use vortex_decimal_byte_parts::DecimalByteParts; +use vortex_decimal_byte_parts::DecimalBytePartsArraySlotsExt; +use vortex_edition::DEFAULT_CORE_EDITION; +use vortex_edition::EDITION_DECLARATIONS; use vortex_edition::EditionSession; +use vortex_edition::EditionSessionExt; use vortex_error::VortexExpect; use vortex_error::VortexResult; use vortex_io::session::RuntimeSession; @@ -99,6 +105,7 @@ use crate::V1_FOOTER_FBS_SIZE; use crate::VERSION; use crate::VortexFile; use crate::WriteOptionsSessionExt; +use crate::WriteStrategyBuilder; use crate::flatbuffers::footer as fb; use crate::footer::SegmentSpec; static SESSION: LazyLock = LazyLock::new(|| { @@ -173,6 +180,87 @@ async fn test_read_simple() { assert_eq!(row_count, 8); } +#[rstest] +#[case::default_writer(false, false)] +#[case::custom_layout(true, false)] +#[case::explicit_compressor(true, true)] +#[tokio::test] +#[cfg_attr(miri, ignore)] +async fn decimal_writer_uses_default_or_session_permissions( + #[values(false, true)] use_i256: bool, + #[case] custom_strategy: bool, + #[case] explicit_compressor: bool, +) -> VortexResult<()> { + let session = array_session() + .with::() + .with::() + .with::(); + crate::register_default_encodings(&session); + for declaration in EDITION_DECLARATIONS { + session.register_edition(declaration)?; + } + session.enable_edition(DEFAULT_CORE_EDITION)?; + + let array = if use_i256 { + DecimalArray::new( + (0..1024u128) + .map(|i| i256::from_parts(i * 17, 1i128 << 70)) + .collect::>(), + DecimalDType::new(76, 2), + Validity::NonNullable, + ) + } else { + DecimalArray::new( + (0..1024i128) + .map(|i| (1i128 << 70) + i * 17) + .collect::>(), + DecimalDType::new(38, 2), + Validity::NonNullable, + ) + } + .into_array(); + let strategy = WriteStrategyBuilder::default() + .with_row_block_size(256) + .with_data_block_target_bytes(None); + let strategy = if explicit_compressor { + strategy.with_btrblocks_builder(BtrBlocksCompressorBuilder::default()) + } else { + strategy + } + .build(); + + for disable_editions in [false, true] { + let mut options = session.write_options(); + if custom_strategy { + options = options.with_strategy(Arc::clone(&strategy)); + } + let options = if disable_editions { + options.disable_editions() + } else { + options + }; + let mut buffer = ByteBufferMut::empty(); + options + .write(&mut buffer, array.clone().to_array_stream()) + .await?; + let actual = session + .open_options() + .open_buffer(buffer)? + .scan()? + .into_array_stream()? + .read_all() + .await?; + let uses_v2 = actual.depth_first_traversal().any(|array| { + array + .as_opt::() + .is_some_and(|parts| !parts.lower_parts().is_empty()) + }); + assert_eq!(uses_v2, disable_editions && !custom_strategy); + assert_arrays_eq!(array, actual, &mut session.create_execution_ctx()); + } + Ok(()) +} + #[tokio::test] #[cfg_attr(miri, ignore)] async fn test_round_trip_many_types() { @@ -1874,9 +1962,7 @@ async fn write_read_roundtrip_with_layout( array: ArrayRef, use_list_layout: bool, ) -> VortexResult { - let strategy = crate::strategy::WriteStrategyBuilder::default() - .with_list_layout() - .build(); + let strategy = WriteStrategyBuilder::default().with_list_layout().build(); let mut buf = ByteBufferMut::empty(); if use_list_layout { SESSION @@ -2260,7 +2346,7 @@ async fn timestamp_unit_mismatch_errors_with_constant_children() .into_array(); let temporal = TemporalArray::new_timestamp(ts_array, TimeUnit::Milliseconds, None); - let strategy = crate::strategy::WriteStrategyBuilder::default() + let strategy = WriteStrategyBuilder::default() .with_compressor(compressor) .build(); @@ -2578,7 +2664,7 @@ async fn dict_probe_honours_configured_compressor() -> VortexResult<()> { let mut buf = ByteBufferMut::empty(); let summary = SESSION .write_options() - .with_strategy(crate::strategy::WriteStrategyBuilder::default().build()) + .with_strategy(WriteStrategyBuilder::default().build()) .write(&mut buf, strings.clone().to_array_stream()) .await?; assert!( @@ -2592,7 +2678,7 @@ async fn dict_probe_honours_configured_compressor() -> VortexResult<()> { let summary = SESSION .write_options() .with_strategy( - crate::strategy::WriteStrategyBuilder::default() + WriteStrategyBuilder::default() .with_btrblocks_builder(no_string_dict) .build(), ) @@ -2622,7 +2708,7 @@ async fn probe_compressor_override_is_independent() -> VortexResult<()> { let summary = SESSION .write_options() .with_strategy( - crate::strategy::WriteStrategyBuilder::default() + WriteStrategyBuilder::default() .with_probe_compressor(probe_without_dict) .build(), ) diff --git a/vortex-file/src/writer.rs b/vortex-file/src/writer.rs index 874d08be306..5d6504dd7c3 100644 --- a/vortex-file/src/writer.rs +++ b/vortex-file/src/writer.rs @@ -252,10 +252,7 @@ impl VortexWriteOptions { let strategy = match self.strategy { Some(strategy) => strategy, None => WriteStrategyBuilder::default() - .with_btrblocks_builder( - BtrBlocksCompressorBuilder::default() - .retain_allowed_encodings(&allowed_serialized_ids), - ) + .with_btrblocks_builder(BtrBlocksCompressorBuilder::new(allowed_serialized_ids)) .build(), }; let dtype = stream.dtype().clone(); diff --git a/vortex-python/src/io.rs b/vortex-python/src/io.rs index 7288fc5d1e3..97c3f4e898d 100644 --- a/vortex-python/src/io.rs +++ b/vortex-python/src/io.rs @@ -384,12 +384,11 @@ impl PyVortexWriteOptions { .enabled_component_ids(ComponentKind::Array) .into_iter() .collect(); - let mut compressor = BtrBlocksCompressorBuilder::default(); + let mut compressor = BtrBlocksCompressorBuilder::new(allowed_encodings); if self.use_compact_encodings { compressor = compressor.with_compact(); } - let strategy = WriteStrategyBuilder::default() - .with_btrblocks_builder(compressor.retain_allowed_encodings(&allowed_encodings)); + let strategy = WriteStrategyBuilder::default().with_btrblocks_builder(compressor); let strategy = strategy.build(); current_runtime().block_on(async move { match resolve_store(path, store.map(|x| x.into_inner()))? { diff --git a/vortex-tui/src/convert.rs b/vortex-tui/src/convert.rs index ab316982b27..10eb3bbf2be 100644 --- a/vortex-tui/src/convert.rs +++ b/vortex-tui/src/convert.rs @@ -102,12 +102,11 @@ pub async fn exec_convert(session: &VortexSession, flags: ConvertArgs) -> anyhow .enabled_component_ids(ComponentKind::Array) .into_iter() .collect(); - let mut compressor = BtrBlocksCompressorBuilder::default(); + let mut compressor = BtrBlocksCompressorBuilder::new(allowed_encodings); if matches!(flags.strategy, Strategy::Compact) { compressor = compressor.with_compact(); } - let strategy = WriteStrategyBuilder::default() - .with_btrblocks_builder(compressor.retain_allowed_encodings(&allowed_encodings)); + let strategy = WriteStrategyBuilder::default().with_btrblocks_builder(compressor); let mut file = File::create(output_path).await?; session diff --git a/vortex/src/editions/mod.rs b/vortex/src/editions/mod.rs index 464cce2197b..6cd38268ed5 100644 --- a/vortex/src/editions/mod.rs +++ b/vortex/src/editions/mod.rs @@ -22,6 +22,7 @@ mod tests; pub use vortex_edition::ComponentKind; +pub use vortex_edition::DEFAULT_CORE_EDITION; pub use vortex_edition::EDITION_DECLARATIONS; pub use vortex_edition::EDITION_FAMILIES; pub use vortex_edition::Edition; @@ -47,9 +48,6 @@ use vortex_error::VortexExpect; use vortex_error::vortex_err; use vortex_session::VortexSession; -/// The `core` edition enabled for writing by the default Vortex session. -pub const DEFAULT_CORE_EDITION: EditionId = CORE_2026_08_3; - /// The newest `preview` edition. The default Vortex session registers it but does not enable it. pub const DEFAULT_PREVIEW_EDITION: EditionId = PREVIEW_2026_08_0;