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() + ); +}