diff --git a/CHANGELOG.md b/CHANGELOG.md index acd24d6..4f65734 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -33,6 +33,61 @@ ok (baseline 600), 0 our-error, 0 mismatch, 3 ref-bug, 92 h5py-cannot-read, no panic/hang/crash/oom. +### Remote files in the browser: fewer round trips to list a group or open a dataset (2026-09-27) +- **Opening one dataset of a v1 (symbol table) group no longer reads the + whole group.** A name is looked up down the group's B-tree, as + libhdf5's `H5G__stab_lookup` does (binary search on the node keys, names + in the local heap compared bytewise, then one symbol table node); only + when that finds no hard link of that name (a soft link, or a B-tree out + of name order, where libhdf5 would report it missing) is every entry + read, as before. Local files benefit too (a lookup read O(entries)). + In a group holding two entries of one name, the B-tree's is now the one + found, as in libhdf5. +- **`Storage::hint(offset, len)`** (clawhdf5-format): a parser says what + it reads next — a B-tree node's or symbol table node's body, an object + header's first chunk and continuation chunks, the symbol table nodes a + B-tree leaf names, a dense group's name index header and heap blocks, + and in a listing every child's object header. Every backend ignores it + but the browser's restartable reader (`clawhdf5_wasm::lazy`), which + fetches the hinted blocks it lacks together with the blocks a pass + missed, within the call's `maxFetch` budget; a pass that misses + nothing ignores them, so a hint never adds a round trip, and results + never depend on hints. +- The v1 and v2 B-tree walks of a listing descend into every child after + one fails (before, the siblings were only read, so their subtrees came + a pass later), then return the first error: same results and errors. +- Counted on tank, 2026-09-27, with `CLAWHDF5_WASM_LIST_FILE= + CLAWHDF5_WASM_READ=/d1500 cargo test --release -p clawhdf5-wasm --test + lazy listing_cost_of_a_given_file -- --nocapture` on an h5py file like + the reviewer's (3000 datasets of 16384 `f32`, 198 MB, h5py 3.16 / + HDF5 2.0), passes / requests / bytes, before -> after: + + | file, block size | `list('/')` | open + read one dataset | + |---|---|---| + | earliest, 1 MiB | 6 / 73 / 192.5 MB -> 4 / 68 / 192.5 MB | 7 / 74 / 193.6 MB -> 6 / 5 / 5.2 MB | + | earliest, 64 KiB | 8 / 531 / 35.2 MB -> 5 / 530 / 35.3 MB | 9 / 515 / 34.1 MB -> 8 / 7 / 0.52 MB | + | latest, 1 MiB | 9 / 98 / 196.5 MB -> 5 / 86 / 196.5 MB | 8 / 7 / 6.7 MB -> 7 / 7 / 6.7 MB | + | latest, 64 KiB | 11 / 452 / 29.6 MB -> 6 / 454 / 30.5 MB | 9 / 8 / 0.58 MB -> 8 / 8 / 0.58 MB | + + The listing's passes now follow the depth of the group's index (the + chain index levels -> symbol table nodes or heap objects -> child + headers); its bytes are the child headers, which h5py spreads through + the file (at 1 MiB blocks most of it). The whole corpus read lazily + (`CLAWHDF5_WASM_CORPUS`, 656 files, every object listed, described and + read): 27 513 -> 27 496 passes, 3 001 -> 2 962 requests and 194.8 -> + 195.3 MB at 64 KiB blocks; 41 342 -> 37 517 passes, 32 578 -> 29 511 + requests, 65.4 -> 65.7 MB at 512 B. The Node and Chromium suite + (`examples/wasm-viewer/test/run.sh`) passes unchanged (the 200 MB file + still takes 5 requests, 6 MiB); its corpus comparison fetched 33.54 -> + 33.61 MB. +- Tests: the listing budgets (`listing_a_large_group_takes_a_few_passes`, + 512-byte blocks) are tightened to the new counts (FileBuilder, 600 + children: 5 -> 4 passes; h5py, 2000 children: 8 -> 5 and 11 -> 6); new + `reading_one_dataset_of_a_large_group_fetches_a_few_blocks` (h5py + earliest: 529 requests, 333 kB -> at most 6 requests, 27 kB), v1 + lookups against the listing (and with a name moved out of B-tree + order), hints riding only on misses and within the fetch budget. + ### Deterministic errors on damaged chunked datasets (2026-09-27) - A read through the file's chunk cache listed a damaged dataset's chunks in hash-map order, seeded per `File`, so two opens of the same file could diff --git a/crates/clawhdf5-format/src/btree_v1.rs b/crates/clawhdf5-format/src/btree_v1.rs index ea4e5b1..78b2329 100644 --- a/crates/clawhdf5-format/src/btree_v1.rs +++ b/crates/clawhdf5-format/src/btree_v1.rs @@ -92,6 +92,8 @@ impl BTreeV1Node { // + left_sibling(offset_size) + right_sibling(offset_size) let os = offset_size as usize; let header_size = 8 + os * 2; + // The body is read once the header says how long it is. + file.hint(offset, NODE_HINT_LEN); let header = read_exact_at(file, offset, header_size)?; let file_data: &[u8] = &header; // The header's read checked that `offset + header_size` fits. @@ -156,7 +158,18 @@ impl BTreeV1Node { } /// Maximum recursion depth for B-tree traversal (malformed data protection). -const MAX_BTREE_DEPTH: usize = 64; +pub(crate) const MAX_BTREE_DEPTH: usize = 64; + +/// What a symbol table node takes with libhdf5's default group leaf K (4): +/// its 8-byte header and 2K entries of 40 bytes (8-byte offsets). Hinted +/// before one is read ([`Storage::hint`]); a node of another size is read +/// all the same. +const SNOD_HINT_LEN: usize = 8 + 8 * 40; + +/// What a group B-tree node takes with libhdf5's default internal K (16): +/// its header (24 bytes with 8-byte offsets), 2K + 1 keys and 2K children +/// of 8 bytes. Hinted before one is read. +const NODE_HINT_LEN: usize = 24 + (2 * 16 + 1 + 2 * 16) * 8; /// Collect all leaf-level child addresses (SNOD addresses) by traversing the B-tree. pub fn collect_symbol_table_nodes( @@ -196,20 +209,22 @@ fn collect_symbol_table_nodes_inner( } if node.node_level == 0 { - // Leaf: children are SNOD addresses + // Leaf: children are SNOD addresses, read next (see + // `Storage::hint`). + for &snod in &node.children { + file.hint(snod, SNOD_HINT_LEN); + } Ok(node.children) } else { - // Internal: recurse into children. After the first child that - // fails, the others are only read (as `storage::touch` does), not - // descended into; that error is returned. + // Internal: recurse into children. A child that fails does not + // stop the walk: the others are still descended into (reading, not + // using, what they hold), then the first error is returned. The + // result and the error are those of stopping at the first failure; + // a storage that records what it lacks (see `storage::touch`) + // learns every node the walk can reach in one attempt. let mut result = Vec::new(); let mut failed = None; for &child_addr in &node.children { - if failed.is_some() { - // Parsing reads the node's header, then its body. - let _ = BTreeV1Node::parse_in(file, child_addr, offset_size, length_size); - continue; - } match collect_symbol_table_nodes_inner( file, child_addr, @@ -217,8 +232,11 @@ fn collect_symbol_table_nodes_inner( length_size, depth + 1, ) { - Ok(child_snods) => result.extend(child_snods), - Err(e) => failed = Some(e), + Ok(child_snods) if failed.is_none() => result.extend(child_snods), + Ok(_) => {} + Err(e) => { + failed.get_or_insert(e); + } } } match failed { diff --git a/crates/clawhdf5-format/src/btree_v2.rs b/crates/clawhdf5-format/src/btree_v2.rs index f721f40..21010f7 100644 --- a/crates/clawhdf5-format/src/btree_v2.rs +++ b/crates/clawhdf5-format/src/btree_v2.rs @@ -194,10 +194,17 @@ const MAX_DEPTH: u16 = 64; /// Take `n` records from the traversal's budget, or refuse the tree. fn spend(budget: &mut usize, n: usize) -> Result<(), FormatError> { - *budget = budget - .checked_sub(n) - .ok_or(FormatError::NestingDepthExceeded)?; - Ok(()) + match budget.checked_sub(n) { + Some(left) => { + *budget = left; + Ok(()) + } + None => { + // Spent: a walk that goes on after a failure stops here. + *budget = 0; + Err(FormatError::NestingDepthExceeded) + } + } } /// Collect all records from a B-tree v2 by traversing from the root. @@ -504,16 +511,18 @@ fn collect_internal_records( // Interleave: child[0], record[0], child[1], record[1], ..., child[nr] // We collect child[0] records, then record[0], then child[1], etc. - // After the first child that fails, the others are only touched (see - // `storage::touch`); that error is returned. + // A child that fails does not stop the walk: the others are still + // descended into (their records are dropped with the result), then the + // first error is returned, as when stopping there. A storage that + // records what it lacks (see `storage::touch`) so learns every node the + // walk can reach in one attempt. The record budget is spent as before, + // so the walk is no longer than a successful one. let mut failed = None; for (i, &(child_addr, child_nrec)) in node.children.iter().enumerate() { - if failed.is_some() { - let len = usize::try_from(node_size) - .unwrap_or(usize::MAX) - .min(1 << 16); - crate::storage::touch(file, child_addr, len); - continue; + if failed.is_some() && *budget == 0 { + // The record budget is spent: the tree is refused, and a walk + // over what is left could be as long as the one it bounds. + break; } if let Err(e) = (|| -> Result<(), FormatError> { if child_depth == 0 { @@ -553,7 +562,7 @@ fn collect_internal_records( } Ok(()) })() { - failed = Some(e); + failed.get_or_insert(e); } } diff --git a/crates/clawhdf5-format/src/fractal_heap.rs b/crates/clawhdf5-format/src/fractal_heap.rs index 83cfd33..9b39c6a 100644 --- a/crates/clawhdf5-format/src/fractal_heap.rs +++ b/crates/clawhdf5-format/src/fractal_heap.rs @@ -10,7 +10,7 @@ use crate::addr::to_usize; use crate::btree_v2::{BTreeV2Header, find_btree_v2_records_in}; use crate::error::FormatError; use crate::filter_pipeline::FilterPipeline; -use crate::storage::{Storage, Window, len_usize, read_exact_at}; +use crate::storage::{Storage, Window, len_usize, read_exact_at, read_upto}; /// Parsed fractal heap header (signature "FRHP"). #[derive(Debug, Clone)] @@ -132,6 +132,9 @@ fn heap_id_type(first: u8) -> Result { const BTREE_HUGE_INDIRECT: u8 = 1; const BTREE_HUGE_INDIRECT_FILTERED: u8 = 2; +/// Most child entries [`FractalHeapHeader::hint_managed_blocks`] walks. +const MAX_HINTED_BLOCKS: usize = 4096; + impl FractalHeapHeader { /// Parse a fractal heap header at the given offset. pub fn parse( @@ -667,15 +670,8 @@ impl FractalHeapHeader { "fractal heap: maximum recursion depth exceeded".into(), )); } - let block_offset_bytes = (self.max_heap_size as usize).div_ceil(8); - let iblock_header = 5 + offset_size as usize + block_offset_bytes; let nrows_usize = nrows as usize; - // Rows below max_direct_rows hold direct blocks; rows at/above hold - // child indirect blocks. (NOT the FRHP "starting rows" field.) - let start_indirect = self.max_direct_rows(); - let max_direct_rows = nrows_usize.min(start_indirect); - // The block up to its last child entry. The walk below reads // entries in order and stops at the one covering the target, which // the geometry alone locates, so the first window ends there: a @@ -685,32 +681,11 @@ impl FractalHeapHeader { // over the whole block. Either window holds what it was asked for or // ends at the end of the file, so its bounds checks are the // whole-file ones. - let direct_entry = usize::from(offset_size) - + if self.filter_pipeline.is_some() { - usize::from(self.length_size) + 4 - } else { - 0 - }; - let direct_entries = max_direct_rows.saturating_mul(usize::from(self.table_width)); - let entries_len = |n: usize| { - n.min(direct_entries) - .saturating_mul(direct_entry) - .saturating_add( - n.saturating_sub(direct_entries) - .saturating_mul(usize::from(offset_size)), - ) - }; - let all_entries = direct_entries.saturating_add( - nrows_usize - .saturating_sub(start_indirect) - .saturating_mul(usize::from(self.table_width)), - ); - let block_len = iblock_header.saturating_add(entries_len(all_entries)); + let layout = self.indirect_layout(nrows_usize, offset_size); + let block_len = layout.len_upto(layout.all_entries); let target_entry = self.indirect_entry_for(nrows_usize, iblock_heap_offset, target_offset); let first_len = target_entry.map_or(block_len, |i| { - iblock_header - .saturating_add(entries_len(i.saturating_add(1))) - .min(block_len) + layout.len_upto(i.saturating_add(1)).min(block_len) }); let mut next = self.walk_indirect_block( &Window::read(file, iblock_addr as u64, first_len)?, @@ -755,6 +730,140 @@ impl FractalHeapHeader { } } + /// Where the child entries of an indirect block of `nrows` rows are. + fn indirect_layout(&self, nrows: usize, offset_size: u8) -> IndirectLayout { + let block_offset_bytes = (self.max_heap_size as usize).div_ceil(8); + // Rows below max_direct_rows hold direct blocks; rows at/above hold + // child indirect blocks. (NOT the FRHP "starting rows" field.) + let start_indirect = self.max_direct_rows(); + let direct_entries = nrows + .min(start_indirect) + .saturating_mul(usize::from(self.table_width)); + IndirectLayout { + header: 5 + usize::from(offset_size) + block_offset_bytes, + direct_entry: usize::from(offset_size) + + if self.filter_pipeline.is_some() { + usize::from(self.length_size) + 4 + } else { + 0 + }, + direct_entries, + indirect_entry: usize::from(offset_size), + all_entries: direct_entries.saturating_add( + nrows + .saturating_sub(start_indirect) + .saturating_mul(usize::from(self.table_width)), + ), + } + } + + /// Hint the heap's root block (see [`Storage::hint`]): every managed + /// object is read through it, and a storage that fetches between + /// attempts can fetch it along with whatever else the attempt missed + /// (the name index read before any object, say). + pub fn hint_root_block(&self, file: &S) { + if is_undefined(self.root_block_address, self.offset_size) { + return; + } + let len = if self.current_rows_in_root_indirect_block == 0 { + if self.filter_pipeline.is_some() { + self.root_direct_block_filtered_size + } else { + self.starting_block_size + } + } else { + let layout = self.indirect_layout( + usize::from(self.current_rows_in_root_indirect_block), + self.offset_size, + ); + layout.len_upto(layout.all_entries) as u64 + }; + file.hint( + self.root_block_address, + usize::try_from(len).unwrap_or(usize::MAX), + ); + } + + /// Hint every managed direct block of the heap (see + /// [`Storage::hint`]), for a caller about to read all of its objects (a + /// dense group's listing). The indirect blocks leading to them are read + /// here, as every object read goes through them, a few levels deep and + /// up to [`MAX_HINTED_BLOCKS`] entries; direct blocks are only hinted. + /// Nothing is returned and no error: a storage that has the file in + /// memory skips it, and the objects are read (and checked) as before. + pub fn hint_managed_blocks(&self, file: &S) { + if file.as_contiguous().is_some() + || is_undefined(self.root_block_address, self.offset_size) + || self.current_rows_in_root_indirect_block == 0 + { + // A direct root block is what `hint_root_block` hints. + return; + } + let mut budget = MAX_HINTED_BLOCKS; + self.hint_indirect_block( + file, + self.root_block_address, + usize::from(self.current_rows_in_root_indirect_block), + 0, + &mut budget, + ); + } + + fn hint_indirect_block( + &self, + file: &S, + addr: u64, + nrows: usize, + depth: usize, + budget: &mut usize, + ) { + let os = self.offset_size; + let layout = self.indirect_layout(nrows, os); + let len = layout.len_upto(layout.all_entries); + if depth > 4 || len > 1 << 20 { + return; + } + let Ok(bytes) = read_upto(file, addr, len) else { + return; + }; + if bytes.len() < len || bytes.get(..4) != Some(b"FHIB".as_slice()) { + return; + } + let start_indirect = self.max_direct_rows(); + let mut pos = layout.header; + for row in 0..nrows { + for _ in 0..self.table_width { + let Some(left) = budget.checked_sub(1) else { + return; + }; + *budget = left; + let Ok(child) = read_offset(&bytes, pos, os) else { + return; + }; + if row < start_indirect { + let size = if self.filter_pipeline.is_some() { + match read_offset(&bytes, pos + usize::from(os), self.length_size) { + Ok(n) => n, + Err(_) => return, + } + } else { + self.block_size_for_row(row) + }; + pos += layout.direct_entry; + if !is_undefined(child, os) { + file.hint(child, usize::try_from(size).unwrap_or(usize::MAX)); + } + } else { + pos += layout.indirect_entry; + if !is_undefined(child, os) { + let rows = self.rows_for_size(self.block_size_for_row(row)); + self.hint_indirect_block(file, child, usize::from(rows), depth + 1, budget); + } + } + } + } + } + /// Which child entry of an indirect block (numbered in walk order: /// direct rows, then indirect rows) covers `target_offset`, from the /// doubling-table geometry alone — the entry @@ -939,6 +1048,31 @@ impl FractalHeapHeader { } } +/// Where an indirect block's child entries are: after its header, the +/// direct blocks' entries (address, and for a filtered heap the stored +/// size and filter mask), then the child indirect blocks' (address). +struct IndirectLayout { + header: usize, + direct_entry: usize, + direct_entries: usize, + indirect_entry: usize, + all_entries: usize, +} + +impl IndirectLayout { + /// Bytes from the block's start to the end of its first `n` entries. + fn len_upto(&self, n: usize) -> usize { + let direct = n.min(self.direct_entries); + self.header + .saturating_add(direct.saturating_mul(self.direct_entry)) + .saturating_add( + n.min(self.all_entries) + .saturating_sub(direct) + .saturating_mul(self.indirect_entry), + ) + } +} + /// A managed direct block's location, extent and (for a filtered heap) its /// stored size and filter mask. /// The child of an indirect block that covers a heap offset. diff --git a/crates/clawhdf5-format/src/group_v1.rs b/crates/clawhdf5-format/src/group_v1.rs index bf7dcc4..4a930dc 100644 --- a/crates/clawhdf5-format/src/group_v1.rs +++ b/crates/clawhdf5-format/src/group_v1.rs @@ -4,7 +4,7 @@ use alloc::{string::String, vec::Vec}; use crate::addr::checked_addr; -use crate::btree_v1::collect_symbol_table_nodes_in; +use crate::btree_v1::{BTreeV1Node, collect_symbol_table_nodes_in}; use crate::error::FormatError; use crate::local_heap::LocalHeap; use crate::message_type::MessageType; @@ -48,7 +48,7 @@ pub fn resolve_v1_group_entries_in( offset_size: u8, length_size: u8, ) -> Result, FormatError> { - let entries = v1_group_entries(file_data, sym_table_msg, offset_size, length_size)?; + let entries = v1_group_entries(file_data, sym_table_msg, offset_size, length_size, true)?; if entries.iter().any(|e| e.name.is_empty()) { return Err(FormatError::InvalidLinkName); } @@ -57,11 +57,16 @@ pub fn resolve_v1_group_entries_in( /// Every entry of a v1 group, empty names included — for looking a name up, /// which never matches an empty name. +/// +/// With `hint_headers` (a listing, whose children are usually opened +/// next), each entry's object header is hinted (see +/// [`Storage::hint`]) as soon as its symbol table node is read. pub(crate) fn v1_group_entries( file_data: &S, sym_table_msg: &SymbolTableMessage, offset_size: u8, length_size: u8, + hint_headers: bool, ) -> Result, FormatError> { // Parse local heap let heap = LocalHeap::parse_in( @@ -93,12 +98,23 @@ pub(crate) fn v1_group_entries( // `storage::touch` does); that error is returned. let mut failed = None; for snod_addr in snod_addrs { + let snod = checked_addr(snod_addr) + .and_then(|a| SymbolTableNode::parse_in(file_data, a, offset_size)); + if hint_headers && let Ok(snod) = &snod { + for entry in &snod.entries { + if entry.object_header_address != u64::MAX { + file_data.hint( + entry.object_header_address, + crate::object_header::OBJECT_HEADER_HINT_LEN, + ); + } + } + } if failed.is_some() { - let _ = SymbolTableNode::parse_in(file_data, snod_addr, offset_size); continue; } - let mut node = || -> Result<(), FormatError> { - let snod = SymbolTableNode::parse_in(file_data, checked_addr(snod_addr)?, offset_size)?; + let node = || -> Result<(), FormatError> { + let snod = snod?; for entry in &snod.entries { // Like libhdf5, look at the heap's free list only once a name // is needed: an empty group with a damaged heap still lists. @@ -126,6 +142,95 @@ pub(crate) fn v1_group_entries( } } +/// The entry called `name` in a v1 group, looked up as libhdf5 looks it up +/// (`H5G__stab_lookup`: `H5B_find` down the group's B-tree, then +/// `H5G__node_found` in one symbol table node) instead of by reading every +/// entry: at each node, the child whose key interval holds the name +/// (left key < name <= right key, keys being names in the local heap, +/// compared bytewise as `strcmp` does) is found by binary search. Reads +/// O(depth) nodes and names, where listing reads the whole group. +/// +/// `Ok(None)` when the search does not lead to the name. In a group whose +/// B-tree is out of name order (damaged, or made by hand) that does not +/// prove it absent, so callers then fall back to reading every entry; +/// libhdf5 would report it missing. +pub(crate) fn find_v1_entry( + file_data: &S, + sym_table_msg: &SymbolTableMessage, + name: &str, + offset_size: u8, + length_size: u8, +) -> Result, FormatError> { + let heap = LocalHeap::parse_in( + file_data, + checked_addr(sym_table_msg.local_heap_address)?, + offset_size, + length_size, + )?; + // The keys and names are read from the heap's data segment one at a + // time: hint it (see `Storage::hint`). + file_data.hint( + heap.data_segment_address, + usize::try_from(heap.data_segment_size).map_or(1 << 20, |n| n.min(1 << 20)), + ); + // As in listing: the heap's free list is checked before a name is used. + let mut heap_checked = false; + let mut name_at = |offset: u64| -> Result { + if !heap_checked { + heap.validate_free_list_in(file_data, length_size)?; + heap_checked = true; + } + heap.read_string_in(file_data, offset) + }; + let want = name.as_bytes(); + let mut address = sym_table_msg.btree_address; + for _ in 0..=crate::btree_v1::MAX_BTREE_DEPTH { + let node = + BTreeV1Node::parse_in(file_data, checked_addr(address)?, offset_size, length_size)?; + if node.node_type != 0 { + return Err(FormatError::InvalidBTreeNodeType(node.node_type)); + } + // H5B_find's binary search with H5G__node_cmp3: go left when the + // name sorts at or before the left key, right when after the right + // key; otherwise this child holds it. + let (mut lo, mut hi) = (0usize, node.children.len()); + let mut child = None; + while lo < hi { + let i = lo + (hi - lo) / 2; + let (Some(&left), Some(&right)) = (node.keys.get(i), node.keys.get(i + 1)) else { + return Ok(None); + }; + if want <= name_at(left)?.as_bytes() { + hi = i; + } else if want > name_at(right)?.as_bytes() { + lo = i + 1; + } else { + child = Some(node.children[i]); + break; + } + } + let Some(child) = child else { + return Ok(None); + }; + if node.node_level > 0 { + address = child; + continue; + } + let snod = SymbolTableNode::parse_in(file_data, checked_addr(child)?, offset_size)?; + for entry in &snod.entries { + if name_at(entry.link_name_offset)?.as_bytes() == want { + return Ok(Some(GroupEntry { + name: String::from(name), + object_header_address: entry.object_header_address, + cache_type: entry.cache_type, + })); + } + } + return Ok(None); + } + Err(FormatError::NestingDepthExceeded) +} + /// Symbol table cache type for a soft link: the scratch pad's first four bytes /// are the local-heap offset of the link's target path, and the entry's object /// header address is undefined. @@ -298,7 +403,13 @@ pub fn resolve_path_in( let mut current_sym_table = root_sym_table.clone(); for (i, component) in components.iter().enumerate() { - let entries = v1_group_entries(file_data, ¤t_sym_table, offset_size, length_size)?; + let entries = v1_group_entries( + file_data, + ¤t_sym_table, + offset_size, + length_size, + false, + )?; let found = entries.iter().find(|e| e.name == *component); match found { diff --git a/crates/clawhdf5-format/src/group_v2.rs b/crates/clawhdf5-format/src/group_v2.rs index a3e026d..e902239 100644 --- a/crates/clawhdf5-format/src/group_v2.rs +++ b/crates/clawhdf5-format/src/group_v2.rs @@ -101,18 +101,50 @@ fn resolve_compact_entries( Ok(entries) } +/// What a version-2 B-tree header takes with 8-byte offsets and lengths +/// (22 bytes of fields, the root node's address and record count, and the +/// checksum), rounded up: hinted before one is read. +const BTREE_V2_HEADER_HINT_LEN: usize = 64; + +/// The fractal heap of a dense group. The name index's header and the +/// heap's root block are read next, whatever the lookup: they are hinted +/// (see [`Storage::hint`]) so that a storage fetching between attempts +/// gets them in the same round trip as the heap's header. +fn dense_heap( + file_data: &S, + link_info: &LinkInfoMessage, + fh_addr: u64, + offset_size: u8, + length_size: u8, +) -> Result { + if let Some(btree_addr) = link_info.btree_name_index_address { + file_data.hint(btree_addr, BTREE_V2_HEADER_HINT_LEN); + } + let fh = + FractalHeapHeader::parse_in(file_data, checked_addr(fh_addr)?, offset_size, length_size)?; + fh.hint_root_block(file_data); + Ok(fh) +} + /// Visit every link in dense storage (fractal heap + B-tree v2 name index). +/// +/// With `hint_headers` (a listing, whose children are usually opened +/// next), the object header of every hard link is hinted (see +/// [`Storage::hint`]) as soon as the link is read, even after a failure. fn for_each_dense_link( file_data: &S, link_info: &LinkInfoMessage, fh_addr: u64, offset_size: u8, length_size: u8, + hint_headers: bool, mut visit: impl FnMut(LinkMessage), ) -> Result<(), FormatError> { - // Parse fractal heap - let fh = - FractalHeapHeader::parse_in(file_data, checked_addr(fh_addr)?, offset_size, length_size)?; + let fh = dense_heap(file_data, link_info, fh_addr, offset_size, length_size)?; + if hint_headers { + // Every link is read: so is every block of the heap. + fh.hint_managed_blocks(file_data); + } // Parse B-tree v2 for name index let btree_addr = link_info @@ -144,11 +176,27 @@ fn for_each_dense_link( let id_bytes = &record.data[id_offset..id_offset + fh.heap_id_length as usize]; // Read managed object from fractal heap - let link_data = fh.read_managed_object_in(file_data, id_bytes, offset_size); + let link = fh + .read_managed_object_in(file_data, id_bytes, offset_size) + .and_then(|d| parse_link(&d, offset_size)); + if hint_headers + && let Ok(Some(LinkMessage { + link_target: + LinkTarget::Hard { + object_header_address, + }, + .. + })) = &link + { + file_data.hint( + *object_header_address, + crate::object_header::OBJECT_HEADER_HINT_LEN, + ); + } if failed.is_some() { continue; } - match link_data.and_then(|d| parse_link(&d, offset_size)) { + match link { Ok(Some(link)) => visit(link), Ok(None) => {} Err(e) => failed = Some(e), @@ -175,6 +223,7 @@ fn resolve_dense_entries( fh_addr, offset_size, length_size, + true, |link| { if let LinkTarget::Hard { object_header_address, @@ -247,8 +296,7 @@ fn links_named( return Ok(found); }; - let fh = - FractalHeapHeader::parse_in(file_data, checked_addr(fh_addr)?, offset_size, length_size)?; + let fh = dense_heap(file_data, &link_info, fh_addr, offset_size, length_size)?; let btree_addr = link_info .btree_name_index_address .ok_or_else(|| FormatError::PathNotFound(String::from("no B-tree v2 name index")))?; @@ -265,6 +313,7 @@ fn links_named( fh_addr, offset_size, length_size, + false, |link| { if link.name == name { found.push(link); @@ -336,6 +385,27 @@ fn lookup_link( length_size: u8, ) -> Result, FormatError> { if is_v1_group(object_header) { + // Down the group's B-tree, as libhdf5 looks a name up; only when + // that does not find a hard link of that name is every entry read + // (a soft link, a group whose B-tree is out of order). A storage + // error (a read a restartable storage has not fetched yet) is + // returned as is: reading every entry would not get further. + if let Some(sym_msg) = object_header + .messages + .iter() + .find(|m| m.msg_type == MessageType::SymbolTable) + { + let stm = SymbolTableMessage::parse(&sym_msg.data, offset_size)?; + match group_v1::find_v1_entry(file_data, &stm, name, offset_size, length_size) { + Ok(Some(e)) if e.object_header_address != u64::MAX => { + return Ok(Some(LinkTarget::Hard { + object_header_address: e.object_header_address, + })); + } + Err(e @ FormatError::Storage(_)) => return Err(e), + _ => {} + } + } let entries = resolve_group_entries(file_data, object_header, offset_size, length_size)?; if let Some(e) = entries .iter() @@ -409,7 +479,7 @@ fn resolve_child_core( let not_found = || FormatError::PathNotFound(String::from(name)); let header = ObjectHeader::parse_in(file_data, checked_addr(group_address)?, os, ls)?; if !is_v2_group(&header) || is_v1_group(&header) { - return resolve_group_children_in(file_data, superblock, group_address)? + return group_children(file_data, superblock, group_address, false)? .into_iter() .find(|e| e.name == name) .map(|e| e.object_header_address) @@ -575,6 +645,18 @@ fn resolve_group_children_core( file_data: &S, superblock: &Superblock, group_address: u64, +) -> Result, FormatError> { + group_children(file_data, superblock, group_address, true) +} + +/// [`resolve_group_children`]; with `hint_headers`, every child's object +/// header is hinted (see [`Storage::hint`]) as soon as its address is +/// known, for a listing whose children are opened next. +fn group_children( + file_data: &S, + superblock: &Superblock, + group_address: u64, + hint_headers: bool, ) -> Result, FormatError> { let os = superblock.offset_size; let ls = superblock.length_size; @@ -589,7 +671,10 @@ fn resolve_group_children_core( .find(|m| m.msg_type == MessageType::SymbolTable) .ok_or_else(|| FormatError::PathNotFound(String::from("no symbol table message")))?; let stm = SymbolTableMessage::parse(&sym_msg.data, os)?; - let all = group_v1::resolve_v1_group_entries_in(file_data, &stm, os, ls)?; + let all = group_v1::v1_group_entries(file_data, &stm, os, ls, hint_headers)?; + if all.iter().any(|e| e.name.is_empty()) { + return Err(FormatError::InvalidLinkName); + } if all.iter().any(group_v1::is_v1_soft_link) { soft = group_v1::v1_soft_links_in(file_data, &stm, os, ls)?; } @@ -615,7 +700,7 @@ fn resolve_group_children_core( }; let link_info = find_link_info(&header, os)?; if let Some(fh_addr) = link_info.fractal_heap_address { - for_each_dense_link(file_data, &link_info, fh_addr, os, ls, visit)?; + for_each_dense_link(file_data, &link_info, fh_addr, os, ls, hint_headers, visit)?; } else { for msg in &header.messages { if msg.msg_type == MessageType::Link @@ -737,7 +822,7 @@ fn resolve_group_entries( let stm = SymbolTableMessage::parse(&sym_msg.data, offset_size)?; // A lookup: an entry with an empty name (which fails a listing) is // skipped by the name comparison, as in libhdf5. - group_v1::v1_group_entries(file_data, &stm, offset_size, length_size) + group_v1::v1_group_entries(file_data, &stm, offset_size, length_size, false) } else if is_v2_group(object_header) { resolve_v2_group_entries_in(file_data, object_header, offset_size, length_size) } else { @@ -920,6 +1005,86 @@ mod tests { assert_eq!(values, vec![22.5, 23.1, 21.8]); } + /// `v1_groups_400.h5` from its superblock on (it has a user block). + fn v1_groups_400() -> (Vec, Superblock) { + let all: &[u8] = include_bytes!("../tests/fixtures/v1_groups_400.h5"); + let data = all[signature::find_signature(all).unwrap()..].to_vec(); + let sb = Superblock::parse(&data, 0).unwrap(); + (data, sb) + } + + /// Every child of a v1 group resolves by name, down the group's B-tree, + /// to the address the listing gives, reading a small part of what the + /// listing reads; a name the group does not hold is not found. + #[test] + fn v1_lookup_down_the_btree_agrees_with_the_listing() { + let (data, sb) = v1_groups_400(); + let children = resolve_group_children(&data, &sb, sb.root_group_address).unwrap(); + assert_eq!(children.len(), 401); + for c in &children { + let path = format!("/{}", c.name); + assert_eq!( + resolve_path_any(&data, &sb, &path).unwrap(), + c.object_header_address, + "{path}" + ); + } + for missing in ["/g0400", "/a", "/g", "/g00000", "/zz", "/x0"] { + assert!( + matches!( + resolve_path_any(&data, &sb, missing), + Err(FormatError::PathNotFound(_)) + ), + "{missing}" + ); + } + let st = crate::storage::CountingStorage::new(data.clone()); + resolve_group_children_in(&st, &sb, sb.root_group_address).unwrap(); + let listing = st.bytes_read(); + st.reset(); + let last = children.last().unwrap(); + assert_eq!( + resolve_path_any_in(&st, &sb, &format!("/{}", last.name)).unwrap(), + last.object_header_address + ); + assert!( + st.bytes_read() * 8 < listing, + "lookup read {} bytes, listing {listing}", + st.bytes_read() + ); + } + + /// A v1 group whose B-tree is out of name order (a name changed in the + /// heap so that it sorts past every key) is still looked up by reading + /// every entry, as before the lookup went down the B-tree. + #[test] + fn v1_lookup_falls_back_when_the_btree_is_out_of_order() { + let (mut data, sb) = v1_groups_400(); + let at: Vec = data + .windows(6) + .enumerate() + .filter(|(_, w)| *w == b"g0200\0") + .map(|(i, _)| i) + .collect(); + assert_eq!(at.len(), 1, "one heap string"); + data[at[0]] = b'~'; + let children = resolve_group_children(&data, &sb, sb.root_group_address).unwrap(); + let moved = children.iter().find(|c| c.name == "~0200").unwrap(); + assert_eq!( + resolve_path_any(&data, &sb, "/~0200").unwrap(), + moved.object_header_address + ); + assert!(resolve_path_any(&data, &sb, "/g0200").is_err()); + for c in &children { + let path = format!("/{}", c.name); + assert_eq!( + resolve_path_any(&data, &sb, &path).unwrap(), + c.object_header_address, + "{path}" + ); + } + } + #[test] fn path_not_found_v2() { let file_data: &[u8] = include_bytes!("../tests/fixtures/v2_groups.h5"); diff --git a/crates/clawhdf5-format/src/object_header.rs b/crates/clawhdf5-format/src/object_header.rs index 2045a4c..9c9f238 100644 --- a/crates/clawhdf5-format/src/object_header.rs +++ b/crates/clawhdf5-format/src/object_header.rs @@ -163,6 +163,12 @@ impl ObjectHeader { offset_size: u8, length_size: u8, ) -> Result { + // The first chunk is read once the prefix says how long it is: say + // so (see `Storage::hint`), for a storage that fetches between + // attempts. + if file.as_contiguous().is_none() { + file.hint(offset, OBJECT_HEADER_HINT_LEN); + } // The longest prefix of either version, in one read. It holds the // whole prefix or ends at the end of the file, so its bounds checks // are the whole-file ones. @@ -279,10 +285,20 @@ impl ObjectHeader { let mut spans = ChunkSpans::new(file.len(), offset, length)?; let mut chunk0_count = 0usize; let mut next = 0usize; + let hints = file.as_contiguous().is_none(); while let Some((chunk_offset, chunk_length)) = spans.get(next) { let chunk = read_exact_at(file, chunk_offset, chunk_length)?; + let known = spans.len; let count = Self::parse_v1_messages(&chunk, offset_size, length_size, messages, &mut spans)?; + // The continuation chunks this one names are read next. + if hints { + for i in known..spans.len { + if let Some((o, l)) = spans.get(i) { + file.hint(o, l); + } + } + } // Only the first chunk's messages are held to the prefix count. if next == 0 { chunk0_count = count; @@ -483,8 +499,17 @@ impl ObjectHeader { // whenever a message no longer fits), up to the same bound as a // version-1 header. let mut spans = ChunkSpans::new(file.len(), base as u64, chunk0_msg_end.saturating_add(4))?; + // The continuation chunks a chunk names are read next (see + // `Storage::hint`). + let hints = file.as_contiguous().is_none(); + if hints { + for &(o, l) in &continuations { + file.hint(o as u64, l); + } + } while let Some((cont_offset, cont_length)) = continuations.pop() { spans.add(cont_offset as u64, cont_length)?; + let known = continuations.len(); Self::parse_v2_continuation( file, cont_offset as u64, @@ -495,6 +520,11 @@ impl ObjectHeader { &mut messages, &mut continuations, )?; + if hints { + for &(o, l) in &continuations[known..] { + file.hint(o as u64, l); + } + } } Ok(ObjectHeader { @@ -640,6 +670,11 @@ impl ObjectHeader { } } +/// What an object header is hinted to take before its prefix is read (see +/// [`Storage::hint`]): the first chunk of a typical dataset's header. A +/// longer header is read all the same. +pub(crate) const OBJECT_HEADER_HINT_LEN: usize = 512; + /// Longest version-2 object header prefix: signature(4) + version(1) + /// flags(1) + times(16) + attribute phase change(4) + chunk-0 size(8). const V2_PREFIX_MAX: usize = 34; diff --git a/crates/clawhdf5-format/src/storage.rs b/crates/clawhdf5-format/src/storage.rs index 5fe9497..6b74378 100644 --- a/crates/clawhdf5-format/src/storage.rs +++ b/crates/clawhdf5-format/src/storage.rs @@ -89,6 +89,21 @@ pub trait Storage { fn as_contiguous(&self) -> Option<&[u8]> { None } + + /// A hint that `[offset, offset + len)` is about to be read by the + /// same operation: a parser that has just learnt where the next + /// structures are (a node's children, a structure's body) says so + /// before it reads them one at a time. Nothing is read and nothing + /// fails. The default ignores it, as does every backend that reads + /// when asked; the browser's restartable reader, which fetches over + /// the network between attempts, fetches hinted bytes along with the + /// bytes an attempt actually missed, so structures a parser only + /// reaches after a miss arrive in the same round trip. Results never + /// depend on hints. + #[inline] + fn hint(&self, offset: u64, len: usize) { + let _ = (offset, len); + } } impl Storage for [u8] { @@ -148,6 +163,11 @@ impl Storage for &T { fn as_contiguous(&self) -> Option<&[u8]> { (**self).as_contiguous() } + + #[inline] + fn hint(&self, offset: u64, len: usize) { + (**self).hint(offset, len) + } } impl Storage for Box { @@ -170,6 +190,11 @@ impl Storage for Box { fn as_contiguous(&self) -> Option<&[u8]> { (**self).as_contiguous() } + + #[inline] + fn hint(&self, offset: u64, len: usize) { + (**self).hint(offset, len) + } } #[cfg(feature = "std")] @@ -193,6 +218,11 @@ impl Storage for std::sync::Arc { fn as_contiguous(&self) -> Option<&[u8]> { (**self).as_contiguous() } + + #[inline] + fn hint(&self, offset: u64, len: usize) { + (**self).hint(offset, len) + } } /// `storage.len()` as the `usize` the parsers' end-of-file errors report diff --git a/crates/clawhdf5-format/src/symbol_table.rs b/crates/clawhdf5-format/src/symbol_table.rs index 467f43b..4243106 100644 --- a/crates/clawhdf5-format/src/symbol_table.rs +++ b/crates/clawhdf5-format/src/symbol_table.rs @@ -90,6 +90,10 @@ impl SymbolTableNode { offset: u64, offset_size: u8, ) -> Result { + // The entries are read once the header says how many there are: + // say so (see `Storage::hint`), for libhdf5's default node size + // (group leaf K = 4: 8 entries of 40 bytes with 8-byte offsets). + file.hint(offset, 8 + 8 * (2 * usize::from(offset_size) + 24)); // signature(4) + version(1) + reserved(1) + number_of_symbols(2) = 8 let header = read_exact_at(file, offset, 8)?; diff --git a/crates/clawhdf5-wasm/src/lazy.rs b/crates/clawhdf5-wasm/src/lazy.rs index dadb939..e45dce4 100644 --- a/crates/clawhdf5-wasm/src/lazy.rs +++ b/crates/clawhdf5-wasm/src/lazy.rs @@ -25,6 +25,14 @@ //! pass per block it needs. In practice it is one pass per *wave* of //! misses: a chunked read asks for all the chunks of a batch at once. //! +//! Parsers also say what they are about to read ([`Storage::hint`]: a +//! node's children, a structure's body, a listed group's child headers). +//! A pass that misses fetches the hinted blocks it lacks too, as far as the +//! operation's fetch budget allows, so what the parser would only have +//! reached on the next pass arrives in the same round trip; a pass that +//! misses nothing ignores its hints, so they never add a round trip, and +//! results never depend on them. +//! //! Blocks are kept in an LRU cache with a byte budget, trimmed only when no //! operation is in flight. Blocks fetched for bulk reads (raw data: a //! `read_ranges` call, or a read longer than a block) go first, so reading @@ -50,6 +58,14 @@ pub const DEFAULT_MAX_FETCH: u64 = 512 << 20; /// the caller of [`LazyStorage::attempt`]: a pass that missed is re-run. pub const NEED_BYTES: &str = "bytes not fetched yet (restartable read)"; +/// Most blocks one pass records as hinted (see [`Storage::hint`]). +const MAX_HINTED: usize = 1 << 16; + +/// Total length of `ranges`. +fn ranges_len(ranges: &[Range]) -> u64 { + ranges.iter().map(|r| r.end - r.start).sum() +} + /// Settings of a [`LazyStorage`]. #[derive(Debug, Clone, PartialEq, Eq)] pub struct LazyConfig { @@ -94,6 +110,9 @@ pub struct LazyStats { pub bytes_fetched: u64, /// Blocks evicted to stay within the budget. pub evictions: u64, + /// Blocks asked for because a parser hinted it would read them + /// ([`Storage::hint`]), not because a pass missed them. + pub hinted_blocks: u64, /// Bytes cached now. pub cached_bytes: u64, } @@ -125,6 +144,9 @@ struct State { /// Blocks the current pass missed, and whether a small read wanted /// them (metadata). missing: HashMap, + /// Blocks the current pass was told it is about to read and that are + /// not cached ([`Storage::hint`]), in file order. + hinted: BTreeSet, /// Blocks a bulk read missed that have not been supplied yet: kept /// as bulk when they arrive. bulk_pending: BTreeSet, @@ -155,6 +177,19 @@ pub struct Operation<'a> { } impl Operation<'_> { + /// One pass of `f`, as [`LazyStorage::attempt`], but a pass that misses + /// also asks for the blocks it was hinted it would read, as many as fit + /// in what is left of the operation's budget ([`LazyConfig::max_fetch`]) + /// after the blocks it missed. + pub fn attempt(&self, f: impl FnOnce() -> T) -> Step { + let left = self + .storage + .config + .max_fetch + .saturating_sub(self.fetched.get()); + self.storage.attempt_within(f, left) + } + /// Count `ranges` against the operation's budget /// ([`LazyConfig::max_fetch`]) before they are fetched: an error, and /// nothing counted, if they would take it past the budget. @@ -225,20 +260,70 @@ impl LazyStorage { /// Run one pass of `f` over this storage. `Done` when `f` read nothing /// that is missing; otherwise `Need` with the ranges to fetch, and `f`'s /// result is dropped (it may be an error caused by the miss, or a - /// result built around one). + /// result built around one). Hints are not followed; see + /// [`Operation::attempt`]. pub fn attempt(&self, f: impl FnOnce() -> T) -> Step { + self.attempt_within(f, 0) + } + + /// [`attempt`](Self::attempt), adding to a pass that misses the hinted + /// blocks it lacks while everything asked for stays within `budget` + /// bytes. + fn attempt_within(&self, f: impl FnOnce() -> T, budget: u64) -> Step { { let mut st = lock(&self.state); st.missing.clear(); + st.hinted.clear(); st.stats.passes += 1; } let out = f(); - let missing = std::mem::take(&mut lock(&self.state).missing); + let (missing, hinted) = { + let mut st = lock(&self.state); + ( + std::mem::take(&mut st.missing), + std::mem::take(&mut st.hinted), + ) + }; if missing.is_empty() { return Step::Done(out); } drop(out); - Step::Need(self.runs(missing)) + let need = self.runs(&missing); + if hinted.is_empty() || budget == 0 { + return Step::Need(need); + } + // The hinted blocks still missing, in file order, while they fit. + let bs = self.config.block_size; + let mut total = ranges_len(&need); + let mut wanted = missing.clone(); + { + let st = lock(&self.state); + for i in hinted { + if wanted.contains_key(&i) || st.blocks.contains_key(&i) { + continue; + } + let n = (self.len - i * bs).min(bs); + if total + n > budget { + break; + } + total += n; + wanted.insert(i, true); + } + } + if wanted.len() == missing.len() { + return Step::Need(need); + } + let with_hints = self.runs(&wanted); + // Filling holes between runs can add blocks: never let the hints + // take the pass past the budget. + if ranges_len(&with_hints) > budget { + return Step::Need(need); + } + { + let mut st = lock(&self.state); + st.stats.hinted_blocks += (wanted.len() - missing.len()) as u64; + } + Step::Need(with_hints) } /// The bytes of the file at `offset`, fetched for a range a pass asked @@ -309,7 +394,7 @@ impl LazyStorage { ) -> Result { let op = self.operation(); loop { - match self.attempt(&mut f) { + match op.attempt(&mut f) { Step::Done(v) => return Ok(v), Step::Need(ranges) => { op.charge(&ranges)?; @@ -350,14 +435,14 @@ impl LazyStorage { /// blocks, a one-block hole between two runs filled so they merge /// (unless the hole is cached: it would be fetched again), each at /// most `max_request` long. - fn runs(&self, missing: HashMap) -> Vec> { + fn runs(&self, missing: &HashMap) -> Vec> { let bs = self.config.block_size; let mut wanted: Vec = missing.keys().copied().collect(); wanted.sort_unstable(); let mut st = lock(&self.state); // Remember which blocks only bulk reads asked for: they are kept // as bulk once supplied. - for (&i, &metadata) in &missing { + for (&i, &metadata) in missing { if metadata { st.bulk_pending.remove(&i); } else { @@ -499,6 +584,25 @@ impl Storage for LazyStorage { self.len } + fn hint(&self, offset: u64, len: usize) { + // A hint is about a structure, not bulk data: at most a block or + // 1 MiB of it is followed, and at most `MAX_HINTED` blocks a pass, + // whatever a hostile file makes a parser hint. + let len = (len as u64).min(self.config.block_size.max(1 << 20)); + let Some(span) = self.span(offset, len) else { + return; + }; + let mut st = lock(&self.state); + for i in span { + if st.hinted.len() >= MAX_HINTED { + return; + } + if !st.blocks.contains_key(&i) { + st.hinted.insert(i); + } + } + } + fn read_ranges(&self, ranges: &[Range]) -> Result>, FormatError> { let mut spans = Vec::with_capacity(ranges.len()); let mut total = 0u64; @@ -797,6 +901,89 @@ mod tests { } } + #[test] + fn hinted_blocks_come_with_a_miss_and_never_alone() { + let data = file(16 * 1024); + let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20)); + serve(&s, &data, &[3072..4096]); + let op = s.operation(); + // Hints alone: the pass is done, nothing is fetched. + let step = op.attempt(|| { + s.hint(8 * 1024, 100); + owned(s.read_at(3072, 8)) + }); + assert!(matches!(step, Step::Done(Ok(_))), "{step:?}"); + // With a miss, the hinted blocks not cached come too (block 3 is + // cached; 10..=11 is one run, 14 another). + let pass = || { + s.hint(3072, 10); + s.hint(10 * 1024 + 1000, 100); + s.hint(14 * 1024, 1); + owned(s.read_at(0, 8)) + }; + let Step::Need(need) = op.attempt(pass) else { + panic!("block 0 is missing"); + }; + assert_eq!( + need, + vec![0..1024, 10 * 1024..12 * 1024, 14 * 1024..15 * 1024] + ); + // A plain attempt does not follow hints. + let Step::Need(plain) = s.attempt(pass) else { + panic!("block 0 is missing"); + }; + assert_eq!(plain, vec![0..1024]); + serve(&s, &data, &need); + let Step::Done(got) = op.attempt(pass) else { + panic!("everything was supplied"); + }; + assert_eq!(got.unwrap(), &data[..8]); + assert_eq!(s.stats().hinted_blocks, 3); + } + + #[test] + fn hints_stay_within_the_fetch_budget() { + // A budget of 3 blocks: the missed block and the first two hinted + // ones fit, the rest are left out; a hint never makes a call fail. + // (Hinted blocks two apart would be merged with the hole between + // them, which would not fit: then no hint is followed.) + let data = file(64 * 1024); + let mut c = config(1024, 1 << 20); + c.max_fetch = 3 * 1024; + let s = LazyStorage::new(data.len() as u64, c); + let got = s + .run_blocking( + || { + for i in 0..20 { + s.hint(20 * 1024 + i * 3072, 1); + } + owned(s.read_at(0, 8)) + }, + |r| Ok(data[r.start as usize..r.end as usize].to_vec()), + ) + .unwrap() + .unwrap(); + assert_eq!(got, &data[..8]); + let st = s.stats(); + assert_eq!( + (st.bytes_fetched, st.hinted_blocks), + (3 * 1024, 2), + "{st:?}" + ); + // A hint longer than the file, or past its end, is harmless. + let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20)); + s.run_blocking( + || { + s.hint(0, usize::MAX); + s.hint(u64::MAX, 10); + owned(s.read_at(0, 8)) + }, + |r| Ok(data[r.start as usize..r.end as usize].to_vec()), + ) + .unwrap() + .unwrap(); + } + #[test] fn a_failed_fetch_is_an_error_not_data() { let data = file(4096); diff --git a/crates/clawhdf5-wasm/src/lib.rs b/crates/clawhdf5-wasm/src/lib.rs index 6a512a2..85271d2 100644 --- a/crates/clawhdf5-wasm/src/lib.rs +++ b/crates/clawhdf5-wasm/src/lib.rs @@ -391,7 +391,7 @@ impl Http { ) -> Result { let op = storage.operation(); loop { - match storage.attempt(&mut f) { + match op.attempt(&mut f) { Step::Done(v) => return Ok(v), Step::Need(ranges) => { op.charge(&ranges).map_err(js_err)?; diff --git a/crates/clawhdf5-wasm/tests/lazy.rs b/crates/clawhdf5-wasm/tests/lazy.rs index 0d0b536..26e10af 100644 --- a/crates/clawhdf5-wasm/tests/lazy.rs +++ b/crates/clawhdf5-wasm/tests/lazy.rs @@ -146,8 +146,8 @@ fn transcript(api: &impl Api) -> Vec { /// takes), and agrees with the in-memory one: the same values, and an error /// wherever it has one (a malformed file can fail at a different check, /// with a different message, when read by ranges). Returns what the lazy -/// reader fetched and its transcript. -fn check_equal(name: &str, data: &[u8], block: u64) -> (u64, u64, Vec) { +/// reader fetched (requests, bytes, passes) and its transcript. +fn check_equal(name: &str, data: &[u8], block: u64) -> (u64, u64, u64, Vec) { let ctx = format!("{name} (blocks of {block} B)"); let ranged = Reader::open_storage(Arc::new(CountingStorage::new(data.to_vec()))); let local = Reader::open(data.to_vec()); @@ -156,7 +156,7 @@ fn check_equal(name: &str, data: &[u8], block: u64) -> (u64, u64, Vec) { (Ok(r), Ok(l), Ok(z)) => (r, l, z), (Err(r), Err(_), Err(z)) => { assert_eq!(z, r, "{ctx}: open error"); - return (0, 0, Vec::new()); + return (0, 0, 0, Vec::new()); } (r, l, z) => panic!( "{ctx}: opens differently: ranged {:?}, in memory {:?}, lazily {:?}", @@ -184,7 +184,7 @@ fn check_equal(name: &str, data: &[u8], block: u64) -> (u64, u64, Vec) { } assert_eq!(got.len(), local.len(), "{ctx}: transcript length"); let st = lazy.storage.stats(); - (st.requests, st.bytes_fetched, got) + (st.requests, st.bytes_fetched, st.passes, got) } fn config(block: u64) -> LazyConfig { @@ -223,7 +223,7 @@ fn builder_file() -> Vec { fn builder_files_read_the_same_at_every_block_size() { let data = builder_file(); for block in [512, 4096, 1 << 20] { - let (requests, _, lines) = check_equal("builder", &data, block); + let (requests, _, _, lines) = check_equal("builder", &data, block); assert!(requests > 0); // The transcript covers every object, values included. assert!(lines.iter().any(|l| l.starts_with("/grid read: Ok"))); @@ -315,7 +315,7 @@ fn h5py_and_netcdf4_files_read_the_same_lazily() { for name in ["fixture.h5", "fixture.nc"] { let data = std::fs::read(dir.path().join(name)).unwrap(); for block in [512, 64 * 1024] { - let (_, _, lines) = check_equal(name, &data, block); + let (_, _, _, lines) = check_equal(name, &data, block); assert!(lines.iter().filter(|l| l.contains(" read: Ok")).count() >= 2); } } @@ -460,25 +460,32 @@ fn corpus_files_read_the_same_lazily() { } files.sort(); assert!(!files.is_empty(), "no HDF5 files under {dirs}"); - let (mut requests, mut bytes, mut total) = (0u64, 0u64, 0u64); + let (mut requests, mut bytes, mut passes, mut total) = (0u64, 0u64, 0u64, 0u64); for f in &files { let data = std::fs::read(f).unwrap(); total += data.len() as u64; - let (r, b, _) = check_equal(&f.display().to_string(), &data, 64 * 1024); + let (r, b, p, _) = check_equal(&f.display().to_string(), &data, 64 * 1024); requests += r; bytes += b; + passes += p; } eprintln!( - "{} files ({total} bytes): {requests} requests, {bytes} bytes fetched", + "{} files ({total} bytes): {passes} passes, {requests} requests, {bytes} bytes fetched", files.len() ); } -/// Passes and requests `list(path)` takes on a file opened lazily at -/// `block`-byte blocks (the open not counted), checking the listing against -/// the in-memory one. -fn listing_cost(data: &[u8], path: &str, block: u64) -> (u64, u64) { - let want = Reader::open(data.to_vec()).unwrap().list(path).unwrap(); +/// What one call cost on a file opened lazily (the open not counted). +#[derive(Debug, Clone, Copy)] +struct Cost { + passes: u64, + requests: u64, + bytes: u64, +} + +/// Open `data` lazily at `block`-byte blocks, run `op`, and return its +/// result and what it cost (the open not counted). +fn cost_of(data: &[u8], block: u64, op: impl Fn(&Reader) -> T) -> (T, Cost) { let lazy = Lazy::open( data.to_vec(), LazyConfig { @@ -488,21 +495,73 @@ fn listing_cost(data: &[u8], path: &str, block: u64) -> (u64, u64) { ) .unwrap(); let before = lazy.storage.stats(); - assert_eq!(lazy.call(|r| r.list(path)).unwrap(), want); + let out = lazy.call(op); let after = lazy.storage.stats(); ( - after.passes - before.passes, - after.requests - before.requests, + out, + Cost { + passes: after.passes - before.passes, + requests: after.requests - before.requests, + bytes: after.bytes_fetched - before.bytes_fetched, + }, ) } +/// Passes and requests `list(path)` takes on a file opened lazily at +/// `block`-byte blocks (the open not counted), checking the listing against +/// the in-memory one. +fn listing_cost(data: &[u8], path: &str, block: u64) -> (u64, u64) { + let want = Reader::open(data.to_vec()).unwrap().list(path).unwrap(); + let (got, cost) = cost_of(data, block, |r| r.list(path)); + assert_eq!(got.unwrap(), want); + (cost.passes, cost.requests) +} + +/// What reading the dataset at `path` whole takes on a file opened lazily +/// at `block`-byte blocks (the open not counted), checking the values +/// against the in-memory read. +fn read_cost(data: &[u8], path: &str, block: u64) -> Cost { + let want = Reader::open(data.to_vec()) + .unwrap() + .read(path, None) + .unwrap(); + let (got, cost) = cost_of(data, block, |r| r.read(path, None)); + assert_eq!(got.unwrap(), want, "{path}"); + cost +} + +/// An h5py file of `n` datasets of 256 `f32` each (`d0` ... ) in the root +/// group, written with `libver`. +fn h5py_many(dir: &Path, libver: &str, n: usize) -> Vec { + let path = dir.join(format!("{libver}_{n}.h5")); + let script = format!( + "import h5py, numpy as np\n\ + with h5py.File({:?}, 'w', libver='{libver}') as f:\n\ + \x20 for i in range({n}):\n\ + \x20 f.create_dataset('d%d' % i, data=np.full(256, i, np.float32))\n", + path.display().to_string() + ); + let out = Command::new(python()) + .args(["-c", &script]) + .output() + .unwrap(); + assert!( + out.status.success(), + "{}", + String::from_utf8_lossy(&out.stderr) + ); + std::fs::read(&path).unwrap() +} + /// Listing a group reads every child's object header, and its index (B-tree /// and symbol table nodes, or B-tree v2 and heap blocks) before that. Each -/// pass asks for every node of a level it is missing, not the first one -/// only, so the passes (network round trips) grow with the depth of the -/// index, not with the number of children: 2000 children with headers -/// scattered over 512-byte blocks list in a handful of passes, where each -/// header block used to cost its own. +/// pass asks for every node of the index it can reach (a failed node does +/// not stop the walk), the blocks it has been told it reads next (a node's +/// body, a symbol table node's entries, the heap's blocks, each child's +/// header: `Storage::hint`), so the passes (network round trips) follow the +/// depth of the index, not the number of children: 2000 children with +/// headers scattered over 512-byte blocks list in a handful of passes, +/// where each header block used to cost its own. #[test] fn listing_a_large_group_takes_a_few_passes() { let mut b = FileBuilder::new(); @@ -514,39 +573,63 @@ fn listing_a_large_group_takes_a_few_passes() { let data = b.finish().unwrap(); let (passes, requests) = listing_cost(&data, "/many", 512); eprintln!("FileBuilder, 600 children: {passes} passes, {requests} requests"); - assert!(passes <= 6, "{passes} passes"); + // 5 before hints (2026-09-27), 102 before the walks went on past a miss. + assert!(passes <= 4, "{passes} passes"); if !python_available() { + assert!( + !std::env::var("CLAWHDF5_REQUIRE_INTEROP").is_ok_and(|v| v == "1"), + "CLAWHDF5_REQUIRE_INTEROP=1 but {} lacks h5py/netCDF4/numpy", + python() + ); return; } let dir = tempfile::tempdir().unwrap(); - for libver in ["earliest", "latest"] { - let path = dir.path().join(format!("{libver}.h5")); - let script = format!( - "import h5py, numpy as np\n\ - with h5py.File({:?}, 'w', libver='{libver}') as f:\n\ - \x20 for i in range(2000):\n\ - \x20 f.create_dataset('d%d' % i, data=np.full(256, i, np.float32))\n", - path.display().to_string() - ); - let out = Command::new(python()) - .args(["-c", &script]) - .output() - .unwrap(); - assert!( - out.status.success(), - "{}", - String::from_utf8_lossy(&out.stderr) - ); - let data = std::fs::read(&path).unwrap(); + // Most passes each may take: 8 and 11 before hints (2026-09-27). + for (libver, most) in [("earliest", 5), ("latest", 6)] { + let data = h5py_many(dir.path(), libver, 2000); let (passes, requests) = listing_cost(&data, "/", 512); eprintln!("h5py libver={libver}, 2000 children: {passes} passes, {requests} requests"); - assert!(passes <= 12, "{libver}: {passes} passes"); + assert!(passes <= most, "{libver}: {passes} passes"); + } +} + +/// Opening one dataset of a large group looks its name up, not the whole +/// group: in a v1 (symbol table) group down its B-tree as libhdf5 does, in +/// a dense group down its name index. Reading one small dataset of 2000 +/// fetches a few blocks, where it used to read every symbol table node and +/// name of a v1 group (529 requests, 333 kB at 512-byte blocks, before +/// 2026-09-27). +#[test] +fn reading_one_dataset_of_a_large_group_fetches_a_few_blocks() { + if !python_available() { + assert!( + !std::env::var("CLAWHDF5_REQUIRE_INTEROP").is_ok_and(|v| v == "1"), + "CLAWHDF5_REQUIRE_INTEROP=1 but {} lacks h5py/netCDF4/numpy", + python() + ); + eprintln!("skipping: {} lacks h5py/netCDF4/numpy", python()); + return; + } + let dir = tempfile::tempdir().unwrap(); + for (libver, most_passes, most_requests) in [("earliest", 7, 6), ("latest", 8, 9)] { + let data = h5py_many(dir.path(), libver, 2000); + for path in ["/d0", "/d1234", "/d1999"] { + let cost = read_cost(&data, path, 512); + eprintln!("h5py libver={libver}, 2000 children, read {path}: {cost:?}"); + assert!( + cost.passes <= most_passes && cost.requests <= most_requests, + "{libver} {path}: {cost:?}" + ); + assert!(cost.bytes <= 64 * 1024, "{libver} {path}: {cost:?}"); + } } } /// `CLAWHDF5_WASM_LIST_FILE=file.h5`: what listing the root group of that -/// file costs lazily, at 1 MiB and 64 KiB blocks (a measurement, printed). +/// file costs lazily, and opening it and reading one dataset whole +/// (`CLAWHDF5_WASM_READ`, by default the middle dataset of the listing), +/// at 1 MiB and 64 KiB blocks (a measurement, printed). #[test] fn listing_cost_of_a_given_file() { let Ok(path) = std::env::var("CLAWHDF5_WASM_LIST_FILE") else { @@ -563,16 +646,44 @@ fn listing_cost_of_a_given_file() { ) .unwrap(); let open = lazy.storage.stats(); - let n = lazy.call(|r| r.list("/")).unwrap().len(); + let list = lazy.call(|r| r.list("/")).unwrap(); let st = lazy.storage.stats(); eprintln!( - "{path} ({} bytes), {block}-byte blocks: open {} requests / {} passes; list('/') of {n}: {} passes, {} requests, {} bytes", + "{path} ({} bytes), {block}-byte blocks: open {} requests / {} passes; list('/') of {}: {} passes, {} requests, {} bytes ({} blocks hinted)", data.len(), open.requests, open.passes, + list.len(), st.passes - open.passes, st.requests - open.requests, - st.bytes_fetched - open.bytes_fetched + st.bytes_fetched - open.bytes_fetched, + st.hinted_blocks - open.hinted_blocks, + ); + let name = std::env::var("CLAWHDF5_WASM_READ").unwrap_or_else(|_| { + let datasets: Vec<_> = list.iter().filter(|c| c.kind == Kind::Dataset).collect(); + format!("/{}", datasets[datasets.len() / 2].name) + }); + // Open and read on a fresh cache: the open's own cost (the probe + // block and its passes) and then the read's. + let fresh = Lazy::open( + data.clone(), + LazyConfig { + block_size: block, + ..LazyConfig::default() + }, + ) + .unwrap(); + let open = fresh.storage.stats(); + let values = fresh.call(|r| r.read(&name, None)).unwrap().data.len(); + let st = fresh.storage.stats(); + eprintln!( + " open + read('{name}') ({values} values): {} passes, {} requests, {} bytes (the open: {} passes, {} requests, {} bytes)", + st.passes, + st.requests, + st.bytes_fetched, + open.passes, + open.requests, + open.bytes_fetched, ); } } diff --git a/crates/clawhdf5/src/reader.rs b/crates/clawhdf5/src/reader.rs index 4be5f93..7e2ceb5 100644 --- a/crates/clawhdf5/src/reader.rs +++ b/crates/clawhdf5/src/reader.rs @@ -372,6 +372,21 @@ impl Storage for FileData { fn as_contiguous(&self) -> Option<&[u8]> { self.contiguous() } + + fn hint(&self, offset: u64, len: usize) { + if self.contiguous().is_some() { + return; + } + if let Backing::Storage(s) = &self.backing { + // Within the HDF5 data, as a read would be clamped. + let size = Storage::len(self); + let start = offset.min(size); + let len = usize::try_from(size - start).map_or(len, |avail| avail.min(len)); + if len > 0 { + s.hint(self.base + start, len); + } + } + } } /// `bytes`, a backend's answer to a read of `len` bytes, without anything diff --git a/docs/design/range-reads.md b/docs/design/range-reads.md index 4ab9ef5..9dea022 100644 --- a/docs/design/range-reads.md +++ b/docs/design/range-reads.md @@ -607,6 +607,16 @@ fast path within benchmark noise. missing blocks: listing 3000 datasets went from 185 passes to 6. The 32-bit risk below is covered by a Node test that reads data at 3 GiB from a mock server and is refused a 4 GiB file. + - Fewer round trips (2026-09-27, later): the walks descend into every + child after a failure (not only read the siblings), and parsers + call `Storage::hint` for what they read next (node bodies, object + header chunks, a dense group's heap blocks, a listing's child + headers); `LazyStorage` fetches hinted blocks only along with a + pass's real misses and within `maxFetch`, so hints never add a + round trip nor change a result. A v1 group's name is looked up down + its B-tree (`H5G__stab_lookup`), not by listing it. 3000 datasets + list in 4 passes (earliest) and 5 (latest) at 1 MiB blocks, and + opening one of them costs 5 requests, not 74 (CHANGELOG). **M5 — SWMR and growth (later, separate design).** `Storage::len()` may grow; add `File::refresh()` that re-reads the superblock/EOF and invalidates cached diff --git a/docs/known-issues.md b/docs/known-issues.md index b8e47d6..0b25635 100644 --- a/docs/known-issues.md +++ b/docs/known-issues.md @@ -1104,10 +1104,29 @@ cache, but: header, and every node of a level of the group's index, in one pass (since 2026-09-27; it was one round trip per header block): 3000 datasets of an h5py file took 6 passes at 1 MiB blocks, 9 for a - `libver="latest"` file (dense links). Each pass re-parses what the - call reads (CPU, not network). With headers spread through the file - (h5py writes each next to its data) a listing still fetches most of - the file at 1 MiB blocks; a smaller `blockSize` fetches less. + `libver="latest"` file (dense links). **Since 2026-09-27 (later):** + 4 and 5 passes (5 and 6 at 64 KiB, from 8 and 11): the index walks go + on past a missing node, and parsers hint what they read next + (`Storage::hint`: node bodies, the heap's blocks, each child's + header), which the lazy reader fetches with a pass's misses. That is + the depth of the chain (index levels, then symbol table nodes or + heap objects, then headers) plus the pass that finishes; it cannot + go lower without reading structures before their addresses are + known. Opening one dataset of a v1 group looks its name up down the + group's B-tree (it read every entry: 74 requests, 193 MB at 1 MiB + blocks for one 64 KiB dataset of the 3000; now 5 requests, 5 MB). + Each pass re-parses what the call reads (CPU, not network). With + headers spread through the file (h5py writes each next to its data) + a listing still fetches most of the file at 1 MiB blocks (192 of + 198 MB; 35 MB in 530 requests at 64 KiB); a smaller `blockSize` + fetches less. Merging nearby requests does not help such a file: the + blocks a listing needs are five or six apart at 64 KiB, so fewer requests would + mean fetching most of the file. Listing it a second time is free at + 64 KiB blocks, but at 1 MiB its metadata blocks (192 MB) exceed the + 64 MiB `cacheSize`, so they are fetched again (the earliest file: 4 + passes, 50 requests); a larger `cacheSize` keeps them. A file's paged + aggregation (metadata in pages) is not used to fetch its metadata in + one request. - **Memory:** a call keeps every block it reads until it finishes (the cache budget applies between calls). It may fetch at most `maxFetch` bytes (512 MiB by default, at most 1 GiB), and a single read longer diff --git a/examples/wasm-viewer/README.md b/examples/wasm-viewer/README.md index dc9dc05..7eb0849 100644 --- a/examples/wasm-viewer/README.md +++ b/examples/wasm-viewer/README.md @@ -70,8 +70,10 @@ are fetched (in parallel, adjacent blocks in one request), and the pass is run again, until one completes (`docs/design/range-reads.md`, M4). Opening costs one request (the first block, which also gives the file's size); listing a group whose metadata is in blocks already fetched costs none, -and otherwise a round trip per level of the group's index plus one for -its children's headers, all fetched together; +and otherwise about a round trip per level of the group's index plus one +for its children's headers, all fetched together (the reader fetches +what it knows it reads next along with what a pass missed); opening one +object looks its name up in the group's index, not the whole group; reading a chunked dataset costs a round trip for its chunk index (a few for a deep one) and one batch of requests for its chunks. Every answer is checked — a `206` with exactly the bytes asked for, from the same file