Lazy remote files in the browser (M4), SWMR reader (M5), Python remote reads and editing #19
@@ -199,19 +199,32 @@ fn collect_symbol_table_nodes_inner<S: Storage + ?Sized>(
|
||||
// Leaf: children are SNOD addresses
|
||||
Ok(node.children)
|
||||
} else {
|
||||
// Internal: recurse into children
|
||||
// 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.
|
||||
let mut result = Vec::new();
|
||||
let mut failed = None;
|
||||
for &child_addr in &node.children {
|
||||
let child_snods = collect_symbol_table_nodes_inner(
|
||||
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,
|
||||
offset_size,
|
||||
length_size,
|
||||
depth + 1,
|
||||
)?;
|
||||
result.extend(child_snods);
|
||||
) {
|
||||
Ok(child_snods) => result.extend(child_snods),
|
||||
Err(e) => failed = Some(e),
|
||||
}
|
||||
}
|
||||
match failed {
|
||||
Some(e) => Err(e),
|
||||
None => Ok(result),
|
||||
}
|
||||
Ok(result)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -504,45 +504,63 @@ 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.
|
||||
let mut failed = None;
|
||||
for (i, &(child_addr, child_nrec)) in node.children.iter().enumerate() {
|
||||
if child_depth == 0 {
|
||||
// Before parsing, so a refused tree is not also a large allocation.
|
||||
spend(budget, usize::from(child_nrec))?;
|
||||
let leaf_recs = parse_leaf_records(
|
||||
file,
|
||||
to_usize(child_addr)?,
|
||||
child_nrec,
|
||||
record_size,
|
||||
node_size,
|
||||
)?;
|
||||
out.extend(leaf_recs);
|
||||
} else {
|
||||
collect_internal_records(
|
||||
file,
|
||||
to_usize(child_addr)?,
|
||||
child_nrec,
|
||||
child_depth,
|
||||
record_size,
|
||||
node_size,
|
||||
offset_size,
|
||||
length_size,
|
||||
max_leaf_nrec,
|
||||
budget,
|
||||
out,
|
||||
)?;
|
||||
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 let Err(e) = (|| -> Result<(), FormatError> {
|
||||
if child_depth == 0 {
|
||||
// Before parsing, so a refused tree is not also a large allocation.
|
||||
spend(budget, usize::from(child_nrec))?;
|
||||
let leaf_recs = parse_leaf_records(
|
||||
file,
|
||||
to_usize(child_addr)?,
|
||||
child_nrec,
|
||||
record_size,
|
||||
node_size,
|
||||
)?;
|
||||
out.extend(leaf_recs);
|
||||
} else {
|
||||
collect_internal_records(
|
||||
file,
|
||||
to_usize(child_addr)?,
|
||||
child_nrec,
|
||||
child_depth,
|
||||
record_size,
|
||||
node_size,
|
||||
offset_size,
|
||||
length_size,
|
||||
max_leaf_nrec,
|
||||
budget,
|
||||
out,
|
||||
)?;
|
||||
}
|
||||
|
||||
// Add record[i] (except after the last child)
|
||||
if i < nr {
|
||||
let data = node.record(i, rs)?;
|
||||
spend(budget, 1)?;
|
||||
out.push(BTreeV2Record {
|
||||
data: data.to_vec(),
|
||||
});
|
||||
// Add record[i] (except after the last child)
|
||||
if i < nr {
|
||||
let data = node.record(i, rs)?;
|
||||
spend(budget, 1)?;
|
||||
out.push(BTreeV2Record {
|
||||
data: data.to_vec(),
|
||||
});
|
||||
}
|
||||
Ok(())
|
||||
})() {
|
||||
failed = Some(e);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
match failed {
|
||||
Some(e) => Err(e),
|
||||
None => Ok(()),
|
||||
}
|
||||
}
|
||||
|
||||
/// The records of a B-tree v2 that fall in one key range, found by
|
||||
|
||||
@@ -79,27 +79,51 @@ pub(crate) fn v1_group_entries<S: Storage + ?Sized>(
|
||||
length_size,
|
||||
)?;
|
||||
|
||||
// The names are read one by one from the heap's data segment; read
|
||||
// (up to 1 MiB of) it first, so a storage that records what it lacks
|
||||
// asks for it at once (see `storage::touch`).
|
||||
if !snod_addrs.is_empty() {
|
||||
let len = usize::try_from(heap.data_segment_size).map_or(1 << 20, |n| n.min(1 << 20));
|
||||
crate::storage::touch(file_data, heap.data_segment_address, len);
|
||||
}
|
||||
|
||||
let mut entries = Vec::new();
|
||||
let mut heap_checked = false;
|
||||
// After the first node that fails, the others are only read (as
|
||||
// `storage::touch` does); that error is returned.
|
||||
let mut failed = None;
|
||||
for snod_addr in snod_addrs {
|
||||
let snod = SymbolTableNode::parse_in(file_data, checked_addr(snod_addr)?, offset_size)?;
|
||||
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.
|
||||
if !heap_checked {
|
||||
heap.validate_free_list_in(file_data, length_size)?;
|
||||
heap_checked = true;
|
||||
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)?;
|
||||
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.
|
||||
if !heap_checked {
|
||||
heap.validate_free_list_in(file_data, length_size)?;
|
||||
heap_checked = true;
|
||||
}
|
||||
let name = heap.read_string_in(file_data, entry.link_name_offset)?;
|
||||
entries.push(GroupEntry {
|
||||
name,
|
||||
object_header_address: entry.object_header_address,
|
||||
cache_type: entry.cache_type,
|
||||
});
|
||||
}
|
||||
let name = heap.read_string_in(file_data, entry.link_name_offset)?;
|
||||
entries.push(GroupEntry {
|
||||
name,
|
||||
object_header_address: entry.object_header_address,
|
||||
cache_type: entry.cache_type,
|
||||
});
|
||||
Ok(())
|
||||
};
|
||||
if let Err(e) = node() {
|
||||
failed = Some(e);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(entries)
|
||||
match failed {
|
||||
Some(e) => Err(e),
|
||||
None => Ok(entries),
|
||||
}
|
||||
}
|
||||
|
||||
/// Symbol table cache type for a soft link: the scratch pad's first four bytes
|
||||
|
||||
@@ -126,6 +126,9 @@ fn for_each_dense_link<S: Storage + ?Sized>(
|
||||
)?;
|
||||
let records = collect_btree_v2_records_in(file_data, &btree_hdr, offset_size, length_size)?;
|
||||
|
||||
// After the first link that fails, the others are only read, not
|
||||
// visited (a touch, see `storage::touch`); that error is returned.
|
||||
let mut failed = None;
|
||||
for record in &records {
|
||||
// For type 5 (name index): hash(4) + heap_id(heap_id_length)
|
||||
// For type 6 (creation order): creation_order(8) + heap_id(heap_id_length)
|
||||
@@ -141,12 +144,20 @@ 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)?;
|
||||
if let Some(link) = parse_link(&link_data, offset_size)? {
|
||||
visit(link);
|
||||
let link_data = fh.read_managed_object_in(file_data, id_bytes, offset_size);
|
||||
if failed.is_some() {
|
||||
continue;
|
||||
}
|
||||
match link_data.and_then(|d| parse_link(&d, offset_size)) {
|
||||
Ok(Some(link)) => visit(link),
|
||||
Ok(None) => {}
|
||||
Err(e) => failed = Some(e),
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
match failed {
|
||||
Some(e) => Err(e),
|
||||
None => Ok(()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Resolve entries from dense storage (fractal heap + B-tree v2).
|
||||
|
||||
@@ -202,6 +202,19 @@ pub(crate) fn len_usize<S: Storage + ?Sized>(file: &S) -> usize {
|
||||
usize::try_from(file.len()).unwrap_or(usize::MAX)
|
||||
}
|
||||
|
||||
/// Read `len` bytes at `offset` and drop them, ignoring any error.
|
||||
///
|
||||
/// For a traversal that has failed on one sibling (a B-tree child, a
|
||||
/// symbol table node, a heap object) and would stop there: it first
|
||||
/// touches the siblings it did not get to, so a storage that records what
|
||||
/// it lacks — the browser's restartable reader, which fetches over the
|
||||
/// network between attempts — learns about all of them in one attempt
|
||||
/// instead of one per attempt. Results and errors are unchanged (the first
|
||||
/// error is still the one returned); an in-memory read is free.
|
||||
pub fn touch<S: Storage + ?Sized>(file: &S, offset: u64, len: usize) {
|
||||
let _ = file.read_at(offset, len);
|
||||
}
|
||||
|
||||
/// Bytes `[offset, offset + len)`, all of them.
|
||||
///
|
||||
/// A range that runs past the end of the storage is
|
||||
|
||||
@@ -11,6 +11,8 @@ use std::sync::Arc;
|
||||
use clawhdf5::{AttrValue, File, Selection};
|
||||
use clawhdf5_format::data_read;
|
||||
use clawhdf5_format::datatype::{Datatype, DatatypeByteOrder};
|
||||
use clawhdf5_format::message_type::MessageType;
|
||||
use clawhdf5_format::object_header::ObjectHeader;
|
||||
use clawhdf5_format::storage::Storage;
|
||||
use clawhdf5_format::vl_data::{VlResolver, check_element_size};
|
||||
|
||||
@@ -166,31 +168,54 @@ impl Reader {
|
||||
/// The groups, then the datasets, in the group at `path` (`/` is the
|
||||
/// root). Soft links are listed as their targets; external and dangling
|
||||
/// links, and named datatypes, are left out.
|
||||
///
|
||||
/// What [`Group::groups`](clawhdf5::Group::groups) and `datasets` list,
|
||||
/// but every child's object header is read before an error ends the
|
||||
/// listing (the first error, in listing order, is the one returned, as
|
||||
/// there). Over a [`LazyStorage`](crate::lazy::LazyStorage) that makes
|
||||
/// one pass ask for all the headers it is missing at once, instead of
|
||||
/// one pass, and one round trip, per header.
|
||||
pub fn list(&self, path: &str) -> Result<Vec<Child>> {
|
||||
if self.kind(path)? != Kind::Group {
|
||||
return Err(format!("not a group: {path}"));
|
||||
}
|
||||
let group = self.file.group(path).map_err(err)?;
|
||||
let mut out: Vec<Child> = group
|
||||
.groups()
|
||||
.map_err(err)?
|
||||
.into_iter()
|
||||
.map(|name| Child {
|
||||
name,
|
||||
kind: Kind::Group,
|
||||
})
|
||||
.collect();
|
||||
out.extend(
|
||||
group
|
||||
.datasets()
|
||||
.map_err(err)?
|
||||
.into_iter()
|
||||
.map(|name| Child {
|
||||
name,
|
||||
kind: Kind::Dataset,
|
||||
}),
|
||||
);
|
||||
Ok(out)
|
||||
let entries = group.entries().map_err(err)?;
|
||||
let sb = self.file.superblock();
|
||||
let storage = self.file.storage();
|
||||
let mut groups = Vec::new();
|
||||
let mut datasets = Vec::new();
|
||||
let mut first_error = None;
|
||||
for (name, address) in entries {
|
||||
match ObjectHeader::parse_in(storage, address, sb.offset_size, sb.length_size) {
|
||||
Ok(header) => {
|
||||
let has = |t: MessageType| header.messages.iter().any(|m| m.msg_type == t);
|
||||
if has(MessageType::LinkInfo)
|
||||
|| has(MessageType::Link)
|
||||
|| has(MessageType::SymbolTable)
|
||||
{
|
||||
groups.push(Child {
|
||||
name: name.clone(),
|
||||
kind: Kind::Group,
|
||||
});
|
||||
}
|
||||
if has(MessageType::DataLayout) {
|
||||
datasets.push(Child {
|
||||
name,
|
||||
kind: Kind::Dataset,
|
||||
});
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
first_error.get_or_insert(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(e) = first_error {
|
||||
return Err(err(clawhdf5::Error::from(e)));
|
||||
}
|
||||
groups.extend(datasets);
|
||||
Ok(groups)
|
||||
}
|
||||
|
||||
/// Shape, max shape and datatype of the dataset at `path`.
|
||||
|
||||
@@ -347,32 +347,38 @@ impl LazyStorage {
|
||||
}
|
||||
|
||||
/// Byte ranges covering the missing blocks: runs of consecutive
|
||||
/// blocks, a one-block hole between two runs filled so they merge,
|
||||
/// each at most `max_request` long.
|
||||
/// 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>> {
|
||||
let bs = self.config.block_size;
|
||||
let mut wanted: Vec<u64> = missing.keys().copied().collect();
|
||||
wanted.sort_unstable();
|
||||
{
|
||||
// Remember which blocks only bulk reads asked for: they are
|
||||
// kept as bulk once supplied.
|
||||
let mut st = lock(&self.state);
|
||||
for (&i, &metadata) in &missing {
|
||||
if metadata {
|
||||
st.bulk_pending.remove(&i);
|
||||
} else {
|
||||
st.bulk_pending.insert(i);
|
||||
}
|
||||
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 {
|
||||
if metadata {
|
||||
st.bulk_pending.remove(&i);
|
||||
} else {
|
||||
st.bulk_pending.insert(i);
|
||||
}
|
||||
}
|
||||
let per_request = self.config.max_request / bs;
|
||||
let mut runs: Vec<(u64, u64)> = Vec::new();
|
||||
for i in wanted {
|
||||
match runs.last_mut() {
|
||||
Some((first, last)) if i <= *last + 2 && i - *first < per_request => *last = i,
|
||||
Some((first, last))
|
||||
if (i == *last + 1
|
||||
|| (i == *last + 2 && !st.blocks.contains_key(&(i - 1))))
|
||||
&& i - *first < per_request =>
|
||||
{
|
||||
*last = i
|
||||
}
|
||||
_ => runs.push((i, i)),
|
||||
}
|
||||
}
|
||||
drop(st);
|
||||
runs.into_iter()
|
||||
.map(|(a, b)| a * bs..((b + 1) * bs).min(self.len))
|
||||
.collect()
|
||||
@@ -650,6 +656,21 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_cached_hole_is_not_fetched_again() {
|
||||
let data = file(8 * 1024);
|
||||
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
|
||||
serve(&s, &data, &[1024..2048]);
|
||||
// Blocks 0 and 2 missing, 1 cached: two requests, not 0..3072.
|
||||
let Step::Need(need) = s.attempt(|| {
|
||||
let _ = s.read_at(0, 10);
|
||||
s.read_at(2048, 10).map(|_| ())
|
||||
}) else {
|
||||
panic!("blocks 0 and 2 are missing");
|
||||
};
|
||||
assert_eq!(need, vec![0..1024, 2048..3072]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn supply_refuses_what_was_not_asked_for() {
|
||||
let data = file(10_000);
|
||||
|
||||
@@ -473,3 +473,106 @@ fn corpus_files_read_the_same_lazily() {
|
||||
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();
|
||||
let lazy = Lazy::open(
|
||||
data.to_vec(),
|
||||
LazyConfig {
|
||||
block_size: block,
|
||||
..LazyConfig::default()
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
let before = lazy.storage.stats();
|
||||
assert_eq!(lazy.call(|r| r.list(path)).unwrap(), want);
|
||||
let after = lazy.storage.stats();
|
||||
(
|
||||
after.passes - before.passes,
|
||||
after.requests - before.requests,
|
||||
)
|
||||
}
|
||||
|
||||
/// 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"));
|
||||
let script = format!(
|
||||
"import h5py, numpy as np\n\
|
||||
with h5py.File({:?}, 'w', libver='{libver}') as f:\n\
|
||||
\x20 for i in range(2000):\n\
|
||||
\x20 f.create_dataset('d%d' % i, data=np.full(256, i, np.float32))\n",
|
||||
path.display().to_string()
|
||||
);
|
||||
let out = Command::new(python())
|
||||
.args(["-c", &script])
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"{}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
let data = std::fs::read(&path).unwrap();
|
||||
let (passes, requests) = listing_cost(&data, "/", 512);
|
||||
eprintln!("h5py libver={libver}, 2000 children: {passes} passes, {requests} requests");
|
||||
assert!(passes <= 12, "{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).
|
||||
#[test]
|
||||
fn listing_cost_of_a_given_file() {
|
||||
let Ok(path) = std::env::var("CLAWHDF5_WASM_LIST_FILE") else {
|
||||
return;
|
||||
};
|
||||
let data = std::fs::read(&path).unwrap();
|
||||
for block in [1 << 20, 64 << 10] {
|
||||
let lazy = Lazy::open(
|
||||
data.clone(),
|
||||
LazyConfig {
|
||||
block_size: block,
|
||||
..LazyConfig::default()
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
let open = lazy.storage.stats();
|
||||
let n = lazy.call(|r| r.list("/")).unwrap().len();
|
||||
let st = lazy.storage.stats();
|
||||
eprintln!(
|
||||
"{path} ({} bytes), {block}-byte blocks: open {} requests / {} passes; list('/') of {n}: {} passes, {} requests, {} bytes",
|
||||
data.len(),
|
||||
open.requests,
|
||||
open.passes,
|
||||
st.passes - open.passes,
|
||||
st.requests - open.requests,
|
||||
st.bytes_fetched - open.bytes_fetched
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user