From 076feb089a79e0141bcaf27cd041fa55ed6b4596 Mon Sep 17 00:00:00 2001 From: osobh Date: Sun, 27 Sep 2026 06:35:21 -0500 Subject: [PATCH 01/10] wasm: restartable NeedBytes storage for lazy reads (range-read M4, core) LazyStorage holds the blocks of a remote file fetched so far. An operation runs as passes over it: a read that misses records the missing blocks and fails, the pass's result is dropped whatever it is (a parser may have caught the error and carried on), and the caller fetches the reported ranges and re-runs the pass. No block is evicted while an operation is in flight, so every pass that does not finish asks for at least one new block and the operation ends. Blocks sit in an LRU with a byte budget, trimmed between operations, bulk (raw data) blocks first. Reader::open_storage opens a file through any Storage, and variable- length strings resolve through the file's storage instead of File::as_bytes, which panics for a file not in memory. tests/lazy.rs compares, file by file, what the viewer can show (kinds, listings, attributes, info, whole reads and hyperslabs) read lazily with the facade's range-storage path and the in-memory reader: files built here at 512 B to 1 MiB blocks, the h5py/netCDF4 fixture, and with CLAWHDF5_WASM_CORPUS the conformance corpus (656 files agree). A listing plus a small read of a 48 MB file fetches 3 ranges. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/clawhdf5-wasm/src/core.rs | 28 +- crates/clawhdf5-wasm/src/lazy.rs | 672 +++++++++++++++++++++++++++++ crates/clawhdf5-wasm/src/lib.rs | 1 + crates/clawhdf5-wasm/tests/lazy.rs | 391 +++++++++++++++++ 4 files changed, 1086 insertions(+), 6 deletions(-) create mode 100644 crates/clawhdf5-wasm/src/lazy.rs create mode 100644 crates/clawhdf5-wasm/tests/lazy.rs diff --git a/crates/clawhdf5-wasm/src/core.rs b/crates/clawhdf5-wasm/src/core.rs index 4c6beab..23b4b1d 100644 --- a/crates/clawhdf5-wasm/src/core.rs +++ b/crates/clawhdf5-wasm/src/core.rs @@ -6,9 +6,12 @@ //! 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::storage::Storage; use clawhdf5_format::vl_data::{VlResolver, check_element_size}; /// Errors are reported to JavaScript as messages. @@ -119,7 +122,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 +137,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) -> Result { + 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 { match self.file.dataset(path) { @@ -281,11 +294,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)?) diff --git a/crates/clawhdf5-wasm/src/lazy.rs b/crates/clawhdf5-wasm/src/lazy.rs new file mode 100644 index 0000000..193f61f --- /dev/null +++ b/crates/clawhdf5-wasm/src/lazy.rs @@ -0,0 +1,672 @@ +//! 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; + +/// 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, +} + +impl Default for LazyConfig { + fn default() -> Self { + LazyConfig { + block_size: DEFAULT_BLOCK_SIZE, + capacity: 64 << 20, + max_request: 8 << 20, + } + } +} + +/// 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 { + /// 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>), +} + +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, + /// `(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, + /// Blocks a bulk read missed that have not been supplied yet: kept + /// as bulk when they arrive. + bulk_pending: BTreeSet, + /// 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, +} + +fn lock(m: &Mutex) -> 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, +} + +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; + 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 } + } + + /// 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(&self, f: impl FnOnce() -> T) -> Step { + { + 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, 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( + &self, + mut f: impl FnMut() -> T, + mut fetch: impl FnMut(Range) -> Result, String>, + ) -> Result { + let _op = self.operation(); + loop { + match self.attempt(&mut f) { + Step::Done(v) => return Ok(v), + Step::Need(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, + /// each at most `max_request` long. + fn runs(&self, missing: HashMap) -> Vec> { + let bs = self.config.block_size; + let mut wanted: Vec = 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 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, + _ => runs.push((i, i)), + } + } + 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> { + 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], + metadata: bool, + ) -> Result>, 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) + } + + fn assemble(&self, offset: u64, end: u64, blocks: &HashMap>) -> Vec { + let bs = self.config.block_size; + let mut out = Vec::with_capacity(usize::try_from(end - offset).unwrap_or(0)); + 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; + } + out + } +} + +impl Storage for LazyStorage { + fn read_at(&self, offset: u64, len: usize) -> Result, FormatError> { + let Some(span) = self.span(offset, len as u64) else { + return Ok(Cow::Owned(Vec::new())); + }; + let metadata = len as u64 <= self.config.block_size; + let blocks = self.blocks(std::slice::from_ref(&span), metadata)?; + let end = offset.saturating_add(len as u64).min(self.len); + Ok(Cow::Owned(self.assemble(offset, end, &blocks))) + } + + fn len(&self) -> u64 { + self.len + } + + fn read_ranges(&self, ranges: &[Range]) -> Result>, FormatError> { + let mut spans = Vec::with_capacity(ranges.len()); + for r in ranges { + if r.end < r.start { + return Err(FormatError::Storage( + "read range ends before it starts".into(), + )); + } + spans.extend(self.span(r.start, r.end - r.start)); + } + let blocks = self.blocks(&spans, false)?; + Ok(ranges + .iter() + .map(|r| { + let end = r.end.min(self.len); + if r.start >= end { + Cow::Owned(Vec::new()) + } else { + Cow::Owned(self.assemble(r.start, end, &blocks)) + } + }) + .collect()) + } +} + +#[cfg(test)] +#[allow(clippy::single_range_in_vec_init)] +mod tests { + use super::*; + + fn file(n: usize) -> Vec { + (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, + } + } + + /// Supply every range of `need` from `data`. + fn serve(s: &LazyStorage, data: &[u8], need: &[Range]) { + for r in need { + s.supply_range(r, &data[r.start as usize..r.end as usize]) + .unwrap(); + } + } + + fn owned(r: Result, FormatError>) -> Result, 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::>()) + }; + 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 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::, _>>() + }, + |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| 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_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); + } +} diff --git a/crates/clawhdf5-wasm/src/lib.rs b/crates/clawhdf5-wasm/src/lib.rs index 7c6b4a9..2fed048 100644 --- a/crates/clawhdf5-wasm/src/lib.rs +++ b/crates/clawhdf5-wasm/src/lib.rs @@ -21,6 +21,7 @@ //! The logic lives in [`core`], which is plain Rust and tested natively. pub mod core; +pub mod lazy; use clawhdf5::AttrValue; use js_sys::{Array, Object, Reflect}; diff --git a/crates/clawhdf5-wasm/tests/lazy.rs b/crates/clawhdf5-wasm/tests/lazy.rs new file mode 100644 index 0000000..eed0ed5 --- /dev/null +++ b/crates/clawhdf5-wasm/tests/lazy.rs @@ -0,0 +1,391 @@ +//! 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(&self, op: impl Fn(&Reader) -> T) -> T; +} + +struct Local(Reader); + +impl Api for Local { + fn call(&self, op: impl Fn(&Reader) -> T) -> T { + op(&self.0) + } +} + +/// A lazily read file and the "server" it fetches from. +struct Lazy { + data: Arc>, + storage: Arc, + reader: Reader, +} + +fn fetch(data: &[u8], r: Range) -> Result, 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, config: LazyConfig) -> Result { + 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(&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 { + 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) { + 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, + } +} + +fn builder_file() -> Vec { + let mut b = FileBuilder::new(); + b.create_dataset("grid") + .with_f64_data(&(0..20_000).map(f64::from).collect::>()) + .with_shape(&[100, 200]) + .with_chunks(&[10, 25]) + .with_deflate(4); + b.create_dataset("contiguous") + .with_i32_data(&(0..50_000).collect::>()); + 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::>()); + 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 = 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) + ); + 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); + } + } +} + +fn hdf5_files(dir: &Path, out: &mut Vec) { + 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() + ); +} From 5107583b97c8e722b9e6c51e422091f3f4a2018c Mon Sep 17 00:00:00 2001 From: osobh Date: Sun, 27 Sep 2026 06:44:45 -0500 Subject: [PATCH 02/10] wasm: openUrl reads remote files by HTTP range requests (range-read M4) openUrl(url, opts) returns a RemoteFile with the methods of H5File (kind, list, info, attrs, attrErrors, read, readHyperslab), each a promise, and stats(). It runs every call through the restartable LazyStorage: a pass that misses reports the byte ranges, js/remote.js fetches them with fetch() and Range headers (six at a time), and the pass is re-run. This keeps the main thread free without a Worker or synchronous XHR (h5wasm's lazy files need both), as the design doc recommends; the cost is re-running a pass per wave of misses. Every answer is checked: a 206 with exactly the bytes asked for, and the same ETag/Last-Modified and length as at open, else an error (never data). A server that ignores Range (200) is downloaded whole, up to maxDownload (512 MiB), unless fallback: "error". Options: blockSize, cacheSize, headers, credentials, parallel, fetch. test/serve.py is a range-capable static server with request counting (and /norange/ for a server without range support). test.mjs repeats every fixture check on files opened by URL (1 MiB and 512 B blocks), checks the request budget on a 200 MB h5py file (list, three small reads and a window of the big dataset: 5 requests, 6 MiB), the download fallback, and HTTP errors, changed files and wrong answers; with CLAWHDF5_WASM_CORPUS every corpus file is compared with open(bytes). Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/clawhdf5-wasm/Cargo.toml | 2 + crates/clawhdf5-wasm/js/remote.js | 182 ++++++++ crates/clawhdf5-wasm/src/lib.rs | 498 +++++++++++++++++++--- examples/wasm-viewer/test/make_fixture.py | 37 ++ examples/wasm-viewer/test/run.sh | 31 +- examples/wasm-viewer/test/serve.py | 201 +++++++++ examples/wasm-viewer/test/test.mjs | 338 +++++++++++++-- examples/wasm-viewer/viewer-lib.js | 21 + 8 files changed, 1209 insertions(+), 101 deletions(-) create mode 100644 crates/clawhdf5-wasm/js/remote.js create mode 100644 examples/wasm-viewer/test/serve.py diff --git a/crates/clawhdf5-wasm/Cargo.toml b/crates/clawhdf5-wasm/Cargo.toml index 303837a..6d222bc 100644 --- a/crates/clawhdf5-wasm/Cargo.toml +++ b/crates/clawhdf5-wasm/Cargo.toml @@ -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" diff --git a/crates/clawhdf5-wasm/js/remote.js b/crates/clawhdf5-wasm/js/remote.js new file mode 100644 index 0000000..fdbe619 --- /dev/null +++ b/crates/clawhdf5-wasm/js/remote.js @@ -0,0 +1,182 @@ +// 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; +} + +function init(opts, extra, method = "GET") { + return { method, headers: { ...(opts?.headers ?? {}), ...extra }, credentials: opts?.credentials }; +} + +// "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 whole body, refusing more than `limit` bytes as they arrive. +async function readAll(resp, limit, url) { + const tooBig = (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"); + const declared = resp.headers.get("Content-Length"); + if (declared !== null && Number(declared) > limit) { + await discard(resp); + throw tooBig(declared); + } + // A declared length was checked above: read the body at once. (Only an + // undeclared length is streamed, to stop at the limit; stream reads + // also stalled in the headless Chromium test under --virtual-time-budget.) + if (!resp.body || declared !== null) { + const all = new Uint8Array(await resp.arrayBuffer()); + if (all.length > limit) throw tooBig(all.length); + return all; + } + const reader = resp.body.getReader(); + const parts = []; + let n = 0; + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + n += value.length; + if (n > limit) { + await reader.cancel(); + throw tooBig(`over ${limit}`); + } + parts.push(value); + } + const all = new Uint8Array(n); + let at = 0; + for (const p of parts) { + all.set(p, at); + at += p.length; + } + return all; +} + +/** + * 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); + 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 = new Uint8Array(await resp.arrayBuffer()); + 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) || 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). + */ +export async function fetchRanges(url, ranges, opts, validator, length) { + const f = fetcher(opts); + const n = ranges.length / 2; + const out = new Array(n); + let next = 0; + async function worker() { + while (next < n) { + const i = next++; + const start = ranges[2 * i]; + const end = ranges[2 * i + 1]; + const resp = await f(url, init(opts, { Range: `bytes=${start}-${end - 1}` })); + 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 = new Uint8Array(await resp.arrayBuffer()); + if (body.length !== end - start) { + throw new Error(`${url}: asked for ${end - start} bytes at offset ${start}, got ${body.length}`); + } + out[i] = body; + } + } + const workers = Math.max(1, Math.min(opts?.parallel ?? DEFAULT_PARALLEL, n)); + await Promise.all(Array.from({ length: workers }, worker)); + return out; +} diff --git a/crates/clawhdf5-wasm/src/lib.rs b/crates/clawhdf5-wasm/src/lib.rs index 2fed048..a421573 100644 --- a/crates/clawhdf5-wasm/src/lib.rs +++ b/crates/clawhdf5-wasm/src/lib.rs @@ -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,16 +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 clawhdf5::AttrValue; -use js_sys::{Array, Object, Reflect}; -use wasm_bindgen::prelude::*; +use std::ops::Range; +use std::rc::Rc; +use std::sync::Arc; -use crate::core::{Data, Hyperslab, Reader}; +use clawhdf5::AttrValue; +use js_sys::{Array, Object, Promise, Reflect, Uint8Array}; +use wasm_bindgen::prelude::*; +use wasm_bindgen_futures::future_to_promise; + +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; @@ -36,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::() + .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) { // 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() } @@ -59,6 +91,20 @@ fn indices_from_js(what: &str, v: &[f64]) -> Result, JsError> { .collect() } +fn slab_from_js( + start: &[f64], + count: &[f64], + stride: Option>, + block: Option>, +) -> Result { + 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(), @@ -111,6 +157,72 @@ fn attr_to_js(value: AttrValue) -> (JsValue, Option) { } } +fn list_to_js(children: Vec) -> 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::() + .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) -> 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) -> 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 { @@ -146,38 +258,13 @@ impl H5File { /// The group's members: `[{ name, kind }]`, groups first. pub fn list(&self, path: &str) -> Result { - 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 { - 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::() - .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` @@ -186,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 { - 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 { - 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 { - 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` @@ -224,22 +301,315 @@ impl H5File { stride: Option>, block: Option>, ) -> Result { - 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 { - 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; + + #[wasm_bindgen(catch, js_name = fetchRanges)] + async fn fetch_ranges( + url: &str, + ranges: Vec, + opts: &JsValue, + validator: &JsValue, + length: f64, + ) -> Result; +} + +/// 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, + 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]) -> Result<(), JsError> { + let flat: Vec = 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( + &self, + storage: &LazyStorage, + mut f: impl FnMut() -> T, + ) -> Result { + let _op = storage.operation(); + loop { + match storage.attempt(&mut f) { + Step::Done(v) => return Ok(v), + Step::Need(ranges) => self.fetch(storage, &ranges).await?, + } + } + } +} + +impl Source { + async fn run(&self, op: impl Fn(&Reader) -> core::Result) -> Result { + 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, 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}" + ))), + } +} + +fn config_from(opts: &JsValue) -> Result { + 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; + } + 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); +/// - `fallback` — `"download"` (default) reads the whole file when the +/// server ignores `Range` (answers 200), up to `maxDownload` bytes +/// (default 512 MiB); `"error"` refuses such a server; +/// - `headers`, `credentials` — passed to every `fetch`; +/// - `parallel` — range requests in flight at once (default 6); +/// - `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`). +#[wasm_bindgen(js_name = openUrl)] +pub async fn open_url(url: String, opts: JsValue) -> Result { + 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")))?; + 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, +} + +impl RemoteFile { + /// Run `op` (fetching what it needs) and convert its result. + fn call( + &self, + op: impl Fn(&Reader) -> core::Result + '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")] + 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>")] + 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")] + 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>")] + 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>")] + 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")] + 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")] + pub fn read_hyperslab( + &self, + path: String, + start: Vec, + count: Vec, + stride: Option>, + block: Option>, + ) -> Result { + 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 } } diff --git a/examples/wasm-viewer/test/make_fixture.py b/examples/wasm-viewer/test/make_fixture.py index 8c7722c..1e79065 100644 --- a/examples/wasm-viewer/test/make_fixture.py +++ b/examples/wasm-viewer/test/make_fixture.py @@ -14,6 +14,7 @@ encoded as strings so JSON.parse keeps 64-bit values exact. """ import json +import os import sys import warnings from pathlib import Path @@ -201,3 +202,39 @@ def slab_for(obj): json.dump({"fixture.h5": describe(h5), "fixture.nc": describe(nc)}, open(out / "expected.json", "w"), indent=1, ensure_ascii=False) + + +def write_big(path, megabytes): + """A large file for the range-request tests (`openUrl`): `/big`, about + `megabytes` MB of float64 in 1 MiB chunks, written after a small + dataset and a group, so listing and reading `/small` touch a few blocks + of the file and a window of `/big` one chunk. Returns what h5py reads + back.""" + n = megabytes * 1_000_000 // 8 + chunk = 1 << 17 + with h5py.File(path, "w") as f: + f.attrs["note"] = "large file for range reads" + f.create_dataset("small", data=np.array([1.5, -2.0, 3.25])) + g = f.create_group("meta") + g.attrs["units"] = "m" + g.create_dataset("ids", data=np.arange(10, dtype=" 0: + json.dump(write_big(out / "big.h5", big_mb), open(out / "big.json", "w"), indent=1) diff --git a/examples/wasm-viewer/test/run.sh b/examples/wasm-viewer/test/run.sh index c1cdae0..68c4677 100755 --- a/examples/wasm-viewer/test/run.sh +++ b/examples/wasm-viewer/test/run.sh @@ -1,11 +1,17 @@ #!/usr/bin/env bash # Build the wasm package (../build.sh), test it under Node against files -# written by h5py and netCDF4 (make_fixture.py), then load the viewer page -# in headless Chromium if one is found (browser.sh). +# written by h5py and netCDF4 (make_fixture.py) — opened from bytes, and +# opened by URL from a local range-capable HTTP server (serve.py) — then +# load the viewer page in headless Chromium if one is found (browser.sh). # # Needs node, the wasm-bindgen CLI (see ../build.sh) and a Python with h5py, # netCDF4 and numpy: CLAWHDF5_PYTHON names it (default python3). Without that # Python the test is skipped, unless CLAWHDF5_REQUIRE_INTEROP=1. +# +# WASM_BIG_MB (default 200) sizes the large file of the range-request +# budget test (0 leaves it out); it is written under TMPDIR. +# CLAWHDF5_WASM_CORPUS=DIR also compares every HDF5 file under DIR (up to +# 16 MiB) read over HTTP with the same file read from bytes. set -euo pipefail HERE="$(cd "$(dirname "$0")" && pwd)" @@ -23,9 +29,24 @@ fi bash "$HERE/../build.sh" fix="$(mktemp -d)" -trap 'rm -rf "$fix"' EXIT -"$PY" "$HERE/make_fixture.py" "$fix" -node "$HERE/test.mjs" "$HERE/../pkg" "$fix" +server="" +cleanup() { + [ -n "$server" ] && kill "$server" 2>/dev/null || true + rm -rf "$fix" +} +trap cleanup EXIT +WASM_BIG_MB="${WASM_BIG_MB:-200}" "$PY" "$HERE/make_fixture.py" "$fix" + +# The fixtures (and the corpus) over HTTP with range support. +roots=(--root "fix=$fix") +[ -n "${CLAWHDF5_WASM_CORPUS:-}" ] && roots+=(--root "corpus=$CLAWHDF5_WASM_CORPUS") +"$PY" "$HERE/serve.py" "${roots[@]}" > "$fix/port" & +server=$! +for _ in $(seq 50); do + [ -s "$fix/port" ] && break + sleep 0.1 +done +node "$HERE/test.mjs" "$HERE/../pkg" "$fix" "http://127.0.0.1:$(head -1 "$fix/port")" # The page itself, in headless Chromium when one is available. status=0 diff --git a/examples/wasm-viewer/test/serve.py b/examples/wasm-viewer/test/serve.py new file mode 100644 index 0000000..f4a6e09 --- /dev/null +++ b/examples/wasm-viewer/test/serve.py @@ -0,0 +1,201 @@ +"""A static HTTP server for the wasm tests, with HTTP Range support and +request counting. + + python serve.py [--root PREFIX=DIR ...] + +Serves each DIR under URL PREFIX (the first match wins; PREFIX "" is the +site root), prints the port on its first line of stdout, and runs until +killed. No symlinks or copies are made: files are read where they are. + +- `Range: bytes=a-b`, `bytes=a-` and `bytes=-n` get 206 with Content-Range, + an unsatisfiable range 416; every file answer carries an ETag, and CORS + headers exposing Content-Range, so a page on another origin can use it. +- Under `/norange/...` the same files are served but Range is ignored + (200 with the whole file), as by a server without range support. +- `GET /__stats` returns `{"requests": n, "bytes": n, "log": [...]}` for + file requests since the last `GET /__reset`, which zeroes them. +""" + +import argparse +import hashlib +import json +import os +import posixpath +import sys +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from urllib.parse import unquote, urlsplit + +TYPES = { + ".html": "text/html; charset=utf-8", + ".js": "text/javascript; charset=utf-8", + ".mjs": "text/javascript; charset=utf-8", + ".wasm": "application/wasm", + ".json": "application/json", + ".ts": "text/plain; charset=utf-8", +} + +lock = threading.Lock() +stats = {"requests": 0, "bytes": 0, "log": []} + + +def resolve(roots, path): + """The file for URL `path`, or None. `..` never leaves a root.""" + parts = [p for p in posixpath.normpath(unquote(path)).split("/") if p] + if any(p in (".", "..") for p in parts): + return None + for prefix, root in roots: + pre = [p for p in prefix.split("/") if p] + if parts[: len(pre)] == pre: + rest = parts[len(pre):] or ["index.html"] + f = os.path.join(root, *rest) + if os.path.isfile(f): + return f + return None + + +def parse_range(header, size): + """(start, end exclusive) for a single `bytes=` range, "bad" when + unsatisfiable, None when absent or unparsable (served whole).""" + if not header or not header.startswith("bytes=") or "," in header: + return None + a, _, b = header[len("bytes="):].strip().partition("-") + try: + if a == "": + n = int(b) + return (max(0, size - n), size) if n > 0 and size > 0 else "bad" + start = int(a) + end = int(b) + 1 if b else size + except ValueError: + return None + if start >= size or end <= start: + return "bad" + return start, min(end, size) + + +def make_handler(roots): + class Handler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def log_message(self, *args): + pass + + def cors(self): + self.send_header("Access-Control-Allow-Origin", "*") + self.send_header("Access-Control-Expose-Headers", + "Content-Range, Content-Length, ETag, Accept-Ranges") + + def do_OPTIONS(self): + self.send_response(204) + self.cors() + self.send_header("Access-Control-Allow-Headers", "Range") + self.send_header("Content-Length", "0") + self.end_headers() + + def do_HEAD(self): + self.serve(head=True) + + def do_GET(self): + self.serve(head=False) + + def json(self, obj): + body = json.dumps(obj).encode() + self.send_response(200) + self.cors() + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.send_header("Cache-Control", "no-store") + self.end_headers() + self.wfile.write(body) + + def serve(self, head): + path = urlsplit(self.path).path + if path == "/__stats": + with lock: + return self.json(stats) + if path == "/__reset": + with lock: + stats.update(requests=0, bytes=0, log=[]) + return self.json({}) + ranges = True + if path.startswith("/norange/"): + ranges = False + path = path[len("/norange"):] + f = resolve(roots, path) + if f is None: + self.send_response(404) + self.cors() + self.send_header("Content-Length", "0") + self.end_headers() + return + size = os.path.getsize(f) + st = os.stat(f) + etag = '"%s"' % hashlib.sha1( + f"{f}:{size}:{st.st_mtime_ns}".encode()).hexdigest()[:16] + r = parse_range(self.headers.get("Range"), size) if ranges else None + if r == "bad": + self.send_response(416) + self.cors() + self.send_header("Content-Range", f"bytes */{size}") + self.send_header("Content-Length", "0") + self.end_headers() + return + start, end = r if r else (0, size) + self.send_response(206 if r else 200) + self.cors() + ext = os.path.splitext(f)[1] + self.send_header("Content-Type", TYPES.get(ext, "application/octet-stream")) + self.send_header("Content-Length", str(end - start)) + self.send_header("ETag", etag) + self.send_header("Cache-Control", "no-store") + if ranges: + self.send_header("Accept-Ranges", "bytes") + if r: + self.send_header("Content-Range", f"bytes {start}-{end - 1}/{size}") + self.end_headers() + if not head: + with lock: + stats["requests"] += 1 + stats["bytes"] += end - start + stats["log"].append([path, start, end, 206 if r else 200]) + with open(f, "rb") as fh: + fh.seek(start) + left = end - start + try: + while left: + buf = fh.read(min(left, 1 << 20)) + if not buf: + break + self.wfile.write(buf) + left -= len(buf) + except (BrokenPipeError, ConnectionResetError): + pass + + return Handler + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--root", action="append", default=[], + help="PREFIX=DIR: serve DIR under URL PREFIX") + args = ap.parse_args() + roots = [] + for spec in args.root: + prefix, _, d = spec.partition("=") + roots.append((prefix, os.path.abspath(d))) + class Server(ThreadingHTTPServer): + def handle_error(self, request, client_address): + # A client that drops a connection (a cancelled download) is + # not an error of the server. + if not isinstance(sys.exc_info()[1], (ConnectionError, TimeoutError)): + super().handle_error(request, client_address) + + httpd = Server(("127.0.0.1", 0), make_handler(roots)) + httpd.daemon_threads = True + print(httpd.server_address[1], flush=True) + sys.stdout.close() + httpd.serve_forever() + + +if __name__ == "__main__": + main() diff --git a/examples/wasm-viewer/test/test.mjs b/examples/wasm-viewer/test/test.mjs index 80ca30a..7a264dc 100644 --- a/examples/wasm-viewer/test/test.mjs +++ b/examples/wasm-viewer/test/test.mjs @@ -1,20 +1,31 @@ // Node test of the built wasm package (the exact pkg/ the viewer page loads) // and the viewer's DOM-free helpers. Run by test/run.sh: -// node test.mjs PKG_DIR FIXTURE_DIR +// node test.mjs PKG_DIR FIXTURE_DIR [SERVER_URL] // FIXTURE_DIR holds fixture.h5, fixture.nc and expected.json from -// make_fixture.py (values as libhdf5 reads them back). +// make_fixture.py (values as libhdf5 reads them back), and big.h5/big.json +// when it was run with WASM_BIG_MB. SERVER_URL is test/serve.py serving +// FIXTURE_DIR under /fix (and $CLAWHDF5_WASM_CORPUS under /corpus): with +// it, every check is repeated on files opened with openUrl (HTTP range +// requests), and the request budget, the full-download fallback and the +// error paths of openUrl are tested. import assert from "node:assert/strict"; -import { readFileSync } from "node:fs"; -import { join } from "node:path"; +import { existsSync, readFileSync, readdirSync, statSync } from "node:fs"; +import { join, relative } from "node:path"; import { pathToFileURL } from "node:url"; -const [pkgDir, fixDir] = process.argv.slice(2); +const [pkgDir, fixDir, base] = process.argv.slice(2); const pkg = await import(pathToFileURL(join(pkgDir, "clawhdf5_wasm.js"))); pkg.initSync({ module: readFileSync(join(pkgDir, "clawhdf5_wasm_bg.wasm")) }); const lib = await import(pathToFileURL(join(import.meta.dirname, "..", "viewer-lib.js"))); let checks = 0; const eq = (a, b, msg) => { assert.deepEqual(a, b, msg); checks++; }; +// A call that must fail: a thrown Error (open) or a rejected promise +// (openUrl), whose message matches `re`. +const fails = async (fn, re, msg) => { + await assert.rejects(async () => fn(), (e) => e instanceof Error && re.test(e.message), msg); + checks++; +}; const ARRAY_TYPES = { f32: Float32Array, f64: Float64Array, i8: Int8Array, i16: Int16Array, i32: Int32Array, @@ -54,13 +65,13 @@ function checkAttr(ctx, a, want) { assert.fail(`${ctx}: unknown expectation ${JSON.stringify(want)}`); } -const expected = JSON.parse(readFileSync(join(fixDir, "expected.json"), "utf8")); -for (const [name, exp] of Object.entries(expected)) { - const file = pkg.open(new Uint8Array(readFileSync(join(fixDir, name)))); - +// Every listing, dataset, error and attribute of `file` against what +// libhdf5 reads (expected.json). `file` is an H5File (synchronous methods) +// or a RemoteFile (promises): every call is awaited. +async function checkFile(name, exp, file) { for (const [path, want] of Object.entries(exp.lists)) { - eq(file.kind(path), "group", `${name}:${path} kind`); - const list = file.list(path); + eq(await file.kind(path), "group", `${name}:${path} kind`); + const list = await file.list(path); for (const [kind, key] of [["group", "groups"], ["dataset", "datasets"]]) { eq(list.filter((c) => c.kind === kind).map((c) => c.name).sort(), want[key], `${name}:${path} ${key}`); } @@ -68,58 +79,67 @@ for (const [name, exp] of Object.entries(expected)) { for (const [path, want] of Object.entries(exp.datasets)) { const ctx = `${name}:${path}`; - eq(file.kind(path), "dataset", `${ctx} kind`); + eq(await file.kind(path), "dataset", `${ctx} kind`); if (want.unavailable) { // The wasm build has no zstd (it links C): a clear error, no data. - assert.throws(() => file.read(path), (e) => e.message.includes(want.unavailable), ctx); - checks++; + await fails(() => file.read(path), new RegExp(want.unavailable), ctx); continue; } - const info = file.info(path); + const info = await file.info(path); eq([...info.shape, ...info.elementShape], want.shape, `${ctx} info shape`); - const r = file.read(path); + const r = await file.read(path); eq(r.shape, want.shape, `${ctx} shape`); eq(r.dtype, info.dtype, `${ctx} dtype`); assert.ok(r.data instanceof ARRAY_TYPES[want.kind], `${ctx}: ${r.data.constructor.name} for ${want.kind}`); eq(values(want.kind, r.data), want.values, ctx); if (want.slab) { const s = want.slab; - const part = file.readHyperslab(path, s.start, s.count, s.stride); + const part = await file.readHyperslab(path, s.start, s.count, s.stride); eq(part.shape, s.shape, `${ctx} slab shape`); eq(values(want.kind, part.data), s.values, `${ctx} slab`); } } for (const [path, what] of Object.entries(exp.errors)) { - assert.throws(() => file.read(path), (e) => e instanceof Error && e.message.includes(what), `${name}:${path}`); - checks++; + await fails(() => file.read(path), new RegExp(what), `${name}:${path}`); } for (const [path, want] of Object.entries(exp.attrs)) { - const attrs = file.attrs(path); - eq(file.attrErrors(path), [], `${name}:${path} attr errors`); + const attrs = await file.attrs(path); + eq(await file.attrErrors(path), [], `${name}:${path} attr errors`); const seen = attrs.filter((a) => !a.name.startsWith("_") && !exp.skip_attrs.includes(a.name)); eq(seen.map((a) => a.name).sort(), Object.keys(want).sort(), `${name}:${path} attr names`); for (const a of seen) checkAttr(`${name}:${path}@${a.name}`, a, want[a.name]); } +} + +// The error paths of the reader, for an H5File or a RemoteFile of +// fixture.h5. +async function checkErrors(h5) { + await fails(() => h5.read("/nope"), /./); + await fails(() => h5.list("/grid"), /not a group/); + await fails(() => h5.readHyperslab("/grid", [0], [1]), /dimensions/); + await fails(() => h5.readHyperslab("/grid", [5, 0], [2, 1]), /exceeds/); + await fails(() => h5.readHyperslab("/grid", [-1, 0], [1, 1]), /non-negative integers/); + await fails(() => h5.readHyperslab("/grid", [0.5, 0], [1, 1]), /non-negative integers/); + // Big integers stay exact. + eq((await h5.read("/u64")).data[0], 18446744073709551615n, "u64 max"); +} + +const expected = JSON.parse(readFileSync(join(fixDir, "expected.json"), "utf8")); +for (const [name, exp] of Object.entries(expected)) { + const file = pkg.open(new Uint8Array(readFileSync(join(fixDir, name)))); + await checkFile(name, exp, file); file.free(); } // Errors reach JavaScript as thrown Errors, never as data. const h5 = pkg.open(new Uint8Array(readFileSync(join(fixDir, "fixture.h5")))); -const throwsMsg = (fn, re) => { assert.throws(fn, (e) => e instanceof Error && re.test(e.message)); checks++; }; -throwsMsg(() => pkg.open(new Uint8Array(64)), /./); -throwsMsg(() => h5.read("/nope"), /./); -throwsMsg(() => h5.list("/grid"), /not a group/); -throwsMsg(() => h5.readHyperslab("/grid", [0], [1]), /dimensions/); -throwsMsg(() => h5.readHyperslab("/grid", [5, 0], [2, 1]), /exceeds/); -throwsMsg(() => h5.readHyperslab("/grid", [-1, 0], [1, 1]), /non-negative integers/); -throwsMsg(() => h5.readHyperslab("/grid", [0.5, 0], [1, 1]), /non-negative integers/); +await fails(() => pkg.open(new Uint8Array(64)), /./); +await checkErrors(h5); // Info for a dataset with an unlimited dimension (netCDF "time"). const nc = pkg.open(new Uint8Array(readFileSync(join(fixDir, "fixture.nc")))); eq(nc.info("/time").maxshape, [null], "unlimited dimension is null"); -// Big integers stay exact. -eq(h5.read("/u64").data[0], 18446744073709551615n, "u64 max"); eq(typeof pkg.version(), "string", "version"); // Viewer helpers. @@ -137,7 +157,261 @@ eq(lib.toRows(pairs.data, 2, 1, lib.perElement([2])), [["[2, 3]"], ["[4, 5]"]], eq(lib.formatValue(0.1 + 0.2), "0.3", "float formatting"); eq(lib.formatValue(2n ** 64n - 1n), "18446744073709551615", "bigint formatting"); eq(lib.formatValue("x"), '"x"', "string formatting"); +eq(lib.formatBytes(0), "0 B", "bytes"); +eq(lib.formatBytes(1536), "1.5 KiB", "KiB"); +eq(lib.formatBytes(200 * 1024 * 1024), "200 MiB", "MiB"); +eq(lib.formatStats({ lazy: true, requests: 3, bytesFetched: 2 << 20, size: 200 << 20 }), + "3 requests, 2 MiB of 200 MiB fetched (1.0%)", "stats line"); +eq(lib.formatStats({ lazy: false, requests: 1, bytesFetched: 1024, size: 1024 }), + "downloaded whole (1 KiB): the server does not support range requests", "stats line, no ranges"); h5.free(); nc.free(); - console.log(`wasm package: ${checks} checks passed`); + +if (base) await remoteTests(); + +async function serverStats() { + return (await fetch(`${base}/__stats`)).json(); +} + +async function remoteTests() { + const before = checks; + await fetch(`${base}/__reset`); + + // The same checks over HTTP range requests, at the default block size + // and at 512-byte blocks with a 4 KiB cache (almost every structure read + // a miss, evictions between calls). + for (const opts of [undefined, { blockSize: 512, cacheSize: 4096 }]) { + for (const [name, exp] of Object.entries(expected)) { + const f = await pkg.openUrl(`${base}/fix/${name}`, opts); + await checkFile(`${name} (openUrl ${JSON.stringify(opts ?? {})})`, exp, f); + const st = f.stats(); + eq(st.lazy, true, "read by ranges"); + eq(st.size, statSync(join(fixDir, name)).size, "size"); + f.free(); + } + } + const remote = await pkg.openUrl(`${base}/fix/fixture.h5`); + await checkErrors(remote); + eq((await (await pkg.openUrl(`${base}/fix/fixture.nc`)).info("/time")).maxshape, [null], "remote unlimited"); + + // What the page counts is what the server served. + await fetch(`${base}/__reset`); + const counted = await pkg.openUrl(`${base}/fix/fixture.nc`, { blockSize: 1024 }); + await counted.read("/temp"); + const server = await serverStats(); + eq(counted.stats().requests, server.requests, "requests counted"); + eq(counted.stats().bytesFetched, server.bytes, "bytes counted"); + + // A custom fetch is used for every request, with the caller's headers. + let calls = 0; + const seen = new Set(); + const viaCustom = await pkg.openUrl(`${base}/fix/fixture.h5`, { + blockSize: 4096, + headers: { "X-Test": "1" }, + fetch: (url, init) => { calls++; seen.add(init.headers["X-Test"]); return fetch(url, init); }, + }); + eq((await viaCustom.read("/sensors/temp")).data[0], 21.5, "custom fetch values"); + eq(calls, viaCustom.stats().requests, "custom fetch calls"); + eq([...seen], ["1"], "headers passed"); + + // Calls in flight at once share the cache (a 1 KiB budget: nothing is + // evicted while any of them runs) and each gets its own answer. + const both = await pkg.openUrl(`${base}/fix/fixture.h5`, { blockSize: 512, cacheSize: 1024 }); + const [g, t, l, a] = await Promise.all([ + both.read("/grid"), both.read("/sensors/temp"), both.list("/sensors"), both.attrs("/"), + ]); + const exp = expected["fixture.h5"]; + eq(Array.from(g.data), exp.datasets["/grid"].values, "concurrent /grid"); + eq(Array.from(t.data), exp.datasets["/sensors/temp"].values, "concurrent /sensors/temp"); + eq(l.map((c) => c.name).sort(), [...exp.lists["/sensors"].groups, ...exp.lists["/sensors"].datasets].sort(), "concurrent list"); + eq(a.length > 0, true, "concurrent attrs"); + assert.ok(both.stats().cachedBytes <= 1024, "trimmed to the budget when idle"); + + // Listing and reading small things of a large file fetches a few blocks, + // not the file. + if (existsSync(join(fixDir, "big.json"))) { + const big = JSON.parse(readFileSync(join(fixDir, "big.json"), "utf8")); + await fetch(`${base}/__reset`); + const f = await pkg.openUrl(`${base}/fix/big.h5`); + const list = await f.list("/"); + eq(list.filter((c) => c.kind === "group").map((c) => c.name), big.list.groups, "big: groups"); + eq(list.filter((c) => c.kind === "dataset").map((c) => c.name).sort(), big.list.datasets, "big: datasets"); + eq(Array.from((await f.read("/small")).data), big.small, "big: /small"); + eq(Array.from((await f.read("/meta/ids")).data, String), big.ids, "big: /meta/ids"); + eq((await f.info("/big")).shape, big.big_shape, "big: shape"); + eq((await f.attrs("/meta"))[0].value, "m", "big: group attribute"); + const win = await f.readHyperslab("/big", [big.window.start], [big.window.count]); + eq(Array.from(win.data), big.window.values, "big: window"); + const st = f.stats(); + const server = await serverStats(); + eq(st.requests, server.requests, "big: requests counted"); + eq(st.bytesFetched, server.bytes, "big: bytes counted"); + console.log(`big.h5 (${big.size} bytes): listed, 3 small reads and a window in ` + + `${server.requests} requests, ${server.bytes} bytes (${(100 * server.bytes / big.size).toFixed(2)}%)`); + assert.ok(server.requests <= 8, `big: ${server.requests} requests`); + assert.ok(server.bytes * 20 < big.size, `big: ${server.bytes} bytes fetched`); + checks += 2; + } + + // A server without range support: downloaded whole (the default), or + // refused. + const whole = await pkg.openUrl(`${base}/norange/fix/fixture.h5`); + await checkFile("fixture.h5 (no range support)", expected["fixture.h5"], whole); + eq(whole.stats().lazy, false, "downloaded whole"); + eq(whole.stats().requests, 1, "one request"); + await fails(() => pkg.openUrl(`${base}/norange/fix/fixture.h5`, { fallback: "error" }), + /does not support HTTP range requests/, "fallback: error"); + await fails(() => pkg.openUrl(`${base}/norange/fix/fixture.h5`, { maxDownload: 1000 }), + /more than maxDownload/, "maxDownload"); + // Without a Content-Length the body is streamed, and stopped at the limit. + const undeclared = async (url, init) => { + const r = await fetch(url, init); + return new Response(r.body, { status: r.status }); + }; + await fails(() => pkg.openUrl(`${base}/norange/fix/fixture.h5`, { maxDownload: 1000, fetch: undeclared }), + /more than maxDownload/, "maxDownload, streamed"); + const streamed = await pkg.openUrl(`${base}/norange/fix/fixture.h5`, { fetch: undeclared }); + eq(Array.from((await streamed.read("/sensors/temp")).data), [21.5, 22, 22.25], "streamed download"); + + // Errors: HTTP status, not HDF5, bad options, a file that changes, a + // server that answers with the wrong bytes. + await fails(() => pkg.openUrl(`${base}/fix/missing.h5`), /HTTP 404/, "404"); + await fails(() => pkg.openUrl(`${base}/fix/expected.json`), /./, "not HDF5"); + await fails(() => pkg.openUrl(`${base}/fix/fixture.h5`, { blockSize: 100 }), /blockSize/, "blockSize"); + const tamper = (edit) => async (url, init) => { + const r = await fetch(url, init); + return init.headers.Range === "bytes=0-511" ? r : edit(r); + }; + const withHeaders = async (r, headers) => { + const h = new Headers(r.headers); + for (const [k, v] of Object.entries(headers)) h.set(k, v); + return new Response(await r.arrayBuffer(), { status: r.status, headers: h }); + }; + await fails(async () => { + const f = await pkg.openUrl(`${base}/fix/fixture.h5`, { blockSize: 512, fetch: tamper((r) => withHeaders(r, { ETag: '"other"' })) }); + await f.read("/grid"); + }, /changed on the server/, "changed file"); + await fails(async () => { + const f = await pkg.openUrl(`${base}/fix/fixture.h5`, { + blockSize: 512, + fetch: tamper(async (r) => new Response((await r.arrayBuffer()).slice(1), { status: 206 })), + }); + await f.read("/grid"); + }, /got \d+/, "short answer"); + await fails(async () => { + const f = await pkg.openUrl(`${base}/fix/fixture.h5`, { + blockSize: 512, + fetch: tamper(async (r) => new Response(await r.arrayBuffer(), { status: 200 })), + }); + await f.read("/grid"); + }, /stopped honouring range requests/, "200 mid-file"); + await fails(async () => { + const f = await pkg.openUrl(`${base}/fix/fixture.h5`, { + blockSize: 512, + fetch: tamper(async (r) => withHeaders(r, { "Content-Range": "bytes 0-511/25752" })), + }); + await f.read("/grid"); + }, /the server sent 0-511/, "wrong range"); + await fails(async () => { + const f = await pkg.openUrl(`${base}/fix/fixture.h5`, { + blockSize: 512, + fetch: tamper(async (r) => withHeaders(r, { "Content-Range": "bytes */25752" })), + }); + await f.read("/grid"); + }, /unusable Content-Range/, "unusable Content-Range"); + + // Corpus files: what the viewer can show of each is the same read by + // ranges as in memory (an error wherever it gives one). + const corpus = process.env.CLAWHDF5_WASM_CORPUS; + if (corpus) await corpusTests(corpus); + console.log(`openUrl: ${checks - before} checks passed`); +} + +function hdf5Files(dir, out) { + for (const e of readdirSync(dir, { withFileTypes: true })) { + const p = join(dir, e.name); + let st; + try { + st = statSync(p); + } catch { + continue; + } + if (st.isDirectory()) hdf5Files(p, out); + else if (st.size <= 16 << 20) { + const b = readFileSync(p); + const sig = [0x89, 0x48, 0x44, 0x46, 0x0d, 0x0a, 0x1a, 0x0a]; + for (let at = 0; at + 8 <= b.length; at = at === 0 ? 512 : at * 2) { + if (sig.every((x, i) => b[at + i] === x)) { out.push(p); break; } + } + } + } + return out; +} + +function show(v) { + return JSON.stringify(v, (_, x) => { + if (typeof x === "bigint") return `${x}n`; + if (ArrayBuffer.isView(x)) return Array.from(x, (y) => (typeof y === "bigint" ? `${y}n` : Number.isNaN(y) ? "NaN" : y)); + return x; + }); +} + +async function transcript(f) { + const out = []; + const call = async (what, fn) => { + try { + out.push(`${what}: ${show(await fn())}`); + } catch { + out.push(`${what}: Err`); + } + }; + const todo = [["/", 0]]; + while (todo.length && out.length < 3000) { + const [path, depth] = todo.pop(); + let kind = null; + await call(`${path} kind`, async () => (kind = await f.kind(path))); + await call(`${path} attrs`, () => f.attrs(path)); + if (kind === "group") { + let list = []; + await call(`${path} list`, async () => (list = await f.list(path))); + if (depth < 12) for (const c of list.reverse()) todo.push([lib.joinPath(path, c.name), depth + 1]); + } else if (kind === "dataset") { + let info = null; + await call(`${path} info`, async () => (info = await f.info(path))); + if (info && [...info.shape, ...info.elementShape].reduce((a, b) => a * b, 1) <= 1 << 20) { + await call(`${path} read`, () => f.read(path)); + } + } + } + return out; +} + +async function corpusTests(corpus) { + const files = hdf5Files(corpus, []).sort(); + let opened = 0; + let bytes = 0; + let size = 0; + for (const p of files) { + const rel = relative(corpus, p).split("/").map(encodeURIComponent).join("/"); + let local; + try { + local = pkg.open(new Uint8Array(readFileSync(p))); + } catch { + await fails(() => pkg.openUrl(`${base}/corpus/${rel}`, { blockSize: 65536 }), /./, `${rel}: opens in neither`); + continue; + } + const remote = await pkg.openUrl(`${base}/corpus/${rel}`, { blockSize: 65536 }); + const want = await transcript(local); + const got = await transcript(remote); + // Errors are compared as errors: a malformed file can fail at another + // check, with another message, when read by ranges. + eq(got, want, `${rel}: transcript`); + opened++; + bytes += remote.stats().bytesFetched; + size += remote.stats().size; + local.free(); + remote.free(); + } + console.log(`corpus: ${opened} of ${files.length} files agree over HTTP (${bytes} of ${size} bytes fetched)`); +} diff --git a/examples/wasm-viewer/viewer-lib.js b/examples/wasm-viewer/viewer-lib.js index 72c5bf3..fece1ab 100644 --- a/examples/wasm-viewer/viewer-lib.js +++ b/examples/wasm-viewer/viewer-lib.js @@ -75,3 +75,24 @@ export function toRows(data, rows, cols, per = 1) { export function perElement(elementShape) { return elementShape.reduce((a, b) => a * b, 1); } + +/** A byte count for people: `1.5 KiB`, `200 MiB`. */ +export function formatBytes(n) { + const units = ["B", "KiB", "MiB", "GiB", "TiB"]; + let i = 0; + let v = n; + while (v >= 1024 && i < units.length - 1) { + v /= 1024; + i++; + } + const s = i === 0 ? String(v) : v < 10 ? String(Number(v.toFixed(1))) : String(Math.round(v)); + return `${s} ${units[i]}`; +} + +/** What reading a remote file has cost, from `RemoteFile.stats()`. */ +export function formatStats({ lazy, requests, bytesFetched, size }) { + if (!lazy) return `downloaded whole (${formatBytes(size)}): the server does not support range requests`; + const pct = size > 0 ? (100 * bytesFetched) / size : 0; + return `${requests} request${requests === 1 ? "" : "s"}, ${formatBytes(bytesFetched)} of ` + + `${formatBytes(size)} fetched (${pct.toFixed(1)}%)`; +} From b8c85f7627a612c3e376fac76ffad1f828e746bd Mon Sep 17 00:00:00 2001 From: osobh Date: Sun, 27 Sep 2026 06:44:51 -0500 Subject: [PATCH 03/10] wasm-viewer: open by URL, lazily, with a request counter The page gets a URL box; ?file= now opens with openUrl instead of downloading the file, and the header shows the requests made and bytes fetched so far ("5 requests, 5 MiB of 191 MiB fetched (2.6%)"). Every file call is awaited, so local files (open) and remote ones share the code; a selection that finishes after another was made is not shown. browser.sh serves the page and fixtures with test/serve.py (no symlinks into the repository any more) and also checks a server without range support and, on the 200 MB file, that a small dataset and a window of the big one render with a single-digit percentage fetched. Co-Authored-By: Claude Opus 5.5 (1M context) --- examples/wasm-viewer/index.html | 138 ++++++++++++++++++--------- examples/wasm-viewer/test/browser.sh | 76 +++++++++------ 2 files changed, 140 insertions(+), 74 deletions(-) diff --git a/examples/wasm-viewer/index.html b/examples/wasm-viewer/index.html index d477e49..279c9b8 100644 --- a/examples/wasm-viewer/index.html +++ b/examples/wasm-viewer/index.html @@ -56,27 +56,39 @@ border: 1px solid var(--line); border-radius: 4px; } .error { color: var(--bad); font-family: var(--mono); white-space: pre-wrap; } .muted { color: var(--muted); } + form.url { display: flex; gap: 6px; margin-left: auto; flex: 1 1 320px; max-width: 560px; } + form.url input { flex: 1; min-width: 0; font: 13px var(--mono); padding: 4px 8px; background: var(--panel); + color: var(--ink); border: 1px solid var(--line); border-radius: 6px; } + .netstats { flex-basis: 100%; color: var(--muted); font-size: 12.5px; font-family: var(--mono); } + .netstats:empty { display: none; }

HDF5 Viewer

no file +
+ + +
+
-

Drop an HDF5 or NetCDF-4 file here, or use “Open file…”.

-

The file is read in this page by clawhdf5 compiled to WebAssembly; it is not uploaded anywhere.

+

Drop an HDF5 or NetCDF-4 file here, use “Open file…”, or give a URL.

+

The file is read in this page by clawhdf5 compiled to WebAssembly; a local file is not uploaded anywhere. + A file given by URL is not downloaded: only the byte ranges each view needs are fetched (HTTP range requests), + and the requests and bytes it has cost are shown at the top.