diff --git a/CHANGELOG.md b/CHANGELOG.md index 4ddf711..d72b5e4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,14 @@ All notable changes to WASM-OJ are recorded here. Releases follow ## Unreleased +- Browser runner, compiler, compiler stage (rustc, Go, Java) and interactive side Workers that die + without an `error` event, for example when the browser terminates them, now reject their + operation promptly as a runner or compiler failure. A silently killed Worker used to leave the operation waiting for its wall-time + limit, which reported the student's program as `wall-time-limit`. Each module Worker holds a Web + Lock for its lifetime and its owner treats the lock's release as a crash; without Web Locks + nothing changes. Interactive pipes now wake every 100 ms while waiting, so a terminated side + Worker stops promptly in WebKit, which otherwise keeps it blocked in `Atomics.wait`. + ## 0.2.4 - 2026-10-09 - Fixed a host process crash (`Uncaught Error: write EPIPE`) when `ServerRunner` cancelled or diff --git a/docs/library-contract.md b/docs/library-contract.md index 214a5fa..3a18cf7 100644 --- a/docs/library-contract.md +++ b/docs/library-contract.md @@ -125,6 +125,14 @@ nested stages with bounded lifetime. Python and JavaScript package source files crossing a stage budget, cancellation, timeout, restart, cache clearing, disposal, or infrastructure failure establishes a complete Worker-generation boundary. +A module Worker that dies without an `error` event, for example because the browser terminated it, +is reported like a crash: its operation rejects with a `runner-failure` or `compiler-failure` +instead of waiting for its wall-time or build deadline, so a killed Worker is never reported as a +time limit. Each Worker holds a Web Lock for its lifetime, and its owner learns of the death when +that lock frees. Without Web Locks the previous behaviour remains. In WebKit a terminated Worker +keeps its lock while it runs Wasm, so a Worker killed in the middle of a computation is reported +only when it next returns to JavaScript; Chromium reports a busy Worker about 2 s after the kill. + Wasmer secondary Workers are host implementation details. They use the SDK's supported `workerUrl` protocol and do not grant guest thread-spawn capability. The host page must be cross-origin isolated. diff --git a/scripts/verify-browser-csp.mjs b/scripts/verify-browser-csp.mjs index 5f9a545..ee1eedc 100644 --- a/scripts/verify-browser-csp.mjs +++ b/scripts/verify-browser-csp.mjs @@ -27,6 +27,9 @@ const bootstrap = `import { createBrowserEngine, WASM_OJ_LIBCXX_PCH_HEADER } fro window.cspViolations = []; addEventListener('securitypolicyviolation', e => window.cspViolations.push({directive:e.effectiveDirective, blockedURI:e.blockedURI, source:e.sourceFile})); try { new Function('return 1')(); window.evalBlocked = false; } catch { window.evalBlocked = true; } window.header = WASM_OJ_LIBCXX_PCH_HEADER; +const NativeWorker = Worker; window.createdWorkers = []; +window.Worker = class extends NativeWorker { constructor(url, options) { super(url, options); window.createdWorkers.push({ name:options?.name, worker:this }); } }; +window.killWorker = name => NativeWorker.prototype.terminate.call(window.createdWorkers.findLast(entry => entry.name === name).worker); window.engine = await createBrowserEngine({ artifactCache:false, toolchains: ${JSON.stringify(sources)} }); window.ready = true;`; const server = createServer(async (req, res) => { @@ -219,6 +222,104 @@ try { await writeFile(path.join(output,"results.json"),JSON.stringify(record,null,2)+"\n"); console.log(JSON.stringify({ label:fixture.label, pass, elapsedMs:outcome.elapsedMs, error:outcome.error, summary })); } + if (selected.length === 0 || selected.includes("liveness")) { + record.liveness = []; + const sources = { + readOne:{ language:"c", source:'#include \nint main(void){int x;return scanf("%d",&x)==1?0:1;}' }, + yieldLoop:{ language:"c", source:'#include \nint main(void){for(;;)sched_yield();}' }, + computeLoop:{ language:"cpp", source:'int main(){volatile unsigned long long spin=0;for(;;)spin=spin+1;}' }, + guessContestant:guessC, + guessInteractor, + }; + const prepared = await page.evaluate(async sources => { + window.livenessBuilds = {}; + for (const [name, { language, source }] of Object.entries(sources)) { + const entry = language === "cpp" ? "main.cpp" : "main.c"; + const files = { [entry]:source }; + if (language === "cpp") files["src/bits/stdc++.h"] = window.header; + const built = await window.engine.compile({ language, target:"wasip1", optimization:"release", entry, files, projectId:`csp-liveness-${name}` }, { cache:false }); + if (!built.success || !built.artifact) throw new Error(`liveness build ${name} failed: ${built.stderr}`); + window.livenessBuilds[name] = built.artifact; + } + const summary = value => value.termination ?? (value.contestant ? `${value.contestant.termination}/${value.interactor.code}` : `compiled:${value.success}`); + window.settle = promise => promise.then(value => ({ ok:true, summary:summary(value), stderr:value.stderr, at:performance.now() }), error => ({ ok:false, error:String(error), at:performance.now() })); + const blocked = { contestant:{ resources:{ wallTimeLimitMs:20000 } }, interactor:{ resources:{ wallTimeLimitMs:20000 } } }; + window.livenessOperation = operation => { + const builds = window.livenessBuilds; + if (operation === "interact") return window.engine.interact(builds.readOne, builds.readOne, blocked); + if (operation === "run-yielding") return window.engine.run(builds.yieldLoop, { resources:{ instructionBudget:1e15, wallTimeLimitMs:20000 } }); + if (operation === "run-compute") return window.engine.run(builds.computeLoop, { resources:{ instructionBudget:1e15, wallTimeLimitMs:15000 } }); + if (operation === "compile-rust") return window.engine.compile({ language:"rust", target:"wasip1", optimization:"release", entry:"main.rs", files:{ "main.rs":'fn main(){println!("{}", 42);}' }, projectId:"csp-liveness-rust" }, { cache:false }); + return window.engine.compile({ language:"cpp", target:"wasip1", optimization:"release", entry:"main.cpp", files:{ "main.cpp":"#include \n#include \nint main(){std::regex r(\"a+\");std::cout< undefined, error => String(error)); + if (prepared) record.liveness.push({ label:"liveness-preparation", pass:false, error:prepared }); + const workerNamed = async name => { + const nameOf = worker => Promise.race([worker.evaluate(() => self.name).catch(() => ""), new Promise(resolve => setTimeout(resolve, 3000, ""))]); + for (let attempt = 0; attempt < 5; attempt++) { + for (const worker of page.workers().reverse()) if (await nameOf(worker) === name) return worker; + await page.waitForTimeout(500); + } + throw new Error(`The ${name} Worker is not visible to Playwright.`); + }; + const nestedKiller = async (parentName, childName, settleMs = 0) => { + const parent = await workerNamed(parentName); + await parent.evaluate(() => { + if (self.livenessWorkers) return; + const Native = self.Worker; self.livenessWorkers = []; self.nativeTerminate = Native.prototype.terminate; + self.Worker = class extends Native { constructor(url, options) { super(url, options); self.livenessWorkers.push({ name:options?.name, worker:this }); } }; + }); + return () => parent.evaluate(async ([name, settleMs]) => { + for (let attempt = 0; attempt < 600 && !self.livenessWorkers.some(entry => entry.name === name); attempt++) await new Promise(resolve => setTimeout(resolve, 100)); + await new Promise(resolve => setTimeout(resolve, settleMs)); + self.nativeTerminate.call(self.livenessWorkers.findLast(entry => entry.name === name).worker); + }, [childName, settleMs]); + }; + const pageKiller = name => () => page.evaluate(name => window.killWorker(name), name); + const browserName = process.env.WASM_OJ_BROWSER ?? "chromium"; + const crashed = (outcome, killMs, limitMs) => /stopped without reporting an error/.test(outcome.ok ? outcome.stderr ?? "" : outcome.error) && killMs < limitMs; + const livenessCases = [ + { label:"liveness-interactive-contestant", operation:"interact", killer:() => nestedKiller("wasm-oj-runner", "wasm-oj-interactive-contestant"), check:(outcome, killMs) => crashed(outcome, killMs, 3000) && outcome.error.includes("interactive contestant Worker") }, + { label:"liveness-interactive-interactor", operation:"interact", killer:() => nestedKiller("wasm-oj-runner", "wasm-oj-interactive-interactor"), check:(outcome, killMs) => crashed(outcome, killMs, 3000) && outcome.error.includes("interactive interactor Worker") }, + { label:"liveness-runner-interact", operation:"interact", killer:async () => pageKiller("wasm-oj-runner"), check:(outcome, killMs) => crashed(outcome, killMs, 3000) }, + { label:"liveness-runner-run-yielding", operation:"run-yielding", killer:async () => pageKiller("wasm-oj-runner"), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + { label:"liveness-runner-run-compute", operation:"run-compute", killer:async () => pageKiller("wasm-oj-runner"), check:(outcome, killMs) => browserName === "webkit" ? outcome.ok && outcome.summary === "wall-time-limit" : crashed(outcome, killMs, 4000) }, + { label:"liveness-compiler", operation:"compile", killAfterMs:300, killer:async () => pageKiller("wasm-oj-compiler"), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + { label:"liveness-compiler-stage", operation:"compile-rust", killAfterMs:0, killer:() => nestedKiller("wasm-oj-compiler", "wasm-oj-rustc-stage", 1000), check:(outcome, killMs) => crashed(outcome, killMs, 4000) }, + ]; + for (const fixture of prepared ? [] : livenessCases) { + console.log(`START ${fixture.label}`); + let outcome; let killMs; let error; + try { + const kill = await fixture.killer(); + await page.evaluate(operation => { window.livenessPending = window.settle(window.livenessOperation(operation)); }, fixture.operation); + await page.waitForTimeout(fixture.killAfterMs ?? 1500); + await kill(); + const killedAt = await page.evaluate(() => performance.now()); + outcome = await page.evaluate(() => window.livenessPending); + killMs = Math.round(outcome.at - killedAt); + } catch (caught) { error = String(caught); } + const recovery = await page.evaluate(() => window.settle(window.engine.run(window.livenessBuilds.readOne, { stdin:"7\n" }))); + const pass = !error && fixture.check(outcome, killMs) && recovery.summary === "exited"; + record.liveness.push({ label:fixture.label, pass, killMs, outcome, error, recovery }); + await writeFile(path.join(output,"results.json"),JSON.stringify(record,null,2)+"\n"); + console.log(JSON.stringify({ label:fixture.label, pass, killMs, error, outcome, recovery:recovery.summary ?? recovery.error })); + } + if (!prepared) { + console.log("START liveness-no-false-positive"); + const steady = await page.evaluate(async () => { + const outcomes = []; + const builds = window.livenessBuilds; + for (let index = 0; index < 20; index++) outcomes.push(await window.settle(window.engine.run(builds.readOne, { stdin:`${index}\n` }))); + for (let index = 0; index < 5; index++) outcomes.push(await window.settle(window.engine.interact(builds.guessContestant, builds.guessInteractor, { interactor:{ args:["/judge/input.txt"], files:{ "/judge/input.txt":`${index + 1} 25\n` } } }))); + outcomes.push(await window.settle(window.engine.compile({ language:"c", target:"wasip1", optimization:"release", entry:"main.c", files:{ "main.c":"int main(void){return 0;}" }, projectId:"csp-liveness-steady" }, { cache:false }))); + return outcomes.map(outcome => outcome.summary ?? outcome.error); + }); + const steadyPass = steady.length === 26 && steady.slice(0, 20).every(value => value === "exited") && steady.slice(20, 25).every(value => value === "exited/0") && steady[25] === "compiled:true"; + record.liveness.push({ label:"liveness-no-false-positive", pass:steadyPass, outcomes:steady }); + console.log(JSON.stringify({ label:"liveness-no-false-positive", pass:steadyPass, outcomes:steady })); + } else console.log(JSON.stringify({ label:"liveness-preparation", pass:false, error:prepared })); + } record.capabilities = []; for (const invoke of [false, true]) { const wasmPath = path.join(output, `capability-${invoke}.wasm`); @@ -251,5 +352,5 @@ finally { await browser?.close(); await new Promise(resolve => server.close(resolve)); } -if(record.results.some(result=>!result.pass)||record.capabilities?.some(result=>!result.pass)||record.executionTiming?.pass===false||record.interactive?.some(result=>!result.pass))process.exitCode=1; +if(record.results.some(result=>!result.pass)||record.capabilities?.some(result=>!result.pass)||record.executionTiming?.pass===false||record.interactive?.some(result=>!result.pass)||record.liveness?.some(result=>!result.pass))process.exitCode=1; console.log(`EVIDENCE ${path.join(output,"results.json")}`); diff --git a/src/runtime/client-lifecycle.test.ts b/src/runtime/client-lifecycle.test.ts index bcb901e..f1ee543 100644 --- a/src/runtime/client-lifecycle.test.ts +++ b/src/runtime/client-lifecycle.test.ts @@ -44,6 +44,7 @@ const TEST_TOOLCHAINS = Object.freeze([{ interface FakeWorker { readonly messages: unknown[]; readonly listeners: Map void>>; + readonly lostListeners: Set<(error: Error) => void>; terminated: boolean; addEventListener(type: string, listener: (event: unknown) => void): void; postMessage(message: unknown): void; @@ -65,6 +66,7 @@ vi.mock("./module-worker", () => ({ const worker: FakeWorker = { messages: [], listeners: new Map(), + lostListeners: new Set(), terminated: false, addEventListener: () => undefined, postMessage: () => undefined, @@ -87,6 +89,9 @@ vi.mock("./module-worker", () => ({ collection.push(worker); return worker; }, + onModuleWorkerLost(worker: FakeWorker, listener: (error: Error) => void): void { + worker.lostListeners.add(listener); + }, })); beforeEach(() => { @@ -153,6 +158,46 @@ describe("browser client lifecycle", () => { compiler.dispose(); }); + it("rejects a running execution at once when the runner Worker is lost", async () => { + vi.useFakeTimers(); + const runner = new BrowserRunner({ toolchains: TEST_TOOLCHAINS, additionalCostBaselines: { [TEST_COST_PROFILE]: 0 } }); + try { + const worker = workerState.runners[0]!; + respondToInitialization(worker); + await runner.ready(); + const pending = runner.run(wasmArtifact(), runConfig()); + await Promise.resolve(); + const { requestId } = requestOfType(worker, "run"); + respond(worker, { type: "progress", requestId, progress: { phase: "running", label: "guest" } }); + await vi.advanceTimersByTimeAsync(10); + + lose(worker, "The wasm-oj-runner Worker stopped without reporting an error."); + + await expect(pending).rejects.toThrow("The wasm-oj-runner Worker stopped without reporting an error."); + expect(worker.terminated).toBe(true); + expect(workerState.runners).toHaveLength(2); + } finally { + runner.dispose(); + vi.useRealTimers(); + } + }); + + it("rejects a build at once when the compiler Worker is lost", async () => { + const compiler = new BrowserCompiler({ toolchains: TEST_TOOLCHAINS }); + const worker = workerState.compilers[0]!; + respondToInitialization(worker); + await compiler.ready(); + const pending = compiler.build(javascriptProject(), "cache-key"); + await vi.waitFor(() => expect(requestsOfType(worker, "build")).toHaveLength(1)); + + lose(worker, "The wasm-oj-compiler Worker stopped without reporting an error."); + + await expect(pending).rejects.toThrow("The wasm-oj-compiler Worker stopped without reporting an error."); + expect(worker.terminated).toBe(true); + expect(workerState.compilers).toHaveLength(2); + compiler.dispose(); + }); + it("rejects malformed direct compiler inputs before crossing the Worker boundary", async () => { const compiler = new BrowserCompiler({ toolchains: TEST_TOOLCHAINS }); const worker = workerState.compilers[0]!; @@ -711,6 +756,10 @@ function dispatch(worker: FakeWorker, type: string, event: unknown): void { for (const listener of worker.listeners.get(type) ?? []) listener(event); } +function lose(worker: FakeWorker, message: string): void { + for (const listener of worker.lostListeners) listener(new Error(message)); +} + function respond(worker: FakeWorker, data: unknown): void { for (const listener of worker.listeners.get("message") ?? []) listener({ data }); } diff --git a/src/runtime/compiler-client.ts b/src/runtime/compiler-client.ts index e7c84e7..5a6d1a7 100644 --- a/src/runtime/compiler-client.ts +++ b/src/runtime/compiler-client.ts @@ -29,7 +29,7 @@ import { maximumOutputReadyRustStages, } from "../compiler/browser-rust-policy"; import CompilerWorkerUrl from "./compiler.worker?worker&url"; -import { createModuleWorker } from "./module-worker"; +import { createModuleWorker, onModuleWorkerLost } from "./module-worker"; import { prefetchBrowserToolchain } from "./toolchain-prefetch"; import { clearClangBuildGraphCache } from "../compiler/indexeddb-build-graph-cache"; @@ -273,13 +273,14 @@ export class BrowserCompiler implements Compiler { worker.addEventListener("message", (event: MessageEvent) => { if (!this.disposed && !this.workerDormant && this.worker === worker) this.handleMessage(event.data); }); - worker.addEventListener("error", (event) => { - const error = new Error(event.message || "The compiler worker crashed."); + const crashed = (error: Error) => { if (this.disposed || this.workerDormant || this.worker !== worker) return; const canRecover = this.workerInitialized; this.stopWorker(error); if (canRecover) this.installWorker(); - }); + }; + worker.addEventListener("error", (event) => crashed(new Error(event.message || "The compiler worker crashed."))); + onModuleWorkerLost(worker, crashed); return worker; } diff --git a/src/runtime/interactive-pipe.ts b/src/runtime/interactive-pipe.ts index 1512990..fe16c34 100644 --- a/src/runtime/interactive-pipe.ts +++ b/src/runtime/interactive-pipe.ts @@ -7,6 +7,12 @@ const WRITER_CLOSED = 2; const READER_CLOSED = 3; const SEQUENCE = 4; const HEADER_BYTES = 32; +/** + * WebKit does not stop a Worker that `terminate()` catches inside `Atomics.wait` until the wait + * returns, and Chromium waits up to 2 s. Waking periodically lets a terminated side Worker stop, + * and release its liveness lock, promptly. + */ +const WAIT_SLICE_MS = 100; /** * The smallest ring that holds a writer's whole output budget. Only budgeted stdout bytes enter @@ -49,7 +55,7 @@ class InteractivePipeEnd { } protected sleep(sequence: number): void { - Atomics.wait(this.header, SEQUENCE, sequence); + Atomics.wait(this.header, SEQUENCE, sequence, WAIT_SLICE_MS); } } diff --git a/src/runtime/isolated-stage.ts b/src/runtime/isolated-stage.ts index b207bc5..6d4d729 100644 --- a/src/runtime/isolated-stage.ts +++ b/src/runtime/isolated-stage.ts @@ -1,3 +1,5 @@ +import { onModuleWorkerLost } from "./module-worker"; + export type IsolatedStageResponse = | { type: "result"; result: Result } | { type: "shutdown-complete" } @@ -51,6 +53,7 @@ export class PersistentIsolatedStage { this.worker.addEventListener("message", this.onMessage); this.worker.addEventListener("error", this.onError); this.worker.addEventListener("messageerror", this.onMessageError); + onModuleWorkerLost(this.worker, (error) => this.fail(error)); } run(request: Request): Promise { @@ -247,6 +250,7 @@ export function runIsolatedStage( worker.addEventListener("message", onMessage); worker.addEventListener("error", onError); worker.addEventListener("messageerror", onMessageError); + onModuleWorkerLost(worker, (error) => finish(() => reject(error))); try { worker.postMessage(request); } catch (error) { diff --git a/src/runtime/module-worker.test.ts b/src/runtime/module-worker.test.ts index 36be42a..4481b1a 100644 --- a/src/runtime/module-worker.test.ts +++ b/src/runtime/module-worker.test.ts @@ -3,6 +3,7 @@ import { createModuleWorker, createModuleWorkerBootstrap, moduleWorkerBaseUrl, + onModuleWorkerLost, } from "./module-worker"; interface WorkerConstruction { @@ -11,15 +12,42 @@ interface WorkerConstruction { } const constructions: WorkerConstruction[] = []; +const workers: FakeWorker[] = []; + +class FakeWorker extends EventTarget { + terminated = false; -class FakeWorker { constructor(url: string | URL, options?: WorkerOptions) { + super(); constructions.push({ url, options }); + workers.push(this); + } + + terminate(): void { + this.terminated = true; } } +interface LockRequest { + name: string; + signal: AbortSignal; + grant(): void; +} + +const lockRequests: LockRequest[] = []; +const fakeLocks = { + request(name: string, options: { signal: AbortSignal }, callback: () => void): Promise { + return new Promise((resolve) => { + lockRequests.push({ name, signal: options.signal, grant: () => resolve(callback()) }); + }); + }, +}; + beforeEach(() => { constructions.length = 0; + workers.length = 0; + lockRequests.length = 0; + vi.stubGlobal("navigator", { locks: fakeLocks }); vi.stubGlobal("location", { href: "https://wasm-oj.example/judge", origin: "https://wasm-oj.example", @@ -50,6 +78,7 @@ describe("module Worker bootstrap", () => { "const queueMessage = (event) => { event.stopImmediatePropagation(); pendingMessages.push(event.data); };", 'globalThis.addEventListener("message", queueMessage);', 'Object.defineProperty(globalThis, "__wasmOjModuleWorkerBaseUrl", { value: "https://wasm-oj.example/judge" });', + 'try { const lock = "wasm-oj-worker-" + crypto.randomUUID(); navigator.locks.request(lock, () => { postMessage({ __wasmOjWorkerLiveness: lock }); return new Promise(() => {}); }).catch(() => {}); } catch {}', 'try { await import("https://wasm-oj.example/assets/compiler.worker.js"); } finally { globalThis.removeEventListener("message", queueMessage); }', 'for (const data of pendingMessages) globalThis.dispatchEvent(new MessageEvent("message", { data }));', "", @@ -90,6 +119,7 @@ describe("module Worker bootstrap", () => { expect(await (source as Blob).text()).toContain( 'await import("https://wasm-oj.example/assets/wasmer-thread.worker.js")', ); + expect(await (source as Blob).text()).not.toContain("__wasmOjWorkerLiveness"); bootstrap.revoke(); bootstrap.revoke(); @@ -97,6 +127,52 @@ describe("module Worker bootstrap", () => { expect(URL.revokeObjectURL).toHaveBeenCalledWith("blob:https://wasm-oj.example/bootstrap"); }); + it("reports a Worker whose liveness lock frees before its owner terminates it", async () => { + const worker = createModuleWorker("/assets/runner.worker.js", { name: "wasm-oj-runner" }); + const messages: unknown[] = []; + const lost: string[] = []; + worker.addEventListener("message", (event) => messages.push((event as MessageEvent).data)); + onModuleWorkerLost(worker, (error) => lost.push(error.message)); + + worker.dispatchEvent(new MessageEvent("message", { data: { __wasmOjWorkerLiveness: "wasm-oj-worker-1" } })); + worker.dispatchEvent(new MessageEvent("message", { data: { type: "ready" } })); + expect(messages).toEqual([{ type: "ready" }]); + expect(lockRequests.map(({ name }) => name)).toEqual(["wasm-oj-worker-1"]); + expect(lost).toEqual([]); + + lockRequests[0]!.grant(); + expect(lost).toEqual(["The wasm-oj-runner Worker stopped without reporting an error."]); + + const late = await new Promise((resolve) => onModuleWorkerLost(worker, (error) => resolve(error.message))); + expect(late).toBe("The wasm-oj-runner Worker stopped without reporting an error."); + }); + + it("stops watching a Worker once its owner terminates it", () => { + const worker = createModuleWorker("/assets/runner.worker.js", { name: "wasm-oj-runner" }); + const lost: Error[] = []; + onModuleWorkerLost(worker, (error) => lost.push(error)); + worker.dispatchEvent(new MessageEvent("message", { data: { __wasmOjWorkerLiveness: "wasm-oj-worker-2" } })); + + worker.terminate(); + lockRequests[0]!.grant(); + + expect(workers[0]?.terminated).toBe(true); + expect(lockRequests[0]?.signal.aborted).toBe(true); + expect(lost).toEqual([]); + }); + + it("keeps the previous behaviour when Web Locks are unavailable", () => { + vi.stubGlobal("navigator", {}); + const worker = createModuleWorker("/assets/runner.worker.js"); + const messages: unknown[] = []; + worker.addEventListener("message", (event) => messages.push((event as MessageEvent).data)); + + worker.dispatchEvent(new MessageEvent("message", { data: { __wasmOjWorkerLiveness: "wasm-oj-worker-3" } })); + + expect(messages).toEqual([]); + expect(lockRequests).toEqual([]); + }); + it("uses the injected browser base inside a blob Worker", () => { vi.stubGlobal("location", { href: "blob:https://wasm-oj.example/bootstrap", diff --git a/src/runtime/module-worker.ts b/src/runtime/module-worker.ts index 235fd4d..08bece9 100644 --- a/src/runtime/module-worker.ts +++ b/src/runtime/module-worker.ts @@ -24,19 +24,27 @@ export function createModuleWorker( scriptUrl: string | URL, options: ModuleWorkerOptions = {}, ): Worker { - const bootstrap = createModuleWorkerBootstrap(scriptUrl); + const bootstrap = createModuleWorkerBootstrap(scriptUrl, { liveness: true }); + let worker: Worker; try { - return new Worker(bootstrap.url, { ...options, type: "module" }); + worker = new Worker(bootstrap.url, { ...options, type: "module" }); } finally { bootstrap.revoke(); } + superviseLiveness(worker, options.name); + return worker; } /** * Creates a reusable blob bootstrap for APIs such as the Wasmer SDK that own * Worker construction and may defer it until a command is instantiated. + * `liveness` makes the Worker hold and report its liveness lock; only + * `createModuleWorker`, which consumes that report, enables it. */ -export function createModuleWorkerBootstrap(scriptUrl: string | URL): ModuleWorkerBootstrap { +export function createModuleWorkerBootstrap( + scriptUrl: string | URL, + { liveness = false }: { liveness?: boolean } = {}, +): ModuleWorkerBootstrap { const baseUrl = moduleWorkerBaseUrl(); const absoluteScriptUrl = resolveModuleWorkerUrl(scriptUrl); const bootstrap = new Blob( @@ -45,6 +53,7 @@ export function createModuleWorkerBootstrap(scriptUrl: string | URL): ModuleWork "const queueMessage = (event) => { event.stopImmediatePropagation(); pendingMessages.push(event.data); };\n", 'globalThis.addEventListener("message", queueMessage);\n', `Object.defineProperty(globalThis, "__wasmOjModuleWorkerBaseUrl", { value: ${JSON.stringify(baseUrl.href)} });\n`, + ...(liveness ? [LIVENESS_BOOTSTRAP] : []), `try { await import(${JSON.stringify(absoluteScriptUrl)}); } finally { globalThis.removeEventListener("message", queueMessage); }\n`, 'for (const data of pendingMessages) globalThis.dispatchEvent(new MessageEvent("message", { data }));\n', ], @@ -62,6 +71,76 @@ export function createModuleWorkerBootstrap(scriptUrl: string | URL): ModuleWork }; } +const LIVENESS_MESSAGE_KEY = "__wasmOjWorkerLiveness"; + +const LIVENESS_BOOTSTRAP = "try { const lock = \"wasm-oj-worker-\" + crypto.randomUUID(); " + + `navigator.locks.request(lock, () => { postMessage({ ${LIVENESS_MESSAGE_KEY}: lock }); return new Promise(() => {}); }).catch(() => {}); } catch {}\n`; + +interface WorkerLiveness { + lost?: Error; + readonly listeners: Set<(error: Error) => void>; +} + +const liveness = new WeakMap(); + +/** + * Calls `listener` once if `worker`, created by `createModuleWorker`, stops + * without an `error` event before its owner calls `terminate()`, for example + * because the browser terminated it. Owners treat this like a crash. Without + * Web Locks the listener is never called and the previous behaviour remains. + */ +export function onModuleWorkerLost(worker: Worker, listener: (error: Error) => void): void { + const state = liveness.get(worker); + if (!state) return; + if (state.lost) { + const lost = state.lost; + queueMicrotask(() => listener(lost)); + return; + } + state.listeners.add(listener); +} + +/** + * The bootstrap holds a uniquely named Web Lock until its context is destroyed + * and reports the name. Requesting the same lock here is granted only once the + * Worker is gone; unless its owner terminated it first, that is reported to + * `onModuleWorkerLost` listeners. They are called directly rather than through + * an `error` event because WebKit drops events dispatched on a Worker after + * `terminate()`. + */ +function superviseLiveness(worker: Worker, name: string | undefined): void { + const state: WorkerLiveness = { listeners: new Set() }; + liveness.set(worker, state); + const owner = new AbortController(); + const terminate = worker.terminate.bind(worker); + worker.terminate = () => { + owner.abort(); + state.listeners.clear(); + terminate(); + }; + worker.addEventListener("message", (event: MessageEvent) => { + const lock = livenessLock(event.data); + if (lock === undefined) return; + event.stopImmediatePropagation(); + const locks = globalThis.navigator?.locks; + if (!locks) return; + void locks.request(lock, { signal: owner.signal }, () => { + if (owner.signal.aborted) return; + const lost = new Error(`The ${name ?? "module"} Worker stopped without reporting an error.`); + state.lost = lost; + const listeners = [...state.listeners]; + state.listeners.clear(); + for (const listener of listeners) listener(lost); + }).catch(() => undefined); + }); +} + +function livenessLock(data: unknown): string | undefined { + if (typeof data !== "object" || data === null) return undefined; + const lock = (data as Record)[LIVENESS_MESSAGE_KEY]; + return typeof lock === "string" ? lock : undefined; +} + export function moduleWorkerBaseUrl(): URL { const locationHref = globalThis.location?.href; if (locationHref) { diff --git a/src/runtime/runner-client.ts b/src/runtime/runner-client.ts index 5130218..a2bf7a8 100644 --- a/src/runtime/runner-client.ts +++ b/src/runtime/runner-client.ts @@ -20,7 +20,7 @@ import type { CostBaselineRegistry, } from "@wasm-oj/core"; import RunnerWorkerUrl from "./runner.worker?worker&url"; -import { createModuleWorker } from "./module-worker"; +import { createModuleWorker, onModuleWorkerLost } from "./module-worker"; import { validateBrowserRuntimeDriverPlugins } from "./browser-runtime-plugin"; import { runtimePreparationTimeoutMs } from "../runner/preparation-timeout-policy"; @@ -278,13 +278,14 @@ export class BrowserRunner implements Runner { worker.addEventListener("message", (event: MessageEvent) => { if (!this.disposed && this.worker === worker) this.handleMessage(event.data); }); - worker.addEventListener("error", (event) => { - const error = new Error(event.message || "The runner worker crashed."); + const crashed = (error: Error) => { if (this.disposed || this.worker !== worker) return; const canRecover = this.workerInitialized; this.stopWorker(error); if (canRecover) this.installWorker(); - }); + }; + worker.addEventListener("error", (event) => crashed(new Error(event.message || "The runner worker crashed."))); + onModuleWorkerLost(worker, crashed); return worker; } diff --git a/src/runtime/runner.worker.ts b/src/runtime/runner.worker.ts index 2c32fde..ddb6a53 100644 --- a/src/runtime/runner.worker.ts +++ b/src/runtime/runner.worker.ts @@ -60,6 +60,7 @@ import { createModuleWorkerBootstrap, type ModuleWorkerBootstrap, moduleWorkerBaseUrl, + onModuleWorkerLost, } from "./module-worker"; import { createInteractivePipe, interactivePipeCapacity } from "./interactive-pipe"; import type { InteractiveSideMessage, InteractiveSideStart } from "./interactive-side.worker"; @@ -515,10 +516,14 @@ function startInteractiveSide(role: InteractiveRole, start: InteractiveSideStart } } }); + const crashed = (detail: string) => { + reject(Object.assign(new Error(`The interactive ${role} Worker crashed${detail ? `: ${detail}` : "."}`), { code: "RUNTIME_ERROR" })); + }; worker.addEventListener("error", (event) => { event.preventDefault(); - reject(Object.assign(new Error(`The interactive ${role} Worker crashed${event.message ? `: ${event.message}` : "."}`), { code: "RUNTIME_ERROR" })); + crashed(event.message); }); + onModuleWorkerLost(worker, (error) => crashed(error.message)); }); try { worker.postMessage(start);