wasm: openUrl reads remote files by HTTP range requests (range-read M4)
openUrl(url, opts) returns a RemoteFile with the methods of H5File (kind, list, info, attrs, attrErrors, read, readHyperslab), each a promise, and stats(). It runs every call through the restartable LazyStorage: a pass that misses reports the byte ranges, js/remote.js fetches them with fetch() and Range headers (six at a time), and the pass is re-run. This keeps the main thread free without a Worker or synchronous XHR (h5wasm's lazy files need both), as the design doc recommends; the cost is re-running a pass per wave of misses. Every answer is checked: a 206 with exactly the bytes asked for, and the same ETag/Last-Modified and length as at open, else an error (never data). A server that ignores Range (200) is downloaded whole, up to maxDownload (512 MiB), unless fallback: "error". Options: blockSize, cacheSize, headers, credentials, parallel, fetch. test/serve.py is a range-capable static server with request counting (and /norange/ for a server without range support). test.mjs repeats every fixture check on files opened by URL (1 MiB and 512 B blocks), checks the request budget on a 200 MB h5py file (list, three small reads and a window of the big dataset: 5 requests, 6 MiB), the download fallback, and HTTP errors, changed files and wrong answers; with CLAWHDF5_WASM_CORPUS every corpus file is compared with open(bytes). Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -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"
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
+434
-64
@@ -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::<js_sys::Error>()
|
||||
.map(|e| String::from(e.message()))
|
||||
.or_else(|| e.as_string())
|
||||
.unwrap_or_else(|| format!("{e:?}"));
|
||||
JsError::new(&msg)
|
||||
}
|
||||
|
||||
fn set(obj: &Object, key: &str, value: impl Into<JsValue>) {
|
||||
// Defining a property on a fresh plain object cannot fail.
|
||||
Reflect::set(obj, &JsValue::from_str(key), &value.into()).unwrap_throw();
|
||||
}
|
||||
|
||||
fn get(obj: &JsValue, key: &str) -> JsValue {
|
||||
if obj.is_object() {
|
||||
Reflect::get(obj, &JsValue::from_str(key)).unwrap_or(JsValue::UNDEFINED)
|
||||
} else {
|
||||
JsValue::UNDEFINED
|
||||
}
|
||||
}
|
||||
|
||||
fn shape_to_js(shape: &[u64]) -> Array {
|
||||
shape.iter().map(|&d| JsValue::from_f64(d as f64)).collect()
|
||||
}
|
||||
@@ -59,6 +91,20 @@ fn indices_from_js(what: &str, v: &[f64]) -> Result<Vec<u64>, JsError> {
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn slab_from_js(
|
||||
start: &[f64],
|
||||
count: &[f64],
|
||||
stride: Option<Vec<f64>>,
|
||||
block: Option<Vec<f64>>,
|
||||
) -> Result<Hyperslab, JsError> {
|
||||
Ok(Hyperslab {
|
||||
start: indices_from_js("start", start)?,
|
||||
count: indices_from_js("count", count)?,
|
||||
stride: stride.map(|s| indices_from_js("stride", &s)).transpose()?,
|
||||
block: block.map(|b| indices_from_js("block", &b)).transpose()?,
|
||||
})
|
||||
}
|
||||
|
||||
fn data_to_js(data: Data) -> JsValue {
|
||||
match data {
|
||||
Data::F32(v) => js_sys::Float32Array::from(&v[..]).into(),
|
||||
@@ -111,6 +157,72 @@ fn attr_to_js(value: AttrValue) -> (JsValue, Option<String>) {
|
||||
}
|
||||
}
|
||||
|
||||
fn list_to_js(children: Vec<Child>) -> Array {
|
||||
children
|
||||
.into_iter()
|
||||
.map(|c| {
|
||||
let o = Object::new();
|
||||
set(&o, "name", c.name);
|
||||
set(&o, "kind", c.kind.as_str());
|
||||
JsValue::from(o)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn info_to_js(i: DatasetInfo) -> Object {
|
||||
let o = Object::new();
|
||||
set(&o, "shape", shape_to_js(&i.shape));
|
||||
let max: JsValue = match i.maxshape {
|
||||
None => JsValue::NULL,
|
||||
Some(dims) => dims
|
||||
.into_iter()
|
||||
.map(|d| d.map_or(JsValue::NULL, |d| JsValue::from_f64(d as f64)))
|
||||
.collect::<Array>()
|
||||
.into(),
|
||||
};
|
||||
set(&o, "maxshape", max);
|
||||
set(&o, "dtype", i.dtype);
|
||||
set(&o, "elementShape", shape_to_js(&i.element_shape));
|
||||
o
|
||||
}
|
||||
|
||||
fn attrs_to_js(attrs: Vec<Attr>) -> Array {
|
||||
attrs
|
||||
.into_iter()
|
||||
.map(|a| {
|
||||
let o = Object::new();
|
||||
set(&o, "name", a.name);
|
||||
let (value, dtype) = attr_to_js(a.value);
|
||||
set(&o, "value", value);
|
||||
set(&o, "dtype", dtype.map_or(JsValue::NULL, JsValue::from));
|
||||
JsValue::from(o)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn errors_to_js(errors: Vec<String>) -> Array {
|
||||
errors.into_iter().map(JsValue::from).collect()
|
||||
}
|
||||
|
||||
/// A read's values with the dataset's datatype.
|
||||
fn values_to_js((dtype, a): (String, core::Array)) -> Object {
|
||||
let o = Object::new();
|
||||
set(&o, "shape", shape_to_js(&a.shape));
|
||||
set(&o, "dtype", dtype);
|
||||
set(&o, "data", data_to_js(a.data));
|
||||
o
|
||||
}
|
||||
|
||||
/// Read the dataset at `path` (whole, or `slab`) with its datatype.
|
||||
fn read_values(
|
||||
r: &Reader,
|
||||
path: &str,
|
||||
slab: Option<&Hyperslab>,
|
||||
) -> core::Result<(String, core::Array)> {
|
||||
let dtype = r.info(path)?.dtype;
|
||||
Ok((dtype, r.read(path, slab)?))
|
||||
}
|
||||
|
||||
/// An open HDF5 (or NetCDF-4) file.
|
||||
#[wasm_bindgen]
|
||||
pub struct H5File {
|
||||
@@ -146,38 +258,13 @@ impl H5File {
|
||||
|
||||
/// The group's members: `[{ name, kind }]`, groups first.
|
||||
pub fn list(&self, path: &str) -> Result<Array, JsError> {
|
||||
Ok(self
|
||||
.inner
|
||||
.list(path)
|
||||
.map_err(js_err)?
|
||||
.into_iter()
|
||||
.map(|c| {
|
||||
let o = Object::new();
|
||||
set(&o, "name", c.name);
|
||||
set(&o, "kind", c.kind.as_str());
|
||||
JsValue::from(o)
|
||||
})
|
||||
.collect())
|
||||
Ok(list_to_js(self.inner.list(path).map_err(js_err)?))
|
||||
}
|
||||
|
||||
/// `{ shape, maxshape, dtype, elementShape }`. `maxshape` is `null`
|
||||
/// when not recorded, with `null` for each unlimited dimension.
|
||||
pub fn info(&self, path: &str) -> Result<Object, JsError> {
|
||||
let i = self.inner.info(path).map_err(js_err)?;
|
||||
let o = Object::new();
|
||||
set(&o, "shape", shape_to_js(&i.shape));
|
||||
let max: JsValue = match i.maxshape {
|
||||
None => JsValue::NULL,
|
||||
Some(dims) => dims
|
||||
.into_iter()
|
||||
.map(|d| d.map_or(JsValue::NULL, |d| JsValue::from_f64(d as f64)))
|
||||
.collect::<Array>()
|
||||
.into(),
|
||||
};
|
||||
set(&o, "maxshape", max);
|
||||
set(&o, "dtype", i.dtype);
|
||||
set(&o, "elementShape", shape_to_js(&i.element_shape));
|
||||
Ok(o)
|
||||
Ok(info_to_js(self.inner.info(path).map_err(js_err)?))
|
||||
}
|
||||
|
||||
/// `[{ name, value, dtype }]`, sorted by name. Scalars are `number`
|
||||
@@ -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<Array, JsError> {
|
||||
let (attrs, _) = self.inner.attrs(path).map_err(js_err)?;
|
||||
Ok(attrs
|
||||
.into_iter()
|
||||
.map(|a| {
|
||||
let o = Object::new();
|
||||
set(&o, "name", a.name);
|
||||
let (value, dtype) = attr_to_js(a.value);
|
||||
set(&o, "value", value);
|
||||
set(&o, "dtype", dtype.map_or(JsValue::NULL, JsValue::from));
|
||||
JsValue::from(o)
|
||||
})
|
||||
.collect())
|
||||
Ok(attrs_to_js(self.inner.attrs(path).map_err(js_err)?.0))
|
||||
}
|
||||
|
||||
/// Messages for attributes that could not be read.
|
||||
#[wasm_bindgen(js_name = attrErrors)]
|
||||
pub fn attr_errors(&self, path: &str) -> Result<Array, JsError> {
|
||||
let (_, errors) = self.inner.attrs(path).map_err(js_err)?;
|
||||
Ok(errors.into_iter().map(JsValue::from).collect())
|
||||
Ok(errors_to_js(self.inner.attrs(path).map_err(js_err)?.1))
|
||||
}
|
||||
|
||||
/// The whole dataset: `{ shape, dtype, data }`, `data` in row-major
|
||||
/// order.
|
||||
pub fn read(&self, path: &str) -> Result<Object, JsError> {
|
||||
self.read_impl(path, None)
|
||||
Ok(values_to_js(
|
||||
read_values(&self.inner, path, None).map_err(js_err)?,
|
||||
))
|
||||
}
|
||||
|
||||
/// A regular hyperslab (`H5Sselect_hyperslab`): `stride` and `block`
|
||||
@@ -224,22 +301,315 @@ impl H5File {
|
||||
stride: Option<Vec<f64>>,
|
||||
block: Option<Vec<f64>>,
|
||||
) -> Result<Object, JsError> {
|
||||
let slab = Hyperslab {
|
||||
start: indices_from_js("start", &start)?,
|
||||
count: indices_from_js("count", &count)?,
|
||||
stride: stride.map(|s| indices_from_js("stride", &s)).transpose()?,
|
||||
block: block.map(|b| indices_from_js("block", &b)).transpose()?,
|
||||
};
|
||||
self.read_impl(path, Some(&slab))
|
||||
}
|
||||
|
||||
fn read_impl(&self, path: &str, slab: Option<&Hyperslab>) -> Result<Object, JsError> {
|
||||
let dtype = self.inner.info(path).map_err(js_err)?.dtype;
|
||||
let a = self.inner.read(path, slab).map_err(js_err)?;
|
||||
let slab = slab_from_js(&start, &count, stride, block)?;
|
||||
Ok(values_to_js(
|
||||
read_values(&self.inner, path, Some(&slab)).map_err(js_err)?,
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Remote files: openUrl.
|
||||
|
||||
#[wasm_bindgen(module = "/js/remote.js")]
|
||||
extern "C" {
|
||||
#[wasm_bindgen(catch)]
|
||||
async fn probe(url: &str, first_len: f64, opts: &JsValue) -> Result<JsValue, JsValue>;
|
||||
|
||||
#[wasm_bindgen(catch, js_name = fetchRanges)]
|
||||
async fn fetch_ranges(
|
||||
url: &str,
|
||||
ranges: Vec<f64>,
|
||||
opts: &JsValue,
|
||||
validator: &JsValue,
|
||||
length: f64,
|
||||
) -> Result<JsValue, JsValue>;
|
||||
}
|
||||
|
||||
/// Where a remote file's bytes come from.
|
||||
struct Http {
|
||||
url: String,
|
||||
opts: JsValue,
|
||||
/// ETag or Last-Modified at open (`null` if the server sent neither).
|
||||
validator: JsValue,
|
||||
length: u64,
|
||||
/// Requests the probe made (1, or 2 with a HEAD for the length).
|
||||
probe_requests: u64,
|
||||
}
|
||||
|
||||
enum Source {
|
||||
/// Read by range requests through a restartable cache.
|
||||
Lazy {
|
||||
http: Http,
|
||||
storage: Arc<LazyStorage>,
|
||||
reader: Reader,
|
||||
},
|
||||
/// The server ignored `Range`: the whole file, downloaded at open.
|
||||
Whole {
|
||||
reader: Reader,
|
||||
size: u64,
|
||||
requests: u64,
|
||||
},
|
||||
}
|
||||
|
||||
impl Http {
|
||||
/// Fetch `ranges` and hand them to `storage`.
|
||||
async fn fetch(&self, storage: &LazyStorage, ranges: &[Range<u64>]) -> Result<(), JsError> {
|
||||
let flat: Vec<f64> = ranges
|
||||
.iter()
|
||||
.flat_map(|r| [r.start as f64, r.end as f64])
|
||||
.collect();
|
||||
let got = fetch_ranges(
|
||||
&self.url,
|
||||
flat,
|
||||
&self.opts,
|
||||
&self.validator,
|
||||
self.length as f64,
|
||||
)
|
||||
.await
|
||||
.map_err(js_exception)?;
|
||||
let got = Array::from(&got);
|
||||
if got.length() as usize != ranges.len() {
|
||||
return Err(js_err(format!(
|
||||
"fetchRanges returned {} ranges for {}",
|
||||
got.length(),
|
||||
ranges.len()
|
||||
)));
|
||||
}
|
||||
for (r, bytes) in ranges.iter().zip(got.iter()) {
|
||||
let bytes = Uint8Array::new(&bytes).to_vec();
|
||||
storage.supply_range(r, &bytes).map_err(js_err)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Run `f` over `storage` until it has every byte it reads.
|
||||
async fn drive<T>(
|
||||
&self,
|
||||
storage: &LazyStorage,
|
||||
mut f: impl FnMut() -> T,
|
||||
) -> Result<T, JsError> {
|
||||
let _op = storage.operation();
|
||||
loop {
|
||||
match storage.attempt(&mut f) {
|
||||
Step::Done(v) => return Ok(v),
|
||||
Step::Need(ranges) => self.fetch(storage, &ranges).await?,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Source {
|
||||
async fn run<T>(&self, op: impl Fn(&Reader) -> core::Result<T>) -> Result<T, JsError> {
|
||||
match self {
|
||||
Source::Whole { reader, .. } => op(reader).map_err(js_err),
|
||||
Source::Lazy {
|
||||
http,
|
||||
storage,
|
||||
reader,
|
||||
} => http.drive(storage, || op(reader)).await?.map_err(js_err),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A non-negative integer option, or `None` when not given.
|
||||
fn int_opt(opts: &JsValue, key: &str, min: f64, max: f64) -> Result<Option<u64>, JsError> {
|
||||
let v = get(opts, key);
|
||||
if v.is_undefined() || v.is_null() {
|
||||
return Ok(None);
|
||||
}
|
||||
match v.as_f64() {
|
||||
Some(x) if x.fract() == 0.0 && (min..=max).contains(&x) => Ok(Some(x as u64)),
|
||||
_ => Err(js_err(format!(
|
||||
"openUrl: {key} must be an integer from {min} to {max}"
|
||||
))),
|
||||
}
|
||||
}
|
||||
|
||||
fn config_from(opts: &JsValue) -> Result<LazyConfig, JsError> {
|
||||
let mut c = LazyConfig::default();
|
||||
if let Some(b) = int_opt(opts, "blockSize", 512.0, (64u64 << 20) as f64)? {
|
||||
c.block_size = b;
|
||||
}
|
||||
if let Some(n) = int_opt(opts, "cacheSize", 0.0, MAX_SAFE_INTEGER)? {
|
||||
c.capacity = n;
|
||||
}
|
||||
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<RemoteFile, JsError> {
|
||||
let config = config_from(&opts)?;
|
||||
let p = probe(&url, config.block_size as f64, &opts)
|
||||
.await
|
||||
.map_err(js_exception)?;
|
||||
let requests = get(&p, "requests").as_f64().unwrap_or(1.0) as u64;
|
||||
let whole = get(&p, "whole");
|
||||
if !whole.is_undefined() {
|
||||
let bytes = Uint8Array::new(&whole).to_vec();
|
||||
let size = bytes.len() as u64;
|
||||
let reader = Reader::open(bytes).map_err(js_err)?;
|
||||
return Ok(RemoteFile {
|
||||
inner: Rc::new(Source::Whole {
|
||||
reader,
|
||||
size,
|
||||
requests,
|
||||
}),
|
||||
});
|
||||
}
|
||||
let length = get(&p, "length")
|
||||
.as_f64()
|
||||
.filter(|x| x.fract() == 0.0 && (0.0..=MAX_SAFE_INTEGER).contains(x))
|
||||
.ok_or_else(|| js_err(format!("{url}: the server gave no usable file size")))?;
|
||||
let http = Http {
|
||||
url,
|
||||
opts,
|
||||
validator: get(&p, "validator"),
|
||||
length: length as u64,
|
||||
probe_requests: requests,
|
||||
};
|
||||
let storage = Arc::new(LazyStorage::new(http.length, config));
|
||||
let first = Uint8Array::new(&get(&p, "first")).to_vec();
|
||||
storage.supply(0, &first).map_err(js_err)?;
|
||||
let s = storage.clone();
|
||||
let reader = http
|
||||
.drive(&storage, || Reader::open_storage(s.clone()))
|
||||
.await?
|
||||
.map_err(js_err)?;
|
||||
Ok(RemoteFile {
|
||||
inner: Rc::new(Source::Lazy {
|
||||
http,
|
||||
storage,
|
||||
reader,
|
||||
}),
|
||||
})
|
||||
}
|
||||
|
||||
/// A file opened with [`openUrl`](open_url): the methods of [`H5File`],
|
||||
/// each returning a `Promise` (it may have to fetch bytes first).
|
||||
#[wasm_bindgen]
|
||||
pub struct RemoteFile {
|
||||
inner: Rc<Source>,
|
||||
}
|
||||
|
||||
impl RemoteFile {
|
||||
/// Run `op` (fetching what it needs) and convert its result.
|
||||
fn call<T: 'static>(
|
||||
&self,
|
||||
op: impl Fn(&Reader) -> core::Result<T> + 'static,
|
||||
to_js: impl FnOnce(T) -> JsValue + 'static,
|
||||
) -> Promise {
|
||||
let inner = self.inner.clone();
|
||||
future_to_promise(async move {
|
||||
let v = inner.run(op).await.map_err(JsValue::from)?;
|
||||
Ok(to_js(v))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[wasm_bindgen]
|
||||
impl RemoteFile {
|
||||
/// `"group"` or `"dataset"`.
|
||||
#[wasm_bindgen(unchecked_return_type = "Promise<string>")]
|
||||
pub fn kind(&self, path: String) -> Promise {
|
||||
self.call(move |r| r.kind(&path), |k| k.as_str().into())
|
||||
}
|
||||
|
||||
/// The group's members: `[{ name, kind }]`, groups first.
|
||||
#[wasm_bindgen(unchecked_return_type = "Promise<Array<any>>")]
|
||||
pub fn list(&self, path: String) -> Promise {
|
||||
self.call(move |r| r.list(&path), |c| list_to_js(c).into())
|
||||
}
|
||||
|
||||
/// `{ shape, maxshape, dtype, elementShape }`, as [`H5File::info`].
|
||||
#[wasm_bindgen(unchecked_return_type = "Promise<any>")]
|
||||
pub fn info(&self, path: String) -> Promise {
|
||||
self.call(move |r| r.info(&path), |i| info_to_js(i).into())
|
||||
}
|
||||
|
||||
/// `[{ name, value, dtype }]`, as [`H5File::attrs`].
|
||||
#[wasm_bindgen(unchecked_return_type = "Promise<Array<any>>")]
|
||||
pub fn attrs(&self, path: String) -> Promise {
|
||||
self.call(move |r| r.attrs(&path), |a| attrs_to_js(a.0).into())
|
||||
}
|
||||
|
||||
/// Messages for attributes that could not be read.
|
||||
#[wasm_bindgen(js_name = attrErrors, unchecked_return_type = "Promise<Array<string>>")]
|
||||
pub fn attr_errors(&self, path: String) -> Promise {
|
||||
self.call(move |r| r.attrs(&path), |a| errors_to_js(a.1).into())
|
||||
}
|
||||
|
||||
/// The whole dataset: `{ shape, dtype, data }`, as [`H5File::read`].
|
||||
#[wasm_bindgen(unchecked_return_type = "Promise<any>")]
|
||||
pub fn read(&self, path: String) -> Promise {
|
||||
self.call(
|
||||
move |r| read_values(r, &path, None),
|
||||
|v| values_to_js(v).into(),
|
||||
)
|
||||
}
|
||||
|
||||
/// A regular hyperslab, as [`H5File::read_hyperslab`]. Only the chunks
|
||||
/// (or the contiguous runs) the selection touches are fetched.
|
||||
#[wasm_bindgen(js_name = readHyperslab, unchecked_return_type = "Promise<any>")]
|
||||
pub fn read_hyperslab(
|
||||
&self,
|
||||
path: String,
|
||||
start: Vec<f64>,
|
||||
count: Vec<f64>,
|
||||
stride: Option<Vec<f64>>,
|
||||
block: Option<Vec<f64>>,
|
||||
) -> Result<Promise, JsError> {
|
||||
let slab = slab_from_js(&start, &count, stride, block)?;
|
||||
Ok(self.call(
|
||||
move |r| read_values(r, &path, Some(&slab)),
|
||||
|v| values_to_js(v).into(),
|
||||
))
|
||||
}
|
||||
|
||||
/// What reading this file has cost so far: `{ lazy, size, requests,
|
||||
/// bytesFetched, cachedBytes, passes }`. `lazy` is false when the
|
||||
/// server ignored `Range` and the file was downloaded whole.
|
||||
pub fn stats(&self) -> Object {
|
||||
let o = Object::new();
|
||||
set(&o, "shape", shape_to_js(&a.shape));
|
||||
set(&o, "dtype", dtype);
|
||||
set(&o, "data", data_to_js(a.data));
|
||||
Ok(o)
|
||||
match &*self.inner {
|
||||
Source::Lazy { http, storage, .. } => {
|
||||
let st = storage.stats();
|
||||
set(&o, "lazy", true);
|
||||
set(&o, "size", http.length as f64);
|
||||
set(
|
||||
&o,
|
||||
"requests",
|
||||
(st.requests.saturating_sub(1) + http.probe_requests) as f64,
|
||||
);
|
||||
set(&o, "bytesFetched", st.bytes_fetched as f64);
|
||||
set(&o, "cachedBytes", st.cached_bytes as f64);
|
||||
set(&o, "passes", st.passes as f64);
|
||||
}
|
||||
Source::Whole { size, requests, .. } => {
|
||||
set(&o, "lazy", false);
|
||||
set(&o, "size", *size as f64);
|
||||
set(&o, "requests", *requests as f64);
|
||||
set(&o, "bytesFetched", *size as f64);
|
||||
set(&o, "cachedBytes", *size as f64);
|
||||
set(&o, "passes", 0.0);
|
||||
}
|
||||
}
|
||||
o
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user