Download tests/test_reader_request_manager.js from VoiceOfML/Search: direct link, hf CLI and curl.
- Browser
- Download file 11.1 kB
-
https://huggingface.co/spaces/VoiceOfML/Search/resolve/main/tests/test_reader_request_manager.js
- Command line
-
hf download hf://spaces/VoiceOfML/Search/tests/test_reader_request_manager.js
-
curl -L -o test_reader_request_manager.js https://huggingface.co/spaces/VoiceOfML/Search/resolve/main/tests/test_reader_request_manager.js
11.1 kB
| const assert = require("assert"); | |
| const fs = require("fs"); | |
| const path = require("path"); | |
| const vm = require("vm"); | |
| function loadManager(fetchImpl) { | |
| const timers = new Map(); | |
| let nextTimer = 1; | |
| const sandbox = { | |
| self: {}, | |
| DOMException, | |
| Response, | |
| ReadableStream, | |
| AbortController, | |
| fetch: fetchImpl, | |
| setTimeout(callback) { | |
| const id = nextTimer++; | |
| timers.set(id, callback); | |
| return id; | |
| }, | |
| clearTimeout(id) { | |
| timers.delete(id); | |
| } | |
| }; | |
| sandbox.self = sandbox; | |
| vm.runInNewContext( | |
| fs.readFileSync(path.join(__dirname, "../static/reader-request-manager.js"), "utf8"), | |
| sandbox | |
| ); | |
| return { manager: sandbox.VoiceOfMLReaderRequests.createReaderRequestManager(), timers }; | |
| } | |
| async function main() { | |
| let resolveFetch, | |
| calls = 0; | |
| const shared = loadManager(() => { | |
| calls += 1; | |
| return new Promise((resolve) => { | |
| resolveFetch = resolve; | |
| }); | |
| }); | |
| const first = shared.manager.request("book", 1000), | |
| second = shared.manager.request("book", 1000); | |
| await Promise.resolve(); | |
| assert.notStrictEqual(first, second); | |
| assert.strictEqual(calls, 1); | |
| assert.strictEqual(shared.timers.size, 1); | |
| resolveFetch(new Response("shared body", { status: 200 })); | |
| const [firstResponse, secondResponse] = await Promise.all([first, second]); | |
| assert.notStrictEqual(firstResponse, secondResponse); | |
| assert.strictEqual(await firstResponse.text(), "shared body"); | |
| assert.strictEqual(await secondResponse.text(), "shared body"); | |
| assert.strictEqual(shared.manager.pendingCount, 0); | |
| assert.strictEqual(shared.manager.activeCount, 0); | |
| assert.strictEqual(shared.timers.size, 0); | |
| const timedOut = loadManager( | |
| (_url, { signal }) => | |
| new Promise((_resolve, reject) => { | |
| signal.onabort = () => reject(new DOMException("timed out", "AbortError")); | |
| }) | |
| ); | |
| const timedRequest = timedOut.manager.request("timeout", 25); | |
| await Promise.resolve(); | |
| [...timedOut.timers.values()][0](); | |
| await assert.rejects(timedRequest, (error) => error.name === "TimeoutError"); | |
| assert.strictEqual(timedOut.manager.pendingCount, 0); | |
| assert.strictEqual(timedOut.timers.size, 0); | |
| let attempts = 0; | |
| const retry = loadManager(() => | |
| ++attempts === 1 ? Promise.reject(new Error("temporary")) : Promise.resolve({ ok: true }) | |
| ); | |
| await assert.rejects(retry.manager.request("retry", 1000), /temporary/); | |
| assert.strictEqual((await retry.manager.request("retry", 1000)).ok, true); | |
| assert.strictEqual(attempts, 2); | |
| const disposed = loadManager( | |
| (_url, { signal }) => | |
| new Promise((_resolve, reject) => { | |
| signal.onabort = () => reject(new DOMException("aborted", "AbortError")); | |
| }) | |
| ); | |
| const pending = disposed.manager.request("slow", 1000); | |
| await Promise.resolve(); | |
| disposed.manager.dispose(); | |
| await assert.rejects(pending, (error) => error.name === "AbortError"); | |
| assert.strictEqual(disposed.manager.pendingCount, 0); | |
| assert.strictEqual(disposed.timers.size, 0); | |
| await assert.rejects( | |
| disposed.manager.request("late", 1000), | |
| (error) => error.name === "AbortError" | |
| ); | |
| let streamController; | |
| const bodyTimeout = loadManager((_url, { signal }) => | |
| Promise.resolve( | |
| new Response( | |
| new ReadableStream({ | |
| start(controller) { | |
| streamController = controller; | |
| signal.onabort = () => controller.error(new DOMException("timed out", "AbortError")); | |
| } | |
| }), | |
| { status: 200 } | |
| ) | |
| ) | |
| ); | |
| const bodyRequest = bodyTimeout.manager.request("body", 25); | |
| const bodyResponse = await bodyRequest; | |
| assert.strictEqual(bodyTimeout.timers.size, 1); | |
| [...bodyTimeout.timers.values()][0](); | |
| await assert.rejects(bodyResponse.text(), (error) => error.name === "TimeoutError"); | |
| assert.ok(streamController); | |
| assert.strictEqual(bodyTimeout.timers.size, 0); | |
| assert.strictEqual(bodyTimeout.manager.activeCount, 0); | |
| // Fetch headers have settled, but both independent bodies still need disposal. | |
| let bodyCancels = 0; | |
| const signals = []; | |
| const streaming = loadManager((_url, { signal }) => { | |
| signals.push(signal); | |
| return Promise.resolve( | |
| new Response( | |
| new ReadableStream({ | |
| cancel(reason) { | |
| assert.strictEqual(reason.name, "AbortError"); | |
| bodyCancels++; | |
| } | |
| }) | |
| ) | |
| ); | |
| }); | |
| const firstBody = await streaming.manager.request("stream", 1000); | |
| const secondBody = await streaming.manager.request("stream", 1000); | |
| assert.strictEqual(streaming.manager.pendingCount, 0); | |
| assert.strictEqual(streaming.manager.activeCount, 2); | |
| const firstReading = assert.rejects(firstBody.text(), (error) => error.name === "AbortError"); | |
| streaming.manager.dispose(); | |
| await firstReading; | |
| await assert.rejects(secondBody.text(), (error) => error.name === "AbortError"); | |
| assert.strictEqual(bodyCancels, 2); | |
| assert.ok(signals.every((signal) => signal.aborted)); | |
| assert.strictEqual(streaming.timers.size, 0); | |
| assert.strictEqual(streaming.manager.activeCount, 0); | |
| streaming.manager.dispose(); | |
| assert.strictEqual(bodyCancels, 2); | |
| let finishClonedFetch, | |
| cloneCancels = 0; | |
| const cloned = loadManager( | |
| () => | |
| new Promise((resolve) => { | |
| finishClonedFetch = resolve; | |
| }) | |
| ); | |
| const cloneOne = cloned.manager.request("clones", 1000), | |
| cloneTwo = cloned.manager.request("clones", 1000); | |
| await Promise.resolve(); | |
| finishClonedFetch( | |
| new Response( | |
| new ReadableStream({ | |
| cancel() { | |
| cloneCancels++; | |
| } | |
| }) | |
| ) | |
| ); | |
| const cloneResponses = await Promise.all([cloneOne, cloneTwo]); | |
| const cloneReads = cloneResponses.map((response) => | |
| assert.rejects(response.text(), (error) => error.name === "AbortError") | |
| ); | |
| cloned.manager.dispose(); | |
| await Promise.all(cloneReads); | |
| assert.strictEqual(cloneCancels, 1); | |
| assert.strictEqual(cloned.manager.activeCount, 0); | |
| const cancelled = loadManager(() => Promise.resolve(new Response(new ReadableStream()))); | |
| const cancelResponse = await cancelled.manager.request("cancel", 1000); | |
| await cancelResponse.body.cancel(); | |
| assert.strictEqual(cancelled.manager.activeCount, 0); | |
| assert.strictEqual(cancelled.timers.size, 0); | |
| let queuedCalls = 0; | |
| const queued = loadManager(() => { | |
| queuedCalls++; | |
| return Promise.resolve(new Response("unused")); | |
| }); | |
| const queuedRequest = queued.manager.request("queued", 1000); | |
| queued.manager.dispose(); | |
| await assert.rejects(queuedRequest, (error) => error.name === "AbortError"); | |
| assert.strictEqual(queuedCalls, 0); | |
| assert.strictEqual(queued.manager.activeCount, 0); | |
| let finishLate, | |
| lateCancels = 0; | |
| const late = loadManager( | |
| () => | |
| new Promise((resolve) => { | |
| finishLate = resolve; | |
| }) | |
| ); | |
| const lateRequest = late.manager.request("late-response", 1000); | |
| await Promise.resolve(); | |
| late.manager.dispose(); | |
| finishLate( | |
| new Response( | |
| new ReadableStream({ | |
| cancel() { | |
| lateCancels++; | |
| } | |
| }) | |
| ) | |
| ); | |
| await assert.rejects(lateRequest, (error) => error.name === "AbortError"); | |
| assert.strictEqual(lateCancels, 1); | |
| assert.strictEqual(late.manager.activeCount, 0); | |
| // Exercise the same cancellation boundaries for deadlines and lifecycle disposal. | |
| for (const timeout of [true, false]) { | |
| const name = timeout ? "TimeoutError" : "AbortError"; | |
| const stop = (harness) => { | |
| if (timeout) [...harness.timers.values()][0](); | |
| else harness.manager.dispose(); | |
| }; | |
| const assertSettled = (harness) => { | |
| assert.strictEqual(harness.manager.pendingCount, 0); | |
| assert.strictEqual(harness.manager.activeCount, 0); | |
| assert.strictEqual(harness.timers.size, 0); | |
| assert.strictEqual(harness.manager.disposed, !timeout); | |
| }; | |
| let starts = 0; | |
| const beforeFetch = loadManager(() => { | |
| starts++; | |
| return Promise.resolve(new Response("unused")); | |
| }); | |
| const queued = beforeFetch.manager.request("queued", 25); | |
| stop(beforeFetch); | |
| await assert.rejects(queued, (error) => error.name === name); | |
| assert.strictEqual(starts, 0); | |
| assertSettled(beforeFetch); | |
| let headerSignal; | |
| const headers = loadManager((_url, { signal }) => { | |
| headerSignal = signal; | |
| return new Promise((_resolve, reject) => { | |
| // Some transports reject with AbortError even when given a custom reason. | |
| signal.addEventListener( | |
| "abort", | |
| () => reject(new DOMException("transport aborted", "AbortError")), | |
| { once: true } | |
| ); | |
| }); | |
| }); | |
| const headerRequests = [ | |
| headers.manager.request("headers", 25), | |
| headers.manager.request("headers", 25) | |
| ]; | |
| await Promise.resolve(); | |
| stop(headers); | |
| await Promise.all( | |
| headerRequests.map((request) => | |
| assert.rejects(request, (error) => error === headerSignal.reason && error.name === name) | |
| ) | |
| ); | |
| assertSettled(headers); | |
| let finish, lateSignal, cancelledWith; | |
| const delayedHeaders = loadManager((_url, { signal }) => { | |
| lateSignal = signal; | |
| return new Promise((resolve) => { | |
| finish = resolve; | |
| }); | |
| }); | |
| const delayedRequest = delayedHeaders.manager.request("late", 25); | |
| await Promise.resolve(); | |
| stop(delayedHeaders); | |
| finish( | |
| new Response( | |
| new ReadableStream({ | |
| cancel(reason) { | |
| cancelledWith = reason; | |
| } | |
| }) | |
| ) | |
| ); | |
| await assert.rejects( | |
| delayedRequest, | |
| (error) => error === lateSignal.reason && error.name === name | |
| ); | |
| assert.strictEqual(cancelledWith, lateSignal.reason); | |
| assertSettled(delayedHeaders); | |
| let bodySignal, | |
| cancels = 0; | |
| const bodies = loadManager((_url, { signal }) => { | |
| bodySignal = signal; | |
| return Promise.resolve( | |
| new Response( | |
| new ReadableStream({ | |
| cancel(reason) { | |
| assert.strictEqual(reason, signal.reason); | |
| assert.strictEqual(reason.name, name); | |
| cancels++; | |
| } | |
| }) | |
| ) | |
| ); | |
| }); | |
| const responses = await Promise.all([ | |
| bodies.manager.request("clones", 25), | |
| bodies.manager.request("clones", 25) | |
| ]); | |
| assert.notStrictEqual(responses[0], responses[1]); | |
| assert.strictEqual(bodies.manager.pendingCount, 0); | |
| assert.strictEqual(bodies.manager.activeCount, 1); | |
| const reads = responses.map((response) => | |
| assert.rejects(response.text(), (error) => error.name === name) | |
| ); | |
| stop(bodies); | |
| const firstReason = bodySignal.reason; | |
| bodies.manager.dispose(); | |
| await Promise.all(reads); | |
| assert.strictEqual(bodySignal.reason, firstReason); | |
| assert.strictEqual(bodySignal.reason.name, name); | |
| assert.strictEqual(cancels, 1); | |
| assert.strictEqual(bodies.manager.activeCount, 0); | |
| assert.strictEqual(bodies.timers.size, 0); | |
| } | |
| console.log("reader request manager contracts passed"); | |
| } | |
| main().catch((error) => { | |
| console.error(error); | |
| process.exitCode = 1; | |
| }); | |