From e4b3a76dc7ec5c8876e7701fecb74165fc2c6860 Mon Sep 17 00:00:00 2001 From: osobh Date: Tue, 29 Sep 2026 20:43:14 -0500 Subject: [PATCH] Chunk dimensions of 2^32 or more; unfiltered 4 GiB chunks written without copies - DataLayout::Chunked::chunk_dimensions is Vec (was Vec), and the chunk index readers, writers and serializers take &[u64]: layout messages of version 4/5 store each dimension in up to 8 bytes, and libhdf5 2.x writes dimensions of 2^32 or more (layout version 5). Such dimensions were refused on read (InvalidChunkDimensions) and write. A version-3 layout (4-byte dimensions) is never written for them; a chunk whose size overflows 64 bits is refused when opened. - The file writer lays chunked datasets out as pieces referring to the chunks instead of copying them into one buffer per pass, and an unfiltered chunk that is a contiguous run of the dataset's data (a dataset stored as one chunk of its shape, row blocks) borrows it; a filtered one is compressed straight from it. Contiguous datasets are not copied either. FileWriter::finish_with streams the file to a callback; FileBuilder::write uses it, so the file is never assembled in memory. DatasetBuilder::with_u8_data_owned takes the data without a copy. Peak RSS writing one unfiltered 1 GiB chunk (with_u8_data_owned + write): 5.0 GiB before, 1.0 GiB after; with_u8_data: 6.0 -> 2.0 GiB. - extract_chunk no longer panics on data shorter than the shape. - Tests: huge_chunk_dims.h5 fixture (libhdf5 2.0.0 via h5py 3.16, u8 chunks of 2^32 + 7), and opt-in end-to-end tests of chunk dims >= 2^32 (filtered and unfiltered, read and written, h5py and h5dump 2.2.0), of an unfiltered 4 GiB+ chunk written by clawhdf5, and of LZ4/Zstd chunks of that size; example write_one_chunk for memory measurements. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/clawhdf5-format/src/chunked_read.rs | 72 +-- crates/clawhdf5-format/src/chunked_write.rs | 453 ++++++++++++++---- crates/clawhdf5-format/src/data_layout.rs | 42 +- crates/clawhdf5-format/src/ea_writer.rs | 2 +- .../clawhdf5-format/src/extensible_array.rs | 16 +- crates/clawhdf5-format/src/file_writer.rs | 291 ++++++++--- crates/clawhdf5-format/src/fixed_array.rs | 18 +- crates/clawhdf5-format/src/type_builders.rs | 8 +- crates/clawhdf5-py/src/node.rs | 7 +- crates/clawhdf5-tools/src/check.rs | 8 +- crates/clawhdf5-tools/src/ls.rs | 2 +- crates/clawhdf5/examples/write_one_chunk.rs | 46 ++ crates/clawhdf5/src/edit/mod.rs | 5 +- crates/clawhdf5/src/writer.rs | 34 +- .../tests/fixtures/gen_huge_chunks.py | 59 ++- .../tests/fixtures/huge_chunk_dims.h5 | Bin 0 -> 31199 bytes crates/clawhdf5/tests/huge_chunks_interop.rs | 349 +++++++++++++- 17 files changed, 1156 insertions(+), 256 deletions(-) create mode 100644 crates/clawhdf5/examples/write_one_chunk.rs create mode 100644 crates/clawhdf5/tests/fixtures/huge_chunk_dims.h5 diff --git a/crates/clawhdf5-format/src/chunked_read.rs b/crates/clawhdf5-format/src/chunked_read.rs index 5f0528a..593be30 100644 --- a/crates/clawhdf5-format/src/chunked_read.rs +++ b/crates/clawhdf5-format/src/chunked_read.rs @@ -555,7 +555,7 @@ pub(crate) fn checked_byte_len(elements: u64, elem_size: usize) -> Result 0, dim = {d}" ))); } - let bytes = spatial - .iter() - .fold(elem_size as u128, |acc, &c| acc * u128::from(c)); + // Saturating: 33 dimensions of up to 64 bits each can overflow u128. + let bytes = spatial.iter().fold(elem_size as u128, |acc, &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) { return Err(FormatError::InvalidChunkDimensions(format!( "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::, _>>()?; + Ok((rank, spatial)) } /// 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(()); }; let expected = stored_element_size(datatype, offset_size); - if u64::from(stored) != expected { + if stored != expected { return Err(FormatError::InvalidChunkDimensions(format!( "stored datatype size in chunk layout does not match datatype description \ (layout {stored} bytes, datatype {expected})" @@ -759,7 +770,7 @@ pub fn collect_chunk_info_in( pub fn collect_chunk_info_checked( file_data: &[u8], btree_address: u64, - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], offset_size: u8, length_size: u8, ) -> Result, FormatError> { @@ -776,7 +787,7 @@ pub fn collect_chunk_info_checked( pub fn collect_chunk_info_checked_in( file_data: &S, btree_address: u64, - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], offset_size: u8, length_size: u8, ) -> Result, FormatError> { @@ -808,7 +819,7 @@ pub fn collect_chunk_info_checked_in( .iter_mut() .zip(chunk.offsets.iter().zip(chunk_dimensions)) { - *w = o / u64::from(d); + *w = o / d; } wanted[ndims - 1] = 0; 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()) { - *k /= u64::from(d); + *k /= d; } if let ChunkChildren::Nodes(nodes) = &mut self.children { 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 /// [`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) { - if dim == 0 || offset % u64::from(dim) != 0 { + if dim == 0 || offset % dim != 0 { return Err(FormatError::ChunkedReadError(format!( "bad coordinate offset {offsets:?} for chunk dimensions {chunk_dimensions:?}" ))); @@ -937,7 +948,7 @@ fn read_key_offsets( file_data: &[u8], pos: usize, ndims: usize, - chunk_dimensions: Option<&[u32]>, + chunk_dimensions: Option<&[u64]>, out: &mut Vec, ) -> Result<(), FormatError> { let start = out.len(); @@ -966,7 +977,7 @@ fn parse_chunk_node( file_data: &S, btree_address: u64, ndims: usize, - chunk_dimensions: Option<&[u32]>, + chunk_dimensions: Option<&[u64]>, offset_size: u8, depth: usize, stored: &mut Vec, @@ -1075,7 +1086,7 @@ fn parse_chunk_node( pub fn generate_implicit_chunks( base_address: u64, dataset_dims: &[u64], - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], element_size: u32, ) -> Vec { generate_implicit_chunks_in_grid( @@ -1097,17 +1108,18 @@ pub fn generate_implicit_chunks_in_grid( base_address: u64, dataset_dims: &[u64], max_dims: &[u64], - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], element_size: u32, ) -> Vec { let rank = chunk_dimensions.len(); - let chunk_byte_size: u64 = - chunk_dimensions.iter().map(|&d| d as u64).product::() * element_size as u64; + let chunk_byte_size: u64 = chunk_dimensions + .iter() + .fold(u64::from(element_size), |acc, &d| acc.saturating_mul(d)); let mut num_chunks_per_dim = Vec::with_capacity(rank); let mut grid_per_dim = Vec::with_capacity(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); num_chunks_per_dim.push(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 chunk_idx = 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)); down = down.saturating_mul(grid_per_dim[d]); } @@ -1350,7 +1362,7 @@ pub fn list_chunks_in( } (4, Some(2)) => { // 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( addr, &dataspace.dimensions, @@ -1364,7 +1376,7 @@ pub fn list_chunks_in( } (4, Some(3)) => { // 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( file_data, checked_addr(addr)?, @@ -1384,7 +1396,7 @@ pub fn list_chunks_in( } (4, Some(4)) => { // 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( file_data, checked_addr(addr)?, @@ -2435,12 +2447,12 @@ mod tests { /// 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 /// element size last). - fn build_chunk_btree_leaf_dims(chunks: &[ChunkInfo], dims: &[u32], offset_size: u8) -> Vec { + fn build_chunk_btree_leaf_dims(chunks: &[ChunkInfo], dims: &[u64], offset_size: u8) -> Vec { let last = &chunks.last().expect("a chunk").offsets; let end: Vec = dims .iter() .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(); build_chunk_btree_leaf_to(chunks, &end, offset_size) } @@ -2771,7 +2783,7 @@ mod tests { } // 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_addr = 0x100usize; 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 } - 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_addr = 0x100usize; 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_addr = 0x100usize; file_data[btree_addr..btree_addr + btree.len()].copy_from_slice(&btree); @@ -3408,7 +3420,7 @@ mod tests { } 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), version: 4, chunk_index_type: Some(1), diff --git a/crates/clawhdf5-format/src/chunked_write.rs b/crates/clawhdf5-format/src/chunked_write.rs index d750fd8..0c61aaf 100644 --- a/crates/clawhdf5-format/src/chunked_write.rs +++ b/crates/clawhdf5-format/src/chunked_write.rs @@ -5,12 +5,14 @@ extern crate alloc; use crate::addr::saturating_usize; #[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_v2_write::{BTreeV2Params, build_btree_v2}; use crate::checksum::jenkins_lookup3; -use crate::chunk_cache::{CACHE_LINE_SIZE, align_to_cache_line}; +use crate::chunk_cache::CACHE_LINE_SIZE; use crate::chunk_grid::ChunkGrid; use crate::ea_writer; use crate::error::FormatError; @@ -464,7 +466,9 @@ fn extract_chunk( let dst: usize = (0..rank).map(|d| idx[d] * chunk_strides[d]).sum::() * 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]); + if n > 0 { + 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 { @@ -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> { + 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. /// /// 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 /// rayon threads; otherwise it is sequential. Output order matches chunk /// order, so per-chunk bytes are identical to the sequential path. -fn compress_all_chunks( - raw_data: &[u8], +fn compress_all_chunks<'a>( + raw_data: &'a [u8], shape: &[u64], chunk_dims: &[u64], element_size: usize, chunk_bytes: u64, pipeline: &Option, -) -> Result, u32)>, FormatError> { - let one = |i: u64| -> Result<(u64, Vec, 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) - }; +) -> Result>, FormatError> { + let one = |i: u64| -> Result, FormatError> { + let raw = chunk_bytes_of(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)) + Ok((raw_size, Cow::Owned(stored), mask)) } 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. /// Serialize a v4 single chunk layout message (public for OH size estimation). pub fn serialize_v4_single_chunk_pub( - chunk_dims: &[u32], + chunk_dims: &[u64], chunk_address: u64, filtered_size: Option, filter_mask: Option, @@ -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 /// layout message. fn serialize_v4_single_chunk( - chunk_dims: &[u32], + chunk_dims: &[u64], chunk_address: u64, filtered_size: Option, filter_mask: Option, @@ -622,7 +682,7 @@ fn serialize_v4_single_chunk( /// Serialize a v4 Fixed Array layout message. fn serialize_v4_fixed_array( - chunk_dims: &[u32], + chunk_dims: &[u64], fixed_array_address: u64, offset_size: u8, element_size: u32, @@ -653,7 +713,8 @@ fn serialize_v4_fixed_array( /// dimensions, then the element size). Each takes the fewest bytes that hold /// the largest, as libhdf5 computes it (`H5D__chunk_set_sizes`: /// `(log2(dim) + 8) / 8`); HDF5 2.0.0 refuses any other width. -pub(crate) fn push_v4_chunk_dims(buf: &mut Vec, chunk_dims: &[u32], element_size: u32) { +pub(crate) fn push_v4_chunk_dims(buf: &mut Vec, chunk_dims: &[u64], element_size: u32) { + let element_size = u64::from(element_size); let max_dim = chunk_dims .iter() .copied() @@ -661,14 +722,14 @@ pub(crate) fn push_v4_chunk_dims(buf: &mut Vec, chunk_dims: &[u32], element_ .max() .unwrap_or(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); for &d in chunk_dims.iter().chain(core::iter::once(&element_size)) { buf.extend_from_slice(&d.to_le_bytes()[..width]); } } -fn layout_v4_chunked_prefix(chunk_dims: &[u32], element_size: u32, version: u8) -> Vec { +fn layout_v4_chunked_prefix(chunk_dims: &[u64], element_size: u32, version: u8) -> Vec { let mut buf = Vec::new(); buf.push(version); buf.push(2); // class = chunked @@ -884,7 +945,66 @@ pub fn precompress_chunks( element_size: usize, options: &ChunkOptions, ) -> Result { - 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>, + pub(crate) has_filters: bool, + pub(crate) element_size: usize, + pub(crate) shape: Vec, + pub(crate) chunk_dims: Vec, + pub(crate) pipeline_message: Option>, +} + +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, FormatError> { + let chunk_bytes = checked_chunk_dims(chunk_dims, element_size)?; if chunk_bytes > MAX_V4_CHUNK_BYTES { check_huge_chunk_filters(options, chunk_bytes)?; } @@ -902,7 +1022,7 @@ pub fn precompress_chunks( &pipeline, )?; - Ok(PrecompressedChunks { + Ok(PreparedChunks { chunks, has_filters, element_size, @@ -912,26 +1032,87 @@ 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, 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::, _>>()?; +/// A piece of a chunked dataset's data blob as laid out by +/// [`place_chunked_data`]: zero padding, a chunk's stored bytes (by index +/// into [`PreparedChunks::chunks`]), or index structures. +#[derive(Debug)] +pub(crate) enum Piece { + Zeros(usize), + Chunk(usize), + Bytes(Vec), +} + +/// A chunked dataset laid out at an address, without its chunk bytes +/// copied: the pieces of its data blob, in order, and its messages. +pub(crate) struct ChunkedPlacement { + pub(crate) pieces: Vec, + /// Length of the data blob in bytes. + pub(crate) len: u64, + pub(crate) layout_message: Vec, + pub(crate) pipeline_message: Option>, +} + +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, + 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) { + 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 { + if chunk_dims.contains(&0) { + return Err(FormatError::InvalidChunkDimensions(format!( + "chunk dimensions {chunk_dims:?} include 0" + ))); + } let bytes = chunk_dims .iter() .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" )) })?; - Ok((dims, bytes)) + Ok(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, high: LibVer, ) -> Result { - 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 { + let chunk_bytes = checked_chunk_dims(&pre.chunk_dims, pre.element_size)?; if chunk_bytes > MAX_V4_CHUNK_BYTES { if high < LibVer::V200 { return Err(FormatError::LibverBound { @@ -1020,7 +1216,7 @@ pub fn build_chunked_data_from_precompressed_libver( }); } } 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 offset_size: u8 = 8; @@ -1028,17 +1224,14 @@ pub fn build_chunked_data_from_precompressed_libver( let num_chunks = pre.chunks.len(); 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); - for (raw_size, compressed, filter_mask) in &pre.chunks { - let aligned_offset = align_to_cache_line(data_buf.len()); - if aligned_offset > data_buf.len() { - data_buf.resize(aligned_offset, 0u8); - } - let address = base_address + data_buf.len() as u64; + for (i, (raw_size, compressed, filter_mask)) in pre.chunks.iter().enumerate() { + blob.align(); + let address = base_address + blob.len; let compressed_size = compressed.len() as u64; - data_buf.extend_from_slice(compressed); + blob.chunk(i, compressed.len()); written_chunks.push(WrittenChunk { address, 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 aligned_idx = align_to_cache_line(data_buf.len()); - if aligned_idx > data_buf.len() { - data_buf.resize(aligned_idx, 0u8); - } + blob.align(); let layout_message = match &index { 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 ea_bytes = ea_writer::build_extensible_array_at( + blob.bytes(ea_writer::build_extensible_array_at( &slots, offset_size, length_size, pre.has_filters, ea_address, - ); - data_buf.extend_from_slice(&ea_bytes); + )); ea_writer::serialize_v4_extensible_array( - &chunk_dims_u32, + &pre.chunk_dims, ea_address, offset_size, 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); serialize_v4_single_chunk( - &chunk_dims_u32, + &pre.chunk_dims, chunk_addr, filtered_size, filter_mask, @@ -1094,7 +1281,7 @@ pub fn build_chunked_data_from_precompressed_libver( ) } 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( grid, &pre.shape, @@ -1102,16 +1289,15 @@ pub fn build_chunked_data_from_precompressed_libver( &written_chunks, Some(*nslots), )?; - let fa_bytes = build_fixed_array_at( + blob.bytes(build_fixed_array_at( &slots, offset_size, length_size, pre.has_filters, fa_address, - ); - data_buf.extend_from_slice(&fa_bytes); + )); serialize_v4_fixed_array( - &chunk_dims_u32, + &pre.chunk_dims, fa_address, offset_size, element_size as u32, @@ -1120,7 +1306,7 @@ pub fn build_chunked_data_from_precompressed_libver( ) } ChunkIndexPlan::BTreeV2 => { - let bt_address = base_address + data_buf.len() as u64; + let bt_address = base_address + blob.len; let records: Vec<(Vec, &WrittenChunk)> = written_chunks .iter() .enumerate() @@ -1134,9 +1320,9 @@ pub fn build_chunked_data_from_precompressed_libver( pre.has_filters, bt_address, )?; - data_buf.extend_from_slice(&bt_bytes); + blob.bytes(bt_bytes); serialize_v4_btree_v2( - &chunk_dims_u32, + &pre.chunk_dims, bt_address, offset_size, element_size as u32, @@ -1146,8 +1332,9 @@ pub fn build_chunked_data_from_precompressed_libver( } }; - Ok(ChunkedDataResult { - data_bytes: data_buf, + Ok(ChunkedPlacement { + pieces: blob.pieces, + len: blob.len, layout_message, 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 /// B-tree chunk index, with a version-3 layout message: what libhdf5 writes /// for a chunked dataset under a low bound of 1.8. -fn build_btree_v1_chunked_data( - pre: &PrecompressedChunks, +fn place_btree_v1_chunked_data( + pre: &PreparedChunks<'_>, base_address: u64, maxshape: Option<&[u64]>, -) -> Result { +) -> Result { if let Some(ms) = maxshape { let bad = |what: &str| FormatError::ChunkedReadError(format!("maxshape: {what}")); if ms.len() != pre.shape.len() { @@ -1171,20 +1358,17 @@ fn build_btree_v1_chunked_data( } } 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()); for (i, (_raw_size, stored, filter_mask)) in pre.chunks.iter().enumerate() { - let aligned_offset = align_to_cache_line(data_buf.len()); - if aligned_offset > data_buf.len() { - data_buf.resize(aligned_offset, 0u8); - } + blob.align(); entries.push(btree_v1_write::ChunkEntry { scaled: scaled_coords(&pre.shape, &pre.chunk_dims, i), nbytes: stored.len() as u64, filter_mask: *filter_mask, - address: base_address + data_buf.len() as u64, + 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) .map_err(|_| FormatError::Overflow("element size".into()))?; @@ -1193,11 +1377,8 @@ fn build_btree_v1_chunked_data( let btree_address = if entries.is_empty() { u64::MAX } else { - let aligned_idx = align_to_cache_line(data_buf.len()); - if aligned_idx > data_buf.len() { - data_buf.resize(aligned_idx, 0u8); - } - let addr = base_address + data_buf.len() as u64; + blob.align(); + let addr = base_address + blob.len; let tree = btree_v1_write::build_chunk_btree_v1_at( &entries, &pre.chunk_dims, @@ -1205,13 +1386,14 @@ fn build_btree_v1_chunked_data( addr, offset_size, )?; - data_buf.extend_from_slice(&tree); + blob.bytes(tree); addr }; let layout_message = serialize_v3_chunked(&pre.chunk_dims, btree_address, offset_size, element_size)?; - Ok(ChunkedDataResult { - data_bytes: data_buf, + Ok(ChunkedPlacement { + pieces: blob.pieces, + len: blob.len, layout_message, 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. fn serialize_v4_btree_v2( - chunk_dims: &[u32], + chunk_dims: &[u64], btree_address: u64, offset_size: u8, element_size: u32, @@ -2184,7 +2366,7 @@ mod tests { /// 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 dims = [HUGE_DIM]; let parsed = |msg: &[u8]| { assert_eq!(msg[0], 5, "layout message version"); // Class chunked, then (after the flags) 2 dimensions of 4 bytes. @@ -2197,7 +2379,7 @@ mod tests { 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() } other => panic!("{other:?}"), @@ -2245,22 +2427,89 @@ mod tests { 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 = (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] fn huge_chunks_refused_where_unsupported() { - // A chunk dimension of 2^32 or more. + // A chunk dimension of 0. assert!(matches!( - checked_chunk_dims(&[1 << 32], 1), - Err(FormatError::InvalidChunkDimensions(m)) if m.contains("2^32") + checked_chunk_dims(&[4, 0], 1), + 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. 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) - ); + assert_eq!(checked_chunk_dims(&[HUGE_DIM], 8).unwrap(), HUGE_BYTES); // Filters that cannot take a chunk that large. for plugin in [ PluginFilter::Lzf, diff --git a/crates/clawhdf5-format/src/data_layout.rs b/crates/clawhdf5-format/src/data_layout.rs index cf439d6..8310cd4 100644 --- a/crates/clawhdf5-format/src/data_layout.rs +++ b/crates/clawhdf5-format/src/data_layout.rs @@ -34,7 +34,7 @@ const MAX_LAYOUT_NDIMS: usize = 33; /// (`H5O__layout_decode`): at most [`MAX_LAYOUT_NDIMS`], no dimension 0, and /// 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. -fn check_chunk_dims(dims: Vec, layout_version: u8) -> Result, FormatError> { +fn check_chunk_dims(dims: Vec, layout_version: u8) -> Result, FormatError> { if dims.len() > MAX_LAYOUT_NDIMS { return Err(FormatError::InvalidChunkDimensions( "dimensionality is too large".into(), @@ -72,7 +72,7 @@ pub enum DataLayout { /// Chunked: data stored in chunks via a B-tree. Chunked { /// Chunk dimension sizes. - chunk_dimensions: Vec, + chunk_dimensions: Vec, /// B-tree address, or `None` if undefined. btree_address: Option, /// 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)), }; ensure_len(data, p, dimensionality * 4)?; - let dims: Vec = data[p..p + dimensionality * 4] + let dims: Vec = data[p..p + dimensionality * 4] .as_chunks::<4>() .0 .iter() - .map(|c| u32::from_le_bytes(*c)) + .map(|c| u64::from(u32::from_le_bytes(*c))) .collect(); p += dimensionality * 4; match layout_class { @@ -424,7 +424,7 @@ impl DataLayout { 1 => { let size = dims .iter() - .try_fold(1u64, |acc, &d| acc.checked_mul(d as u64)) + .try_fold(1u64, |acc, &d| acc.checked_mul(d)) .ok_or_else(|| { FormatError::Overflow(format!("contiguous layout size {dims:?}")) })?; @@ -489,7 +489,12 @@ impl DataLayout { ensure_len(data, p, dimensionality * 4)?; let mut chunk_dimensions = Vec::with_capacity(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); p += 4; } @@ -565,14 +570,8 @@ impl DataLayout { .iter() .rev() .fold(0u64, |acc, &b| (acc << 8) | u64::from(b)); - // Chunk dimensions are held as u32; HDF5 2.0 can write - // larger ones (layout version 5), which are refused - // 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" - )) - })?; + // HDF5 2.0 writes dimensions of 2^32 or more (layout + // version 5, chunks over 4 GiB) in 5 to 8 bytes. chunk_dimensions.push(val); p += dim_size_encoded_length; } @@ -858,7 +857,7 @@ mod tests { .unwrap_or_else(|e| panic!("width {width}: {e:?}")); assert!( matches!(&layout, DataLayout::Chunked { chunk_dimensions, .. } - if chunk_dimensions.iter().map(|&d| u64::from(d)).eq(dims)), + if chunk_dimensions[..] == dims), "width {width}: {layout:?}" ); } @@ -871,12 +870,13 @@ mod tests { ) ); } - // A dimension past u32 cannot be represented and is refused, not - // truncated. - assert!(matches!( - DataLayout::parse(&v4_chunked_msg(5, &[1 << 32, 8]), 8, 8), - Err(FormatError::InvalidChunkDimensions(m)) if m.contains("2^32") - )); + // HDF5 2.0 writes dimensions of 2^32 or more (up to 8 bytes each). + for (width, dim) in [(5u8, (1u64 << 32) + 3), (8, u64::MAX)] { + assert!(matches!( + DataLayout::parse(&v4_chunked_msg(width, &[dim, 1]), 8, 8), + Ok(DataLayout::Chunked { chunk_dimensions, .. }) if chunk_dimensions == [dim, 1] + )); + } } #[test] diff --git a/crates/clawhdf5-format/src/ea_writer.rs b/crates/clawhdf5-format/src/ea_writer.rs index 5951c50..04a996c 100644 --- a/crates/clawhdf5-format/src/ea_writer.rs +++ b/crates/clawhdf5-format/src/ea_writer.rs @@ -14,7 +14,7 @@ use crate::chunked_write::{ /// Serialize a v4 Extensible Array layout message. pub(crate) fn serialize_v4_extensible_array( - chunk_dims: &[u32], + chunk_dims: &[u64], ea_address: u64, offset_size: u8, element_size: u32, diff --git a/crates/clawhdf5-format/src/extensible_array.rs b/crates/clawhdf5-format/src/extensible_array.rs index 889f2cf..45f8d82 100644 --- a/crates/clawhdf5-format/src/extensible_array.rs +++ b/crates/clawhdf5-format/src/extensible_array.rs @@ -454,7 +454,7 @@ pub fn read_extensible_array_chunks( header: &ExtensibleArrayHeader, dataset_dims: &[u64], max_dims: Option<&[u64]>, - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], element_size: u32, offset_size: u8, length_size: u8, @@ -480,7 +480,7 @@ pub fn read_extensible_array_chunks_in( header: &ExtensibleArrayHeader, dataset_dims: &[u64], max_dims: Option<&[u64]>, - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], element_size: u32, offset_size: u8, _length_size: u8, @@ -489,12 +489,12 @@ pub fn read_extensible_array_chunks_in( // Linear indexes follow the maximum dimensions, with the unlimited // dimension swizzled to the slowest position (see `chunk_grid`). - let dims_u64: Vec = chunk_dimensions.iter().map(|&d| d as u64).collect(); - let grid = ChunkGrid::extensible_array(dataset_dims, max_dims, &dims_u64)?; + let grid = ChunkGrid::extensible_array(dataset_dims, max_dims, chunk_dimensions)?; let grid = &grid; - let chunk_byte_size: u64 = - chunk_dimensions.iter().map(|&d| d as u64).product::() * element_size as u64; + let chunk_byte_size: u64 = chunk_dimensions + .iter() + .fold(u64::from(element_size), |acc, &d| acc.saturating_mul(d)); // Parse index block (EAIB): signature(4) + version(1) + client_id(1) // + 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 ds_dims = vec![40u64]; // 2 chunks × 20 elements - let chunk_dims = vec![20u32]; + let chunk_dims = vec![20u64]; let chunks = read_extensible_array_chunks( &file_data, &header, @@ -1057,7 +1057,7 @@ mod tests { let file_data = build_inline_plus_data_blocks(); let header = ExtensibleArrayHeader::parse(&file_data, 0x100, os, ls).unwrap(); let ds_dims = vec![40u64]; - let chunk_dims = vec![10u32]; + let chunk_dims = vec![10u64]; let chunks = read_extensible_array_chunks( &file_data, &header, diff --git a/crates/clawhdf5-format/src/file_writer.rs b/crates/clawhdf5-format/src/file_writer.rs index 4eb59ee..a1025ad 100644 --- a/crates/clawhdf5-format/src/file_writer.rs +++ b/crates/clawhdf5-format/src/file_writer.rs @@ -3,15 +3,14 @@ //! Produces valid HDF5 files with v3 superblock, v2 object headers, //! link messages, contiguous datasets, inline and dense attributes. -use crate::addr::saturating_usize; +use crate::addr::{saturating_usize, to_usize}; #[cfg(not(feature = "std"))] use alloc::{format, string::String, vec, vec::Vec}; use crate::attribute::AttributeMessage; use crate::btree_v2_write::{BTreeV2Params, build_btree_v2}; use crate::chunked_write::{ - ChunkOptions, PrecompressedChunks, build_chunked_data_from_precompressed_libver, - precompress_chunks, + ChunkOptions, Piece, PreparedChunks, place_chunked_data, prepare_chunks, }; use crate::data_layout::VdsMapping; use crate::dataspace::{Dataspace, DataspaceType}; @@ -1612,7 +1611,30 @@ impl FileWriter { 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, 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>( + self, + put: impl FnMut(&[u8]) -> Result<(), E>, + ) -> Result<(), E> { + self.finish_into(&mut FnSink(put)) + } + + fn finish_into(self, sink: &mut S) -> Result<(), S::Error> { let page_size = self.page_size; if let Some(ps) = page_size && !(MIN_FILE_SPACE_PAGE_SIZE..=MAX_FILE_SPACE_PAGE_SIZE).contains(&ps) @@ -1620,7 +1642,8 @@ impl FileWriter { return Err(FormatError::SerializationError(format!( "file space page size {ps} is outside libhdf5's \ {MIN_FILE_SPACE_PAGE_SIZE}..={MAX_FILE_SPACE_PAGE_SIZE} bytes" - ))); + )) + .into()); } let (low, high) = (self.low, self.high); @@ -1774,15 +1797,16 @@ impl FileWriter { }) .collect::>()?; - struct DataBlob { + struct DataBlob<'a> { data: Vec, oh_bytes: Vec, /// Cached compressed chunks for chunked datasets; reused in Pass 2 - /// to avoid re-compressing the same data. - precompressed: Option, + /// to avoid re-compressing the same data. Unfiltered chunks + /// borrow the dataset's data where they can. + precompressed: Option>, } - let mut dummy_blobs: Vec = Vec::new(); + let mut dummy_blobs: Vec> = Vec::new(); let mut dummy_cursor = 0u64; for (i, d) in all_ds.iter().enumerate() { let dense_blob = ds_dense[i] @@ -1820,21 +1844,16 @@ impl FileWriter { .resolve_chunk_dims_for(&d.ds.dimensions, elem_size); // Compress once in Pass 1; cache the result so Pass 2 can skip // re-compression and just rebuild the index with real addresses. - let pre = precompress_chunks( + let pre = prepare_chunks( &d.raw, &d.ds.dimensions, &chunk_dims, elem_size, &d.chunk_options, )?; - let result = build_chunked_data_from_precompressed_libver( - &pre, - dummy_cursor, - d.maxshape.as_deref(), - low, - high, - )?; - dummy_cursor += result.data_bytes.len() as u64; + let result = + place_chunked_data(&pre, dummy_cursor, d.maxshape.as_deref(), low, high)?; + dummy_cursor += result.len; let oh = build_chunked_dataset_oh( &d.dt, &d.ds, @@ -1849,7 +1868,7 @@ impl FileWriter { d.refcount, )?; dummy_blobs.push(DataBlob { - data: result.data_bytes, + data: Vec::new(), oh_bytes: oh, precompressed: Some(pre), }); @@ -1955,7 +1974,7 @@ impl FileWriter { }) .collect::>()?; - let mut ds_blobs2: Vec = Vec::new(); + let mut ds_blobs2: Vec> = Vec::new(); let global_align_threshold = self.alignment_threshold; let global_align_bytes = self.alignment_bytes; for (i, d) in all_ds.iter().enumerate() { @@ -1977,26 +1996,21 @@ impl FileWriter { &d.fill_message, d.refcount, )?; - ds_blobs2.push(DataBlob { - data: gcol_bytes.clone(), + ds_blobs2.push(OutBlob { + data: vec![Out::Slice(gcol_bytes)], oh_bytes: oh, - precompressed: None, }); } else if is_chunked[i] { let base_address = cursor2 as u64; // Reuse precompressed chunks from Pass 1 — avoids re-compressing // the same data a second time. - let result = build_chunked_data_from_precompressed_libver( - dummy_blobs[i] - .precompressed - .as_ref() - .expect("chunked dataset missing precompressed cache"), - base_address, - d.maxshape.as_deref(), - low, - high, - )?; - cursor2 += result.data_bytes.len(); + let pre = dummy_blobs[i] + .precompressed + .as_ref() + .expect("chunked dataset missing precompressed cache"); + let result = + place_chunked_data(pre, base_address, d.maxshape.as_deref(), low, high)?; + cursor2 += to_usize(result.len)?; let oh = build_chunked_dataset_oh( &d.dt, &d.ds, @@ -2010,11 +2024,16 @@ impl FileWriter { &d.fill_message, d.refcount, )?; - ds_blobs2.push(DataBlob { - data: result.data_bytes, - oh_bytes: oh, - precompressed: None, - }); + let data = result + .pieces + .into_iter() + .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] { // Compact: data is inline in the object header, no external blob let oh = build_compact_dataset_oh( @@ -2030,10 +2049,9 @@ impl FileWriter { d.refcount, layout_version, )?; - ds_blobs2.push(DataBlob { + ds_blobs2.push(OutBlob { data: vec![], oh_bytes: oh, - precompressed: None, }); } else { // Determine alignment: per-dataset overrides global @@ -2060,13 +2078,10 @@ impl FileWriter { d.refcount, layout_version, )?; - let mut data = vec![0u8; padding]; - data.extend_from_slice(&d.raw); cursor2 += d.raw.len(); - ds_blobs2.push(DataBlob { - data, + ds_blobs2.push(OutBlob { + data: vec![Out::Zeros(padding), Out::Slice(&d.raw)], oh_bytes: oh, - precompressed: None, }); } } @@ -2080,7 +2095,8 @@ impl FileWriter { cursor2 = cursor2.next_multiple_of(ps as usize); } 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 { version: superblock_version, @@ -2103,9 +2119,9 @@ impl FileWriter { checksum: None, page_size: None, }; - buf.extend_from_slice(&sb.serialize()); + out.put(&sb.serialize())?; 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) @@ -2132,32 +2148,120 @@ impl FileWriter { g.refcount, )?; debug_assert_eq!(oh.len(), group_oh_sizes[gi]); - debug_assert_eq!(buf.len() as u64, group_addrs2[gi]); - buf.extend_from_slice(&oh); + debug_assert_eq!(out.written as u64, group_addrs2[gi]); + out.put(&oh)?; 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] { - buf.extend_from_slice(&blob.blob); + out.put(&blob.blob)?; } } // Dataset OHs + dense blobs 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] { - buf.extend_from_slice(&dense.blob); + out.put(&dense.blob)?; } } // Data 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); - buf.resize(cursor2, 0); - Ok(buf) + debug_assert_eq!(out.written, data_end); + out.zeros(cursor2 - data_end)?; + Ok(()) + } +} + +/// A piece of a dataset's data as [`FileWriter`] writes it out. +enum Out<'a> { + Bytes(Vec), + /// 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, + data: Vec>, +} + +/// Where [`FileWriter`] puts a file's bytes, in order. +trait Sink { + type Error: From; + + /// 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 { + 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); + +impl, F: FnMut(&[u8]) -> Result<(), E>> Sink for FnSink { + 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 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!(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 = (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); + } } diff --git a/crates/clawhdf5-format/src/fixed_array.rs b/crates/clawhdf5-format/src/fixed_array.rs index dbe00ae..b28f60c 100644 --- a/crates/clawhdf5-format/src/fixed_array.rs +++ b/crates/clawhdf5-format/src/fixed_array.rs @@ -155,7 +155,7 @@ pub fn read_fixed_array_chunks( header: &FixedArrayHeader, dataset_dims: &[u64], max_dims: Option<&[u64]>, - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], element_size: u32, offset_size: u8, length_size: u8, @@ -180,7 +180,7 @@ pub fn read_fixed_array_chunks_in( header: &FixedArrayHeader, dataset_dims: &[u64], max_dims: Option<&[u64]>, - chunk_dimensions: &[u32], + chunk_dimensions: &[u64], element_size: u32, offset_size: u8, _length_size: u8, @@ -227,11 +227,11 @@ pub fn read_fixed_array_chunks_in( // The index is laid out over the chunk grid of the *maximum* dimensions // (row-major), so a dataset smaller than its maxshape has gaps. - let dims_u64: Vec = chunk_dimensions.iter().map(|&d| d as u64).collect(); - let grid = ChunkGrid::fixed_array(dataset_dims, max_dims, &dims_u64)?; + let grid = ChunkGrid::fixed_array(dataset_dims, max_dims, chunk_dimensions)?; - let chunk_byte_size: u64 = - chunk_dimensions.iter().map(|&d| d as u64).product::() * element_size as u64; + let chunk_byte_size: u64 = chunk_dimensions + .iter() + .fold(u64::from(element_size), |acc, &d| acc.saturating_mul(d)); let mut chunks = Vec::new(); // `rel` is relative to the data block, whose bytes are in `w`. @@ -670,7 +670,7 @@ mod tests { let header = FixedArrayHeader::parse(&file_data, fahd_offset, offset_size, length_size).unwrap(); let ds_dims = vec![100u64]; - let chunk_dims = vec![20u32]; + let chunk_dims = vec![20u64]; let chunks = read_fixed_array_chunks( &file_data, &header, @@ -747,7 +747,7 @@ mod tests { let header = FixedArrayHeader::parse(&file_data, fahd_offset, offset_size, length_size).unwrap(); let ds_dims = vec![60u64]; - let chunk_dims = vec![20u32]; + let chunk_dims = vec![20u64]; let chunks = read_fixed_array_chunks( &file_data, &header, @@ -848,7 +848,7 @@ mod tests { assert_eq!(header.num_elements, 11); let ds_dims = vec![11u64 * 20]; - let chunk_dims = vec![20u32]; + let chunk_dims = vec![20u64]; let chunks = read_fixed_array_chunks( &file_data, &header, diff --git a/crates/clawhdf5-format/src/type_builders.rs b/crates/clawhdf5-format/src/type_builders.rs index ba8c292..f9776b9 100644 --- a/crates/clawhdf5-format/src/type_builders.rs +++ b/crates/clawhdf5-format/src/type_builders.rs @@ -663,11 +663,17 @@ impl DatasetBuilder { } 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) -> &mut Self { self.datatype = Some(make_u8_type()); - self.data = Some(data.to_vec()); if self.shape.is_none() { self.shape = Some(vec![data.len() as u64]); } + self.data = Some(data); self } diff --git a/crates/clawhdf5-py/src/node.rs b/crates/clawhdf5-py/src/node.rs index 35e5dec..9ef9290 100644 --- a/crates/clawhdf5-py/src/node.rs +++ b/crates/clawhdf5-py/src/node.rs @@ -189,12 +189,7 @@ pub(crate) fn chunk_shape(file: &File, hdr: &ObjectHeader, rank: usize) -> Optio { clawhdf5_format::data_layout::DataLayout::Chunked { chunk_dimensions, .. - } if chunk_dimensions.len() >= rank => Some( - chunk_dimensions[..rank] - .iter() - .map(|&d| u64::from(d)) - .collect(), - ), + } if chunk_dimensions.len() >= rank => Some(chunk_dimensions[..rank].to_vec()), _ => None, } } diff --git a/crates/clawhdf5-tools/src/check.rs b/crates/clawhdf5-tools/src/check.rs index 1ba8239..aa7729d 100644 --- a/crates/clawhdf5-tools/src/check.rs +++ b/crates/clawhdf5-tools/src/check.rs @@ -873,9 +873,9 @@ impl Checker<'_> { self.counts.chunk_index_checksummed += 1; } 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| { - a.checked_mul(u64::from(d)) - }); + let chunk_bytes = cdims + .iter() + .try_fold(u64::from(dt.type_size()), |a, &d| a.checked_mul(d)); let max = ds.max_dimensions.clone(); let mut seen: HashSet> = HashSet::with_capacity(chunks.len().min(1 << 20)); let mut reported = 0usize; @@ -893,7 +893,7 @@ impl Checker<'_> { )); } else { for (i, (&o, &cd)) in c.offsets.iter().zip(cdims).enumerate() { - if o % u64::from(cd) != 0 { + if o % cd != 0 { bad.push(format!( "offset {o} in dimension {i} is not a multiple of the chunk size {cd}" )); diff --git a/crates/clawhdf5-tools/src/ls.rs b/crates/clawhdf5-tools/src/ls.rs index 203cce3..8a466a2 100644 --- a/crates/clawhdf5-tools/src/ls.rs +++ b/crates/clawhdf5-tools/src/ls.rs @@ -292,7 +292,7 @@ impl Ls<'_> { .unwrap_or(0); let bytes = dims .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()) .unwrap_or_else(|| "?".into()); writeln!( diff --git a/crates/clawhdf5/examples/write_one_chunk.rs b/crates/clawhdf5/examples/write_one_chunk.rs new file mode 100644 index 0000000..2761e8b --- /dev/null +++ b/crates/clawhdf5/examples/write_one_chunk.rs @@ -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 = 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 = (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(); +} diff --git a/crates/clawhdf5/src/edit/mod.rs b/crates/clawhdf5/src/edit/mod.rs index 169c6d3..a550edc 100644 --- a/crates/clawhdf5/src/edit/mod.rs +++ b/crates/clawhdf5/src/edit/mod.rs @@ -408,10 +408,7 @@ impl<'t> ChunkedEdit<'t> { "chunk rank differs from the dataspace".into(), )); } - let cd: Vec = chunk_dimensions[..rank] - .iter() - .map(|&c| u64::from(c)) - .collect(); + let cd: Vec = chunk_dimensions[..rank].to_vec(); // 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 diff --git a/crates/clawhdf5/src/writer.rs b/crates/clawhdf5/src/writer.rs index 9fcf6fc..eb05b97 100644 --- a/crates/clawhdf5/src/writer.rs +++ b/crates/clawhdf5/src/writer.rs @@ -137,10 +137,15 @@ impl FileBuilder { 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>(self, path: P) -> Result<(), Error> { - let bytes = self.finish()?; - write_file_atomically(path.as_ref(), &bytes).map_err(Error::Io) + use std::io::Write; + 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) -> Result, Erro /// 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 /// existing file. +#[cfg(test)] fn write_file_atomically(path: &std::path::Path, bytes: &[u8]) -> std::io::Result<()> { 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>( + path: &std::path::Path, + write: impl FnOnce(&mut std::io::BufWriter) -> Result<(), E>, +) -> Result<(), E> { + use std::io::Write; // Same directory as the target, so the rename stays on one filesystem. 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())); let tmp_path = path.with_file_name(tmp_name); - let result = (|| { - let mut f = std::fs::File::create(&tmp_path)?; - f.write_all(bytes)?; - f.sync_all()?; - std::fs::rename(&tmp_path, path) + let result = (|| -> Result<(), E> { + let mut f = std::io::BufWriter::with_capacity(1 << 20, std::fs::File::create(&tmp_path)?); + write(&mut f)?; + f.flush()?; + f.get_ref().sync_all()?; + std::fs::rename(&tmp_path, path)?; + Ok(()) })(); if result.is_err() { let _ = std::fs::remove_file(&tmp_path); diff --git a/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py b/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py index f901fd5..6cdf311 100644 --- a/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py +++ b/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py @@ -2,6 +2,12 @@ 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 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 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 `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 @@ -92,12 +99,60 @@ def make(fid, name, index, filtered): 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(): mode, out = sys.argv[1], sys.argv[2] - filtered = {"filtered": True, "unfiltered": False}[mode] fapl = h5py.h5p.create(h5py.h5p.FILE_ACCESS) fapl.set_libver_bounds(h5py.h5f.LIBVER_EARLIEST, h5py.h5f.LIBVER_V200) 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"] if not filtered: indexes.insert(1, "implicit") diff --git a/crates/clawhdf5/tests/fixtures/huge_chunk_dims.h5 b/crates/clawhdf5/tests/fixtures/huge_chunk_dims.h5 new file mode 100644 index 0000000000000000000000000000000000000000..34d97d5a00c7d994d66ae61a244d564f15c09950 GIT binary patch literal 31199 zcmeHPZA?>F7(TbOSlJX9j8#M})66Av8~bo0fG`$tkP&gZh(cUJsq!Ti3bX=+Ml)c6 zg=L1!=#VYT+@fZH7{G}@B@Ep{wgep_8gz<8ffZ%C)k3>_?uWZ7Vd2-7@LZC6&Uw#A z+vh&-^PbZC-jAXpR?lc-6pn#iyhUt_{fR<_5y~>5{8q+S7()LD zQ3-ZWVf{zPy}vFplCUGi^~F+{Mgrt~O(_~s&ME|=BE#3x54=(M)FPkG2s{(P(5Gf@ z&(IR0P0G$rVmT}F@-1<|3`>S6`89-v7=g(Q!yD$X>(QE6y0HAFZR$n*2%fPD!7g)= zaWm<7ddP+K#WkKFf!sXif&b9gXGpSkTSiijmdxd+v6iXrO7c86LLrmM6x^HQ-UU$?#EG5~8g$fmTU*lp@l?rpNp=02*a(Km1TQAUsE~9{Q59cG?J`B@EH?ccuRK|7 z=rs{_@W1Dkn!8#G8Ikfzu9nL^+{#?JK6-kbPwEGK(^VqFqaq07t(JSp`B*f~VA17; zkU)k*_`yLtHKEf?L)R1t7%z@q@=vc~R3%rfs*uvEN%PWc9Sf54%C)KI>(bubSaEuH z&rA2cDxP(j=a!!7T5btl8y@rOhE7GLvFHbbCN6(`Va?94py06T?vkumD(p2cFIl3q zE?#6hSzPbOzMeTcc;Z;@j`~(dVe{l4D_KiG?!r0lU3Y!P^tp{=^I5~bmg@1I<-;#t z>AYpR_kxp7I$IE2qp`IqMjWQ`k*^J`sr|ai`$^j6ZDG$O#gl8OqzxZw13knH41fW3 z8SnrYfCq>eh?>|rV2@a)H$3ttdrk!#+8iG^norc*-1Yl2^SEE1GJ6D4^{<9$yze)62YELwipaKYc>wY|%{>U;y1`7zOFKJ`~xR#tEJDeUylw(Lpy?4HW` z=dN?#gxS2T{dcdWg)Htq^i^NqWv8KS&5hn(cb#ML$CS0<1x|bC?VRGtgx2^%`U7;y z?_+niG!Dod8NKVRbcChA*8KfN9aGaZv|KYuhqzkp$=P-NvKt2u#A;f2%Z329gRGO& z))$dnV~eyAFKwWQ5MTfdsLOx{zyLfz#6Z-<&H;NQf_)+Z2CNc*2f%=Ci@*b503Pt$ zWB3Am0lolV5Wgxy51SyA z6GqHc8;+ptB1D=PnX2J!d6HyFF#(V7;m_0A#J<^5o-l8>NwGm42mk>f00e-*{|f=6 zIsA~X@22u*|@Kog};6Pu-RVMeLW>F2Wix^4%zBD2R3i}<#@oa zD|Xi%QybNxl~qZbPOV5k;`(0K>060D3*P}b-&KyBRpLES{V)s&N^Z&XztT;;c@ z*w$n)*@j1d8ob%mF;Hq>%JMPsLg=N&xKwLmv vKwLmvKwNk@{}r(jvGMWPIGEfzp3|U+(e)m&Hh5~SU7BYSdYoGCqaJ?)U%R8< literal 0 HcmV?d00001 diff --git a/crates/clawhdf5/tests/huge_chunks_interop.rs b/crates/clawhdf5/tests/huge_chunks_interop.rs index 52885e7..4a3b5cf 100644 --- a/crates/clawhdf5/tests/huge_chunks_interop.rs +++ b/crates/clawhdf5/tests/huge_chunks_interop.rs @@ -81,8 +81,15 @@ fn layout_of( let lm = msg(MessageType::DataLayout); 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(); + // 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, _) = - 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) } @@ -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 { + 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) -> Vec { + 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 = 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 { + 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 = (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 = (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, + ); +}