wasm: openUrl takes Headers, aborts sibling requests on failure, checks parallel
- opts.headers is read as fetch reads it (a Headers, [name, value] pairs or a plain object); it was spread as an object, which silently dropped a Headers instance (a common way to pass Authorization). A caller's Range is not sent. - When one range request of a batch fails, the others in flight are aborted (one AbortController per batch, its signal passed to fetch) and no new ones start; the first failure is the error. The workers used to go on issuing requests nobody waited for. - parallel must be an integer from 1 to 1024 (openUrl) or a positive integer (fetchRanges): a non-number gave NaN workers, so none ran and fetchRanges returned nothing. Tests (test.mjs): headers as object, Headers and pairs reach the fetch; bad parallel values are option errors; fetchRanges with a 500 on the third of 20 ranges at parallel 3 starts 3 requests and aborts the 2 in flight. Before: the Headers case sent no header, 17 requests started after the failure with none aborted, and parallel "x" returned. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -17,8 +17,28 @@ function fetcher(opts) {
|
|||||||
return f;
|
return f;
|
||||||
}
|
}
|
||||||
|
|
||||||
function init(opts, extra, method = "GET") {
|
// The caller's headers (`opts.headers`: a Headers, [name, value] pairs or a
|
||||||
return { method, headers: { ...(opts?.headers ?? {}), ...extra }, credentials: opts?.credentials };
|
// 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
|
// "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) {
|
export async function probe(url, firstLen, opts) {
|
||||||
const f = fetcher(opts);
|
const f = fetcher(opts);
|
||||||
|
parallelism(opts);
|
||||||
const resp = await f(url, init(opts, { Range: `bytes=0-${firstLen - 1}` }));
|
const resp = await f(url, init(opts, { Range: `bytes=0-${firstLen - 1}` }));
|
||||||
if (resp.status === 206) {
|
if (resp.status === 206) {
|
||||||
const cr = contentRange(resp, url);
|
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
|
* Fetch `ranges` ([start0, end0, start1, end1, ...], ends exclusive) of a
|
||||||
* file opened by `probe`, at most `opts.parallel` (default 6) at a time.
|
* 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
|
* 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) {
|
export async function fetchRanges(url, ranges, opts, validator, length) {
|
||||||
const f = fetcher(opts);
|
const f = fetcher(opts);
|
||||||
|
const parallel = parallelism(opts);
|
||||||
const n = ranges.length / 2;
|
const n = ranges.length / 2;
|
||||||
const out = new Array(n);
|
const out = new Array(n);
|
||||||
|
const abort = new AbortController();
|
||||||
let next = 0;
|
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() {
|
async function worker() {
|
||||||
while (next < n) {
|
while (next < n && !abort.signal.aborted) {
|
||||||
const i = next++;
|
try {
|
||||||
const start = ranges[2 * i];
|
await one(next++);
|
||||||
const end = ranges[2 * i + 1];
|
} catch (e) {
|
||||||
const resp = await f(url, init(opts, { Range: `bytes=${start}-${end - 1}` }));
|
// The first failure stops the rest: requests in flight are aborted
|
||||||
if (resp.status !== 206) {
|
// (their AbortErrors are not reported) and no new ones start.
|
||||||
await discard(resp);
|
if (!abort.signal.aborted) {
|
||||||
throw new Error(resp.status === 200
|
abort.abort();
|
||||||
? `${url}: the server stopped honouring range requests`
|
throw e;
|
||||||
: `${url}: HTTP ${resp.status} ${resp.statusText ?? ""}`.trim());
|
}
|
||||||
|
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: Math.min(parallel, n) }, worker));
|
||||||
await Promise.all(Array.from({ length: workers }, worker));
|
|
||||||
return out;
|
return out;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -461,6 +461,8 @@ fn config_from(opts: &JsValue) -> Result<LazyConfig, JsError> {
|
|||||||
// Read by remote.js; checked here so a value wasm32 cannot hold is an
|
// Read by remote.js; checked here so a value wasm32 cannot hold is an
|
||||||
// option error rather than a download that cannot be kept.
|
// option error rather than a download that cannot be kept.
|
||||||
int_opt(opts, "maxDownload", 0.0, MAX_FETCH_LIMIT as f64)?;
|
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)
|
Ok(c)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -477,8 +479,10 @@ fn config_from(opts: &JsValue) -> Result<LazyConfig, JsError> {
|
|||||||
/// - `fallback` — `"download"` (default) reads the whole file when the
|
/// - `fallback` — `"download"` (default) reads the whole file when the
|
||||||
/// server ignores `Range` (answers 200), up to `maxDownload` bytes
|
/// server ignores `Range` (answers 200), up to `maxDownload` bytes
|
||||||
/// (default 512 MiB, at most 1 GiB); `"error"` refuses such a server;
|
/// (default 512 MiB, at most 1 GiB); `"error"` refuses such a server;
|
||||||
/// - `headers`, `credentials` — passed to every `fetch`;
|
/// - `headers`, `credentials` — passed to every `fetch` (`headers` as
|
||||||
/// - `parallel` — range requests in flight at once (default 6);
|
/// `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.
|
/// - `fetch` — a `fetch`-compatible function to use instead of the global.
|
||||||
///
|
///
|
||||||
/// Cross-origin servers must allow CORS and expose `Content-Range` (or
|
/// Cross-origin servers must allow CORS and expose `Content-Range` (or
|
||||||
|
|||||||
@@ -203,17 +203,61 @@ async function remoteTests() {
|
|||||||
eq(counted.stats().requests, server.requests, "requests counted");
|
eq(counted.stats().requests, server.requests, "requests counted");
|
||||||
eq(counted.stats().bytesFetched, server.bytes, "bytes counted");
|
eq(counted.stats().bytesFetched, server.bytes, "bytes counted");
|
||||||
|
|
||||||
// A custom fetch is used for every request, with the caller's headers.
|
// A custom fetch is used for every request, with the caller's headers,
|
||||||
let calls = 0;
|
// given in any form fetch takes (a caller's Range is not sent).
|
||||||
const seen = new Set();
|
for (const headers of [{ "X-Test": "1" }, new Headers({ "X-Test": "1", Range: "bytes=0-0" }), [["X-Test", "1"]]]) {
|
||||||
const viaCustom = await pkg.openUrl(`${base}/fix/fixture.h5`, {
|
let calls = 0;
|
||||||
blockSize: 4096,
|
const seen = new Set();
|
||||||
headers: { "X-Test": "1" },
|
const viaCustom = await pkg.openUrl(`${base}/fix/fixture.h5`, {
|
||||||
fetch: (url, init) => { calls++; seen.add(init.headers["X-Test"]); return fetch(url, init); },
|
blockSize: 4096,
|
||||||
});
|
headers,
|
||||||
eq((await viaCustom.read("/sensors/temp")).data[0], 21.5, "custom fetch values");
|
fetch: (url, init) => {
|
||||||
eq(calls, viaCustom.stats().requests, "custom fetch calls");
|
calls++;
|
||||||
eq([...seen], ["1"], "headers passed");
|
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
|
// 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.
|
// evicted while any of them runs) and each gets its own answer.
|
||||||
|
|||||||
Reference in New Issue
Block a user