Lazy remote files in the browser (M4), SWMR reader (M5), Python remote reads and editing #19

Merged
osobh merged 33 commits from feat/p3-wasm-swmr-python into main 2026-09-27 14:47:26 +00:00
4 changed files with 1086 additions and 6 deletions
Showing only changes of commit 076feb089a - Show all commits
+22 -6
View File
@@ -6,9 +6,12 @@
//! with no typed-array mapping (compound, reference, opaque, ...) is refused //! with no typed-array mapping (compound, reference, opaque, ...) is refused
//! with a message naming it, never returned as reinterpreted bytes. //! with a message naming it, never returned as reinterpreted bytes.
use std::sync::Arc;
use clawhdf5::{AttrValue, File, Selection}; use clawhdf5::{AttrValue, File, Selection};
use clawhdf5_format::data_read; use clawhdf5_format::data_read;
use clawhdf5_format::datatype::{Datatype, DatatypeByteOrder}; use clawhdf5_format::datatype::{Datatype, DatatypeByteOrder};
use clawhdf5_format::storage::Storage;
use clawhdf5_format::vl_data::{VlResolver, check_element_size}; use clawhdf5_format::vl_data::{VlResolver, check_element_size};
/// Errors are reported to JavaScript as messages. /// Errors are reported to JavaScript as messages.
@@ -119,7 +122,9 @@ pub struct Attr {
pub value: AttrValue, 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 { pub struct Reader {
file: File, 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<dyn Storage + Send + Sync>) -> Result<Self> {
Ok(Self {
file: File::open_storage(storage).map_err(err)?,
})
}
/// Whether `path` names a group or a dataset. /// Whether `path` names a group or a dataset.
pub fn kind(&self, path: &str) -> Result<Kind> { pub fn kind(&self, path: &str) -> Result<Kind> {
match self.file.dataset(path) { match self.file.dataset(path) {
@@ -281,11 +294,14 @@ impl Reader {
// string ends at its first NUL and a heap object of the // string ends at its first NUL and a heap object of the
// wrong size is an error, as in libhdf5 and h5py. // wrong size is an error, as in libhdf5 and h5py.
let sb = self.file.superblock(); let sb = self.file.superblock();
Data::Strings( let strings = match self.file.contiguous_bytes() {
VlResolver::new(self.file.as_bytes(), sb.offset_size, sb.length_size) Some(bytes) => {
.strings(raw) VlResolver::new(bytes, sb.offset_size, sb.length_size).strings(raw)
.map_err(err)?, }
) 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 => { Datatype::Enumeration { .. } if !is_array => {
Data::Strings(data_read::read_enum_names(raw, dt).map_err(err)?) Data::Strings(data_read::read_enum_names(raw, dt).map_err(err)?)
+672
View File
@@ -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<T> {
/// The pass read only bytes that were present: its result stands.
Done(T),
/// The pass missed: fetch these byte ranges (sorted, disjoint, block
/// aligned), [`supply`](LazyStorage::supply) them and run it again.
Need(Vec<Range<u64>>),
}
struct Block {
data: Arc<[u8]>,
/// Eviction order: bulk blocks (`false`) before metadata (`true`),
/// then least recently used first.
key: (bool, u64),
}
#[derive(Default)]
struct State {
blocks: HashMap<u64, Block>,
/// `(metadata?, tick, block index)`, in eviction order.
order: BTreeSet<(bool, u64, u64)>,
tick: u64,
bytes: u64,
/// Blocks the current pass missed, and whether a small read wanted
/// them (metadata).
missing: HashMap<u64, bool>,
/// Blocks a bulk read missed that have not been supplied yet: kept
/// as bulk when they arrive.
bulk_pending: BTreeSet<u64>,
/// Operations in flight: no eviction while any is.
active: u32,
stats: LazyStats,
}
/// A [`Storage`] over the blocks of a file fetched so far; a read of
/// anything else fails and is recorded, so the pass can be re-run once the
/// bytes arrive. See the [module documentation](self).
pub struct LazyStorage {
len: u64,
config: LazyConfig,
state: Mutex<State>,
}
fn lock(m: &Mutex<State>) -> MutexGuard<'_, State> {
m.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
}
/// Keeps an operation's blocks cached until it is dropped; see
/// [`LazyStorage::operation`].
pub struct Operation<'a> {
storage: &'a LazyStorage,
}
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<T>(&self, f: impl FnOnce() -> T) -> Step<T> {
{
let mut st = lock(&self.state);
st.missing.clear();
st.stats.passes += 1;
}
let out = f();
let missing = std::mem::take(&mut lock(&self.state).missing);
if missing.is_empty() {
return Step::Done(out);
}
drop(out);
Step::Need(self.runs(missing))
}
/// The bytes of the file at `offset`, fetched for a range a pass asked
/// for. `offset` must be block aligned and the bytes whole blocks (the
/// last block of the file may be short) inside the file, or this is an
/// error and nothing is kept. Blocks already present are left alone.
pub fn supply(&self, offset: u64, bytes: &[u8]) -> Result<(), String> {
let bs = self.config.block_size;
let end = offset
.checked_add(bytes.len() as u64)
.filter(|&e| e <= self.len)
.ok_or_else(|| {
format!(
"{} bytes at offset {offset} run past the end of the {}-byte file",
bytes.len(),
self.len
)
})?;
if !offset.is_multiple_of(bs) || (!end.is_multiple_of(bs) && end != self.len) {
return Err(format!(
"{} bytes at offset {offset} are not whole {bs}-byte blocks",
bytes.len()
));
}
let mut st = lock(&self.state);
st.stats.requests += 1;
st.stats.bytes_fetched += bytes.len() as u64;
let mut start = offset;
while start < end {
let i = start / bs;
let stop = (start + bs).min(end);
if !st.blocks.contains_key(&i) {
let rel = (start - offset) as usize..(stop - offset) as usize;
let metadata = !st.bulk_pending.remove(&i);
self.keep(&mut st, i, Arc::from(&bytes[rel]), metadata);
}
start = stop;
}
if st.active == 0 {
self.evict(&mut st);
}
Ok(())
}
/// [`supply`](Self::supply) the bytes fetched for `range`, one of the
/// ranges a [`Step::Need`] asked for: anything but exactly its length
/// (a server that answered with more or less) is an error.
pub fn supply_range(&self, range: &Range<u64>, bytes: &[u8]) -> Result<(), String> {
let want = range.end.saturating_sub(range.start);
if bytes.len() as u64 != want {
return Err(format!(
"asked for {want} bytes at offset {}, got {}",
range.start,
bytes.len()
));
}
self.supply(range.start, bytes)
}
/// Run `f` to completion, fetching what its passes miss with `fetch`
/// (a byte range to its bytes). The blocking driver, for native code
/// and tests; the browser's is the same loop with an `await` between
/// passes.
pub fn run_blocking<T>(
&self,
mut f: impl FnMut() -> T,
mut fetch: impl FnMut(Range<u64>) -> Result<Vec<u8>, String>,
) -> Result<T, String> {
let _op = self.operation();
loop {
match self.attempt(&mut f) {
Step::Done(v) => return Ok(v),
Step::Need(ranges) => {
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<u64, bool>) -> Vec<Range<u64>> {
let bs = self.config.block_size;
let mut wanted: Vec<u64> = missing.keys().copied().collect();
wanted.sort_unstable();
{
// Remember which blocks only bulk reads asked for: they are
// kept as bulk once supplied.
let mut st = lock(&self.state);
for (&i, &metadata) in &missing {
if metadata {
st.bulk_pending.remove(&i);
} else {
st.bulk_pending.insert(i);
}
}
}
let 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<Range<u64>> {
let end = offset.saturating_add(len).min(self.len);
if offset >= end {
return None;
}
let bs = self.config.block_size;
Some(offset / bs..(end - 1) / bs + 1)
}
/// The blocks of `spans` if all are present (touching them), else
/// record the missing ones and fail.
fn blocks(
&self,
spans: &[Range<u64>],
metadata: bool,
) -> Result<HashMap<u64, Arc<[u8]>>, FormatError> {
let mut st = lock(&self.state);
let mut have = HashMap::new();
let mut missed = false;
for span in spans {
for i in span.clone() {
if have.contains_key(&i) {
continue;
}
match st.blocks.get(&i) {
Some(b) => {
have.insert(i, b.data.clone());
}
None => {
missed = true;
let m = st.missing.entry(i).or_insert(metadata);
*m |= metadata;
}
}
}
}
if missed {
return Err(FormatError::Storage(NEED_BYTES.into()));
}
// Touch: most recently used last; a small read promotes a bulk
// block to metadata.
for &i in have.keys() {
st.tick += 1;
let tick = st.tick;
let Some(b) = st.blocks.get_mut(&i) else {
continue;
};
let old = b.key;
b.key = (old.0 || metadata, tick);
let new = b.key;
st.order.remove(&(old.0, old.1, i));
st.order.insert((new.0, new.1, i));
}
Ok(have)
}
fn assemble(&self, offset: u64, end: u64, blocks: &HashMap<u64, Arc<[u8]>>) -> Vec<u8> {
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<Cow<'_, [u8]>, 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<u64>]) -> Result<Vec<Cow<'_, [u8]>>, 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<u8> {
(0..n).map(|i| (i * 7 + i / 251) as u8).collect()
}
fn config(block: u64, capacity: u64) -> LazyConfig {
LazyConfig {
block_size: block,
capacity,
max_request: 4 * block,
}
}
/// Supply every range of `need` from `data`.
fn serve(s: &LazyStorage, data: &[u8], need: &[Range<u64>]) {
for r in need {
s.supply_range(r, &data[r.start as usize..r.end as usize])
.unwrap();
}
}
fn owned(r: Result<Cow<'_, [u8]>, FormatError>) -> Result<Vec<u8>, FormatError> {
r.map(Cow::into_owned)
}
#[test]
fn a_miss_asks_for_whole_blocks_then_the_rerun_reads_them() {
let data = file(10_000);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let Step::Need(need) = s.attempt(|| owned(s.read_at(1500, 1000))) else {
panic!("nothing is cached yet");
};
assert_eq!(need, vec![1024..3072]);
serve(&s, &data, &need);
let Step::Done(got) = s.attempt(|| owned(s.read_at(1500, 1000))) else {
panic!("the blocks were supplied");
};
assert_eq!(got.unwrap(), &data[1500..2500]);
// Past the end: short, then empty, as for a slice.
serve(&s, &data, &[9216..10_000]);
let Step::Done(tail) = s.attempt(|| owned(s.read_at(9_990, 100))) else {
panic!("the last block was supplied");
};
assert_eq!(tail.unwrap(), &data[9_990..]);
assert!(matches!(
s.attempt(|| s.read_at(20_000, 10).map(|c| c.len())),
Step::Done(Ok(0))
));
let st = s.stats();
assert_eq!(
(st.passes, st.requests, st.bytes_fetched),
(4, 2, 2048 + 784)
);
}
#[test]
fn a_pass_that_swallowed_the_miss_is_still_rerun() {
// A parser that catches the error and returns something anyway
// (a listing skipping a link it cannot resolve) must not have its
// result used.
let data = file(4096);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let step = s.attempt(|| s.read_at(0, 4).map(|b| b.len()).unwrap_or(0));
assert!(
matches!(step, Step::Need(ref n) if n == &vec![0..1024]),
"{step:?}"
);
}
#[test]
fn read_ranges_asks_for_every_missing_block_in_one_pass() {
let data = file(64 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
let ranges = [
100..200,
5000..5100,
5200..5300,
30_000..33_000,
60_000..60_010,
];
let read = || {
s.read_ranges(&ranges)
.map(|v| v.into_iter().map(Cow::into_owned).collect::<Vec<_>>())
};
let Step::Need(need) = s.attempt(read) else {
panic!("nothing is cached yet");
};
// 5000..5300 is blocks 4 and 5; 30_000..33_000 is blocks 29..=32.
assert_eq!(
need,
vec![
0..1024,
4096..6144,
29 * 1024..33 * 1024,
58 * 1024..59 * 1024
]
);
serve(&s, &data, &need);
let Step::Done(got) = s.attempt(read) else {
panic!("one pass fetched everything");
};
for (r, g) in ranges.iter().zip(&got.unwrap()) {
assert_eq!(g, &data[r.start as usize..r.end as usize]);
}
}
#[test]
fn runs_merge_one_block_holes_and_split_long_runs() {
let data = file(32 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
// Blocks 0 and 2 (a hole of one: merged), 5..=14 (split in fours).
let Step::Need(need) = s.attempt(|| {
let _ = s.read_at(0, 10);
let _ = s.read_at(2048, 10);
s.read_ranges(&[5120..15 * 1024]).map(|_| ())
}) else {
panic!("nothing is cached yet");
};
assert_eq!(
need,
vec![0..3072, 5120..9216, 9216..13_312, 13_312..15_360]
);
}
#[test]
fn supply_refuses_what_was_not_asked_for() {
let data = file(10_000);
let s = LazyStorage::new(data.len() as u64, config(1024, 1 << 20));
assert!(s.supply(1, &data[1..1025]).unwrap_err().contains("whole"));
assert!(s.supply(0, &data[..1000]).unwrap_err().contains("whole"));
assert!(
s.supply(9216, &[0u8; 1024])
.unwrap_err()
.contains("past the end")
);
for wrong in [&data[..1023], &data[..2048]] {
assert!(
s.supply_range(&(0..1024), wrong)
.unwrap_err()
.contains("asked for 1024 bytes")
);
}
assert_eq!(s.stats().cached_bytes, 0);
// The file's last, short block is whole.
s.supply(9216, &data[9216..]).unwrap();
assert_eq!(s.stats().cached_bytes, 784);
}
#[test]
fn an_operation_keeps_its_blocks_whatever_the_budget() {
// A budget of one block, an operation that needs eight: without
// the operation guard every pass would evict what the last one
// fetched and never finish.
let data = file(8 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 1024));
let got = s
.run_blocking(
|| {
(0..8)
.map(|i| owned(s.read_at(i * 1024 + 3, 10)))
.collect::<Result<Vec<_>, _>>()
},
|r| Ok(data[r.start as usize..r.end as usize].to_vec()),
)
.unwrap()
.unwrap();
for (i, g) in got.iter().enumerate() {
assert_eq!(g, &data[i * 1024 + 3..i * 1024 + 13]);
}
// Trimmed to the budget once the operation is over.
let st = s.stats();
assert_eq!(st.cached_bytes, 1024);
assert_eq!(st.evictions, 7);
assert_eq!(st.passes, 9, "one pass per block, then the one that ends");
}
#[test]
fn bulk_blocks_are_evicted_before_metadata() {
let data = file(16 * 1024);
let s = LazyStorage::new(data.len() as u64, config(1024, 4 * 1024));
let fetch = |r: Range<u64>| Ok(data[r.start as usize..r.end as usize].to_vec());
// Metadata: a small read of block 0.
s.run_blocking(|| s.read_at(0, 16).map(|_| ()), fetch)
.unwrap()
.unwrap();
// Bulk: raw data over blocks 4..12, more than the budget.
s.run_blocking(|| s.read_ranges(&[4096..12 * 1024]).map(|_| ()), fetch)
.unwrap()
.unwrap();
// The metadata block survived: reading it again is a hit.
let before = s.stats();
assert!(before.cached_bytes <= 4 * 1024);
assert!(matches!(
s.attempt(|| s.read_at(0, 16).map(|_| ())),
Step::Done(Ok(()))
));
assert_eq!(s.stats().requests, before.requests);
}
#[test]
fn a_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);
}
}
+1
View File
@@ -21,6 +21,7 @@
//! The logic lives in [`core`], which is plain Rust and tested natively. //! The logic lives in [`core`], which is plain Rust and tested natively.
pub mod core; pub mod core;
pub mod lazy;
use clawhdf5::AttrValue; use clawhdf5::AttrValue;
use js_sys::{Array, Object, Reflect}; use js_sys::{Array, Object, Reflect};
+391
View File
@@ -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<T>(&self, op: impl Fn(&Reader) -> T) -> T;
}
struct Local(Reader);
impl Api for Local {
fn call<T>(&self, op: impl Fn(&Reader) -> T) -> T {
op(&self.0)
}
}
/// A lazily read file and the "server" it fetches from.
struct Lazy {
data: Arc<Vec<u8>>,
storage: Arc<LazyStorage>,
reader: Reader,
}
fn fetch(data: &[u8], r: Range<u64>) -> Result<Vec<u8>, String> {
Ok(data[r.start as usize..r.end as usize].to_vec())
}
impl Lazy {
/// Open as `openUrl` does: the first block comes with the probe that
/// learns the length, then the open is run until it has its bytes.
fn open(data: Vec<u8>, config: LazyConfig) -> Result<Lazy, String> {
let data = Arc::new(data);
let storage = Arc::new(LazyStorage::new(data.len() as u64, config));
let first = (storage.config().block_size as usize).min(data.len());
storage.supply(0, &data[..first])?;
let s = storage.clone();
let reader =
storage.run_blocking(|| Reader::open_storage(s.clone()), |r| fetch(&data, r))??;
Ok(Lazy {
data,
storage,
reader,
})
}
}
impl Api for Lazy {
fn call<T>(&self, op: impl Fn(&Reader) -> T) -> T {
self.storage
.run_blocking(|| op(&self.reader), |r| fetch(&self.data, r))
.expect("serving from memory cannot fail")
}
}
/// Everything the viewer can show of a file, as text: each object's kind,
/// listing, attributes (and attribute errors), dataset info, whole value
/// and a hyperslab — or the error each gives.
fn transcript(api: &impl Api) -> Vec<String> {
let mut out = Vec::new();
let mut todo = vec![("/".to_string(), 0usize)];
while let Some((path, depth)) = todo.pop() {
if out.len() > 4000 {
out.push("... (truncated)".into());
break;
}
let kind = api.call(|r| r.kind(&path));
out.push(format!("{path}: {kind:?}"));
out.push(format!("{path} attrs: {:?}", api.call(|r| r.attrs(&path))));
match kind {
Ok(Kind::Group) => {
let list = api.call(|r| r.list(&path));
out.push(format!("{path} list: {list:?}"));
if let Ok(children) = list
&& depth < 12
{
for c in children.into_iter().rev() {
let child = if path == "/" {
format!("/{}", c.name)
} else {
format!("{path}/{}", c.name)
};
todo.push((child, depth + 1));
}
}
}
Ok(Kind::Dataset) => {
let info = api.call(|r| r.info(&path));
out.push(format!("{path} info: {info:?}"));
let Ok(info) = info else { continue };
let n = info
.shape
.iter()
.chain(&info.element_shape)
.try_fold(1u64, |a, &d| a.checked_mul(d));
if n.is_none_or(|n| n > 4_000_000) {
out.push(format!("{path}: not read ({n:?} values)"));
continue;
}
out.push(format!(
"{path} read: {:?}",
api.call(|r| r.read(&path, None))
));
if !info.shape.is_empty() && info.shape.iter().all(|&d| d > 1) {
let slab = Hyperslab {
start: info.shape.iter().map(|_| 1).collect(),
count: info.shape.iter().map(|&d| d / 2).collect(),
stride: None,
block: None,
};
let part = api.call(|r| r.read(&path, Some(&slab)));
out.push(format!("{path} slab: {part:?}"));
}
}
Err(_) => {}
}
}
out
}
/// The lazy transcript of `data` at `block` bytes per block equals the
/// transcript of the same file through a range storage that has every byte
/// (`CountingStorage`: the facade's `Storage` path, the one the lazy reader
/// takes), and agrees with the in-memory one: the same values, and an error
/// wherever it has one (a malformed file can fail at a different check,
/// with a different message, when read by ranges). Returns what the lazy
/// reader fetched and its transcript.
fn check_equal(name: &str, data: &[u8], block: u64) -> (u64, u64, Vec<String>) {
let ctx = format!("{name} (blocks of {block} B)");
let ranged = Reader::open_storage(Arc::new(CountingStorage::new(data.to_vec())));
let local = Reader::open(data.to_vec());
let lazy = Lazy::open(data.to_vec(), config(block));
let (ranged, local, lazy) = match (ranged, local, lazy) {
(Ok(r), Ok(l), Ok(z)) => (r, l, z),
(Err(r), Err(_), Err(z)) => {
assert_eq!(z, r, "{ctx}: open error");
return (0, 0, Vec::new());
}
(r, l, z) => panic!(
"{ctx}: opens differently: ranged {:?}, in memory {:?}, lazily {:?}",
r.err(),
l.err(),
z.err()
),
};
let got = transcript(&lazy);
let want = transcript(&Local(ranged));
for (i, (w, g)) in want.iter().zip(&got).enumerate() {
assert_eq!(g, w, "{ctx}, line {i}");
}
assert_eq!(got.len(), want.len(), "{ctx}: transcript length");
let local = transcript(&Local(local));
for (i, (l, g)) in local.iter().zip(&got).enumerate() {
let both_errors = match (l.split_once("Err("), g.split_once("Err(")) {
(Some((a, _)), Some((b, _))) => a == b,
_ => false,
};
assert!(
l == g || both_errors,
"{ctx}, line {i}: in memory\n {l}\nlazily\n {g}"
);
}
assert_eq!(got.len(), local.len(), "{ctx}: transcript length");
let st = lazy.storage.stats();
(st.requests, st.bytes_fetched, got)
}
fn config(block: u64) -> LazyConfig {
LazyConfig {
block_size: block,
// A small budget, so eviction between operations is exercised.
capacity: 16 * block,
max_request: 8 * block,
}
}
fn builder_file() -> Vec<u8> {
let mut b = FileBuilder::new();
b.create_dataset("grid")
.with_f64_data(&(0..20_000).map(f64::from).collect::<Vec<_>>())
.with_shape(&[100, 200])
.with_chunks(&[10, 25])
.with_deflate(4);
b.create_dataset("contiguous")
.with_i32_data(&(0..50_000).collect::<Vec<_>>());
b.create_dataset("bytes").with_u8_data(&[1, 2, 250]);
let mut g = b.create_group("sensors");
for i in 0..40 {
g.create_dataset(&format!("t{i}"))
.with_f32_data(&[i as f32, 1.5, -2.25]);
}
g.set_attr("location", AttrValue::String("lab".into()));
b.add_group(g.finish());
b.set_attr("version", AttrValue::I64(3));
b.set_attr("scale", AttrValue::F64Array(vec![0.5, 2.0]));
b.finish().unwrap()
}
#[test]
fn builder_files_read_the_same_at_every_block_size() {
let data = builder_file();
for block in [512, 4096, 1 << 20] {
let (requests, _, lines) = check_equal("builder", &data, block);
assert!(requests > 0);
// The transcript covers every object, values included.
assert!(lines.iter().any(|l| l.starts_with("/grid read: Ok")));
assert!(lines.iter().any(|l| l.starts_with("/grid slab: Ok")));
assert!(lines.iter().any(|l| l.starts_with("/sensors/t39 read: Ok")));
}
}
#[test]
fn garbage_fails_to_open_as_in_memory() {
check_equal("zeros", &[0u8; 5000], 512);
check_equal("empty", &[], 512);
let mut cut = builder_file();
cut.truncate(cut.len() / 3);
check_equal("truncated", &cut, 512);
}
/// Listing a large file and reading one small dataset fetches a few blocks,
/// not the file.
#[test]
fn a_small_read_of_a_large_file_fetches_a_few_blocks() {
let mut b = FileBuilder::new();
b.create_dataset("small").with_f64_data(&[1.0, 2.0, 3.0]);
// 48 MB of raw data, written after the small dataset's metadata.
b.create_dataset("big")
.with_f64_data(&(0..6_000_000).map(f64::from).collect::<Vec<_>>());
let mut g = b.create_group("group");
g.create_dataset("inner").with_i32_data(&[7, 8]);
b.add_group(g.finish());
let data = b.finish().unwrap();
let lazy = Lazy::open(data.clone(), LazyConfig::default()).unwrap();
let list = lazy.call(|r| r.list("/")).unwrap();
assert_eq!(list.len(), 3);
assert_eq!(
format!("{:?}", lazy.call(|r| r.read("/small", None)).unwrap().data),
"F64([1.0, 2.0, 3.0])"
);
assert_eq!(
format!(
"{:?}",
lazy.call(|r| r.read("/group/inner", None)).unwrap().data
),
"I32([7, 8])"
);
// A window of the big dataset reads only its block(s).
let slab = Hyperslab {
start: vec![3_000_000],
count: vec![4],
stride: None,
block: None,
};
assert_eq!(
format!(
"{:?}",
lazy.call(|r| r.read("/big", Some(&slab))).unwrap().data
),
"F64([3000000.0, 3000001.0, 3000002.0, 3000003.0])"
);
let st = lazy.storage.stats();
eprintln!("{} bytes: {st:?}", data.len());
assert!(st.requests <= 6, "{st:?}");
assert!(st.bytes_fetched <= 6 << 20, "{st:?}");
assert!(st.bytes_fetched * 8 < data.len() as u64, "{st:?}");
}
fn python() -> String {
std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string())
}
fn python_available() -> bool {
Command::new(python())
.args(["-c", "import h5py, netCDF4, numpy"])
.output()
.is_ok_and(|o| o.status.success())
}
#[test]
fn h5py_and_netcdf4_files_read_the_same_lazily() {
if !python_available() {
assert!(
!std::env::var("CLAWHDF5_REQUIRE_INTEROP").is_ok_and(|v| v == "1"),
"CLAWHDF5_REQUIRE_INTEROP=1 but {} lacks h5py/netCDF4/numpy",
python()
);
eprintln!("skipping: {} lacks h5py/netCDF4/numpy", python());
return;
}
let dir = 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<PathBuf>) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
for e in entries.flatten() {
let p = e.path();
if p.is_dir() {
hdf5_files(&p, out);
} else if std::fs::read(&p)
.ok()
.is_some_and(|b| b.len() <= 64 << 20 && is_hdf5(&b))
{
out.push(p);
}
}
}
/// The HDF5 signature at 0 or a power-of-two user-block offset.
fn is_hdf5(b: &[u8]) -> bool {
const SIG: &[u8] = b"\x89HDF\r\n\x1a\n";
let mut at = 0usize;
loop {
if b.get(at..at + 8) == Some(SIG) {
return true;
}
at = if at == 0 { 512 } else { at * 2 };
if at >= b.len() {
return false;
}
}
}
#[test]
fn corpus_files_read_the_same_lazily() {
let Ok(dirs) = std::env::var("CLAWHDF5_WASM_CORPUS") else {
eprintln!("CLAWHDF5_WASM_CORPUS not set; skipping the corpus");
return;
};
let mut files = Vec::new();
for d in std::env::split_paths(&dirs) {
hdf5_files(&d, &mut files);
}
files.sort();
assert!(!files.is_empty(), "no HDF5 files under {dirs}");
let (mut requests, mut bytes, mut total) = (0u64, 0u64, 0u64);
for f in &files {
let data = std::fs::read(f).unwrap();
total += data.len() as u64;
let (r, b, _) = check_equal(&f.display().to_string(), &data, 64 * 1024);
requests += r;
bytes += b;
}
eprintln!(
"{} files ({total} bytes): {requests} requests, {bytes} bytes fetched",
files.len()
);
}