diff --git a/crates/clawhdf5-wasm/js/remote.js b/crates/clawhdf5-wasm/js/remote.js index fb78c55..cbe3990 100644 --- a/crates/clawhdf5-wasm/js/remote.js +++ b/crates/clawhdf5-wasm/js/remote.js @@ -17,8 +17,28 @@ function fetcher(opts) { return f; } -function init(opts, extra, method = "GET") { - return { method, headers: { ...(opts?.headers ?? {}), ...extra }, credentials: opts?.credentials }; +// The caller's headers (`opts.headers`: a Headers, [name, value] pairs or a +// plain object, as fetch takes them; names come out in lower case) with +// `extra` over them. A Range of the caller's is dropped: this file asks for +// the ranges. +function init(opts, extra, method = "GET", signal = undefined) { + const headers = {}; + if (opts?.headers != null) { + for (const [name, value] of new Headers(opts.headers)) { + if (name !== "range") headers[name] = value; + } + } + Object.assign(headers, extra); + return { method, headers, credentials: opts?.credentials, signal }; +} + +// `opts.parallel`: range requests in flight at once. +function parallelism(opts) { + const p = opts?.parallel ?? DEFAULT_PARALLEL; + if (!Number.isSafeInteger(p) || p < 1) { + throw new Error(`openUrl: parallel must be a positive integer, got ${String(p)}`); + } + return p; } // "bytes a-b/total" -> { start, end (exclusive), total | null }; null when @@ -108,6 +128,7 @@ function readLimited(resp, limit, url, start) { */ export async function probe(url, firstLen, opts) { const f = fetcher(opts); + parallelism(opts); const resp = await f(url, init(opts, { Range: `bytes=0-${firstLen - 1}` })); if (resp.status === 206) { const cr = contentRange(resp, url); @@ -157,44 +178,58 @@ export async function probe(url, firstLen, opts) { * 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). + * file (validator and length). When one request fails, the others in + * flight are aborted and no more are made; that failure is the error. */ export async function fetchRanges(url, ranges, opts, validator, length) { const f = fetcher(opts); + const parallel = parallelism(opts); const n = ranges.length / 2; const out = new Array(n); + const abort = new AbortController(); let next = 0; + async function one(i) { + const start = ranges[2 * i]; + const end = ranges[2 * i + 1]; + const resp = await f(url, init(opts, { Range: `bytes=${start}-${end - 1}` }, "GET", abort.signal)); + 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 = await readLimited(resp, end - start, url, start); + if (body.length !== end - start) { + throw new Error(`${url}: asked for ${end - start} bytes at offset ${start}, got ${body.length}`); + } + out[i] = body; + } 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()); + while (next < n && !abort.signal.aborted) { + try { + await one(next++); + } catch (e) { + // The first failure stops the rest: requests in flight are aborted + // (their AbortErrors are not reported) and no new ones start. + if (!abort.signal.aborted) { + abort.abort(); + throw e; + } + return; } - 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 = await readLimited(resp, end - start, url, start); - 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)); + await Promise.all(Array.from({ length: Math.min(parallel, n) }, worker)); return out; } diff --git a/crates/clawhdf5-wasm/src/lib.rs b/crates/clawhdf5-wasm/src/lib.rs index 1fc2521..6a512a2 100644 --- a/crates/clawhdf5-wasm/src/lib.rs +++ b/crates/clawhdf5-wasm/src/lib.rs @@ -461,6 +461,8 @@ fn config_from(opts: &JsValue) -> Result { // Read by remote.js; checked here so a value wasm32 cannot hold is an // option error rather than a download that cannot be kept. int_opt(opts, "maxDownload", 0.0, MAX_FETCH_LIMIT as f64)?; + // Also read by remote.js (which checks it too, for direct callers). + int_opt(opts, "parallel", 1.0, 1024.0)?; Ok(c) } @@ -477,8 +479,10 @@ fn config_from(opts: &JsValue) -> Result { /// - `fallback` — `"download"` (default) reads the whole file when the /// server ignores `Range` (answers 200), up to `maxDownload` bytes /// (default 512 MiB, at most 1 GiB); `"error"` refuses such a server; -/// - `headers`, `credentials` — passed to every `fetch`; -/// - `parallel` — range requests in flight at once (default 6); +/// - `headers`, `credentials` — passed to every `fetch` (`headers` as +/// `fetch` takes them: a `Headers`, `[name, value]` pairs or an object); +/// - `parallel` — range requests in flight at once, 1 to 1024 (default +/// 6); when one fails the others are aborted; /// - `fetch` — a `fetch`-compatible function to use instead of the global. /// /// Cross-origin servers must allow CORS and expose `Content-Range` (or diff --git a/examples/wasm-viewer/test/test.mjs b/examples/wasm-viewer/test/test.mjs index d53b23e..e76ecdc 100644 --- a/examples/wasm-viewer/test/test.mjs +++ b/examples/wasm-viewer/test/test.mjs @@ -203,17 +203,61 @@ async function remoteTests() { 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"); + // A custom fetch is used for every request, with the caller's headers, + // given in any form fetch takes (a caller's Range is not sent). + for (const headers of [{ "X-Test": "1" }, new Headers({ "X-Test": "1", Range: "bytes=0-0" }), [["X-Test", "1"]]]) { + let calls = 0; + const seen = new Set(); + const viaCustom = await pkg.openUrl(`${base}/fix/fixture.h5`, { + blockSize: 4096, + headers, + fetch: (url, init) => { + calls++; + const h = new Headers(init.headers); + seen.add(`${h.get("X-Test")} ${h.get("Range").startsWith("bytes=0-0") ? "caller's range" : "ours"}`); + return fetch(url, init); + }, + }); + const kind = headers.constructor.name; + eq((await viaCustom.read("/sensors/temp")).data[0], 21.5, `custom fetch values (${kind})`); + eq(calls, viaCustom.stats().requests, `custom fetch calls (${kind})`); + eq([...seen], ["1 ours"], `headers passed (${kind})`); + } + + // parallel must be a positive integer. + for (const parallel of [0, -1, 1.5, "4", NaN, 5000]) { + await fails(() => pkg.openUrl(`${base}/fix/fixture.h5`, { parallel }), /parallel/, `parallel ${String(parallel)}`); + } + + // When one range request fails, the others in flight are aborted and no + // more start (fetchRanges, called directly: remote.js is in the package). + { + const snippets = join(pkgDir, "snippets"); + const dir = readdirSync(snippets).find((d) => existsSync(join(snippets, d, "js", "remote.js"))); + const remote = await import(pathToFileURL(join(snippets, dir, "js", "remote.js"))); + let started = 0; + const aborted = []; + const slowFetch = (url, init) => { + const k = started++; + if (k === 2) return Promise.resolve(new Response(null, { status: 500 })); + return new Promise((resolve, reject) => { + const t = setTimeout(() => resolve(new Response(new Uint8Array(10), { status: 206 })), 200); + init.signal?.addEventListener("abort", () => { + clearTimeout(t); + aborted.push(k); + reject(new DOMException("aborted", "AbortError")); + }); + }); + }; + const ranges = Array.from({ length: 20 }, (_, i) => [i * 10, i * 10 + 10]).flat(); + await fails(() => remote.fetchRanges("http://x.invalid/f.h5", ranges, { fetch: slowFetch, parallel: 3 }, null, 200), + /HTTP 500/, "a failed range request is the error"); + await new Promise((r) => setTimeout(r, 300)); + eq(started, 3, "no request starts after a failure"); + eq(aborted.sort(), [0, 1], "requests in flight are aborted"); + await fails(() => remote.fetchRanges("http://x.invalid/f.h5", ranges, { fetch: slowFetch, parallel: "x" }, null, 200), + /parallel must be a positive integer/, "fetchRanges checks parallel"); + } // 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.