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