Merge branch 'feat/p3-m4-wasm-lazy' into feat/p3-wasm-swmr-python

# Conflicts:
#	CHANGELOG.md
#	docs/design/range-reads.md
#	docs/known-issues.md
This commit is contained in:
osobh
2026-09-27 07:59:14 -05:00
24 changed files with 3752 additions and 293 deletions
+18 -5
View File
@@ -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)
}
}
+51 -33
View File
@@ -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
+38 -14
View File
@@ -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
+15 -4
View File
@@ -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).
+13
View File
@@ -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
+2
View File
@@ -24,6 +24,8 @@ clawhdf5-format = { path = "../clawhdf5-format", version = "2.7.0" }
# Must match the wasm-bindgen CLI exactly; build.sh checks.
wasm-bindgen = "0.2.129"
js-sys = "0.3.106"
# Promises for openUrl and RemoteFile (pure Rust over js-sys).
wasm-bindgen-futures = "0.4.79"
[dev-dependencies]
serde_json = "1"
+235
View File
@@ -0,0 +1,235 @@
// HTTP for clawhdf5-wasm's openUrl (see src/lib.rs and src/lazy.rs).
//
// The Rust side decides which byte ranges a read needs; this file fetches
// them with `fetch` and `Range` headers and checks every answer, so a server
// that ignores the range, answers with other bytes, or serves a file that
// changed since it was opened is an error, never data. wasm-bindgen copies
// it into the package (pkg/snippets/...).
const DEFAULT_MAX_DOWNLOAD = 512 * 1024 * 1024;
const DEFAULT_PARALLEL = 6;
function fetcher(opts) {
const f = opts?.fetch ?? globalThis.fetch;
if (typeof f !== "function") {
throw new Error("openUrl: no fetch() in this environment (pass opts.fetch)");
}
return f;
}
// The caller's headers (`opts.headers`: a Headers, [name, value] pairs or a
// plain object, as fetch takes them; names come out in lower case) with
// `extra` over them. A Range of the caller's is dropped: this file asks for
// the ranges.
function init(opts, extra, method = "GET", signal = undefined) {
const headers = {};
if (opts?.headers != null) {
for (const [name, value] of new Headers(opts.headers)) {
if (name !== "range") headers[name] = value;
}
}
Object.assign(headers, extra);
return { method, headers, credentials: opts?.credentials, signal };
}
// `opts.parallel`: range requests in flight at once.
function parallelism(opts) {
const p = opts?.parallel ?? DEFAULT_PARALLEL;
if (!Number.isSafeInteger(p) || p < 1) {
throw new Error(`openUrl: parallel must be a positive integer, got ${String(p)}`);
}
return p;
}
// "bytes a-b/total" -> { start, end (exclusive), total | null }; null when
// the page cannot see the header (cross-origin, not exposed).
function contentRange(resp, url) {
const v = resp.headers.get("Content-Range");
if (v === null) return null;
const m = /^bytes (\d+)-(\d+)\/(\d+|\*)$/.exec(v.trim());
if (!m) throw new Error(`${url}: the server sent an unusable Content-Range: ${v}`);
return { start: Number(m[1]), end: Number(m[2]) + 1, total: m[3] === "*" ? null : Number(m[3]) };
}
// What pins the file: its ETag, else its Last-Modified (null if neither is
// visible to this page).
function validatorOf(resp) {
return resp.headers.get("ETag") ?? resp.headers.get("Last-Modified");
}
async function discard(resp) {
try {
await resp.body?.cancel();
} catch {
// Nothing to release.
}
}
// The body, refusing more than `limit` bytes as they arrive: it is piped
// through a TransformStream that errors the moment the count passes the
// limit, which cancels the body and so aborts the request. Whatever the
// server declares or sends, the page never holds more than `limit` bytes of
// it. (`tooBig(n)` makes the error; n is the count so far.) A reader loop
// would do the same, but it stalls on small bodies in headless Chromium
// under --virtual-time-budget, which the page test uses; a pipe does not.
async function readCapped(resp, limit, tooBig) {
const declared = resp.headers.get("Content-Length");
if (declared !== null && Number(declared) > limit) {
await discard(resp);
throw tooBig(declared);
}
if (!resp.body) {
// No stream to read from (some fetch implementations): all at once.
const all = new Uint8Array(await resp.arrayBuffer());
if (all.length > limit) throw tooBig(all.length);
return all;
}
let n = 0;
let over = null;
const capped = resp.body.pipeThrough(new TransformStream({
transform(chunk, ctl) {
n += chunk.length;
if (n > limit) {
over = tooBig(`over ${limit}`);
ctl.error(over);
return;
}
ctl.enqueue(chunk);
},
}));
try {
return new Uint8Array(await new Response(capped).arrayBuffer());
} catch (e) {
throw over ?? e;
}
}
// The whole body of a 200 answer (a server without range support), at
// most `limit` (maxDownload) bytes.
function readAll(resp, limit, url) {
return readCapped(resp, limit, (n) =>
new Error(`${url} is ${n} bytes, more than maxDownload (${limit}); ` +
"the server does not support range requests, so the whole file would have to be downloaded"));
}
// The body of a 206 answer, which may not be longer than the `limit` bytes
// asked for at `start` (the caller checks the exact length).
function readLimited(resp, limit, url, start) {
return readCapped(resp, limit, () =>
new Error(`${url}: asked for ${limit} bytes at offset ${start}, the server sent more`));
}
/**
* Ask for the file's first `firstLen` bytes. A server that honours the
* range (206) gives `{ length, first, validator, requests }`; one that
* answers 200 sends the whole file, which is kept (`{ whole, requests }`)
* when `opts.fallback` is "download" (the default) and the file is at most
* `opts.maxDownload` bytes, and is an error otherwise.
*/
export async function probe(url, firstLen, opts) {
const f = fetcher(opts);
parallelism(opts);
const resp = await f(url, init(opts, { Range: `bytes=0-${firstLen - 1}` }));
if (resp.status === 206) {
const cr = contentRange(resp, url);
if (cr && cr.start !== 0) {
await discard(resp);
throw new Error(`${url}: asked for bytes from 0, the server sent bytes from ${cr.start}`);
}
const first = await readLimited(resp, firstLen, url, 0);
let length = cr?.total ?? null;
let requests = 1;
if (length === null) {
// Content-Range is not readable here: a cross-origin server that does
// not list it in Access-Control-Expose-Headers. Content-Length of a
// HEAD request is always readable.
const head = await f(url, init(opts, {}, "HEAD"));
requests++;
const cl = head.headers.get("Content-Length");
if (!head.ok || cl === null) {
throw new Error(`${url}: cannot learn the file's size (a cross-origin server must send ` +
"Access-Control-Expose-Headers: Content-Range, or answer HEAD with Content-Length)");
}
length = Number(cl);
}
if (!Number.isSafeInteger(length) || length < 0) {
throw new Error(`${url}: the server gave a file size of ${length} bytes; openUrl reads files ` +
"of up to 2^53 - 1 bytes (the largest offset a JavaScript number holds exactly)");
}
if (first.length !== Math.min(firstLen, length)) {
throw new Error(`${url}: asked for the first ${firstLen} bytes of ${length}, got ${first.length}`);
}
return { length, first, validator: validatorOf(resp), requests };
}
if (resp.status === 200) {
if ((opts?.fallback ?? "download") !== "download") {
await discard(resp);
throw new Error(`${url}: the server does not support HTTP range requests (it answered 200 ` +
"to a Range request); open it with { fallback: \"download\" } to download the whole file");
}
const whole = await readAll(resp, opts?.maxDownload ?? DEFAULT_MAX_DOWNLOAD, url);
return { whole, requests: 1 };
}
await discard(resp);
throw new Error(`${url}: HTTP ${resp.status} ${resp.statusText ?? ""}`.trim());
}
/**
* Fetch `ranges` ([start0, end0, start1, end1, ...], ends exclusive) of a
* file opened by `probe`, at most `opts.parallel` (default 6) at a time.
* Every answer must be a 206 with exactly the bytes asked for, from the same
* file (validator and length). When one request fails, the others in
* flight are aborted and no more are made; that failure is the error.
*/
export async function fetchRanges(url, ranges, opts, validator, length) {
const f = fetcher(opts);
const parallel = parallelism(opts);
const n = ranges.length / 2;
const out = new Array(n);
const abort = new AbortController();
let next = 0;
async function one(i) {
const start = ranges[2 * i];
const end = ranges[2 * i + 1];
const resp = await f(url, init(opts, { Range: `bytes=${start}-${end - 1}` }, "GET", abort.signal));
if (resp.status !== 206) {
await discard(resp);
throw new Error(resp.status === 200
? `${url}: the server stopped honouring range requests`
: `${url}: HTTP ${resp.status} ${resp.statusText ?? ""}`.trim());
}
const cr = contentRange(resp, url);
const v = validatorOf(resp);
if ((validator != null && v !== null && v !== validator) ||
(cr?.total != null && cr.total !== length)) {
await discard(resp);
throw new Error(`${url} changed on the server since it was opened`);
}
if (cr && (cr.start !== start || cr.end !== end)) {
await discard(resp);
throw new Error(`${url}: asked for bytes ${start}-${end - 1}, the server sent ${cr.start}-${cr.end - 1}`);
}
const body = await readLimited(resp, end - start, url, start);
if (body.length !== end - start) {
throw new Error(`${url}: asked for ${end - start} bytes at offset ${start}, got ${body.length}`);
}
out[i] = body;
}
async function worker() {
while (next < n && !abort.signal.aborted) {
try {
await one(next++);
} catch (e) {
// The first failure stops the rest: requests in flight are aborted
// (their AbortErrors are not reported) and no new ones start.
if (!abort.signal.aborted) {
abort.abort();
throw e;
}
return;
}
}
}
await Promise.all(Array.from({ length: Math.min(parallel, n) }, worker));
return out;
}
+113 -28
View File
@@ -6,11 +6,25 @@
//! with no typed-array mapping (compound, reference, opaque, ...) is refused
//! with a message naming it, never returned as reinterpreted bytes.
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};
/// The most memory one read may use while it decodes: the stored bytes,
/// the values at 64 bits (integers are widened first) and the values
/// returned. A larger read fails with an error naming `readHyperslab`,
/// before anything is read: on wasm32 a buffer past 2 GiB cannot be
/// allocated at all, and failing to allocate aborts the module (every open
/// file on the page with it). 1 GiB leaves room in wasm32's 4 GiB for the
/// file's cached blocks and the JavaScript copy of the result.
pub const MAX_READ_BYTES: u64 = 1 << 30;
/// Errors are reported to JavaScript as messages.
pub type Result<T> = std::result::Result<T, String>;
@@ -119,7 +133,9 @@ pub struct Attr {
pub value: AttrValue,
}
/// An open file, held in memory.
/// An open file: held in memory ([`Reader::open`]) or read through a
/// [`Storage`] ([`Reader::open_storage`], such as a
/// [`LazyStorage`](crate::lazy::LazyStorage)).
pub struct Reader {
file: File,
}
@@ -132,6 +148,14 @@ impl Reader {
})
}
/// Open a file read through `storage` (the file's bytes from offset 0,
/// user block included, as [`File::open_storage`] takes them).
pub fn open_storage(storage: Arc<dyn Storage + Send + Sync>) -> Result<Self> {
Ok(Self {
file: File::open_storage(storage).map_err(err)?,
})
}
/// Whether `path` names a group or a dataset.
pub fn kind(&self, path: &str) -> Result<Kind> {
match self.file.dataset(path) {
@@ -144,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`.
@@ -228,13 +275,27 @@ impl Reader {
if let Datatype::VariableLength { size, .. } = array_base(&dt) {
check_element_size(*size, self.file.superblock().offset_size).map_err(err)?;
}
let raw = ds.read_selection(&selection).map_err(err)?;
let data = self.decode(&raw, &dt)?;
out_shape.extend(element_shape(&dt));
let expected = out_shape
.iter()
.try_fold(1u64, |acc, &d| acc.checked_mul(d))
.ok_or("selection size overflows")?;
let cost = expected.saturating_mul(bytes_per_value(&dt));
if cost > MAX_READ_BYTES {
return Err(format!(
"reading {path}{} would take about {} MiB of memory, more than the {} MiB \
one read may use; read it in parts (readHyperslab)",
if slab.is_some() {
" (this selection)"
} else {
" whole"
},
cost >> 20,
MAX_READ_BYTES >> 20
));
}
let raw = ds.read_selection(&selection).map_err(err)?;
let data = self.decode(&raw, &dt)?;
if data.len() as u64 != expected {
return Err(format!(
"read {} values for shape {out_shape:?} ({expected} expected)",
@@ -281,11 +342,14 @@ impl Reader {
// string ends at its first NUL and a heap object of the
// wrong size is an error, as in libhdf5 and h5py.
let sb = self.file.superblock();
Data::Strings(
VlResolver::new(self.file.as_bytes(), sb.offset_size, sb.length_size)
.strings(raw)
.map_err(err)?,
)
let strings = match self.file.contiguous_bytes() {
Some(bytes) => {
VlResolver::new(bytes, sb.offset_size, sb.length_size).strings(raw)
}
None => VlResolver::new_in(self.file.storage(), sb.offset_size, sb.length_size)
.strings(raw),
};
Data::Strings(strings.map_err(err)?)
}
Datatype::Enumeration { .. } if !is_array => {
Data::Strings(data_read::read_enum_names(raw, dt).map_err(err)?)
@@ -300,6 +364,27 @@ impl Reader {
}
}
/// Memory one value of type `dt` takes while [`Reader::read`] decodes it
/// (an array type's elements count as values): its stored bytes, plus what
/// [`Reader::decode`] builds from them. A string counts its `String` (24
/// bytes on 64-bit targets, less on wasm32) and, for a fixed-length one,
/// its text; a variable-length string's text lives in the heap and is
/// bounded by the storage's own read limit.
fn bytes_per_value(dt: &Datatype) -> u64 {
let base = array_base(dt);
let stored = u64::from(base.type_size());
stored
+ match base {
Datatype::FloatingPoint { size, .. } if *size <= 4 => 4,
Datatype::FloatingPoint { .. } => 8,
// Widened to 64 bits, then narrowed to a new vector.
Datatype::FixedPoint { .. } => 8 + stored,
Datatype::String { .. } => 24 + stored,
Datatype::VariableLength { .. } | Datatype::Enumeration { .. } => 24,
_ => 0,
}
}
/// Narrow integers read at 64 bits to the dataset's own width. The source is
/// that width, so this cannot fail on correct input; it is checked anyway.
fn narrow<S: Copy + std::fmt::Display, T: TryFrom<S>>(v: Vec<S>) -> Result<Vec<T>> {
+815
View File
@@ -0,0 +1,815 @@
//! Reading a file that is not all here, when no read may wait for the
//! network: the restartable "NeedBytes" mode of `docs/design/range-reads.md`
//! (milestone M4).
//!
//! A browser's main thread cannot block on `fetch`, and the parsers are
//! synchronous. So an operation (open, list a group, read a dataset) runs
//! as a *pass* over a [`LazyStorage`] that holds the blocks fetched so far:
//!
//! 1. [`LazyStorage::attempt`] runs the operation. A read whose blocks are
//! all present is served; a read that misses records the missing blocks
//! and fails with a storage error.
//! 2. If the pass missed anything, its result is thrown away — whatever it
//! is, since a parser may have caught the error and carried on (a
//! listing skips a link it cannot resolve) — and the caller gets the
//! byte ranges to fetch ([`Step::Need`]).
//! 3. The caller fetches them (asynchronously, with HTTP `Range` requests),
//! hands them over with [`LazyStorage::supply`] and runs the operation
//! again.
//!
//! A pass is pure over the storage: the facade only caches what completed
//! reads decoded (its chunk cache), so re-running it is safe. Every pass
//! that does not finish asks for at least one block not yet present, and no
//! block is evicted while an operation is in flight
//! ([`LazyStorage::operation`]), so an operation finishes after at most one
//! 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.
//!
//! 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
//! a large dataset does not evict the metadata.
use std::borrow::Cow;
use std::collections::{BTreeSet, HashMap};
use std::ops::Range;
use std::sync::{Arc, Mutex, MutexGuard};
use clawhdf5_format::error::FormatError;
use clawhdf5_format::storage::Storage;
/// Default block size: 1 MiB, as `clawhdf5-remote`'s block cache (the size
/// `docs/design/range-reads.md` §2 measured).
pub const DEFAULT_BLOCK_SIZE: u64 = 1 << 20;
/// Default of [`LazyConfig::max_fetch`]: 512 MiB, the same as `openUrl`'s
/// `maxDownload` for a server without range support.
pub const DEFAULT_MAX_FETCH: u64 = 512 << 20;
/// The message of the error a read that misses returns. It never reaches
/// the caller of [`LazyStorage::attempt`]: a pass that missed is re-run.
pub const NEED_BYTES: &str = "bytes not fetched yet (restartable read)";
/// Settings of a [`LazyStorage`].
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LazyConfig {
/// Size of a block in bytes (at least 512); fetches are whole, aligned
/// blocks (the file's last block is shorter).
pub block_size: u64,
/// Byte budget of cached blocks between operations. An operation keeps
/// every block it needs until it finishes, whatever the budget.
pub capacity: u64,
/// Largest single range asked for, in bytes (whole blocks, at least
/// one); longer runs are split so they can be fetched in parallel.
pub max_request: u64,
/// Most bytes one operation may fetch (at least one block), and so the
/// longest single read: a read longer than this fails at once, before
/// anything is fetched, and so does an operation whose passes would
/// fetch more. The file's length comes from the server, so without
/// this a hostile file (a heap "collection" claiming 2 GiB) makes the
/// reader fetch and hold whatever it names; on wasm32 a buffer past
/// 2 GiB cannot even be allocated.
pub max_fetch: u64,
}
impl Default for LazyConfig {
fn default() -> Self {
LazyConfig {
block_size: DEFAULT_BLOCK_SIZE,
capacity: 64 << 20,
max_request: 8 << 20,
max_fetch: DEFAULT_MAX_FETCH,
}
}
}
/// What a [`LazyStorage`] has done so far.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct LazyStats {
/// Passes run by [`LazyStorage::attempt`].
pub passes: u64,
/// Ranges handed to [`LazyStorage::supply`]: one HTTP request each.
pub requests: u64,
/// Bytes handed to [`LazyStorage::supply`].
pub bytes_fetched: u64,
/// Blocks evicted to stay within the budget.
pub evictions: u64,
/// Bytes cached now.
pub cached_bytes: u64,
}
/// The outcome of one pass.
#[derive(Debug)]
pub enum Step<T> {
/// The pass read only bytes that were present: its result stands.
Done(T),
/// The pass missed: fetch these byte ranges (sorted, disjoint, block
/// aligned), [`supply`](LazyStorage::supply) them and run it again.
Need(Vec<Range<u64>>),
}
struct Block {
data: Arc<[u8]>,
/// Eviction order: bulk blocks (`false`) before metadata (`true`),
/// then least recently used first.
key: (bool, u64),
}
#[derive(Default)]
struct State {
blocks: HashMap<u64, Block>,
/// `(metadata?, tick, block index)`, in eviction order.
order: BTreeSet<(bool, u64, u64)>,
tick: u64,
bytes: u64,
/// Blocks the current pass missed, and whether a small read wanted
/// them (metadata).
missing: HashMap<u64, bool>,
/// Blocks a bulk read missed that have not been supplied yet: kept
/// as bulk when they arrive.
bulk_pending: BTreeSet<u64>,
/// Operations in flight: no eviction while any is.
active: u32,
stats: LazyStats,
}
/// A [`Storage`] over the blocks of a file fetched so far; a read of
/// anything else fails and is recorded, so the pass can be re-run once the
/// bytes arrive. See the [module documentation](self).
pub struct LazyStorage {
len: u64,
config: LazyConfig,
state: Mutex<State>,
}
fn lock(m: &Mutex<State>) -> MutexGuard<'_, State> {
m.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
}
/// Keeps an operation's blocks cached until it is dropped; see
/// [`LazyStorage::operation`].
pub struct Operation<'a> {
storage: &'a LazyStorage,
/// Bytes fetched for this operation so far.
fetched: std::cell::Cell<u64>,
}
impl Operation<'_> {
/// 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.
pub fn charge(&self, ranges: &[Range<u64>]) -> Result<(), String> {
let max = self.storage.config.max_fetch;
let total = ranges.iter().fold(self.fetched.get(), |n, r| {
n.saturating_add(r.end.saturating_sub(r.start))
});
if total > max {
return Err(format!(
"this call would fetch more than {max} bytes of the file (the maxFetch limit); \
read less at a time (readHyperslab) or raise maxFetch"
));
}
self.fetched.set(total);
Ok(())
}
}
impl Drop for Operation<'_> {
fn drop(&mut self) {
let mut st = lock(&self.storage.state);
st.active = st.active.saturating_sub(1);
if st.active == 0 {
self.storage.evict(&mut st);
}
}
}
impl LazyStorage {
/// An empty cache for a file of `len` bytes.
pub fn new(len: u64, mut config: LazyConfig) -> Self {
config.block_size = config.block_size.max(512);
config.max_request = (config.max_request / config.block_size).max(1) * config.block_size;
config.max_fetch = config.max_fetch.max(config.block_size);
LazyStorage {
len,
config,
state: Mutex::new(State::default()),
}
}
/// The settings in use (after rounding).
pub fn config(&self) -> &LazyConfig {
&self.config
}
/// Counters since the storage was made.
pub fn stats(&self) -> LazyStats {
let st = lock(&self.state);
LazyStats {
cached_bytes: st.bytes,
..st.stats
}
}
/// Mark an operation in flight until the guard is dropped: no block is
/// evicted meanwhile, so re-running its passes always makes progress.
/// Hold it across every pass of one operation.
pub fn operation(&self) -> Operation<'_> {
lock(&self.state).active += 1;
Operation {
storage: self,
fetched: std::cell::Cell::new(0),
}
}
/// 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).
pub fn attempt<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
{
let mut st = lock(&self.state);
st.missing.clear();
st.stats.passes += 1;
}
let out = f();
let missing = std::mem::take(&mut lock(&self.state).missing);
if missing.is_empty() {
return Step::Done(out);
}
drop(out);
Step::Need(self.runs(missing))
}
/// The bytes of the file at `offset`, fetched for a range a pass asked
/// for. `offset` must be block aligned and the bytes whole blocks (the
/// last block of the file may be short) inside the file, or this is an
/// error and nothing is kept. Blocks already present are left alone.
pub fn supply(&self, offset: u64, bytes: &[u8]) -> Result<(), String> {
let bs = self.config.block_size;
let end = offset
.checked_add(bytes.len() as u64)
.filter(|&e| e <= self.len)
.ok_or_else(|| {
format!(
"{} bytes at offset {offset} run past the end of the {}-byte file",
bytes.len(),
self.len
)
})?;
if !offset.is_multiple_of(bs) || (!end.is_multiple_of(bs) && end != self.len) {
return Err(format!(
"{} bytes at offset {offset} are not whole {bs}-byte blocks",
bytes.len()
));
}
let mut st = lock(&self.state);
st.stats.requests += 1;
st.stats.bytes_fetched += bytes.len() as u64;
let mut start = offset;
while start < end {
let i = start / bs;
let stop = (start + bs).min(end);
if !st.blocks.contains_key(&i) {
let rel = (start - offset) as usize..(stop - offset) as usize;
let metadata = !st.bulk_pending.remove(&i);
self.keep(&mut st, i, Arc::from(&bytes[rel]), metadata);
}
start = stop;
}
if st.active == 0 {
self.evict(&mut st);
}
Ok(())
}
/// [`supply`](Self::supply) the bytes fetched for `range`, one of the
/// ranges a [`Step::Need`] asked for: anything but exactly its length
/// (a server that answered with more or less) is an error.
pub fn supply_range(&self, range: &Range<u64>, bytes: &[u8]) -> Result<(), String> {
let want = range.end.saturating_sub(range.start);
if bytes.len() as u64 != want {
return Err(format!(
"asked for {want} bytes at offset {}, got {}",
range.start,
bytes.len()
));
}
self.supply(range.start, bytes)
}
/// Run `f` to completion, fetching what its passes miss with `fetch`
/// (a byte range to its bytes). The blocking driver, for native code
/// and tests; the browser's is the same loop with an `await` between
/// passes.
pub fn run_blocking<T>(
&self,
mut f: impl FnMut() -> T,
mut fetch: impl FnMut(Range<u64>) -> Result<Vec<u8>, String>,
) -> Result<T, String> {
let op = self.operation();
loop {
match self.attempt(&mut f) {
Step::Done(v) => return Ok(v),
Step::Need(ranges) => {
op.charge(&ranges)?;
for r in ranges {
let bytes = fetch(r.clone())?;
self.supply_range(&r, &bytes)?;
}
}
}
}
}
/// Cache block `i`.
fn keep(&self, st: &mut State, i: u64, data: Arc<[u8]>, metadata: bool) {
st.tick += 1;
let key = (metadata, st.tick);
st.bytes += <[u8]>::len(&data) as u64;
st.order.insert((key.0, key.1, i));
if let Some(old) = st.blocks.insert(i, Block { data, key }) {
st.order.remove(&(old.key.0, old.key.1, i));
st.bytes -= <[u8]>::len(&old.data) as u64;
}
}
fn evict(&self, st: &mut State) {
while st.bytes > self.config.capacity {
let Some((_, _, i)) = st.order.pop_first() else {
break;
};
if let Some(b) = st.blocks.remove(&i) {
st.bytes -= <[u8]>::len(&b.data) as u64;
st.stats.evictions += 1;
}
}
}
/// Byte ranges covering the missing blocks: runs of consecutive
/// 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();
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 + 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()
}
/// Block indices covering `[offset, offset + len)`, clamped to the file.
fn span(&self, offset: u64, len: u64) -> Option<Range<u64>> {
let end = offset.saturating_add(len).min(self.len);
if offset >= end {
return None;
}
let bs = self.config.block_size;
Some(offset / bs..(end - 1) / bs + 1)
}
/// The blocks of `spans` if all are present (touching them), else
/// record the missing ones and fail.
fn blocks(
&self,
spans: &[Range<u64>],
metadata: bool,
) -> Result<HashMap<u64, Arc<[u8]>>, FormatError> {
let mut st = lock(&self.state);
let mut have = HashMap::new();
let mut missed = false;
for span in spans {
for i in span.clone() {
if have.contains_key(&i) {
continue;
}
match st.blocks.get(&i) {
Some(b) => {
have.insert(i, b.data.clone());
}
None => {
missed = true;
let m = st.missing.entry(i).or_insert(metadata);
*m |= metadata;
}
}
}
}
if missed {
return Err(FormatError::Storage(NEED_BYTES.into()));
}
// Touch: most recently used last; a small read promotes a bulk
// block to metadata.
for &i in have.keys() {
st.tick += 1;
let tick = st.tick;
let Some(b) = st.blocks.get_mut(&i) else {
continue;
};
let old = b.key;
b.key = (old.0 || metadata, tick);
let new = b.key;
st.order.remove(&(old.0, old.1, i));
st.order.insert((new.0, new.1, i));
}
Ok(have)
}
/// Refuse a read of `n` bytes longer than an operation may fetch
/// ([`LazyConfig::max_fetch`]), before its blocks are asked for.
fn check_len(&self, n: u64) -> Result<(), FormatError> {
let max = self.config.max_fetch;
if n > max {
return Err(FormatError::Storage(format!(
"a read of {n} bytes is more than one call may fetch ({max} bytes, the maxFetch limit)"
)));
}
Ok(())
}
/// The bytes `offset..end` from `blocks`, which hold every block of
/// that span. The buffer is reserved fallibly: a length the address
/// space cannot hold (past `isize::MAX` on wasm32) is an error, never
/// an abort.
fn assemble(
&self,
offset: u64,
end: u64,
blocks: &HashMap<u64, Arc<[u8]>>,
) -> Result<Vec<u8>, FormatError> {
let bs = self.config.block_size;
let n = end - offset;
let too_long =
|| FormatError::Storage(format!("cannot hold a read of {n} bytes in memory"));
let mut out = Vec::new();
out.try_reserve_exact(usize::try_from(n).map_err(|_| too_long())?)
.map_err(|_| too_long())?;
let mut pos = offset;
while pos < end {
let i = pos / bs;
let block = &blocks[&i];
let from = (pos - i * bs) as usize;
let to = ((end - i * bs) as usize).min(<[u8]>::len(block));
out.extend_from_slice(&block[from..to]);
pos = i * bs + to as u64;
}
Ok(out)
}
}
impl Storage for LazyStorage {
fn read_at(&self, offset: u64, len: usize) -> Result<Cow<'_, [u8]>, FormatError> {
let Some(span) = self.span(offset, len as u64) else {
return Ok(Cow::Owned(Vec::new()));
};
let end = offset.saturating_add(len as u64).min(self.len);
self.check_len(end - offset)?;
let metadata = len as u64 <= self.config.block_size;
let blocks = self.blocks(std::slice::from_ref(&span), metadata)?;
Ok(Cow::Owned(self.assemble(offset, end, &blocks)?))
}
fn len(&self) -> u64 {
self.len
}
fn read_ranges(&self, ranges: &[Range<u64>]) -> Result<Vec<Cow<'_, [u8]>>, FormatError> {
let mut spans = Vec::with_capacity(ranges.len());
let mut total = 0u64;
for r in ranges {
if r.end < r.start {
return Err(FormatError::Storage(
"read range ends before it starts".into(),
));
}
total = total.saturating_add(r.end.min(self.len).saturating_sub(r.start));
spans.extend(self.span(r.start, r.end - r.start));
}
self.check_len(total)?;
let blocks = self.blocks(&spans, false)?;
ranges
.iter()
.map(|r| {
let end = r.end.min(self.len);
if r.start >= end {
Ok(Cow::Owned(Vec::new()))
} else {
self.assemble(r.start, end, &blocks).map(Cow::Owned)
}
})
.collect()
}
}
#[cfg(test)]
#[allow(clippy::single_range_in_vec_init)]
mod tests {
use super::*;
fn file(n: usize) -> Vec<u8> {
(0..n).map(|i| (i * 7 + i / 251) as u8).collect()
}
fn config(block: u64, capacity: u64) -> LazyConfig {
LazyConfig {
block_size: block,
capacity,
max_request: 4 * block,
max_fetch: DEFAULT_MAX_FETCH,
}
}
/// Supply every range of `need` from `data`.
fn serve(s: &LazyStorage, data: &[u8], need: &[Range<u64>]) {
for r in need {
s.supply_range(r, &data[r.start as usize..r.end as usize])
.unwrap();
}
}
fn owned(r: Result<Cow<'_, [u8]>, FormatError>) -> Result<Vec<u8>, FormatError> {
r.map(Cow::into_owned)
}
#[test]
fn a_miss_asks_for_whole_blocks_then_the_rerun_reads_them() {
let data = file(10_000);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let Step::Need(need) = s.attempt(|| owned(s.read_at(1500, 1000))) else {
panic!("nothing is cached yet");
};
assert_eq!(need, vec![1024..3072]);
serve(&s, &data, &need);
let Step::Done(got) = s.attempt(|| owned(s.read_at(1500, 1000))) else {
panic!("the blocks were supplied");
};
assert_eq!(got.unwrap(), &data[1500..2500]);
// Past the end: short, then empty, as for a slice.
serve(&s, &data, &[9216..10_000]);
let Step::Done(tail) = s.attempt(|| owned(s.read_at(9_990, 100))) else {
panic!("the last block was supplied");
};
assert_eq!(tail.unwrap(), &data[9_990..]);
assert!(matches!(
s.attempt(|| s.read_at(20_000, 10).map(|c| c.len())),
Step::Done(Ok(0))
));
let st = s.stats();
assert_eq!(
(st.passes, st.requests, st.bytes_fetched),
(4, 2, 2048 + 784)
);
}
#[test]
fn a_pass_that_swallowed_the_miss_is_still_rerun() {
// A parser that catches the error and returns something anyway
// (a listing skipping a link it cannot resolve) must not have its
// result used.
let data = file(4096);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let step = s.attempt(|| s.read_at(0, 4).map(|b| b.len()).unwrap_or(0));
assert!(
matches!(step, Step::Need(ref n) if n == &vec![0..1024]),
"{step:?}"
);
}
#[test]
fn read_ranges_asks_for_every_missing_block_in_one_pass() {
let data = file(64 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let ranges = [
100..200,
5000..5100,
5200..5300,
30_000..33_000,
60_000..60_010,
];
let read = || {
s.read_ranges(&ranges)
.map(|v| v.into_iter().map(Cow::into_owned).collect::<Vec<_>>())
};
let Step::Need(need) = s.attempt(read) else {
panic!("nothing is cached yet");
};
// 5000..5300 is blocks 4 and 5; 30_000..33_000 is blocks 29..=32.
assert_eq!(
need,
vec![
0..1024,
4096..6144,
29 * 1024..33 * 1024,
58 * 1024..59 * 1024
]
);
serve(&s, &data, &need);
let Step::Done(got) = s.attempt(read) else {
panic!("one pass fetched everything");
};
for (r, g) in ranges.iter().zip(&got.unwrap()) {
assert_eq!(g, &data[r.start as usize..r.end as usize]);
}
}
#[test]
fn runs_merge_one_block_holes_and_split_long_runs() {
let data = file(32 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
// Blocks 0 and 2 (a hole of one: merged), 5..=14 (split in fours).
let Step::Need(need) = s.attempt(|| {
let _ = s.read_at(0, 10);
let _ = s.read_at(2048, 10);
s.read_ranges(&[5120..15 * 1024]).map(|_| ())
}) else {
panic!("nothing is cached yet");
};
assert_eq!(
need,
vec![0..3072, 5120..9216, 9216..13_312, 13_312..15_360]
);
}
#[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);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
assert!(s.supply(1, &data[1..1025]).unwrap_err().contains("whole"));
assert!(s.supply(0, &data[..1000]).unwrap_err().contains("whole"));
assert!(
s.supply(9216, &[0u8; 1024])
.unwrap_err()
.contains("past the end")
);
for wrong in [&data[..1023], &data[..2048]] {
assert!(
s.supply_range(&(0..1024), wrong)
.unwrap_err()
.contains("asked for 1024 bytes")
);
}
assert_eq!(s.stats().cached_bytes, 0);
// The file's last, short block is whole.
s.supply(9216, &data[9216..]).unwrap();
assert_eq!(s.stats().cached_bytes, 784);
}
#[test]
fn an_operation_keeps_its_blocks_whatever_the_budget() {
// A budget of one block, an operation that needs eight: without
// the operation guard every pass would evict what the last one
// fetched and never finish.
let data = file(8 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 1024));
let got = s
.run_blocking(
|| {
(0..8)
.map(|i| owned(s.read_at(i * 1024 + 3, 10)))
.collect::<Result<Vec<_>, _>>()
},
|r| Ok(data[r.start as usize..r.end as usize].to_vec()),
)
.unwrap()
.unwrap();
for (i, g) in got.iter().enumerate() {
assert_eq!(g, &data[i * 1024 + 3..i * 1024 + 13]);
}
// Trimmed to the budget once the operation is over.
let st = s.stats();
assert_eq!(st.cached_bytes, 1024);
assert_eq!(st.evictions, 7);
assert_eq!(st.passes, 9, "one pass per block, then the one that ends");
}
#[test]
fn bulk_blocks_are_evicted_before_metadata() {
let data = file(16 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 4 * 1024));
let fetch = |r: Range<u64>| Ok(data[r.start as usize..r.end as usize].to_vec());
// Metadata: a small read of block 0.
s.run_blocking(|| s.read_at(0, 16).map(|_| ()), fetch)
.unwrap()
.unwrap();
// Bulk: raw data over blocks 4..12, more than the budget.
s.run_blocking(|| s.read_ranges(&[4096..12 * 1024]).map(|_| ()), fetch)
.unwrap()
.unwrap();
// The metadata block survived: reading it again is a hit.
let before = s.stats();
assert!(before.cached_bytes <= 4 * 1024);
assert!(matches!(
s.attempt(|| s.read_at(0, 16).map(|_| ())),
Step::Done(Ok(()))
));
assert_eq!(s.stats().requests, before.requests);
}
#[test]
fn a_read_longer_than_max_fetch_fails_without_fetching() {
// A hostile file names a 2 GiB heap collection in a "file" the
// server claims is 1 TiB: the read is refused before any block is
// asked for (on wasm32 its buffer could not even be allocated).
let s = LazyStorage::new(1 << 40, LazyConfig::default());
let step = s.attempt(|| s.read_at(4096, (1usize << 31) + 4096).map(|b| b.len()));
match step {
Step::Done(Err(e)) => assert!(e.to_string().contains("maxFetch"), "{e}"),
other => panic!("expected a refusal, got {other:?}"),
}
let step = s.attempt(|| s.read_ranges(&[0..(600 << 20)]).map(|v| v.len()));
assert!(matches!(step, Step::Done(Err(_))), "{step:?}");
assert_eq!(s.stats().requests, 0);
// At the limit it is an ordinary miss.
let s = LazyStorage::new(1 << 40, config(1024, 1 << 20));
let step = s.attempt(|| s.read_at(0, DEFAULT_MAX_FETCH as usize).map(|b| b.len()));
assert!(matches!(step, Step::Need(_)), "{step:?}");
}
#[test]
fn an_operation_stops_at_its_fetch_budget() {
// Many small reads, none over the limit, that together would fetch
// more than the budget: the operation fails before fetching past it.
let data = file(64 * 1024);
let mut c = config(1024, 1 << 20);
c.max_fetch = 8 * 1024;
let s = LazyStorage::new(data.len() as u64, c);
let e = s
.run_blocking(
|| {
(0..64)
.map(|i| owned(s.read_at(i * 1024, 8)))
.collect::<Result<Vec<_>, _>>()
},
|r| Ok(data[r.start as usize..r.end as usize].to_vec()),
)
.unwrap_err();
assert!(e.contains("maxFetch"), "{e}");
assert!(s.stats().bytes_fetched <= 8 * 1024, "{:?}", s.stats());
// Within the budget it completes, and the budget is per operation.
for _ in 0..3 {
s.run_blocking(
|| owned(s.read_at(10 * 1024, 3000)),
|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);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let e = s
.run_blocking(|| owned(s.read_at(0, 8)), |_| Err("HTTP 500".into()))
.unwrap_err();
assert_eq!(e, "HTTP 500");
let e = s
.run_blocking(|| owned(s.read_at(0, 8)), |_| Ok(vec![0; 10]))
.unwrap_err();
assert!(e.contains("asked for 1024 bytes"), "{e}");
assert_eq!(s.stats().cached_bytes, 0);
drop(data);
}
}
+476 -62
View File
@@ -1,7 +1,7 @@
//! clawhdf5's HDF5 reader for JavaScript, via `wasm-bindgen`.
//!
//! ```js
//! import init, { open } from "./pkg/clawhdf5_wasm.js";
//! import init, { open, openUrl } from "./pkg/clawhdf5_wasm.js";
//! await init();
//! const file = open(new Uint8Array(await blob.arrayBuffer()));
//! file.list("/"); // [{ name, kind: "group" | "dataset" }]
@@ -10,6 +10,12 @@
//! file.read("/x"); // { shape, dtype, data: Float64Array | ... | string[] }
//! file.readHyperslab("/x", [0, 0], [10, 10]); // stride, block optional
//! file.free();
//!
//! // A file on a web server, read by HTTP range requests as needed: the
//! // same methods, returning promises.
//! const remote = await openUrl("https://example.org/data.h5");
//! await remote.list("/");
//! remote.stats(); // { requests, bytesFetched, size, ... }
//! ```
//!
//! Numeric data comes back in the typed array of the stored width
@@ -18,15 +24,23 @@
//! Anything else is a thrown `Error` naming the datatype. Only the reader is
//! exposed: nothing here writes files.
//!
//! The logic lives in [`core`], which is plain Rust and tested natively.
//! The logic lives in [`core`] and [`lazy`], which are plain Rust and tested
//! natively; `js/remote.js` does the HTTP.
pub mod core;
pub mod lazy;
use std::ops::Range;
use std::rc::Rc;
use std::sync::Arc;
use clawhdf5::AttrValue;
use js_sys::{Array, Object, Reflect};
use js_sys::{Array, Object, Promise, Reflect, Uint8Array};
use wasm_bindgen::prelude::*;
use wasm_bindgen_futures::future_to_promise;
use crate::core::{Data, Hyperslab, Reader};
use crate::core::{Attr, Child, Data, DatasetInfo, Hyperslab, Reader};
use crate::lazy::{LazyConfig, LazyStorage, Step};
/// JavaScript numbers are exact up to 2^53.
const MAX_SAFE_INTEGER: f64 = 9_007_199_254_740_991.0;
@@ -35,11 +49,30 @@ fn js_err(msg: String) -> JsError {
JsError::new(&msg)
}
/// A JavaScript exception (from `fetch`, or `js/remote.js`) as a `JsError`
/// with its message.
fn js_exception(e: JsValue) -> JsError {
let msg = e
.dyn_ref::<js_sys::Error>()
.map(|e| String::from(e.message()))
.or_else(|| e.as_string())
.unwrap_or_else(|| format!("{e:?}"));
JsError::new(&msg)
}
fn set(obj: &Object, key: &str, value: impl Into<JsValue>) {
// Defining a property on a fresh plain object cannot fail.
Reflect::set(obj, &JsValue::from_str(key), &value.into()).unwrap_throw();
}
fn get(obj: &JsValue, key: &str) -> JsValue {
if obj.is_object() {
Reflect::get(obj, &JsValue::from_str(key)).unwrap_or(JsValue::UNDEFINED)
} else {
JsValue::UNDEFINED
}
}
fn shape_to_js(shape: &[u64]) -> Array {
shape.iter().map(|&d| JsValue::from_f64(d as f64)).collect()
}
@@ -58,6 +91,20 @@ fn indices_from_js(what: &str, v: &[f64]) -> Result<Vec<u64>, JsError> {
.collect()
}
fn slab_from_js(
start: &[f64],
count: &[f64],
stride: Option<Vec<f64>>,
block: Option<Vec<f64>>,
) -> Result<Hyperslab, JsError> {
Ok(Hyperslab {
start: indices_from_js("start", start)?,
count: indices_from_js("count", count)?,
stride: stride.map(|s| indices_from_js("stride", &s)).transpose()?,
block: block.map(|b| indices_from_js("block", &b)).transpose()?,
})
}
fn data_to_js(data: Data) -> JsValue {
match data {
Data::F32(v) => js_sys::Float32Array::from(&v[..]).into(),
@@ -110,6 +157,72 @@ fn attr_to_js(value: AttrValue) -> (JsValue, Option<String>) {
}
}
fn list_to_js(children: Vec<Child>) -> Array {
children
.into_iter()
.map(|c| {
let o = Object::new();
set(&o, "name", c.name);
set(&o, "kind", c.kind.as_str());
JsValue::from(o)
})
.collect()
}
fn info_to_js(i: DatasetInfo) -> Object {
let o = Object::new();
set(&o, "shape", shape_to_js(&i.shape));
let max: JsValue = match i.maxshape {
None => JsValue::NULL,
Some(dims) => dims
.into_iter()
.map(|d| d.map_or(JsValue::NULL, |d| JsValue::from_f64(d as f64)))
.collect::<Array>()
.into(),
};
set(&o, "maxshape", max);
set(&o, "dtype", i.dtype);
set(&o, "elementShape", shape_to_js(&i.element_shape));
o
}
fn attrs_to_js(attrs: Vec<Attr>) -> Array {
attrs
.into_iter()
.map(|a| {
let o = Object::new();
set(&o, "name", a.name);
let (value, dtype) = attr_to_js(a.value);
set(&o, "value", value);
set(&o, "dtype", dtype.map_or(JsValue::NULL, JsValue::from));
JsValue::from(o)
})
.collect()
}
fn errors_to_js(errors: Vec<String>) -> Array {
errors.into_iter().map(JsValue::from).collect()
}
/// A read's values with the dataset's datatype.
fn values_to_js((dtype, a): (String, core::Array)) -> Object {
let o = Object::new();
set(&o, "shape", shape_to_js(&a.shape));
set(&o, "dtype", dtype);
set(&o, "data", data_to_js(a.data));
o
}
/// Read the dataset at `path` (whole, or `slab`) with its datatype.
fn read_values(
r: &Reader,
path: &str,
slab: Option<&Hyperslab>,
) -> core::Result<(String, core::Array)> {
let dtype = r.info(path)?.dtype;
Ok((dtype, r.read(path, slab)?))
}
/// An open HDF5 (or NetCDF-4) file.
#[wasm_bindgen]
pub struct H5File {
@@ -145,38 +258,13 @@ impl H5File {
/// The group's members: `[{ name, kind }]`, groups first.
pub fn list(&self, path: &str) -> Result<Array, JsError> {
Ok(self
.inner
.list(path)
.map_err(js_err)?
.into_iter()
.map(|c| {
let o = Object::new();
set(&o, "name", c.name);
set(&o, "kind", c.kind.as_str());
JsValue::from(o)
})
.collect())
Ok(list_to_js(self.inner.list(path).map_err(js_err)?))
}
/// `{ shape, maxshape, dtype, elementShape }`. `maxshape` is `null`
/// when not recorded, with `null` for each unlimited dimension.
pub fn info(&self, path: &str) -> Result<Object, JsError> {
let i = self.inner.info(path).map_err(js_err)?;
let o = Object::new();
set(&o, "shape", shape_to_js(&i.shape));
let max: JsValue = match i.maxshape {
None => JsValue::NULL,
Some(dims) => dims
.into_iter()
.map(|d| d.map_or(JsValue::NULL, |d| JsValue::from_f64(d as f64)))
.collect::<Array>()
.into(),
};
set(&o, "maxshape", max);
set(&o, "dtype", i.dtype);
set(&o, "elementShape", shape_to_js(&i.element_shape));
Ok(o)
Ok(info_to_js(self.inner.info(path).map_err(js_err)?))
}
/// `[{ name, value, dtype }]`, sorted by name. Scalars are `number`
@@ -185,31 +273,21 @@ impl H5File {
/// and its `dtype`; one that could not be read at all is reported by
/// [`attrErrors`](Self::attr_errors).
pub fn attrs(&self, path: &str) -> Result<Array, JsError> {
let (attrs, _) = self.inner.attrs(path).map_err(js_err)?;
Ok(attrs
.into_iter()
.map(|a| {
let o = Object::new();
set(&o, "name", a.name);
let (value, dtype) = attr_to_js(a.value);
set(&o, "value", value);
set(&o, "dtype", dtype.map_or(JsValue::NULL, JsValue::from));
JsValue::from(o)
})
.collect())
Ok(attrs_to_js(self.inner.attrs(path).map_err(js_err)?.0))
}
/// Messages for attributes that could not be read.
#[wasm_bindgen(js_name = attrErrors)]
pub fn attr_errors(&self, path: &str) -> Result<Array, JsError> {
let (_, errors) = self.inner.attrs(path).map_err(js_err)?;
Ok(errors.into_iter().map(JsValue::from).collect())
Ok(errors_to_js(self.inner.attrs(path).map_err(js_err)?.1))
}
/// The whole dataset: `{ shape, dtype, data }`, `data` in row-major
/// order.
pub fn read(&self, path: &str) -> Result<Object, JsError> {
self.read_impl(path, None)
Ok(values_to_js(
read_values(&self.inner, path, None).map_err(js_err)?,
))
}
/// A regular hyperslab (`H5Sselect_hyperslab`): `stride` and `block`
@@ -223,22 +301,358 @@ impl H5File {
stride: Option<Vec<f64>>,
block: Option<Vec<f64>>,
) -> Result<Object, JsError> {
let slab = Hyperslab {
start: indices_from_js("start", &start)?,
count: indices_from_js("count", &count)?,
stride: stride.map(|s| indices_from_js("stride", &s)).transpose()?,
block: block.map(|b| indices_from_js("block", &b)).transpose()?,
};
self.read_impl(path, Some(&slab))
}
fn read_impl(&self, path: &str, slab: Option<&Hyperslab>) -> Result<Object, JsError> {
let dtype = self.inner.info(path).map_err(js_err)?.dtype;
let a = self.inner.read(path, slab).map_err(js_err)?;
let slab = slab_from_js(&start, &count, stride, block)?;
Ok(values_to_js(
read_values(&self.inner, path, Some(&slab)).map_err(js_err)?,
))
}
}
// ---------------------------------------------------------------------------
// Remote files: openUrl.
#[wasm_bindgen(module = "/js/remote.js")]
extern "C" {
#[wasm_bindgen(catch)]
async fn probe(url: &str, first_len: f64, opts: &JsValue) -> Result<JsValue, JsValue>;
#[wasm_bindgen(catch, js_name = fetchRanges)]
async fn fetch_ranges(
url: &str,
ranges: Vec<f64>,
opts: &JsValue,
validator: &JsValue,
length: f64,
) -> Result<JsValue, JsValue>;
}
/// Where a remote file's bytes come from.
struct Http {
url: String,
opts: JsValue,
/// ETag or Last-Modified at open (`null` if the server sent neither).
validator: JsValue,
length: u64,
/// Requests the probe made (1, or 2 with a HEAD for the length).
probe_requests: u64,
}
enum Source {
/// Read by range requests through a restartable cache.
Lazy {
http: Http,
storage: Arc<LazyStorage>,
reader: Reader,
},
/// The server ignored `Range`: the whole file, downloaded at open.
Whole {
reader: Reader,
size: u64,
requests: u64,
},
}
impl Http {
/// Fetch `ranges` and hand them to `storage`.
async fn fetch(&self, storage: &LazyStorage, ranges: &[Range<u64>]) -> Result<(), JsError> {
let flat: Vec<f64> = ranges
.iter()
.flat_map(|r| [r.start as f64, r.end as f64])
.collect();
let got = fetch_ranges(
&self.url,
flat,
&self.opts,
&self.validator,
self.length as f64,
)
.await
.map_err(js_exception)?;
let got = Array::from(&got);
if got.length() as usize != ranges.len() {
return Err(js_err(format!(
"fetchRanges returned {} ranges for {}",
got.length(),
ranges.len()
)));
}
for (r, bytes) in ranges.iter().zip(got.iter()) {
let bytes = Uint8Array::new(&bytes).to_vec();
storage.supply_range(r, &bytes).map_err(js_err)?;
}
Ok(())
}
/// Run `f` over `storage` until it has every byte it reads.
async fn drive<T>(
&self,
storage: &LazyStorage,
mut f: impl FnMut() -> T,
) -> Result<T, JsError> {
let op = storage.operation();
loop {
match storage.attempt(&mut f) {
Step::Done(v) => return Ok(v),
Step::Need(ranges) => {
op.charge(&ranges).map_err(js_err)?;
self.fetch(storage, &ranges).await?
}
}
}
}
}
impl Source {
async fn run<T>(&self, op: impl Fn(&Reader) -> core::Result<T>) -> Result<T, JsError> {
match self {
Source::Whole { reader, .. } => op(reader).map_err(js_err),
Source::Lazy {
http,
storage,
reader,
} => http.drive(storage, || op(reader)).await?.map_err(js_err),
}
}
}
/// A non-negative integer option, or `None` when not given.
fn int_opt(opts: &JsValue, key: &str, min: f64, max: f64) -> Result<Option<u64>, JsError> {
let v = get(opts, key);
if v.is_undefined() || v.is_null() {
return Ok(None);
}
match v.as_f64() {
Some(x) if x.fract() == 0.0 && (min..=max).contains(&x) => Ok(Some(x as u64)),
_ => Err(js_err(format!(
"openUrl: {key} must be an integer from {min} to {max}"
))),
}
}
/// Most bytes `maxFetch` and `maxDownload` may allow: 1 GiB. wasm32 has
/// 4 GiB of memory and no buffer past 2 GiB, and what is fetched is held
/// while it is decoded.
const MAX_FETCH_LIMIT: u64 = 1 << 30;
/// The largest file `openUrl` reads by ranges: on wasm32, 4 GiB - 1 bytes.
/// The format code turns file offsets into `usize` to use them (with a
/// clean error past it, see scripts/check-32bit-casts.sh), so on a 32-bit
/// target nothing at 4 GiB or beyond can be read; a larger file is refused
/// at open rather than failing on whichever read reaches past 4 GiB. On
/// 64-bit targets it is 2^53 - 1, the largest offset a JavaScript number
/// holds exactly.
const MAX_REMOTE_LENGTH: u64 = if (usize::MAX as u64) < MAX_SAFE_INTEGER as u64 {
usize::MAX as u64
} else {
MAX_SAFE_INTEGER as u64
};
fn config_from(opts: &JsValue) -> Result<LazyConfig, JsError> {
let mut c = LazyConfig::default();
if let Some(b) = int_opt(opts, "blockSize", 512.0, (64u64 << 20) as f64)? {
c.block_size = b;
}
if let Some(n) = int_opt(opts, "cacheSize", 0.0, MAX_SAFE_INTEGER)? {
c.capacity = n;
}
if let Some(n) = int_opt(opts, "maxFetch", 512.0, MAX_FETCH_LIMIT as f64)? {
c.max_fetch = n;
}
// Read by remote.js; checked here so a value wasm32 cannot hold is an
// option error rather than a download that cannot be kept.
int_opt(opts, "maxDownload", 0.0, MAX_FETCH_LIMIT as f64)?;
// Also read by remote.js (which checks it too, for direct callers).
int_opt(opts, "parallel", 1.0, 1024.0)?;
Ok(c)
}
/// Open the HDF5 file at `url` without downloading it: its bytes are
/// fetched with HTTP `Range` requests as the methods of the returned
/// [`RemoteFile`] need them, through a block cache.
///
/// `opts` (all optional):
/// - `blockSize` — bytes per request block, 512 to 64 MiB (default 1 MiB);
/// - `cacheSize` — bytes of blocks kept between calls (default 64 MiB);
/// - `maxFetch` — most bytes one call may fetch, and so the longest single
/// read, up to 1 GiB (default 512 MiB): a call that would fetch more
/// fails before fetching it;
/// - `fallback` — `"download"` (default) reads the whole file when the
/// server ignores `Range` (answers 200), up to `maxDownload` bytes
/// (default 512 MiB, at most 1 GiB); `"error"` refuses such a server;
/// - `headers`, `credentials` — passed to every `fetch` (`headers` as
/// `fetch` takes them: a `Headers`, `[name, value]` pairs or an object);
/// - `parallel` — range requests in flight at once, 1 to 1024 (default
/// 6); when one fails the others are aborted;
/// - `fetch` — a `fetch`-compatible function to use instead of the global.
///
/// Cross-origin servers must allow CORS and expose `Content-Range` (or
/// answer `HEAD` with `Content-Length`). A file may be up to 4 GiB - 1
/// bytes long (wasm32 offsets); a longer one is refused at open. A whole-dataset `read` that would use more than 1 GiB of
/// memory ([`core::MAX_READ_BYTES`]) is refused: read it in parts with
/// `readHyperslab`.
#[wasm_bindgen(js_name = openUrl)]
pub async fn open_url(url: String, opts: JsValue) -> Result<RemoteFile, JsError> {
let config = config_from(&opts)?;
let p = probe(&url, config.block_size as f64, &opts)
.await
.map_err(js_exception)?;
let requests = get(&p, "requests").as_f64().unwrap_or(1.0) as u64;
let whole = get(&p, "whole");
if !whole.is_undefined() {
let bytes = Uint8Array::new(&whole).to_vec();
let size = bytes.len() as u64;
let reader = Reader::open(bytes).map_err(js_err)?;
return Ok(RemoteFile {
inner: Rc::new(Source::Whole {
reader,
size,
requests,
}),
});
}
let length = get(&p, "length")
.as_f64()
.filter(|x| x.fract() == 0.0 && (0.0..=MAX_SAFE_INTEGER).contains(x))
.ok_or_else(|| js_err(format!("{url}: the server gave no usable file size")))?;
if length as u64 > MAX_REMOTE_LENGTH {
return Err(js_err(format!(
"{url} is {length} bytes; openUrl reads files of up to {MAX_REMOTE_LENGTH} bytes \
(4 GiB - 1: the WebAssembly reader addresses a file with 32-bit offsets)"
)));
}
let http = Http {
url,
opts,
validator: get(&p, "validator"),
length: length as u64,
probe_requests: requests,
};
let storage = Arc::new(LazyStorage::new(http.length, config));
let first = Uint8Array::new(&get(&p, "first")).to_vec();
storage.supply(0, &first).map_err(js_err)?;
let s = storage.clone();
let reader = http
.drive(&storage, || Reader::open_storage(s.clone()))
.await?
.map_err(js_err)?;
Ok(RemoteFile {
inner: Rc::new(Source::Lazy {
http,
storage,
reader,
}),
})
}
/// A file opened with [`openUrl`](open_url): the methods of [`H5File`],
/// each returning a `Promise` (it may have to fetch bytes first).
#[wasm_bindgen]
pub struct RemoteFile {
inner: Rc<Source>,
}
impl RemoteFile {
/// Run `op` (fetching what it needs) and convert its result.
fn call<T: 'static>(
&self,
op: impl Fn(&Reader) -> core::Result<T> + 'static,
to_js: impl FnOnce(T) -> JsValue + 'static,
) -> Promise {
let inner = self.inner.clone();
future_to_promise(async move {
let v = inner.run(op).await.map_err(JsValue::from)?;
Ok(to_js(v))
})
}
}
#[wasm_bindgen]
impl RemoteFile {
/// `"group"` or `"dataset"`.
#[wasm_bindgen(unchecked_return_type = "Promise<string>")]
pub fn kind(&self, path: String) -> Promise {
self.call(move |r| r.kind(&path), |k| k.as_str().into())
}
/// The group's members: `[{ name, kind }]`, groups first.
#[wasm_bindgen(unchecked_return_type = "Promise<Array<any>>")]
pub fn list(&self, path: String) -> Promise {
self.call(move |r| r.list(&path), |c| list_to_js(c).into())
}
/// `{ shape, maxshape, dtype, elementShape }`, as [`H5File::info`].
#[wasm_bindgen(unchecked_return_type = "Promise<any>")]
pub fn info(&self, path: String) -> Promise {
self.call(move |r| r.info(&path), |i| info_to_js(i).into())
}
/// `[{ name, value, dtype }]`, as [`H5File::attrs`].
#[wasm_bindgen(unchecked_return_type = "Promise<Array<any>>")]
pub fn attrs(&self, path: String) -> Promise {
self.call(move |r| r.attrs(&path), |a| attrs_to_js(a.0).into())
}
/// Messages for attributes that could not be read.
#[wasm_bindgen(js_name = attrErrors, unchecked_return_type = "Promise<Array<string>>")]
pub fn attr_errors(&self, path: String) -> Promise {
self.call(move |r| r.attrs(&path), |a| errors_to_js(a.1).into())
}
/// The whole dataset: `{ shape, dtype, data }`, as [`H5File::read`].
#[wasm_bindgen(unchecked_return_type = "Promise<any>")]
pub fn read(&self, path: String) -> Promise {
self.call(
move |r| read_values(r, &path, None),
|v| values_to_js(v).into(),
)
}
/// A regular hyperslab, as [`H5File::read_hyperslab`]. Only the chunks
/// (or the contiguous runs) the selection touches are fetched.
#[wasm_bindgen(js_name = readHyperslab, unchecked_return_type = "Promise<any>")]
pub fn read_hyperslab(
&self,
path: String,
start: Vec<f64>,
count: Vec<f64>,
stride: Option<Vec<f64>>,
block: Option<Vec<f64>>,
) -> Result<Promise, JsError> {
let slab = slab_from_js(&start, &count, stride, block)?;
Ok(self.call(
move |r| read_values(r, &path, Some(&slab)),
|v| values_to_js(v).into(),
))
}
/// What reading this file has cost so far: `{ lazy, size, requests,
/// bytesFetched, cachedBytes, passes }`. `lazy` is false when the
/// server ignored `Range` and the file was downloaded whole.
pub fn stats(&self) -> Object {
let o = Object::new();
set(&o, "shape", shape_to_js(&a.shape));
set(&o, "dtype", dtype);
set(&o, "data", data_to_js(a.data));
Ok(o)
match &*self.inner {
Source::Lazy { http, storage, .. } => {
let st = storage.stats();
set(&o, "lazy", true);
set(&o, "size", http.length as f64);
set(
&o,
"requests",
(st.requests.saturating_sub(1) + http.probe_requests) as f64,
);
set(&o, "bytesFetched", st.bytes_fetched as f64);
set(&o, "cachedBytes", st.cached_bytes as f64);
set(&o, "passes", st.passes as f64);
}
Source::Whole { size, requests, .. } => {
set(&o, "lazy", false);
set(&o, "size", *size as f64);
set(&o, "requests", *requests as f64);
set(&o, "bytesFetched", *size as f64);
set(&o, "cachedBytes", *size as f64);
set(&o, "passes", 0.0);
}
}
o
}
}
+578
View File
@@ -0,0 +1,578 @@
//! The restartable ("NeedBytes") reader against the in-memory one: every
//! file must list, describe and read the same through a [`LazyStorage`]
//! that starts empty and is fed only the ranges its passes ask for, as the
//! browser's `openUrl` feeds it from HTTP range requests.
//!
//! - Files written here with `FileBuilder`, at several block sizes (512 B
//! blocks make almost every structure read a miss).
//! - The h5py/netCDF4 fixture of `examples/wasm-viewer/test/make_fixture.py`
//! (skipped without h5py, unless `CLAWHDF5_REQUIRE_INTEROP=1`;
//! `CLAWHDF5_PYTHON` names the interpreter).
//! - `CLAWHDF5_WASM_CORPUS=dir[:dir...]`: every HDF5 file under those
//! directories up to 64 MiB (e.g. `conformance/.cache/corpus`).
//!
//! Also the request budget: listing and reading one small dataset of a large
//! file fetches a few blocks, not the file.
use std::ops::Range;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
use clawhdf5::{AttrValue, FileBuilder};
use clawhdf5_format::storage::CountingStorage;
use clawhdf5_wasm::core::{Hyperslab, Kind, Reader};
use clawhdf5_wasm::lazy::{LazyConfig, LazyStorage};
/// The API the JavaScript side calls, one operation at a time.
trait Api {
fn call<T>(&self, op: impl Fn(&Reader) -> T) -> T;
}
struct Local(Reader);
impl Api for Local {
fn call<T>(&self, op: impl Fn(&Reader) -> T) -> T {
op(&self.0)
}
}
/// A lazily read file and the "server" it fetches from.
struct Lazy {
data: Arc<Vec<u8>>,
storage: Arc<LazyStorage>,
reader: Reader,
}
fn fetch(data: &[u8], r: Range<u64>) -> Result<Vec<u8>, String> {
Ok(data[r.start as usize..r.end as usize].to_vec())
}
impl Lazy {
/// Open as `openUrl` does: the first block comes with the probe that
/// learns the length, then the open is run until it has its bytes.
fn open(data: Vec<u8>, config: LazyConfig) -> Result<Lazy, String> {
let data = Arc::new(data);
let storage = Arc::new(LazyStorage::new(data.len() as u64, config));
let first = (storage.config().block_size as usize).min(data.len());
storage.supply(0, &data[..first])?;
let s = storage.clone();
let reader =
storage.run_blocking(|| Reader::open_storage(s.clone()), |r| fetch(&data, r))??;
Ok(Lazy {
data,
storage,
reader,
})
}
}
impl Api for Lazy {
fn call<T>(&self, op: impl Fn(&Reader) -> T) -> T {
self.storage
.run_blocking(|| op(&self.reader), |r| fetch(&self.data, r))
.expect("serving from memory cannot fail")
}
}
/// Everything the viewer can show of a file, as text: each object's kind,
/// listing, attributes (and attribute errors), dataset info, whole value
/// and a hyperslab — or the error each gives.
fn transcript(api: &impl Api) -> Vec<String> {
let mut out = Vec::new();
let mut todo = vec![("/".to_string(), 0usize)];
while let Some((path, depth)) = todo.pop() {
if out.len() > 4000 {
out.push("... (truncated)".into());
break;
}
let kind = api.call(|r| r.kind(&path));
out.push(format!("{path}: {kind:?}"));
out.push(format!("{path} attrs: {:?}", api.call(|r| r.attrs(&path))));
match kind {
Ok(Kind::Group) => {
let list = api.call(|r| r.list(&path));
out.push(format!("{path} list: {list:?}"));
if let Ok(children) = list
&& depth < 12
{
for c in children.into_iter().rev() {
let child = if path == "/" {
format!("/{}", c.name)
} else {
format!("{path}/{}", c.name)
};
todo.push((child, depth + 1));
}
}
}
Ok(Kind::Dataset) => {
let info = api.call(|r| r.info(&path));
out.push(format!("{path} info: {info:?}"));
let Ok(info) = info else { continue };
let n = info
.shape
.iter()
.chain(&info.element_shape)
.try_fold(1u64, |a, &d| a.checked_mul(d));
if n.is_none_or(|n| n > 4_000_000) {
out.push(format!("{path}: not read ({n:?} values)"));
continue;
}
out.push(format!(
"{path} read: {:?}",
api.call(|r| r.read(&path, None))
));
if !info.shape.is_empty() && info.shape.iter().all(|&d| d > 1) {
let slab = Hyperslab {
start: info.shape.iter().map(|_| 1).collect(),
count: info.shape.iter().map(|&d| d / 2).collect(),
stride: None,
block: None,
};
let part = api.call(|r| r.read(&path, Some(&slab)));
out.push(format!("{path} slab: {part:?}"));
}
}
Err(_) => {}
}
}
out
}
/// The lazy transcript of `data` at `block` bytes per block equals the
/// transcript of the same file through a range storage that has every byte
/// (`CountingStorage`: the facade's `Storage` path, the one the lazy reader
/// 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>) {
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());
let lazy = Lazy::open(data.to_vec(), config(block));
let (ranged, local, lazy) = match (ranged, local, lazy) {
(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());
}
(r, l, z) => panic!(
"{ctx}: opens differently: ranged {:?}, in memory {:?}, lazily {:?}",
r.err(),
l.err(),
z.err()
),
};
let got = transcript(&lazy);
let want = transcript(&Local(ranged));
for (i, (w, g)) in want.iter().zip(&got).enumerate() {
assert_eq!(g, w, "{ctx}, line {i}");
}
assert_eq!(got.len(), want.len(), "{ctx}: transcript length");
let local = transcript(&Local(local));
for (i, (l, g)) in local.iter().zip(&got).enumerate() {
let both_errors = match (l.split_once("Err("), g.split_once("Err(")) {
(Some((a, _)), Some((b, _))) => a == b,
_ => false,
};
assert!(
l == g || both_errors,
"{ctx}, line {i}: in memory\n {l}\nlazily\n {g}"
);
}
assert_eq!(got.len(), local.len(), "{ctx}: transcript length");
let st = lazy.storage.stats();
(st.requests, st.bytes_fetched, got)
}
fn config(block: u64) -> LazyConfig {
LazyConfig {
block_size: block,
// A small budget, so eviction between operations is exercised.
capacity: 16 * block,
max_request: 8 * block,
..LazyConfig::default()
}
}
fn builder_file() -> Vec<u8> {
let mut b = FileBuilder::new();
b.create_dataset("grid")
.with_f64_data(&(0..20_000).map(f64::from).collect::<Vec<_>>())
.with_shape(&[100, 200])
.with_chunks(&[10, 25])
.with_deflate(4);
b.create_dataset("contiguous")
.with_i32_data(&(0..50_000).collect::<Vec<_>>());
b.create_dataset("bytes").with_u8_data(&[1, 2, 250]);
let mut g = b.create_group("sensors");
for i in 0..40 {
g.create_dataset(&format!("t{i}"))
.with_f32_data(&[i as f32, 1.5, -2.25]);
}
g.set_attr("location", AttrValue::String("lab".into()));
b.add_group(g.finish());
b.set_attr("version", AttrValue::I64(3));
b.set_attr("scale", AttrValue::F64Array(vec![0.5, 2.0]));
b.finish().unwrap()
}
#[test]
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);
assert!(requests > 0);
// The transcript covers every object, values included.
assert!(lines.iter().any(|l| l.starts_with("/grid read: Ok")));
assert!(lines.iter().any(|l| l.starts_with("/grid slab: Ok")));
assert!(lines.iter().any(|l| l.starts_with("/sensors/t39 read: Ok")));
}
}
#[test]
fn garbage_fails_to_open_as_in_memory() {
check_equal("zeros", &[0u8; 5000], 512);
check_equal("empty", &[], 512);
let mut cut = builder_file();
cut.truncate(cut.len() / 3);
check_equal("truncated", &cut, 512);
}
/// Listing a large file and reading one small dataset fetches a few blocks,
/// not the file.
#[test]
fn a_small_read_of_a_large_file_fetches_a_few_blocks() {
let mut b = FileBuilder::new();
b.create_dataset("small").with_f64_data(&[1.0, 2.0, 3.0]);
// 48 MB of raw data, written after the small dataset's metadata.
b.create_dataset("big")
.with_f64_data(&(0..6_000_000).map(f64::from).collect::<Vec<_>>());
let mut g = b.create_group("group");
g.create_dataset("inner").with_i32_data(&[7, 8]);
b.add_group(g.finish());
let data = b.finish().unwrap();
let lazy = Lazy::open(data.clone(), LazyConfig::default()).unwrap();
let list = lazy.call(|r| r.list("/")).unwrap();
assert_eq!(list.len(), 3);
assert_eq!(
format!("{:?}", lazy.call(|r| r.read("/small", None)).unwrap().data),
"F64([1.0, 2.0, 3.0])"
);
assert_eq!(
format!(
"{:?}",
lazy.call(|r| r.read("/group/inner", None)).unwrap().data
),
"I32([7, 8])"
);
// A window of the big dataset reads only its block(s).
let slab = Hyperslab {
start: vec![3_000_000],
count: vec![4],
stride: None,
block: None,
};
assert_eq!(
format!(
"{:?}",
lazy.call(|r| r.read("/big", Some(&slab))).unwrap().data
),
"F64([3000000.0, 3000001.0, 3000002.0, 3000003.0])"
);
let st = lazy.storage.stats();
eprintln!("{} bytes: {st:?}", data.len());
assert!(st.requests <= 6, "{st:?}");
assert!(st.bytes_fetched <= 6 << 20, "{st:?}");
assert!(st.bytes_fetched * 8 < data.len() as u64, "{st:?}");
}
fn python() -> String {
std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string())
}
fn python_available() -> bool {
Command::new(python())
.args(["-c", "import h5py, netCDF4, numpy"])
.output()
.is_ok_and(|o| o.status.success())
}
#[test]
fn h5py_and_netcdf4_files_read_the_same_lazily() {
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 = fixture_dir();
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);
assert!(lines.iter().filter(|l| l.contains(" read: Ok")).count() >= 2);
}
}
}
/// The bytes of `data` at `r`, zero past its end: a server that claims
/// the file is longer than it is.
fn fetch_padded(data: &[u8], r: Range<u64>) -> Result<Vec<u8>, String> {
let mut out = vec![0u8; (r.end - r.start) as usize];
let len = data.len() as u64;
if r.start < len {
let end = r.end.min(len);
out[..(end - r.start) as usize].copy_from_slice(&data[r.start as usize..end as usize]);
}
Ok(out)
}
/// Sizes a hostile server or a large dataset can name are errors, never
/// allocations that abort the wasm module: make_fixture.py's limits.h5 and
/// hostile_vl.h5 (see write_limits there).
#[test]
fn size_limits_are_errors_not_aborts() {
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 = fixture_dir();
// Read whole, /huge_u8 would widen 2^28 values to 64 bits (2 GiB): an
// error naming readHyperslab, before its chunks are read. A window of
// it reads.
let data = std::fs::read(dir.path().join("limits.h5")).unwrap();
let n = (1u64 << 28) + 1024;
let window = Hyperslab {
start: vec![n - 4],
count: vec![4],
stride: None,
block: None,
};
let local = Reader::open(data.clone()).unwrap();
let lazy = Lazy::open(data, LazyConfig::default()).unwrap();
let before = lazy.storage.stats().requests;
for e in [
local.read("/huge_u8", None).unwrap_err(),
lazy.call(|r| r.read("/huge_u8", None)).unwrap_err(),
] {
assert!(e.contains("readHyperslab"), "{e}");
}
assert_eq!(lazy.storage.stats().requests, before, "nothing fetched");
for part in [
local.read("/huge_u8", Some(&window)).unwrap(),
lazy.call(|r| r.read("/huge_u8", Some(&window))).unwrap(),
] {
assert_eq!(format!("{:?}", part.data), "U8([0, 0, 0, 7])");
}
// A server that claims 3 GiB and a heap collection of 2 GiB + 4 KiB:
// reading the strings fails at once, fetching a few blocks.
let data = std::fs::read(dir.path().join("hostile_vl.h5")).unwrap();
let storage = Arc::new(LazyStorage::new(3 << 30, LazyConfig::default()));
let s = storage.clone();
let reader = storage
.run_blocking(
|| Reader::open_storage(s.clone()),
|r| fetch_padded(&data, r),
)
.unwrap()
.unwrap();
let e = storage
.run_blocking(|| reader.read("/a", None), |r| fetch_padded(&data, r))
.unwrap()
.unwrap_err();
assert!(e.contains("maxFetch"), "{e}");
let st = storage.stats();
assert!(st.requests <= 4 && st.bytes_fetched <= 4 << 20, "{st:?}");
}
/// make_fixture.py's files, written to a temporary directory.
fn fixture_dir() -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
let generator = Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../../examples/wasm-viewer/test/make_fixture.py");
let out = Command::new(python())
.arg(&generator)
.arg(dir.path())
.output()
.unwrap();
assert!(
out.status.success(),
"{}",
String::from_utf8_lossy(&out.stderr)
);
dir
}
fn hdf5_files(dir: &Path, out: &mut Vec<PathBuf>) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
for e in entries.flatten() {
let p = e.path();
if p.is_dir() {
hdf5_files(&p, out);
} else if std::fs::read(&p)
.ok()
.is_some_and(|b| b.len() <= 64 << 20 && is_hdf5(&b))
{
out.push(p);
}
}
}
/// The HDF5 signature at 0 or a power-of-two user-block offset.
fn is_hdf5(b: &[u8]) -> bool {
const SIG: &[u8] = b"\x89HDF\r\n\x1a\n";
let mut at = 0usize;
loop {
if b.get(at..at + 8) == Some(SIG) {
return true;
}
at = if at == 0 { 512 } else { at * 2 };
if at >= b.len() {
return false;
}
}
}
#[test]
fn corpus_files_read_the_same_lazily() {
let Ok(dirs) = std::env::var("CLAWHDF5_WASM_CORPUS") else {
eprintln!("CLAWHDF5_WASM_CORPUS not set; skipping the corpus");
return;
};
let mut files = Vec::new();
for d in std::env::split_paths(&dirs) {
hdf5_files(&d, &mut files);
}
files.sort();
assert!(!files.is_empty(), "no HDF5 files under {dirs}");
let (mut requests, mut bytes, mut total) = (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);
requests += r;
bytes += b;
}
eprintln!(
"{} files ({total} bytes): {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();
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
);
}
}