Chunk dimensions of 2^32 or more; copy-free unfiltered 4 GiB writes #32

Open
osobh wants to merge 7 commits from feat/huge-chunk-dims into main
17 changed files with 1156 additions and 256 deletions
Showing only changes of commit e4b3a76dc7 - Show all commits
+42 -30
View File
@@ -555,7 +555,7 @@ pub(crate) fn checked_byte_len(elements: u64, elem_size: usize) -> Result<usize,
/// HDF5 2.0 writes them with layout version 5. A zero chunk dimension used /// HDF5 2.0 writes them with layout version 5. A zero chunk dimension used
/// to read as all fill values, and a huge one to hang the reader. /// to read as all fill values, and a huge one to hang the reader.
pub(crate) fn chunk_geometry( pub(crate) fn chunk_geometry(
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
layout_version: u8, layout_version: u8,
dataspace: &Dataspace, dataspace: &Dataspace,
elem_size: usize, elem_size: usize,
@@ -576,15 +576,26 @@ pub(crate) fn chunk_geometry(
"chunk size must be > 0, dim = {d}" "chunk size must be > 0, dim = {d}"
))); )));
} }
let bytes = spatial // Saturating: 33 dimensions of up to 64 bits each can overflow u128.
.iter() let bytes = spatial.iter().fold(elem_size as u128, |acc, &c| {
.fold(elem_size as u128, |acc, &c| acc * u128::from(c)); acc.saturating_mul(u128::from(c))
});
// So the indexes can compute a chunk's size in 64 bits.
if bytes > u128::from(u64::MAX) {
return Err(FormatError::InvalidChunkDimensions(format!(
"chunk size overflows 64 bits (chunk {spatial:?} of {elem_size}-byte elements)"
)));
}
if layout_version < 4 && bytes > u128::from(u32::MAX) { if layout_version < 4 && bytes > u128::from(u32::MAX) {
return Err(FormatError::InvalidChunkDimensions(format!( return Err(FormatError::InvalidChunkDimensions(format!(
"chunk size must be < 4GB with v1 b-tree index (chunk {spatial:?} of {elem_size}-byte elements)" "chunk size must be < 4GB with v1 b-tree index (chunk {spatial:?} of {elem_size}-byte elements)"
))); )));
} }
Ok((rank, spatial.iter().map(|&c| c as usize).collect())) let spatial = spatial
.iter()
.map(|&c| to_usize(c))
.collect::<Result<Vec<usize>, _>>()?;
Ok((rank, spatial))
} }
/// The size of one element of `dt` as stored in the file: a /// The size of one element of `dt` as stored in the file: a
@@ -625,7 +636,7 @@ pub fn check_chunk_element_size(
return Ok(()); return Ok(());
}; };
let expected = stored_element_size(datatype, offset_size); let expected = stored_element_size(datatype, offset_size);
if u64::from(stored) != expected { if stored != expected {
return Err(FormatError::InvalidChunkDimensions(format!( return Err(FormatError::InvalidChunkDimensions(format!(
"stored datatype size in chunk layout does not match datatype description \ "stored datatype size in chunk layout does not match datatype description \
(layout {stored} bytes, datatype {expected})" (layout {stored} bytes, datatype {expected})"
@@ -759,7 +770,7 @@ pub fn collect_chunk_info_in<S: Storage + ?Sized>(
pub fn collect_chunk_info_checked( pub fn collect_chunk_info_checked(
file_data: &[u8], file_data: &[u8],
btree_address: u64, btree_address: u64,
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
) -> Result<Vec<ChunkInfo>, FormatError> { ) -> Result<Vec<ChunkInfo>, FormatError> {
@@ -776,7 +787,7 @@ pub fn collect_chunk_info_checked(
pub fn collect_chunk_info_checked_in<S: Storage + ?Sized>( pub fn collect_chunk_info_checked_in<S: Storage + ?Sized>(
file_data: &S, file_data: &S,
btree_address: u64, btree_address: u64,
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
) -> Result<Vec<ChunkInfo>, FormatError> { ) -> Result<Vec<ChunkInfo>, FormatError> {
@@ -808,7 +819,7 @@ pub fn collect_chunk_info_checked_in<S: Storage + ?Sized>(
.iter_mut() .iter_mut()
.zip(chunk.offsets.iter().zip(chunk_dimensions)) .zip(chunk.offsets.iter().zip(chunk_dimensions))
{ {
*w = o / u64::from(d); *w = o / d;
} }
wanted[ndims - 1] = 0; wanted[ndims - 1] = 0;
if let Some(i) = root.find(&wanted) { if let Some(i) = root.find(&wanted) {
@@ -849,9 +860,9 @@ impl ChunkNode {
} }
} }
fn scale_keys(&mut self, dims: &[u32]) { fn scale_keys(&mut self, dims: &[u64]) {
for (k, &d) in self.keys.iter_mut().zip(dims.iter().cycle()) { for (k, &d) in self.keys.iter_mut().zip(dims.iter().cycle()) {
*k /= u64::from(d); *k /= d;
} }
if let ChunkChildren::Nodes(nodes) = &mut self.children { if let ChunkChildren::Nodes(nodes) = &mut self.children {
for n in nodes { for n in nodes {
@@ -919,9 +930,9 @@ fn btree_cmp3(lt: &[u64], scaled: &[u64], rt: &[u64]) -> core::cmp::Ordering {
/// Check one v1 B-tree chunk key's offsets (see /// Check one v1 B-tree chunk key's offsets (see
/// [`collect_chunk_info_checked`]). /// [`collect_chunk_info_checked`]).
fn check_key_offsets(offsets: &[u64], chunk_dimensions: &[u32]) -> Result<(), FormatError> { fn check_key_offsets(offsets: &[u64], chunk_dimensions: &[u64]) -> Result<(), FormatError> {
for (&offset, &dim) in offsets.iter().zip(chunk_dimensions) { for (&offset, &dim) in offsets.iter().zip(chunk_dimensions) {
if dim == 0 || offset % u64::from(dim) != 0 { if dim == 0 || offset % dim != 0 {
return Err(FormatError::ChunkedReadError(format!( return Err(FormatError::ChunkedReadError(format!(
"bad coordinate offset {offsets:?} for chunk dimensions {chunk_dimensions:?}" "bad coordinate offset {offsets:?} for chunk dimensions {chunk_dimensions:?}"
))); )));
@@ -937,7 +948,7 @@ fn read_key_offsets(
file_data: &[u8], file_data: &[u8],
pos: usize, pos: usize,
ndims: usize, ndims: usize,
chunk_dimensions: Option<&[u32]>, chunk_dimensions: Option<&[u64]>,
out: &mut Vec<u64>, out: &mut Vec<u64>,
) -> Result<(), FormatError> { ) -> Result<(), FormatError> {
let start = out.len(); let start = out.len();
@@ -966,7 +977,7 @@ fn parse_chunk_node<S: Storage + ?Sized>(
file_data: &S, file_data: &S,
btree_address: u64, btree_address: u64,
ndims: usize, ndims: usize,
chunk_dimensions: Option<&[u32]>, chunk_dimensions: Option<&[u64]>,
offset_size: u8, offset_size: u8,
depth: usize, depth: usize,
stored: &mut Vec<ChunkInfo>, stored: &mut Vec<ChunkInfo>,
@@ -1075,7 +1086,7 @@ fn parse_chunk_node<S: Storage + ?Sized>(
pub fn generate_implicit_chunks( pub fn generate_implicit_chunks(
base_address: u64, base_address: u64,
dataset_dims: &[u64], dataset_dims: &[u64],
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
element_size: u32, element_size: u32,
) -> Vec<ChunkInfo> { ) -> Vec<ChunkInfo> {
generate_implicit_chunks_in_grid( generate_implicit_chunks_in_grid(
@@ -1097,17 +1108,18 @@ pub fn generate_implicit_chunks_in_grid(
base_address: u64, base_address: u64,
dataset_dims: &[u64], dataset_dims: &[u64],
max_dims: &[u64], max_dims: &[u64],
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
element_size: u32, element_size: u32,
) -> Vec<ChunkInfo> { ) -> Vec<ChunkInfo> {
let rank = chunk_dimensions.len(); let rank = chunk_dimensions.len();
let chunk_byte_size: u64 = let chunk_byte_size: u64 = chunk_dimensions
chunk_dimensions.iter().map(|&d| d as u64).product::<u64>() * element_size as u64; .iter()
.fold(u64::from(element_size), |acc, &d| acc.saturating_mul(d));
let mut num_chunks_per_dim = Vec::with_capacity(rank); let mut num_chunks_per_dim = Vec::with_capacity(rank);
let mut grid_per_dim = Vec::with_capacity(rank); let mut grid_per_dim = Vec::with_capacity(rank);
for d in 0..rank { for d in 0..rank {
let ch = chunk_dimensions[d] as u64; let ch = chunk_dimensions[d];
let n = dataset_dims[d].div_ceil(ch); let n = dataset_dims[d].div_ceil(ch);
num_chunks_per_dim.push(n); num_chunks_per_dim.push(n);
grid_per_dim.push(max_dims.get(d).map_or(n, |m| m.div_ceil(ch)).max(n)); grid_per_dim.push(max_dims.get(d).map_or(n, |m| m.div_ceil(ch)).max(n));
@@ -1125,7 +1137,7 @@ pub fn generate_implicit_chunks_in_grid(
let nchunks = num_chunks_per_dim[d]; let nchunks = num_chunks_per_dim[d];
let chunk_idx = remaining % nchunks; let chunk_idx = remaining % nchunks;
remaining /= nchunks; remaining /= nchunks;
offsets[d] = chunk_idx * chunk_dimensions[d] as u64; offsets[d] = chunk_idx * chunk_dimensions[d];
grid_idx = grid_idx.saturating_add(chunk_idx.saturating_mul(down)); grid_idx = grid_idx.saturating_add(chunk_idx.saturating_mul(down));
down = down.saturating_mul(grid_per_dim[d]); down = down.saturating_mul(grid_per_dim[d]);
} }
@@ -1350,7 +1362,7 @@ pub fn list_chunks_in<S: Storage + ?Sized>(
} }
(4, Some(2)) => { (4, Some(2)) => {
// Implicit index — use spatial chunk dims only // Implicit index — use spatial chunk dims only
let spatial_chunk_dims: &[u32] = &chunk_dimensions[..rank]; let spatial_chunk_dims: &[u64] = &chunk_dimensions[..rank];
generate_implicit_chunks_in_grid( generate_implicit_chunks_in_grid(
addr, addr,
&dataspace.dimensions, &dataspace.dimensions,
@@ -1364,7 +1376,7 @@ pub fn list_chunks_in<S: Storage + ?Sized>(
} }
(4, Some(3)) => { (4, Some(3)) => {
// Fixed Array — use spatial chunk dims only // Fixed Array — use spatial chunk dims only
let spatial_chunk_dims: &[u32] = &chunk_dimensions[..rank]; let spatial_chunk_dims: &[u64] = &chunk_dimensions[..rank];
let header = FixedArrayHeader::parse_in( let header = FixedArrayHeader::parse_in(
file_data, file_data,
checked_addr(addr)?, checked_addr(addr)?,
@@ -1384,7 +1396,7 @@ pub fn list_chunks_in<S: Storage + ?Sized>(
} }
(4, Some(4)) => { (4, Some(4)) => {
// Extensible Array — use spatial chunk dims only // Extensible Array — use spatial chunk dims only
let spatial_chunk_dims: &[u32] = &chunk_dimensions[..rank]; let spatial_chunk_dims: &[u64] = &chunk_dimensions[..rank];
let header = ExtensibleArrayHeader::parse_in( let header = ExtensibleArrayHeader::parse_in(
file_data, file_data,
checked_addr(addr)?, checked_addr(addr)?,
@@ -2435,12 +2447,12 @@ mod tests {
/// A leaf whose final key is the one libhdf5 writes past the last chunk: /// A leaf whose final key is the one libhdf5 writes past the last chunk:
/// each coordinate of the last chunk plus its chunk dimension (the /// each coordinate of the last chunk plus its chunk dimension (the
/// element size last). /// element size last).
fn build_chunk_btree_leaf_dims(chunks: &[ChunkInfo], dims: &[u32], offset_size: u8) -> Vec<u8> { fn build_chunk_btree_leaf_dims(chunks: &[ChunkInfo], dims: &[u64], offset_size: u8) -> Vec<u8> {
let last = &chunks.last().expect("a chunk").offsets; let last = &chunks.last().expect("a chunk").offsets;
let end: Vec<u64> = dims let end: Vec<u64> = dims
.iter() .iter()
.enumerate() .enumerate()
.map(|(d, &c)| last.get(d).copied().unwrap_or(0) + u64::from(c)) .map(|(d, &c)| last.get(d).copied().unwrap_or(0) + c)
.collect(); .collect();
build_chunk_btree_leaf_to(chunks, &end, offset_size) build_chunk_btree_leaf_to(chunks, &end, offset_size)
} }
@@ -2771,7 +2783,7 @@ mod tests {
} }
// Build B-tree at offset 0x100 // Build B-tree at offset 0x100
let dims = [chunk_size_elems as u32, elem_size as u32]; let dims = [chunk_size_elems as u64, elem_size as u64];
let btree = build_chunk_btree_leaf_dims(&chunk_infos, &dims, os); let btree = build_chunk_btree_leaf_dims(&chunk_infos, &dims, os);
let btree_addr = 0x100usize; let btree_addr = 0x100usize;
file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree); file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree);
@@ -3145,7 +3157,7 @@ mod tests {
data_offset += compressed.len() + 16; // some padding data_offset += compressed.len() + 16; // some padding
} }
let dims = [chunk_elems as u32, elem_size as u32]; let dims = [chunk_elems as u64, elem_size as u64];
let btree = build_chunk_btree_leaf_dims(&chunk_infos, &dims, os); let btree = build_chunk_btree_leaf_dims(&chunk_infos, &dims, os);
let btree_addr = 0x100usize; let btree_addr = 0x100usize;
file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree); file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree);
@@ -3228,7 +3240,7 @@ mod tests {
} }
} }
let dims = [chunk_dims[0] as u32, chunk_dims[1] as u32, elem_size as u32]; let dims = [chunk_dims[0] as u64, chunk_dims[1] as u64, elem_size as u64];
let btree = build_chunk_btree_leaf_dims(&chunk_infos, &dims, os); let btree = build_chunk_btree_leaf_dims(&chunk_infos, &dims, os);
let btree_addr = 0x100usize; let btree_addr = 0x100usize;
file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree); file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree);
@@ -3408,7 +3420,7 @@ mod tests {
} }
let layout = DataLayout::Chunked { let layout = DataLayout::Chunked {
chunk_dimensions: vec![chunk_elems as u32, elem_size as u32], chunk_dimensions: vec![chunk_elems as u64, elem_size as u64],
btree_address: Some(data_addr as u64), btree_address: Some(data_addr as u64),
version: 4, version: 4,
chunk_index_type: Some(1), chunk_index_type: Some(1),
+351 -102
View File
@@ -5,12 +5,14 @@ extern crate alloc;
use crate::addr::saturating_usize; use crate::addr::saturating_usize;
#[cfg(not(feature = "std"))] #[cfg(not(feature = "std"))]
use alloc::{format, vec, vec::Vec}; use alloc::{borrow::Cow, format, vec, vec::Vec};
#[cfg(feature = "std")]
use std::borrow::Cow;
use crate::btree_v1_write; use crate::btree_v1_write;
use crate::btree_v2_write::{BTreeV2Params, build_btree_v2}; use crate::btree_v2_write::{BTreeV2Params, build_btree_v2};
use crate::checksum::jenkins_lookup3; use crate::checksum::jenkins_lookup3;
use crate::chunk_cache::{CACHE_LINE_SIZE, align_to_cache_line}; use crate::chunk_cache::CACHE_LINE_SIZE;
use crate::chunk_grid::ChunkGrid; use crate::chunk_grid::ChunkGrid;
use crate::ea_writer; use crate::ea_writer;
use crate::error::FormatError; use crate::error::FormatError;
@@ -464,7 +466,9 @@ fn extract_chunk(
let dst: usize = (0..rank).map(|d| idx[d] * chunk_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. // Whole elements only, as far as `raw_data` reaches.
let n = row.min(raw_data.len().saturating_sub(src)) / element_size * element_size; 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]); if n > 0 {
chunk[dst..dst + n].copy_from_slice(&raw_data[src..src + n]);
}
// Next row: advance every dimension but the last. // Next row: advance every dimension but the last.
let mut d = rank - 1; let mut d = rank - 1;
loop { loop {
@@ -481,6 +485,66 @@ fn extract_chunk(
} }
} }
/// The bytes of the `linear_idx`-th chunk (see [`extract_chunk`]),
/// borrowed from `raw_data` when the chunk is one contiguous run of it: a
/// chunk inside the dataset (not past its edge) whose dimensions after the
/// first one above 1 equal the dataset's. A dataset stored as one chunk of
/// its own shape is such a run, so its chunk is never copied. Other chunks
/// are copied out, as [`extract_chunk`] does.
fn chunk_bytes_of<'a>(
raw_data: &'a [u8],
shape: &[u64],
chunk_dims: &[u64],
element_size: usize,
linear_idx: u64,
) -> Cow<'a, [u8]> {
if shape.is_empty() {
return Cow::Borrowed(raw_data);
}
if let Some(range) = contiguous_chunk(shape, chunk_dims, element_size, linear_idx)
&& let (Ok(start), Ok(end)) = (usize::try_from(range.start), usize::try_from(range.end))
&& end <= raw_data.len()
{
return Cow::Borrowed(&raw_data[start..end]);
}
Cow::Owned(extract_chunk(raw_data, shape, chunk_dims, element_size, linear_idx).1)
}
/// The byte range of the `linear_idx`-th chunk in the row-major dataset,
/// when the chunk is a contiguous run of it (see [`chunk_bytes_of`]).
fn contiguous_chunk(
shape: &[u64],
chunk_dims: &[u64],
element_size: usize,
linear_idx: u64,
) -> Option<core::ops::Range<u64>> {
let rank = shape.len();
// Past the first chunk dimension above 1, every chunk dimension must
// span the dataset's.
let first_wide = chunk_dims.iter().position(|&c| c > 1).unwrap_or(rank);
if (first_wide + 1..rank).any(|d| chunk_dims[d] != shape[d]) {
return None;
}
let mut remaining = linear_idx;
let mut start = 0u64;
let mut stride = element_size as u64;
for d in (0..rank).rev() {
let n = shape[d].div_ceil(chunk_dims[d]);
let offset = (remaining % n).checked_mul(chunk_dims[d])?;
remaining /= n;
// A chunk past the dataset's edge is padded: not a run of it.
if offset.checked_add(chunk_dims[d])? > shape[d] {
return None;
}
start = start.checked_add(offset.checked_mul(stride)?)?;
stride = stride.checked_mul(shape[d])?;
}
let len = chunk_dims
.iter()
.try_fold(element_size as u64, |acc, &c| acc.checked_mul(c))?;
Some(start..start.checked_add(len)?)
}
/// Parallel compression threshold: use rayon when chunk count exceeds this. /// Parallel compression threshold: use rayon when chunk count exceeds this.
/// ///
/// Lowered to 2 to enable parallel compression for typical 4-chunk workloads /// Lowered to 2 to enable parallel compression for typical 4-chunk workloads
@@ -508,25 +572,21 @@ const PARALLEL_COMPRESS_MAX_CHUNK_BYTES: u64 = 64 << 20;
/// at most [`PARALLEL_COMPRESS_MAX_CHUNK_BYTES`], compression runs across /// at most [`PARALLEL_COMPRESS_MAX_CHUNK_BYTES`], compression runs across
/// rayon threads; otherwise it is sequential. Output order matches chunk /// rayon threads; otherwise it is sequential. Output order matches chunk
/// order, so per-chunk bytes are identical to the sequential path. /// order, so per-chunk bytes are identical to the sequential path.
fn compress_all_chunks( fn compress_all_chunks<'a>(
raw_data: &[u8], raw_data: &'a [u8],
shape: &[u64], shape: &[u64],
chunk_dims: &[u64], chunk_dims: &[u64],
element_size: usize, element_size: usize,
chunk_bytes: u64, chunk_bytes: u64,
pipeline: &Option<FilterPipeline>, pipeline: &Option<FilterPipeline>,
) -> Result<Vec<(u64, Vec<u8>, u32)>, FormatError> { ) -> Result<Vec<PreparedChunk<'a>>, FormatError> {
let one = |i: u64| -> Result<(u64, Vec<u8>, u32), FormatError> { let one = |i: u64| -> Result<PreparedChunk<'a>, FormatError> {
let (_, raw) = if shape.is_empty() { let raw = chunk_bytes_of(raw_data, shape, chunk_dims, element_size, i);
(Vec::new(), raw_data.to_vec())
} else {
extract_chunk(raw_data, shape, chunk_dims, element_size, i)
};
let raw_size = raw.len() as u64; let raw_size = raw.len() as u64;
match pipeline { match pipeline {
Some(pl) => { Some(pl) => {
let (stored, mask) = compress_chunk_masked(&raw, pl, element_size as u32)?; let (stored, mask) = compress_chunk_masked(&raw, pl, element_size as u32)?;
Ok((raw_size, stored, mask)) Ok((raw_size, Cow::Owned(stored), mask))
} }
None => Ok((raw_size, raw, 0)), None => Ok((raw_size, raw, 0)),
} }
@@ -557,7 +617,7 @@ fn compress_all_chunks(
/// layout/pipeline messages. `base_address` is where the blob will be placed in the file. /// 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). /// Serialize a v4 single chunk layout message (public for OH size estimation).
pub fn serialize_v4_single_chunk_pub( pub fn serialize_v4_single_chunk_pub(
chunk_dims: &[u32], chunk_dims: &[u64],
chunk_address: u64, chunk_address: u64,
filtered_size: Option<u64>, filtered_size: Option<u64>,
filter_mask: Option<u32>, filter_mask: Option<u32>,
@@ -578,7 +638,7 @@ pub fn serialize_v4_single_chunk_pub(
/// Serialize a v4 (or, for a chunk of 4 GiB or more, v5) single chunk /// Serialize a v4 (or, for a chunk of 4 GiB or more, v5) single chunk
/// layout message. /// layout message.
fn serialize_v4_single_chunk( fn serialize_v4_single_chunk(
chunk_dims: &[u32], chunk_dims: &[u64],
chunk_address: u64, chunk_address: u64,
filtered_size: Option<u64>, filtered_size: Option<u64>,
filter_mask: Option<u32>, filter_mask: Option<u32>,
@@ -622,7 +682,7 @@ fn serialize_v4_single_chunk(
/// Serialize a v4 Fixed Array layout message. /// Serialize a v4 Fixed Array layout message.
fn serialize_v4_fixed_array( fn serialize_v4_fixed_array(
chunk_dims: &[u32], chunk_dims: &[u64],
fixed_array_address: u64, fixed_array_address: u64,
offset_size: u8, offset_size: u8,
element_size: u32, element_size: u32,
@@ -653,7 +713,8 @@ fn serialize_v4_fixed_array(
/// dimensions, then the element size). Each takes the fewest bytes that hold /// dimensions, then the element size). Each takes the fewest bytes that hold
/// the largest, as libhdf5 computes it (`H5D__chunk_set_sizes`: /// the largest, as libhdf5 computes it (`H5D__chunk_set_sizes`:
/// `(log2(dim) + 8) / 8`); HDF5 2.0.0 refuses any other width. /// `(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) { pub(crate) fn push_v4_chunk_dims(buf: &mut Vec<u8>, chunk_dims: &[u64], element_size: u32) {
let element_size = u64::from(element_size);
let max_dim = chunk_dims let max_dim = chunk_dims
.iter() .iter()
.copied() .copied()
@@ -661,14 +722,14 @@ pub(crate) fn push_v4_chunk_dims(buf: &mut Vec<u8>, chunk_dims: &[u32], element_
.max() .max()
.unwrap_or(1) .unwrap_or(1)
.max(1); .max(1);
let width = (32 - max_dim.leading_zeros()).div_ceil(8) as usize; let width = (64 - max_dim.leading_zeros()).div_ceil(8) as usize;
buf.push(width as u8); buf.push(width as u8);
for &d in chunk_dims.iter().chain(core::iter::once(&element_size)) { for &d in chunk_dims.iter().chain(core::iter::once(&element_size)) {
buf.extend_from_slice(&d.to_le_bytes()[..width]); buf.extend_from_slice(&d.to_le_bytes()[..width]);
} }
} }
fn layout_v4_chunked_prefix(chunk_dims: &[u32], element_size: u32, version: u8) -> Vec<u8> { fn layout_v4_chunked_prefix(chunk_dims: &[u64], element_size: u32, version: u8) -> Vec<u8> {
let mut buf = Vec::new(); let mut buf = Vec::new();
buf.push(version); buf.push(version);
buf.push(2); // class = chunked buf.push(2); // class = chunked
@@ -884,7 +945,66 @@ pub fn precompress_chunks(
element_size: usize, element_size: usize,
options: &ChunkOptions, options: &ChunkOptions,
) -> Result<PrecompressedChunks, FormatError> { ) -> Result<PrecompressedChunks, FormatError> {
let (_, chunk_bytes) = checked_chunk_dims(chunk_dims, element_size)?; let p = prepare_chunks(raw_data, shape, chunk_dims, element_size, options)?;
Ok(PrecompressedChunks {
chunks: p
.chunks
.into_iter()
.map(|(raw, stored, mask)| (raw, stored.into_owned(), mask))
.collect(),
has_filters: p.has_filters,
element_size: p.element_size,
shape: p.shape,
chunk_dims: p.chunk_dims,
pipeline_message: p.pipeline_message,
})
}
/// One chunk of [`PreparedChunks`]: its raw size, its stored bytes and its
/// filter mask.
pub(crate) type PreparedChunk<'a> = (u64, Cow<'a, [u8]>, u32);
/// [`PrecompressedChunks`] whose stored bytes may borrow the dataset's data:
/// an unfiltered chunk that is a contiguous run of it (a dataset stored as
/// one chunk, say) is not copied. What the file writer lays out.
pub(crate) struct PreparedChunks<'a> {
pub(crate) chunks: Vec<PreparedChunk<'a>>,
pub(crate) has_filters: bool,
pub(crate) element_size: usize,
pub(crate) shape: Vec<u64>,
pub(crate) chunk_dims: Vec<u64>,
pub(crate) pipeline_message: Option<Vec<u8>>,
}
impl PrecompressedChunks {
/// A view of these chunks as [`PreparedChunks`] (their bytes borrowed).
fn prepared(&self) -> PreparedChunks<'_> {
PreparedChunks {
chunks: self
.chunks
.iter()
.map(|(raw, stored, mask)| (*raw, Cow::Borrowed(&stored[..]), *mask))
.collect(),
has_filters: self.has_filters,
element_size: self.element_size,
shape: self.shape.clone(),
chunk_dims: self.chunk_dims.clone(),
pipeline_message: self.pipeline_message.clone(),
}
}
}
/// [`precompress_chunks`], borrowing what it can (see [`PreparedChunks`]).
/// A filtered chunk that is a contiguous run of `raw_data` is compressed
/// straight from it, without a copy of the raw chunk.
pub(crate) fn prepare_chunks<'a>(
raw_data: &'a [u8],
shape: &[u64],
chunk_dims: &[u64],
element_size: usize,
options: &ChunkOptions,
) -> Result<PreparedChunks<'a>, FormatError> {
let chunk_bytes = checked_chunk_dims(chunk_dims, element_size)?;
if chunk_bytes > MAX_V4_CHUNK_BYTES { if chunk_bytes > MAX_V4_CHUNK_BYTES {
check_huge_chunk_filters(options, chunk_bytes)?; check_huge_chunk_filters(options, chunk_bytes)?;
} }
@@ -902,7 +1022,7 @@ pub fn precompress_chunks(
&pipeline, &pipeline,
)?; )?;
Ok(PrecompressedChunks { Ok(PreparedChunks {
chunks, chunks,
has_filters, has_filters,
element_size, element_size,
@@ -912,26 +1032,87 @@ pub fn precompress_chunks(
}) })
} }
/// The chunk dimensions as the layout message stores them (each below /// A piece of a chunked dataset's data blob as laid out by
/// 2^32), and one chunk's size in bytes. A chunk dimension of 2^32 or more /// [`place_chunked_data`]: zero padding, a chunk's stored bytes (by index
/// (which HDF5 2.0 can store, in wider fields) is refused, as it is when /// into [`PreparedChunks::chunks`]), or index structures.
/// read: clawhdf5 holds chunk dimensions as `u32`. So is a chunk whose size #[derive(Debug)]
/// overflows 64 bits, or that this platform cannot hold in memory (a chunk pub(crate) enum Piece {
/// of 4 GiB or more on a 32-bit target). Zeros(usize),
fn checked_chunk_dims( Chunk(usize),
chunk_dims: &[u64], Bytes(Vec<u8>),
element_size: usize, }
) -> Result<(Vec<u32>, u64), FormatError> {
let dims = chunk_dims /// A chunked dataset laid out at an address, without its chunk bytes
.iter() /// copied: the pieces of its data blob, in order, and its messages.
.map(|&d| { pub(crate) struct ChunkedPlacement {
u32::try_from(d).map_err(|_| { pub(crate) pieces: Vec<Piece>,
FormatError::InvalidChunkDimensions(format!( /// Length of the data blob in bytes.
"chunk dimension {d} is 2^32 or more, which clawhdf5 does not support" pub(crate) len: u64,
)) pub(crate) layout_message: Vec<u8>,
}) pub(crate) pipeline_message: Option<Vec<u8>>,
}) }
.collect::<Result<Vec<u32>, _>>()?;
impl ChunkedPlacement {
/// The data blob as one buffer (the chunks copied in).
fn into_result(self, pre: &PreparedChunks<'_>) -> ChunkedDataResult {
let mut data_bytes = Vec::with_capacity(saturating_usize(self.len));
for piece in &self.pieces {
match piece {
Piece::Zeros(n) => data_bytes.resize(data_bytes.len() + n, 0),
Piece::Chunk(i) => data_bytes.extend_from_slice(&pre.chunks[*i].1),
Piece::Bytes(b) => data_bytes.extend_from_slice(b),
}
}
ChunkedDataResult {
data_bytes,
layout_message: self.layout_message,
pipeline_message: self.pipeline_message,
}
}
}
/// Builds the pieces of a data blob, tracking its length.
#[derive(Default)]
struct Pieces {
pieces: Vec<Piece>,
len: u64,
}
impl Pieces {
/// Pad to the next cache line.
fn align(&mut self) {
let aligned = align_chunk_offset(self.len);
if aligned > self.len {
self.pieces
.push(Piece::Zeros(saturating_usize(aligned - self.len)));
self.len = aligned;
}
}
fn chunk(&mut self, i: usize, len: usize) {
self.pieces.push(Piece::Chunk(i));
self.len += len as u64;
}
fn bytes(&mut self, b: Vec<u8>) {
self.len += b.len() as u64;
self.pieces.push(Piece::Bytes(b));
}
}
/// One chunk's size in bytes. A chunk dimension of 0 is refused, and 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). Chunk
/// dimensions of 2^32 or more are allowed, as in libhdf5 2.x: such a chunk
/// is over 4 GiB, so it takes layout message version 5, whose dimension
/// fields are up to 8 bytes wide (a version-3 layout, with 4-byte fields,
/// is never written for it).
fn checked_chunk_dims(chunk_dims: &[u64], element_size: usize) -> Result<u64, FormatError> {
if chunk_dims.contains(&0) {
return Err(FormatError::InvalidChunkDimensions(format!(
"chunk dimensions {chunk_dims:?} include 0"
)));
}
let bytes = chunk_dims let bytes = chunk_dims
.iter() .iter()
.try_fold(element_size as u64, |acc, &d| acc.checked_mul(d)) .try_fold(element_size as u64, |acc, &d| acc.checked_mul(d))
@@ -941,7 +1122,7 @@ fn checked_chunk_dims(
"a chunk of {chunk_dims:?} x {element_size} bytes exceeds this platform's address space" "a chunk of {chunk_dims:?} x {element_size} bytes exceeds this platform's address space"
)) ))
})?; })?;
Ok((dims, bytes)) Ok(bytes)
} }
/// The filters clawhdf5 can apply to a chunk of more than `u32::MAX` bytes: /// The filters clawhdf5 can apply to a chunk of more than `u32::MAX` bytes:
@@ -1010,7 +1191,22 @@ pub fn build_chunked_data_from_precompressed_libver(
low: LibVer, low: LibVer,
high: LibVer, high: LibVer,
) -> Result<ChunkedDataResult, FormatError> { ) -> Result<ChunkedDataResult, FormatError> {
let (_, chunk_bytes) = checked_chunk_dims(&pre.chunk_dims, pre.element_size)?; let pre = pre.prepared();
let placement = place_chunked_data(&pre, base_address, maxshape, low, high)?;
Ok(placement.into_result(&pre))
}
/// [`build_chunked_data_from_precompressed_libver`] without copying the
/// chunks into one buffer: the data blob as [`Piece`]s referring to
/// `pre`'s chunks. The file writer writes the pieces out one by one.
pub(crate) fn place_chunked_data(
pre: &PreparedChunks<'_>,
base_address: u64,
maxshape: Option<&[u64]>,
low: LibVer,
high: LibVer,
) -> Result<ChunkedPlacement, FormatError> {
let chunk_bytes = checked_chunk_dims(&pre.chunk_dims, pre.element_size)?;
if chunk_bytes > MAX_V4_CHUNK_BYTES { if chunk_bytes > MAX_V4_CHUNK_BYTES {
if high < LibVer::V200 { if high < LibVer::V200 {
return Err(FormatError::LibverBound { return Err(FormatError::LibverBound {
@@ -1020,7 +1216,7 @@ pub fn build_chunked_data_from_precompressed_libver(
}); });
} }
} else if low < LibVer::V110 { } else if low < LibVer::V110 {
return build_btree_v1_chunked_data(pre, base_address, maxshape); return place_btree_v1_chunked_data(pre, base_address, maxshape);
} }
let index = ChunkIndexPlan::new(&pre.shape, maxshape, &pre.chunk_dims)?; let index = ChunkIndexPlan::new(&pre.shape, maxshape, &pre.chunk_dims)?;
let offset_size: u8 = 8; let offset_size: u8 = 8;
@@ -1028,17 +1224,14 @@ pub fn build_chunked_data_from_precompressed_libver(
let num_chunks = pre.chunks.len(); let num_chunks = pre.chunks.len();
let element_size = pre.element_size; let element_size = pre.element_size;
let mut data_buf = Vec::new(); let mut blob = Pieces::default();
let mut written_chunks = Vec::with_capacity(num_chunks); let mut written_chunks = Vec::with_capacity(num_chunks);
for (raw_size, compressed, filter_mask) in &pre.chunks { for (i, (raw_size, compressed, filter_mask)) in pre.chunks.iter().enumerate() {
let aligned_offset = align_to_cache_line(data_buf.len()); blob.align();
if aligned_offset > data_buf.len() { let address = base_address + blob.len;
data_buf.resize(aligned_offset, 0u8);
}
let address = base_address + data_buf.len() as u64;
let compressed_size = compressed.len() as u64; let compressed_size = compressed.len() as u64;
data_buf.extend_from_slice(compressed); blob.chunk(i, compressed.len());
written_chunks.push(WrittenChunk { written_chunks.push(WrittenChunk {
address, address,
compressed_size, compressed_size,
@@ -1047,28 +1240,22 @@ pub fn build_chunked_data_from_precompressed_libver(
}); });
} }
let (chunk_dims_u32, chunk_bytes) = checked_chunk_dims(&pre.chunk_dims, element_size)?;
let version = layout_version_for(chunk_bytes); let version = layout_version_for(chunk_bytes);
blob.align();
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 { let layout_message = match &index {
ChunkIndexPlan::ExtensibleArray(grid) => { ChunkIndexPlan::ExtensibleArray(grid) => {
let ea_address = base_address + data_buf.len() as u64; let ea_address = base_address + blob.len;
let slots = index_slots(grid, &pre.shape, &pre.chunk_dims, &written_chunks, None)?; let slots = index_slots(grid, &pre.shape, &pre.chunk_dims, &written_chunks, None)?;
let ea_bytes = ea_writer::build_extensible_array_at( blob.bytes(ea_writer::build_extensible_array_at(
&slots, &slots,
offset_size, offset_size,
length_size, length_size,
pre.has_filters, pre.has_filters,
ea_address, ea_address,
); ));
data_buf.extend_from_slice(&ea_bytes);
ea_writer::serialize_v4_extensible_array( ea_writer::serialize_v4_extensible_array(
&chunk_dims_u32, &pre.chunk_dims,
ea_address, ea_address,
offset_size, offset_size,
element_size as u32, element_size as u32,
@@ -1084,7 +1271,7 @@ pub fn build_chunked_data_from_precompressed_libver(
}; };
let filter_mask = pre.has_filters.then_some(written_chunks[0].filter_mask); let filter_mask = pre.has_filters.then_some(written_chunks[0].filter_mask);
serialize_v4_single_chunk( serialize_v4_single_chunk(
&chunk_dims_u32, &pre.chunk_dims,
chunk_addr, chunk_addr,
filtered_size, filtered_size,
filter_mask, filter_mask,
@@ -1094,7 +1281,7 @@ pub fn build_chunked_data_from_precompressed_libver(
) )
} }
ChunkIndexPlan::FixedArray(grid, nslots) => { ChunkIndexPlan::FixedArray(grid, nslots) => {
let fa_address = base_address + data_buf.len() as u64; let fa_address = base_address + blob.len;
let slots = index_slots( let slots = index_slots(
grid, grid,
&pre.shape, &pre.shape,
@@ -1102,16 +1289,15 @@ pub fn build_chunked_data_from_precompressed_libver(
&written_chunks, &written_chunks,
Some(*nslots), Some(*nslots),
)?; )?;
let fa_bytes = build_fixed_array_at( blob.bytes(build_fixed_array_at(
&slots, &slots,
offset_size, offset_size,
length_size, length_size,
pre.has_filters, pre.has_filters,
fa_address, fa_address,
); ));
data_buf.extend_from_slice(&fa_bytes);
serialize_v4_fixed_array( serialize_v4_fixed_array(
&chunk_dims_u32, &pre.chunk_dims,
fa_address, fa_address,
offset_size, offset_size,
element_size as u32, element_size as u32,
@@ -1120,7 +1306,7 @@ pub fn build_chunked_data_from_precompressed_libver(
) )
} }
ChunkIndexPlan::BTreeV2 => { ChunkIndexPlan::BTreeV2 => {
let bt_address = base_address + data_buf.len() as u64; let bt_address = base_address + blob.len;
let records: Vec<(Vec<u64>, &WrittenChunk)> = written_chunks let records: Vec<(Vec<u64>, &WrittenChunk)> = written_chunks
.iter() .iter()
.enumerate() .enumerate()
@@ -1134,9 +1320,9 @@ pub fn build_chunked_data_from_precompressed_libver(
pre.has_filters, pre.has_filters,
bt_address, bt_address,
)?; )?;
data_buf.extend_from_slice(&bt_bytes); blob.bytes(bt_bytes);
serialize_v4_btree_v2( serialize_v4_btree_v2(
&chunk_dims_u32, &pre.chunk_dims,
bt_address, bt_address,
offset_size, offset_size,
element_size as u32, element_size as u32,
@@ -1146,8 +1332,9 @@ pub fn build_chunked_data_from_precompressed_libver(
} }
}; };
Ok(ChunkedDataResult { Ok(ChunkedPlacement {
data_bytes: data_buf, pieces: blob.pieces,
len: blob.len,
layout_message, layout_message,
pipeline_message: pre.pipeline_message.clone(), pipeline_message: pre.pipeline_message.clone(),
}) })
@@ -1156,11 +1343,11 @@ pub fn build_chunked_data_from_precompressed_libver(
/// Lay out precompressed chunks at `base_address` followed by a version-1 /// 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 /// B-tree chunk index, with a version-3 layout message: what libhdf5 writes
/// for a chunked dataset under a low bound of 1.8. /// for a chunked dataset under a low bound of 1.8.
fn build_btree_v1_chunked_data( fn place_btree_v1_chunked_data(
pre: &PrecompressedChunks, pre: &PreparedChunks<'_>,
base_address: u64, base_address: u64,
maxshape: Option<&[u64]>, maxshape: Option<&[u64]>,
) -> Result<ChunkedDataResult, FormatError> { ) -> Result<ChunkedPlacement, FormatError> {
if let Some(ms) = maxshape { if let Some(ms) = maxshape {
let bad = |what: &str| FormatError::ChunkedReadError(format!("maxshape: {what}")); let bad = |what: &str| FormatError::ChunkedReadError(format!("maxshape: {what}"));
if ms.len() != pre.shape.len() { if ms.len() != pre.shape.len() {
@@ -1171,20 +1358,17 @@ fn build_btree_v1_chunked_data(
} }
} }
let offset_size: u8 = 8; let offset_size: u8 = 8;
let mut data_buf = Vec::new(); let mut blob = Pieces::default();
let mut entries = Vec::with_capacity(pre.chunks.len()); let mut entries = Vec::with_capacity(pre.chunks.len());
for (i, (_raw_size, stored, filter_mask)) in pre.chunks.iter().enumerate() { for (i, (_raw_size, stored, filter_mask)) in pre.chunks.iter().enumerate() {
let aligned_offset = align_to_cache_line(data_buf.len()); blob.align();
if aligned_offset > data_buf.len() {
data_buf.resize(aligned_offset, 0u8);
}
entries.push(btree_v1_write::ChunkEntry { entries.push(btree_v1_write::ChunkEntry {
scaled: scaled_coords(&pre.shape, &pre.chunk_dims, i), scaled: scaled_coords(&pre.shape, &pre.chunk_dims, i),
nbytes: stored.len() as u64, nbytes: stored.len() as u64,
filter_mask: *filter_mask, filter_mask: *filter_mask,
address: base_address + data_buf.len() as u64, address: base_address + blob.len,
}); });
data_buf.extend_from_slice(stored); blob.chunk(i, stored.len());
} }
let element_size = u32::try_from(pre.element_size) let element_size = u32::try_from(pre.element_size)
.map_err(|_| FormatError::Overflow("element size".into()))?; .map_err(|_| FormatError::Overflow("element size".into()))?;
@@ -1193,11 +1377,8 @@ fn build_btree_v1_chunked_data(
let btree_address = if entries.is_empty() { let btree_address = if entries.is_empty() {
u64::MAX u64::MAX
} else { } else {
let aligned_idx = align_to_cache_line(data_buf.len()); blob.align();
if aligned_idx > data_buf.len() { let addr = base_address + blob.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( let tree = btree_v1_write::build_chunk_btree_v1_at(
&entries, &entries,
&pre.chunk_dims, &pre.chunk_dims,
@@ -1205,13 +1386,14 @@ fn build_btree_v1_chunked_data(
addr, addr,
offset_size, offset_size,
)?; )?;
data_buf.extend_from_slice(&tree); blob.bytes(tree);
addr addr
}; };
let layout_message = let layout_message =
serialize_v3_chunked(&pre.chunk_dims, btree_address, offset_size, element_size)?; serialize_v3_chunked(&pre.chunk_dims, btree_address, offset_size, element_size)?;
Ok(ChunkedDataResult { Ok(ChunkedPlacement {
data_bytes: data_buf, pieces: blob.pieces,
len: blob.len,
layout_message, layout_message,
pipeline_message: pre.pipeline_message.clone(), pipeline_message: pre.pipeline_message.clone(),
}) })
@@ -1424,7 +1606,7 @@ fn build_btree_v2_chunk_index_at(
/// Serialize a v4 layout message for a version-2 B-tree chunk index. /// Serialize a v4 layout message for a version-2 B-tree chunk index.
fn serialize_v4_btree_v2( fn serialize_v4_btree_v2(
chunk_dims: &[u32], chunk_dims: &[u64],
btree_address: u64, btree_address: u64,
offset_size: u8, offset_size: u8,
element_size: u32, element_size: u32,
@@ -2184,7 +2366,7 @@ mod tests {
/// filtered index element stores the chunk's size in 8 bytes. /// filtered index element stores the chunk's size in 8 bytes.
#[test] #[test]
fn huge_chunk_layout_messages_and_index_elements() { fn huge_chunk_layout_messages_and_index_elements() {
let dims = [HUGE_DIM as u32]; let dims = [HUGE_DIM];
let parsed = |msg: &[u8]| { let parsed = |msg: &[u8]| {
assert_eq!(msg[0], 5, "layout message version"); assert_eq!(msg[0], 5, "layout message version");
// Class chunked, then (after the flags) 2 dimensions of 4 bytes. // Class chunked, then (after the flags) 2 dimensions of 4 bytes.
@@ -2197,7 +2379,7 @@ mod tests {
chunk_index_type, chunk_index_type,
.. ..
} => { } => {
assert_eq!(chunk_dimensions, vec![HUGE_DIM as u32, 8]); assert_eq!(chunk_dimensions, vec![HUGE_DIM, 8]);
chunk_index_type.unwrap() chunk_index_type.unwrap()
} }
other => panic!("{other:?}"), other => panic!("{other:?}"),
@@ -2245,22 +2427,89 @@ mod tests {
assert_eq!(u16::from_le_bytes([bt[10], bt[11]]), 28); assert_eq!(u16::from_le_bytes([bt[10], bt[11]]), 28);
} }
/// A chunk dimension of 2^32 or more (a 1-byte element, 2^32 + 7 per
/// chunk): the dimensions take 5 bytes each, as libhdf5 2.x writes
/// them (`(log2(dim) + 8) / 8`), and parse back whole.
#[test]
fn chunk_dimensions_of_2_pow_32_or_more_take_5_bytes() {
const DIM: u64 = (1 << 32) + 7;
let msg = serialize_v4_single_chunk(&[DIM], 0x800, None, None, 8, 1, 5);
assert_eq!(&msg[..5], &[5, 2, 0, 2, 5]);
assert_eq!(&msg[5..10], &[7, 0, 0, 0, 1]);
assert_eq!(&msg[10..15], &[1, 0, 0, 0, 0]);
assert_eq!(msg[15], 1, "Single Chunk index");
match DataLayout::parse(&msg, 8, 8).unwrap() {
DataLayout::Chunked {
chunk_dimensions, ..
} => assert_eq!(chunk_dimensions, [DIM, 1]),
other => panic!("{other:?}"),
}
// The largest dimension takes all 8 bytes.
let mut buf = Vec::new();
push_v4_chunk_dims(&mut buf, &[u64::MAX, 3], 4);
assert_eq!(buf[0], 8);
assert_eq!(buf.len(), 1 + 3 * 8);
// A version-3 layout has 4-byte dimensions: refused, not truncated.
assert!(serialize_v3_chunked(&[DIM], 0x800, 8, 1).is_err());
}
/// Which chunks are contiguous runs of the dataset's data (borrowed,
/// not copied, when unfiltered).
#[test]
fn contiguous_chunks_are_borrowed() {
// Rows of a 2-D dataset: chunk (2, 5) of shape (5, 5).
assert_eq!(contiguous_chunk(&[5, 5], &[2, 5], 4, 0), Some(0..40));
assert_eq!(contiguous_chunk(&[5, 5], &[2, 5], 4, 1), Some(40..80));
// The last row chunk runs past the edge: padded, so copied.
assert_eq!(contiguous_chunk(&[5, 5], &[2, 5], 4, 2), None);
// Columns split: not contiguous.
assert_eq!(contiguous_chunk(&[4, 6], &[2, 3], 1, 0), None);
// Leading dimensions of 1, then a partial row: contiguous.
assert_eq!(contiguous_chunk(&[2, 3, 8], &[1, 1, 4], 2, 5), Some(40..48));
// The whole dataset as one chunk.
assert_eq!(contiguous_chunk(&[3, 4], &[3, 4], 8, 0), Some(0..96));
let raw: Vec<u8> = (0..100).collect();
let whole = chunk_bytes_of(&raw, &[10, 10], &[10, 10], 1, 0);
assert!(matches!(whole, Cow::Borrowed(b) if b.as_ptr() == raw.as_ptr()));
let rows = chunk_bytes_of(&raw, &[10, 10], &[3, 10], 1, 1);
assert!(matches!(rows, Cow::Borrowed(b) if b == &raw[30..60]));
// Data shorter than the shape is copied (and zero-padded) as before.
let short = chunk_bytes_of(&raw[..50], &[10, 10], &[10, 10], 1, 0);
assert!(matches!(&short, Cow::Owned(v) if v.len() == 100));
// Borrowed or copied, the bytes are what extract_chunk gives.
for (shape, chunk) in [
(&[10u64, 10][..], &[3u64, 10][..]),
(&[10, 10], &[3, 4]),
(&[4, 5, 5], &[1, 5, 5]),
(&[4, 5, 5], &[1, 2, 5]),
] {
for i in 0..chunk_count(shape, chunk) {
let got = chunk_bytes_of(&raw, shape, chunk, 1, i);
assert_eq!(&got[..], &extract_chunk(&raw, shape, chunk, 1, i).1[..]);
}
}
}
#[test] #[test]
fn huge_chunks_refused_where_unsupported() { fn huge_chunks_refused_where_unsupported() {
// A chunk dimension of 2^32 or more. // A chunk dimension of 0.
assert!(matches!( assert!(matches!(
checked_chunk_dims(&[1 << 32], 1), checked_chunk_dims(&[4, 0], 1),
Err(FormatError::InvalidChunkDimensions(m)) if m.contains("2^32") Err(FormatError::InvalidChunkDimensions(_))
)); ));
// A chunk dimension of 2^32 or more is fine (on a 64-bit target).
#[cfg(target_pointer_width = "64")]
assert_eq!(
checked_chunk_dims(&[(1 << 32) + 5], 1).unwrap(),
(1 << 32) + 5
);
// A chunk size that overflows u64. // A chunk size that overflows u64.
assert!(matches!( assert!(matches!(
checked_chunk_dims(&[u32::MAX.into(), u32::MAX.into(), 2], 8), checked_chunk_dims(&[u32::MAX.into(), u32::MAX.into(), 2], 8),
Err(FormatError::Overflow(_)) Err(FormatError::Overflow(_))
)); ));
assert_eq!( assert_eq!(checked_chunk_dims(&[HUGE_DIM], 8).unwrap(), HUGE_BYTES);
checked_chunk_dims(&[HUGE_DIM], 8).unwrap(),
(vec![HUGE_DIM as u32], HUGE_BYTES)
);
// Filters that cannot take a chunk that large. // Filters that cannot take a chunk that large.
for plugin in [ for plugin in [
PluginFilter::Lzf, PluginFilter::Lzf,
+21 -21
View File
@@ -34,7 +34,7 @@ const MAX_LAYOUT_NDIMS: usize = 33;
/// (`H5O__layout_decode`): at most [`MAX_LAYOUT_NDIMS`], no dimension 0, and /// (`H5O__layout_decode`): at most [`MAX_LAYOUT_NDIMS`], no dimension 0, and
/// before version 4 at least one dataspace dimension plus the element size. /// before version 4 at least one dataspace dimension plus the element size.
/// A zero chunk dimension used to read the dataset as all fill values. /// A zero chunk dimension used to read the dataset as all fill values.
fn check_chunk_dims(dims: Vec<u32>, layout_version: u8) -> Result<Vec<u32>, FormatError> { fn check_chunk_dims(dims: Vec<u64>, layout_version: u8) -> Result<Vec<u64>, FormatError> {
if dims.len() > MAX_LAYOUT_NDIMS { if dims.len() > MAX_LAYOUT_NDIMS {
return Err(FormatError::InvalidChunkDimensions( return Err(FormatError::InvalidChunkDimensions(
"dimensionality is too large".into(), "dimensionality is too large".into(),
@@ -72,7 +72,7 @@ pub enum DataLayout {
/// Chunked: data stored in chunks via a B-tree. /// Chunked: data stored in chunks via a B-tree.
Chunked { Chunked {
/// Chunk dimension sizes. /// Chunk dimension sizes.
chunk_dimensions: Vec<u32>, chunk_dimensions: Vec<u64>,
/// B-tree address, or `None` if undefined. /// B-tree address, or `None` if undefined.
btree_address: Option<u64>, btree_address: Option<u64>,
/// Layout version (3 or 4). Version 1/2 messages (HDF5 1.4/1.6-era) /// Layout version (3 or 4). Version 1/2 messages (HDF5 1.4/1.6-era)
@@ -404,11 +404,11 @@ impl DataLayout {
_ => return Err(FormatError::InvalidLayoutClass(layout_class)), _ => return Err(FormatError::InvalidLayoutClass(layout_class)),
}; };
ensure_len(data, p, dimensionality * 4)?; ensure_len(data, p, dimensionality * 4)?;
let dims: Vec<u32> = data[p..p + dimensionality * 4] let dims: Vec<u64> = data[p..p + dimensionality * 4]
.as_chunks::<4>() .as_chunks::<4>()
.0 .0
.iter() .iter()
.map(|c| u32::from_le_bytes(*c)) .map(|c| u64::from(u32::from_le_bytes(*c)))
.collect(); .collect();
p += dimensionality * 4; p += dimensionality * 4;
match layout_class { match layout_class {
@@ -424,7 +424,7 @@ impl DataLayout {
1 => { 1 => {
let size = dims let size = dims
.iter() .iter()
.try_fold(1u64, |acc, &d| acc.checked_mul(d as u64)) .try_fold(1u64, |acc, &d| acc.checked_mul(d))
.ok_or_else(|| { .ok_or_else(|| {
FormatError::Overflow(format!("contiguous layout size {dims:?}")) FormatError::Overflow(format!("contiguous layout size {dims:?}"))
})?; })?;
@@ -489,7 +489,12 @@ impl DataLayout {
ensure_len(data, p, dimensionality * 4)?; ensure_len(data, p, dimensionality * 4)?;
let mut chunk_dimensions = Vec::with_capacity(dimensionality); let mut chunk_dimensions = Vec::with_capacity(dimensionality);
for _ in 0..dimensionality { for _ in 0..dimensionality {
let dim = u32::from_le_bytes([data[p], data[p + 1], data[p + 2], data[p + 3]]); let dim = u64::from(u32::from_le_bytes([
data[p],
data[p + 1],
data[p + 2],
data[p + 3],
]));
chunk_dimensions.push(dim); chunk_dimensions.push(dim);
p += 4; p += 4;
} }
@@ -565,14 +570,8 @@ impl DataLayout {
.iter() .iter()
.rev() .rev()
.fold(0u64, |acc, &b| (acc << 8) | u64::from(b)); .fold(0u64, |acc, &b| (acc << 8) | u64::from(b));
// Chunk dimensions are held as u32; HDF5 2.0 can write // HDF5 2.0 writes dimensions of 2^32 or more (layout
// larger ones (layout version 5), which are refused // version 5, chunks over 4 GiB) in 5 to 8 bytes.
// rather than truncated.
let val = u32::try_from(val).map_err(|_| {
FormatError::InvalidChunkDimensions(format!(
"chunk dimension {val} is larger than 2^32 - 1, which is not supported"
))
})?;
chunk_dimensions.push(val); chunk_dimensions.push(val);
p += dim_size_encoded_length; p += dim_size_encoded_length;
} }
@@ -858,7 +857,7 @@ mod tests {
.unwrap_or_else(|e| panic!("width {width}: {e:?}")); .unwrap_or_else(|e| panic!("width {width}: {e:?}"));
assert!( assert!(
matches!(&layout, DataLayout::Chunked { chunk_dimensions, .. } matches!(&layout, DataLayout::Chunked { chunk_dimensions, .. }
if chunk_dimensions.iter().map(|&d| u64::from(d)).eq(dims)), if chunk_dimensions[..] == dims),
"width {width}: {layout:?}" "width {width}: {layout:?}"
); );
} }
@@ -871,12 +870,13 @@ mod tests {
) )
); );
} }
// A dimension past u32 cannot be represented and is refused, not // HDF5 2.0 writes dimensions of 2^32 or more (up to 8 bytes each).
// truncated. for (width, dim) in [(5u8, (1u64 << 32) + 3), (8, u64::MAX)] {
assert!(matches!( assert!(matches!(
DataLayout::parse(&v4_chunked_msg(5, &[1 << 32, 8]), 8, 8), DataLayout::parse(&v4_chunked_msg(width, &[dim, 1]), 8, 8),
Err(FormatError::InvalidChunkDimensions(m)) if m.contains("2^32") Ok(DataLayout::Chunked { chunk_dimensions, .. }) if chunk_dimensions == [dim, 1]
)); ));
}
} }
#[test] #[test]
+1 -1
View File
@@ -14,7 +14,7 @@ use crate::chunked_write::{
/// Serialize a v4 Extensible Array layout message. /// Serialize a v4 Extensible Array layout message.
pub(crate) fn serialize_v4_extensible_array( pub(crate) fn serialize_v4_extensible_array(
chunk_dims: &[u32], chunk_dims: &[u64],
ea_address: u64, ea_address: u64,
offset_size: u8, offset_size: u8,
element_size: u32, element_size: u32,
@@ -454,7 +454,7 @@ pub fn read_extensible_array_chunks(
header: &ExtensibleArrayHeader, header: &ExtensibleArrayHeader,
dataset_dims: &[u64], dataset_dims: &[u64],
max_dims: Option<&[u64]>, max_dims: Option<&[u64]>,
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
element_size: u32, element_size: u32,
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
@@ -480,7 +480,7 @@ pub fn read_extensible_array_chunks_in<S: Storage + ?Sized>(
header: &ExtensibleArrayHeader, header: &ExtensibleArrayHeader,
dataset_dims: &[u64], dataset_dims: &[u64],
max_dims: Option<&[u64]>, max_dims: Option<&[u64]>,
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
element_size: u32, element_size: u32,
offset_size: u8, offset_size: u8,
_length_size: u8, _length_size: u8,
@@ -489,12 +489,12 @@ pub fn read_extensible_array_chunks_in<S: Storage + ?Sized>(
// Linear indexes follow the maximum dimensions, with the unlimited // Linear indexes follow the maximum dimensions, with the unlimited
// dimension swizzled to the slowest position (see `chunk_grid`). // dimension swizzled to the slowest position (see `chunk_grid`).
let dims_u64: Vec<u64> = chunk_dimensions.iter().map(|&d| d as u64).collect(); let grid = ChunkGrid::extensible_array(dataset_dims, max_dims, chunk_dimensions)?;
let grid = ChunkGrid::extensible_array(dataset_dims, max_dims, &dims_u64)?;
let grid = &grid; let grid = &grid;
let chunk_byte_size: u64 = let chunk_byte_size: u64 = chunk_dimensions
chunk_dimensions.iter().map(|&d| d as u64).product::<u64>() * element_size as u64; .iter()
.fold(u64::from(element_size), |acc, &d| acc.saturating_mul(d));
// Parse index block (EAIB): signature(4) + version(1) + client_id(1) // Parse index block (EAIB): signature(4) + version(1) + client_id(1)
// + header address(offset_size), then the inline elements, then the // + header address(offset_size), then the inline elements, then the
@@ -930,7 +930,7 @@ mod tests {
let header = ExtensibleArrayHeader::parse(&file_data, aehd_offset, os, ls).unwrap(); let header = ExtensibleArrayHeader::parse(&file_data, aehd_offset, os, ls).unwrap();
let ds_dims = vec![40u64]; // 2 chunks × 20 elements let ds_dims = vec![40u64]; // 2 chunks × 20 elements
let chunk_dims = vec![20u32]; let chunk_dims = vec![20u64];
let chunks = read_extensible_array_chunks( let chunks = read_extensible_array_chunks(
&file_data, &file_data,
&header, &header,
@@ -1057,7 +1057,7 @@ mod tests {
let file_data = build_inline_plus_data_blocks(); let file_data = build_inline_plus_data_blocks();
let header = ExtensibleArrayHeader::parse(&file_data, 0x100, os, ls).unwrap(); let header = ExtensibleArrayHeader::parse(&file_data, 0x100, os, ls).unwrap();
let ds_dims = vec![40u64]; let ds_dims = vec![40u64];
let chunk_dims = vec![10u32]; let chunk_dims = vec![10u64];
let chunks = read_extensible_array_chunks( let chunks = read_extensible_array_chunks(
&file_data, &file_data,
&header, &header,
+233 -58
View File
@@ -3,15 +3,14 @@
//! Produces valid HDF5 files with v3 superblock, v2 object headers, //! Produces valid HDF5 files with v3 superblock, v2 object headers,
//! link messages, contiguous datasets, inline and dense attributes. //! link messages, contiguous datasets, inline and dense attributes.
use crate::addr::saturating_usize; use crate::addr::{saturating_usize, to_usize};
#[cfg(not(feature = "std"))] #[cfg(not(feature = "std"))]
use alloc::{format, string::String, vec, vec::Vec}; use alloc::{format, string::String, vec, vec::Vec};
use crate::attribute::AttributeMessage; use crate::attribute::AttributeMessage;
use crate::btree_v2_write::{BTreeV2Params, build_btree_v2}; use crate::btree_v2_write::{BTreeV2Params, build_btree_v2};
use crate::chunked_write::{ use crate::chunked_write::{
ChunkOptions, PrecompressedChunks, build_chunked_data_from_precompressed_libver, ChunkOptions, Piece, PreparedChunks, place_chunked_data, prepare_chunks,
precompress_chunks,
}; };
use crate::data_layout::VdsMapping; use crate::data_layout::VdsMapping;
use crate::dataspace::{Dataspace, DataspaceType}; use crate::dataspace::{Dataspace, DataspaceType};
@@ -1612,7 +1611,30 @@ impl FileWriter {
self self
} }
/// The file, in memory. Its datasets' data is copied into the buffer
/// once (a dataset stored as one chunk of its own shape, or unfiltered
/// chunks that are contiguous runs of its data, are not copied before);
/// to write a file without holding it all, use [`Self::finish_with`].
pub fn finish(self) -> Result<Vec<u8>, FormatError> { pub fn finish(self) -> Result<Vec<u8>, FormatError> {
let mut buf = Vec::new();
self.finish_into(&mut buf)?;
Ok(buf)
}
/// The file, handed to `put` in order, a piece at a time, without
/// assembling it in memory: an unfiltered chunk or contiguous dataset
/// is passed straight from the data the dataset was given. Writing a
/// dataset stored as one unfiltered chunk this way needs little memory
/// beyond the dataset's own data. `put`'s first error stops the write
/// and is returned.
pub fn finish_with<E: From<FormatError>>(
self,
put: impl FnMut(&[u8]) -> Result<(), E>,
) -> Result<(), E> {
self.finish_into(&mut FnSink(put))
}
fn finish_into<S: Sink>(self, sink: &mut S) -> Result<(), S::Error> {
let page_size = self.page_size; let page_size = self.page_size;
if let Some(ps) = page_size if let Some(ps) = page_size
&& !(MIN_FILE_SPACE_PAGE_SIZE..=MAX_FILE_SPACE_PAGE_SIZE).contains(&ps) && !(MIN_FILE_SPACE_PAGE_SIZE..=MAX_FILE_SPACE_PAGE_SIZE).contains(&ps)
@@ -1620,7 +1642,8 @@ impl FileWriter {
return Err(FormatError::SerializationError(format!( return Err(FormatError::SerializationError(format!(
"file space page size {ps} is outside libhdf5's \ "file space page size {ps} is outside libhdf5's \
{MIN_FILE_SPACE_PAGE_SIZE}..={MAX_FILE_SPACE_PAGE_SIZE} bytes" {MIN_FILE_SPACE_PAGE_SIZE}..={MAX_FILE_SPACE_PAGE_SIZE} bytes"
))); ))
.into());
} }
let (low, high) = (self.low, self.high); let (low, high) = (self.low, self.high);
@@ -1774,15 +1797,16 @@ impl FileWriter {
}) })
.collect::<Result<_, _>>()?; .collect::<Result<_, _>>()?;
struct DataBlob { struct DataBlob<'a> {
data: Vec<u8>, data: Vec<u8>,
oh_bytes: Vec<u8>, oh_bytes: Vec<u8>,
/// Cached compressed chunks for chunked datasets; reused in Pass 2 /// Cached compressed chunks for chunked datasets; reused in Pass 2
/// to avoid re-compressing the same data. /// to avoid re-compressing the same data. Unfiltered chunks
precompressed: Option<PrecompressedChunks>, /// borrow the dataset's data where they can.
precompressed: Option<PreparedChunks<'a>>,
} }
let mut dummy_blobs: Vec<DataBlob> = Vec::new(); let mut dummy_blobs: Vec<DataBlob<'_>> = Vec::new();
let mut dummy_cursor = 0u64; let mut dummy_cursor = 0u64;
for (i, d) in all_ds.iter().enumerate() { for (i, d) in all_ds.iter().enumerate() {
let dense_blob = ds_dense[i] let dense_blob = ds_dense[i]
@@ -1820,21 +1844,16 @@ impl FileWriter {
.resolve_chunk_dims_for(&d.ds.dimensions, elem_size); .resolve_chunk_dims_for(&d.ds.dimensions, elem_size);
// Compress once in Pass 1; cache the result so Pass 2 can skip // Compress once in Pass 1; cache the result so Pass 2 can skip
// re-compression and just rebuild the index with real addresses. // re-compression and just rebuild the index with real addresses.
let pre = precompress_chunks( let pre = prepare_chunks(
&d.raw, &d.raw,
&d.ds.dimensions, &d.ds.dimensions,
&chunk_dims, &chunk_dims,
elem_size, elem_size,
&d.chunk_options, &d.chunk_options,
)?; )?;
let result = build_chunked_data_from_precompressed_libver( let result =
&pre, place_chunked_data(&pre, dummy_cursor, d.maxshape.as_deref(), low, high)?;
dummy_cursor, dummy_cursor += result.len;
d.maxshape.as_deref(),
low,
high,
)?;
dummy_cursor += result.data_bytes.len() as u64;
let oh = build_chunked_dataset_oh( let oh = build_chunked_dataset_oh(
&d.dt, &d.dt,
&d.ds, &d.ds,
@@ -1849,7 +1868,7 @@ impl FileWriter {
d.refcount, d.refcount,
)?; )?;
dummy_blobs.push(DataBlob { dummy_blobs.push(DataBlob {
data: result.data_bytes, data: Vec::new(),
oh_bytes: oh, oh_bytes: oh,
precompressed: Some(pre), precompressed: Some(pre),
}); });
@@ -1955,7 +1974,7 @@ impl FileWriter {
}) })
.collect::<Result<_, FormatError>>()?; .collect::<Result<_, FormatError>>()?;
let mut ds_blobs2: Vec<DataBlob> = Vec::new(); let mut ds_blobs2: Vec<OutBlob<'_>> = Vec::new();
let global_align_threshold = self.alignment_threshold; let global_align_threshold = self.alignment_threshold;
let global_align_bytes = self.alignment_bytes; let global_align_bytes = self.alignment_bytes;
for (i, d) in all_ds.iter().enumerate() { for (i, d) in all_ds.iter().enumerate() {
@@ -1977,26 +1996,21 @@ impl FileWriter {
&d.fill_message, &d.fill_message,
d.refcount, d.refcount,
)?; )?;
ds_blobs2.push(DataBlob { ds_blobs2.push(OutBlob {
data: gcol_bytes.clone(), data: vec![Out::Slice(gcol_bytes)],
oh_bytes: oh, oh_bytes: oh,
precompressed: None,
}); });
} else if is_chunked[i] { } else if is_chunked[i] {
let base_address = cursor2 as u64; let base_address = cursor2 as u64;
// Reuse precompressed chunks from Pass 1 — avoids re-compressing // Reuse precompressed chunks from Pass 1 — avoids re-compressing
// the same data a second time. // the same data a second time.
let result = build_chunked_data_from_precompressed_libver( let pre = dummy_blobs[i]
dummy_blobs[i] .precompressed
.precompressed .as_ref()
.as_ref() .expect("chunked dataset missing precompressed cache");
.expect("chunked dataset missing precompressed cache"), let result =
base_address, place_chunked_data(pre, base_address, d.maxshape.as_deref(), low, high)?;
d.maxshape.as_deref(), cursor2 += to_usize(result.len)?;
low,
high,
)?;
cursor2 += result.data_bytes.len();
let oh = build_chunked_dataset_oh( let oh = build_chunked_dataset_oh(
&d.dt, &d.dt,
&d.ds, &d.ds,
@@ -2010,11 +2024,16 @@ impl FileWriter {
&d.fill_message, &d.fill_message,
d.refcount, d.refcount,
)?; )?;
ds_blobs2.push(DataBlob { let data = result
data: result.data_bytes, .pieces
oh_bytes: oh, .into_iter()
precompressed: None, .map(|p| match p {
}); Piece::Zeros(n) => Out::Zeros(n),
Piece::Chunk(c) => Out::Slice(&pre.chunks[c].1),
Piece::Bytes(b) => Out::Bytes(b),
})
.collect();
ds_blobs2.push(OutBlob { data, oh_bytes: oh });
} else if is_compact[i] { } else if is_compact[i] {
// Compact: data is inline in the object header, no external blob // Compact: data is inline in the object header, no external blob
let oh = build_compact_dataset_oh( let oh = build_compact_dataset_oh(
@@ -2030,10 +2049,9 @@ impl FileWriter {
d.refcount, d.refcount,
layout_version, layout_version,
)?; )?;
ds_blobs2.push(DataBlob { ds_blobs2.push(OutBlob {
data: vec![], data: vec![],
oh_bytes: oh, oh_bytes: oh,
precompressed: None,
}); });
} else { } else {
// Determine alignment: per-dataset overrides global // Determine alignment: per-dataset overrides global
@@ -2060,13 +2078,10 @@ impl FileWriter {
d.refcount, d.refcount,
layout_version, layout_version,
)?; )?;
let mut data = vec![0u8; padding];
data.extend_from_slice(&d.raw);
cursor2 += d.raw.len(); cursor2 += d.raw.len();
ds_blobs2.push(DataBlob { ds_blobs2.push(OutBlob {
data, data: vec![Out::Zeros(padding), Out::Slice(&d.raw)],
oh_bytes: oh, oh_bytes: oh,
precompressed: None,
}); });
} }
} }
@@ -2080,7 +2095,8 @@ impl FileWriter {
cursor2 = cursor2.next_multiple_of(ps as usize); cursor2 = cursor2.next_multiple_of(ps as usize);
} }
let eof_addr2 = cursor2 as u64; let eof_addr2 = cursor2 as u64;
let mut buf = Vec::with_capacity(cursor2); sink.reserve(cursor2);
let mut out = Counted { sink, written: 0 };
let sb = Superblock { let sb = Superblock {
version: superblock_version, version: superblock_version,
@@ -2103,9 +2119,9 @@ impl FileWriter {
checksum: None, checksum: None,
page_size: None, page_size: None,
}; };
buf.extend_from_slice(&sb.serialize()); out.put(&sb.serialize())?;
if let Some(ref ext) = sb_ext { if let Some(ref ext) = sb_ext {
buf.extend_from_slice(ext); out.put(ext)?;
} }
// Group OHs + dense blobs (link blob, then attr blob, matching pass 2) // Group OHs + dense blobs (link blob, then attr blob, matching pass 2)
@@ -2132,32 +2148,120 @@ impl FileWriter {
g.refcount, g.refcount,
)?; )?;
debug_assert_eq!(oh.len(), group_oh_sizes[gi]); debug_assert_eq!(oh.len(), group_oh_sizes[gi]);
debug_assert_eq!(buf.len() as u64, group_addrs2[gi]); debug_assert_eq!(out.written as u64, group_addrs2[gi]);
buf.extend_from_slice(&oh); out.put(&oh)?;
if let Some(ref b) = link_blob { if let Some(ref b) = link_blob {
buf.extend_from_slice(&b.blob); out.put(&b.blob)?;
} }
if let Some(ref blob) = group_dense_blobs[gi] { if let Some(ref blob) = group_dense_blobs[gi] {
buf.extend_from_slice(&blob.blob); out.put(&blob.blob)?;
} }
} }
// Dataset OHs + dense blobs // Dataset OHs + dense blobs
for (i, blob) in ds_blobs2.iter().enumerate() { for (i, blob) in ds_blobs2.iter().enumerate() {
buf.extend_from_slice(&blob.oh_bytes); out.put(&blob.oh_bytes)?;
if let Some(ref dense) = ds_dense_blobs[i] { if let Some(ref dense) = ds_dense_blobs[i] {
buf.extend_from_slice(&dense.blob); out.put(&dense.blob)?;
} }
} }
// Data // Data
for blob in &ds_blobs2 { for blob in &ds_blobs2 {
buf.extend_from_slice(&blob.data); for piece in &blob.data {
match piece {
Out::Bytes(b) => out.put(b)?,
Out::Slice(s) => out.put(s)?,
Out::Zeros(n) => out.zeros(*n)?,
}
}
} }
debug_assert_eq!(buf.len(), data_end); debug_assert_eq!(out.written, data_end);
buf.resize(cursor2, 0); out.zeros(cursor2 - data_end)?;
Ok(buf) Ok(())
}
}
/// A piece of a dataset's data as [`FileWriter`] writes it out.
enum Out<'a> {
Bytes(Vec<u8>),
/// Borrowed from the dataset's data or its prepared chunks.
Slice(&'a [u8]),
Zeros(usize),
}
/// A dataset's object header and the pieces of its data.
struct OutBlob<'a> {
oh_bytes: Vec<u8>,
data: Vec<Out<'a>>,
}
/// Where [`FileWriter`] puts a file's bytes, in order.
trait Sink {
type Error: From<FormatError>;
/// The file will be `_len` bytes long.
fn reserve(&mut self, _len: usize) {}
fn put(&mut self, bytes: &[u8]) -> Result<(), Self::Error>;
fn zeros(&mut self, n: usize) -> Result<(), Self::Error> {
const ZEROS: [u8; 4096] = [0; 4096];
let mut left = n;
while left > 0 {
let k = left.min(ZEROS.len());
self.put(&ZEROS[..k])?;
left -= k;
}
Ok(())
}
}
impl Sink for Vec<u8> {
type Error = FormatError;
fn reserve(&mut self, len: usize) {
self.reserve_exact(len.saturating_sub(self.len()));
}
fn put(&mut self, bytes: &[u8]) -> Result<(), FormatError> {
self.extend_from_slice(bytes);
Ok(())
}
fn zeros(&mut self, n: usize) -> Result<(), FormatError> {
self.resize(self.len() + n, 0);
Ok(())
}
}
/// A [`Sink`] calling a function ([`FileWriter::finish_with`]).
struct FnSink<F>(F);
impl<E: From<FormatError>, F: FnMut(&[u8]) -> Result<(), E>> Sink for FnSink<F> {
type Error = E;
fn put(&mut self, bytes: &[u8]) -> Result<(), E> {
(self.0)(bytes)
}
}
/// A [`Sink`] and the number of bytes put into it.
struct Counted<'s, S> {
sink: &'s mut S,
written: usize,
}
impl<S: Sink> Counted<'_, S> {
fn put(&mut self, bytes: &[u8]) -> Result<(), S::Error> {
self.written += bytes.len();
self.sink.put(bytes)
}
fn zeros(&mut self, n: usize) -> Result<(), S::Error> {
self.written += n;
self.sink.zeros(n)
} }
} }
@@ -3015,4 +3119,75 @@ mod tests {
assert_eq!(Superblock::parse(&bytes, 0).unwrap().version, 3); assert_eq!(Superblock::parse(&bytes, 0).unwrap().version, 3);
assert_eq!(layout_of(&bytes, "x")[0], 3); assert_eq!(layout_of(&bytes, "x")[0], 3);
} }
/// `finish_with` hands out exactly the bytes `finish` returns, for every
/// storage (contiguous, compact, chunked filtered and not, chunks
/// borrowed and copied, a version-1 B-tree, a paged file), and stops at
/// the sink's first error.
#[test]
fn finish_with_matches_finish() {
let build = |low: LibVer, paged: bool| {
let mut fw = FileWriter::new();
fw.libver_bounds(low, LibVer::Latest);
if paged {
fw.with_page_size(4096);
}
let v: Vec<f64> = (0..300).map(f64::from).collect();
fw.create_dataset("contig").with_f64_data(&v);
fw.create_dataset("compact")
.with_f64_data(&[1.0, 2.0])
.compact();
fw.create_dataset("one_chunk")
.with_f64_data(&v)
.with_shape(&[30, 10])
.with_chunks(&[30, 10]);
fw.create_dataset("rows")
.with_f64_data(&v)
.with_shape(&[30, 10])
.with_chunks(&[7, 10]);
fw.create_dataset("tiles")
.with_f64_data(&v)
.with_shape(&[30, 10])
.with_chunks(&[7, 3]);
fw.create_dataset("deflated")
.with_f64_data(&v)
.with_chunks(&[64])
.with_shuffle()
.with_deflate(4);
fw.create_dataset("grows")
.with_u8_data_owned((0..=255).collect())
.with_chunks(&[100])
.with_maxshape(&[u64::MAX]);
fw
};
for (low, paged) in [
(LibVer::V18, false),
(LibVer::V110, false),
(LibVer::Latest, true),
] {
let bytes = build(low, paged).finish().unwrap();
let mut streamed = Vec::new();
build(low, paged)
.finish_with(|b: &[u8]| {
streamed.extend_from_slice(b);
Ok::<(), FormatError>(())
})
.unwrap();
assert_eq!(streamed, bytes, "{low:?} paged {paged}");
}
let mut calls = 0;
let err = build(LibVer::V110, false)
.finish_with(|_| {
calls += 1;
if calls == 3 {
Err(FormatError::SerializationError("full".into()))
} else {
Ok(())
}
})
.unwrap_err();
assert_eq!(err, FormatError::SerializationError("full".into()));
assert_eq!(calls, 3);
}
} }
+9 -9
View File
@@ -155,7 +155,7 @@ pub fn read_fixed_array_chunks(
header: &FixedArrayHeader, header: &FixedArrayHeader,
dataset_dims: &[u64], dataset_dims: &[u64],
max_dims: Option<&[u64]>, max_dims: Option<&[u64]>,
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
element_size: u32, element_size: u32,
offset_size: u8, offset_size: u8,
length_size: u8, length_size: u8,
@@ -180,7 +180,7 @@ pub fn read_fixed_array_chunks_in<S: Storage + ?Sized>(
header: &FixedArrayHeader, header: &FixedArrayHeader,
dataset_dims: &[u64], dataset_dims: &[u64],
max_dims: Option<&[u64]>, max_dims: Option<&[u64]>,
chunk_dimensions: &[u32], chunk_dimensions: &[u64],
element_size: u32, element_size: u32,
offset_size: u8, offset_size: u8,
_length_size: u8, _length_size: u8,
@@ -227,11 +227,11 @@ pub fn read_fixed_array_chunks_in<S: Storage + ?Sized>(
// The index is laid out over the chunk grid of the *maximum* dimensions // The index is laid out over the chunk grid of the *maximum* dimensions
// (row-major), so a dataset smaller than its maxshape has gaps. // (row-major), so a dataset smaller than its maxshape has gaps.
let dims_u64: Vec<u64> = chunk_dimensions.iter().map(|&d| d as u64).collect(); let grid = ChunkGrid::fixed_array(dataset_dims, max_dims, chunk_dimensions)?;
let grid = ChunkGrid::fixed_array(dataset_dims, max_dims, &dims_u64)?;
let chunk_byte_size: u64 = let chunk_byte_size: u64 = chunk_dimensions
chunk_dimensions.iter().map(|&d| d as u64).product::<u64>() * element_size as u64; .iter()
.fold(u64::from(element_size), |acc, &d| acc.saturating_mul(d));
let mut chunks = Vec::new(); let mut chunks = Vec::new();
// `rel` is relative to the data block, whose bytes are in `w`. // `rel` is relative to the data block, whose bytes are in `w`.
@@ -670,7 +670,7 @@ mod tests {
let header = let header =
FixedArrayHeader::parse(&file_data, fahd_offset, offset_size, length_size).unwrap(); FixedArrayHeader::parse(&file_data, fahd_offset, offset_size, length_size).unwrap();
let ds_dims = vec![100u64]; let ds_dims = vec![100u64];
let chunk_dims = vec![20u32]; let chunk_dims = vec![20u64];
let chunks = read_fixed_array_chunks( let chunks = read_fixed_array_chunks(
&file_data, &file_data,
&header, &header,
@@ -747,7 +747,7 @@ mod tests {
let header = let header =
FixedArrayHeader::parse(&file_data, fahd_offset, offset_size, length_size).unwrap(); FixedArrayHeader::parse(&file_data, fahd_offset, offset_size, length_size).unwrap();
let ds_dims = vec![60u64]; let ds_dims = vec![60u64];
let chunk_dims = vec![20u32]; let chunk_dims = vec![20u64];
let chunks = read_fixed_array_chunks( let chunks = read_fixed_array_chunks(
&file_data, &file_data,
&header, &header,
@@ -848,7 +848,7 @@ mod tests {
assert_eq!(header.num_elements, 11); assert_eq!(header.num_elements, 11);
let ds_dims = vec![11u64 * 20]; let ds_dims = vec![11u64 * 20];
let chunk_dims = vec![20u32]; let chunk_dims = vec![20u64];
let chunks = read_fixed_array_chunks( let chunks = read_fixed_array_chunks(
&file_data, &file_data,
&header, &header,
+7 -1
View File
@@ -663,11 +663,17 @@ impl DatasetBuilder {
} }
pub fn with_u8_data(&mut self, data: &[u8]) -> &mut Self { pub fn with_u8_data(&mut self, data: &[u8]) -> &mut Self {
self.with_u8_data_owned(data.to_vec())
}
/// [`Self::with_u8_data`] without the copy: the builder takes the
/// vector. For large data this halves the memory a write needs.
pub fn with_u8_data_owned(&mut self, data: Vec<u8>) -> &mut Self {
self.datatype = Some(make_u8_type()); self.datatype = Some(make_u8_type());
self.data = Some(data.to_vec());
if self.shape.is_none() { if self.shape.is_none() {
self.shape = Some(vec![data.len() as u64]); self.shape = Some(vec![data.len() as u64]);
} }
self.data = Some(data);
self self
} }
+1 -6
View File
@@ -189,12 +189,7 @@ pub(crate) fn chunk_shape(file: &File, hdr: &ObjectHeader, rank: usize) -> Optio
{ {
clawhdf5_format::data_layout::DataLayout::Chunked { clawhdf5_format::data_layout::DataLayout::Chunked {
chunk_dimensions, .. chunk_dimensions, ..
} if chunk_dimensions.len() >= rank => Some( } if chunk_dimensions.len() >= rank => Some(chunk_dimensions[..rank].to_vec()),
chunk_dimensions[..rank]
.iter()
.map(|&d| u64::from(d))
.collect(),
),
_ => None, _ => None,
} }
} }
+4 -4
View File
@@ -873,9 +873,9 @@ impl Checker<'_> {
self.counts.chunk_index_checksummed += 1; self.counts.chunk_index_checksummed += 1;
} }
let filtered = matches!(&info.filters, Ok(Some(p)) if !p.filters.is_empty()); let filtered = matches!(&info.filters, Ok(Some(p)) if !p.filters.is_empty());
let chunk_bytes = cdims.iter().try_fold(u64::from(dt.type_size()), |a, &d| { let chunk_bytes = cdims
a.checked_mul(u64::from(d)) .iter()
}); .try_fold(u64::from(dt.type_size()), |a, &d| a.checked_mul(d));
let max = ds.max_dimensions.clone(); let max = ds.max_dimensions.clone();
let mut seen: HashSet<Vec<u64>> = HashSet::with_capacity(chunks.len().min(1 << 20)); let mut seen: HashSet<Vec<u64>> = HashSet::with_capacity(chunks.len().min(1 << 20));
let mut reported = 0usize; let mut reported = 0usize;
@@ -893,7 +893,7 @@ impl Checker<'_> {
)); ));
} else { } else {
for (i, (&o, &cd)) in c.offsets.iter().zip(cdims).enumerate() { for (i, (&o, &cd)) in c.offsets.iter().zip(cdims).enumerate() {
if o % u64::from(cd) != 0 { if o % cd != 0 {
bad.push(format!( bad.push(format!(
"offset {o} in dimension {i} is not a multiple of the chunk size {cd}" "offset {o} in dimension {i} is not a multiple of the chunk size {cd}"
)); ));
+1 -1
View File
@@ -292,7 +292,7 @@ impl Ls<'_> {
.unwrap_or(0); .unwrap_or(0);
let bytes = dims let bytes = dims
.iter() .iter()
.try_fold(esize, |a, &d| a.checked_mul(u64::from(d))) .try_fold(esize, |a, &d| a.checked_mul(d))
.map(|b| b.to_string()) .map(|b| b.to_string())
.unwrap_or_else(|| "?".into()); .unwrap_or_else(|| "?".into());
writeln!( writeln!(
@@ -0,0 +1,46 @@
//! Write one dataset of `u8` stored as a single chunk, for measuring the
//! writer's peak memory (`/usr/bin/time -v`):
//!
//! ```text
//! cargo run --release -p clawhdf5 --example write_one_chunk -- BYTES OUT.h5 [slice|owned] [deflate|lz4|zstd]
//! ```
//!
//! `slice` passes the data by reference (`with_u8_data`, which copies it),
//! `owned` hands the builder the vector (`with_u8_data_owned`). The data is
//! `i % 251`, one byte per element; with no filter the chunk is stored raw.
fn main() {
let args: Vec<String> = std::env::args().collect();
let n: u64 = args[1].parse().expect("BYTES");
let out = &args[2];
let mode = args.get(3).map_or("slice", String::as_str);
let filter = args.get(4).map(String::as_str);
let data: Vec<u8> = (0..n).map(|i| (i % 251) as u8).collect();
let mut b = clawhdf5::FileBuilder::new();
let ds = b.create_dataset("x");
match mode {
"slice" => {
ds.with_u8_data(&data);
}
"owned" => {
ds.with_u8_data_owned(data);
}
m => panic!("mode {m}"),
}
ds.with_chunks(&[n]);
match filter {
None => {}
Some("deflate") => {
ds.with_deflate(1);
}
#[cfg(feature = "lz4")]
Some("lz4") => {
ds.with_lz4();
}
#[cfg(feature = "zstd")]
Some("zstd") => {
ds.with_zstd(1);
}
Some(f) => panic!("filter {f} (not built in?)"),
}
b.write(out).unwrap();
}
+1 -4
View File
@@ -408,10 +408,7 @@ impl<'t> ChunkedEdit<'t> {
"chunk rank differs from the dataspace".into(), "chunk rank differs from the dataspace".into(),
)); ));
} }
let cd: Vec<u64> = chunk_dimensions[..rank] let cd: Vec<u64> = chunk_dimensions[..rank].to_vec();
.iter()
.map(|&c| u64::from(c))
.collect();
// Chunks of 4 GiB or more (HDF5 2.0 writes them with layout message // 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 // version 5) are read, but not rewritten: the editor holds a chunk
// it rewrites in memory, compressed and not. Writing values, and a // it rewrites in memory, compressed and not. Writing values, and a
+26 -8
View File
@@ -137,10 +137,15 @@ impl FileBuilder {
Ok(self.writer.finish()?) Ok(self.writer.finish()?)
} }
/// Serialize and write the file to the given path. /// Serialize and write the file to the given path. The file is written
/// as it is serialized, not assembled in memory first: unfiltered
/// chunks and contiguous data go straight from the datasets' data.
pub fn write<P: AsRef<std::path::Path>>(self, path: P) -> Result<(), Error> { pub fn write<P: AsRef<std::path::Path>>(self, path: P) -> Result<(), Error> {
let bytes = self.finish()?; use std::io::Write;
write_file_atomically(path.as_ref(), &bytes).map_err(Error::Io) write_file_atomically_with(path.as_ref(), |f| {
self.writer
.finish_with(|bytes| f.write_all(bytes).map_err(Error::Io))
})
} }
} }
@@ -259,8 +264,19 @@ pub fn create_datasets_parallel(specs: Vec<DatasetSpec>) -> Result<Vec<u8>, Erro
/// file or the complete new one — never a truncated mix. `std::fs::write` /// file or the complete new one — never a truncated mix. `std::fs::write`
/// truncates the destination first, so dying mid-write used to destroy the /// truncates the destination first, so dying mid-write used to destroy the
/// existing file. /// existing file.
#[cfg(test)]
fn write_file_atomically(path: &std::path::Path, bytes: &[u8]) -> std::io::Result<()> { fn write_file_atomically(path: &std::path::Path, bytes: &[u8]) -> std::io::Result<()> {
use std::io::Write; use std::io::Write;
write_file_atomically_with(path, |f| f.write_all(bytes))
}
/// [`write_file_atomically`] with the contents written by `write` (into a
/// buffered writer over the temporary file).
fn write_file_atomically_with<E: From<std::io::Error>>(
path: &std::path::Path,
write: impl FnOnce(&mut std::io::BufWriter<std::fs::File>) -> Result<(), E>,
) -> Result<(), E> {
use std::io::Write;
// Same directory as the target, so the rename stays on one filesystem. // Same directory as the target, so the rename stays on one filesystem.
let mut tmp_name = path let mut tmp_name = path
@@ -272,11 +288,13 @@ fn write_file_atomically(path: &std::path::Path, bytes: &[u8]) -> std::io::Resul
tmp_name.push(format!(".tmp-{}", std::process::id())); tmp_name.push(format!(".tmp-{}", std::process::id()));
let tmp_path = path.with_file_name(tmp_name); let tmp_path = path.with_file_name(tmp_name);
let result = (|| { let result = (|| -> Result<(), E> {
let mut f = std::fs::File::create(&tmp_path)?; let mut f = std::io::BufWriter::with_capacity(1 << 20, std::fs::File::create(&tmp_path)?);
f.write_all(bytes)?; write(&mut f)?;
f.sync_all()?; f.flush()?;
std::fs::rename(&tmp_path, path) f.get_ref().sync_all()?;
std::fs::rename(&tmp_path, path)?;
Ok(())
})(); })();
if result.is_err() { if result.is_err() {
let _ = std::fs::remove_file(&tmp_path); let _ = std::fs::remove_file(&tmp_path);
+57 -2
View File
@@ -2,6 +2,12 @@
python gen_huge_chunks.py filtered OUT.h5 # huge_chunks_filtered.h5 python gen_huge_chunks.py filtered OUT.h5 # huge_chunks_filtered.h5
python gen_huge_chunks.py unfiltered OUT.h5 # generated at test time python gen_huge_chunks.py unfiltered OUT.h5 # generated at test time
python gen_huge_chunks.py dims OUT.h5 # huge_chunk_dims.h5
python gen_huge_chunks.py dims-unfiltered OUT.h5 # generated at test time
`dims` and `dims-unfiltered` (see `make_dims`) hold `u1` datasets whose
chunk dimension is N8 = 2**32 + 7: a chunk dimension of 2**32 or more,
which a layout message of version 5 stores in 5 bytes.
libhdf5 2.x writes a chunk of more than 0xFFFFFFFF bytes with layout message libhdf5 2.x writes a chunk of more than 0xFFFFFFFF bytes with layout message
version 5 (`H5D__chunk_construct`: "chunk size > 4GB requires version 5 (`H5D__chunk_construct`: "chunk size > 4GB requires
@@ -38,7 +44,8 @@ allocated early (see `make`).
libhdf5 holds a whole filtered chunk in memory while it writes it, so the libhdf5 holds a whole filtered chunk in memory while it writes it, so the
`filtered` run needs about 4 GiB of memory and 100 s (tank). `filtered` run needs about 4 GiB of memory and 100 s (tank).
Generated 2026-09-28 with h5py 3.16.0 (libhdf5 2.0.0). huge_chunks_filtered.h5 was generated 2026-09-28 and huge_chunk_dims.h5
2026-09-29, both with h5py 3.16.0 (libhdf5 2.0.0).
""" """
import sys import sys
@@ -92,12 +99,60 @@ def make(fid, name, index, filtered):
dsid.close() dsid.close()
N8 = 2**32 + 7
def make_dims(fid, name, index, filtered):
"""A `u1` dataset whose chunks are N8 elements (a chunk dimension of
2**32 or more). `d[0:10] = 0..9` and `d[N8-10:N8] = 100..109`; where
there is a second chunk (`farray`, `earray`: shape N8 + 10),
`d[N8:N8+10] = 200..209`.
filtered (`dims`, committed): deflate level 9 twice, fill value 7;
`single` (Single Chunk) and `earray` (Extensible Array, maxshape None).
Writing it holds a 4 GiB chunk: 4.1 GiB peak resident memory and about
a minute (tank, 2026-09-29).
unfiltered (`dims-unfiltered`): fill time "never", sparse; `single`
(allocated early, as in `make`) and `farray` (Fixed Array).
"""
dcpl = h5py.h5p.create(h5py.h5p.DATASET_CREATE)
if index == "single":
shape, maxshape = (N8,), (N8,)
elif index == "farray":
shape, maxshape = (N8 + 10,), (N8 + 10,)
elif index == "earray":
shape, maxshape = (N8 + 10,), (UNLIM,)
dcpl.set_chunk((N8,))
if filtered:
dcpl.set_deflate(9)
dcpl.set_deflate(9)
dcpl.set_fill_value(np.array(7, dtype="u1"))
else:
dcpl.set_fill_time(h5py.h5d.FILL_TIME_NEVER)
if index == "single":
dcpl.set_alloc_time(h5py.h5d.ALLOC_TIME_EARLY)
space = h5py.h5s.create_simple(shape, maxshape)
dsid = h5py.h5d.create(fid, name.encode(), h5py.h5t.STD_U8LE, space, dcpl=dcpl)
ds = h5py.Dataset(dsid)
ds[0:10] = np.arange(10, dtype="u1")
ds[N8 - 10 : N8] = np.arange(100, 110, dtype="u1")
if shape[0] > N8:
ds[N8 : N8 + 10] = np.arange(200, 210, dtype="u1")
dsid.close()
def main(): def main():
mode, out = sys.argv[1], sys.argv[2] mode, out = sys.argv[1], sys.argv[2]
filtered = {"filtered": True, "unfiltered": False}[mode]
fapl = h5py.h5p.create(h5py.h5p.FILE_ACCESS) fapl = h5py.h5p.create(h5py.h5p.FILE_ACCESS)
fapl.set_libver_bounds(h5py.h5f.LIBVER_EARLIEST, h5py.h5f.LIBVER_V200) fapl.set_libver_bounds(h5py.h5f.LIBVER_EARLIEST, h5py.h5f.LIBVER_V200)
fid = h5py.h5f.create(out.encode(), h5py.h5f.ACC_TRUNC, fapl=fapl) fid = h5py.h5f.create(out.encode(), h5py.h5f.ACC_TRUNC, fapl=fapl)
if mode in ("dims", "dims-unfiltered"):
filtered = mode == "dims"
for index in ["single", "earray"] if filtered else ["single", "farray"]:
make_dims(fid, index, index, filtered)
fid.close()
return
filtered = {"filtered": True, "unfiltered": False}[mode]
indexes = ["single", "farray", "earray", "btree2"] indexes = ["single", "farray", "earray", "btree2"]
if not filtered: if not filtered:
indexes.insert(1, "implicit") indexes.insert(1, "implicit")
Binary file not shown.
+348 -1
View File
@@ -81,8 +81,15 @@ fn layout_of(
let lm = msg(MessageType::DataLayout); let lm = msg(MessageType::DataLayout);
let layout = DataLayout::parse(&lm.data, sb.offset_size, sb.length_size).unwrap(); let layout = DataLayout::parse(&lm.data, sb.offset_size, sb.length_size).unwrap();
let space = Dataspace::parse(&msg(MessageType::Dataspace).data, sb.length_size).unwrap(); let space = Dataspace::parse(&msg(MessageType::Dataspace).data, sb.length_size).unwrap();
// The element size is the layout's last dimension.
let elem = match &layout {
DataLayout::Chunked {
chunk_dimensions, ..
} => *chunk_dimensions.last().unwrap() as usize,
_ => 8,
};
let (chunks, _) = let (chunks, _) =
list_chunks(data, &layout, &space, 8, sb.offset_size, sb.length_size).unwrap(); list_chunks(data, &layout, &space, elem, sb.offset_size, sb.length_size).unwrap();
(lm.data[0], layout, space, chunks) (lm.data[0], layout, space, chunks)
} }
@@ -516,3 +523,343 @@ print("ok")
} }
} }
} }
// ---- Chunk dimensions of 2^32 or more ----
/// Elements per chunk of the `u8` datasets below: a chunk dimension of
/// 2^32 or more (4 GiB + 7 bytes), which a version-5 layout message stores
/// in 5 bytes.
const N8: u64 = (1 << 32) + 7;
fn dims_fixture() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/huge_chunk_dims.h5")
}
/// `u8` values of `sel` in dataset `name` of `file`.
fn sel_u8(file: &File, name: &str, start: u64, count: u64) -> Vec<u8> {
let ds = file.dataset(name).unwrap();
let s = Selection::Hyperslab {
start: vec![start],
stride: vec![1],
count: vec![count],
block: vec![1],
};
ds.read_selection(&s).unwrap()
}
fn u8s(range: std::ops::Range<u8>) -> Vec<u8> {
range.collect()
}
/// `fixtures/huge_chunk_dims.h5` (libhdf5 2.0.0 through h5py 3.16,
/// `gen_huge_chunks.py dims`): `u8` datasets chunked by N8 = 2^32 + 7, in a
/// Single Chunk and an Extensible Array index. Their layouts parse to the
/// whole dimension (it was refused: chunk dimensions were `u32`) and their
/// chunks list at the right origins.
#[test]
fn huge_chunk_dims_list() {
let data = std::fs::read(dims_fixture()).unwrap();
for (name, index, shape, origins) in [
("single", 1u8, N8, &[0u64][..]),
("earray", 4, N8 + 10, &[0, N8]),
] {
let (version, layout, space, mut chunks) = layout_of(&data, name);
assert_eq!(version, 5, "{name}");
let DataLayout::Chunked {
chunk_dimensions,
chunk_index_type,
..
} = &layout
else {
panic!("{name}: {layout:?}");
};
assert_eq!(chunk_dimensions, &[N8, 1], "{name}");
assert_eq!(*chunk_index_type, Some(index), "{name}");
assert_eq!(space.dimensions, [shape], "{name}");
chunks.sort_by_key(|c| c.offsets[0]);
let got: Vec<u64> = chunks.iter().map(|c| c.offsets[0]).collect();
assert_eq!(got, origins, "{name}");
}
}
/// Decoding the fixture's chunks (4 GiB each, one at a time): written
/// values, the fill value (7) around them, and the second chunk.
#[test]
fn huge_chunk_dims_read() {
if !heavy() {
return;
}
let file = File::open(dims_fixture()).unwrap();
let mut first = u8s(0..10);
first.extend([7, 7]);
let mut last = vec![7, 7];
last.extend(u8s(100..110));
for name in ["single", "earray"] {
assert_eq!(sel_u8(&file, name, 0, 12), first, "{name}");
assert_eq!(sel_u8(&file, name, N8 - 12, 12), last, "{name}");
}
let mut across = u8s(108..110);
across.extend(u8s(200..210));
assert_eq!(sel_u8(&file, "earray", N8 - 2, 12), across);
}
/// Unfiltered chunks of N8 bytes that libhdf5 2.0.0 wrote (a sparse file,
/// `gen_huge_chunks.py dims-unfiltered`), in a Single Chunk and a Fixed
/// Array index: read through a memory map, and through positioned reads
/// that fetch only the rows selected.
#[test]
fn huge_chunk_dims_unfiltered_read() {
if !heavy() {
return;
}
let dir = scratch();
let Some(path) = generate(dir.path(), "dims-unfiltered") else {
return;
};
let check = |file: &File| {
let mut first = u8s(0..10);
first.extend([0, 0]);
let mut last = vec![0, 0];
last.extend(u8s(100..110));
for name in ["single", "farray"] {
assert_eq!(sel_u8(file, name, 0, 12), first, "{name}");
assert_eq!(sel_u8(file, name, N8 - 12, 12), last, "{name}");
}
let mut across = u8s(108..110);
across.extend(u8s(200..210));
assert_eq!(sel_u8(file, "farray", N8 - 2, 12), across);
};
check(&File::open(&path).unwrap());
let file = std::fs::File::open(&path).unwrap();
let len = file.metadata().unwrap().len();
let storage = std::sync::Arc::new(Counting {
file,
len,
read: 0.into(),
});
check(&File::open_storage(storage.clone()).unwrap());
let read = storage.read.load(std::sync::atomic::Ordering::Relaxed);
assert!(read < 1 << 20, "{read} bytes read");
}
/// Run `script` under h5py with `path` as `sys.argv[1]`; skipped without
/// h5py (a failure under `CLAWHDF5_REQUIRE_INTEROP=1`).
fn h5py_check(script: &str, path: &Path) {
if !python_available() {
assert!(!interop_required(), "h5py not available");
eprintln!("SKIP: python3 with h5py not available");
return;
}
let out = Command::new(python())
.arg("-c")
.arg(script)
.arg(path)
.output()
.unwrap();
assert!(
out.status.success(),
"h5py:\n{}",
String::from_utf8_lossy(&out.stderr)
);
}
/// `h5dump` 2.x's output for `args` on `path`, checked to succeed without
/// errors; `None` without `CLAWHDF5_H5DUMP2`.
fn h5dump2_text(args: &[&str], path: &Path) -> Option<String> {
let h5dump = h5dump2()?;
let out = Command::new(&h5dump).args(args).arg(path).output().unwrap();
let text = format!(
"{}{}",
String::from_utf8_lossy(&out.stdout),
String::from_utf8_lossy(&out.stderr)
);
assert!(out.status.success(), "h5dump {args:?}:\n{text}");
assert!(!text.contains("rror"), "h5dump {args:?}:\n{text}");
Some(text)
}
/// clawhdf5 writes deflated `u8` chunks of N8 elements (a chunk dimension
/// of 2^32 or more) in a Single Chunk and an Extensible Array index; it,
/// h5py (libhdf5 2.0.0) and h5dump 2.2.0 (layouts only: it cannot inflate
/// chunks over 4 GiB, see known-issues) read them.
#[test]
fn writer_huge_chunk_dims_round_trip() {
if !heavy() {
return;
}
let dir = scratch();
let path = dir.path().join("dims.h5");
let values = u8s(0..10);
let mut b = clawhdf5::FileBuilder::new();
for (name, max) in [("single", N8), ("earray", u64::MAX)] {
b.create_dataset(name)
.with_u8_data(&values)
.with_maxshape(&[max])
.with_chunks(&[N8])
.with_deflate(1);
}
b.write(&path).unwrap();
let data = std::fs::read(&path).unwrap();
// Deflate level 1 leaves about 40 MiB of each 4 GiB chunk of zeros.
assert!(data.len() < 256 << 20, "{} bytes", data.len());
for (name, index) in [("single", 1), ("earray", 4)] {
let (version, layout, _, chunks) = layout_of(&data, name);
assert_eq!(version, 5, "{name}");
assert!(
matches!(&layout, DataLayout::Chunked { chunk_dimensions, chunk_index_type: Some(t), .. }
if *t == index && chunk_dimensions == &[N8, 1]),
"{name}: {layout:?}"
);
assert_eq!(chunks.len(), 1, "{name}");
}
let file = File::open(&path).unwrap();
for name in ["single", "earray"] {
assert_eq!(sel_u8(&file, name, 0, 10), values, "{name}");
}
drop(file);
h5py_check(
r#"
import sys, h5py
with h5py.File(sys.argv[1], "r") as f:
for name in ("single", "earray"):
d = f[name]
assert d.chunks == (2**32 + 7,), (name, d.chunks)
assert list(d[:]) == list(range(10)), (name, d[:])
"#,
&path,
);
for name in ["single", "earray"] {
if let Some(text) = h5dump2_text(&["-H", "-p", "-d", name], &path) {
assert!(text.contains("CHUNKED ( 4294967303 )"), "{name}:\n{text}");
}
}
}
/// An unfiltered chunk of N8 bytes (a chunk dimension of 2^32 or more)
/// written by clawhdf5 from data it was handed (`with_u8_data_owned`):
/// `FileBuilder::write` streams it to the file without copying it, so the
/// write needs about the data's own 4 GiB. Read back by clawhdf5, h5py
/// (a few elements) and h5dump 2.2.0.
#[test]
fn writer_unfiltered_huge_chunk_round_trip() {
if !heavy() {
return;
}
let dir = scratch();
let path = dir.path().join("raw.h5");
let value = |i: u64| (i % 251) as u8;
let mut b = clawhdf5::FileBuilder::new();
b.create_dataset("single")
.with_u8_data_owned((0..N8).map(value).collect())
.with_chunks(&[N8]);
b.write(&path).unwrap();
let file = File::open(&path).unwrap();
let data = file.as_bytes();
let (version, layout, space, chunks) = layout_of(data, "single");
assert_eq!(version, 5);
assert!(
matches!(&layout, DataLayout::Chunked { chunk_dimensions, chunk_index_type: Some(1), .. }
if chunk_dimensions == &[N8, 1]),
"{layout:?}"
);
assert_eq!(space.dimensions, [N8]);
assert_eq!(chunks.len(), 1);
assert_eq!(chunks[0].chunk_size, N8);
for start in [0, (1 << 32) - 3, N8 - 6] {
let want: Vec<u8> = (start..start + 6).map(value).collect();
assert_eq!(sel_u8(&file, "single", start, 6), want, "at {start}");
}
drop(file);
h5py_check(
r#"
import sys, h5py, numpy as np
n = 2**32 + 7
with h5py.File(sys.argv[1], "r") as f:
d = f["single"]
assert d.chunks == (n,) and d.shape == (n,) and d.compression is None
for start in (0, 2**32 - 3, n - 6):
want = [i % 251 for i in range(start, start + 6)]
assert list(d[start:start + 6]) == want, (start, d[start:start + 6])
"#,
&path,
);
let start = (1u64 << 32) - 3;
if let Some(text) = h5dump2_text(
&["-d", "single", "-s", &start.to_string(), "-c", "6"],
&path,
) {
let want: Vec<String> = (start..start + 6).map(|i| value(i).to_string()).collect();
assert!(text.contains(&want.join(", ")), "{text}");
}
}
/// LZ4 and Zstandard chunks of 4 GiB or more: written by clawhdf5, read
/// back by it and by h5py (through `hdf5plugin`'s LZ4 and Zstandard
/// filters when installed).
#[cfg(any(feature = "lz4", feature = "zstd"))]
#[test]
fn writer_huge_chunks_lz4_zstd_round_trip() {
if !heavy() {
return;
}
let dir = scratch();
let path = dir.path().join("codecs.h5");
let values = u8s(0..10);
let mut b = clawhdf5::FileBuilder::new();
let mut names = Vec::new();
#[cfg(feature = "lz4")]
{
b.create_dataset("lz4")
.with_u8_data(&values)
.with_maxshape(&[N8])
.with_chunks(&[N8])
.with_lz4();
names.push("lz4");
}
#[cfg(feature = "zstd")]
{
b.create_dataset("zstd")
.with_u8_data(&values)
.with_maxshape(&[N8])
.with_chunks(&[N8])
.with_zstd(3);
names.push("zstd");
}
b.write(&path).unwrap();
let data = std::fs::read(&path).unwrap();
assert!(data.len() < 256 << 20, "{} bytes", data.len());
for name in &names {
let (version, _, _, chunks) = layout_of(&data, name);
assert_eq!(version, 5, "{name}");
assert_eq!(chunks.len(), 1, "{name}");
}
drop(data);
let file = File::open(&path).unwrap();
for name in &names {
assert_eq!(sel_u8(&file, name, 0, 10), values, "{name}");
}
drop(file);
h5py_check(
&format!(
r#"
import sys, h5py
try:
import hdf5plugin
except ImportError:
print("SKIP: no hdf5plugin")
sys.exit(0)
with h5py.File(sys.argv[1], "r") as f:
for name in {names:?}:
d = f[name]
assert d.chunks == (2**32 + 7,), (name, d.chunks)
assert list(d[0:10]) == list(range(10)), (name, d[0:10])
"#,
names = names
),
&path,
);
}