Chunks of 4 GiB or more: read in every index, write as libhdf5 2.x does #29

Open
osobh wants to merge 10 commits from feat/huge-chunks into main
8 changed files with 788 additions and 175 deletions
Showing only changes of commit 5a20cf04e8 - Show all commits
+32 -3
View File
@@ -327,24 +327,53 @@ fn inflate_bounded(data: &[u8], size_hint: usize, limit: usize) -> Result<Vec<u8
}
}
/// Largest compression output reserved at its worst-case size up front.
const DEFLATE_EXACT_BOUND: usize = 64 << 20;
/// Compress data using flate2 (zlib-ng, zlib-rs or miniz_oxide; see module docs).
pub(crate) fn flate2_compress(data: &[u8], level: u32) -> Result<Vec<u8>, String> {
use flate2::{Compress, Compression, FlushCompress, Status};
// zlib's compressBound, plus the zlib header and trailer.
let bound = data.len() + (data.len() >> 12) + (data.len() >> 14) + (data.len() >> 25) + 13 + 6;
// flate2's Rust backends (zlib-rs, miniz_oxide) zero the whole spare
// capacity on each call, so a large input's worst-case bound would be
// memory held for nothing (4 GiB for a 4 GiB chunk that deflates to a
// few MiB): past 64 MiB the output starts at 1/16 of the bound and
// doubles as needed.
let first = if bound <= DEFLATE_EXACT_BOUND {
bound
} else {
bound / 16
};
let mut out = Vec::new();
out.try_reserve_exact(bound)
out.try_reserve_exact(first)
.map_err(|e| format!("deflate: cannot allocate output: {e}"))?;
let mut deflater = Compress::new(Compression::new(level), true);
loop {
let (in_before, out_before) = (deflater.total_in(), deflater.total_out());
let rest = &data[in_before as usize..];
// zlib takes at most u32::MAX input bytes per call, and `Finish`
// ends the stream after the bytes it took: input of 4 GiB or more
// was cut at 4 GiB - 1. Finish only once the rest fits one call.
let flush = if rest.len() > u32::MAX as usize {
FlushCompress::None
} else {
FlushCompress::Finish
};
let status = deflater
.compress_vec(&data[in_before as usize..], &mut out, FlushCompress::Finish)
.compress_vec(rest, &mut out, flush)
.map_err(|e| format!("deflate: {e}"))?;
match status {
Status::StreamEnd => return Ok(out),
Status::StreamEnd => {
if out.capacity() - out.len() > DEFLATE_EXACT_BOUND {
out.shrink_to_fit();
}
return Ok(out);
}
// Out of room (the bound makes it unreachable below
// `DEFLATE_EXACT_BOUND`): grow rather than fail.
Status::Ok | Status::BufError if out.len() == out.capacity() => out
.try_reserve(out.capacity().max(4096))
.map_err(|e| format!("deflate: cannot allocate output: {e}"))?,
+452 -128
View File
@@ -402,86 +402,83 @@ pub fn split_into_chunks(
chunk_dims: &[u64],
element_size: usize,
) -> Vec<(Vec<u64>, Vec<u8>)> {
let rank = shape.len();
if rank == 0 {
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()
}
// Compute number of chunks per dimension
let mut num_chunks_per_dim = Vec::with_capacity(rank);
for d in 0..rank {
num_chunks_per_dim.push(shape[d].div_ceil(chunk_dims[d]));
/// 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 total_chunks: u64 = num_chunks_per_dim.iter().product();
let chunk_elements: usize = chunk_dims.iter().map(|&d| saturating_usize(d)).product();
let mut chunk = vec![0u8; chunk_elements.saturating_mul(element_size)];
// Dataset strides (row-major)
// 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];
for i in (0..rank.saturating_sub(1)).rev() {
ds_strides[i] = ds_strides[i + 1] * saturating_usize(shape[i + 1]);
}
// Chunk strides
let mut chunk_strides = vec![1usize; rank];
for i in (0..rank.saturating_sub(1)).rev() {
chunk_strides[i] = chunk_strides[i + 1] * saturating_usize(chunk_dims[i + 1]);
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 chunk_total_elements: usize = chunk_dims.iter().map(|&d| saturating_usize(d)).product();
let mut result = Vec::with_capacity(saturating_usize(total_chunks));
for linear_idx in 0..total_chunks {
// Convert linear index to chunk grid coordinates
let mut chunk_grid_coords = vec![0u64; rank];
let mut remaining = linear_idx;
for d in (0..rank).rev() {
chunk_grid_coords[d] = remaining % num_chunks_per_dim[d];
remaining /= num_chunks_per_dim[d];
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;
}
// Chunk offset in dataset space
let offsets: Vec<u64> = (0..rank)
.map(|d| chunk_grid_coords[d] * chunk_dims[d])
.collect();
// Extract chunk data
let mut chunk_bytes = vec![0u8; chunk_total_elements * element_size];
for flat_idx in 0..chunk_total_elements {
let mut remaining_idx = flat_idx;
let mut ds_flat = 0usize;
let mut out_of_bounds = false;
for d in 0..rank {
let coord_in_chunk = remaining_idx / chunk_strides[d];
remaining_idx %= chunk_strides[d];
let global_coord = saturating_usize(offsets[d]) + coord_in_chunk;
if global_coord >= saturating_usize(shape[d]) {
out_of_bounds = true;
break;
}
ds_flat += global_coord * ds_strides[d];
}
if out_of_bounds {
// Zero-filled (already initialized)
continue;
}
let src_start = ds_flat * element_size;
let dst_start = flat_idx * element_size;
if src_start + element_size <= raw_data.len() {
chunk_bytes[dst_start..dst_start + element_size]
.copy_from_slice(&raw_data[src_start..src_start + element_size]);
}
}
result.push((offsets, chunk_bytes));
}
result
}
/// Parallel compression threshold: use rayon when chunk count exceeds this.
@@ -492,46 +489,68 @@ pub fn split_into_chunks(
#[cfg(feature = "parallel")]
const PARALLEL_COMPRESS_THRESHOLD: usize = 2;
/// Compress all chunks, using parallel compression when beneficial, and
/// return each chunk's stored bytes with its filter mask.
/// 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.
///
/// With the `parallel` feature and more than [`PARALLEL_COMPRESS_THRESHOLD`]
/// filtered chunks, compression runs across rayon threads; otherwise it is
/// sequential. Output order matches input order, so per-chunk bytes are
/// identical to the sequential path.
/// 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(
chunks: &[(Vec<u64>, Vec<u8>)],
raw_data: &[u8],
shape: &[u64],
chunk_dims: &[u64],
element_size: usize,
chunk_bytes: u64,
pipeline: &Option<FilterPipeline>,
element_size: u32,
) -> Result<Vec<(Vec<u8>, u32)>, FormatError> {
) -> 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 let Some(pl) = pipeline
&& chunks.len() > PARALLEL_COMPRESS_THRESHOLD
if pipeline.is_some()
&& n > PARALLEL_COMPRESS_THRESHOLD as u64
&& chunk_bytes <= PARALLEL_COMPRESS_MAX_CHUNK_BYTES
{
use rayon::prelude::*;
return chunks
.par_iter()
.map(|(_offsets, chunk_bytes)| compress_chunk_masked(chunk_bytes, pl, element_size))
.collect();
return (0..n).into_par_iter().map(one).collect();
}
}
#[cfg(not(feature = "parallel"))]
let _ = chunk_bytes;
// Sequential fallback
chunks
.iter()
.map(|(_offsets, chunk_bytes)| {
if let Some(pl) = pipeline {
compress_chunk_masked(chunk_bytes, pl, element_size)
} else {
Ok((chunk_bytes.clone(), 0))
}
})
.collect()
(0..n).map(one).collect()
}
/// Build the complete chunked dataset blob (chunk data + index) and return
@@ -552,10 +571,12 @@ pub fn serialize_v4_single_chunk_pub(
filter_mask,
offset_size,
element_size,
4,
)
}
/// Serialize a v4 single chunk layout message.
/// 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,
@@ -563,9 +584,10 @@ fn serialize_v4_single_chunk(
filter_mask: Option<u32>,
offset_size: u8,
element_size: u32,
version: u8,
) -> Vec<u8> {
let mut buf = Vec::new();
buf.push(4); // version
buf.push(version);
buf.push(2); // class = chunked
// flags: bit 0 = unknown meaning in some files, bit 1 = filters for single chunk
@@ -605,8 +627,9 @@ fn serialize_v4_fixed_array(
offset_size: u8,
element_size: u32,
max_bits: u8,
version: u8,
) -> Vec<u8> {
let mut buf = layout_v4_chunked_prefix(chunk_dims, element_size);
let mut buf = layout_v4_chunked_prefix(chunk_dims, element_size, version);
// chunk index type = 3 (Fixed Array)
buf.push(3);
@@ -645,9 +668,9 @@ pub(crate) fn push_v4_chunk_dims(buf: &mut Vec<u8>, chunk_dims: &[u32], element_
}
}
fn layout_v4_chunked_prefix(chunk_dims: &[u32], element_size: u32) -> Vec<u8> {
fn layout_v4_chunked_prefix(chunk_dims: &[u32], element_size: u32, version: u8) -> Vec<u8> {
let mut buf = Vec::new();
buf.push(4); // version
buf.push(version);
buf.push(2); // class = chunked
let flags: u8 = 0x00;
@@ -673,21 +696,53 @@ pub(crate) fn push_addr(buf: &mut Vec<u8>, addr: u64, offset_size: u8) {
/// 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):
/// `1 + ((log2(unfiltered chunk bytes) + 8) / 8)`, capped at 8.
pub(crate) fn filtered_chunk_size_len(slots: &[Option<WrittenChunk>]) -> usize {
/// 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);
let log2_val = if max_raw <= 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 - max_raw.leading_zeros()
63 - chunk_bytes.leading_zeros()
};
(1 + ((log2_val + 8) / 8) as usize).min(8)
(1 + ((log2 + 8) / 8) as usize).min(8)
}
/// Append one chunk index element: the chunk's address, plus its stored size
@@ -734,7 +789,7 @@ pub fn build_fixed_array_at(
let os = offset_size as usize;
let num_elements = slots.len();
let chunk_size_bytes = has_filters.then(|| filtered_chunk_size_len(slots));
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 };
@@ -829,23 +884,23 @@ pub fn precompress_chunks(
element_size: usize,
options: &ChunkOptions,
) -> Result<PrecompressedChunks, FormatError> {
let chunk_bytes = chunk_dims
.iter()
.try_fold(element_size as u64, |acc, &d| acc.checked_mul(d))
.and_then(|b| u32::try_from(b).ok())
.unwrap_or(0);
let pipeline = options.build_pipeline_for_chunk(element_size as u32, chunk_bytes);
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 raw_chunks = split_into_chunks(raw_data, shape, chunk_dims, element_size);
let compressed = compress_all_chunks(&raw_chunks, &pipeline, element_size as u32)?;
let chunks = raw_chunks
.into_iter()
.zip(compressed)
.map(|((_offsets, raw_bytes), (c, mask))| (raw_bytes.len() as u64, c, mask))
.collect();
let chunks = compress_all_chunks(
raw_data,
shape,
chunk_dims,
element_size,
chunk_bytes,
&pipeline,
)?;
Ok(PrecompressedChunks {
chunks,
@@ -857,6 +912,67 @@ pub fn precompress_chunks(
})
}
/// 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
@@ -910,7 +1026,8 @@ pub fn build_chunked_data_from_precompressed_libver(
});
}
let chunk_dims_u32: Vec<u32> = pre.chunk_dims.iter().map(|&d| d as u32).collect();
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() {
@@ -934,6 +1051,7 @@ pub fn build_chunked_data_from_precompressed_libver(
ea_address,
offset_size,
element_size as u32,
version,
)
}
ChunkIndexPlan::SingleChunk => {
@@ -951,6 +1069,7 @@ pub fn build_chunked_data_from_precompressed_libver(
filter_mask,
offset_size,
element_size as u32,
version,
)
}
ChunkIndexPlan::FixedArray(grid, nslots) => {
@@ -976,6 +1095,7 @@ pub fn build_chunked_data_from_precompressed_libver(
offset_size,
element_size as u32,
FA_PAGE_BITS,
version,
)
}
ChunkIndexPlan::BTreeV2 => {
@@ -1000,6 +1120,7 @@ pub fn build_chunked_data_from_precompressed_libver(
offset_size,
element_size as u32,
node_size,
version,
)
}
};
@@ -1237,7 +1358,7 @@ fn build_btree_v2_chunk_index_at(
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)
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)
@@ -1287,8 +1408,9 @@ fn serialize_v4_btree_v2(
offset_size: u8,
element_size: u32,
node_size: u32,
version: u8,
) -> Vec<u8> {
let mut buf = layout_v4_chunked_prefix(chunk_dims, element_size);
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);
@@ -1868,7 +1990,7 @@ mod tests {
#[test]
fn serialize_v4_single_chunk_no_filters_roundtrip() {
let msg = serialize_v4_single_chunk(&[20], 0x1000, None, None, 8, 8);
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 {
@@ -1893,7 +2015,7 @@ mod tests {
#[test]
fn serialize_v4_single_chunk_with_filters_roundtrip() {
let msg = serialize_v4_single_chunk(&[100], 0x2000, Some(500), Some(0), 8, 8);
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 {
@@ -1912,7 +2034,7 @@ mod tests {
#[test]
fn serialize_v4_fixed_array_roundtrip() {
let msg = serialize_v4_fixed_array(&[20], 0x3000, 8, 8, 4);
let msg = serialize_v4_fixed_array(&[20], 0x3000, 8, 8, 4, 4);
let layout = DataLayout::parse(&msg, 8, 8).unwrap();
match layout {
DataLayout::Chunked {
@@ -1956,11 +2078,213 @@ mod tests {
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);
}
/// `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.chunks_exact(2).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);
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 {
+3 -2
View File
@@ -18,9 +18,10 @@ pub(crate) fn serialize_v4_extensible_array(
ea_address: u64,
offset_size: u8,
element_size: u32,
version: u8,
) -> Vec<u8> {
let mut buf = Vec::new();
buf.push(4); // version
buf.push(version);
buf.push(2); // class = chunked
buf.push(0x00); // flags
@@ -87,7 +88,7 @@ pub fn build_extensible_array_at(
ea_base_address: u64,
) -> Vec<u8> {
let os = offset_size as usize;
let chunk_size_bytes = has_filters.then(|| filtered_chunk_size_len(slots));
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 };
let arr_off_size = (MAX_NELMTS_BITS as usize).div_ceil(8);
+60 -26
View File
@@ -343,10 +343,10 @@ pub fn decompress_chunk_exact_with<'s>(
// The bound is the output's size when only size-preserving
// filters (shuffle, Fletcher32) remain to be undone; before
// any other filter (a second deflate, N-Bit, ...) it is only
// a ceiling, and the output starts smaller and grows: the
// inflater writes every byte it reserves, so reserving a
// 4 GiB chunk's bound for a stage a few MiB long would hold
// twice the chunk's memory.
// a ceiling, and the output starts smaller and grows.
// Reserving the bound of a 4 GiB chunk for a stage a few MiB
// long doubled the peak memory of reading it (8.5 GiB, now
// 4.0, for `huge_chunks_filtered.h5`'s double deflate).
let exact = pipeline.filters[..i].iter().enumerate().all(|(j, f)| {
filter_skipped(filter_mask, j)
|| matches!(f.filter_id, FILTER_SHUFFLE | FILTER_FLETCHER32)
@@ -442,7 +442,10 @@ pub fn compress_chunk_masked(
"more than 32 filters in a pipeline".into(),
));
}
let mut result = data.to_vec();
// The input is not copied: the first filter reads it where it is (a
// chunk may be 4 GiB or more), and each filter's output replaces the
// previous one.
let mut owned: Option<Vec<u8>> = None;
let mut mask = 0u32;
for (i, filter) in pipeline.filters.iter().enumerate() {
let ctx = FilterContext {
@@ -450,7 +453,8 @@ pub fn compress_chunk_masked(
element_size: element_size as usize,
max_output: 0,
};
let out = match filter_registry::encode(&result, &ctx) {
let result: &[u8] = owned.as_deref().unwrap_or(data);
let out = match filter_registry::encode(result, &ctx) {
Ok(out)
if FAIL_UNLESS_SMALLER.contains(&filter.filter_id) && out.len() >= result.len() =>
{
@@ -462,13 +466,13 @@ pub fn compress_chunk_masked(
r => r,
};
match out {
Ok(out) => result = out,
Ok(out) => owned = Some(out),
Err(e @ FormatError::UnsupportedFilter(_)) => return Err(e),
Err(_) if filter.flags & FILTER_FLAG_OPTIONAL != 0 => mask |= 1 << i,
Err(e) => return Err(e),
}
}
Ok((result, mask))
Ok((owned.unwrap_or_else(|| data.to_vec()), mask))
}
/// The filters compiled into this build, sorted by ID (see
@@ -1352,6 +1356,9 @@ fn deflate_compress(data: &[u8], level: u32) -> Result<Vec<u8>, FormatError> {
deflate_bounded(data, level).map_err(FormatError::CompressionError)
}
/// Largest compression output reserved at its worst-case size up front.
const DEFLATE_EXACT_BOUND: usize = 64 << 20;
/// Deflate `data` into a zlib stream in one pass, into a buffer sized for the
/// worst case up front (the same reasoning as [`inflate_bounded`]).
#[cfg(feature = "deflate")]
@@ -1360,24 +1367,44 @@ pub(crate) fn deflate_bounded(data: &[u8], level: u32) -> Result<Vec<u8>, String
// zlib's compressBound, plus the zlib header and trailer.
let bound = data.len() + (data.len() >> 12) + (data.len() >> 14) + (data.len() >> 25) + 13 + 6;
// flate2's Rust backends (zlib-rs, miniz_oxide) zero the whole spare
// capacity on each call, so a large input's worst-case bound would be
// memory held for nothing (4 GiB for a 4 GiB chunk that deflates to a
// few MiB): past 64 MiB the output starts at 1/16 of the bound and
// doubles as needed.
let first = if bound <= DEFLATE_EXACT_BOUND {
bound
} else {
bound / 16
};
let mut out = Vec::new();
out.try_reserve_exact(bound)
out.try_reserve_exact(first)
.map_err(|e| format!("deflate: cannot allocate output: {e}"))?;
let mut deflater = Compress::new(Compression::new(level), true);
loop {
let (in_before, out_before) = (deflater.total_in(), deflater.total_out());
let rest = &data[saturating_usize(in_before)..];
// zlib takes at most u32::MAX input bytes per call, and `Finish`
// ends the stream after the bytes it took: a chunk of 4 GiB or more
// was cut at 4 GiB - 1. Finish only once the rest fits one call.
let flush = if rest.len() > u32::MAX as usize {
FlushCompress::None
} else {
FlushCompress::Finish
};
let status = deflater
.compress_vec(
&data[saturating_usize(in_before)..],
&mut out,
FlushCompress::Finish,
)
.compress_vec(rest, &mut out, flush)
.map_err(|e| format!("deflate: {e}"))?;
match status {
Status::StreamEnd => return Ok(out),
// The bound should make running out of room unreachable; grow
// rather than fail if it happens.
Status::StreamEnd => {
if out.capacity() - out.len() > DEFLATE_EXACT_BOUND {
out.shrink_to_fit();
}
return Ok(out);
}
// Out of room (the bound makes it unreachable below
// `DEFLATE_EXACT_BOUND`): grow rather than fail.
Status::Ok | Status::BufError if out.len() == out.capacity() => out
.try_reserve(out.capacity().max(4096))
.map_err(|e| format!("deflate: cannot allocate output: {e}"))?,
@@ -1408,10 +1435,12 @@ const LZ4_DEFAULT_BLOCK_SIZE: usize = 1 << 30;
/// * The legacy clawhdf5 framing (up to 2.7.0): a 4-byte little-endian size
/// followed by one raw LZ4 block. libhdf5 cannot read it.
///
/// They are told apart unambiguously: an HDF5 chunk is smaller than 4 GiB, so
/// the registered format's big-endian `u64` size always starts with four zero
/// bytes and the whole chunk is at least 12 bytes; a legacy chunk starts with
/// four zero bytes only when it is empty, and is then 5 bytes long.
/// They are told apart unambiguously: for a chunk under 4 GiB the registered
/// format's big-endian `u64` size starts with four zero bytes and the whole
/// chunk is at least 12 bytes; a legacy chunk starts with four zero bytes
/// only when it is empty, and is then 5 bytes long. A chunk of 4 GiB or more
/// (HDF5 2.0) is always the registered format: clawhdf5 never wrote legacy
/// chunks that large.
///
/// Every size read from the payload is bounded against `expected_bytes` (the
/// pipeline's declared chunk size) before it sizes an allocation, so a crafted
@@ -1429,14 +1458,18 @@ fn lz4_decompress(data: &[u8], expected_bytes: usize) -> Result<Vec<u8>, FormatE
"lz4: declared size exceeds chunk size".into(),
));
}
if size > MAX_DECOMPRESS_SIZE {
// Without a chunk size, a ceiling; with one, the chunk size is the
// bound (a chunk may be 4 GiB or more, like a deflated one).
if expected_bytes == 0 && size > MAX_DECOMPRESS_SIZE {
return Err(FormatError::DecompressionError(
"lz4: declared size exceeds limit".into(),
));
}
Ok(())
};
if data.len() >= 12 && data[..4] == [0, 0, 0, 0] {
// A chunk of 4 GiB or more cannot be legacy (its size field is 32
// bits), and its registered-format size does not start with zeros.
if data.len() >= 12 && (data[..4] == [0, 0, 0, 0] || expected_bytes > u32::MAX as usize) {
return lz4_decompress_hdf5(data, check_size);
}
// Legacy clawhdf5 framing: 4-byte LE size + one LZ4 block.
@@ -1454,9 +1487,10 @@ fn lz4_decompress_hdf5(
) -> Result<Vec<u8>, FormatError> {
let err = |m: &str| FormatError::DecompressionError(format!("lz4: {m}"));
let be32 = |b: &[u8]| u32::from_be_bytes([b[0], b[1], b[2], b[3]]) as usize;
// The first four bytes are zero (checked by the caller), so the size is
// the low 32 bits of the big-endian u64.
let orig_size = be32(&data[4..8]);
let orig_size = u64::from_be_bytes([
data[0], data[1], data[2], data[3], data[4], data[5], data[6], data[7],
]);
let orig_size = usize::try_from(orig_size).map_err(|_| err("chunk too large"))?;
check_size(orig_size)?;
let block_size = be32(&data[8..12]).min(orig_size);
if block_size == 0 && orig_size != 0 {
+4 -12
View File
@@ -49,17 +49,9 @@ pub(crate) fn encode_elem(
/// element for chunks of `chunk_bytes` bytes (`H5D__earray_idx_create`,
/// `H5D__farray_idx_create`): one byte more than the nominal size needs —
/// except under layout message version 5 (HDF5 2.0's own format), which
/// always uses 8 bytes.
pub(crate) fn chunk_size_len(chunk_bytes: u64, layout_version: u8) -> usize {
if layout_version >= 5 {
return 8;
}
let log2 = if chunk_bytes <= 1 {
0
} else {
63 - chunk_bytes.leading_zeros()
};
(1 + ((log2 + 8) / 8) as usize).min(8)
/// uses the file's size of lengths (`length_size`).
pub(crate) fn chunk_size_len(chunk_bytes: u64, layout_version: u8, length_size: u8) -> usize {
clawhdf5_format::chunked_write::chunk_size_len(chunk_bytes, layout_version, length_size)
}
/// Creation parameters, in the layout message's order.
@@ -226,7 +218,7 @@ impl Ea {
) -> Result<Self, Error> {
let os = img.os;
let elem_size = if filtered {
os as usize + chunk_size_len(chunk_bytes, layout_version) + 4
os as usize + chunk_size_len(chunk_bytes, layout_version, img.ls) + 4
} else {
os as usize
};
+1 -1
View File
@@ -83,7 +83,7 @@ impl Fa {
let os = img.os;
let osz = os as usize;
let elem_size = if filtered {
osz + chunk_size_len(chunk_bytes, layout_version) + 4
osz + chunk_size_len(chunk_bytes, layout_version, img.ls) + 4
} else {
osz
};
+16 -3
View File
@@ -412,12 +412,24 @@ impl<'t> ChunkedEdit<'t> {
.iter()
.map(|&c| u64::from(c))
.collect();
// Chunks of 4 GiB or more (HDF5 2.0 writes them with layout message
// version 5) are read, but not rewritten: the editor holds a chunk
// it rewrites in memory, compressed and not. Writing values, and a
// resize that prunes or allocates chunks, are refused here, before
// anything is written; growing the extent (late allocation) and
// setting attributes do not come here.
let chunk_bytes = cd
.iter()
.try_fold(t.es as u64, |a, &c| a.checked_mul(c))
.filter(|&b| b <= u64::from(u32::MAX))
.and_then(|b| usize::try_from(b).ok())
.filter(|&b| b <= u32::MAX as usize)
.ok_or_else(|| Error::Unsupported("chunk larger than 4 GiB".into()))?;
.ok_or_else(|| {
Error::Unsupported(
"rewriting chunks of 4 GiB or more (writing values, or a resize that \
prunes or allocates chunks)"
.into(),
)
})?;
let hdr = Header::load(img, t.addr)?;
let layout_msg = hdr.find(MSG_LAYOUT).ok_or(Error::MissingMessage(
clawhdf5_format::message_type::MessageType::DataLayout,
@@ -600,7 +612,8 @@ impl<'t> ChunkedEdit<'t> {
let d = self.hdr.data(img, self.layout_msg)?;
let node_size = u32::from_le_bytes([d[at], d[at + 1], d[at + 2], d[at + 3]]);
let size_len = if filtered {
earray::chunk_size_len(self.chunk_bytes as u64, self.lpos.version) + 4
earray::chunk_size_len(self.chunk_bytes as u64, self.lpos.version, img.ls)
+ 4
} else {
0
};
@@ -294,3 +294,223 @@ fn unfiltered_huge_chunks_read() {
let read = storage.read.load(std::sync::atomic::Ordering::Relaxed);
assert!(read < 1 << 20, "{read} bytes read");
}
/// `FileEditor` does not rewrite chunks of 4 GiB or more: writing values,
/// or a resize that prunes or allocates chunks, is refused before anything
/// is written. Growing the extent and setting attributes still work.
#[test]
fn editor_refuses_rewriting_huge_chunks() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("huge.h5");
std::fs::copy(fixture(), &path).unwrap();
let before = std::fs::read(&path).unwrap();
let mut ed = clawhdf5::FileEditor::open(&path).unwrap();
let one = Selection::Hyperslab {
start: vec![0],
stride: vec![1],
count: vec![1],
block: vec![1],
};
for name in ["single", "farray", "earray"] {
let err = ed.write_values(name, &one, &[5.0f64]).unwrap_err();
assert!(
matches!(&err, clawhdf5::Error::Unsupported(m) if m.contains("4 GiB")),
"{name}: {err:?}"
);
}
// Shrinking prunes and fills chunks: refused.
for (name, shape) in [("earray", vec![10]), ("btree2", vec![1, N + 10])] {
let err = ed.resize(name, &shape).unwrap_err();
assert!(
matches!(&err, clawhdf5::Error::Unsupported(m) if m.contains("4 GiB")),
"{name}: {err:?}"
);
}
assert_eq!(
std::fs::read(&path).unwrap(),
before,
"a refused edit wrote"
);
// Growing without early allocation touches no chunk: the dataspace
// changes, as libhdf5's `H5Dset_extent` changes it.
ed.resize("earray", &[N + 20]).unwrap();
ed.resize("btree2", &[3, N + 10]).unwrap();
ed.set_attr("earray", "note", &clawhdf5::AttrValue::F64(1.5))
.unwrap();
let file = ed.reader().unwrap();
assert_eq!(file.dataset("earray").unwrap().shape().unwrap(), [N + 20]);
assert_eq!(
file.dataset("btree2").unwrap().shape().unwrap(),
[3, N + 10]
);
assert!(matches!(
file.dataset("earray").unwrap().attr("note").unwrap(),
Some(clawhdf5::AttrValue::F64(v)) if v == 1.5
));
}
/// The Fixed and Extensible Array structures clawhdf5 builds for 4 GiB
/// chunks are libhdf5's byte for byte: built at the fixture's addresses
/// from the fixture's chunks, they match what libhdf5 2.0.0 wrote (the
/// header's own address fields aside, which point where each writer put
/// the next block).
#[test]
fn huge_chunk_array_indexes_match_libhdf5() {
use clawhdf5_format::chunked_write::WrittenChunk;
let data = std::fs::read(fixture()).unwrap();
let u64_at = |at: usize| u64::from_le_bytes(data[at..at + 8].try_into().unwrap());
for name in ["farray", "earray"] {
let (_, layout, _, mut chunks) = layout_of(&data, name);
let DataLayout::Chunked {
btree_address: Some(hdr),
..
} = layout
else {
panic!("{name}: {layout:?}");
};
let hdr = hdr as usize;
chunks.sort_by_key(|c| c.offsets[0]);
let slots: Vec<Option<WrittenChunk>> = chunks
.iter()
.map(|c| {
Some(WrittenChunk {
address: c.address,
compressed_size: c.chunk_size,
raw_size: N * 8,
filter_mask: c.filter_mask,
})
})
.collect();
if name == "farray" {
let ours = clawhdf5_format::chunked_write::build_fixed_array_at(
&slots, 8, 8, true, hdr as u64,
);
// FAHD up to (not including) the data block address.
assert_eq!(&ours[..16], &data[hdr..hdr + 16], "FAHD");
let dblk = u64_at(hdr + 16) as usize;
let fadb = &ours[28..];
assert_eq!(fadb, &data[dblk..dblk + fadb.len()], "FADB");
} else {
let ours = clawhdf5_format::ea_writer::build_extensible_array_at(
&slots, 8, 8, true, hdr as u64,
);
// EAHD up to (not including) the index block address.
assert_eq!(&ours[..60], &data[hdr..hdr + 60], "EAHD");
let iblk = u64_at(hdr + 60) as usize;
let eaib = &ours[72..];
assert_eq!(eaib, &data[iblk..iblk + eaib.len()], "EAIB");
}
}
}
/// h5dump from libhdf5 2.x, when `CLAWHDF5_H5DUMP2` names one (Debian's
/// h5dump is 1.14, which cannot read layout message version 5).
fn h5dump2() -> Option<PathBuf> {
std::env::var_os("CLAWHDF5_H5DUMP2").map(PathBuf::from)
}
/// clawhdf5 writes chunks of 4 GiB + 8 bytes (deflated) in every index it
/// uses — Single Chunk, Fixed Array, Extensible Array, v2 B-tree — with
/// layout message version 5 and 8-byte stored sizes, as libhdf5 2.x does;
/// h5py (libhdf5 2.x) and h5dump 2.x read them. Each dataset holds ten
/// values in one 4 GiB chunk, so writing it holds one 4 GiB chunk.
#[test]
fn writer_huge_chunks_round_trip() {
if !heavy() {
return;
}
const UNLIM: u64 = u64::MAX;
let dir = scratch();
let path = dir.path().join("ours.h5");
let values: Vec<f64> = (0..10).map(f64::from).collect();
let mut b = clawhdf5::FileBuilder::new();
for (name, shape, max, chunk) in [
("single", vec![10], vec![N], vec![N]),
("farray", vec![10], vec![N + 10], vec![N]),
("earray", vec![10], vec![UNLIM], vec![N]),
("btree2", vec![1, 10], vec![UNLIM, UNLIM], vec![1, N]),
] {
b.create_dataset(name)
.with_f64_data(&values)
.with_shape(&shape)
.with_maxshape(&max)
.with_chunks(&chunk)
.with_deflate(6);
}
b.write(&path).unwrap();
let data = std::fs::read(&path).unwrap();
assert!(data.len() < 64 << 20, "{} bytes", data.len());
let mut stored = Vec::new();
for (name, index) in [("single", 1), ("farray", 3), ("earray", 4), ("btree2", 5)] {
let (version, layout, _, chunks) = layout_of(&data, name);
assert_eq!(version, 5, "{name}");
assert!(
matches!(layout, DataLayout::Chunked { chunk_index_type: Some(t), .. } if t == index),
"{name}: {layout:?}"
);
assert_eq!(chunks.len(), 1, "{name}");
stored.push((name, chunks[0].chunk_size));
}
let file = File::open(&path).unwrap();
let mut expect = values.clone();
expect.extend([0.0; 2]);
for name in ["single", "farray", "earray"] {
let got = file.dataset(name).unwrap().read_f64().unwrap();
assert_eq!(got, values, "{name}");
assert_eq!(sel(&file, name, &[0], &[10]), values, "{name}");
}
assert_eq!(sel(&file, "btree2", &[0, 3], &[1, 7]), f(3..10));
if python_available() {
let out = Command::new(python())
.arg("-c")
.arg(
r#"
import sys, h5py, numpy as np
with h5py.File(sys.argv[1], "r") as f:
for name in ("single", "farray", "earray", "btree2"):
d = f[name]
assert d.chunks[-1] == 2**29 + 1, (name, d.chunks)
got = d[0] if d.ndim == 2 else d[:]
assert list(got) == list(np.arange(10.0)), (name, got)
print("ok")
"#,
)
.arg(&path)
.output()
.unwrap();
assert!(
out.status.success(),
"h5py:\n{}",
String::from_utf8_lossy(&out.stderr)
);
} else {
assert!(!interop_required(), "h5py not available");
}
// h5dump 2.2.0 reads the layouts and walks every index: the storage
// size it reports is the chunk's. (It cannot print the values: its
// deflate filter fails on any chunk over 4 GiB, libhdf5's own included,
// with "memory allocation failed for deflate uncompression".)
if let Some(h5dump) = h5dump2() {
for (name, size) in stored {
let out = Command::new(&h5dump)
.args(["-H", "-p", "-d", name])
.arg(&path)
.output()
.unwrap();
let text = format!(
"{}{}",
String::from_utf8_lossy(&out.stdout),
String::from_utf8_lossy(&out.stderr)
);
assert!(out.status.success(), "h5dump {name}:\n{text}");
assert!(text.contains("536870913 )"), "{name}:\n{text}");
assert!(text.contains(&format!("SIZE {size} ")), "{name}:\n{text}");
assert!(!text.contains("rror"), "{name}:\n{text}");
}
}
}