Merge branch 'perf/wasm-listing-passes' into feat/listing-header-last-files

# Conflicts:
#	CHANGELOG.md
This commit is contained in:
osobh
2026-09-27 23:16:27 -05:00
16 changed files with 1039 additions and 134 deletions
+55
View File
@@ -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=<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
+30 -12
View File
@@ -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<S: Storage + ?Sized>(
}
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<S: Storage + ?Sized>(
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 {
+21 -12
View File
@@ -194,11 +194,18 @@ 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)?;
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.
pub fn collect_btree_v2_records(
@@ -504,16 +511,18 @@ fn collect_internal_records<S: Storage + ?Sized>(
// 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<S: Storage + ?Sized>(
}
Ok(())
})() {
failed = Some(e);
failed.get_or_insert(e);
}
}
+166 -32
View File
@@ -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<u8, FormatError> {
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<S: Storage + ?Sized>(&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<S: Storage + ?Sized>(&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<S: Storage + ?Sized>(
&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.
+117 -6
View File
@@ -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<S: Storage + ?Sized>(
offset_size: u8,
length_size: u8,
) -> Result<Vec<GroupEntry>, 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<S: Storage + ?Sized>(
/// 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<S: Storage + ?Sized>(
file_data: &S,
sym_table_msg: &SymbolTableMessage,
offset_size: u8,
length_size: u8,
hint_headers: bool,
) -> Result<Vec<GroupEntry>, FormatError> {
// Parse local heap
let heap = LocalHeap::parse_in(
@@ -93,12 +98,23 @@ pub(crate) fn v1_group_entries<S: Storage + ?Sized>(
// `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<S: Storage + ?Sized>(
}
}
/// 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<S: Storage + ?Sized>(
file_data: &S,
sym_table_msg: &SymbolTableMessage,
name: &str,
offset_size: u8,
length_size: u8,
) -> Result<Option<GroupEntry>, 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<String, FormatError> {
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<S: Storage + ?Sized>(
let mut current_sym_table = root_sym_table.clone();
for (i, component) in components.iter().enumerate() {
let entries = v1_group_entries(file_data, &current_sym_table, offset_size, length_size)?;
let entries = v1_group_entries(
file_data,
&current_sym_table,
offset_size,
length_size,
false,
)?;
let found = entries.iter().find(|e| e.name == *component);
match found {
+176 -11
View File
@@ -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<S: Storage + ?Sized>(
file_data: &S,
link_info: &LinkInfoMessage,
fh_addr: u64,
offset_size: u8,
length_size: u8,
) -> Result<FractalHeapHeader, FormatError> {
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<S: Storage + ?Sized>(
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<S: Storage + ?Sized>(
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<S: Storage + ?Sized>(
fh_addr,
offset_size,
length_size,
true,
|link| {
if let LinkTarget::Hard {
object_header_address,
@@ -247,8 +296,7 @@ fn links_named<S: Storage + ?Sized>(
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<S: Storage + ?Sized>(
fh_addr,
offset_size,
length_size,
false,
|link| {
if link.name == name {
found.push(link);
@@ -336,6 +385,27 @@ fn lookup_link<S: Storage + ?Sized>(
length_size: u8,
) -> Result<Option<LinkTarget>, 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<S: Storage + ?Sized>(
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<S: Storage + ?Sized>(
file_data: &S,
superblock: &Superblock,
group_address: u64,
) -> Result<Vec<GroupEntry>, 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<S: Storage + ?Sized>(
file_data: &S,
superblock: &Superblock,
group_address: u64,
hint_headers: bool,
) -> Result<Vec<GroupEntry>, FormatError> {
let os = superblock.offset_size;
let ls = superblock.length_size;
@@ -589,7 +671,10 @@ fn resolve_group_children_core<S: Storage + ?Sized>(
.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<S: Storage + ?Sized>(
};
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<S: Storage + ?Sized>(
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<u8>, 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<usize> = 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");
@@ -163,6 +163,12 @@ impl ObjectHeader {
offset_size: u8,
length_size: u8,
) -> Result<ObjectHeader, FormatError> {
// 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;
+30
View File
@@ -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<T: Storage + ?Sized> 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<T: Storage + ?Sized> Storage for Box<T> {
@@ -170,6 +190,11 @@ impl<T: Storage + ?Sized> Storage for Box<T> {
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<T: Storage + ?Sized> Storage for std::sync::Arc<T> {
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
@@ -90,6 +90,10 @@ impl SymbolTableNode {
offset: u64,
offset_size: u8,
) -> Result<SymbolTableNode, FormatError> {
// 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)?;
+193 -6
View File
@@ -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>]) -> 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<u64, bool>,
/// Blocks the current pass was told it is about to read and that are
/// not cached ([`Storage::hint`]), in file order.
hinted: BTreeSet<u64>,
/// Blocks a bulk read missed that have not been supplied yet: kept
/// as bulk when they arrive.
bulk_pending: BTreeSet<u64>,
@@ -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<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
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<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
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<T>(&self, f: impl FnOnce() -> T, budget: u64) -> Step<T> {
{
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<T, String> {
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<u64, bool>) -> Vec<Range<u64>> {
fn runs(&self, missing: &HashMap<u64, bool>) -> Vec<Range<u64>> {
let bs = self.config.block_size;
let mut wanted: Vec<u64> = 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<u64>]) -> Result<Vec<Cow<'_, [u8]>>, 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);
+1 -1
View File
@@ -391,7 +391,7 @@ impl Http {
) -> Result<T, JsError> {
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)?;
+158 -47
View File
@@ -146,8 +146,8 @@ fn transcript(api: &impl Api) -> Vec<String> {
/// 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<String>) {
/// reader fetched (requests, bytes, passes) and its transcript.
fn check_equal(name: &str, data: &[u8], block: u64) -> (u64, u64, u64, Vec<String>) {
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<String>) {
(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<String>) {
}
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<u8> {
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<T>(data: &[u8], block: u64, op: impl Fn(&Reader) -> T) -> (T, Cost) {
let lazy = Lazy::open(
data.to_vec(),
LazyConfig {
@@ -488,44 +495,49 @@ 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,
},
)
}
/// 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.
#[test]
fn listing_a_large_group_takes_a_few_passes() {
let mut b = FileBuilder::new();
let mut g = b.create_group("many");
for i in 0..600 {
g.create_dataset(&format!("d{i}")).with_i32_data(&[i; 64]);
/// 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)
}
b.add_group(g.finish());
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");
if !python_available() {
return;
/// 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
}
let dir = tempfile::tempdir().unwrap();
for libver in ["earliest", "latest"] {
let path = dir.path().join(format!("{libver}.h5"));
/// 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<u8> {
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(2000):\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()
);
@@ -538,15 +550,86 @@ fn listing_a_large_group_takes_a_few_passes() {
"{}",
String::from_utf8_lossy(&out.stderr)
);
let data = std::fs::read(&path).unwrap();
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 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();
let mut g = b.create_group("many");
for i in 0..600 {
g.create_dataset(&format!("d{i}")).with_i32_data(&[i; 64]);
}
b.add_group(g.finish());
let data = b.finish().unwrap();
let (passes, requests) = listing_cost(&data, "/many", 512);
eprintln!("FileBuilder, 600 children: {passes} passes, {requests} requests");
// 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();
// 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,
);
}
}
+15
View File
@@ -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
+10
View File
@@ -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
+23 -4
View File
@@ -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
+4 -2
View File
@@ -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