format: group walks go on past a failed node and hint what they read next
Listing a large group over openUrl still took 6-11 passes (network round
trips) for the reviewer's 3000-dataset h5py file: each pass only found
the structures the walk reached before its first miss.
- The v1 and v2 B-tree collectors descend into every child of a node
after one fails (they only read the siblings before, so a sibling's
subtree came a pass later), then return the first error: results and
errors unchanged. The v2 walk stops once its record budget is spent,
so a shared-subtree tree still cannot multiply the work.
- Hints (`Storage::hint`, a no-op for every backend but the lazy one):
a group B-tree node's and a symbol table node's body (read once their
header gives a length, a round trip later when the body is in the
next block), an object header's first chunk and its continuation
chunks, the symbol table nodes a B-tree leaf names, a dense group's
name index header and the heap's root block (both read right after
the heap header). A listing also hints every child's object header as
its entry is read, even after a failure, and every direct block of a
dense group's heap (reading the indirect blocks, at most 4096 entries
and 4 levels deep); a lookup does not.
- The fractal heap's indirect-block layout (entry sizes, where the first
n entries end) is one helper used by the object reads and the hints.
Measured with tests/lazy.rs listing_cost_of_a_given_file on an h5py file
like the reviewer's (3000 datasets of 64 KiB, 198 MB), list('/'),
passes/requests/bytes, before -> after:
earliest, 1 MiB: 6/73/192.5 MB -> 4/68/192.5 MB
earliest, 64 KiB: 8/531/35.2 MB -> 5/530/35.3 MB
latest, 1 MiB: 9/98/196.5 MB -> 5/86/196.5 MB
latest, 64 KiB: 11/452/29.6 MB -> 6/454/30.5 MB
listing_a_large_group_takes_a_few_passes (512-byte blocks), budgets
tightened to the new counts: FileBuilder 600 children 5 -> 4 passes,
h5py 2000 children earliest 8 -> 5, latest 11 -> 6.
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -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.
|
||||
@@ -158,6 +160,17 @@ impl BTreeV1Node {
|
||||
/// Maximum recursion depth for B-tree traversal (malformed data protection).
|
||||
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(
|
||||
file_data: &[u8],
|
||||
@@ -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 {
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
@@ -298,7 +314,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, ¤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 {
|
||||
|
||||
@@ -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);
|
||||
@@ -409,7 +458,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 +624,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 +650,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 +679,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 +801,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 {
|
||||
|
||||
@@ -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;
|
||||
@@ -477,8 +493,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,
|
||||
@@ -489,6 +514,11 @@ impl ObjectHeader {
|
||||
&mut messages,
|
||||
&mut continuations,
|
||||
)?;
|
||||
if hints {
|
||||
for &(o, l) in &continuations[known..] {
|
||||
file.hint(o as u64, l);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(ObjectHeader {
|
||||
@@ -634,6 +664,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;
|
||||
|
||||
@@ -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)?;
|
||||
|
||||
|
||||
@@ -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,16 +460,17 @@ 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()
|
||||
);
|
||||
}
|
||||
@@ -496,36 +497,14 @@ fn listing_cost(data: &[u8], path: &str, block: u64) -> (u64, u64) {
|
||||
)
|
||||
}
|
||||
|
||||
/// 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]);
|
||||
}
|
||||
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;
|
||||
}
|
||||
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 +517,54 @@ 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");
|
||||
}
|
||||
}
|
||||
|
||||
/// `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 +581,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,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user