diff --git a/.gitignore b/.gitignore index cf90c2c..d0a22c0 100644 --- a/.gitignore +++ b/.gitignore @@ -7,3 +7,6 @@ weights/ .venv __pycache__/ .pytest_cache/ + +# Scratch files the heavy tests generate (huge_chunks_interop) +crates/*/tests/scratch/ diff --git a/crates/clawhdf5-format/src/chunk_cache.rs b/crates/clawhdf5-format/src/chunk_cache.rs index 6e3fee8..cca133d 100644 --- a/crates/clawhdf5-format/src/chunk_cache.rs +++ b/crates/clawhdf5-format/src/chunk_cache.rs @@ -899,7 +899,7 @@ mod tests { fn make_chunk(offsets: Vec, address: u64, size: u32) -> ChunkInfo { ChunkInfo { - chunk_size: size, + chunk_size: u64::from(size), filter_mask: 0, offsets, address, diff --git a/crates/clawhdf5-format/src/chunk_index.rs b/crates/clawhdf5-format/src/chunk_index.rs index 706bad0..1e083a0 100644 --- a/crates/clawhdf5-format/src/chunk_index.rs +++ b/crates/clawhdf5-format/src/chunk_index.rs @@ -109,7 +109,7 @@ pub struct ChunkMapping { /// File byte address of the compressed chunk. pub file_offset: u64, /// Size of the compressed chunk in the file. - pub file_size: u32, + pub file_size: u64, /// Filter mask (0 = all filters applied). pub filter_mask: u32, /// Pre-computed row-copy operations for assembling this chunk into output. @@ -338,7 +338,7 @@ mod tests { fn make_chunk(offsets: Vec, address: u64, size: u32) -> ChunkInfo { ChunkInfo { - chunk_size: size, + chunk_size: u64::from(size), filter_mask: 0, offsets, address, diff --git a/crates/clawhdf5-format/src/chunked_read.rs b/crates/clawhdf5-format/src/chunked_read.rs index 4d97adb..fa3f016 100644 --- a/crates/clawhdf5-format/src/chunked_read.rs +++ b/crates/clawhdf5-format/src/chunked_read.rs @@ -314,7 +314,7 @@ pub(crate) fn chunk_req( chunk_bytes: usize, wanted: bool, ) -> ExtentReq { - let len = c.chunk_size as usize; + let len = crate::addr::saturating_usize(c.chunk_size); ExtentReq { addr: c.address, len, @@ -521,7 +521,7 @@ pub fn decompress_all_chunks_with_stats_in( #[derive(Debug, Clone)] pub struct ChunkInfo { /// Size of chunk data in the file (after compression). - pub chunk_size: u32, + pub chunk_size: u64, /// Bitmask of filters that were NOT applied (0 = all applied). pub filter_mask: u32, /// N-dimensional offset of this chunk in dataset space. @@ -610,7 +610,7 @@ fn stored_element_size(dt: &Datatype, offset_size: u8) -> u64 { /// (`H5D__chunk_set_sizes`: "stored datatype size in chunk layout does not /// match datatype description"). Reading it anyway laid the chunks out with /// the wrong element size. -pub(crate) fn check_chunk_element_size( +pub fn check_chunk_element_size( layout: &DataLayout, datatype: &Datatype, offset_size: u8, @@ -1028,7 +1028,7 @@ fn parse_chunk_node( if node_level == 0 { chunks.push(stored.len()); stored.push(ChunkInfo { - chunk_size, + chunk_size: u64::from(chunk_size), filter_mask, offsets: keys[k..].to_vec(), address, @@ -1131,7 +1131,7 @@ pub fn generate_implicit_chunks_in_grid( } chunks.push(ChunkInfo { - chunk_size: chunk_byte_size as u32, + chunk_size: chunk_byte_size, filter_mask: 0, offsets, address: base_address.saturating_add(grid_idx.saturating_mul(chunk_byte_size)), @@ -1190,9 +1190,15 @@ fn read_btree_v2_chunks( } _ => return Err(bad("tree is not a chunk index")), }; - let unfiltered_bytes = checked_chunk_byte_len(chunk_dims, elem_size)?; - let unfiltered_bytes = - u32::try_from(unfiltered_bytes).map_err(|_| bad("chunk larger than 4 GiB"))?; + // A u64 whatever the platform: an unfiltered chunk's size is only + // recorded here; a chunk this platform cannot address fails when read. + let unfiltered_bytes = u64::try_from( + chunk_dims + .iter() + .try_fold(elem_size as u128, |acc, &c| acc.checked_mul(c as u128)) + .ok_or_else(|| bad("chunk size overflows"))?, + ) + .map_err(|_| bad("chunk size overflows"))?; let records = collect_btree_v2_records_in(file_data, &header, offset_size, length_size)?; let mut chunks = Vec::with_capacity(records.len()); @@ -1213,10 +1219,7 @@ fn read_btree_v2_chunks( pos += size_len; let mask = u32::from_le_bytes([data[pos], data[pos + 1], data[pos + 2], data[pos + 3]]); pos += 4; - ( - u32::try_from(size).map_err(|_| bad("stored chunk larger than 4 GiB"))?, - mask, - ) + (size, mask) }; let mut offsets = Vec::with_capacity(rank); for &dim in chunk_dims { @@ -1334,9 +1337,9 @@ pub fn list_chunks_in( // Single chunk — one chunk covering the entire dataset let chunk_byte_size = checked_chunk_byte_len(&chunk_dims, elem_size)?; let (csize, fmask) = if let Some(fs) = single_filtered_size { - (fs as u32, single_filter_mask.unwrap_or(0)) + (fs, single_filter_mask.unwrap_or(0)) } else { - (chunk_byte_size as u32, 0) + (chunk_byte_size as u64, 0) }; vec![ChunkInfo { chunk_size: csize, @@ -2124,7 +2127,7 @@ pub fn read_chunked_data_indexed_in( .iter() .zip(&hits) .map(|(m, hit)| { - let len = m.file_size as usize; + let len = crate::addr::saturating_usize(m.file_size); ExtentReq { addr: m.file_offset, len, @@ -2757,7 +2760,7 @@ mod tests { } chunk_infos.push(ChunkInfo { - chunk_size: chunk_bytes as u32, + chunk_size: chunk_bytes as u64, filter_mask: 0, offsets: vec![start as u64, 0], address: data_offset as u64, @@ -2869,7 +2872,7 @@ mod tests { .collect(); let stored = crate::filters::compress_chunk(&chunk, &pipeline, 4).unwrap(); chunks.push(ChunkInfo { - chunk_size: stored.len() as u32, + chunk_size: stored.len() as u64, filter_mask: 0, offsets: vec![r0 as u64, c0 as u64, 0], address: file.len() as u64, @@ -2943,7 +2946,7 @@ mod tests { let short = crate::filters::compress_chunk(&[1u8; 64], &pipeline, 4).unwrap(); for bad in [5usize, 11, 40] { chunks[bad].address = file.len() as u64; - chunks[bad].chunk_size = short.len() as u32; + chunks[bad].chunk_size = short.len() as u64; file.extend_from_slice(&short); } for _ in 0..20 { @@ -3133,7 +3136,7 @@ mod tests { file_data[data_offset..data_offset + compressed.len()].copy_from_slice(&compressed); chunk_infos.push(ChunkInfo { - chunk_size: compressed.len() as u32, + chunk_size: compressed.len() as u64, filter_mask: 0, offsets: vec![start as u64, 0], address: data_offset as u64, @@ -3215,7 +3218,7 @@ mod tests { file_data[data_offset..data_offset + chunk_size].copy_from_slice(&chunk_bytes); chunk_infos.push(ChunkInfo { - chunk_size: chunk_size as u32, + chunk_size: chunk_size as u64, filter_mask: 0, offsets: vec![row_start as u64, col_start as u64, 0], address: data_offset as u64, @@ -3340,7 +3343,7 @@ mod tests { assert_eq!(c.address, 0x1000 + i as u64 * chunk_byte_size as u64); assert_eq!(c.offsets, vec![i as u64 * 20]); assert_eq!(c.filter_mask, 0); - assert_eq!(c.chunk_size, chunk_byte_size as u32); + assert_eq!(c.chunk_size, chunk_byte_size as u64); } } diff --git a/crates/clawhdf5-format/src/chunked_write.rs b/crates/clawhdf5-format/src/chunked_write.rs index 10f1099..746a492 100644 --- a/crates/clawhdf5-format/src/chunked_write.rs +++ b/crates/clawhdf5-format/src/chunked_write.rs @@ -2125,7 +2125,7 @@ mod tests { for info in &infos { // Skipped chunks are stored at the chunk's size (shuffled). assert_eq!( - info.chunk_size == (c * 8) as u32, + info.chunk_size == (c * 8) as u64, info.filter_mask != 0, "{info:?}" ); diff --git a/crates/clawhdf5-format/src/extensible_array.rs b/crates/clawhdf5-format/src/extensible_array.rs index 42cf595..a17367b 100644 --- a/crates/clawhdf5-format/src/extensible_array.rs +++ b/crates/clawhdf5-format/src/extensible_array.rs @@ -224,7 +224,7 @@ fn read_element( }; Ok(( Some(ChunkInfo { - chunk_size: chunk_byte_size as u32, + chunk_size: chunk_byte_size, filter_mask: 0, offsets, address, @@ -259,7 +259,7 @@ fn read_element( }; Ok(( Some(ChunkInfo { - chunk_size: chunk_size as u32, + chunk_size, filter_mask, offsets, address, @@ -946,7 +946,7 @@ mod tests { assert_eq!(chunks.len(), 2); assert_eq!(chunks[0].address, base_addr); assert_eq!(chunks[0].offsets, vec![0]); - assert_eq!(chunks[0].chunk_size, chunk_byte_size as u32); + assert_eq!(chunks[0].chunk_size, chunk_byte_size as u64); assert_eq!(chunks[1].address, base_addr + chunk_byte_size); assert_eq!(chunks[1].offsets, vec![20]); } diff --git a/crates/clawhdf5-format/src/filters.rs b/crates/clawhdf5-format/src/filters.rs index 28070d4..48553c1 100644 --- a/crates/clawhdf5-format/src/filters.rs +++ b/crates/clawhdf5-format/src/filters.rs @@ -340,10 +340,23 @@ pub fn decompress_chunk_exact_with<'s>( } else { MAX_DECOMPRESS_SIZE }; - let size_hint = if ctx.max_output != 0 { + // The bound is the output's size when only size-preserving + // filters (shuffle, Fletcher32) remain to be undone; before + // any other filter (a second deflate, N-Bit, ...) it is only + // a ceiling, and the output starts smaller and grows: the + // inflater writes every byte it reserves, so reserving a + // 4 GiB chunk's bound for a stage a few MiB long would hold + // twice the chunk's memory. + let exact = pipeline.filters[..i].iter().enumerate().all(|(j, f)| { + filter_skipped(filter_mask, j) + || matches!(f.filter_id, FILTER_SHUFFLE | FILTER_FLETCHER32) + }); + let size_hint = if ctx.max_output == 0 { + input.len().saturating_mul(4).min(1 << 20) + } else if exact { ctx.max_output } else { - input.len().saturating_mul(4).min(1 << 20) + input.len().saturating_mul(4).min(ctx.max_output) }; let inflater = scratch .inflater diff --git a/crates/clawhdf5-format/src/fixed_array.rs b/crates/clawhdf5-format/src/fixed_array.rs index de80b7a..dbe00ae 100644 --- a/crates/clawhdf5-format/src/fixed_array.rs +++ b/crates/clawhdf5-format/src/fixed_array.rs @@ -383,7 +383,7 @@ fn parse_fa_element( offset_size: u8, element_size: u8, chunk_byte_size: u64, -) -> Result, FormatError> { +) -> Result, FormatError> { let os = offset_size as usize; if client_id == 0 { // Non-filtered: element is just the chunk address. @@ -393,7 +393,7 @@ fn parse_fa_element( return Ok(None); } let address = read_offset(file_data, abs, offset_size)?; - Ok(Some((address, chunk_byte_size as u32, 0))) + Ok(Some((address, chunk_byte_size, 0))) } else { // Filtered: address(offset_size) + chunk_size(variable) + filter_mask(4) let es = element_size as usize; @@ -418,7 +418,7 @@ fn parse_fa_element( file_data[fm_off + 2], file_data[fm_off + 3], ]); - Ok(Some((address, chunk_size as u32, filter_mask))) + Ok(Some((address, chunk_size, filter_mask))) } } @@ -688,7 +688,7 @@ mod tests { assert_eq!(c.address, base_addr + i as u64 * chunk_byte_size as u64); assert_eq!(c.offsets, vec![i as u64 * 20]); assert_eq!(c.filter_mask, 0); - assert_eq!(c.chunk_size, chunk_byte_size as u32); + assert_eq!(c.chunk_size, chunk_byte_size as u64); } } diff --git a/crates/clawhdf5-format/src/parallel_read.rs b/crates/clawhdf5-format/src/parallel_read.rs index 606173c..0594e42 100644 --- a/crates/clawhdf5-format/src/parallel_read.rs +++ b/crates/clawhdf5-format/src/parallel_read.rs @@ -426,7 +426,7 @@ mod tests { for i in 0..8u64 { let len = if short && i == 5 { 16 } else { 32 }; infos.push(ChunkInfo { - chunk_size: len as u32, + chunk_size: len as u64, filter_mask: 0, offsets: vec![i * 8], address: file.len() as u64, diff --git a/crates/clawhdf5-format/src/partial_read.rs b/crates/clawhdf5-format/src/partial_read.rs index 3361485..269c45b 100644 --- a/crates/clawhdf5-format/src/partial_read.rs +++ b/crates/clawhdf5-format/src/partial_read.rs @@ -243,6 +243,58 @@ fn copy_overlap( } } +/// Copy the part of unfiltered chunk `chunk` (`chunk_bytes` long, shape +/// `chunk_shape`) that overlaps the box into `out`, reading only the runs of +/// the overlap from the file. The whole chunk must still lie inside the +/// file, as it must when it is fetched whole. +#[allow(clippy::too_many_arguments)] +fn read_unfiltered_overlap( + file_data: &S, + chunk: &crate::chunked_read::ChunkInfo, + chunk_bytes: usize, + chunk_shape: &[u64], + elem_size: usize, + out: &mut [u8], + box_start: &[u64], + box_extent: &[u64], +) -> Result<(), FormatError> { + let rank = chunk_shape.len(); + let origin = &chunk.offsets[..rank]; + let (lo, extent): (Vec, Vec) = (0..rank) + .map(|d| { + let lo = origin[d].max(box_start[d]); + let hi = origin[d] + .saturating_add(chunk_shape[d]) + .min(box_start[d] + box_extent[d]); + (lo, hi.saturating_sub(lo)) + }) + .unzip(); + let file_len = crate::storage::len_usize(file_data); + let base = crate::addr::to_usize(chunk.address)?; + if base > file_len || chunk_bytes > file_len - base { + return Err(FormatError::UnexpectedEof { + expected: base.saturating_add(chunk_bytes), + available: file_len, + }); + } + let overlap = Selection::Hyperslab { + start: lo.iter().zip(origin).map(|(l, o)| l - o).collect(), + stride: vec![1; rank], + count: extent.clone(), + block: vec![1; rank], + }; + let rows = crate::gather::gather_storage( + file_data, + chunk.address, + chunk_bytes, + chunk_shape, + elem_size, + &overlap, + )?; + copy_overlap(&rows, &lo, &extent, out, box_start, box_extent, elem_size); + Ok(()) +} + /// Read `selection` without materialising the whole dataset, when that is /// possible and worthwhile. `Ok(None)` means "use the full-read path": an /// `All`/`None`/invalid selection, a layout this doesn't handle (compact, @@ -281,6 +333,39 @@ pub fn read_selection_in( offset_size: u8, length_size: u8, selection: &Selection, +) -> Result>, FormatError> { + read_selection_filled_in( + file_data, + layout, + dataspace, + elem_size, + pipeline, + offset_size, + length_size, + selection, + None, + ) +} + +/// [`read_selection_in`] for a dataset whose fill value is `fill` (one +/// element's bytes; `None` or all zeros is the default fill): the elements +/// of a chunked dataset's selection that lie in chunks never written read as +/// `fill`, as they do in a full read. Only the chunks the selection's +/// bounding box overlaps are read, so a selection of a few elements of a +/// dataset whose chunks are 4 GiB or more costs one decoded chunk (a +/// filtered chunk has to be decoded whole) or, unfiltered, only the bytes +/// it selects. +#[allow(clippy::too_many_arguments)] +pub fn read_selection_filled_in( + file_data: &S, + layout: &DataLayout, + dataspace: &Dataspace, + elem_size: usize, + pipeline: Option<&FilterPipeline>, + offset_size: u8, + length_size: u8, + selection: &Selection, + fill: Option<&[u8]>, ) -> Result>, FormatError> { let dims = &dataspace.dimensions; if dims.is_empty() || elem_size == 0 { @@ -341,8 +426,18 @@ pub fn read_selection_in( return Ok(None); } let mut boxed = alloc_output(checked_byte_len(box_elements, elem_size)?)?; + if let Some(fill) = fill.filter(|f| <[u8]>::len(f) == elem_size && f.iter().any(|&b| b != 0)) { + for element in boxed.chunks_exact_mut(elem_size) { + element.copy_from_slice(fill); + } + } match layout { + // No chunk was ever written: every element is the fill value. + DataLayout::Chunked { + btree_address: None, + .. + } if fill.is_some() => {} DataLayout::Chunked { btree_address: Some(_), .. @@ -373,6 +468,25 @@ pub fn read_selection_in( }) }) .collect(); + // An unfiltered chunk of a file that is not in memory: fetch + // only the rows the box needs, not the whole chunk (which may be + // 4 GiB or more). + let (direct, wanted): (Vec<_>, Vec<_>) = wanted.into_iter().partition(|c| { + file_data.as_contiguous().is_none() + && pipeline.is_none_or(|pl| all_filters_skipped(pl, c.filter_mask)) + }); + for chunk in direct { + read_unfiltered_overlap( + file_data, + chunk, + chunk_bytes, + &chunk_shape, + elem_size, + &mut boxed, + &box_start, + &box_extent, + )?; + } // Their stored bytes, batch by batch when the file is not in // memory; each batch's chunks are decoded into this thread's // reusable buffers before the next batch is fetched. diff --git a/crates/clawhdf5-format/tests/raw_fetch_bounds.rs b/crates/clawhdf5-format/tests/raw_fetch_bounds.rs index 8217323..a65d611 100644 --- a/crates/clawhdf5-format/tests/raw_fetch_bounds.rs +++ b/crates/clawhdf5-format/tests/raw_fetch_bounds.rs @@ -173,7 +173,7 @@ fn crafted() -> (Vec, Chunked, Vec) { assert!( chunks .iter() - .all(|c| c.chunk_size == HUGE && c.address == blob) + .all(|c| c.chunk_size == u64::from(HUGE) && c.address == blob) ); (bytes, ds, chunks) } @@ -313,7 +313,7 @@ fn large_reads_are_fetched_in_batches() { let data: Vec = (0..2 * CHUNK).map(|i| (i % 251) as u8).collect(); let chunks: Vec = (0..40u64) .map(|i| ChunkInfo { - chunk_size: CHUNK as u32, + chunk_size: CHUNK as u64, filter_mask: 0, offsets: vec![i * CHUNK as u64], address: (i % 2) * CHUNK as u64, diff --git a/crates/clawhdf5/src/reader.rs b/crates/clawhdf5/src/reader.rs index 4456d66..4426d39 100644 --- a/crates/clawhdf5/src/reader.rs +++ b/crates/clawhdf5/src/reader.rs @@ -1528,11 +1528,11 @@ impl<'f> Dataset<'f> { let ds = self.dataspace()?; let dl = self.data_layout()?; let pipeline = self.filter_pipeline()?; - // The selection reader knows nothing about fill values. When they - // matter — no storage at all, or a non-zero fill on a chunked (possibly - // sparse) dataset — select from a fill-aware full read instead. (The - // selection reader currently decodes the full dataset too, so this - // costs nothing extra.) + // Fill values matter when there is no storage at all, or a non-zero + // fill on a chunked (possibly sparse) dataset: a chunked dataset's + // selection is then read over a box of fill values; anything else + // (and a selection whose box is most of the dataset) is selected + // from a fill-aware full read. let fill = clawhdf5_format::fill_value::dataset_fill_value_from_storage( &self.file.data, &self.header.messages, @@ -1545,6 +1545,29 @@ impl<'f> Dataset<'f> { && !clawhdf5_format::fill_value::is_default(fill.as_deref())); if fill_matters { clawhdf5_format::partial_read::validate(selection, &ds.dimensions)?; + // A chunked dataset: read the chunks the selection touches over + // a box of fill values, not the whole dataset (which may be many + // GiB when its chunks are large). + if matches!(dl, DataLayout::Chunked { .. }) { + clawhdf5_format::chunked_read::check_chunk_element_size( + &dl, + &dt, + self.file.offset_size(), + )?; + if let Some(selected) = clawhdf5_format::partial_read::read_selection_filled_in( + &self.file.data, + &dl, + &ds, + dt.type_size() as usize, + pipeline.as_ref(), + self.file.offset_size(), + self.file.length_size(), + selection, + Some(fill.as_deref().unwrap_or(&[])), + )? { + return Ok(selected); + } + } let full = self.read_raw()?; return Ok(data_read::extract_selection_from_buffer( &full, diff --git a/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py b/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py new file mode 100644 index 0000000..16c18c4 --- /dev/null +++ b/crates/clawhdf5/tests/fixtures/gen_huge_chunks.py @@ -0,0 +1,109 @@ +"""Write HDF5 files whose chunks are 4 GiB or more, one dataset per chunk index. + + python gen_huge_chunks.py filtered OUT.h5 # huge_chunks_filtered.h5 + python gen_huge_chunks.py unfiltered OUT.h5 # generated at test time + +libhdf5 2.x writes a chunk of more than 0xFFFFFFFF bytes with layout message +version 5 (`H5D__chunk_construct`: "chunk size > 4GB requires +H5F_LIBVER_V200"), so never with a version-1 B-tree; a filtered chunk index +element of a version-5 layout stores the chunk's size in "size of lengths" +bytes (8). Every dataset is ` N: + ds[N : N + 10] = np.arange(100.0, 110.0) + if index == "implicit": + ds[2 * N - 10 : 2 * N] = np.arange(500.0, 510.0) + 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) + indexes = ["single", "farray", "earray", "btree2"] + if not filtered: + indexes.insert(1, "implicit") + for index in indexes: + make(fid, index, index, filtered) + fid.close() + + +if __name__ == "__main__": + main() diff --git a/crates/clawhdf5/tests/fixtures/huge_chunks_filtered.h5 b/crates/clawhdf5/tests/fixtures/huge_chunks_filtered.h5 new file mode 100644 index 0000000..60a8931 Binary files /dev/null and b/crates/clawhdf5/tests/fixtures/huge_chunks_filtered.h5 differ diff --git a/crates/clawhdf5/tests/huge_chunks_interop.rs b/crates/clawhdf5/tests/huge_chunks_interop.rs new file mode 100644 index 0000000..bf2ff19 --- /dev/null +++ b/crates/clawhdf5/tests/huge_chunks_interop.rs @@ -0,0 +1,296 @@ +//! Chunks of 4 GiB or more, which HDF5 2.0 writes (layout message version 5, +//! `H5F_LIBVER_V200`), in every chunk index libhdf5 uses for them. +//! +//! `fixtures/huge_chunks_filtered.h5` (written by libhdf5 2.0.0 through h5py +//! 3.16, `fixtures/gen_huge_chunks.py filtered`) holds one dataset per +//! filtered index — Single Chunk, Fixed Array, Extensible Array, v2 B-tree — +//! whose chunks are 2^29 + 1 `f64` (4 GiB + 8 bytes), stored deflated twice +//! so each takes about 20 KiB. Listing their chunks is cheap and always runs; +//! decoding one inflates 4 GiB, so those tests are opt-in: +//! `CLAWHDF5_HUGE_CHUNKS=1` (run them one at a time, `--test-threads=1`: +//! each needs about 4.5 GiB of memory). The same variable enables the +//! unfiltered tests, which have h5py write a sparse file (about 44 GiB long, +//! a few blocks on disk) under `tests/scratch/`, and the writer tests. +//! +//! libhdf5 never writes such a chunk with a version-1 B-tree (a chunk of more +//! than 0xFFFFFFFF bytes forces layout version 5 and with it the newer +//! indexes) and refuses to open one; `chunked_read`'s unit tests cover that. + +use std::path::{Path, PathBuf}; +use std::process::Command; + +use clawhdf5::File; +use clawhdf5_format::chunked_read::list_chunks; +use clawhdf5_format::data_layout::DataLayout; +use clawhdf5_format::dataspace::Dataspace; +use clawhdf5_format::message_type::MessageType; +use clawhdf5_format::object_header::ObjectHeader; +use clawhdf5_format::selection::Selection; +use clawhdf5_format::superblock::Superblock; + +/// Elements per chunk along the chunked axis: 4 GiB + 8 bytes of `f64`. +const N: u64 = (1 << 29) + 1; + +fn fixture() -> PathBuf { + Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/huge_chunks_filtered.h5") +} + +fn heavy() -> bool { + if std::env::var("CLAWHDF5_HUGE_CHUNKS").is_ok_and(|v| v == "1") { + return true; + } + eprintln!("SKIP: set CLAWHDF5_HUGE_CHUNKS=1 to decode 4 GiB chunks"); + false +} + +fn python() -> String { + std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string()) +} + +fn python_available() -> bool { + Command::new(python()) + .args(["-c", "import h5py"]) + .output() + .is_ok_and(|o| o.status.success()) +} + +fn interop_required() -> bool { + std::env::var("CLAWHDF5_REQUIRE_INTEROP").is_ok_and(|v| v == "1") +} + +/// The raw layout message version, the parsed layout, the dataspace and the +/// chunks of `name` in the in-memory file `data`. +fn layout_of( + data: &[u8], + name: &str, +) -> ( + u8, + DataLayout, + Dataspace, + Vec, +) { + let sb = Superblock::parse(data, 0).unwrap(); + let addr = clawhdf5_format::group_v2::resolve_path_any(data, &sb, name).unwrap(); + let hdr = ObjectHeader::parse(data, addr as usize, sb.offset_size, sb.length_size).unwrap(); + let msg = |t| { + hdr.messages + .iter() + .find(|m| m.msg_type == t) + .unwrap_or_else(|| panic!("{name}: no {t:?} message")) + }; + 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(); + let (chunks, _) = + list_chunks(data, &layout, &space, 8, sb.offset_size, sb.length_size).unwrap(); + (lm.data[0], layout, space, chunks) +} + +/// The fixture's datasets: name, chunk index type, shape, and the scaled +/// origins of the chunks libhdf5 wrote. +const FILTERED: &[(&str, u8, &[u64], &[&[u64]])] = &[ + ("single", 1, &[N], &[&[0]]), + ("farray", 3, &[N + 10], &[&[0], &[N]]), + ("earray", 4, &[N + 10], &[&[0], &[N]]), + ( + "btree2", + 5, + &[2, N + 10], + &[&[0, 0], &[0, N], &[1, 0], &[1, N]], + ), +]; + +/// Every index of the fixture: layout version 5, its chunks listed at the +/// right origins with their stored (deflated) sizes. The index elements +/// store a size in 8 bytes (libhdf5's `H5F_SIZEOF_SIZE` under layout +/// version 5), where version 4 would use 6 for a chunk this size. +#[test] +fn filtered_huge_chunk_indexes_list() { + let data = std::fs::read(fixture()).unwrap(); + for &(name, index, shape, origins) in FILTERED { + let (version, layout, space, mut chunks) = layout_of(&data, name); + assert_eq!(version, 5, "{name}: layout message version"); + let DataLayout::Chunked { + chunk_dimensions, + chunk_index_type, + .. + } = &layout + else { + panic!("{name}: not chunked: {layout:?}"); + }; + assert_eq!(*chunk_index_type, Some(index), "{name}"); + assert_eq!(chunk_dimensions.last(), Some(&8), "{name}"); + assert_eq!(space.dimensions, shape, "{name}"); + chunks.sort_by(|a, b| a.offsets.cmp(&b.offsets)); + let got: Vec<&[u64]> = chunks.iter().map(|c| &c.offsets[..shape.len()]).collect(); + assert_eq!(got, origins, "{name}"); + for c in &chunks { + assert!( + (10_000..40_000).contains(&c.chunk_size), + "{name}: stored size {} of chunk {:?}", + c.chunk_size, + c.offsets + ); + assert_eq!(c.filter_mask, 0); + } + } +} + +/// `f64` values of `sel` in dataset `name` of `file`. +fn sel(file: &File, name: &str, start: &[u64], count: &[u64]) -> Vec { + let ds = file.dataset(name).unwrap(); + let s = Selection::Hyperslab { + start: start.to_vec(), + stride: vec![1; start.len()], + count: count.to_vec(), + block: vec![1; start.len()], + }; + ds.read_f64_selection(&s).unwrap() +} + +fn f(range: std::ops::Range) -> Vec { + range.map(f64::from).collect() +} + +/// Reads that touch a few elements of each 4 GiB chunk: written values, the +/// fill value (-1) next to them, and the edge of the dataset. +#[test] +fn filtered_huge_chunks_read() { + if !heavy() { + return; + } + let file = File::open(fixture()).unwrap(); + let mut first = f(0..10); + first.extend([-1.0; 2]); + for name in ["single", "farray", "earray"] { + assert_eq!(sel(&file, name, &[0], &[12]), first, "{name}"); + } + assert_eq!(sel(&file, "single", &[N - 2], &[2]), [-1.0; 2]); + let mut edge = vec![-1.0; 2]; + edge.extend(f(100..110)); + for name in ["farray", "earray"] { + assert_eq!(sel(&file, name, &[N - 2], &[12]), edge, "{name}"); + } + assert_eq!(sel(&file, "btree2", &[0, 0], &[1, 10]), f(0..10)); + assert_eq!(sel(&file, "btree2", &[1, N], &[1, 10]), f(300..310)); + assert_eq!( + sel(&file, "btree2", &[0, N - 1], &[2, 2]), + [-1.0, 100.0, -1.0, 300.0] + ); +} + +/// A directory under `tests/scratch/` (on disk: the sparse files must not +/// land on a tmpfs `/tmp`), removed when dropped. +fn scratch() -> tempfile::TempDir { + let root = Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/scratch"); + std::fs::create_dir_all(&root).unwrap(); + tempfile::tempdir_in(root).unwrap() +} + +/// Have h5py (libhdf5 2.x) write `fixtures/gen_huge_chunks.py`'s file for +/// `mode` into `dir`; `None` when there is no h5py (a failure under +/// `CLAWHDF5_REQUIRE_INTEROP=1`). +fn generate(dir: &Path, mode: &str) -> Option { + if !python_available() { + assert!( + !interop_required(), + "CLAWHDF5_REQUIRE_INTEROP=1 but python3 with h5py is not available" + ); + eprintln!("SKIP: python3 with h5py not available"); + return None; + } + let script = Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/gen_huge_chunks.py"); + let path = dir.join(format!("{mode}.h5")); + let out = Command::new(python()) + .arg(&script) + .arg(mode) + .arg(&path) + .output() + .unwrap(); + assert!( + out.status.success(), + "gen_huge_chunks.py {mode} failed:\n{}", + String::from_utf8_lossy(&out.stderr) + ); + Some(path) +} + +/// Positioned reads of a file (no mmap), counting the bytes read. +struct Counting { + file: std::fs::File, + len: u64, + read: std::sync::atomic::AtomicU64, +} + +impl clawhdf5_format::storage::Storage for Counting { + fn read_at( + &self, + offset: u64, + len: usize, + ) -> Result, clawhdf5_format::error::FormatError> { + use std::os::unix::fs::FileExt; + let len = len.min(usize::try_from(self.len.saturating_sub(offset)).unwrap_or(usize::MAX)); + let mut buf = vec![0; len]; + self.file.read_exact_at(&mut buf, offset).unwrap(); + self.read + .fetch_add(len as u64, std::sync::atomic::Ordering::Relaxed); + Ok(buf.into()) + } + + fn len(&self) -> u64 { + self.len + } +} + +fn check_unfiltered(file: &File) { + let mut first = f(0..10); + first.extend([0.0; 2]); + for name in ["single", "implicit", "farray", "earray"] { + assert_eq!(sel(file, name, &[0], &[12]), first, "{name}"); + } + assert_eq!(sel(file, "single", &[N - 2], &[2]), [0.0; 2]); + let mut edge = vec![0.0; 2]; + edge.extend(f(100..110)); + for name in ["implicit", "farray", "earray"] { + assert_eq!(sel(file, name, &[N - 2], &[12]), edge, "{name}"); + } + let mut last = vec![0.0; 2]; + last.extend(f(500..510)); + assert_eq!(sel(file, "implicit", &[2 * N - 12], &[12]), last); + assert_eq!(sel(file, "btree2", &[0, 0], &[1, 10]), f(0..10)); + assert_eq!(sel(file, "btree2", &[1, N], &[1, 10]), f(300..310)); + assert_eq!( + sel(file, "btree2", &[0, N - 1], &[2, 2]), + [0.0, 100.0, 0.0, 300.0] + ); +} + +/// Unfiltered chunks of 4 GiB + 8 bytes in every index libhdf5 gives them +/// (Single Chunk, Implicit, Fixed Array, Extensible Array, v2 B-tree), in a +/// sparse file h5py writes: read through a memory map, and through +/// positioned reads, where a selection reads only the rows it needs. +#[test] +fn unfiltered_huge_chunks_read() { + if !heavy() { + return; + } + let dir = scratch(); + let Some(path) = generate(dir.path(), "unfiltered") else { + return; + }; + check_unfiltered(&File::open(&path).unwrap()); + + let file = std::fs::File::open(&path).unwrap(); + let len = file.metadata().unwrap().len(); + assert!(len > 40 << 30, "{len}"); + let storage = std::sync::Arc::new(Counting { + file, + len, + read: 0.into(), + }); + let positioned = File::open_storage(storage.clone()).unwrap(); + check_unfiltered(&positioned); + // Every read above together: metadata and a few rows, not 4 GiB chunks. + let read = storage.read.load(std::sync::atomic::Ordering::Relaxed); + assert!(read < 1 << 20, "{read} bytes read"); +} diff --git a/crates/clawhdf5/tests/parallel_integration.rs b/crates/clawhdf5/tests/parallel_integration.rs index bdcf4f0..964d3c8 100644 --- a/crates/clawhdf5/tests/parallel_integration.rs +++ b/crates/clawhdf5/tests/parallel_integration.rs @@ -283,7 +283,7 @@ mod parallel_tests { for (i, chunk) in compressed_chunks.iter().enumerate() { file_data[offset..offset + chunk.len()].copy_from_slice(chunk); chunk_infos.push(ChunkInfo { - chunk_size: chunk.len() as u32, + chunk_size: chunk.len() as u64, filter_mask: 0, offsets: vec![(i * chunk_elems) as u64, 0], address: offset as u64,