Stacking feat/libver-v18 under feat/huge-chunks sent every chunked dataset of a file with a 1.8 low bound to a version-1 B-tree, whose key holds a 32-bit chunk size. As in libhdf5, such a chunk now takes layout version 5 whatever the low bound, and a high bound below 2.0 is FormatError::LibverBound. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
2711 lines
97 KiB
Rust
2711 lines
97 KiB
Rust
//! Chunked dataset writing: chunk splitting, compression, index building.
|
||
|
||
#[cfg(not(feature = "std"))]
|
||
extern crate alloc;
|
||
|
||
use crate::addr::saturating_usize;
|
||
#[cfg(not(feature = "std"))]
|
||
use alloc::{format, vec, vec::Vec};
|
||
|
||
use crate::btree_v1_write;
|
||
use crate::btree_v2_write::{BTreeV2Params, build_btree_v2};
|
||
use crate::checksum::jenkins_lookup3;
|
||
use crate::chunk_cache::{CACHE_LINE_SIZE, align_to_cache_line};
|
||
use crate::chunk_grid::ChunkGrid;
|
||
use crate::ea_writer;
|
||
use crate::error::FormatError;
|
||
use crate::filter_pipeline::{
|
||
FILTER_BITSHUFFLE, FILTER_BLOSC, FILTER_BZIP2, FILTER_DEFLATE, FILTER_FLETCHER32, FILTER_LZ4,
|
||
FILTER_LZF, FILTER_PCODEC, FILTER_PCODEC_NAME, FILTER_SHUFFLE, FILTER_ZSTD, FilterDescription,
|
||
FilterPipeline,
|
||
};
|
||
use crate::filters::compress_chunk_masked;
|
||
use crate::libver::LibVer;
|
||
/// Round a file offset up to the next cache-line boundary.
|
||
///
|
||
/// This ensures chunk data starts at an address that is a multiple of the
|
||
/// architecture's cache line size, enabling aligned loads in SIMD paths.
|
||
#[inline]
|
||
pub fn align_chunk_offset(offset: u64) -> u64 {
|
||
let align = CACHE_LINE_SIZE as u64;
|
||
(offset + align - 1) & !(align - 1)
|
||
}
|
||
|
||
/// Options for chunked dataset creation.
|
||
#[derive(Debug, Clone, Default)]
|
||
pub struct ChunkOptions {
|
||
/// Chunk dimensions (one per dataset dimension).
|
||
pub chunk_dims: Option<Vec<u64>>,
|
||
/// Deflate compression level (0-9), None = no deflate.
|
||
pub deflate_level: Option<u32>,
|
||
/// Whether to apply shuffle filter before compression.
|
||
/// If `false` AND compression is enabled AND `no_shuffle` is `false`,
|
||
/// shuffle is auto-applied (matches h5py default behavior).
|
||
pub shuffle: bool,
|
||
/// Disable the automatic shuffle pre-filter. Set via `without_shuffle()`.
|
||
pub no_shuffle: bool,
|
||
/// Whether to apply fletcher32 checksum.
|
||
pub fletcher32: bool,
|
||
/// Whether to use LZ4 compression (filter ID 32004).
|
||
pub lz4: bool,
|
||
/// Zstandard compression level (1-22), None = no zstd. Filter ID 32015.
|
||
pub zstd_level: Option<u32>,
|
||
/// Pcodec lossless numerical compression. Private, unregistered filter
|
||
/// ID [`FILTER_PCODEC`] (480): only clawhdf5 can read it.
|
||
pub pcodec: bool,
|
||
/// A plugin compression filter (LZF, ...). Takes priority over the
|
||
/// codecs above. Each needs its cargo feature to be written.
|
||
pub plugin: Option<PluginFilter>,
|
||
}
|
||
|
||
/// A compression filter from the common HDF5 plugin set, written in the
|
||
/// format the libhdf5 plugin (h5py / hdf5plugin) reads.
|
||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||
#[non_exhaustive]
|
||
pub enum PluginFilter {
|
||
/// LZF (filter 32000), h5py's built-in `compression="lzf"`. Needs the
|
||
/// `lzf` feature.
|
||
Lzf,
|
||
/// Bitshuffle (filter 32008): a bit transpose of each block of
|
||
/// `block_size` elements (0 = bitshuffle's default, else a multiple of
|
||
/// 8), optionally compressed. Needs the `bitshuffle` feature.
|
||
Bitshuffle {
|
||
/// Block size in elements; 0 for the default.
|
||
block_size: u32,
|
||
/// Compression after the transpose.
|
||
compression: BitshuffleCompression,
|
||
},
|
||
/// bzip2 (filter 307) at block size `level` (1-9). Needs the `bzip2`
|
||
/// feature.
|
||
Bzip2 {
|
||
/// Block size 1-9 (9 = hdf5plugin's default).
|
||
level: u32,
|
||
},
|
||
/// Blosc 1 (filter 32001): `codec` at `level` (0-9; 0 stores), after
|
||
/// `shuffle`. Needs the `blosc` feature.
|
||
Blosc {
|
||
/// The codec inside the Blosc frame.
|
||
codec: BloscCodec,
|
||
/// Compression level 0-9 (0 stores the data uncompressed).
|
||
level: u32,
|
||
/// The shuffle Blosc applies first.
|
||
shuffle: BloscShuffle,
|
||
},
|
||
}
|
||
|
||
/// The codec inside a Blosc frame that clawhdf5 can write. (It reads
|
||
/// BloscLZ too, but cannot write it.)
|
||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||
pub enum BloscCodec {
|
||
/// LZ4.
|
||
Lz4,
|
||
/// Snappy.
|
||
Snappy,
|
||
/// Zlib, at the Blosc level.
|
||
Zlib,
|
||
/// Zstandard (clawhdf5's pure-Rust encoder has one level, about zstd 1).
|
||
Zstd,
|
||
}
|
||
|
||
/// The shuffle Blosc applies before compressing.
|
||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||
pub enum BloscShuffle {
|
||
/// None.
|
||
None,
|
||
/// Byte shuffle (Blosc's default).
|
||
Byte,
|
||
/// Bit shuffle.
|
||
Bit,
|
||
}
|
||
|
||
/// What bitshuffle compresses its blocks with.
|
||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||
pub enum BitshuffleCompression {
|
||
/// Transpose only.
|
||
None,
|
||
/// LZ4 (bitshuffle's `cname="lz4"`, the common choice).
|
||
Lz4,
|
||
/// Zstandard. clawhdf5's pure-Rust encoder has a single level (about
|
||
/// zstd's level 1); `level` is recorded in the file for other writers.
|
||
Zstd {
|
||
/// Level recorded in `cd_values[5]`.
|
||
level: u32,
|
||
},
|
||
}
|
||
|
||
impl PluginFilter {
|
||
/// Whether the filter reorders bytes itself, so the automatic shuffle
|
||
/// pre-filter would only get in its way.
|
||
fn shuffles_itself(&self) -> bool {
|
||
match self {
|
||
PluginFilter::Lzf => false,
|
||
PluginFilter::Bitshuffle { .. } => true,
|
||
PluginFilter::Bzip2 { .. } => false,
|
||
PluginFilter::Blosc { .. } => true,
|
||
}
|
||
}
|
||
|
||
/// The pipeline entry for this filter. `chunk_bytes` is one chunk's
|
||
/// uncompressed size (0 if unknown).
|
||
fn description(&self, element_size: u32, chunk_bytes: u32) -> FilterDescription {
|
||
match self {
|
||
// h5py's lzf_set_local: filter version, liblzf version, chunk
|
||
// size in bytes. Optional, as h5py flags it: a chunk the filter
|
||
// cannot shrink may then be stored unfiltered.
|
||
PluginFilter::Lzf => FilterDescription {
|
||
filter_id: FILTER_LZF,
|
||
name: Some("lzf".into()),
|
||
flags: 1,
|
||
client_data: vec![4, 0x0105, chunk_bytes],
|
||
},
|
||
// bshuf_h5_set_local: version 0.4, element size, block size,
|
||
// compression (0 none, 2 LZ4, 3 Zstandard), Zstandard level.
|
||
// hdf5-blosc's blosc_set_local: filter revision 2, Blosc format
|
||
// 2, type size, chunk size, then level, shuffle, compressor.
|
||
PluginFilter::Blosc {
|
||
codec,
|
||
level,
|
||
shuffle,
|
||
} => FilterDescription {
|
||
filter_id: FILTER_BLOSC,
|
||
name: Some("blosc".into()),
|
||
flags: 1,
|
||
client_data: vec![
|
||
2,
|
||
2,
|
||
element_size,
|
||
chunk_bytes,
|
||
(*level).min(9),
|
||
match shuffle {
|
||
BloscShuffle::None => 0,
|
||
BloscShuffle::Byte => 1,
|
||
BloscShuffle::Bit => 2,
|
||
},
|
||
match codec {
|
||
BloscCodec::Lz4 => 1,
|
||
BloscCodec::Snappy => 3,
|
||
BloscCodec::Zlib => 4,
|
||
BloscCodec::Zstd => 5,
|
||
},
|
||
],
|
||
},
|
||
PluginFilter::Bzip2 { level } => FilterDescription {
|
||
filter_id: FILTER_BZIP2,
|
||
name: Some("bzip2".into()),
|
||
flags: 1,
|
||
client_data: vec![(*level).clamp(1, 9)],
|
||
},
|
||
PluginFilter::Bitshuffle {
|
||
block_size,
|
||
compression,
|
||
} => {
|
||
let mut cd = vec![0, 4, element_size, *block_size];
|
||
match compression {
|
||
BitshuffleCompression::None => cd.push(0),
|
||
BitshuffleCompression::Lz4 => cd.push(2),
|
||
BitshuffleCompression::Zstd { level } => cd.extend([3, *level]),
|
||
}
|
||
FilterDescription {
|
||
filter_id: FILTER_BITSHUFFLE,
|
||
name: Some("bitshuffle; see https://github.com/kiyo-masui/bitshuffle".into()),
|
||
flags: 1,
|
||
client_data: cd,
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Largest chunk the automatic choice produces, in bytes.
|
||
const AUTO_CHUNK_TARGET_BYTES: u64 = 1 << 20;
|
||
|
||
/// Extent assumed for a dimension that is currently empty (an unlimited
|
||
/// dimension not yet written to) — the same stand-in h5py uses.
|
||
const AUTO_CHUNK_EMPTY_DIM: u64 = 1024;
|
||
|
||
/// Choose chunk dimensions for a dataset nobody specified them for.
|
||
///
|
||
/// Asking for compression (or any filter) without chunk dimensions used to
|
||
/// make the whole dataset one chunk. That defeats the point of chunking: any
|
||
/// read — even a single row — must decompress everything, and a large dataset
|
||
/// cannot be decompressed in parallel. Datasets up to the target size stay a
|
||
/// single chunk, exactly as before; larger ones are split by halving the
|
||
/// dimensions in turn (so chunks keep roughly the dataset's proportions, the
|
||
/// approach h5py takes) until a chunk fits the target.
|
||
pub fn auto_chunk_dims(shape: &[u64], elem_size: usize) -> Vec<u64> {
|
||
let mut dims: Vec<u64> = shape
|
||
.iter()
|
||
.map(|&d| if d == 0 { AUTO_CHUNK_EMPTY_DIM } else { d })
|
||
.collect();
|
||
let elem = elem_size.max(1) as u64;
|
||
let bytes = |dims: &[u64]| dims.iter().fold(elem, |acc, &d| acc.saturating_mul(d));
|
||
let mut axis = 0;
|
||
while bytes(&dims) > AUTO_CHUNK_TARGET_BYTES && dims.iter().any(|&d| d > 1) {
|
||
let i = axis % dims.len();
|
||
dims[i] = dims[i].div_ceil(2);
|
||
axis += 1;
|
||
}
|
||
dims
|
||
}
|
||
|
||
impl ChunkOptions {
|
||
/// Whether any chunking option is enabled.
|
||
pub fn is_chunked(&self) -> bool {
|
||
self.chunk_dims.is_some()
|
||
|| self.deflate_level.is_some()
|
||
|| self.shuffle
|
||
|| self.fletcher32
|
||
|| self.lz4
|
||
|| self.zstd_level.is_some()
|
||
|| self.pcodec
|
||
|| self.plugin.is_some()
|
||
}
|
||
|
||
/// Build a FilterPipeline from the options.
|
||
pub fn build_pipeline(&self, element_size: u32) -> Option<FilterPipeline> {
|
||
self.build_pipeline_for_chunk(element_size, 0)
|
||
}
|
||
|
||
/// Build a FilterPipeline for chunks of `chunk_bytes` uncompressed bytes
|
||
/// (0 if unknown). Some plugin filters record the chunk size in their
|
||
/// client data.
|
||
pub fn build_pipeline_for_chunk(
|
||
&self,
|
||
element_size: u32,
|
||
chunk_bytes: u32,
|
||
) -> Option<FilterPipeline> {
|
||
let mut filters = Vec::new();
|
||
|
||
let plugin_shuffles = self
|
||
.plugin
|
||
.as_ref()
|
||
.is_some_and(PluginFilter::shuffles_itself);
|
||
let has_compression = self.deflate_level.is_some()
|
||
|| self.zstd_level.is_some()
|
||
|| self.lz4
|
||
|| self.pcodec
|
||
|| (self.plugin.is_some() && !plugin_shuffles);
|
||
|
||
// Shuffle before compression. Applied if explicitly requested OR if compression
|
||
// is active and the caller hasn't disabled it — matches h5py default behavior
|
||
// and implements TDT byte-grouping (arXiv:2506.18062) for free.
|
||
if self.shuffle || (has_compression && !self.no_shuffle) {
|
||
filters.push(FilterDescription {
|
||
filter_id: FILTER_SHUFFLE,
|
||
name: None,
|
||
flags: 0,
|
||
client_data: vec![element_size],
|
||
});
|
||
}
|
||
|
||
// Compression filters (mutually exclusive, priority: plugin > pcodec >
|
||
// zstd > lz4 > deflate)
|
||
if let Some(plugin) = &self.plugin {
|
||
filters.push(plugin.description(element_size, chunk_bytes));
|
||
} else if self.pcodec {
|
||
filters.push(FilterDescription {
|
||
filter_id: FILTER_PCODEC,
|
||
name: Some(FILTER_PCODEC_NAME.into()),
|
||
flags: 0,
|
||
client_data: vec![element_size],
|
||
});
|
||
} else if let Some(level) = self.zstd_level {
|
||
filters.push(FilterDescription {
|
||
filter_id: FILTER_ZSTD,
|
||
name: Some("zstd".into()),
|
||
flags: 0,
|
||
client_data: vec![level],
|
||
});
|
||
} else if self.lz4 {
|
||
filters.push(FilterDescription {
|
||
filter_id: FILTER_LZ4,
|
||
name: Some("lz4".into()),
|
||
flags: 0,
|
||
client_data: vec![],
|
||
});
|
||
} else if let Some(level) = self.deflate_level {
|
||
filters.push(FilterDescription {
|
||
filter_id: FILTER_DEFLATE,
|
||
name: None,
|
||
flags: 0,
|
||
client_data: vec![level],
|
||
});
|
||
}
|
||
|
||
if self.fletcher32 {
|
||
filters.push(FilterDescription {
|
||
filter_id: FILTER_FLETCHER32,
|
||
name: None,
|
||
flags: 0,
|
||
client_data: vec![],
|
||
});
|
||
}
|
||
|
||
// Note: h5py sets flags=0x0001 (optional) on filters, but this is not required
|
||
// for read compatibility.
|
||
|
||
if filters.is_empty() {
|
||
None
|
||
} else {
|
||
Some(FilterPipeline {
|
||
version: 2,
|
||
filters,
|
||
})
|
||
}
|
||
}
|
||
|
||
/// Determine chunk dimensions, using user-specified or auto-computing.
|
||
pub fn resolve_chunk_dims(&self, shape: &[u64]) -> Vec<u64> {
|
||
// Without the element size, assume 8 bytes (the widest common scalar);
|
||
// the writer uses `resolve_chunk_dims_for`.
|
||
self.resolve_chunk_dims_for(shape, 8)
|
||
}
|
||
|
||
/// Chunk dimensions for a dataset of `shape` whose elements are `elem_size`
|
||
/// bytes: the caller's if given, otherwise chosen automatically.
|
||
pub fn resolve_chunk_dims_for(&self, shape: &[u64], elem_size: usize) -> Vec<u64> {
|
||
match self.chunk_dims {
|
||
Some(ref dims) => dims.clone(),
|
||
None => auto_chunk_dims(shape, elem_size),
|
||
}
|
||
}
|
||
}
|
||
|
||
/// A chunk that has been written to the file buffer.
|
||
#[derive(Debug, Clone)]
|
||
pub struct WrittenChunk {
|
||
/// Address within the file where chunk data starts.
|
||
pub address: u64,
|
||
/// Size of the (possibly compressed) chunk data in bytes.
|
||
pub compressed_size: u64,
|
||
/// Original uncompressed size in bytes.
|
||
pub raw_size: u64,
|
||
/// Filter mask (0 = all filters applied).
|
||
pub filter_mask: u32,
|
||
}
|
||
|
||
/// Result of building a chunked dataset.
|
||
pub struct ChunkedDataResult {
|
||
/// Raw bytes containing all chunk data + index structures.
|
||
pub data_bytes: Vec<u8>,
|
||
/// The DataLayout v4 message bytes.
|
||
pub layout_message: Vec<u8>,
|
||
/// The FilterPipeline message bytes, if any.
|
||
pub pipeline_message: Option<Vec<u8>>,
|
||
}
|
||
|
||
/// Split raw data into chunk-sized pieces based on shape and chunk dimensions.
|
||
/// Returns a Vec of (chunk_offset_per_dim, chunk_raw_bytes).
|
||
pub fn split_into_chunks(
|
||
raw_data: &[u8],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
) -> Vec<(Vec<u64>, Vec<u8>)> {
|
||
if shape.is_empty() {
|
||
return vec![(vec![], raw_data.to_vec())];
|
||
}
|
||
(0..chunk_count(shape, chunk_dims))
|
||
.map(|i| extract_chunk(raw_data, shape, chunk_dims, element_size, i))
|
||
.collect()
|
||
}
|
||
|
||
/// Number of chunks of the current extent `shape`.
|
||
fn chunk_count(shape: &[u64], chunk_dims: &[u64]) -> u64 {
|
||
shape
|
||
.iter()
|
||
.zip(chunk_dims)
|
||
.map(|(&s, &c)| s.div_ceil(c))
|
||
.product()
|
||
}
|
||
|
||
/// The `linear_idx`-th chunk (row-major over the chunks of the current
|
||
/// extent) of the row-major dataset `raw_data`: its offset in dataset space
|
||
/// and its bytes, a whole chunk with the part past the dataset's edge zero.
|
||
/// Copied one row (a run along the last dimension) at a time.
|
||
fn extract_chunk(
|
||
raw_data: &[u8],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
linear_idx: u64,
|
||
) -> (Vec<u64>, Vec<u8>) {
|
||
let rank = shape.len();
|
||
let mut offsets = vec![0u64; rank];
|
||
let mut remaining = linear_idx;
|
||
for d in (0..rank).rev() {
|
||
let n = shape[d].div_ceil(chunk_dims[d]);
|
||
offsets[d] = (remaining % n) * chunk_dims[d];
|
||
remaining /= n;
|
||
}
|
||
let chunk_elements: usize = chunk_dims.iter().map(|&d| saturating_usize(d)).product();
|
||
let mut chunk = vec![0u8; chunk_elements.saturating_mul(element_size)];
|
||
|
||
// Elements of the chunk inside the dataset, per dimension.
|
||
let valid: Vec<usize> = (0..rank)
|
||
.map(|d| saturating_usize(shape[d].saturating_sub(offsets[d]).min(chunk_dims[d])))
|
||
.collect();
|
||
if valid.contains(&0) {
|
||
return (offsets, chunk);
|
||
}
|
||
let mut ds_strides = vec![1usize; rank];
|
||
let mut chunk_strides = vec![1usize; rank];
|
||
for d in (0..rank - 1).rev() {
|
||
ds_strides[d] = ds_strides[d + 1] * saturating_usize(shape[d + 1]);
|
||
chunk_strides[d] = chunk_strides[d + 1] * saturating_usize(chunk_dims[d + 1]);
|
||
}
|
||
let row = valid[rank - 1] * element_size;
|
||
let mut idx = vec![0usize; rank];
|
||
loop {
|
||
let src: usize = (0..rank)
|
||
.map(|d| (saturating_usize(offsets[d]) + idx[d]) * ds_strides[d])
|
||
.sum::<usize>()
|
||
* element_size;
|
||
let dst: usize = (0..rank).map(|d| idx[d] * chunk_strides[d]).sum::<usize>() * element_size;
|
||
// Whole elements only, as far as `raw_data` reaches.
|
||
let n = row.min(raw_data.len().saturating_sub(src)) / element_size * element_size;
|
||
chunk[dst..dst + n].copy_from_slice(&raw_data[src..src + n]);
|
||
// Next row: advance every dimension but the last.
|
||
let mut d = rank - 1;
|
||
loop {
|
||
if d == 0 {
|
||
return (offsets, chunk);
|
||
}
|
||
d -= 1;
|
||
idx[d] += 1;
|
||
if idx[d] < valid[d] {
|
||
break;
|
||
}
|
||
idx[d] = 0;
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Parallel compression threshold: use rayon when chunk count exceeds this.
|
||
///
|
||
/// Lowered to 2 to enable parallel compression for typical 4-chunk workloads
|
||
/// (e.g., 128×128 matrix with 32-row chunks = 4 chunks). Rayon's overhead is
|
||
/// ~2 µs, worthwhile at ≥2 chunks with any real compression (arXiv:2206.14761).
|
||
#[cfg(feature = "parallel")]
|
||
const PARALLEL_COMPRESS_THRESHOLD: usize = 2;
|
||
|
||
/// Largest chunk compressed in parallel: every thread holds a chunk and its
|
||
/// compressed copy at once, so chunks larger than this (up to 4 GiB and
|
||
/// more) are compressed one after another.
|
||
#[cfg(feature = "parallel")]
|
||
const PARALLEL_COMPRESS_MAX_CHUNK_BYTES: u64 = 64 << 20;
|
||
|
||
/// Extract and compress every chunk of the dataset, returning each chunk's
|
||
/// raw size, stored bytes and filter mask, in chunk order.
|
||
///
|
||
/// Chunks run through the pipeline as libhdf5 runs them
|
||
/// ([`compress_chunk_masked`]): an optional filter that fails — LZF or Blosc
|
||
/// output no smaller than its input — is skipped and its mask bit set.
|
||
///
|
||
/// Each chunk is extracted just before it is compressed and dropped after,
|
||
/// so at most one raw chunk per thread is held. With the `parallel` feature,
|
||
/// more than [`PARALLEL_COMPRESS_THRESHOLD`] filtered chunks, and chunks of
|
||
/// at most [`PARALLEL_COMPRESS_MAX_CHUNK_BYTES`], compression runs across
|
||
/// rayon threads; otherwise it is sequential. Output order matches chunk
|
||
/// order, so per-chunk bytes are identical to the sequential path.
|
||
fn compress_all_chunks(
|
||
raw_data: &[u8],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
chunk_bytes: u64,
|
||
pipeline: &Option<FilterPipeline>,
|
||
) -> Result<Vec<(u64, Vec<u8>, u32)>, FormatError> {
|
||
let one = |i: u64| -> Result<(u64, Vec<u8>, u32), FormatError> {
|
||
let (_, raw) = if shape.is_empty() {
|
||
(Vec::new(), raw_data.to_vec())
|
||
} else {
|
||
extract_chunk(raw_data, shape, chunk_dims, element_size, i)
|
||
};
|
||
let raw_size = raw.len() as u64;
|
||
match pipeline {
|
||
Some(pl) => {
|
||
let (stored, mask) = compress_chunk_masked(&raw, pl, element_size as u32)?;
|
||
Ok((raw_size, stored, mask))
|
||
}
|
||
None => Ok((raw_size, raw, 0)),
|
||
}
|
||
};
|
||
let n = if shape.is_empty() {
|
||
1
|
||
} else {
|
||
chunk_count(shape, chunk_dims)
|
||
};
|
||
#[cfg(feature = "parallel")]
|
||
{
|
||
if pipeline.is_some()
|
||
&& n > PARALLEL_COMPRESS_THRESHOLD as u64
|
||
&& chunk_bytes <= PARALLEL_COMPRESS_MAX_CHUNK_BYTES
|
||
{
|
||
use rayon::prelude::*;
|
||
return (0..n).into_par_iter().map(one).collect();
|
||
}
|
||
}
|
||
#[cfg(not(feature = "parallel"))]
|
||
let _ = chunk_bytes;
|
||
|
||
// Sequential fallback
|
||
(0..n).map(one).collect()
|
||
}
|
||
|
||
/// Build the complete chunked dataset blob (chunk data + index) and return
|
||
/// layout/pipeline messages. `base_address` is where the blob will be placed in the file.
|
||
/// Serialize a v4 single chunk layout message (public for OH size estimation).
|
||
pub fn serialize_v4_single_chunk_pub(
|
||
chunk_dims: &[u32],
|
||
chunk_address: u64,
|
||
filtered_size: Option<u64>,
|
||
filter_mask: Option<u32>,
|
||
offset_size: u8,
|
||
element_size: u32,
|
||
) -> Vec<u8> {
|
||
serialize_v4_single_chunk(
|
||
chunk_dims,
|
||
chunk_address,
|
||
filtered_size,
|
||
filter_mask,
|
||
offset_size,
|
||
element_size,
|
||
4,
|
||
)
|
||
}
|
||
|
||
/// Serialize a v4 (or, for a chunk of 4 GiB or more, v5) single chunk
|
||
/// layout message.
|
||
fn serialize_v4_single_chunk(
|
||
chunk_dims: &[u32],
|
||
chunk_address: u64,
|
||
filtered_size: Option<u64>,
|
||
filter_mask: Option<u32>,
|
||
offset_size: u8,
|
||
element_size: u32,
|
||
version: u8,
|
||
) -> Vec<u8> {
|
||
let mut buf = Vec::new();
|
||
buf.push(version);
|
||
buf.push(2); // class = chunked
|
||
|
||
// flags: bit 0 = unknown meaning in some files, bit 1 = filters for single chunk
|
||
let flags: u8 = if filtered_size.is_some() { 0x02 } else { 0x00 };
|
||
buf.push(flags);
|
||
|
||
// dimensionality = rank + 1 (chunk dims + element size dim)
|
||
let ndims = chunk_dims.len() as u8 + 1;
|
||
buf.push(ndims);
|
||
|
||
push_v4_chunk_dims(&mut buf, chunk_dims, element_size);
|
||
|
||
// chunk index type = 1 (single chunk)
|
||
buf.push(1);
|
||
|
||
// Index-specific fields
|
||
if let (Some(fs), Some(fm)) = (filtered_size, filter_mask) {
|
||
// filtered_size (length_size bytes)
|
||
buf.extend_from_slice(&fs.to_le_bytes()); // 8 bytes for length_size=8
|
||
buf.extend_from_slice(&fm.to_le_bytes()); // 4 bytes
|
||
}
|
||
|
||
// chunk address
|
||
match offset_size {
|
||
4 => buf.extend_from_slice(&(chunk_address as u32).to_le_bytes()),
|
||
8 => buf.extend_from_slice(&chunk_address.to_le_bytes()),
|
||
_ => {}
|
||
}
|
||
|
||
buf
|
||
}
|
||
|
||
/// Serialize a v4 Fixed Array layout message.
|
||
fn serialize_v4_fixed_array(
|
||
chunk_dims: &[u32],
|
||
fixed_array_address: u64,
|
||
offset_size: u8,
|
||
element_size: u32,
|
||
max_bits: u8,
|
||
version: u8,
|
||
) -> Vec<u8> {
|
||
let mut buf = layout_v4_chunked_prefix(chunk_dims, element_size, version);
|
||
|
||
// chunk index type = 3 (Fixed Array)
|
||
buf.push(3);
|
||
|
||
// max_dblk_page_nelmts_bits — must match FAHD max_nelmts_bits
|
||
buf.push(max_bits);
|
||
|
||
// Fixed Array header address
|
||
match offset_size {
|
||
4 => buf.extend_from_slice(&(fixed_array_address as u32).to_le_bytes()),
|
||
8 => buf.extend_from_slice(&fixed_array_address.to_le_bytes()),
|
||
_ => {}
|
||
}
|
||
|
||
buf
|
||
}
|
||
|
||
/// The part of a v4 chunked layout message before the chunk index type:
|
||
/// version, class, flags and the chunk dimensions (plus the element size).
|
||
/// Append a v4 layout's dimension width and its dimensions (the chunk
|
||
/// dimensions, then the element size). Each takes the fewest bytes that hold
|
||
/// the largest, as libhdf5 computes it (`H5D__chunk_set_sizes`:
|
||
/// `(log2(dim) + 8) / 8`); HDF5 2.0.0 refuses any other width.
|
||
pub(crate) fn push_v4_chunk_dims(buf: &mut Vec<u8>, chunk_dims: &[u32], element_size: u32) {
|
||
let max_dim = chunk_dims
|
||
.iter()
|
||
.copied()
|
||
.chain(core::iter::once(element_size))
|
||
.max()
|
||
.unwrap_or(1)
|
||
.max(1);
|
||
let width = (32 - max_dim.leading_zeros()).div_ceil(8) as usize;
|
||
buf.push(width as u8);
|
||
for &d in chunk_dims.iter().chain(core::iter::once(&element_size)) {
|
||
buf.extend_from_slice(&d.to_le_bytes()[..width]);
|
||
}
|
||
}
|
||
|
||
fn layout_v4_chunked_prefix(chunk_dims: &[u32], element_size: u32, version: u8) -> Vec<u8> {
|
||
let mut buf = Vec::new();
|
||
buf.push(version);
|
||
buf.push(2); // class = chunked
|
||
|
||
let flags: u8 = 0x00;
|
||
buf.push(flags);
|
||
|
||
let ndims = chunk_dims.len() as u8 + 1;
|
||
buf.push(ndims);
|
||
|
||
push_v4_chunk_dims(&mut buf, chunk_dims, element_size);
|
||
buf
|
||
}
|
||
|
||
/// log2 of the elements per Fixed Array data block page (the library's
|
||
/// default, `H5D_FARRAY_MAX_DBLK_PAGE_NELMTS_BITS`).
|
||
const FA_PAGE_BITS: u8 = 10;
|
||
|
||
pub(crate) fn push_addr(buf: &mut Vec<u8>, addr: u64, offset_size: u8) {
|
||
match offset_size {
|
||
4 => buf.extend_from_slice(&(addr as u32).to_le_bytes()),
|
||
_ => buf.extend_from_slice(&addr.to_le_bytes()),
|
||
}
|
||
}
|
||
|
||
/// Width of the chunk-size field of a filtered chunk index element. Must
|
||
/// match the library's `H5D_FARRAY_FILT_COMPUTE_CHUNK_SIZE_LEN` (the EA and
|
||
/// B-tree v2 indexes use the same formula): see [`chunk_size_len`]. Chunks
|
||
/// of more than `u32::MAX` bytes are written with layout version 5
|
||
/// ([`layout_version_for`]).
|
||
pub(crate) fn filtered_chunk_size_len(slots: &[Option<WrittenChunk>], length_size: u8) -> usize {
|
||
let max_raw = slots
|
||
.iter()
|
||
.flatten()
|
||
.map(|c| c.raw_size)
|
||
.max()
|
||
.unwrap_or(1);
|
||
chunk_size_len(max_raw, layout_version_for(max_raw), length_size)
|
||
}
|
||
|
||
/// Largest chunk, in bytes, a layout message of version 4 or lower may
|
||
/// describe: libhdf5 writes a larger one with version 5
|
||
/// (`H5D__chunk_construct`: "chunk size > 4GB requires H5F_LIBVER_V200"),
|
||
/// which libhdf5 before 2.0 cannot read.
|
||
pub const MAX_V4_CHUNK_BYTES: u64 = u32::MAX as u64;
|
||
|
||
/// The layout message version clawhdf5 writes for chunks of `chunk_bytes`
|
||
/// bytes: 4, or 5 for a chunk larger than [`MAX_V4_CHUNK_BYTES`] (what
|
||
/// libhdf5 2.x writes for it; the chunk index is chosen as for version 4).
|
||
pub fn layout_version_for(chunk_bytes: u64) -> u8 {
|
||
if chunk_bytes > MAX_V4_CHUNK_BYTES {
|
||
5
|
||
} else {
|
||
4
|
||
}
|
||
}
|
||
|
||
/// Width libhdf5 gives the stored-size field of a filtered chunk index
|
||
/// element (Fixed Array, Extensible Array, v2 B-tree) for chunks of
|
||
/// `chunk_bytes` bytes under layout message `layout_version`
|
||
/// (`H5D_FARRAY_FILT_COMPUTE_CHUNK_SIZE_LEN` and its EA and B-tree twins):
|
||
/// up to version 4, one byte more than the chunk's size needs,
|
||
/// `1 + ((log2(chunk_bytes) + 8) / 8)` capped at 8; from version 5 (HDF5
|
||
/// 2.0), the file's size of lengths (`length_size`), whatever the chunk.
|
||
pub fn chunk_size_len(chunk_bytes: u64, layout_version: u8, length_size: u8) -> usize {
|
||
if layout_version >= 5 {
|
||
return usize::from(length_size);
|
||
}
|
||
let log2 = if chunk_bytes <= 1 {
|
||
0
|
||
} else {
|
||
63 - chunk_bytes.leading_zeros()
|
||
};
|
||
(1 + ((log2 + 8) / 8) as usize).min(8)
|
||
}
|
||
|
||
/// Append one chunk index element: the chunk's address, plus its stored size
|
||
/// and filter mask when the dataset is filtered. `None` is an unallocated
|
||
/// chunk (undefined address, zero size and mask).
|
||
pub(crate) fn push_index_element(
|
||
buf: &mut Vec<u8>,
|
||
slot: Option<&WrittenChunk>,
|
||
offset_size: u8,
|
||
chunk_size_bytes: Option<usize>,
|
||
) {
|
||
match slot {
|
||
Some(c) => {
|
||
push_addr(buf, c.address, offset_size);
|
||
if let Some(n) = chunk_size_bytes {
|
||
buf.extend_from_slice(&c.compressed_size.to_le_bytes()[..n]);
|
||
buf.extend_from_slice(&c.filter_mask.to_le_bytes());
|
||
}
|
||
}
|
||
None => {
|
||
buf.extend(core::iter::repeat_n(0xFF, offset_size as usize));
|
||
if let Some(n) = chunk_size_bytes {
|
||
buf.extend(core::iter::repeat_n(0x00, n + 4));
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Build a complete Fixed Array at a known absolute address.
|
||
///
|
||
/// `slots` holds one entry per element of the array, i.e. per chunk of the
|
||
/// dataset's *maximum* extent in the order [`crate::chunk_grid`] defines;
|
||
/// `None` marks a chunk that is not allocated. An array with more elements
|
||
/// than fit in one page (`2^FA_PAGE_BITS`) gets a paged data block: a
|
||
/// page-init bitmap after the prefix, then one checksummed page per
|
||
/// `2^FA_PAGE_BITS` elements, the last one short (`H5FA__dblock_create`).
|
||
pub fn build_fixed_array_at(
|
||
slots: &[Option<WrittenChunk>],
|
||
offset_size: u8,
|
||
length_size: u8,
|
||
has_filters: bool,
|
||
fa_base_address: u64,
|
||
) -> Vec<u8> {
|
||
let os = offset_size as usize;
|
||
let num_elements = slots.len();
|
||
|
||
let chunk_size_bytes = has_filters.then(|| filtered_chunk_size_len(slots, length_size));
|
||
let elem_size = os + chunk_size_bytes.map_or(0, |n| n + 4);
|
||
let client_id: u8 = if has_filters { 1 } else { 0 };
|
||
|
||
// FAHD total size
|
||
let fahd_total_size = 4 + 1 + 1 + 1 + 1 + length_size as usize + os + 4;
|
||
let fadb_address = fa_base_address + fahd_total_size as u64;
|
||
|
||
let mut fahd = Vec::with_capacity(fahd_total_size);
|
||
fahd.extend_from_slice(b"FAHD");
|
||
fahd.push(0); // version
|
||
fahd.push(client_id);
|
||
fahd.push(elem_size as u8);
|
||
fahd.push(FA_PAGE_BITS);
|
||
match length_size {
|
||
4 => fahd.extend_from_slice(&(num_elements as u32).to_le_bytes()),
|
||
_ => fahd.extend_from_slice(&(num_elements as u64).to_le_bytes()),
|
||
}
|
||
push_addr(&mut fahd, fadb_address, offset_size);
|
||
let checksum = jenkins_lookup3(&fahd);
|
||
fahd.extend_from_slice(&checksum.to_le_bytes());
|
||
assert_eq!(fahd.len(), fahd_total_size);
|
||
|
||
// FADB prefix
|
||
let mut fadb = Vec::new();
|
||
fadb.extend_from_slice(b"FADB");
|
||
fadb.push(0); // version
|
||
fadb.push(client_id);
|
||
push_addr(&mut fadb, fa_base_address, offset_size);
|
||
|
||
let page_nelmts = 1usize << FA_PAGE_BITS;
|
||
if num_elements <= page_nelmts {
|
||
// Unpaged: the elements follow the prefix, one checksum over both.
|
||
for slot in slots {
|
||
push_index_element(&mut fadb, slot.as_ref(), offset_size, chunk_size_bytes);
|
||
}
|
||
let fadb_checksum = jenkins_lookup3(&fadb);
|
||
fadb.extend_from_slice(&fadb_checksum.to_le_bytes());
|
||
} else {
|
||
// Paged: every page is written, so every page-init bit is set
|
||
// (MSB-first, as `H5VM_bit_set` packs them). The prefix and bitmap
|
||
// share a checksum; each page carries its own.
|
||
let npages = num_elements.div_ceil(page_nelmts);
|
||
let mut bitmap = vec![0u8; npages.div_ceil(8)];
|
||
for p in 0..npages {
|
||
bitmap[p / 8] |= 0x80 >> (p % 8);
|
||
}
|
||
fadb.extend_from_slice(&bitmap);
|
||
let prefix_checksum = jenkins_lookup3(&fadb);
|
||
fadb.extend_from_slice(&prefix_checksum.to_le_bytes());
|
||
for page in slots.chunks(page_nelmts) {
|
||
let start = fadb.len();
|
||
for slot in page {
|
||
push_index_element(&mut fadb, slot.as_ref(), offset_size, chunk_size_bytes);
|
||
}
|
||
let page_checksum = jenkins_lookup3(&fadb[start..]);
|
||
fadb.extend_from_slice(&page_checksum.to_le_bytes());
|
||
}
|
||
}
|
||
|
||
let mut combined = fahd;
|
||
combined.extend_from_slice(&fadb);
|
||
combined
|
||
}
|
||
|
||
/// Compressed chunks ready to be laid out at any file address.
|
||
///
|
||
/// Created by [`precompress_chunks`] and consumed by
|
||
/// [`build_chunked_data_from_precompressed`]. Caching this between the two
|
||
/// writer passes eliminates the double-compression that the two-pass layout
|
||
/// algorithm previously performed.
|
||
pub struct PrecompressedChunks {
|
||
/// Per-chunk: (raw_size_bytes, stored_bytes, filter_mask). Bit `i` of
|
||
/// the mask is set when filter `i` was skipped (an optional filter that
|
||
/// failed); 0 for every chunk of an unfiltered dataset.
|
||
pub chunks: Vec<(u64, Vec<u8>, u32)>,
|
||
pub has_filters: bool,
|
||
pub element_size: usize,
|
||
pub shape: Vec<u64>,
|
||
pub chunk_dims: Vec<u64>,
|
||
pub pipeline_message: Option<Vec<u8>>,
|
||
}
|
||
|
||
/// Compress all chunks of a dataset without laying them out at a file address.
|
||
///
|
||
/// Call this once per dataset in Pass 1, cache the result, then call
|
||
/// [`build_chunked_data_from_precompressed`] in both Pass 1 (dummy address
|
||
/// for sizing) and Pass 2 (real address) to avoid re-compressing.
|
||
pub fn precompress_chunks(
|
||
raw_data: &[u8],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
options: &ChunkOptions,
|
||
) -> Result<PrecompressedChunks, FormatError> {
|
||
let (_, chunk_bytes) = checked_chunk_dims(chunk_dims, element_size)?;
|
||
if chunk_bytes > MAX_V4_CHUNK_BYTES {
|
||
check_huge_chunk_filters(options, chunk_bytes)?;
|
||
}
|
||
let pipeline = options
|
||
.build_pipeline_for_chunk(element_size as u32, u32::try_from(chunk_bytes).unwrap_or(0));
|
||
let has_filters = pipeline.is_some();
|
||
let pipeline_message = pipeline.as_ref().map(|pl| pl.serialize());
|
||
|
||
let chunks = compress_all_chunks(
|
||
raw_data,
|
||
shape,
|
||
chunk_dims,
|
||
element_size,
|
||
chunk_bytes,
|
||
&pipeline,
|
||
)?;
|
||
|
||
Ok(PrecompressedChunks {
|
||
chunks,
|
||
has_filters,
|
||
element_size,
|
||
shape: shape.to_vec(),
|
||
chunk_dims: chunk_dims.to_vec(),
|
||
pipeline_message,
|
||
})
|
||
}
|
||
|
||
/// The chunk dimensions as the layout message stores them (each below
|
||
/// 2^32), and one chunk's size in bytes. A chunk dimension of 2^32 or more
|
||
/// (which HDF5 2.0 can store, in wider fields) is refused, as it is when
|
||
/// read: clawhdf5 holds chunk dimensions as `u32`. So is a chunk whose size
|
||
/// overflows 64 bits, or that this platform cannot hold in memory (a chunk
|
||
/// of 4 GiB or more on a 32-bit target).
|
||
fn checked_chunk_dims(
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
) -> Result<(Vec<u32>, u64), FormatError> {
|
||
let dims = chunk_dims
|
||
.iter()
|
||
.map(|&d| {
|
||
u32::try_from(d).map_err(|_| {
|
||
FormatError::InvalidChunkDimensions(format!(
|
||
"chunk dimension {d} is 2^32 or more, which clawhdf5 does not support"
|
||
))
|
||
})
|
||
})
|
||
.collect::<Result<Vec<u32>, _>>()?;
|
||
let bytes = chunk_dims
|
||
.iter()
|
||
.try_fold(element_size as u64, |acc, &d| acc.checked_mul(d))
|
||
.filter(|&b| usize::try_from(b).is_ok())
|
||
.ok_or_else(|| {
|
||
FormatError::Overflow(format!(
|
||
"a chunk of {chunk_dims:?} x {element_size} bytes exceeds this platform's address space"
|
||
))
|
||
})?;
|
||
Ok((dims, bytes))
|
||
}
|
||
|
||
/// The filters clawhdf5 can apply to a chunk of more than `u32::MAX` bytes:
|
||
/// shuffle, deflate, Zstandard, LZ4 (whose HDF5 framing records the size in
|
||
/// 64 bits and splits the chunk into blocks) and Fletcher32. The others
|
||
/// record the chunk size or their block lengths in 32 bits, or cannot take
|
||
/// a buffer that large (h5py's LZF, bitshuffle, bzip2, Blosc), and pcodec is
|
||
/// clawhdf5's own; they are refused rather than written into a chunk
|
||
/// libhdf5 could not decode.
|
||
fn check_huge_chunk_filters(options: &ChunkOptions, chunk_bytes: u64) -> Result<(), FormatError> {
|
||
let refused = if let Some(plugin) = &options.plugin {
|
||
Some(match plugin {
|
||
PluginFilter::Lzf => "LZF",
|
||
PluginFilter::Bitshuffle { .. } => "bitshuffle",
|
||
PluginFilter::Bzip2 { .. } => "bzip2",
|
||
PluginFilter::Blosc { .. } => "Blosc",
|
||
})
|
||
} else if options.pcodec {
|
||
Some("pcodec")
|
||
} else {
|
||
None
|
||
};
|
||
match refused {
|
||
Some(name) => Err(FormatError::FilterError(format!(
|
||
"{name} cannot compress a chunk of {chunk_bytes} bytes (4 GiB or more); \
|
||
use smaller chunks, or deflate, Zstandard or LZ4"
|
||
))),
|
||
None => Ok(()),
|
||
}
|
||
}
|
||
|
||
/// Lay out precompressed chunks at `base_address` and build index structures.
|
||
///
|
||
/// This is the address-dependent half of chunk writing. Call it in Pass 1
|
||
/// with a dummy address (to get the blob size), and again in Pass 2 with the
|
||
/// real address — both times reusing the same [`PrecompressedChunks`] so
|
||
/// compression happens only once.
|
||
pub fn build_chunked_data_from_precompressed(
|
||
pre: &PrecompressedChunks,
|
||
base_address: u64,
|
||
maxshape: Option<&[u64]>,
|
||
) -> Result<ChunkedDataResult, FormatError> {
|
||
build_chunked_data_from_precompressed_libver(
|
||
pre,
|
||
base_address,
|
||
maxshape,
|
||
LibVer::Latest,
|
||
LibVer::Latest,
|
||
)
|
||
}
|
||
|
||
/// [`build_chunked_data_from_precompressed`] for a file whose low library
|
||
/// version bound is `low`: below [`LibVer::V110`] (that is, for HDF5 1.8)
|
||
/// every chunked dataset gets a version-3 layout message and a version-1
|
||
/// B-tree chunk index, whatever its maximum shape, as libhdf5 writes it;
|
||
/// otherwise the version-4 layout and the index libhdf5 picks for it.
|
||
///
|
||
/// A chunk of 4 GiB or more (over [`MAX_V4_CHUNK_BYTES`]) takes layout
|
||
/// message version 5 whatever `low` is, as in libhdf5 (a version-1 B-tree
|
||
/// key holds a 32-bit size), and needs a `high` bound of at least
|
||
/// [`LibVer::V200`] ([`FormatError::LibverBound`] otherwise).
|
||
pub fn build_chunked_data_from_precompressed_libver(
|
||
pre: &PrecompressedChunks,
|
||
base_address: u64,
|
||
maxshape: Option<&[u64]>,
|
||
low: LibVer,
|
||
high: LibVer,
|
||
) -> Result<ChunkedDataResult, FormatError> {
|
||
let (_, chunk_bytes) = checked_chunk_dims(&pre.chunk_dims, pre.element_size)?;
|
||
if chunk_bytes > MAX_V4_CHUNK_BYTES {
|
||
if high < LibVer::V200 {
|
||
return Err(FormatError::LibverBound {
|
||
what: format!("a chunk of {chunk_bytes} bytes (4 GiB or more)"),
|
||
needs: LibVer::V200,
|
||
high,
|
||
});
|
||
}
|
||
} else if low < LibVer::V110 {
|
||
return build_btree_v1_chunked_data(pre, base_address, maxshape);
|
||
}
|
||
let index = ChunkIndexPlan::new(&pre.shape, maxshape, &pre.chunk_dims)?;
|
||
let offset_size: u8 = 8;
|
||
let length_size: u8 = 8;
|
||
let num_chunks = pre.chunks.len();
|
||
let element_size = pre.element_size;
|
||
|
||
let mut data_buf = Vec::new();
|
||
let mut written_chunks = Vec::with_capacity(num_chunks);
|
||
|
||
for (raw_size, compressed, filter_mask) in &pre.chunks {
|
||
let aligned_offset = align_to_cache_line(data_buf.len());
|
||
if aligned_offset > data_buf.len() {
|
||
data_buf.resize(aligned_offset, 0u8);
|
||
}
|
||
let address = base_address + data_buf.len() as u64;
|
||
let compressed_size = compressed.len() as u64;
|
||
data_buf.extend_from_slice(compressed);
|
||
written_chunks.push(WrittenChunk {
|
||
address,
|
||
compressed_size,
|
||
raw_size: *raw_size,
|
||
filter_mask: *filter_mask,
|
||
});
|
||
}
|
||
|
||
let (chunk_dims_u32, chunk_bytes) = checked_chunk_dims(&pre.chunk_dims, element_size)?;
|
||
let version = layout_version_for(chunk_bytes);
|
||
|
||
let aligned_idx = align_to_cache_line(data_buf.len());
|
||
if aligned_idx > data_buf.len() {
|
||
data_buf.resize(aligned_idx, 0u8);
|
||
}
|
||
|
||
let layout_message = match &index {
|
||
ChunkIndexPlan::ExtensibleArray(grid) => {
|
||
let ea_address = base_address + data_buf.len() as u64;
|
||
let slots = index_slots(grid, &pre.shape, &pre.chunk_dims, &written_chunks, None)?;
|
||
let ea_bytes = ea_writer::build_extensible_array_at(
|
||
&slots,
|
||
offset_size,
|
||
length_size,
|
||
pre.has_filters,
|
||
ea_address,
|
||
);
|
||
data_buf.extend_from_slice(&ea_bytes);
|
||
ea_writer::serialize_v4_extensible_array(
|
||
&chunk_dims_u32,
|
||
ea_address,
|
||
offset_size,
|
||
element_size as u32,
|
||
version,
|
||
)
|
||
}
|
||
ChunkIndexPlan::SingleChunk => {
|
||
let chunk_addr = written_chunks[0].address;
|
||
let filtered_size = if pre.has_filters {
|
||
Some(written_chunks[0].compressed_size)
|
||
} else {
|
||
None
|
||
};
|
||
let filter_mask = pre.has_filters.then_some(written_chunks[0].filter_mask);
|
||
serialize_v4_single_chunk(
|
||
&chunk_dims_u32,
|
||
chunk_addr,
|
||
filtered_size,
|
||
filter_mask,
|
||
offset_size,
|
||
element_size as u32,
|
||
version,
|
||
)
|
||
}
|
||
ChunkIndexPlan::FixedArray(grid, nslots) => {
|
||
let fa_address = base_address + data_buf.len() as u64;
|
||
let slots = index_slots(
|
||
grid,
|
||
&pre.shape,
|
||
&pre.chunk_dims,
|
||
&written_chunks,
|
||
Some(*nslots),
|
||
)?;
|
||
let fa_bytes = build_fixed_array_at(
|
||
&slots,
|
||
offset_size,
|
||
length_size,
|
||
pre.has_filters,
|
||
fa_address,
|
||
);
|
||
data_buf.extend_from_slice(&fa_bytes);
|
||
serialize_v4_fixed_array(
|
||
&chunk_dims_u32,
|
||
fa_address,
|
||
offset_size,
|
||
element_size as u32,
|
||
FA_PAGE_BITS,
|
||
version,
|
||
)
|
||
}
|
||
ChunkIndexPlan::BTreeV2 => {
|
||
let bt_address = base_address + data_buf.len() as u64;
|
||
let records: Vec<(Vec<u64>, &WrittenChunk)> = written_chunks
|
||
.iter()
|
||
.enumerate()
|
||
.map(|(i, c)| (scaled_coords(&pre.shape, &pre.chunk_dims, i), c))
|
||
.collect();
|
||
let (bt_bytes, node_size) = build_btree_v2_chunk_index_at(
|
||
pre.shape.len(),
|
||
&records,
|
||
offset_size,
|
||
length_size,
|
||
pre.has_filters,
|
||
bt_address,
|
||
)?;
|
||
data_buf.extend_from_slice(&bt_bytes);
|
||
serialize_v4_btree_v2(
|
||
&chunk_dims_u32,
|
||
bt_address,
|
||
offset_size,
|
||
element_size as u32,
|
||
node_size,
|
||
version,
|
||
)
|
||
}
|
||
};
|
||
|
||
Ok(ChunkedDataResult {
|
||
data_bytes: data_buf,
|
||
layout_message,
|
||
pipeline_message: pre.pipeline_message.clone(),
|
||
})
|
||
}
|
||
|
||
/// Lay out precompressed chunks at `base_address` followed by a version-1
|
||
/// B-tree chunk index, with a version-3 layout message: what libhdf5 writes
|
||
/// for a chunked dataset under a low bound of 1.8.
|
||
fn build_btree_v1_chunked_data(
|
||
pre: &PrecompressedChunks,
|
||
base_address: u64,
|
||
maxshape: Option<&[u64]>,
|
||
) -> Result<ChunkedDataResult, FormatError> {
|
||
if let Some(ms) = maxshape {
|
||
let bad = |what: &str| FormatError::ChunkedReadError(format!("maxshape: {what}"));
|
||
if ms.len() != pre.shape.len() {
|
||
return Err(bad("rank differs from the shape"));
|
||
}
|
||
if ms.iter().zip(&pre.shape).any(|(&m, &s)| m < s) {
|
||
return Err(bad("smaller than the shape"));
|
||
}
|
||
}
|
||
let offset_size: u8 = 8;
|
||
let mut data_buf = Vec::new();
|
||
let mut entries = Vec::with_capacity(pre.chunks.len());
|
||
for (i, (_raw_size, stored, filter_mask)) in pre.chunks.iter().enumerate() {
|
||
let aligned_offset = align_to_cache_line(data_buf.len());
|
||
if aligned_offset > data_buf.len() {
|
||
data_buf.resize(aligned_offset, 0u8);
|
||
}
|
||
entries.push(btree_v1_write::ChunkEntry {
|
||
scaled: scaled_coords(&pre.shape, &pre.chunk_dims, i),
|
||
nbytes: stored.len() as u64,
|
||
filter_mask: *filter_mask,
|
||
address: base_address + data_buf.len() as u64,
|
||
});
|
||
data_buf.extend_from_slice(stored);
|
||
}
|
||
let element_size = u32::try_from(pre.element_size)
|
||
.map_err(|_| FormatError::Overflow("element size".into()))?;
|
||
// A dataset with no chunks has no tree: its address is undefined, as
|
||
// libhdf5 leaves it until the first chunk is written.
|
||
let btree_address = if entries.is_empty() {
|
||
u64::MAX
|
||
} else {
|
||
let aligned_idx = align_to_cache_line(data_buf.len());
|
||
if aligned_idx > data_buf.len() {
|
||
data_buf.resize(aligned_idx, 0u8);
|
||
}
|
||
let addr = base_address + data_buf.len() as u64;
|
||
let tree = btree_v1_write::build_chunk_btree_v1_at(
|
||
&entries,
|
||
&pre.chunk_dims,
|
||
element_size,
|
||
addr,
|
||
offset_size,
|
||
)?;
|
||
data_buf.extend_from_slice(&tree);
|
||
addr
|
||
};
|
||
let layout_message =
|
||
serialize_v3_chunked(&pre.chunk_dims, btree_address, offset_size, element_size)?;
|
||
Ok(ChunkedDataResult {
|
||
data_bytes: data_buf,
|
||
layout_message,
|
||
pipeline_message: pre.pipeline_message.clone(),
|
||
})
|
||
}
|
||
|
||
/// A version-3 layout message for a chunked dataset: dimensionality (the
|
||
/// rank plus one), the B-tree's address, then each chunk dimension and the
|
||
/// element size, four bytes each.
|
||
fn serialize_v3_chunked(
|
||
chunk_dims: &[u64],
|
||
btree_address: u64,
|
||
offset_size: u8,
|
||
element_size: u32,
|
||
) -> Result<Vec<u8>, FormatError> {
|
||
let ndims = u8::try_from(chunk_dims.len() + 1)
|
||
.map_err(|_| FormatError::Overflow("chunked layout rank".into()))?;
|
||
let mut buf = vec![3u8, 2, ndims];
|
||
push_addr(&mut buf, btree_address, offset_size);
|
||
for &d in chunk_dims {
|
||
let d =
|
||
u32::try_from(d).map_err(|_| FormatError::Overflow(format!("chunk dimension {d}")))?;
|
||
buf.extend_from_slice(&d.to_le_bytes());
|
||
}
|
||
buf.extend_from_slice(&element_size.to_le_bytes());
|
||
Ok(buf)
|
||
}
|
||
|
||
/// Most slots a Fixed Array index may have before we refuse to build it: its
|
||
/// data block holds one element per chunk of the *maximum* extent, so a huge
|
||
/// finite maxshape with small chunks would otherwise exhaust memory.
|
||
const MAX_FIXED_ARRAY_SLOTS: u64 = 1 << 26;
|
||
|
||
/// Which chunk index a dataset gets, following the library's choice in
|
||
/// `H5D__layout_set_latest_indexing`: version-2 B-tree for more than one
|
||
/// unlimited dimension, Extensible Array for exactly one, Fixed Array for a
|
||
/// finite maxshape, Single Chunk when the whole maximum extent is one chunk.
|
||
enum ChunkIndexPlan {
|
||
SingleChunk,
|
||
/// The grid and the number of array elements (chunks of the max extent).
|
||
FixedArray(ChunkGrid, usize),
|
||
ExtensibleArray(ChunkGrid),
|
||
BTreeV2,
|
||
}
|
||
|
||
impl ChunkIndexPlan {
|
||
fn new(
|
||
shape: &[u64],
|
||
maxshape: Option<&[u64]>,
|
||
chunk_dims: &[u64],
|
||
) -> Result<Self, FormatError> {
|
||
let bad = |what: &str| FormatError::ChunkedReadError(format!("maxshape: {what}"));
|
||
if let Some(ms) = maxshape {
|
||
if ms.len() != shape.len() {
|
||
return Err(bad("rank differs from the shape"));
|
||
}
|
||
if ms.iter().zip(shape).any(|(&m, &s)| m < s) {
|
||
return Err(bad("smaller than the shape"));
|
||
}
|
||
}
|
||
let max = maxshape.unwrap_or(shape);
|
||
let nunlim = max.iter().filter(|&&d| d == u64::MAX).count();
|
||
match nunlim {
|
||
0 => {
|
||
let nslots = max
|
||
.iter()
|
||
.zip(chunk_dims)
|
||
.try_fold(1u64, |acc, (&m, &c)| acc.checked_mul(m.div_ceil(c.max(1))))
|
||
.filter(|&n| n <= MAX_FIXED_ARRAY_SLOTS)
|
||
.ok_or_else(|| {
|
||
bad("too many chunks for a Fixed Array index; \
|
||
use larger chunks or an unlimited dimension")
|
||
})?;
|
||
// A Single Chunk index needs that one chunk to exist; an
|
||
// empty dataset gets an all-unallocated Fixed Array instead.
|
||
let empty = shape.contains(&0);
|
||
if nslots == 1 && !empty {
|
||
Ok(Self::SingleChunk)
|
||
} else {
|
||
let grid = ChunkGrid::fixed_array(shape, Some(max), chunk_dims)?;
|
||
Ok(Self::FixedArray(grid, saturating_usize(nslots)))
|
||
}
|
||
}
|
||
1 => Ok(Self::ExtensibleArray(ChunkGrid::extensible_array(
|
||
shape,
|
||
Some(max),
|
||
chunk_dims,
|
||
)?)),
|
||
_ => Ok(Self::BTreeV2),
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Place each written chunk at its linear index in `grid`. `chunks` are in
|
||
/// row-major order over the chunks of the current extent (`split_into_chunks`).
|
||
/// `len` fixes the slot count (Fixed Array); otherwise it is one past the
|
||
/// highest index used.
|
||
fn index_slots(
|
||
grid: &ChunkGrid,
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
chunks: &[WrittenChunk],
|
||
len: Option<usize>,
|
||
) -> Result<Vec<Option<WrittenChunk>>, FormatError> {
|
||
let mut placed: Vec<(usize, &WrittenChunk)> = Vec::with_capacity(chunks.len());
|
||
for (i, chunk) in chunks.iter().enumerate() {
|
||
let scaled = scaled_coords(shape, chunk_dims, i);
|
||
let idx = usize::try_from(grid.linear_index(&scaled))
|
||
.map_err(|_| FormatError::Overflow("chunk index slot".into()))?;
|
||
placed.push((idx, chunk));
|
||
}
|
||
let n = len.unwrap_or_else(|| placed.iter().map(|&(i, _)| i + 1).max().unwrap_or(0));
|
||
let mut slots = vec![None; n];
|
||
for (idx, chunk) in placed {
|
||
*slots
|
||
.get_mut(idx)
|
||
.ok_or_else(|| FormatError::Overflow("chunk index slot".into()))? = Some(chunk.clone());
|
||
}
|
||
Ok(slots)
|
||
}
|
||
|
||
/// Scaled coordinates (`offset / chunk_dim`) of the `i`-th chunk in the
|
||
/// row-major order `split_into_chunks` produces over the current extent.
|
||
fn scaled_coords(shape: &[u64], chunk_dims: &[u64], i: usize) -> Vec<u64> {
|
||
let rank = shape.len();
|
||
let mut scaled = vec![0u64; rank];
|
||
let mut rem = i as u64;
|
||
for d in (0..rank).rev() {
|
||
let n = shape[d].div_ceil(chunk_dims[d]);
|
||
scaled[d] = rem % n;
|
||
rem /= n;
|
||
}
|
||
scaled
|
||
}
|
||
|
||
/// Node size the library gives a chunk index B-tree (`H5D_BT2_NODE_SIZE`),
|
||
/// with its split and merge percentages.
|
||
const BT2_NODE_SIZE: u32 = 2048;
|
||
const BT2_SPLIT_PERCENT: u8 = 100;
|
||
const BT2_MERGE_PERCENT: u8 = 40;
|
||
/// B-tree v2 record types for chunk indexes (`H5B2_CDSET_ID`,
|
||
/// `H5B2_CDSET_FILT_ID`).
|
||
const BT2_CHUNK_UNFILTERED: u8 = 10;
|
||
const BT2_CHUNK_FILTERED: u8 = 11;
|
||
|
||
/// Build a version-2 B-tree chunk index (the library's index for datasets
|
||
/// with more than one unlimited dimension) at a known absolute address.
|
||
///
|
||
/// `records` are `(scaled coordinates, chunk)` in lexicographic order of the
|
||
/// coordinates, which is the order the library's comparator
|
||
/// (`H5VM_vector_cmp_u`) keeps them in. Up to 65 535 chunks go in a single
|
||
/// leaf: the library's 2048-byte node when the records fit, otherwise a leaf
|
||
/// node sized to hold them all (the layout the writer has always used, kept
|
||
/// so those files do not change). More chunks get the library's 2048-byte
|
||
/// nodes with internal nodes above the leaves. Returns the bytes and the
|
||
/// node size the layout message must record.
|
||
fn build_btree_v2_chunk_index_at(
|
||
rank: usize,
|
||
records: &[(Vec<u64>, &WrittenChunk)],
|
||
offset_size: u8,
|
||
length_size: u8,
|
||
has_filters: bool,
|
||
base_address: u64,
|
||
) -> Result<(Vec<u8>, u32), FormatError> {
|
||
let os = offset_size as usize;
|
||
let chunk_size_bytes = has_filters.then(|| {
|
||
let slots: Vec<Option<WrittenChunk>> =
|
||
records.iter().map(|(_, c)| Some((*c).clone())).collect();
|
||
filtered_chunk_size_len(&slots, length_size)
|
||
});
|
||
let record_size = os + chunk_size_bytes.map_or(0, |n| n + 4) + 8 * rank;
|
||
let record_size_u16 = u16::try_from(record_size)
|
||
.map_err(|_| FormatError::Overflow("B-tree v2 record size".into()))?;
|
||
let node_size = if records.len() <= usize::from(u16::MAX) {
|
||
// Leaf: signature, version, type, records, checksum.
|
||
let leaf_len = 4 + 1 + 1 + records.len() * record_size + 4;
|
||
u32::try_from(leaf_len)
|
||
.map_err(|_| FormatError::Overflow("B-tree v2 leaf size".into()))?
|
||
.max(BT2_NODE_SIZE)
|
||
} else {
|
||
BT2_NODE_SIZE
|
||
};
|
||
let tree_type = if has_filters {
|
||
BT2_CHUNK_FILTERED
|
||
} else {
|
||
BT2_CHUNK_UNFILTERED
|
||
};
|
||
|
||
let mut flat = Vec::with_capacity(records.len() * record_size);
|
||
for (scaled, chunk) in records {
|
||
push_index_element(&mut flat, Some(chunk), offset_size, chunk_size_bytes);
|
||
for &c in scaled {
|
||
flat.extend_from_slice(&c.to_le_bytes());
|
||
}
|
||
}
|
||
let out = build_btree_v2(
|
||
BTreeV2Params {
|
||
tree_type,
|
||
node_size,
|
||
record_size: record_size_u16,
|
||
split_percent: BT2_SPLIT_PERCENT,
|
||
merge_percent: BT2_MERGE_PERCENT,
|
||
},
|
||
&flat,
|
||
base_address,
|
||
offset_size,
|
||
length_size,
|
||
)?;
|
||
Ok((out, node_size))
|
||
}
|
||
|
||
/// Serialize a v4 layout message for a version-2 B-tree chunk index.
|
||
fn serialize_v4_btree_v2(
|
||
chunk_dims: &[u32],
|
||
btree_address: u64,
|
||
offset_size: u8,
|
||
element_size: u32,
|
||
node_size: u32,
|
||
version: u8,
|
||
) -> Vec<u8> {
|
||
let mut buf = layout_v4_chunked_prefix(chunk_dims, element_size, version);
|
||
buf.push(5); // chunk index type = 5 (version-2 B-tree)
|
||
buf.extend_from_slice(&node_size.to_le_bytes());
|
||
buf.push(BT2_SPLIT_PERCENT);
|
||
buf.push(BT2_MERGE_PERCENT);
|
||
push_addr(&mut buf, btree_address, offset_size);
|
||
buf
|
||
}
|
||
|
||
/// Build chunked data with absolute addresses.
|
||
/// If `maxshape` has unlimited dims, uses Extensible Array index.
|
||
pub fn build_chunked_data_at(
|
||
raw_data: &[u8],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
options: &ChunkOptions,
|
||
base_address: u64,
|
||
) -> Result<ChunkedDataResult, FormatError> {
|
||
build_chunked_data_at_ext(
|
||
raw_data,
|
||
shape,
|
||
chunk_dims,
|
||
element_size,
|
||
options,
|
||
base_address,
|
||
None,
|
||
)
|
||
}
|
||
|
||
/// Build chunked data with absolute addresses and optional maxshape.
|
||
pub fn build_chunked_data_at_ext(
|
||
raw_data: &[u8],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
element_size: usize,
|
||
options: &ChunkOptions,
|
||
base_address: u64,
|
||
maxshape: Option<&[u64]>,
|
||
) -> Result<ChunkedDataResult, FormatError> {
|
||
let pre = precompress_chunks(raw_data, shape, chunk_dims, element_size, options)?;
|
||
build_chunked_data_from_precompressed(&pre, base_address, maxshape)
|
||
}
|
||
|
||
/// Write selected elements into an existing in-memory dataset buffer.
|
||
///
|
||
/// This performs a read-modify-write on the full buffer: elements matching
|
||
/// the selection are overwritten with the provided `new_data`. The buffer
|
||
/// must be a complete, uncompressed dataset of the given shape.
|
||
///
|
||
/// This is the in-memory equivalent of partial hyperslab writes. The caller
|
||
/// is responsible for re-chunking and re-compressing the buffer afterward.
|
||
pub fn write_selection_to_buffer(
|
||
buffer: &mut [u8],
|
||
dims: &[u64],
|
||
elem_size: usize,
|
||
selection: &crate::selection::Selection,
|
||
new_data: &[u8],
|
||
) {
|
||
use crate::selection::Selection;
|
||
|
||
match selection {
|
||
Selection::All => {
|
||
let len = buffer.len().min(new_data.len());
|
||
buffer[..len].copy_from_slice(&new_data[..len]);
|
||
}
|
||
Selection::None => {}
|
||
Selection::Hyperslab {
|
||
start,
|
||
stride,
|
||
count,
|
||
block,
|
||
} => {
|
||
let rank = dims.len();
|
||
let mut ds_strides = vec![1usize; rank];
|
||
for i in (0..rank.saturating_sub(1)).rev() {
|
||
ds_strides[i] = ds_strides[i + 1] * saturating_usize(dims[i + 1]);
|
||
}
|
||
|
||
let mut src_offset = 0usize;
|
||
|
||
#[allow(clippy::too_many_arguments)]
|
||
fn write_hyperslab(
|
||
d: usize,
|
||
rank: usize,
|
||
start: &[u64],
|
||
stride: &[u64],
|
||
count: &[u64],
|
||
block: &[u64],
|
||
dims: &[u64],
|
||
ds_strides: &[usize],
|
||
elem_size: usize,
|
||
buffer: &mut [u8],
|
||
new_data: &[u8],
|
||
src_offset: &mut usize,
|
||
current_ds_offset: usize,
|
||
) {
|
||
if d == rank {
|
||
let dst = current_ds_offset * elem_size;
|
||
let src = *src_offset * elem_size;
|
||
if dst + elem_size <= buffer.len() && src + elem_size <= new_data.len() {
|
||
buffer[dst..dst + elem_size]
|
||
.copy_from_slice(&new_data[src..src + elem_size]);
|
||
}
|
||
*src_offset += 1;
|
||
return;
|
||
}
|
||
|
||
for bi in 0..count[d] {
|
||
let block_start = start[d] + bi * stride[d];
|
||
for bj in 0..block[d] {
|
||
let coord = block_start + bj;
|
||
if coord < dims[d] {
|
||
write_hyperslab(
|
||
d + 1,
|
||
rank,
|
||
start,
|
||
stride,
|
||
count,
|
||
block,
|
||
dims,
|
||
ds_strides,
|
||
elem_size,
|
||
buffer,
|
||
new_data,
|
||
src_offset,
|
||
current_ds_offset + saturating_usize(coord) * ds_strides[d],
|
||
);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
write_hyperslab(
|
||
0,
|
||
rank,
|
||
start,
|
||
stride,
|
||
count,
|
||
block,
|
||
dims,
|
||
&ds_strides,
|
||
elem_size,
|
||
buffer,
|
||
new_data,
|
||
&mut src_offset,
|
||
0,
|
||
);
|
||
}
|
||
Selection::Points(pts) => {
|
||
let rank = dims.len();
|
||
let mut ds_strides = vec![1usize; rank];
|
||
for i in (0..rank.saturating_sub(1)).rev() {
|
||
ds_strides[i] = ds_strides[i + 1] * saturating_usize(dims[i + 1]);
|
||
}
|
||
|
||
for (pi, pt) in pts.iter().enumerate() {
|
||
let flat: usize = pt
|
||
.iter()
|
||
.zip(ds_strides.iter())
|
||
.map(|(&p, &s)| saturating_usize(p) * s)
|
||
.sum();
|
||
let dst = flat * elem_size;
|
||
let src = pi * elem_size;
|
||
if dst + elem_size <= buffer.len() && src + elem_size <= new_data.len() {
|
||
buffer[dst..dst + elem_size].copy_from_slice(&new_data[src..src + elem_size]);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
#[cfg(test)]
|
||
mod tests {
|
||
|
||
use super::*;
|
||
use crate::chunked_read::read_chunked_data;
|
||
use crate::data_layout::DataLayout;
|
||
use crate::dataspace::{Dataspace, DataspaceType};
|
||
use crate::datatype::{Datatype, DatatypeByteOrder};
|
||
|
||
fn make_f64_type() -> Datatype {
|
||
Datatype::FloatingPoint {
|
||
size: 8,
|
||
byte_order: DatatypeByteOrder::LittleEndian,
|
||
bit_offset: 0,
|
||
bit_precision: 64,
|
||
exponent_location: 52,
|
||
exponent_size: 11,
|
||
mantissa_location: 0,
|
||
mantissa_size: 52,
|
||
exponent_bias: 1023,
|
||
}
|
||
}
|
||
|
||
fn f64_to_bytes(data: &[f64]) -> Vec<u8> {
|
||
let mut b = Vec::with_capacity(data.len() * 8);
|
||
for &v in data {
|
||
b.extend_from_slice(&v.to_le_bytes());
|
||
}
|
||
b
|
||
}
|
||
|
||
fn bytes_to_f64(data: &[u8]) -> Vec<f64> {
|
||
data.chunks(8)
|
||
.map(|c| f64::from_le_bytes(c.try_into().unwrap()))
|
||
.collect()
|
||
}
|
||
|
||
/// Helper: build a chunked file blob and read it back using read_chunked_data
|
||
fn roundtrip_chunked(
|
||
values: &[f64],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
options: &ChunkOptions,
|
||
) -> Vec<f64> {
|
||
let raw = f64_to_bytes(values);
|
||
let base_address = 0x1000u64;
|
||
let result =
|
||
build_chunked_data_at(&raw, shape, chunk_dims, 8, options, base_address).unwrap();
|
||
|
||
// Build a fake file buffer
|
||
let file_size = base_address as usize + result.data_bytes.len();
|
||
let mut file_data = vec![0u8; file_size];
|
||
file_data[base_address as usize..].copy_from_slice(&result.data_bytes);
|
||
|
||
// Parse layout
|
||
let layout = DataLayout::parse(&result.layout_message, 8, 8).unwrap();
|
||
let dataspace = Dataspace {
|
||
space_type: DataspaceType::Simple,
|
||
rank: shape.len() as u8,
|
||
dimensions: shape.to_vec(),
|
||
max_dimensions: None,
|
||
};
|
||
let datatype = make_f64_type();
|
||
|
||
// Parse pipeline if present
|
||
let pipeline = result
|
||
.pipeline_message
|
||
.as_ref()
|
||
.map(|pm| crate::filter_pipeline::FilterPipeline::parse(pm).unwrap());
|
||
|
||
let output = read_chunked_data(
|
||
&file_data,
|
||
&layout,
|
||
&dataspace,
|
||
&datatype,
|
||
pipeline.as_ref(),
|
||
8,
|
||
8,
|
||
)
|
||
.unwrap();
|
||
|
||
bytes_to_f64(&output)
|
||
}
|
||
|
||
#[test]
|
||
fn split_1d_single_chunk() {
|
||
let data = f64_to_bytes(&[1.0, 2.0, 3.0]);
|
||
let result = split_into_chunks(&data, &[3], &[3], 8);
|
||
assert_eq!(result.len(), 1);
|
||
assert_eq!(result[0].0, vec![0]);
|
||
assert_eq!(bytes_to_f64(&result[0].1), vec![1.0, 2.0, 3.0]);
|
||
}
|
||
|
||
#[test]
|
||
fn split_1d_multiple_chunks() {
|
||
let values: Vec<f64> = (0..10).map(|i| i as f64).collect();
|
||
let data = f64_to_bytes(&values);
|
||
let result = split_into_chunks(&data, &[10], &[4], 8);
|
||
assert_eq!(result.len(), 3); // ceil(10/4) = 3
|
||
assert_eq!(result[0].0, vec![0]);
|
||
assert_eq!(result[1].0, vec![4]);
|
||
assert_eq!(result[2].0, vec![8]);
|
||
assert_eq!(bytes_to_f64(&result[0].1), vec![0.0, 1.0, 2.0, 3.0]);
|
||
assert_eq!(bytes_to_f64(&result[1].1), vec![4.0, 5.0, 6.0, 7.0]);
|
||
// Last chunk: 2 valid + 2 padding zeros
|
||
assert_eq!(bytes_to_f64(&result[2].1), vec![8.0, 9.0, 0.0, 0.0]);
|
||
}
|
||
|
||
#[test]
|
||
fn split_2d_chunks() {
|
||
// 4x4 dataset, 2x2 chunks -> 4 chunks
|
||
let values: Vec<f64> = (0..16).map(|i| i as f64).collect();
|
||
let data = f64_to_bytes(&values);
|
||
let result = split_into_chunks(&data, &[4, 4], &[2, 2], 8);
|
||
assert_eq!(result.len(), 4);
|
||
assert_eq!(result[0].0, vec![0, 0]);
|
||
assert_eq!(result[1].0, vec![0, 2]);
|
||
assert_eq!(result[2].0, vec![2, 0]);
|
||
assert_eq!(result[3].0, vec![2, 2]);
|
||
// chunk (0,0): elements [0,1,4,5]
|
||
assert_eq!(bytes_to_f64(&result[0].1), vec![0.0, 1.0, 4.0, 5.0]);
|
||
// chunk (0,2): elements [2,3,6,7]
|
||
assert_eq!(bytes_to_f64(&result[1].1), vec![2.0, 3.0, 6.0, 7.0]);
|
||
}
|
||
|
||
#[test]
|
||
fn roundtrip_1d_single_chunk_no_compression() {
|
||
let values: Vec<f64> = (0..10).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![10]),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[10], &[10], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[cfg(feature = "deflate")]
|
||
#[test]
|
||
fn roundtrip_1d_single_chunk_deflate() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![100]),
|
||
deflate_level: Some(6),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[100], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[test]
|
||
fn roundtrip_1d_multi_chunk_no_compression() {
|
||
let values: Vec<f64> = (0..20).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![8]),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[20], &[8], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[cfg(feature = "deflate")]
|
||
#[test]
|
||
fn roundtrip_1d_multi_chunk_deflate() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![20]),
|
||
deflate_level: Some(6),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[20], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[cfg(feature = "deflate")]
|
||
#[test]
|
||
fn roundtrip_1d_shuffle_deflate() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![50]),
|
||
deflate_level: Some(6),
|
||
shuffle: true,
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[50], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[test]
|
||
fn roundtrip_2d_chunks() {
|
||
// 6x4 dataset, 3x2 chunks
|
||
let values: Vec<f64> = (0..24).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![3, 2]),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[6, 4], &[3, 2], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[test]
|
||
fn align_chunk_offset_values() {
|
||
use super::CACHE_LINE_SIZE;
|
||
use super::align_chunk_offset;
|
||
let cl = CACHE_LINE_SIZE as u64;
|
||
assert_eq!(align_chunk_offset(0), 0);
|
||
assert_eq!(align_chunk_offset(1), cl);
|
||
assert_eq!(align_chunk_offset(cl), cl);
|
||
assert_eq!(align_chunk_offset(cl + 1), cl * 2);
|
||
assert_eq!(align_chunk_offset(cl * 10), cl * 10);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_addresses_are_cache_aligned() {
|
||
use super::align_chunk_offset;
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let raw = f64_to_bytes(&values);
|
||
let base_address = 0x1000u64;
|
||
// Ensure base is aligned for this test
|
||
let base_address = align_chunk_offset(base_address);
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![20]),
|
||
..Default::default()
|
||
};
|
||
let result = build_chunked_data_at(&raw, &[100], &[20], 8, &options, base_address).unwrap();
|
||
|
||
// Parse layout to get chunk addresses (via roundtrip read)
|
||
let file_size = base_address as usize + result.data_bytes.len();
|
||
let mut file_data = vec![0u8; file_size];
|
||
file_data[base_address as usize..].copy_from_slice(&result.data_bytes);
|
||
|
||
let layout = DataLayout::parse(&result.layout_message, 8, 8).unwrap();
|
||
let dataspace = Dataspace {
|
||
space_type: DataspaceType::Simple,
|
||
rank: 1,
|
||
dimensions: vec![100],
|
||
max_dimensions: None,
|
||
};
|
||
let datatype = make_f64_type();
|
||
|
||
// Verify data roundtrips correctly
|
||
let output =
|
||
read_chunked_data(&file_data, &layout, &dataspace, &datatype, None, 8, 8).unwrap();
|
||
assert_eq!(bytes_to_f64(&output), values);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_auto_dims() {
|
||
let options = ChunkOptions {
|
||
chunk_dims: None,
|
||
deflate_level: Some(6),
|
||
..Default::default()
|
||
};
|
||
let dims = options.resolve_chunk_dims(&[100, 50]);
|
||
assert_eq!(dims, vec![100, 50]);
|
||
}
|
||
|
||
#[test]
|
||
fn auto_chunking_splits_only_large_datasets() {
|
||
let bytes = |dims: &[u64], elem: u64| dims.iter().product::<u64>() * elem;
|
||
// Up to the target: one chunk, as before.
|
||
assert_eq!(auto_chunk_dims(&[100, 50], 8), [100, 50]);
|
||
assert_eq!(auto_chunk_dims(&[131_072], 8), [131_072]); // exactly 1 MiB
|
||
// Larger: split, keeping proportions, never above the target.
|
||
let big = auto_chunk_dims(&[4096, 2048], 8);
|
||
assert!(bytes(&big, 8) <= AUTO_CHUNK_TARGET_BYTES, "{big:?}");
|
||
assert!(bytes(&big, 8) > AUTO_CHUNK_TARGET_BYTES / 4, "{big:?}");
|
||
assert_eq!(big[0] / big[1], 2, "proportions kept: {big:?}");
|
||
// Every dimension stays within the dataset and at least 1.
|
||
for shape in [
|
||
vec![10_000_000u64],
|
||
vec![3, 5_000_000],
|
||
vec![1, 1, 9_000_000],
|
||
vec![7; 9],
|
||
] {
|
||
let dims = auto_chunk_dims(&shape, 4);
|
||
assert!(
|
||
dims.iter().zip(&shape).all(|(c, s)| *c >= 1 && c <= s),
|
||
"{shape:?} -> {dims:?}"
|
||
);
|
||
assert!(
|
||
bytes(&dims, 4) <= AUTO_CHUNK_TARGET_BYTES,
|
||
"{shape:?} -> {dims:?}"
|
||
);
|
||
}
|
||
// An empty (unlimited, unwritten) dimension still gets a usable chunk.
|
||
let growable = auto_chunk_dims(&[0, 128], 8);
|
||
assert!(growable[0] >= 1 && bytes(&growable, 8) <= AUTO_CHUNK_TARGET_BYTES);
|
||
// Explicit dimensions always win.
|
||
let explicit = ChunkOptions {
|
||
chunk_dims: Some(vec![10, 10]),
|
||
..Default::default()
|
||
};
|
||
assert_eq!(explicit.resolve_chunk_dims_for(&[4096, 2048], 8), [10, 10]);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_deflate() {
|
||
// Auto-shuffle is applied before compression by default (matches h5py).
|
||
let options = ChunkOptions {
|
||
deflate_level: Some(6),
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(8).unwrap();
|
||
assert_eq!(pl.filters.len(), 2);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_SHUFFLE);
|
||
assert_eq!(pl.filters[1].filter_id, FILTER_DEFLATE);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_deflate_no_shuffle() {
|
||
// Users can opt out of auto-shuffle with no_shuffle = true.
|
||
let options = ChunkOptions {
|
||
deflate_level: Some(6),
|
||
no_shuffle: true,
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(8).unwrap();
|
||
assert_eq!(pl.filters.len(), 1);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_DEFLATE);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_lz4() {
|
||
// Auto-shuffle before LZ4.
|
||
let options = ChunkOptions {
|
||
lz4: true,
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(8).unwrap();
|
||
assert_eq!(pl.filters.len(), 2);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_SHUFFLE);
|
||
assert_eq!(pl.filters[1].filter_id, FILTER_LZ4);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_zstd() {
|
||
// Auto-shuffle before Zstd.
|
||
let options = ChunkOptions {
|
||
zstd_level: Some(3),
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(8).unwrap();
|
||
assert_eq!(pl.filters.len(), 2);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_SHUFFLE);
|
||
assert_eq!(pl.filters[1].filter_id, FILTER_ZSTD);
|
||
assert_eq!(pl.filters[1].client_data, vec![3]);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_lzf() {
|
||
let options = ChunkOptions {
|
||
plugin: Some(PluginFilter::Lzf),
|
||
..Default::default()
|
||
};
|
||
assert!(options.is_chunked());
|
||
let pl = options.build_pipeline_for_chunk(8, 800).unwrap();
|
||
assert_eq!(pl.filters.len(), 2);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_SHUFFLE);
|
||
assert_eq!(pl.filters[1].filter_id, FILTER_LZF);
|
||
assert_eq!(pl.filters[1].client_data, vec![4, 0x0105, 800]);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_bitshuffle_has_no_auto_shuffle() {
|
||
let options = ChunkOptions {
|
||
plugin: Some(PluginFilter::Bitshuffle {
|
||
block_size: 0,
|
||
compression: BitshuffleCompression::Zstd { level: 5 },
|
||
}),
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(4).unwrap();
|
||
assert_eq!(pl.filters.len(), 1);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_BITSHUFFLE);
|
||
assert_eq!(pl.filters[0].client_data, vec![0, 4, 4, 0, 3, 5]);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_zstd_priority_over_deflate() {
|
||
let options = ChunkOptions {
|
||
deflate_level: Some(6),
|
||
zstd_level: Some(3),
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(8).unwrap();
|
||
// shuffle + zstd (deflate is ignored when zstd wins priority)
|
||
assert_eq!(pl.filters.len(), 2);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_SHUFFLE);
|
||
assert_eq!(pl.filters[1].filter_id, FILTER_ZSTD);
|
||
}
|
||
|
||
#[test]
|
||
fn chunk_options_pipeline_shuffle_deflate_fletcher32() {
|
||
let options = ChunkOptions {
|
||
deflate_level: Some(6),
|
||
shuffle: true,
|
||
fletcher32: true,
|
||
..Default::default()
|
||
};
|
||
let pl = options.build_pipeline(8).unwrap();
|
||
assert_eq!(pl.filters.len(), 3);
|
||
assert_eq!(pl.filters[0].filter_id, FILTER_SHUFFLE);
|
||
assert_eq!(pl.filters[1].filter_id, FILTER_DEFLATE);
|
||
assert_eq!(pl.filters[2].filter_id, FILTER_FLETCHER32);
|
||
}
|
||
|
||
#[test]
|
||
fn serialize_v4_single_chunk_no_filters_roundtrip() {
|
||
let msg = serialize_v4_single_chunk(&[20], 0x1000, None, None, 8, 8, 4);
|
||
let layout = DataLayout::parse(&msg, 8, 8).unwrap();
|
||
match layout {
|
||
DataLayout::Chunked {
|
||
chunk_dimensions,
|
||
btree_address,
|
||
version,
|
||
chunk_index_type,
|
||
single_chunk_filtered_size,
|
||
single_chunk_filter_mask,
|
||
..
|
||
} => {
|
||
assert_eq!(version, 4);
|
||
assert_eq!(chunk_index_type, Some(1));
|
||
assert_eq!(chunk_dimensions, vec![20, 8]);
|
||
assert_eq!(btree_address, Some(0x1000));
|
||
assert_eq!(single_chunk_filtered_size, None);
|
||
assert_eq!(single_chunk_filter_mask, None);
|
||
}
|
||
_ => panic!("expected chunked layout"),
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn serialize_v4_single_chunk_with_filters_roundtrip() {
|
||
let msg = serialize_v4_single_chunk(&[100], 0x2000, Some(500), Some(0), 8, 8, 4);
|
||
let layout = DataLayout::parse(&msg, 8, 8).unwrap();
|
||
match layout {
|
||
DataLayout::Chunked {
|
||
btree_address,
|
||
single_chunk_filtered_size,
|
||
single_chunk_filter_mask,
|
||
..
|
||
} => {
|
||
assert_eq!(btree_address, Some(0x2000));
|
||
assert_eq!(single_chunk_filtered_size, Some(500));
|
||
assert_eq!(single_chunk_filter_mask, Some(0));
|
||
}
|
||
_ => panic!("expected chunked layout"),
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn serialize_v4_fixed_array_roundtrip() {
|
||
let msg = serialize_v4_fixed_array(&[20], 0x3000, 8, 8, 4, 4);
|
||
let layout = DataLayout::parse(&msg, 8, 8).unwrap();
|
||
match layout {
|
||
DataLayout::Chunked {
|
||
version,
|
||
chunk_index_type,
|
||
btree_address,
|
||
chunk_dimensions,
|
||
..
|
||
} => {
|
||
assert_eq!(version, 4);
|
||
assert_eq!(chunk_index_type, Some(3));
|
||
assert_eq!(btree_address, Some(0x3000));
|
||
assert_eq!(chunk_dimensions, vec![20, 8]);
|
||
}
|
||
_ => panic!("expected chunked layout"),
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn build_fixed_array_valid_structure() {
|
||
let chunks = vec![
|
||
WrittenChunk {
|
||
address: 0x1000,
|
||
compressed_size: 160,
|
||
raw_size: 160,
|
||
filter_mask: 0,
|
||
},
|
||
WrittenChunk {
|
||
address: 0x10A0,
|
||
compressed_size: 160,
|
||
raw_size: 160,
|
||
filter_mask: 0,
|
||
},
|
||
];
|
||
let slots: Vec<_> = chunks.into_iter().map(Some).collect();
|
||
let fa = build_fixed_array_at(&slots, 8, 8, false, 0x2000);
|
||
// Should start with FAHD
|
||
assert_eq!(&fa[0..4], b"FAHD");
|
||
// FAHD size = 4+1+1+1+1+8+8+4 = 28
|
||
// FADB starts at offset 28
|
||
assert_eq!(&fa[28..32], b"FADB");
|
||
}
|
||
|
||
// ---- Chunks of 4 GiB or more (layout message version 5) ----
|
||
|
||
/// 2^29 + 1 `f64`: 4 GiB + 8 bytes, the smallest `f64` chunk past
|
||
/// `u32::MAX`.
|
||
const HUGE_DIM: u64 = (1 << 29) + 1;
|
||
const HUGE_BYTES: u64 = HUGE_DIM * 8;
|
||
|
||
fn huge_chunk(address: u64, compressed_size: u64) -> WrittenChunk {
|
||
WrittenChunk {
|
||
address,
|
||
compressed_size,
|
||
raw_size: HUGE_BYTES,
|
||
filter_mask: 0,
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn chunks_past_u32_max_take_layout_version_5() {
|
||
assert_eq!(layout_version_for(0), 4);
|
||
assert_eq!(layout_version_for(u64::from(u32::MAX)), 4);
|
||
assert_eq!(layout_version_for(u64::from(u32::MAX) + 1), 5);
|
||
assert_eq!(layout_version_for(HUGE_BYTES), 5);
|
||
}
|
||
|
||
/// A chunk of 4 GiB or more never goes into a version-1 B-tree (its key
|
||
/// holds a 32-bit size): under a 1.8 low bound it still takes layout
|
||
/// version 5, and a high bound below 2.0 refuses it, as in libhdf5.
|
||
#[test]
|
||
fn huge_chunks_ignore_the_v18_low_bound_and_need_v200() {
|
||
// One filtered chunk; the stored bytes stand in for its compression.
|
||
let pre = PrecompressedChunks {
|
||
chunks: vec![(HUGE_BYTES, vec![0u8; 16], 0)],
|
||
has_filters: true,
|
||
element_size: 8,
|
||
shape: vec![HUGE_DIM],
|
||
chunk_dims: vec![HUGE_DIM],
|
||
pipeline_message: None,
|
||
};
|
||
let r = build_chunked_data_from_precompressed_libver(
|
||
&pre,
|
||
4096,
|
||
None,
|
||
LibVer::V18,
|
||
LibVer::Latest,
|
||
)
|
||
.unwrap();
|
||
assert_eq!(r.layout_message[0], 5, "layout message version");
|
||
for high in [LibVer::V18, LibVer::V114] {
|
||
match build_chunked_data_from_precompressed_libver(&pre, 4096, None, LibVer::V18, high)
|
||
{
|
||
Err(FormatError::LibverBound { needs, .. }) => assert_eq!(needs, LibVer::V200),
|
||
other => panic!("high bound {high}: {:?}", other.map(|r| r.layout_message)),
|
||
}
|
||
}
|
||
}
|
||
|
||
/// `H5D_FARRAY_FILT_COMPUTE_CHUNK_SIZE_LEN` (and its EA and v2 B-tree
|
||
/// twins) in libhdf5 2.2.0: one byte more than the chunk size needs up to
|
||
/// layout version 4, the size of lengths from version 5.
|
||
#[test]
|
||
fn chunk_size_len_follows_libhdf5() {
|
||
assert_eq!(chunk_size_len(1, 4, 8), 2);
|
||
assert_eq!(chunk_size_len(160, 4, 8), 2);
|
||
assert_eq!(chunk_size_len(255, 4, 8), 2);
|
||
assert_eq!(chunk_size_len(256, 4, 8), 3);
|
||
assert_eq!(chunk_size_len(u64::from(u32::MAX), 4, 8), 5);
|
||
assert_eq!(chunk_size_len(1 << 32, 4, 8), 6);
|
||
assert_eq!(chunk_size_len(u64::MAX, 4, 8), 8);
|
||
assert_eq!(chunk_size_len(HUGE_BYTES, 5, 8), 8);
|
||
assert_eq!(chunk_size_len(160, 5, 8), 8);
|
||
assert_eq!(chunk_size_len(HUGE_BYTES, 5, 4), 4);
|
||
let slots = [Some(huge_chunk(0x1000, 20_000)), None];
|
||
assert_eq!(filtered_chunk_size_len(&slots, 8), 8);
|
||
let small = [Some(WrittenChunk {
|
||
raw_size: 160,
|
||
..huge_chunk(0x1000, 100)
|
||
})];
|
||
assert_eq!(filtered_chunk_size_len(&small, 8), 2);
|
||
}
|
||
|
||
/// A 4 GiB + 8 byte chunk: every layout message is version 5 with the
|
||
/// dimensions libhdf5 writes (4 bytes each: 0x20000001 and 8), and a
|
||
/// filtered index element stores the chunk's size in 8 bytes.
|
||
#[test]
|
||
fn huge_chunk_layout_messages_and_index_elements() {
|
||
let dims = [HUGE_DIM as u32];
|
||
let parsed = |msg: &[u8]| {
|
||
assert_eq!(msg[0], 5, "layout message version");
|
||
// Class chunked, then (after the flags) 2 dimensions of 4 bytes.
|
||
assert_eq!(&msg[1..2], &[2]);
|
||
assert_eq!(&msg[3..5], &[2, 4]);
|
||
assert_eq!(&msg[5..13], &[1, 0, 0, 0x20, 8, 0, 0, 0]);
|
||
match DataLayout::parse(msg, 8, 8).unwrap() {
|
||
DataLayout::Chunked {
|
||
chunk_dimensions,
|
||
chunk_index_type,
|
||
..
|
||
} => {
|
||
assert_eq!(chunk_dimensions, vec![HUGE_DIM as u32, 8]);
|
||
chunk_index_type.unwrap()
|
||
}
|
||
other => panic!("{other:?}"),
|
||
}
|
||
};
|
||
let single = serialize_v4_single_chunk(&dims, 0x800, Some(20_000), Some(0), 8, 8, 5);
|
||
assert_eq!(parsed(&single), 1);
|
||
// Filtered size (8 bytes), filter mask, address.
|
||
assert_eq!(single.len(), 13 + 1 + 8 + 4 + 8);
|
||
assert_eq!(
|
||
parsed(&serialize_v4_fixed_array(&dims, 0x800, 8, 8, 10, 5)),
|
||
3
|
||
);
|
||
let ea = ea_writer::serialize_v4_extensible_array(&dims, 0x800, 8, 8, 5);
|
||
assert_eq!(parsed(&ea), 4);
|
||
assert_eq!(
|
||
parsed(&serialize_v4_btree_v2(&dims, 0x800, 8, 8, 2048, 5)),
|
||
5
|
||
);
|
||
|
||
let slots = [Some(huge_chunk(0x1000, 20_000)), None];
|
||
// FAHD: element size (address 8 + size 8 + mask 4) at byte 6.
|
||
let fa = build_fixed_array_at(&slots, 8, 8, true, 0x2000);
|
||
assert_eq!(fa[6], 20);
|
||
// FADB element 0: address, then the stored size in 8 bytes.
|
||
let fadb = 28;
|
||
let prefix = 4 + 1 + 1 + 8;
|
||
assert_eq!(
|
||
&fa[fadb + prefix..fadb + prefix + 8],
|
||
&0x1000u64.to_le_bytes()
|
||
);
|
||
assert_eq!(
|
||
&fa[fadb + prefix + 8..fadb + prefix + 16],
|
||
&20_000u64.to_le_bytes()
|
||
);
|
||
// AEHD: element size at byte 6 as well.
|
||
let ea = ea_writer::build_extensible_array_at(&slots, 8, 8, true, 0x2000);
|
||
assert_eq!(&ea[..4], b"EAHD");
|
||
assert_eq!(ea[6], 20);
|
||
// BTHD: record size (element + 8-byte scaled offset) at bytes 10-11.
|
||
let chunk = huge_chunk(0x1000, 20_000);
|
||
let (bt, _) =
|
||
build_btree_v2_chunk_index_at(1, &[(vec![0], &chunk)], 8, 8, true, 0x2000).unwrap();
|
||
assert_eq!(&bt[..4], b"BTHD");
|
||
assert_eq!(u16::from_le_bytes([bt[10], bt[11]]), 28);
|
||
}
|
||
|
||
#[test]
|
||
fn huge_chunks_refused_where_unsupported() {
|
||
// A chunk dimension of 2^32 or more.
|
||
assert!(matches!(
|
||
checked_chunk_dims(&[1 << 32], 1),
|
||
Err(FormatError::InvalidChunkDimensions(m)) if m.contains("2^32")
|
||
));
|
||
// A chunk size that overflows u64.
|
||
assert!(matches!(
|
||
checked_chunk_dims(&[u32::MAX.into(), u32::MAX.into(), 2], 8),
|
||
Err(FormatError::Overflow(_))
|
||
));
|
||
assert_eq!(
|
||
checked_chunk_dims(&[HUGE_DIM], 8).unwrap(),
|
||
(vec![HUGE_DIM as u32], HUGE_BYTES)
|
||
);
|
||
// Filters that cannot take a chunk that large.
|
||
for plugin in [
|
||
PluginFilter::Lzf,
|
||
PluginFilter::Bzip2 { level: 9 },
|
||
PluginFilter::Blosc {
|
||
codec: BloscCodec::Lz4,
|
||
level: 5,
|
||
shuffle: BloscShuffle::Byte,
|
||
},
|
||
PluginFilter::Bitshuffle {
|
||
block_size: 0,
|
||
compression: BitshuffleCompression::Lz4,
|
||
},
|
||
] {
|
||
let options = ChunkOptions {
|
||
plugin: Some(plugin),
|
||
..Default::default()
|
||
};
|
||
assert!(matches!(
|
||
check_huge_chunk_filters(&options, HUGE_BYTES),
|
||
Err(FormatError::FilterError(m)) if m.contains("4 GiB")
|
||
));
|
||
}
|
||
for options in [
|
||
ChunkOptions {
|
||
deflate_level: Some(6),
|
||
shuffle: true,
|
||
fletcher32: true,
|
||
..Default::default()
|
||
},
|
||
ChunkOptions {
|
||
zstd_level: Some(3),
|
||
..Default::default()
|
||
},
|
||
ChunkOptions {
|
||
lz4: true,
|
||
..Default::default()
|
||
},
|
||
] {
|
||
check_huge_chunk_filters(&options, HUGE_BYTES).unwrap();
|
||
}
|
||
}
|
||
|
||
/// Chunks are extracted row by row: the result is the element-by-element
|
||
/// split, edge padding included.
|
||
#[test]
|
||
fn extract_chunk_matches_elementwise_split() {
|
||
let shape = [5u64, 7, 3];
|
||
let chunks = [2u64, 3, 2];
|
||
let data: Vec<u8> = (0..5 * 7 * 3 * 2).map(|i| i as u8).collect();
|
||
let n = chunk_count(&shape, &chunks);
|
||
assert_eq!(n, 3 * 3 * 2);
|
||
for i in 0..n {
|
||
let (offsets, chunk) = extract_chunk(&data, &shape, &chunks, 2, i);
|
||
assert_eq!(chunk.len(), 2 * 3 * 2 * 2);
|
||
for (e, pair) in chunk.as_chunks::<2>().0.iter().enumerate() {
|
||
let c = [e / 6, (e / 2) % 3, e % 2];
|
||
let g: Vec<u64> = (0..3).map(|d| offsets[d] + c[d] as u64).collect();
|
||
let expect = if (0..3).all(|d| g[d] < shape[d]) {
|
||
let flat = ((g[0] * 7 + g[1]) * 3 + g[2]) as usize;
|
||
[data[2 * flat], data[2 * flat + 1]]
|
||
} else {
|
||
[0, 0]
|
||
};
|
||
assert_eq!(*pair, expect, "chunk {i} element {e}");
|
||
}
|
||
}
|
||
// Data shorter than the shape: the missing elements stay zero.
|
||
let (_, chunk) = extract_chunk(&data[..5], &[4], &[4], 2, 0);
|
||
assert_eq!(chunk, [0, 1, 2, 3, 0, 0, 0, 0]);
|
||
}
|
||
|
||
// ---- Extensible Array tests ----
|
||
|
||
#[test]
|
||
fn serialize_v4_extensible_array_roundtrip() {
|
||
let msg = ea_writer::serialize_v4_extensible_array(&[10], 0x4000, 8, 8, 4);
|
||
let layout = DataLayout::parse(&msg, 8, 8).unwrap();
|
||
match layout {
|
||
DataLayout::Chunked {
|
||
version,
|
||
chunk_index_type,
|
||
btree_address,
|
||
chunk_dimensions,
|
||
..
|
||
} => {
|
||
assert_eq!(version, 4);
|
||
assert_eq!(chunk_index_type, Some(4));
|
||
assert_eq!(btree_address, Some(0x4000));
|
||
assert_eq!(chunk_dimensions, vec![10, 8]);
|
||
}
|
||
_ => panic!("expected chunked layout"),
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn build_extensible_array_valid_structure() {
|
||
let chunks = vec![
|
||
WrittenChunk {
|
||
address: 0x1000,
|
||
compressed_size: 80,
|
||
raw_size: 80,
|
||
filter_mask: 0,
|
||
},
|
||
WrittenChunk {
|
||
address: 0x1050,
|
||
compressed_size: 80,
|
||
raw_size: 80,
|
||
filter_mask: 0,
|
||
},
|
||
];
|
||
let slots: Vec<_> = chunks.into_iter().map(Some).collect();
|
||
let ea = ea_writer::build_extensible_array_at(&slots, 8, 8, false, 0x2000);
|
||
assert_eq!(&ea[0..4], b"EAHD");
|
||
// Find EAIB after EAHD: 12 fixed + 6*8 stats + 8 addr + 4 checksum = 72
|
||
let aehd_size = 4 + 1 + 1 + 1 + 1 + 1 + 1 + 1 + 1 + 6 * 8 + 8 + 4;
|
||
assert_eq!(&ea[aehd_size..aehd_size + 4], b"EAIB");
|
||
}
|
||
|
||
/// Helper: roundtrip with EA (maxshape)
|
||
fn roundtrip_ea(
|
||
values: &[f64],
|
||
shape: &[u64],
|
||
chunk_dims: &[u64],
|
||
maxshape: &[u64],
|
||
) -> Vec<f64> {
|
||
let raw = f64_to_bytes(values);
|
||
let base_address = 0x1000u64;
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(chunk_dims.to_vec()),
|
||
..Default::default()
|
||
};
|
||
let result = build_chunked_data_at_ext(
|
||
&raw,
|
||
shape,
|
||
chunk_dims,
|
||
8,
|
||
&options,
|
||
base_address,
|
||
Some(maxshape),
|
||
)
|
||
.unwrap();
|
||
|
||
let file_size = base_address as usize + result.data_bytes.len();
|
||
let mut file_data = vec![0u8; file_size];
|
||
file_data[base_address as usize..].copy_from_slice(&result.data_bytes);
|
||
|
||
let layout = DataLayout::parse(&result.layout_message, 8, 8).unwrap();
|
||
// Verify it uses EA index
|
||
match &layout {
|
||
DataLayout::Chunked {
|
||
chunk_index_type, ..
|
||
} => {
|
||
assert_eq!(*chunk_index_type, Some(4), "expected EA index type");
|
||
}
|
||
_ => panic!("expected chunked layout"),
|
||
}
|
||
|
||
let dataspace = Dataspace {
|
||
space_type: DataspaceType::Simple,
|
||
rank: shape.len() as u8,
|
||
dimensions: shape.to_vec(),
|
||
max_dimensions: Some(maxshape.to_vec()),
|
||
};
|
||
let datatype = make_f64_type();
|
||
|
||
let output =
|
||
read_chunked_data(&file_data, &layout, &dataspace, &datatype, None, 8, 8).unwrap();
|
||
|
||
bytes_to_f64(&output)
|
||
}
|
||
|
||
/// Every chunk index the writer builds records each chunk's real filter
|
||
/// mask: LZF output no smaller than the chunk is skipped (bit 1, behind
|
||
/// shuffle) and the chunk stored shuffled only; compressible chunks keep
|
||
/// mask 0. The data reads back through both kinds of chunk.
|
||
#[cfg(feature = "lzf")]
|
||
#[test]
|
||
fn skipped_lzf_chunks_are_masked_in_every_index() {
|
||
let c = 64usize;
|
||
// Chunks alternate: random bytes (LZF cannot shrink them), then 7s.
|
||
let mut state = 0x1234_5678_u64;
|
||
let data: Vec<f64> = (0..4 * c)
|
||
.map(|i| {
|
||
if (i / c).is_multiple_of(2) {
|
||
state ^= state << 13;
|
||
state ^= state >> 7;
|
||
state ^= state << 17;
|
||
f64::from_bits(state)
|
||
} else {
|
||
7.0
|
||
}
|
||
})
|
||
.collect();
|
||
let raw = f64_to_bytes(&data);
|
||
let options = ChunkOptions {
|
||
plugin: Some(PluginFilter::Lzf),
|
||
..Default::default()
|
||
};
|
||
let c64 = c as u64;
|
||
#[allow(clippy::type_complexity)]
|
||
let cases: [(&[u64], &[u64], Option<&[u64]>, u8, &[u32]); 4] = [
|
||
(&[c64], &[c64], None, 1, &[2]),
|
||
(&[4 * c64], &[c64], None, 3, &[2, 0, 2, 0]),
|
||
(&[4 * c64], &[c64], Some(&[u64::MAX]), 4, &[2, 0, 2, 0]),
|
||
(
|
||
&[2, 2 * c64],
|
||
&[1, c64],
|
||
Some(&[u64::MAX, u64::MAX]),
|
||
5,
|
||
&[2, 0, 2, 0],
|
||
),
|
||
];
|
||
let base = 0x1000u64;
|
||
for (shape, chunks, maxshape, index_type, want_masks) in cases {
|
||
let n: u64 = shape.iter().product();
|
||
let raw = &raw[..n as usize * 8];
|
||
let result =
|
||
build_chunked_data_at_ext(raw, shape, chunks, 8, &options, base, maxshape).unwrap();
|
||
let mut file = vec![0u8; base as usize];
|
||
file.extend_from_slice(&result.data_bytes);
|
||
let layout = DataLayout::parse(&result.layout_message, 8, 8).unwrap();
|
||
assert!(
|
||
matches!(&layout, DataLayout::Chunked { chunk_index_type, .. }
|
||
if *chunk_index_type == Some(index_type)),
|
||
"{layout:?}"
|
||
);
|
||
let dataspace = Dataspace {
|
||
space_type: DataspaceType::Simple,
|
||
rank: shape.len() as u8,
|
||
dimensions: shape.to_vec(),
|
||
max_dimensions: maxshape.map(<[u64]>::to_vec),
|
||
};
|
||
let (mut infos, _) =
|
||
crate::chunked_read::list_chunks(&file, &layout, &dataspace, 8, 8, 8).unwrap();
|
||
infos.sort_by(|a, b| a.offsets.cmp(&b.offsets));
|
||
let masks: Vec<u32> = infos.iter().map(|i| i.filter_mask).collect();
|
||
assert_eq!(masks, want_masks, "index type {index_type}");
|
||
for info in &infos {
|
||
// Skipped chunks are stored at the chunk's size (shuffled).
|
||
assert_eq!(
|
||
info.chunk_size == (c * 8) as u64,
|
||
info.filter_mask != 0,
|
||
"{info:?}"
|
||
);
|
||
}
|
||
let pipeline = crate::filter_pipeline::FilterPipeline::parse(
|
||
result.pipeline_message.as_ref().unwrap(),
|
||
)
|
||
.unwrap();
|
||
let out = read_chunked_data(
|
||
&file,
|
||
&layout,
|
||
&dataspace,
|
||
&make_f64_type(),
|
||
Some(&pipeline),
|
||
8,
|
||
8,
|
||
)
|
||
.unwrap();
|
||
assert_eq!(out, raw, "index type {index_type}");
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn ea_roundtrip_1d_inline_only() {
|
||
let values: Vec<f64> = (0..10).map(|i| i as f64).collect();
|
||
let result = roundtrip_ea(&values, &[10], &[10], &[u64::MAX]);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[test]
|
||
fn ea_roundtrip_1d_multi_chunks() {
|
||
let values: Vec<f64> = (0..20).map(|i| i as f64).collect();
|
||
let result = roundtrip_ea(&values, &[20], &[5], &[u64::MAX]);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[test]
|
||
fn ea_roundtrip_1d_many_chunks() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let result = roundtrip_ea(&values, &[100], &[10], &[u64::MAX]);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
// ---- h5py round-trip tests for chunked writes ----
|
||
|
||
/// The Python interpreter to drive interop checks with.
|
||
///
|
||
/// `CLAWHDF5_PYTHON` lets these run against a virtualenv holding h5py,
|
||
/// which on a PEP 668 "externally managed" system is the only place it
|
||
/// can be installed. Without it the suite silently skips, and a silent
|
||
/// skip here is how a datatype bug once reached a release.
|
||
#[cfg(feature = "std")]
|
||
fn python() -> String {
|
||
std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string())
|
||
}
|
||
|
||
#[cfg(feature = "std")]
|
||
fn h5py_available() -> bool {
|
||
std::process::Command::new(python())
|
||
.args(["-c", "import h5py"])
|
||
.output()
|
||
.map(|o| o.status.success())
|
||
.unwrap_or(false)
|
||
}
|
||
|
||
#[cfg(feature = "std")]
|
||
fn h5py_run(script: &str) -> String {
|
||
if !h5py_available() {
|
||
panic!("h5py not installed — skipping interop test");
|
||
}
|
||
let o = std::process::Command::new(python())
|
||
.args(["-c", script])
|
||
.output()
|
||
.expect("python interpreter");
|
||
if !o.status.success() {
|
||
panic!("h5py: {}", String::from_utf8_lossy(&o.stderr));
|
||
}
|
||
String::from_utf8(o.stdout).unwrap().trim().to_string()
|
||
}
|
||
|
||
#[cfg(feature = "std")]
|
||
#[test]
|
||
#[ignore = "requires Python h5py module"]
|
||
fn h5py_reads_multiple_chunked_datasets() {
|
||
use crate::file_writer::FileWriter;
|
||
let mut fw = FileWriter::new();
|
||
let data1: Vec<f64> = (0..50).map(|i| i as f64).collect();
|
||
let data2: Vec<f64> = (0..30).map(|i| (i * 10) as f64).collect();
|
||
fw.create_dataset("a")
|
||
.with_f64_data(&data1)
|
||
.with_shape(&[50])
|
||
.with_chunks(&[25]);
|
||
fw.create_dataset("b")
|
||
.with_f64_data(&data2)
|
||
.with_shape(&[30])
|
||
.with_chunks(&[10]);
|
||
let bytes = fw.finish().unwrap();
|
||
let path = std::env::temp_dir().join("clawhdf5_chunked_multi.h5");
|
||
std::fs::write(&path, &bytes).unwrap();
|
||
let script = format!(
|
||
"import h5py,json; f=h5py.File('{}','r'); print(json.dumps({{'a':f['a'][:].tolist(),'b':f['b'][:].tolist()}}))",
|
||
path.display()
|
||
);
|
||
let v: serde_json::Value = serde_json::from_str(&h5py_run(&script)).unwrap();
|
||
let va: Vec<f64> = serde_json::from_value(v["a"].clone()).unwrap();
|
||
let vb: Vec<f64> = serde_json::from_value(v["b"].clone()).unwrap();
|
||
assert_eq!(va, data1);
|
||
assert_eq!(vb, data2);
|
||
}
|
||
|
||
#[cfg(feature = "std")]
|
||
#[test]
|
||
#[ignore = "requires Python h5py module"]
|
||
fn h5py_reads_chunked_with_attrs() {
|
||
use crate::file_writer::{AttrValue, FileWriter};
|
||
let mut fw = FileWriter::new();
|
||
let data: Vec<f64> = (0..50).map(|i| i as f64).collect();
|
||
fw.create_dataset("data")
|
||
.with_f64_data(&data)
|
||
.with_shape(&[50])
|
||
.with_chunks(&[25])
|
||
.set_attr("units", AttrValue::String("meters".to_string()));
|
||
let bytes = fw.finish().unwrap();
|
||
let path = std::env::temp_dir().join("clawhdf5_chunked_attrs.h5");
|
||
std::fs::write(&path, &bytes).unwrap();
|
||
let script = format!(
|
||
"import h5py,json; f=h5py.File('{}','r'); d=f['data']; print(json.dumps({{'values':d[:].tolist(),'units':d.attrs['units'].decode() if isinstance(d.attrs['units'],bytes) else str(d.attrs['units'])}}))",
|
||
path.display()
|
||
);
|
||
let v: serde_json::Value = serde_json::from_str(&h5py_run(&script)).unwrap();
|
||
let values: Vec<f64> = serde_json::from_value(v["values"].clone()).unwrap();
|
||
assert_eq!(values, data);
|
||
assert_eq!(v["units"], serde_json::json!("meters"));
|
||
}
|
||
|
||
// --- LZ4 chunked roundtrip tests ---
|
||
|
||
#[cfg(feature = "lz4")]
|
||
#[test]
|
||
fn roundtrip_1d_single_chunk_lz4() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![100]),
|
||
lz4: true,
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[100], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[cfg(feature = "lz4")]
|
||
#[test]
|
||
fn roundtrip_1d_multi_chunk_lz4() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![20]),
|
||
lz4: true,
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[20], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
// --- Zstd chunked roundtrip tests ---
|
||
|
||
#[cfg(feature = "zstd")]
|
||
#[test]
|
||
fn roundtrip_1d_single_chunk_zstd() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![100]),
|
||
zstd_level: Some(3),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[100], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[cfg(feature = "zstd")]
|
||
#[test]
|
||
fn roundtrip_1d_multi_chunk_zstd() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![20]),
|
||
zstd_level: Some(1),
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[20], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
|
||
#[cfg(feature = "zstd")]
|
||
#[test]
|
||
fn roundtrip_1d_shuffle_zstd() {
|
||
let values: Vec<f64> = (0..100).map(|i| i as f64).collect();
|
||
let options = ChunkOptions {
|
||
chunk_dims: Some(vec![50]),
|
||
zstd_level: Some(3),
|
||
shuffle: true,
|
||
..Default::default()
|
||
};
|
||
let result = roundtrip_chunked(&values, &[100], &[50], &options);
|
||
assert_eq!(result, values);
|
||
}
|
||
}
|